From bfe64d63550c5cf32a068e356da09c71f49a0813 Mon Sep 17 00:00:00 2001 From: Ethan <39577870+ethanndickson@users.noreply.github.com> Date: Mon, 15 Jun 2026 20:24:13 +1000 Subject: [PATCH] test(coderd): unskip chatd notification flow tests (#26366) Fixes coder/internal#1519 Fixes CODAGT-353 These nine tests were skipped pending a chatd notification flow refactor that would let workers distinguish stale control `NOTIFY` messages from real interrupts. They now pass consistently, so this drops the `t.Skip` calls and the now-stale `TODO(CODAGT-353)` blocks. While rerunning the package after unskipping them, `TestNewReplicaRecoversStaleChatFromDeadReplica` also surfaced as flaky on `main` because it asserted transient ownership state. This PR keeps that server-level test as a stable end-to-end recovery check and adds a deterministic worker-level stale reacquisition test so we still directly cover lease takeover behavior. ## Tests unskipped - `coderd`: `TestPatchChatMessage/ChangesModel` - `coderd/x/chatd`: - `TestExploreChatSendMessageCannotMutateMCPSnapshot` - `TestAutoPromoteQueuedMessagesPreservesPerTurnModelOrder` - `TestSignalWakeSendMessage` - `TestAdvisorChainMode_SnapshotKeepsFullHistory` - `TestOpenAIResponsesNoStaleWebSearchReplay` - `TestOpenAIResponsesFullReplayPairsReasoningAndWebSearch` - `TestOpenAIResponsesChainModeSkipsWhenLocalCallPending` - `TestOpenAIResponsesChainModeStillFiresForProviderExecutedOnly` ## Stale recovery follow-up - `coderd/x/chatd`: `TestNewReplicaRecoversStaleChatFromDeadReplica` now waits for the stable `waiting` and unowned end state after recovery. - `coderd/x/chatd`: `TestWorker_ReacquiresStaleOwnedChat` blocks the runner after reacquisition and directly asserts the new worker ownership, new runner ID, and fresh heartbeat. I stress-ran the stale recovery tests locally with repeated plain and race runs, and re-ran the nine unskipped tests plus `TestPatchChatMessage/ChangesModel` after these follow-up changes. --- coderd/exp_chats_test.go | 8 --- coderd/x/chatd/chatd_test.go | 62 ++++++-------------- coderd/x/chatd/integration_responses_test.go | 24 -------- coderd/x/chatd/worker_internal_test.go | 29 +++++++++ 4 files changed, 47 insertions(+), 76 deletions(-) diff --git a/coderd/exp_chats_test.go b/coderd/exp_chats_test.go index 06a2fbea9b..b2540edb16 100644 --- a/coderd/exp_chats_test.go +++ b/coderd/exp_chats_test.go @@ -8087,14 +8087,6 @@ func TestPatchChatMessage(t *testing.T) { t.Run("ChangesModel", func(t *testing.T) { t.Parallel() - // TODO(CODAGT-353): Re-enable this test after the chatd notification flow - // refactor gives workers enough causal information to distinguish stale - // control NOTIFY messages from real interrupts. The current design reuses - // the same status notification shape for wake-only and interrupt intents, - // so a stale NOTIFY can cancel a new processChat run. This subtest hits the - // same root cause via the persistInterruptedStep ownership gate, where a - // late insert from the previous turn regresses chats.last_model_config_id. - t.Skip("skipped until chatd notification flow refactor handles stale control notifications") ctx := testutil.Context(t, testutil.WaitLong) client := newChatClient(t) diff --git a/coderd/x/chatd/chatd_test.go b/coderd/x/chatd/chatd_test.go index b25a25fc24..3ed65b6953 100644 --- a/coderd/x/chatd/chatd_test.go +++ b/coderd/x/chatd/chatd_test.go @@ -1123,12 +1123,6 @@ func TestRootExploreChatExcludesWebSearchProviderToolAtRuntime(t *testing.T) { func TestExploreChatSendMessageCannotMutateMCPSnapshot(t *testing.T) { t.Parallel() - // TODO(CODAGT-353): Re-enable this test after the chatd notification flow - // refactor gives workers enough causal information to distinguish stale - // control NOTIFY messages from real interrupts. The current design reuses - // the same status notification shape for wake-only and interrupt intents, - // so a stale NOTIFY can cancel a new processChat run. - t.Skip("skipped until chatd notification flow refactor handles stale control notifications") db, ps := dbtestutil.NewDB(t) ctx := testutil.Context(t, testutil.WaitLong) @@ -2107,12 +2101,6 @@ func TestCreateChatRejectsWhenUsageLimitReached(t *testing.T) { func TestAutoPromoteQueuedMessagesPreservesPerTurnModelOrder(t *testing.T) { t.Parallel() - // TODO(CODAGT-353): Re-enable this test after the chatd notification flow - // refactor gives workers enough causal information to distinguish stale - // control NOTIFY messages from real interrupts. The current design reuses - // the same status notification shape for wake-only and interrupt intents, - // so a stale NOTIFY can cancel a new processChat run. - t.Skip("skipped until chatd notification flow refactor handles stale control notifications") db, ps := dbtestutil.NewDB(t) ctx := testutil.Context(t, testutil.WaitSuperLong) @@ -2934,16 +2922,10 @@ func TestNewReplicaRecoversStaleChatFromDeadReplica(t *testing.T) { if err != nil { return false } - return recovered.Status == database.ChatStatusRunning && - recovered.WorkerID.Valid && recovered.WorkerID.UUID == newWorkerID && - recovered.RunnerID.Valid && recovered.RunnerID.UUID != deadRunnerID + return recovered.Status == database.ChatStatusWaiting && + !recovered.WorkerID.Valid && + !recovered.RunnerID.Valid }, testutil.WaitMedium, testutil.IntervalFast) - - _, err = db.GetChatHeartbeat(ctx, database.GetChatHeartbeatParams{ - ChatID: created.Chat.ID, - RunnerID: recovered.RunnerID.UUID, - }) - require.NoError(t, err) } func TestWaitingChatsAreNotRecoveredAsStale(t *testing.T) { @@ -10956,12 +10938,11 @@ func TestChatTemplateAllowlistEnforcement(t *testing.T) { "create_workspace for blocked template should be rejected") } -// TestSignalWakeImmediateAcquisition verifies that CreateChat triggers -// immediate processing via signalWake without waiting for the polling -// ticker to fire. The ticker interval is set to an hour so it never -// fires during the test. Any processing must come from the wake -// channel. -func TestSignalWakeImmediateAcquisition(t *testing.T) { +// TestCreateChatImmediatelyProcessesNewChat verifies that CreateChat +// starts processing a new chat without waiting for the acquire ticker +// to fire. The ticker interval is set to an hour so it never fires +// during the test. +func TestCreateChatImmediatelyProcessesNewChat(t *testing.T) { t.Parallel() db, ps := dbtestutil.NewDB(t) @@ -10993,7 +10974,8 @@ func TestSignalWakeImmediateAcquisition(t *testing.T) { user, org, model := seedChatDependencies(t, db) setOpenAIProviderBaseURL(ctx, t, db, openAIURL) - // CreateChat sets status=pending and calls signalWake(). + // CreateChat should start the first turn without waiting for the + // acquire ticker. chat, err := server.CreateChat(ctx, chatd.CreateOptions{ OrganizationID: org.ID, OwnerID: user.ID, @@ -11006,8 +10988,8 @@ func TestSignalWakeImmediateAcquisition(t *testing.T) { // The chat should be processed immediately. The LLM handler // closes the `processed` channel when it receives a streaming - // request. Without signalWake this would hang forever because - // the 1-hour ticker never fires. + // request. If CreateChat only relied on the 1-hour ticker, + // this receive would time out. testutil.TryReceive(ctx, t, processed) chatd.WaitUntilIdleForTest(server) @@ -11019,13 +11001,11 @@ func TestSignalWakeImmediateAcquisition(t *testing.T) { "chat should be in waiting status after processing completes") } -// TestSignalWakeSendMessage verifies that SendMessage on an idle chat -// triggers immediate processing via signalWake. -func TestSignalWakeSendMessage(t *testing.T) { +// TestSendMessageImmediatelyProcessesWaitingChat verifies that sending +// a follow-up message to a waiting chat starts the next turn without +// waiting for the acquire ticker. +func TestSendMessageImmediatelyProcessesWaitingChat(t *testing.T) { t.Parallel() - // TODO(CODAGT-353): Re-enable this after the chatd notification - // flow can distinguish stale status notifications from interrupts. - t.Skip("skipped until chatd notification flow refactor handles stale control notifications") db, ps := dbtestutil.NewDB(t) ctx := testutil.Context(t, testutil.WaitSuperLong) @@ -11060,7 +11040,7 @@ func TestSignalWakeSendMessage(t *testing.T) { user, org, model := seedChatDependencies(t, db) setOpenAIProviderBaseURL(ctx, t, db, openAIURL) - // CreateChat triggers wake -> processes first turn. + // CreateChat processes the first turn immediately. chat, err := server.CreateChat(ctx, chatd.CreateOptions{ OrganizationID: org.ID, OwnerID: user.ID, @@ -11078,7 +11058,7 @@ func TestSignalWakeSendMessage(t *testing.T) { chatd.WaitUntilIdleForTest(server) // Now send a follow-up message, which should also be - // processed immediately via signalWake. + // processed immediately without waiting for the acquire ticker. _, err = server.SendMessage(ctx, chatd.SendMessageOptions{ ChatID: chat.ID, APIKeyID: testAPIKeyID(t, db, user.ID), @@ -12165,12 +12145,6 @@ func TestAdvisorGating_ExploreSubagent(t *testing.T) { // message, losing the context the outer model had been building on. func TestAdvisorChainMode_SnapshotKeepsFullHistory(t *testing.T) { t.Parallel() - // TODO(CODAGT-353): Re-enable this test after the chatd notification flow - // refactor gives workers enough causal information to distinguish stale - // control NOTIFY messages from real interrupts. The current design reuses - // the same status notification shape for wake-only and interrupt intents, - // so a stale NOTIFY can cancel a new processChat run. - t.Skip("skipped until chatd notification flow refactor handles stale control notifications") db, ps := dbtestutil.NewDB(t) ctx := testutil.Context(t, testutil.WaitLong) diff --git a/coderd/x/chatd/integration_responses_test.go b/coderd/x/chatd/integration_responses_test.go index e6ca40b1c8..bab9e7aa9c 100644 --- a/coderd/x/chatd/integration_responses_test.go +++ b/coderd/x/chatd/integration_responses_test.go @@ -25,12 +25,6 @@ import ( func TestOpenAIResponsesNoStaleWebSearchReplay(t *testing.T) { t.Parallel() - // TODO(CODAGT-353): Re-enable this test after the chatd notification flow - // refactor gives workers enough causal information to distinguish stale - // control NOTIFY messages from real interrupts. The current design reuses - // the same status notification shape for wake-only and interrupt intents, - // so a stale NOTIFY can cancel a new processChat run. - t.Skip("skipped until chatd notification flow refactor handles stale control notifications") db, ps := dbtestutil.NewDB(t) ctx := testutil.Context(t, testutil.WaitLong) @@ -116,12 +110,6 @@ func TestOpenAIResponsesNoStaleWebSearchReplay(t *testing.T) { func TestOpenAIResponsesFullReplayPairsReasoningAndWebSearch(t *testing.T) { t.Parallel() - // TODO(CODAGT-353): Re-enable this test after the chatd notification flow - // refactor gives workers enough causal information to distinguish stale - // control NOTIFY messages from real interrupts. The current design reuses - // the same status notification shape for wake-only and interrupt intents, - // so a stale NOTIFY can cancel a new processChat run. - t.Skip("skipped until chatd notification flow refactor handles stale control notifications") db, ps := dbtestutil.NewDB(t) ctx := testutil.Context(t, testutil.WaitLong) @@ -206,12 +194,6 @@ func TestOpenAIResponsesFullReplayPairsReasoningAndWebSearch(t *testing.T) { func TestOpenAIResponsesChainModeSkipsWhenLocalCallPending(t *testing.T) { t.Parallel() - // TODO(CODAGT-353): Re-enable this test after the chatd notification flow - // refactor gives workers enough causal information to distinguish stale - // control NOTIFY messages from real interrupts. The current design reuses - // the same status notification shape for wake-only and interrupt intents, - // so a stale NOTIFY can cancel a new processChat run. - t.Skip("skipped until chatd notification flow refactor handles stale control notifications") db, ps := dbtestutil.NewDB(t) ctx := testutil.Context(t, testutil.WaitLong) @@ -280,12 +262,6 @@ func TestOpenAIResponsesChainModeSkipsWhenLocalCallPending(t *testing.T) { func TestOpenAIResponsesChainModeStillFiresForProviderExecutedOnly(t *testing.T) { t.Parallel() - // TODO(CODAGT-353): Re-enable this test after the chatd notification flow - // refactor gives workers enough causal information to distinguish stale - // control NOTIFY messages from real interrupts. The current design reuses - // the same status notification shape for wake-only and interrupt intents, - // so a stale NOTIFY can cancel a new processChat run. - t.Skip("skipped until chatd notification flow refactor handles stale control notifications") db, ps := dbtestutil.NewDB(t) ctx := testutil.Context(t, testutil.WaitLong) diff --git a/coderd/x/chatd/worker_internal_test.go b/coderd/x/chatd/worker_internal_test.go index 6a635d8418..f01bb0d69c 100644 --- a/coderd/x/chatd/worker_internal_test.go +++ b/coderd/x/chatd/worker_internal_test.go @@ -98,6 +98,35 @@ func TestWorker_SkipsFreshlyOwnedChat(t *testing.T) { require.Equal(t, otherRunner, latest.RunnerID.UUID) } +func TestWorker_ReacquiresStaleOwnedChat(t *testing.T) { + t.Parallel() + f := newWorkerTestFixture(t) + chat := f.createRunningChat(t) + deadWorker := uuid.New() + deadRunner := uuid.New() + acquireChat(t, f, chat.ID, deadWorker, deadRunner) + makeHeartbeatStale(t, f, chat.ID, deadRunner) + starter := newBlockingTaskStarter(false) + worker := startWorker(t, testOptions(t, f, starter)) + + call := starter.waitCall(t, taskKindGeneration, chat.ID) + require.Equal(t, worker.chatWorkerID(), call.input.WorkerID) + require.Equal(t, database.ChatStatusRunning, call.input.Status) + require.NotEqual(t, deadRunner, call.input.RunnerID) + + latest, err := f.db.GetChatByID(testutil.Context(t, testutil.WaitShort), chat.ID) + require.NoError(t, err) + require.Equal(t, worker.chatWorkerID(), latest.WorkerID.UUID) + require.Equal(t, call.input.RunnerID, latest.RunnerID.UUID) + require.NotEqual(t, deadWorker, latest.WorkerID.UUID) + require.NotEqual(t, deadRunner, latest.RunnerID.UUID) + _, err = f.db.GetChatHeartbeat(testutil.Context(t, testutil.WaitShort), database.GetChatHeartbeatParams{ + ChatID: chat.ID, + RunnerID: call.input.RunnerID, + }) + require.NoError(t, err) +} + func TestWorker_TwoWorkersRaceSingleOwner(t *testing.T) { t.Parallel() f := newWorkerTestFixture(t)