From 95803311b3cd39820bdd5e17dc74d12e7c97a38f Mon Sep 17 00:00:00 2001 From: Moshe Immerman Date: Tue, 8 Sep 2026 12:22:55 +0300 Subject: [PATCH 1/4] refactor(aichat): collapse the duplicate suspended-seed waits MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit approval_execution.go carried two implementations of awaitSuspendedSeed: a free function with a 15s budget and a Service method with 5s. Production only ever called the method, so the copy the tests exercised was not the copy that ran, and the two had already drifted in both budget and error wording. Keep one — the free function, which takes an explicit store and is therefore testable — with the method's 5s budget and its richer errors, and reduce the method to resolving the store. Claude-Session: https://claude.ai/code/session_01Usd2NB1ZhiU5wMuXg76Y6W Claude-Session-Id: f046afbd-f74b-4513-b309-28a1d5c6b64c --- pkg/aichat/approval_execution.go | 66 ++++++++++---------------------- 1 file changed, 21 insertions(+), 45 deletions(-) diff --git a/pkg/aichat/approval_execution.go b/pkg/aichat/approval_execution.go index 481a8166..a86da331 100644 --- a/pkg/aichat/approval_execution.go +++ b/pkg/aichat/approval_execution.go @@ -11,30 +11,39 @@ import ( ) const ( - suspendedSeedTimeout = 15 * time.Second + // suspendedSeedWait bounds how long an approval resolution waits for the + // suspending turn's assistant message to land in the thread store. + suspendedSeedWait = 5 * time.Second suspendedSeedInterval = 25 * time.Millisecond ) -// awaitSuspendedSeed returns the assistant message the suspended turn ended on. -// The durable suspension (prompt run -> waiting) is committed while the same -// stream's persistence goroutine is still writing that assistant message, so an -// approval resolved the instant the run becomes resumable can observe the run -// before its transcript. Wait for that in-flight write instead of rejecting a -// legitimate approval, and fail loudly when it never lands. +// awaitSuspendedSeed returns the thread's trailing assistant message for the +// suspended turn. +// +// The durable suspension (prompt run -> waiting) is committed from the event +// pipeline while the same stream is still persisting the assistant message it +// suspended on, so an approval resolved the instant the run becomes resumable +// can observe the run before its transcript. The run's durable approval state +// guarantees that message is committed or imminent — wait the in-flight write +// out instead of failing a resolution that has already consumed the approval, +// and fail loudly when it never lands. func awaitSuspendedSeed(ctx context.Context, store ThreadStore, threadID, turnID string) (*UIMessage, error) { - deadline := time.Now().Add(suspendedSeedTimeout) + deadline := time.Now().Add(suspendedSeedWait) for { thread, err := store.Get(ctx, threadID) if err != nil { return nil, err } - if len(thread.Messages) > 0 { - seed := thread.Messages[len(thread.Messages)-1] + if count := len(thread.Messages); count > 0 { + seed := thread.Messages[count-1] if strings.EqualFold(seed.Role, string(api.RoleAssistant)) && seed.TurnID == turnID { return &seed, nil } } if time.Now().After(deadline) { + if len(thread.Messages) == 0 { + return nil, fmt.Errorf("captain chat session %s has no suspended assistant message", threadID) + } return nil, fmt.Errorf("captain chat session %s does not end with the suspended turn %s", threadID, turnID) } select { @@ -155,44 +164,11 @@ func enforceApprovalRuntimeProfile(spec api.Spec, resolved api.ComposedSpec) err return nil } -// suspendedSeedWait bounds how long an approval resolution waits for the -// suspending turn's assistant message to land in the thread store. -const suspendedSeedWait = 5 * time.Second - -// awaitSuspendedSeed returns the thread's trailing assistant message for the -// suspended turn. The prompt run reaches its waiting state from the event -// pipeline before the suspending stream persists that message on its final -// unwind, so an approval resolved from a session poll can arrive while the -// write is still in flight. The run's durable approval state guarantees the -// message is committed or imminent — wait it out instead of failing a -// resolution that has already consumed the approval. +// awaitSuspendedSeed resolves this service's thread store and waits there. func (s *Service) awaitSuspendedSeed(ctx context.Context, threadID, turnID string) (*UIMessage, error) { store, err := s.threads(ctx) if err != nil { return nil, err } - deadline := time.Now().Add(suspendedSeedWait) - for { - thread, err := store.Get(ctx, threadID) - if err != nil { - return nil, err - } - if count := len(thread.Messages); count > 0 { - seed := thread.Messages[count-1] - if strings.EqualFold(seed.Role, string(api.RoleAssistant)) && seed.TurnID == turnID { - return &seed, nil - } - } - if time.Now().After(deadline) { - if len(thread.Messages) == 0 { - return nil, fmt.Errorf("captain chat session %s has no suspended assistant message", threadID) - } - return nil, fmt.Errorf("captain chat session %s does not end with the suspended turn %s", threadID, turnID) - } - select { - case <-ctx.Done(): - return nil, ctx.Err() - case <-time.After(50 * time.Millisecond): - } - } + return awaitSuspendedSeed(ctx, store, threadID, turnID) } From b56593de05ea525b6344da7fc34c6a21f63a6194 Mon Sep 17 00:00:00 2001 From: Moshe Immerman Date: Tue, 8 Sep 2026 12:23:10 +0300 Subject: [PATCH 2/4] fix(aichat): wait for a suspending run to park before refusing its approval MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A provider approval is recorded — and so becomes visible on the session and goes out on the event stream carrying its ID — while the stream that raised it is still finishing the turn and encoding its checkpoint. Only afterwards does the prompt run reach `waiting`, the one state ResolveToolApprovalRequest accepts an answer in. Anything answering the question the moment it is asked lost that race and got a 409: a person clicking Approve promptly saw "Tool approval failed with status 409", and the mocked lifecycle suite failed the same way on CI. The store guard cannot simply be relaxed. A resolution applied before the run parks yields no continuation, and the suspension then parks the run on an already-answered approval that nothing ever resumes. So wait the parking out, bounded, the way awaitSuspendedSeed already waits out the other half of this same window — and fail loudly when it never happens, at once for a run that has already ended rather than burning the whole budget. The projection spec asserted that an early answer is refused, pinning the behaviour being changed; both halves are now covered explicitly instead. The prompt-run-conflict spec paused the first captain_prompt_runs query anywhere in the process, which the new pre-transaction read claims, so its interception is scoped to the read inside the transaction — what it always meant. Claude-Session: https://claude.ai/code/session_01Usd2NB1ZhiU5wMuXg76Y6W Claude-Session-Id: f046afbd-f74b-4513-b309-28a1d5c6b64c --- .../approval_settle_integration_test.go | 149 ++++++++++++++++++ .../database_threads_integration_test.go | 18 +-- pkg/aichat/execution_database_authority.go | 73 +++++++++ .../execution_database_integration_test.go | 7 +- 4 files changed, 229 insertions(+), 18 deletions(-) create mode 100644 pkg/aichat/approval_settle_integration_test.go diff --git a/pkg/aichat/approval_settle_integration_test.go b/pkg/aichat/approval_settle_integration_test.go new file mode 100644 index 00000000..4930932e --- /dev/null +++ b/pkg/aichat/approval_settle_integration_test.go @@ -0,0 +1,149 @@ +package aichat_test + +import ( + "context" + "encoding/json" + "time" + + "github.com/flanksource/captain/pkg/aichat" + "github.com/flanksource/captain/pkg/api" + "github.com/flanksource/captain/pkg/database" + "github.com/flanksource/commons-db/dbtest" + "github.com/google/uuid" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +// A provider approval is answerable — it is on the session, and its ID has gone +// out on the event stream — from the moment the permission frame is observed. +// The run it blocks only reaches `waiting` once the stream has finished the turn +// and encoded its checkpoint, several statements later. These specs pin what +// happens to an answer that arrives inside that window, which is where a person +// clicking Approve promptly, and the mocked lifecycle suite, both landed. +var _ = Describe("Approvals answered while the suspension is still landing", func() { + It("waits for the run to park rather than refusing the answer", func(ctx SpecContext) { + fixture := newApprovalFixture(ctx, "captain_aichat_approval_settles") + + // Nothing has parked the run yet; under the bare store guard this is the + // exact moment that produced "cannot be resolved before its prompt run is + // waiting" and a 409. + suspended := make(chan error, 1) + go func() { + time.Sleep(200 * time.Millisecond) + suspended <- suspendOnAccountsApproval(ctx, fixture.execution) + }() + + continuation, err := fixture.authority.ResolveToolApproval(ctx, aichat.ToolApprovalResolution{ + ThreadID: fixture.thread.ID, ApprovalID: fixture.approvalID, Approved: false, Reason: "not now", + }) + Expect(err).NotTo(HaveOccurred(), "an answer that raced the suspension is still a valid answer") + Expect(continuation).NotTo(BeNil(), "the resolution has to hand back the continuation that resumes the run") + DeferCleanup(continuation.Execution.Close) + Expect(<-suspended).To(Succeed()) + + resolved, err := fixture.store.GetSession(ctx, fixture.thread.ID) + Expect(err).NotTo(HaveOccurred()) + Expect(resolved.Requests).To(HaveLen(1)) + Expect(resolved.Requests[0].State).To(Equal(string(database.TurnRequestStateDenied))) + Expect(resolved.Requests[0].Reason).To(Equal("not now")) + }) + + It("refuses an answer once the run it blocks has already ended", func(ctx SpecContext) { + fixture := newApprovalFixture(ctx, "captain_aichat_approval_run_ended") + + runID, err := uuid.Parse(fixture.execution.PromptRunID()) + Expect(err).NotTo(HaveOccurred()) + run, err := fixture.db.GetPromptRun(ctx, runID) + Expect(err).NotTo(HaveOccurred()) + cancelled := database.PromptRunStateCancelled + _, err = fixture.db.UpdatePromptRun(ctx, database.UpdatePromptRunInput{ + ID: run.ID, ExpectedVersion: run.Version, State: &cancelled, + }) + Expect(err).NotTo(HaveOccurred()) + + // No suspension is coming, so this must fail on the run's state rather + // than burn the whole settle budget waiting for one. + started := time.Now() + _, err = fixture.authority.ResolveToolApproval(ctx, aichat.ToolApprovalResolution{ + ThreadID: fixture.thread.ID, ApprovalID: fixture.approvalID, Approved: true, + }) + Expect(err).To(MatchError(database.ErrTurnRequestConflict)) + Expect(err).To(MatchError(ContainSubstring("already ended (cancelled)"))) + Expect(time.Since(started)).To(BeNumerically("<", time.Second), + "a run that ended is a decided answer, not something to wait out") + }) +}) + +type approvalFixture struct { + db *database.DB + store *aichat.DatabaseThreadStore + authority *aichat.DatabaseExecutionAuthority + thread *aichat.Thread + execution aichat.Execution + approvalID string +} + +// newApprovalFixture drives a chat turn up to the point where the provider has +// asked for permission and the durable approval exists, but the run has not yet +// been parked. +func newApprovalFixture(ctx context.Context, name string) approvalFixture { + GinkgoHelper() + testDB := dbtest.ForGinkgo(dbtest.Options{Name: name}) + db, err := database.Open(ctx, database.WithDSN(testDB.DSN()), database.WithMigrations()) + Expect(err).NotTo(HaveOccurred()) + DeferCleanup(db.Close) + store, err := aichat.NewDatabaseThreadStore(db) + Expect(err).NotTo(HaveOccurred()) + thread, err := store.Create(ctx, "Accounts") + Expect(err).NotTo(HaveOccurred()) + authority, err := aichat.NewDatabaseExecutionAuthority(db) + Expect(err).NotTo(HaveOccurred()) + execution, err := authority.Begin(ctx, aichat.ExecutionRequest{ + ThreadID: thread.ID, RequestID: "user-message-1", Title: thread.Title, + Spec: api.Spec{Model: withCaps(api.Model{Name: "gemini", Mode: api.ModeAPI})}, + }) + Expect(err).NotTo(HaveOccurred()) + Expect(store.AppendMessage(ctx, thread.ID, aichat.UIMessage{ + ID: "user-message-1", TurnID: execution.TurnID(), Role: "user", + Parts: []aichat.UIPart{{Type: "text", Text: "Edit the account"}}, + })).To(Succeed()) + permission, err := execution.Observe(ctx, api.Event{ + Kind: api.EventPermission, ToolCallID: "call-account-1", Tool: "accounts_edit", + Input: map[string]any{"id": "acc-1"}, + }) + Expect(err).NotTo(HaveOccurred()) + Expect(store.AppendMessage(ctx, thread.ID, aichat.UIMessage{ + ID: execution.TurnID() + "-assistant", TurnID: execution.TurnID(), Role: "assistant", + Parts: []aichat.UIPart{{ + Type: "dynamic-tool", ToolName: "accounts_edit", ToolCallID: "call-account-1", + State: "approval-requested", Input: json.RawMessage(`{"id":"acc-1"}`), + Approval: &aichat.Approval{ID: permission.ApprovalID}, + }}, + })).To(Succeed()) + return approvalFixture{ + db: db, store: store, authority: authority, thread: thread, + execution: execution, approvalID: permission.ApprovalID, + } +} + +// suspendOnAccountsApproval completes the turn the way a provider that needs an +// approval does: a terminal result carrying the approval state and the private +// checkpoint the resume replays from. This is what parks the run in `waiting`. +func suspendOnAccountsApproval(ctx context.Context, execution aichat.Execution) error { + _, err := execution.Observe(ctx, api.Event{ + Kind: api.EventResult, Success: true, + ToolApproval: &api.ToolApprovalState{ + Messages: []api.Message{{Role: api.RoleAssistant, Parts: []api.Part{{ + Type: api.PartToolRequest, ToolRequest: &api.ToolRequest{ + ToolCallID: "call-account-1", Name: "accounts_edit", Input: json.RawMessage(`{"id":"acc-1"}`), + }, + }}}}, + Calls: []api.ToolApprovalCall{{Request: api.ToolApprovalRequest{ + ToolCallID: "call-account-1", Tool: "accounts_edit", Input: json.RawMessage(`{"id":"acc-1"}`), + }}}, + ProviderCheckpoint: &api.ProviderCheckpoint{Codec: "test-provider", Version: 1, Payload: []byte("checkpoint")}, + }, + }) + return err +} diff --git a/pkg/aichat/database_threads_integration_test.go b/pkg/aichat/database_threads_integration_test.go index 5d3bba3f..3029dfac 100644 --- a/pkg/aichat/database_threads_integration_test.go +++ b/pkg/aichat/database_threads_integration_test.go @@ -218,8 +218,6 @@ var _ = Describe("Database chat sessions", func() { resolution := aichat.ToolApprovalResolution{ ThreadID: thread.ID, ApprovalID: permission.ApprovalID, Approved: false, Reason: "not now", } - _, err = authority.ResolveToolApproval(ctx, resolution) - Expect(err).To(MatchError(ContainSubstring("cannot be resolved before its prompt run is waiting"))) assistant := aichat.UIMessage{ ID: execution.TurnID() + "-assistant", TurnID: execution.TurnID(), Role: "assistant", Parts: []aichat.UIPart{{ @@ -229,21 +227,7 @@ var _ = Describe("Database chat sessions", func() { }}, } Expect(store.AppendMessage(ctx, thread.ID, assistant)).To(Succeed()) - _, err = execution.Observe(ctx, api.Event{ - Kind: api.EventResult, Success: true, - ToolApproval: &api.ToolApprovalState{ - Messages: []api.Message{{Role: api.RoleAssistant, Parts: []api.Part{{ - Type: api.PartToolRequest, ToolRequest: &api.ToolRequest{ - ToolCallID: "call-account-1", Name: "accounts_edit", Input: json.RawMessage(`{"id":"acc-1"}`), - }, - }}}}, - Calls: []api.ToolApprovalCall{{Request: api.ToolApprovalRequest{ - ToolCallID: "call-account-1", Tool: "accounts_edit", Input: json.RawMessage(`{"id":"acc-1"}`), - }}}, - ProviderCheckpoint: &api.ProviderCheckpoint{Codec: "test-provider", Version: 1, Payload: []byte("checkpoint")}, - }, - }) - Expect(err).NotTo(HaveOccurred()) + Expect(suspendOnAccountsApproval(ctx, execution)).To(Succeed()) aggregate, err := store.GetSession(ctx, thread.ID) Expect(err).NotTo(HaveOccurred()) diff --git a/pkg/aichat/execution_database_authority.go b/pkg/aichat/execution_database_authority.go index 65cb8e8a..16f4b78c 100644 --- a/pkg/aichat/execution_database_authority.go +++ b/pkg/aichat/execution_database_authority.go @@ -6,6 +6,7 @@ import ( "fmt" "maps" "strings" + "time" "github.com/flanksource/captain/pkg/api" "github.com/flanksource/captain/pkg/database" @@ -132,10 +133,82 @@ func (a *DatabaseExecutionAuthority) Begin( return execution, nil } +const ( + // suspendedRunWait bounds how long a resolution waits for the run it answers + // to finish parking, and suspendedRunInterval is how often that is re-read. + // They match the seed wait's budget: the two wait out the two halves of the + // same in-flight suspension. + suspendedRunWait = 5 * time.Second + suspendedRunInterval = 25 * time.Millisecond +) + +// awaitSuspendedRun waits for a provider approval's prompt run to reach +// `waiting`, the one state ResolveToolApprovalRequest accepts an answer in. +// +// A provider approval is recorded — and so becomes visible on the session and +// goes out on the event stream carrying its approval ID — while the stream that +// raised it is still finishing the turn and encoding its checkpoint. The run +// only reaches `waiting` several statements later. So anything that answers the +// question the moment it is asked raced the suspension and got a 409 telling it +// to retry something that was never wrong: a person clicking Approve promptly, +// or a poller in a test. +// +// The guard being waited for is not removable. A resolution applied before the +// run parks yields no continuation (see resolveToolApproval), and the suspension +// then parks the run on an already-answered approval that nothing ever resumes. +// So wait the parking out — the same treatment awaitSuspendedSeed gives the +// other half of this window — and fail loudly when it never happens. +// +// This runs outside the resolving transaction deliberately: a snapshot taken +// inside one would never observe the suspending connection's commit. +func (a *DatabaseExecutionAuthority) awaitSuspendedRun(ctx context.Context, approvalID string) error { + requestID, err := uuid.Parse(approvalID) + if err != nil { + return nil // resolveToolApproval reports a malformed ID, with its own message + } + deadline := time.Now().Add(suspendedRunWait) + for { + request, err := a.db.GetTurnRequest(ctx, requestID) + if err != nil { + return nil // the resolve path owns not-found and read failures alike + } + // A caller-tool approval carries its own authority and is answerable + // whatever its run is doing. Anything already decided, or with no run to + // resume, is likewise the store's answer to give, not this wait's. + if request.CredentialID != nil || request.PromptRunID == nil || + request.State != database.TurnRequestStatePending { + return nil + } + run, err := a.db.GetPromptRun(ctx, *request.PromptRunID) + if err != nil { + return err + } + switch run.State { + case database.PromptRunStateWaiting: + return nil + case database.PromptRunStateSucceeded, database.PromptRunStateFailed, database.PromptRunStateCancelled: + return fmt.Errorf("%w: approval %s cannot be resolved, its prompt run %s already ended (%s)", + database.ErrTurnRequestConflict, request.ID, run.ID, run.State) + } + if time.Now().After(deadline) { + return fmt.Errorf("%w: approval %s is still pending after %s with its prompt run %s in state %q rather than waiting", + database.ErrTurnRequestConflict, request.ID, suspendedRunWait, run.ID, run.State) + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(suspendedRunInterval): + } + } +} + func (a *DatabaseExecutionAuthority) ResolveToolApproval( ctx context.Context, resolution ToolApprovalResolution, ) (*ApprovalContinuation, error) { + if err := a.awaitSuspendedRun(ctx, resolution.ApprovalID); err != nil { + return nil, err + } var continuation *ApprovalContinuation err := a.db.Transaction(ctx, func(tx *database.DB) error { var resolveErr error diff --git a/pkg/aichat/execution_database_integration_test.go b/pkg/aichat/execution_database_integration_test.go index 89d7b9df..8ccb3b32 100644 --- a/pkg/aichat/execution_database_integration_test.go +++ b/pkg/aichat/execution_database_integration_test.go @@ -385,8 +385,13 @@ var _ = Describe("Database execution authority", func() { continueResolution := make(chan struct{}) var intercepted atomic.Bool const callback = "test:pause_approval_after_prompt_run_read" + // Pause the run read the resolution *resumes* from — the one inside its + // transaction, whose version the update then asserts on. Resolution also + // reads the run outside any transaction first, to wait out a suspension + // still landing; pausing there would stall a read this race is not about. Expect(db.Gorm().Callback().Query().After("gorm:query").Register(callback, func(tx *gorm.DB) { - if tx.Statement.Table == "captain_prompt_runs" && intercepted.CompareAndSwap(false, true) { + _, inTransaction := tx.Statement.ConnPool.(gorm.TxCommitter) + if inTransaction && tx.Statement.Table == "captain_prompt_runs" && intercepted.CompareAndSwap(false, true) { close(versionRead) <-continueResolution } From 879990632c60a8b2eeb46adf1a84ecab778c1514 Mon Sep 17 00:00:00 2001 From: Moshe Immerman Date: Tue, 8 Sep 2026 12:23:22 +0300 Subject: [PATCH 3/4] fix(gitagent): publish the agent worktree only once it is complete git clone creates its target directory first and writes the local branch ref and its upstream config last, so cloning straight onto the published path left a window where the directory existed but its HEAD named a branch that did not: `git rev-parse @{u}` there fails with "no such branch", which is how the git-agent e2e cycle intermittently failed on CI. Build the workspace in a staging sibling and move it in with one rename. The path's existence now means what every observer already assumed, and an interrupted dispatch no longer leaves behind a partial worktree that the re-dispatch check mistakes for a finished one. Claude-Session: https://claude.ai/code/session_01Usd2NB1ZhiU5wMuXg76Y6W Claude-Session-Id: f046afbd-f74b-4513-b309-28a1d5c6b64c --- pkg/gitagent/workspace.go | 25 ++++++++++++++++-- pkg/gitagent/workspace_test.go | 46 ++++++++++++++++++++++++++++++++++ 2 files changed, 69 insertions(+), 2 deletions(-) diff --git a/pkg/gitagent/workspace.go b/pkg/gitagent/workspace.go index 173c09b3..8a2157d7 100644 --- a/pkg/gitagent/workspace.go +++ b/pkg/gitagent/workspace.go @@ -39,18 +39,39 @@ func SetupAgentWorkspace(ctx context.Context, sidecarRepo, task, dispatchCommit, if _, err := os.Stat(workdir); err == nil { return workdir, nil // re-dispatch onto an existing workspace is a no-op } + // Build the workspace beside its final path and move it in with one rename. + // `git clone` creates its target directory first and writes the local branch + // ref and its upstream config last, so cloning straight onto workdir + // publishes a directory whose HEAD names a branch that does not exist yet — + // `git rev-parse @{u}` there fails with "no such branch". Anything that + // watches the task directory to know the workspace is ready (the e2e suite, + // an operator, the re-dispatch check above) would be reading a worktree that + // is not one yet, and an interrupted clone would leave a partial workspace + // that every later dispatch mistakes for a finished one. + staging, err := os.MkdirTemp(filepath.Dir(workdir), ".worktree-") + if err != nil { + return "", err + } + defer func() { _ = os.RemoveAll(staging) }() // the rename leaves nothing to remove + // MkdirTemp is 0700; a clone is not. Restore the mode the agent would see. + if err := os.Chmod(staging, 0o755); err != nil { + return "", err + } branchName := "captain/" + task if _, err := runGit(ctx, filepath.Dir(workdir), env, - "clone", "--quiet", "--shared", "--branch", branchName, sidecarRepo, workdir); err != nil { + "clone", "--quiet", "--shared", "--branch", branchName, sidecarRepo, staging); err != nil { return "", err } // Pin the selected runtime so a bare `git commit` needs no global config // and still records which model and effort produced it. for _, kv := range [][2]string{{"user.name", runtimeIdentity}, {"user.email", "agent@captain.local"}} { - if _, err := runGit(ctx, workdir, env, "config", kv[0], kv[1]); err != nil { + if _, err := runGit(ctx, staging, env, "config", kv[0], kv[1]); err != nil { return "", err } } + if err := os.Rename(staging, workdir); err != nil { + return "", fmt.Errorf("publish agent workspace %s: %w", workdir, err) + } return workdir, nil } diff --git a/pkg/gitagent/workspace_test.go b/pkg/gitagent/workspace_test.go index b96037bb..e496aa73 100644 --- a/pkg/gitagent/workspace_test.go +++ b/pkg/gitagent/workspace_test.go @@ -49,6 +49,52 @@ func TestSetupAgentWorkspacePinsRuntimeIdentity(t *testing.T) { } } +// The workspace path is published only once it is a workspace. `git clone` +// creates its target first and writes the branch ref and upstream config last, +// so cloning straight onto the final path let an observer read a directory +// whose HEAD named a branch that did not exist yet — which is how the git-agent +// e2e cycle intermittently failed its `@{u}` check on CI. +func TestSetupAgentWorkspaceIsCompleteWhenItAppears(t *testing.T) { + ctx := context.Background() + repo := filepath.Join(t.TempDir(), "sidecar.git") + if err := InitSidecar(ctx, repo); err != nil { + t.Fatal(err) + } + commit, err := BuildControlCommit(ctx, repo, nil, map[string][]byte{"seed.txt": []byte("seed\n")}) + if err != nil { + t.Fatal(err) + } + const task = "t-complete" + if err := SaveTaskState(repo, &TaskState{Task: task}); err != nil { + t.Fatal(err) + } + workdir, err := SetupAgentWorkspace(ctx, repo, task, commit, "captain-agent") + if err != nil { + t.Fatal(err) + } + // A bare `git push` needs an upstream, and it has to be there the moment the + // path exists — not a few milliseconds later (H17). + env := ScrubGitEnv(os.Environ()) + upstream, err := runGit(ctx, workdir, env, "rev-parse", "--abbrev-ref", "@{u}") + if err != nil { + t.Fatalf("the published workspace has no upstream: %v", err) + } + if !strings.HasSuffix(upstream, task) { + t.Fatalf("upstream = %q, want one ending in %q", upstream, task) + } + // The staging directory the clone was built in is gone, so a later dispatch + // cannot mistake a half-built workspace for a finished one. + entries, err := os.ReadDir(taskStateDir(repo, task)) + if err != nil { + t.Fatal(err) + } + for _, entry := range entries { + if strings.HasPrefix(entry.Name(), ".worktree-") { + t.Fatalf("staging directory %q survived the setup", entry.Name()) + } + } +} + // A dispatch that launches nothing leaves the supervisor waiting out its whole // budget on work that never started — a silence indistinguishable from an // agent still thinking. Empty must therefore be an error, and "no agent" must From 7200771ca2427e647f66995331f3c9ea9f94a8ee Mon Sep 17 00:00:00 2001 From: Moshe Immerman Date: Tue, 8 Sep 2026 12:23:41 +0300 Subject: [PATCH 4/4] test(cli): keep the Go toolchain caches across the test HOME override The package hands itself a throwaway HOME so a developer's ~/.captain.yaml cannot fail specs, including the ones that shell out to a captain subprocess. But the Go toolchain derives GOPATH, GOCACHE and GOMODCACHE from HOME whenever they are unset, so the `go build ./cmd/captain` the git-agent e2e tests run inherited an empty module cache and an empty build cache and re-downloaded and recompiled the entire dependency tree, cgo sqlite3 included, on every run. On CI that was ten minutes for this one package, with the runner's warm caches sitting untouched. Resolve the three variables to absolute paths before HOME is replaced, so the isolation covers captain's config and nothing else. Measured on one e2e test against an already-warm cache: 408s before, 57s after. Claude-Session: https://claude.ai/code/session_01Usd2NB1ZhiU5wMuXg76Y6W Claude-Session-Id: f046afbd-f74b-4513-b309-28a1d5c6b64c --- pkg/cli/main_test.go | 58 ++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 58 insertions(+) diff --git a/pkg/cli/main_test.go b/pkg/cli/main_test.go index 7c58c873..dd211260 100644 --- a/pkg/cli/main_test.go +++ b/pkg/cli/main_test.go @@ -1,8 +1,12 @@ package cli import ( + "encoding/json" + "fmt" "os" + "os/exec" "path/filepath" + "strings" "testing" "github.com/flanksource/captain/pkg/captainconfig" @@ -23,6 +27,9 @@ func TestMain(m *testing.M) { if _, ok := os.LookupEnv("CAPTAIN_SESSION_DB_URL"); !ok { _ = os.Setenv("CAPTAIN_SESSION_DB_URL", "off") } + if err := pinGoToolchainEnv(); err != nil { + panic("pkg/cli tests: " + err.Error()) + } home, err := os.MkdirTemp("", "captain-cli-home") if err != nil { panic("pkg/cli tests: isolate HOME: " + err.Error()) @@ -33,3 +40,54 @@ func TestMain(m *testing.M) { _ = os.RemoveAll(home) os.Exit(code) } + +// pinGoToolchainEnv resolves the Go toolchain's cache locations to absolute +// paths and exports them, so the HOME override above cannot relocate them. +// +// The HOME this package hands itself is for captain's own config. But the Go +// toolchain derives GOPATH, GOCACHE and GOMODCACHE from HOME whenever they are +// unset, so every `go build` a test shells out to — captainBinary builds +// ./cmd/captain — inherited an empty module cache and an empty build cache and +// re-downloaded and recompiled the entire dependency tree, cgo sqlite3 +// included. On CI that was ten minutes for this one package, with the runner's +// warm caches sitting untouched. It must run before HOME is replaced. +func pinGoToolchainEnv() error { + names := []string{"GOCACHE", "GOMODCACHE", "GOPATH"} + out, err := exec.Command("go", append([]string{"env", "-json"}, names...)...).Output() + if err != nil { + return fmt.Errorf("resolve the Go toolchain environment: %w", err) + } + var resolved map[string]string + if err := json.Unmarshal(out, &resolved); err != nil { + return fmt.Errorf("decode `go env -json %v`: %w", names, err) + } + for _, name := range names { + if resolved[name] == "" { + return fmt.Errorf("`go env -json` reported no %s; a toolchain subprocess would derive it from the throwaway HOME", name) + } + if err := os.Setenv(name, resolved[name]); err != nil { + return fmt.Errorf("pin %s: %w", name, err) + } + } + return nil +} + +// A toolchain subprocess started from this package must reach the same caches +// as the rest of the build. Left to inherit the throwaway HOME it reaches none +// of them, which is invisible locally and cost ten minutes per CI run. +func TestGoToolchainCachesSurviveTheHomeOverride(t *testing.T) { + home := os.Getenv("HOME") + if home == "" { + t.Fatal("TestMain is expected to hand this package its own HOME") + } + for _, name := range []string{"GOCACHE", "GOMODCACHE", "GOPATH"} { + value := os.Getenv(name) + if value == "" { + t.Errorf("%s is unset, so `go build` derives it from HOME and gets a cold cache", name) + continue + } + if rel, err := filepath.Rel(home, value); err == nil && !strings.HasPrefix(rel, "..") { + t.Errorf("%s = %q sits inside the throwaway HOME %q, so every toolchain subprocess gets a cold cache", name, value, home) + } + } +}