Skip to content

Commit 088955f

Browse files
authored
fix(mothership): answer a retried task wake whose turn already ran (#8484)
* fix(mothership): answer a retried task wake whose turn already ran The worker retries a task wake under the same run ID until its own run appears. When sim ended that turn without reaching the worker (a usage-limit refusal), every retry reopened the turn, hit the unique stream-id constraint, and failed behind a generic message, so the worker retried forever. Under the chat lock, a wake whose run ID already has a sim run now releases the lock and answers not-found, which the worker treats as a refusal and dismisses the notification. An in-flight turn still holds the lock and answers busy. The headless run-record catch-all now logs the underlying insert error. * fix(mothership): release the wake's chat lock when the run lookup fails The retried-wake check reads copilot_runs after taking the chat lock. If that read threw, the lock stayed held until its TTL because the wake turn that releases it never started. Release with the exact lease on any throw after the acquire. * fix(mothership): release the wake's own chat lease when the run lookup fails * test(mothership): stub the chat lease getter in the wake unit test mock
1 parent dd54d31 commit 088955f

4 files changed

Lines changed: 244 additions & 3 deletions

File tree

‎apps/sim/lib/mothership/request/lifecycle/run.ts‎

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1663,8 +1663,16 @@ async function ensureHeadlessRunIdentity(input: {
16631663
},
16641664
})
16651665
return { executionId, runId, cancelled: run.status === 'cancelled' }
1666-
} catch {
1667-
throw new Error('Chat could not start because its execution record is unavailable')
1666+
} catch (error) {
1667+
logger.error('Headless run record could not be created', {
1668+
chatId: input.chatId,
1669+
streamId: input.messageId,
1670+
error: getErrorMessage(error),
1671+
...causeForLog(error),
1672+
})
1673+
throw new Error('Chat could not start because its execution record is unavailable', {
1674+
cause: error,
1675+
})
16681676
}
16691677
}
16701678

