Skip to content

Commit 7cb813b

Browse files
committed
fix(mothership): release the wake's own chat lease when the run lookup fails
1 parent 0bd25e7 commit 7cb813b

2 files changed

Lines changed: 40 additions & 10 deletions

File tree

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

Lines changed: 36 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ const { redisUrl, inheritedRedisUrl } = await vi.hoisted(async () => {
1212
const url = readTestRedisUrl()
1313
const inheritedRedisUrl = process.env.REDIS_URL
1414
/** The real Redis module reads this at import. */
15-
process.env.REDIS_URL = url
15+
if (url) process.env.REDIS_URL = url
1616
return { redisUrl: url, inheritedRedisUrl }
1717
})
1818

@@ -31,7 +31,7 @@ import { copilotChats, copilotRuns, permissions, user, workspace } from '@sim/db
3131
import { generateId } from '@sim/utils/id'
3232
import { eq, inArray } from 'drizzle-orm'
3333
import { NextRequest } from 'next/server'
34-
import { closeRedisConnection } from '@/lib/core/config/redis'
34+
import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis'
3535
import {
3636
createRunSegment,
3737
getLatestRunForStream,
@@ -40,14 +40,15 @@ import {
4040
import { chatPubSub } from '@/lib/mothership/chat-status'
4141
import {
4242
acquirePendingChatStream,
43+
getLocalChatStreamLease,
4344
releasePendingChatStream,
4445
} from '@/lib/mothership/request/session/abort'
4546
import { POST as wakeRoute } from '@/app/api/mothership/wake/route'
4647

4748
afterAll(async () => {
4849
chatPubSub?.dispose()
4950
await closeRedisConnection()
50-
if (inheritedRedisUrl === undefined) process.env.REDIS_URL = undefined
51+
if (inheritedRedisUrl === undefined) Reflect.deleteProperty(process.env, 'REDIS_URL')
5152
else process.env.REDIS_URL = inheritedRedisUrl
5253
})
5354

@@ -174,4 +175,36 @@ describe.runIf(Boolean(redisUrl))('task wake retries', () => {
174175
expect(await acquirePendingChatStream(chatId, nextTurn, 0)).toBe(true)
175176
await releasePendingChatStream(chatId, nextTurn)
176177
})
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)
177210
})

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

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -44,24 +44,21 @@ export const prepareTaskWake = defineAuthorizedChatUseCase({
4444
if (!(await acquirePendingChatStream(input.chatId, input.runId))) {
4545
throw new OrchestrationError('conflict', 'Another stream holds this chat; retry the wake')
4646
}
47+
const lease = getLocalChatStreamLease(input.chatId, input.runId)
4748
/**
4849
* The worker retries a wake under the same run ID until its own run appears. A turn sim
4950
* already ran under that ID without reaching the worker (a usage-limit refusal) can never
5051
* open again, so answer not-found: the worker dismisses the notification instead of
5152
* retrying forever. Checked under the chat lock, so an in-flight turn still answers busy.
52-
* Any throw here releases the lock just taken, since the wake turn that would release it
53-
* never starts.
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.
5455
*/
5556
try {
5657
if (await getLatestRunForStream(input.runId)) {
5758
throw new OrchestrationError('not_found', 'This wake already ran')
5859
}
5960
} catch (error) {
60-
await releasePendingChatStream(
61-
input.chatId,
62-
input.runId,
63-
getLocalChatStreamLease(input.chatId, input.runId)
64-
)
61+
await releasePendingChatStream(input.chatId, input.runId, lease)
6562
throw error
6663
}
6764
return { accepted: true } as const

0 commit comments

Comments
 (0)