From 98284f4387e940715f3b6828cd653c0c87de3422 Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Mon, 20 Jul 2026 12:01:35 +0100 Subject: [PATCH] fix(workflows): mask output env in persisted logs, preserve cancelled status, run agent step async --- agent/internal/sync/sync.go | 27 +++++++++++++++------ server/internal/services/workflow_runner.go | 24 ++++++++++++------ 2 files changed, 37 insertions(+), 14 deletions(-) diff --git a/agent/internal/sync/sync.go b/agent/internal/sync/sync.go index f9828e8..6d63b8b 100644 --- a/agent/internal/sync/sync.go +++ b/agent/internal/sync/sync.go @@ -14,6 +14,7 @@ import ( "path/filepath" "runtime" "strings" + "sync" "time" "github.com/mrhid6/vantage/agent/internal/config" @@ -169,6 +170,16 @@ func connectAndHandleStream(ctx context.Context, cfg *config.Config) error { log.Println("command stream connected") + // grpc streams are not safe for concurrent Send; RunStep results are sent + // from per-command goroutines, so all sends on this stream must go through + // this mutex-protected helper. + var sendMu sync.Mutex + send := func(msg *pb.AgentMessage) error { + sendMu.Lock() + defer sendMu.Unlock() + return stream.Send(msg) + } + for { cmd, err := stream.Recv() if err != nil { @@ -188,13 +199,15 @@ func connectAndHandleStream(ctx context.Context, cfg *config.Config) error { go handleApplyUpdates(cfg, cmd) } if cmd.RunStep != nil { - res := agentexec.RunStep(cmd.RunStep) - res.CommandId = cmd.CommandId - _ = stream.Send(&pb.AgentMessage{ - ServerId: cfg.ServerID, - AgentToken: cfg.AgentToken, - StepResult: res, - }) + go func(rc *pb.RunStepCmd, cid string) { + res := agentexec.RunStep(rc) + res.CommandId = cid + _ = send(&pb.AgentMessage{ + ServerId: cfg.ServerID, + AgentToken: cfg.AgentToken, + StepResult: res, + }) + }(cmd.RunStep, cmd.CommandId) continue } } diff --git a/server/internal/services/workflow_runner.go b/server/internal/services/workflow_runner.go index de54128..88b3194 100644 --- a/server/internal/services/workflow_runner.go +++ b/server/internal/services/workflow_runner.go @@ -139,7 +139,7 @@ func executeRun(runID string) { now := time.Now() ctx, cancel := wfCtx() defer cancel() - _, _ = db.Col("workflow_runs").UpdateOne(ctx, bson.M{"run_id": runID}, + _, _ = db.Col("workflow_runs").UpdateOne(ctx, bson.M{"run_id": runID, "status": "running"}, bson.M{"$set": bson.M{"status": status, "finished_at": now}}) } @@ -156,6 +156,7 @@ func runServer(runID string, srvIdx int, steps []models.ResolvedStep, serverID s } runEnv := map[string]string{} + allSecrets := map[string]string{} serverFailed := false for i, step := range steps { @@ -169,6 +170,9 @@ func runServer(runID string, srvIdx int, steps []models.ResolvedStep, serverID s // Merge secrets into command env (kept out of persisted logs). secretVals := resolveSecrets(step.SecretRefs) + for k, v := range secretVals { + allSecrets[k] = v + } cmdEnv := map[string]string{} for k, v := range runEnv { cmdEnv[k] = v @@ -193,14 +197,14 @@ func runServer(runID string, srvIdx int, steps []models.ResolvedStep, serverID s // Mask secret values before persisting. stdout, stderr := "", "" exit := 1 - outEnv := map[string]string{} + outEnv := map[string]string{} // masked copy, safe to persist if res != nil { - stdout = maskSecrets(res.Stdout, secretVals) - stderr = maskSecrets(res.Stderr, secretVals) + stdout = maskSecrets(res.Stdout, allSecrets) + stderr = maskSecrets(res.Stderr, allSecrets) exit = res.ExitCode for k, v := range res.OutputEnv { - outEnv[k] = v - runEnv[k] = v // implicit: all outputs flow to all later steps + runEnv[k] = v // real, unmasked value threads forward to later steps + outEnv[k] = maskSecrets(v, allSecrets) } } else { stderr = "[vantage] agent did not return a result" @@ -231,10 +235,16 @@ func runServer(runID string, srvIdx int, steps []models.ResolvedStep, serverID s if serverFailed { status = "failed" } + // Persist only a masked copy of runEnv; the real (unmasked) runEnv was already + // used above to build cmdEnv for each step and must never be written to the DB. + maskedRunEnv := make(map[string]string, len(runEnv)) + for k, v := range runEnv { + maskedRunEnv[k] = maskSecrets(v, allSecrets) + } setServerRun(runID, srvIdx, bson.M{ "server_runs.$.status": status, "server_runs.$.finished_at": fin, - "server_runs.$.run_env": runEnv, + "server_runs.$.run_env": maskedRunEnv, }) }