Lines changed: 210 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,210 @@
1+
/**
2+
* How sim answers the worker's retry of a task wake, against real PostgreSQL and Redis: the
3+
* wake route, the chat stream lock, and the run records are production code. Only `after` is
4+
* stubbed, so the background wake turn never starts; each test writes the run record that
5+
* turn would have written instead. The run lookup passes through to PostgreSQL unless a test
6+
* makes it fail.
7+
*/
8+
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
9+
10+
const { redisUrl, inheritedRedisUrl } = await vi.hoisted(async () => {
11+
const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure')
12+
const url = readTestRedisUrl()
13+
const inheritedRedisUrl = process.env.REDIS_URL
14+
/** The real Redis module reads this at import. */
15+
if (url) process.env.REDIS_URL = url
16+
return { redisUrl: url, inheritedRedisUrl }
17+
})
18+
19+
vi.mock('next/server', async (original) => ({
20+
...(await original<typeof import('next/server')>()),
21+
after: () => {},
22+
}))
23+
24+
vi.mock('@/lib/mothership/async-runs/repository', async (original) => {
25+
const actual = await original<typeof import('@/lib/mothership/async-runs/repository')>()
26+
return { ...actual, getLatestRunForStream: vi.fn(actual.getLatestRunForStream) }
27+
})
28+
29+
import { db } from '@sim/db'
30+
import { copilotChats, copilotRuns, permissions, user, workspace } from '@sim/db/schema'
31+
import { generateId } from '@sim/utils/id'
32+
import { eq, inArray } from 'drizzle-orm'
33+
import { NextRequest } from 'next/server'
34+
import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis'
35+
import {
36+
createRunSegment,
37+
getLatestRunForStream,
38+
updateRunStatus,
39+
} from '@/lib/mothership/async-runs/repository'
40+
import { chatPubSub } from '@/lib/mothership/chat-status'
41+
import {
42+
acquirePendingChatStream,
43+
getLocalChatStreamLease,
44+
releasePendingChatStream,
45+
} from '@/lib/mothership/request/session/abort'
46+
import { POST as wakeRoute } from '@/app/api/mothership/wake/route'
47+
48+
afterAll(async () => {
49+
chatPubSub?.dispose()
50+
await closeRedisConnection()
51+
if (inheritedRedisUrl === undefined) Reflect.deleteProperty(process.env, 'REDIS_URL')
52+
else process.env.REDIS_URL = inheritedRedisUrl
53+
})
54+
55+
describe.runIf(Boolean(redisUrl))('task wake retries', () => {
56+
const userId = generateId()
57+
const workspaceId = generateId()
58+
const chatIds: string[] = []
59+
60+
beforeAll(async () => {
61+
const now = new Date()
62+
await db.insert(user).values({
63+
id: userId,
64+
name: 'Task wake fixture',
65+
email: `${userId}@task-wake.test`,
66+
emailVerified: true,
67+
createdAt: now,
68+
updatedAt: now,
69+
})
70+
await db.insert(workspace).values({
71+
id: workspaceId,
72+
name: 'Task wake fixture',
73+
ownerId: userId,
74+
billedAccountUserId: userId,
75+
})
76+
await db.insert(permissions).values({
77+
id: generateId(),
78+
userId,
79+
entityType: 'workspace',
80+
entityId: workspaceId,
81+
permissionType: 'admin',
82+
})
83+
})
84+
85+
afterAll(async () => {
86+
if (chatIds.length) {
87+
await db.delete(copilotRuns).where(inArray(copilotRuns.chatId, chatIds))
88+
await db.delete(copilotChats).where(inArray(copilotChats.id, chatIds))
89+
}
90+
await db.delete(permissions).where(eq(permissions.userId, userId))
91+
await db.delete(workspace).where(eq(workspace.id, workspaceId))
92+
await db.delete(user).where(eq(user.id, userId))
93+
})
94+
95+
async function idleChat() {
96+
const chatId = generateId()
97+
chatIds.push(chatId)
98+
await db.insert(copilotChats).values({ id: chatId, userId, workspaceId, type: 'mothership' })
99+
return chatId
100+
}
101+
102+
/** The worker's wake call, as `wakeOnSim` sends it. */
103+
function wake(chatId: string, runId: string) {
104+
return wakeRoute(
105+
new NextRequest('http://localhost:3000/api/mothership/wake', {
106+
method: 'POST',
107+
headers: {
108+
'content-type': 'application/json',
109+
'x-api-key': process.env.INTERNAL_API_SECRET ?? '',
110+
'x-mothership-user-id': userId,
111+
'x-mothership-workspace-id': workspaceId,
112+
},
113+
body: JSON.stringify({
114+
taskId: generateId(),
115+
runId,
116+
chatId,
117+
userId,
118+
workspaceId,
119+
message: 'Timer elapsed',
120+
status: 'completed',
121+
summary: 'Timer elapsed',
122+
}),
123+
}),
124+
{ params: Promise.resolve({}) }
125+
)
126+
}
127+
128+
/** The run record the headless wake turn opens under the wake's run ID. */
129+
function openWakeTurn(chatId: string, runId: string) {
130+
return createRunSegment({
131+
executionId: generateId(),
132+
chatId,
133+
userId,
134+
workspaceId,
135+
streamId: runId,
136+
requestContext: { source: 'headless_lifecycle' },
137+
})
138+
}
139+
140+
it('answers not-found to a wake whose turn already ended, and leaves the chat free', async () => {
141+
const chatId = await idleChat()
142+
const runId = generateId()
143+
expect((await wake(chatId, runId)).status).toBe(202)
144+
/** The turn ends inside sim without reaching the worker, as a usage-limit refusal does. */
145+
const turn = await openWakeTurn(chatId, runId)
146+
await updateRunStatus(turn.id, 'complete', { completedAt: new Date() })
147+
await releasePendingChatStream(chatId, runId)
148+
149+
/** The worker saw no run under this ID, so it retries the same wake. */
150+
expect((await wake(chatId, runId)).status).toBe(404)
151+
152+
const nextTurn = generateId()
153+
expect(await acquirePendingChatStream(chatId, nextTurn, 0)).toBe(true)
154+
await releasePendingChatStream(chatId, nextTurn)
155+
})
156+
157+
it('answers busy, never not-found, while the wake turn under that ID still holds the chat', async () => {
158+
const chatId = await idleChat()
159+
const runId = generateId()
160+
expect((await wake(chatId, runId)).status).toBe(202)
161+
await openWakeTurn(chatId, runId)
162+
163+
expect((await wake(chatId, runId)).status).toBe(409)
164+
await releasePendingChatStream(chatId, runId)
165+
}, 15_000)
166+
167+
it('frees the chat when the run lookup fails after the wake took it', async () => {
168+
const chatId = await idleChat()
169+
const runId = generateId()
170+
vi.mocked(getLatestRunForStream).mockRejectedValueOnce(new Error('statement timeout'))
171+
172+
expect((await wake(chatId, runId)).status).toBe(500)
173+
174+
const nextTurn = generateId()
175+
expect(await acquirePendingChatStream(chatId, nextTurn, 0)).toBe(true)
176+
await releasePendingChatStream(chatId, nextTurn)
177+
})
178+
179+
it("keeps a retry's chat lock when an earlier wake's slow lookup fails after its own lock expired", async () => {
180+
const chatId = await idleChat()
181+
const runId = generateId()
182+
let lookupStarted!: () => void
183+
const started = new Promise<void>((resolve) => {
184+
lookupStarted = resolve
185+
})
186+
let failLookup!: () => void
187+
const failed = new Promise<void>((resolve) => {
188+
failLookup = resolve
189+
})
190+
vi.mocked(getLatestRunForStream).mockImplementationOnce(async () => {
191+
lookupStarted()
192+
await failed
193+
throw new Error('statement timeout')
194+
})
195+
196+
const slowWake = wake(chatId, runId)
197+
await started
198+
/** The first wake's lock outlives its TTL while the lookup hangs. */
199+
const firstLease = getLocalChatStreamLease(chatId, runId)
200+
await getRedisClient()?.del(firstLease?.key ?? '')
201+
expect((await wake(chatId, runId)).status).toBe(202)
202+
203+
failLookup()
204+
expect((await slowWake).status).toBe(500)
205+
206+
const nextTurn = generateId()
207+
expect(await acquirePendingChatStream(chatId, nextTurn, 0)).toBe(false)
208+
await releasePendingChatStream(chatId, runId)
209+
}, 15_000)
210+
})

