Bind prepared references to extraction identity

This commit is contained in:
2026-08-29 15:24:31 +00:00
parent 51e0e8c5d0
commit 495f7bcde4
5 changed files with 376 additions and 8 deletions

View File

@@ -80,6 +80,10 @@ func (extractStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
if err != nil {
return nil, fmt.Errorf("extract: resolve final-trimmed transcript identity: %w", err)
}
references, err := resolveExtractReferences(paths, m, notariusConfig, sessionID)
if err != nil {
return nil, fmt.Errorf("extract: resolve Notarius references: %w", err)
}
timeout, err := time.ParseDuration(strings.TrimSpace(notariusConfig.Timeout))
if err != nil || timeout <= 0 {
@@ -125,14 +129,14 @@ func (extractStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
}
}
fingerprint, err := extractionFingerprint(resolvedBinary, configPath, notariusConfig, timeout, workingDirectory, input)
fingerprint, err := extractionFingerprint(resolvedBinary, configPath, notariusConfig, timeout, workingDirectory, input, references.Identities)
if err != nil {
return nil, fmt.Errorf("extract: build configuration fingerprint: %w", err)
}
request := notarius.RunRequest{
Binary: resolvedBinary, ConfigPath: configPath, PipelineID: notariusConfig.PipelineID,
InputPath: inputPath, OutputRoot: outputRoot, WorkingDirectory: workingDirectory,
ReceiptPath: receiptPath, LogPath: logPath, Timeout: timeout,
ReceiptPath: receiptPath, LogPath: logPath, Timeout: timeout, References: references.Bindings,
}
adapterResult, err := env.Notarius.Run(ctx, request)
if err != nil {
@@ -249,6 +253,8 @@ func (extractStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
"narratio_run_id": runID,
"configuration_fingerprint": fingerprint,
"direct_input": input.Metadata(),
"reference_count": len(references.Identities),
"references": extractReferenceMetadata(references.Identities),
"receipt": map[string]any{
"run_id": adapterResult.Receipt.RunID, "pipeline_id": adapterResult.Receipt.PipelineID,
"normalized_output_count": adapterResult.Receipt.NormalizedOutputCount,
@@ -379,6 +385,7 @@ type fingerprintDocument struct {
Timeout string `json:"timeout"`
WorkingDirectory string `json:"working_directory"`
Input artifacts.ExtractionInputIdentity `json:"input"`
References []extractReferenceIdentity `json:"references"`
Outputs []fingerprintOutput `json:"outputs"`
}
@@ -388,6 +395,7 @@ func extractionFingerprint(
timeout time.Duration,
workingDirectory string,
input artifacts.ExtractionInputIdentity,
references []extractReferenceIdentity,
) (string, error) {
keys := make([]string, 0, len(cfg.Outputs))
for key := range cfg.Outputs {
@@ -402,9 +410,20 @@ func extractionFingerprint(
SchemaID: output.SchemaID, SchemaVersion: output.SchemaVersion, ModuleKey: output.ModuleKey,
})
}
sortedReferences := append([]extractReferenceIdentity(nil), references...)
sort.Slice(sortedReferences, func(i, j int) bool {
if sortedReferences[i].Selector != sortedReferences[j].Selector {
return sortedReferences[i].Selector < sortedReferences[j].Selector
}
if sortedReferences[i].SourceID != sortedReferences[j].SourceID {
return sortedReferences[i].SourceID < sortedReferences[j].SourceID
}
return sortedReferences[i].Path < sortedReferences[j].Path
})
payload, err := json.Marshal(fingerprintDocument{
Binary: binary, ConfigPath: configPath, PipelineID: cfg.PipelineID,
Timeout: timeout.String(), WorkingDirectory: workingDirectory, Input: input, Outputs: outputs,
Timeout: timeout.String(), WorkingDirectory: workingDirectory, Input: input,
References: sortedReferences, Outputs: outputs,
})
if err != nil {
return "", err