From 460948d071d2c570737a80d2abce676e50d7e7bb Mon Sep 17 00:00:00 2001 From: dawn Date: Fri, 21 Aug 2026 23:32:51 +0900 Subject: [PATCH] spindle/engine: dont publish failed status before flushing wf log Signed-off-by: dawn --- spindle/engine/engine.go | 58 ++++++++++++++----- spindle/engine/engine_test.go | 104 +++++++++++++++++++++++++++++++++- 2 files changed, 144 insertions(+), 18 deletions(-) diff --git a/spindle/engine/engine.go b/spindle/engine/engine.go index 3fa5a241..cd60f960 100644 --- a/spindle/engine/engine.go +++ b/spindle/engine/engine.go @@ -155,6 +155,7 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, s } var err error var wfLogger models.WorkflowLogger + closeLog := func() {} if p, ok := eng.(workflowLoggerProvider); ok { wfLogger = p.WorkflowLogger(wid) } else if fileLogger, err := models.NewFileWorkflowLogger(cfg.Server.LogDir, wid, secretValues); err != nil { @@ -163,8 +164,16 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, s } else { l.Info("setup step logger; logs will be persisted", "logDir", cfg.Server.LogDir, "wid", wid) wfLogger = fileLogger + var closeOnce sync.Once + closeLog = func() { + closeOnce.Do(func() { + if err := fileLogger.Close(); err != nil { + l.Error("failed to close workflow log", "wid", wid, "err", err) + } + }) + } defer archiveWorkflowLog(l, stores, db, cfg.Server.LogDir, wid) - defer fileLogger.Close() + defer closeLog() } timeoutCtx, timeoutCancel := context.WithTimeout(ctx, workflowTimeout) @@ -186,15 +195,37 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, s l.Info("waiting for slot", "wid", wid) slot := WorkflowSlot(NoopSlot{}) _, remoteStatus := eng.(RemoteStatusEngine) + var publishTerminalStatus func() + destroyWorkflow := false + slotAcquired := false + setTerminalError := func(phase string, workflowErr error) { + publishTerminalStatus = func() { + writeWfError(db, n, l, wfCtx, wid, phase, workflowErr) + } + } + defer func() { + closeLog() + if publishTerminalStatus != nil { + publishTerminalStatus() + } + if destroyWorkflow { + if err := eng.DestroyWorkflow(ctx, wid); err != nil { + l.Error("failed to destroy workflow", "wid", wid, "err", err) + } + } + if slotAcquired { + slot.Release() + } + }() if s, ok := eng.(WorkflowSlotter); ok { slot, err = s.AcquireWorkflowSlot(wfCtx, wid, &w, Wait) if err != nil { - writeWfError(db, n, l, wfCtx, wid, "waiting for slot", err) + setTerminalError("waiting for slot", err) return } + slotAcquired = true } - defer slot.Release() if !remoteStatus { err := db.StatusRunning(wid, n) @@ -206,17 +237,13 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, s err = eng.SetupWorkflow(wfCtx, wid, &w, wfLogger) if err != nil { - if !isCanceled(wfCtx) { - if destroyErr := eng.DestroyWorkflow(ctx, wid); destroyErr != nil { - l.Error("failed to destroy workflow after setup failure", "error", destroyErr) - } - } + destroyWorkflow = !isCanceled(wfCtx) if !remoteStatus { - writeWfError(db, n, l, wfCtx, wid, "setting up workflow", err) + setTerminalError("setting up workflow", err) } return } - defer eng.DestroyWorkflow(ctx, wid) + destroyWorkflow = true for stepIdx, step := range w.Steps { if wfLogger != nil { @@ -235,7 +262,7 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, s if err != nil { if !remoteStatus { - writeWfError(db, n, l, wfCtx, wid, "running step", err) + setTerminalError("running step", err) } return } @@ -243,15 +270,16 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, s if isCanceled(wfCtx) { if !remoteStatus { - writeWfError(db, n, l, wfCtx, wid, "before success", nil) + setTerminalError("before success", nil) } return } if !remoteStatus { - err = db.StatusSuccess(wid, n) - if err != nil { - l.Error("failed to set workflow status to success", "wid", wid, "err", err) + publishTerminalStatus = func() { + if err := db.StatusSuccess(wid, n); err != nil { + l.Error("failed to set workflow status to success", "wid", wid, "err", err) + } } } }) diff --git a/spindle/engine/engine_test.go b/spindle/engine/engine_test.go index 90155334..66a8fc7d 100644 --- a/spindle/engine/engine_test.go +++ b/spindle/engine/engine_test.go @@ -2,13 +2,16 @@ package engine import ( "context" + "errors" "log/slog" "os" "path/filepath" + "strings" "sync" "testing" "time" + "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/tangled" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" @@ -30,7 +33,8 @@ type mockEngine struct { setupCalls []models.WorkflowId runStepCalls []models.WorkflowId setupFunc func(ctx context.Context, wid models.WorkflowId) error - runStepFunc func(ctx context.Context, wid models.WorkflowId, idx int) error + runStepFunc func(ctx context.Context, wid models.WorkflowId, idx int, wfLogger models.WorkflowLogger) error + destroyFunc func(ctx context.Context, wid models.WorkflowId) error timeout time.Duration } @@ -57,6 +61,9 @@ func (m *mockEngine) WorkflowTimeout() time.Duration { } func (m *mockEngine) DestroyWorkflow(ctx context.Context, wid models.WorkflowId) error { + if m.destroyFunc != nil { + return m.destroyFunc(ctx, wid) + } return nil } @@ -67,7 +74,7 @@ func (m *mockEngine) RunStep(ctx context.Context, wid models.WorkflowId, w *mode m.mu.Unlock() if fn != nil { - return fn(ctx, wid, idx) + return fn(ctx, wid, idx, wfLogger) } return nil } @@ -164,7 +171,7 @@ func TestCancelWorkflow_NotOverwritten(t *testing.T) { stepStarted := make(chan struct{}) eng := &mockEngine{ - runStepFunc: func(ctx context.Context, wid models.WorkflowId, idx int) error { + runStepFunc: func(ctx context.Context, wid models.WorkflowId, idx int, wfLogger models.WorkflowLogger) error { close(stepStarted) <-ctx.Done() return ctx.Err() @@ -225,6 +232,97 @@ func TestCancelWorkflow_NotOverwritten(t *testing.T) { } } +func TestStartWorkflows_FlushesLogBeforeTerminalStatus(t *testing.T) { + t.Parallel() + + testDB := newTestDB(t) + logger := slog.New(slog.NewTextHandler(os.Stderr, nil)) + logDir := t.TempDir() + vault, err := secrets.NewSQLiteManager(filepath.Join(t.TempDir(), "secrets.db")) + if err != nil { + t.Fatal(err) + } + + repoDid := syntax.DID("did:plc:test") + if err := vault.AddSecret(context.Background(), secrets.UnlockedSecret{ + Key: "LONG_SECRET", + Value: strings.Repeat("s", 256), + Repo: secrets.RepoIdentifier(repoDid.String()), + }); err != nil { + t.Fatal(err) + } + + destroyStarted := make(chan struct{}) + releaseDestroy := make(chan struct{}) + var releaseOnce sync.Once + release := func() { + releaseOnce.Do(func() { + close(releaseDestroy) + }) + } + t.Cleanup(release) + + eng := &mockEngine{ + runStepFunc: func(ctx context.Context, wid models.WorkflowId, idx int, wfLogger models.WorkflowLogger) error { + if _, err := wfLogger.DataWriter(idx, "stdout").Write([]byte("ssh debug hint\n")); err != nil { + return err + } + return errors.New("step failed") + }, + destroyFunc: func(ctx context.Context, wid models.WorkflowId) error { + close(destroyStarted) + <-releaseDestroy + return nil + }, + } + + pipelineId := models.PipelineId{Knot: "test-knot", Rkey: "test-rkey"} + wid := models.WorkflowId{PipelineId: pipelineId, Name: "failed_job"} + pipeline := &models.Pipeline{ + RepoDid: repoDid, + TrustedSource: true, + Workflows: map[models.Engine][]models.Workflow{ + eng: {{Name: wid.Name, Steps: []models.Step{mockStep{name: "step1"}}}}, + }, + } + cfg := &config.Config{Server: config.Server{LogDir: logDir}} + + done := make(chan struct{}) + go func() { + StartWorkflows(logger, vault, cfg, nil, testDB, nil, context.Background(), pipeline, pipelineId) + close(done) + }() + + select { + case <-destroyStarted: + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for workflow destruction") + } + + status, err := testDB.GetStatus(wid) + if err != nil { + t.Fatal(err) + } + if status.Status != string(models.StatusKindFailed) { + t.Fatalf("expected failed status, got %s", status.Status) + } + + log, err := os.ReadFile(models.LogFilePath(logDir, wid)) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(string(log), "ssh debug hint") { + t.Fatalf("terminal status was visible before the log was flushed: %s", log) + } + + release() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for StartWorkflows to complete") + } +} + func TestSetupTimeout_ReportsTimeout(t *testing.T) { t.Parallel() -- 2.51.2