‎apps/sim/lib/mothership/tasks/application/prepare-wake.ts‎

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,13 @@
11
import { OrchestrationError } from '@/lib/core/orchestration/types'
2+
import { getLatestRunForStream } from '@/lib/mothership/async-runs/repository'
23
import { defineAuthorizedChatUseCase } from '@/lib/mothership/chat/application/authorized-chat-use-case'
34
import { resolveOwnedChatContext } from '@/lib/mothership/chat/application/context'
45
import type { TaskWakeRequest } from '@/lib/mothership/generated/tasks'
5-
import { acquirePendingChatStream } from '@/lib/mothership/request/session/abort'
6+
import {
7+
acquirePendingChatStream,
8+
getLocalChatStreamLease,
9+
releasePendingChatStream,
10+
} from '@/lib/mothership/request/session/abort'
611
import { taskDelegationPolicy } from '@/lib/mothership/tasks/application/context'
712
import {
813
organizationTaskOperations,
@@ -39,6 +44,23 @@ export const prepareTaskWake = defineAuthorizedChatUseCase({
3944
if (!(await acquirePendingChatStream(input.chatId, input.runId))) {
4045
throw new OrchestrationError('conflict', 'Another stream holds this chat; retry the wake')
4146
}
47+
const lease = getLocalChatStreamLease(input.chatId, input.runId)
48+
/**
49+
* The worker retries a wake under the same run ID until its own run appears. A turn sim
50+
* already ran under that ID without reaching the worker (a usage-limit refusal) can never
51+
* open again, so answer not-found: the worker dismisses the notification instead of
52+
* retrying forever. Checked under the chat lock, so an in-flight turn still answers busy.
53+
* Any throw here releases the lock just taken, by its own lease: a slow lookup can outlive
54+
* the lock, and a retry under the same run ID may hold the chat by then.
55+
*/
56+
try {
57+
if (await getLatestRunForStream(input.runId)) {
58+
throw new OrchestrationError('not_found', 'This wake already ran')
59+
}
60+
} catch (error) {
61+
await releasePendingChatStream(input.chatId, input.runId, lease)
62+
throw error
63+
}
4264
return { accepted: true } as const
4365
},
4466
})

‎apps/sim/lib/mothership/tasks/application/tasks.test.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,7 @@ vi.mock('@/lib/workflows/executor/execution-status', () => ({
4747
}))
4848
vi.mock('@/lib/mothership/request/session/abort', () => ({
4949
acquirePendingChatStream: hoisted.acquire,
50+
getLocalChatStreamLease: vi.fn(),
5051
}))
5152

5253
import type { NextRequest } from 'next/server'

0 commit comments

Comments
 (0)