mirror of
https://github.com/coder/coder.git
synced 2026-09-21 20:51:01 +08:00
fix: stabilize flaky chatd subscribe/promote queued tests (#23816)
## Summary Fixes three flaky chatd tests that intermittently fail due to timing races with the background run loop. Closes coder/internal#1428 ## Root Cause `CreateChat` and `PromoteQueued` call `signalWake()` which writes to `wakeCh`, triggering `processOnce` immediately. Even though `newTestServer` sets `PendingChatAcquireInterval: testutil.WaitLong` to prevent ticker-based polling, the wake channel bypasses this. This causes `processOnce` to acquire and process the chat concurrently with the test's manual DB updates and assertions. ### Failing tests | Test | Failure | Cause | |------|---------|-------| | `TestPromoteQueuedAllowsAlreadyQueuedMessageWhenUsageLimitReached` | `expected: "pending", actual: "running"` | Wake from `CreateChat` races with manual `UpdateChatStatus`; wake from `PromoteQueued` acquires the chat before the status assertion | | `TestSendMessageInterruptBehaviorQueuesAndInterruptsWhenBusy` | `should have 1 item(s), but has 2` | Wake from `CreateChat` triggers `processChat` which auto-promotes a queued message, adding an extra row to `chat_messages` | | `TestSubscribeNoPubsubNoDuplicateMessageParts` | `Condition satisfied` (duplicate events) | Pre-existing `WaitGroup.Add/Wait` race in the `Eventually` + `WaitUntilIdleForTest` pattern | ## Fix Introduces a `waitForChatProcessed` helper that: 1. Polls until the chat reaches a **terminal state** (not pending AND not running) 2. Then calls `WaitUntilIdleForTest` to wait for the inflight `WaitGroup` Waiting for a terminal state (not just "not pending") avoids a `sync.WaitGroup` `Add/Wait` race: `AcquireChats` updates the DB status to `running` **before** `processOnce` calls `inflight.Add(1)`. Checking only `status != pending` could return while `Add(1)` hasn't happened yet, causing `Wait()` to return prematurely. ### Per-test changes - **`TestSendMessageInterruptBehaviorQueuesAndInterruptsWhenBusy`**: Call `waitForChatProcessed` after `CreateChat` before manually setting running status - **`TestPromoteQueuedAllowsAlreadyQueuedMessageWhenUsageLimitReached`**: Call `waitForChatProcessed` after `CreateChat`; remove the inherently racy `status == pending` assertion after `PromoteQueued` (the wake immediately acquires the chat). Key assertions on promoted message, queue state, and message count remain. - **`TestSubscribeNoPubsubNoDuplicateMessageParts`**: Replace inline `Eventually` with the safer `waitForChatProcessed` helper ## Verification All three tests pass 150 consecutive executions with `-race -count=10` across 15 runs (0 failures).
This commit is contained in:
@@ -473,6 +473,11 @@ func TestSendMessageInterruptBehaviorQueuesAndInterruptsWhenBusy(t *testing.T) {
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
// CreateChat calls signalWake which triggers processOnce in
|
||||
// the background. Wait for that processing to finish so it
|
||||
// doesn't race with the manual status update below.
|
||||
waitForChatProcessed(ctx, t, db, chat.ID, replica)
|
||||
|
||||
chat, err = db.UpdateChatStatus(ctx, database.UpdateChatStatusParams{
|
||||
ID: chat.ID,
|
||||
Status: database.ChatStatusRunning,
|
||||
@@ -817,6 +822,11 @@ func TestPromoteQueuedAllowsAlreadyQueuedMessageWhenUsageLimitReached(t *testing
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
// CreateChat calls signalWake which triggers processOnce in
|
||||
// the background. Wait for that processing to finish so it
|
||||
// doesn't race with the manual status update below.
|
||||
waitForChatProcessed(ctx, t, db, chat.ID, replica)
|
||||
|
||||
chat, err = db.UpdateChatStatus(ctx, database.UpdateChatStatusParams{
|
||||
ID: chat.ID,
|
||||
Status: database.ChatStatusRunning,
|
||||
@@ -879,10 +889,6 @@ func TestPromoteQueuedAllowsAlreadyQueuedMessageWhenUsageLimitReached(t *testing
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, database.ChatMessageRoleUser, result.PromotedMessage.Role)
|
||||
|
||||
chat, err = db.GetChatByID(ctx, chat.ID)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, database.ChatStatusPending, chat.Status)
|
||||
|
||||
queued, err := db.GetChatQueuedMessages(ctx, chat.ID)
|
||||
require.NoError(t, err)
|
||||
require.Empty(t, queued)
|
||||
@@ -1709,13 +1715,9 @@ func TestSubscribeNoPubsubNoDuplicateMessageParts(t *testing.T) {
|
||||
// subscribing, so the snapshot captures the final state.
|
||||
// The wake signal may trigger processOnce which will fail
|
||||
// (no LLM configured) and set the chat to error status.
|
||||
// Poll until the chat leaves pending status, then wait for
|
||||
// the goroutine to finish.
|
||||
require.Eventually(t, func() bool {
|
||||
c, err := db.GetChatByID(ctx, chat.ID)
|
||||
return err == nil && c.Status != database.ChatStatusPending
|
||||
}, testutil.WaitShort, testutil.IntervalFast)
|
||||
chatd.WaitUntilIdleForTest(replica)
|
||||
// Poll until the chat reaches a terminal state (not pending
|
||||
// and not running), then wait for the goroutine to finish.
|
||||
waitForChatProcessed(ctx, t, db, chat.ID, replica)
|
||||
|
||||
snapshot, events, cancel, ok := replica.Subscribe(ctx, chat.ID, nil, 0)
|
||||
require.True(t, ok)
|
||||
@@ -2598,6 +2600,39 @@ func TestHeartbeatNoWorkspaceNoBump(t *testing.T) {
|
||||
require.Equal(t, 0, count, "expected no workspaces to be flushed when chat has no workspace")
|
||||
}
|
||||
|
||||
// waitForChatProcessed waits for a wake-triggered processOnce to
|
||||
// fully complete for the given chat. It polls until the chat leaves
|
||||
// both pending and running states (meaning processChat has finished
|
||||
// its cleanup and updated the DB), then calls WaitUntilIdleForTest.
|
||||
//
|
||||
// Waiting for a terminal state (not just "not pending") avoids a
|
||||
// WaitGroup Add/Wait race: AcquireChats changes the DB status to
|
||||
// running before processOnce calls inflight.Add(1). If we only
|
||||
// waited for status != pending, we could call Wait() while Add(1)
|
||||
// hasn't happened yet.
|
||||
func waitForChatProcessed(
|
||||
ctx context.Context,
|
||||
t *testing.T,
|
||||
db database.Store,
|
||||
chatID uuid.UUID,
|
||||
server *chatd.Server,
|
||||
) {
|
||||
t.Helper()
|
||||
require.Eventually(t, func() bool {
|
||||
c, err := db.GetChatByID(ctx, chatID)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
// Wait until the chat reaches a terminal state — neither
|
||||
// pending (waiting to be acquired) nor running (being
|
||||
// processed). This guarantees that inflight.Add(1) has
|
||||
// already been called by processOnce.
|
||||
return c.Status != database.ChatStatusPending &&
|
||||
c.Status != database.ChatStatusRunning
|
||||
}, testutil.WaitShort, testutil.IntervalFast)
|
||||
chatd.WaitUntilIdleForTest(server)
|
||||
}
|
||||
|
||||
func newTestServer(
|
||||
t *testing.T,
|
||||
db database.Store,
|
||||
|
||||
Reference in New Issue
Block a user