Write terminal debug reports for failed runs

This commit is contained in:
2026-07-18 14:31:46 +00:00
parent 2111e01142
commit e4ec521bed
6 changed files with 370 additions and 48 deletions

View File

@@ -43,6 +43,7 @@ type Options struct {
UserCacheDir func() (string, error)
ChunkPlanStoreFactory pipeline.ChunkPlanStoreFactory
DebugRecorderFactory func(string) (pipeline.DebugRecorder, error)
DebugTerminalFactory func(*debugbundle.SummaryWriter) DebugTerminalWriter
}
type LLMClientFactory func(ctx context.Context, cfg config.Config, profileID string) (contracts.StructuredLLMClient, []artifacts.LLMProfileManifest, error)
@@ -104,6 +105,9 @@ func normalizeOptions(opts Options) (Options, error) {
if opts.DebugRecorderFactory == nil {
opts.DebugRecorderFactory = frameworkdebug.NewFilesystemRecorder
}
if opts.DebugTerminalFactory == nil {
opts.DebugTerminalFactory = func(writer *debugbundle.SummaryWriter) DebugTerminalWriter { return writer }
}
if isEmptyCatalog(opts.Catalog) && isEmptyRegistries(opts.Registries) {
components, err := newProductionComponents()
if err != nil {
@@ -226,7 +230,10 @@ func runPipelineCommand(args []string, stdout, stderr io.Writer, opts Options) i
fmt.Fprintf(stderr, "notarius: invalid generated run ID: %v\n", err)
return 1
}
runOutputDir := filepath.Join(cfg.Output.Directory, runID)
commandState := newPipelineCommandState(runID, pipelineID, runOutputDir)
var summary *debugbundle.SummaryWriter
var terminalWriter DebugTerminalWriter
debugPath := ""
debugRecorder := pipeline.NoopDebugRecorder()
if *debug {
@@ -236,9 +243,14 @@ func runPipelineCommand(args []string, stdout, stderr io.Writer, opts Options) i
return 1
}
debugPath, summary = bundle.Path(), bundle.Summary()
commandState.setDebugPath(debugPath)
terminalWriter = opts.DebugTerminalFactory(summary)
if terminalWriter == nil {
terminalWriter = summary
}
debugRecorder, err = opts.DebugRecorderFactory(bundle.TraceRoot())
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("create debug recorder: %w", err), true)
return failPipelineCommand(stderr, commandState, terminalWriter, fmt.Errorf("create debug recorder: %w", err))
}
debugRecorder = pipeline.SynchronizedDebugRecorder(debugRecorder)
}
@@ -255,16 +267,16 @@ func runPipelineCommand(args []string, stdout, stderr io.Writer, opts Options) i
StartedAt: startedAt,
}
if err := writeSummary(summary, func() error { return summary.WriteInvocation(invocation) }); err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("write debug invocation metadata: %w", err), false)
return failPipelineCommand(stderr, commandState, terminalWriter, fmt.Errorf("write debug invocation metadata: %w", err))
}
catalog, err := effectiveCatalog(opts)
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, err, true)
return failPipelineCommand(stderr, commandState, terminalWriter, err)
}
referenceOverrides, referenceUnbinds, err := resolveCLIReferenceRequests(cfg, pipelineID, only, catalog, referenceRequests, referenceUnbindRequests)
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, err, true)
return failPipelineCommand(stderr, commandState, terminalWriter, err)
}
effective, err := cfg.Resolve(config.ResolveInput{
PipelineID: pipelineID,
@@ -275,43 +287,43 @@ func runPipelineCommand(args []string, stdout, stderr io.Writer, opts Options) i
ReferenceUnbinds: referenceUnbinds,
})
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, err, true)
return failPipelineCommand(stderr, commandState, terminalWriter, err)
}
profileIDs := effectiveLLMProfileIDs(effective.ResolvedPipeline)
if err := validateExplicitScriptoriumProfiles(context.Background(), effective.Config, profileIDs); err != nil {
return failPipelineCommand(stderr, summary, debugPath, err, true)
return failPipelineCommand(stderr, commandState, terminalWriter, err)
}
workingDir, err := os.Getwd()
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("resolve working directory: %w", err), true)
return failPipelineCommand(stderr, commandState, terminalWriter, fmt.Errorf("resolve working directory: %w", err))
}
materialized, referenceWarnings, err := pipeline.MaterializeReferences(effective.ResolvedPipeline, catalog, pipeline.ReferenceMaterializationOptions{
ConfigPath: loadedConfigPath,
WorkingDir: workingDir,
})
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, err, true)
return failPipelineCommand(stderr, commandState, terminalWriter, err)
}
effective.ResolvedPipeline = materialized
invocation.PipelineDigest = effective.ResolvedPipeline.Digest
if err := writeSummary(summary, func() error { return summary.WriteInvocation(invocation) }); err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("write debug invocation metadata: %w", err), false)
return failPipelineCommand(stderr, commandState, terminalWriter, fmt.Errorf("write debug invocation metadata: %w", err))
}
if err := writeSummary(summary, func() error { return summary.WriteRedactedEffectiveConfig(effective) }); err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("write debug effective config: %w", err), false)
return failPipelineCommand(stderr, commandState, terminalWriter, fmt.Errorf("write debug effective config: %w", err))
}
if err := writeSummary(summary, func() error { return summary.WriteResolvedPipeline(effective) }); err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("write debug resolved pipeline: %w", err), false)
return failPipelineCommand(stderr, commandState, terminalWriter, fmt.Errorf("write debug resolved pipeline: %w", err))
}
if err := writeSummary(summary, func() error {
return summary.WriteResolvedReferences(pipeline.ReferenceProvenance(effective.ResolvedPipeline))
}); err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("write debug resolved references: %w", err), false)
return failPipelineCommand(stderr, commandState, terminalWriter, fmt.Errorf("write debug resolved references: %w", err))
}
registries, err := effectiveRegistries(opts)
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, err, true)
return failPipelineCommand(stderr, commandState, terminalWriter, err)
}
ctx := context.Background()
@@ -321,24 +333,24 @@ func runPipelineCommand(args []string, stdout, stderr io.Writer, opts Options) i
}
llmClient, llmProfiles, err := opts.LLMClientFactory(ctx, effective.Config, factoryProfileID)
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("create LLM client for profile %q: %w", factoryProfileID, err), true)
return failPipelineCommand(stderr, commandState, terminalWriter, fmt.Errorf("create LLM client for profile %q: %w", factoryProfileID, err))
}
llmClient = pipeline.WithDebugLLMRecording(llmClient, debugRecorder)
prepared, err := pipeline.Prepare(effective.ResolvedPipeline, registries, pipeline.ModuleDependencies{LLM: llmClient})
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("prepare pipeline %q: %w", pipelineID, err), true)
return failPipelineCommand(stderr, commandState, terminalWriter, fmt.Errorf("prepare pipeline %q: %w", pipelineID, err))
}
rawInput, err := os.ReadFile(strings.TrimSpace(*inputPath))
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("read input %q: %w", strings.TrimSpace(*inputPath), err), true)
return failPipelineCommand(stderr, commandState, terminalWriter, fmt.Errorf("read input %q: %w", strings.TrimSpace(*inputPath), err))
}
chunkPlans, err := chunkPlanStoreForRun(effective.Config.Cache.ChunkPlans, opts)
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, err, true)
return failPipelineCommand(stderr, commandState, terminalWriter, err)
}
checkpointRecorder, checkpointLoader, err := checkpointHandlersForRun(effective.Config.Cache.Checkpoints, opts, effective.ResolvedPipeline, rawInput, only, llmProfiles, strings.TrimSpace(*llmProfile), strings.TrimSpace(sessionID.value), *resume)
if err != nil {
return failPipelineCommand(stderr, summary, debugPath, err, true)
return failPipelineCommand(stderr, commandState, terminalWriter, err)
}
output, err := pipeline.New().Run(ctx, pipeline.RunInput{
@@ -358,26 +370,25 @@ func runPipelineCommand(args []string, stdout, stderr io.Writer, opts Options) i
Debug: debugRecorder,
ExtractWorkers: cfg.Concurrency.StageWorkers["extract"],
})
commandState.observeOutput(output)
if err != nil {
primaryErr := fmt.Errorf("run pipeline %q: %w", pipelineID, err)
if output.Manifest.PipelineID != "" {
if summaryErr := writePartialSummary(summary, output); summaryErr != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("run pipeline %q: %w; write debug summary: %v", pipelineID, err, summaryErr), false)
return failPipelineCommand(stderr, commandState, terminalWriter, primaryErr, fmt.Errorf("write debug summary: %w", summaryErr))
}
}
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("run pipeline %q: %w", pipelineID, err), true)
return failPipelineCommand(stderr, commandState, terminalWriter, primaryErr)
}
runOutputDir := filepath.Join(effective.Config.Output.Directory, runID)
if err := writePartialSummary(summary, output); err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("write debug summary: %w", err), false)
return failPipelineCommand(stderr, commandState, terminalWriter, fmt.Errorf("write debug summary: %w", err))
}
if err := writeOutputFiles(runOutputDir, output.OutputFiles); err != nil {
return failPipelineCommand(stderr, summary, debugPath, err, true)
return failPipelineCommand(stderr, commandState, terminalWriter, err)
}
if err := writeSummary(summary, func() error {
return summary.WriteRunReport(debugbundle.RunReport{RunID: runID, PipelineID: effective.PipelineID, OutputPath: runOutputDir, DebugPath: debugPath, Succeeded: true, OutputCount: len(output.NormalizeOutputs), RejectedCount: len(output.Rejected), WarningCount: len(output.Warnings), ValidationStatus: output.Manifest.ValidationStatus})
}); err != nil {
return failPipelineCommand(stderr, summary, debugPath, fmt.Errorf("write debug run report: %w", err), false)
if primaryErr, persistenceErr := commandState.terminalize(terminalWriter, nil); primaryErr != nil {
return writePipelineCommandFailure(stderr, commandState, primaryErr, persistenceErr)
}
fmt.Fprintf(stdout, "pipeline %q complete: outputs=%d rejected=%d output=%s\n", effective.PipelineID, len(output.NormalizeOutputs), len(output.Rejected), runOutputDir)
@@ -390,19 +401,6 @@ func runPipelineCommand(args []string, stdout, stderr io.Writer, opts Options) i
return 0
}
func failPipelineCommand(stderr io.Writer, summary *debugbundle.SummaryWriter, debugPath string, err error, recordError bool) int {
fmt.Fprintf(stderr, "notarius: %v\n", err)
if recordError && summary != nil {
if summaryErr := summary.WriteError(err.Error()); summaryErr != nil {
fmt.Fprintf(stderr, "notarius: write debug error log: %v\n", summaryErr)
}
}
if debugPath != "" {
fmt.Fprintf(stderr, "notarius: debug=%s\n", debugPath)
}
return 1
}
func writeSummary(summary *debugbundle.SummaryWriter, write func() error) error {
if summary == nil {
return nil