From 6e3de551732b1cb037882f2cd8c47d77911a359a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Aur=C3=A9lien=20Sibiril?= <81782+aureliensibiril@users.noreply.github.com> Date: Fri, 24 Apr 2026 20:01:06 +0200 Subject: [PATCH] Truncate persisted agent run errors on rune boundary MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Raw tool or LLM errors can embed URLs with credentials, PII, or partial DB records. We now log the full runErr for operator context and persist a sanitized summary into agent_runs.error_message, truncated at 512 bytes with a trailing ellipsis. The cut rewinds to the nearest utf8.RuneStart so we never store a split rune. Also wraps the pending-run loader error in Claim for parity with the other coredata failures. Signed-off-by: Aurélien Sibiril <81782+aureliensibiril@users.noreply.github.com> --- pkg/probo/agent_run_handler.go | 66 +++++++++++++++++++++++++++------- 1 file changed, 53 insertions(+), 13 deletions(-) diff --git a/pkg/probo/agent_run_handler.go b/pkg/probo/agent_run_handler.go index bcaee782a..61c5d8e66 100644 --- a/pkg/probo/agent_run_handler.go +++ b/pkg/probo/agent_run_handler.go @@ -21,6 +21,7 @@ import ( "fmt" "sync" "time" + "unicode/utf8" "go.gearno.de/kit/log" "go.gearno.de/kit/pg" @@ -61,7 +62,7 @@ func (h *agentRunHandler) Claim(ctx context.Context) (coredata.AgentRun, error) ctx, func(ctx context.Context, tx pg.Tx) error { if err := run.LoadNextPendingForUpdateSkipLocked(ctx, tx); err != nil { - return err + return fmt.Errorf("cannot load next pending agent run: %w", err) } run.Status = coredata.AgentRunStatusRunning @@ -90,6 +91,12 @@ func (h *agentRunHandler) Claim(ctx context.Context) (coredata.AgentRun, error) // renew the lease while the run is active, and a forwarder goroutine that // bridges the handler-level shutdown signal to the run's agent stop // channel so the agent checkpoints cleanly at the next turn boundary. +// +// The returned error mirrors the run outcome so the worker kit's +// task metrics and OTel span status reflect actual agent failures. +// nil is returned for both successful runs and graceful exits +// (lease loss, infrastructure suspension) where the row state is +// already consistent. func (h *agentRunHandler) Process(ctx context.Context, run coredata.AgentRun) error { runCtx, cancelRun := context.WithCancelCause(ctx) defer cancelRun(nil) @@ -110,9 +117,7 @@ func (h *agentRunHandler) Process(ctx context.Context, run coredata.AgentRun) er go h.heartbeatLease(heartbeatCtx, run.ID.String(), cancelRun) runCtx = agent.WithStopSignal(runCtx, stopCh) - h.executeRun(runCtx, &run) - - return nil + return h.executeRun(runCtx, &run) } // RecoverStale resets agent runs whose worker lease has expired back to @@ -179,7 +184,29 @@ func (h *agentRunHandler) heartbeatLease( } } -func (h *agentRunHandler) executeRun(ctx context.Context, run *coredata.AgentRun) { +const ( + // agentRunErrorMessageMaxLen caps the error string persisted to the + // agent_runs.error_message column. Raw tool or LLM errors can embed + // URLs with credentials, response snippets containing PII, or partial + // records from failed DB lookups; the full context is logged while + // only a truncated summary is stored for caller-visible state. + agentRunErrorMessageMaxLen = 512 +) + +func sanitizeAgentRunError(err error) string { + msg := err.Error() + if len(msg) <= agentRunErrorMessageMaxLen { + return msg + } + + cut := agentRunErrorMessageMaxLen + for cut > 0 && !utf8.RuneStart(msg[cut]) { + cut-- + } + return msg[:cut] + "…" +} + +func (h *agentRunHandler) executeRun(ctx context.Context, run *coredata.AgentRun) error { runID := run.ID.String() var ( @@ -211,7 +238,8 @@ func (h *agentRunHandler) executeRun(ctx context.Context, run *coredata.AgentRun } // Heartbeat loss: another worker may have taken over. Do not commit - // any status — stale recovery will handle the row. + // any status — stale recovery will handle the row. Surface the cause + // so the worker kit logs and traces a failure for this attempt. if cause := context.Cause(ctx); errors.Is(cause, ErrAgentRunLeaseLost) || errors.Is(cause, ErrAgentRunHeartbeatFailed) { h.logger.WarnCtx( context.WithoutCancel(ctx), @@ -219,13 +247,14 @@ func (h *agentRunHandler) executeRun(ctx context.Context, run *coredata.AgentRun log.String("run_id", runID), log.Error(cause), ) - return + return cause } // Infrastructure-triggered suspension (graceful shutdown): leave the // row as RUNNING so stale recovery resets it to PENDING on restart. // The checkpoint was already saved by coreLoop before returning - // SuspendedError, so Restore will pick up where it left off. + // SuspendedError, so Restore will pick up where it left off. This + // is not a failure from the worker kit's perspective. if runErr != nil { if _, ok := errors.AsType[*agent.SuspendedError](runErr); ok { h.logger.InfoCtx( @@ -233,7 +262,7 @@ func (h *agentRunHandler) executeRun(ctx context.Context, run *coredata.AgentRun "agent run suspended by infrastructure; leaving for stale recovery", log.String("run_id", runID), ) - return + return nil } } @@ -243,21 +272,29 @@ func (h *agentRunHandler) executeRun(ctx context.Context, run *coredata.AgentRun run.LeaseOwner = nil run.LeaseExpiresAt = nil - switch { - case runErr == nil: + if runErr == nil { run.Status = coredata.AgentRunStatusCompleted if result != nil { data, err := json.Marshal(result) if err != nil { h.logger.ErrorCtx(ctx, "cannot marshal agent run result", log.Error(err)) + runErr = fmt.Errorf("cannot marshal agent run result: %w", err) } else { run.Result = data } } + } - default: + if runErr != nil { run.Status = coredata.AgentRunStatusFailed - msg := runErr.Error() + run.Result = nil + h.logger.ErrorCtx( + context.WithoutCancel(ctx), + "agent run failed", + log.String("run_id", runID), + log.Error(runErr), + ) + msg := sanitizeAgentRunError(runErr) run.ErrorMessage = &msg } @@ -280,5 +317,8 @@ func (h *agentRunHandler) executeRun(ctx context.Context, run *coredata.AgentRun }, ); err != nil { h.logger.ErrorCtx(commitCtx, "cannot commit agent run status", log.Error(err)) + return fmt.Errorf("cannot commit agent run status: %w", err) } + + return runErr }