fix(workflows): mask output env in persisted logs, preserve cancelled status, run agent step async
This commit is contained in:
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user