From 85986a768829e709d18b853fadc455c1e8ffc5f0 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 14:47:59 -0700 Subject: [PATCH 1/4] 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. --- .../lib/mothership/request/lifecycle/run.ts | 12 +- .../application/prepare-wake.integration.ts | 155 ++++++++++++++++++ .../tasks/application/prepare-wake.ts | 21 ++- 3 files changed, 185 insertions(+), 3 deletions(-) create mode 100644 apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts diff --git a/apps/sim/lib/mothership/request/lifecycle/run.ts b/apps/sim/lib/mothership/request/lifecycle/run.ts index b141f7d7e06..17319f6df73 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -1663,8 +1663,16 @@ async function ensureHeadlessRunIdentity(input: { }, }) return { executionId, runId, cancelled: run.status === 'cancelled' } - } catch { - throw new Error('Chat could not start because its execution record is unavailable') + } catch (error) { + logger.error('Headless run record could not be created', { + chatId: input.chatId, + streamId: input.messageId, + error: getErrorMessage(error), + ...causeForLog(error), + }) + throw new Error('Chat could not start because its execution record is unavailable', { + cause: error, + }) } } diff --git a/apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts b/apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts new file mode 100644 index 00000000000..a6cc2f3b832 --- /dev/null +++ b/apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts @@ -0,0 +1,155 @@ +/** + * How sim answers the worker's retry of a task wake, against real PostgreSQL and Redis: the + * wake route, the chat stream lock, and the run records are production code. Only `after` is + * stubbed, so the background wake turn never starts; each test writes the run record that + * turn would have written instead. + */ +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' + +const { redisUrl, inheritedRedisUrl } = await vi.hoisted(async () => { + const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') + const url = readTestRedisUrl() + const inheritedRedisUrl = process.env.REDIS_URL + /** The real Redis module reads this at import. */ + process.env.REDIS_URL = url + return { redisUrl: url, inheritedRedisUrl } +}) + +vi.mock('next/server', async (original) => ({ + ...(await original()), + after: () => {}, +})) + +import { db } from '@sim/db' +import { copilotChats, copilotRuns, permissions, user, workspace } from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { eq, inArray } from 'drizzle-orm' +import { NextRequest } from 'next/server' +import { closeRedisConnection } from '@/lib/core/config/redis' +import { createRunSegment, updateRunStatus } from '@/lib/mothership/async-runs/repository' +import { chatPubSub } from '@/lib/mothership/chat-status' +import { + acquirePendingChatStream, + releasePendingChatStream, +} from '@/lib/mothership/request/session/abort' +import { POST as wakeRoute } from '@/app/api/mothership/wake/route' + +afterAll(async () => { + chatPubSub?.dispose() + await closeRedisConnection() + if (inheritedRedisUrl === undefined) process.env.REDIS_URL = undefined + else process.env.REDIS_URL = inheritedRedisUrl +}) + +describe.runIf(Boolean(redisUrl))('task wake retries', () => { + const userId = generateId() + const workspaceId = generateId() + const chatIds: string[] = [] + + beforeAll(async () => { + const now = new Date() + await db.insert(user).values({ + id: userId, + name: 'Task wake fixture', + email: `${userId}@task-wake.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(workspace).values({ + id: workspaceId, + name: 'Task wake fixture', + ownerId: userId, + billedAccountUserId: userId, + }) + await db.insert(permissions).values({ + id: generateId(), + userId, + entityType: 'workspace', + entityId: workspaceId, + permissionType: 'admin', + }) + }) + + afterAll(async () => { + if (chatIds.length) { + await db.delete(copilotRuns).where(inArray(copilotRuns.chatId, chatIds)) + await db.delete(copilotChats).where(inArray(copilotChats.id, chatIds)) + } + await db.delete(permissions).where(eq(permissions.userId, userId)) + await db.delete(workspace).where(eq(workspace.id, workspaceId)) + await db.delete(user).where(eq(user.id, userId)) + }) + + async function idleChat() { + const chatId = generateId() + chatIds.push(chatId) + await db.insert(copilotChats).values({ id: chatId, userId, workspaceId, type: 'mothership' }) + return chatId + } + + /** The worker's wake call, as `wakeOnSim` sends it. */ + function wake(chatId: string, runId: string) { + return wakeRoute( + new NextRequest('http://localhost:3000/api/mothership/wake', { + method: 'POST', + headers: { + 'content-type': 'application/json', + 'x-api-key': process.env.INTERNAL_API_SECRET ?? '', + 'x-mothership-user-id': userId, + 'x-mothership-workspace-id': workspaceId, + }, + body: JSON.stringify({ + taskId: generateId(), + runId, + chatId, + userId, + workspaceId, + message: 'Timer elapsed', + status: 'completed', + summary: 'Timer elapsed', + }), + }), + { params: Promise.resolve({}) } + ) + } + + /** The run record the headless wake turn opens under the wake's run ID. */ + function openWakeTurn(chatId: string, runId: string) { + return createRunSegment({ + executionId: generateId(), + chatId, + userId, + workspaceId, + streamId: runId, + requestContext: { source: 'headless_lifecycle' }, + }) + } + + it('answers not-found to a wake whose turn already ended, and leaves the chat free', async () => { + const chatId = await idleChat() + const runId = generateId() + expect((await wake(chatId, runId)).status).toBe(202) + /** The turn ends inside sim without reaching the worker, as a usage-limit refusal does. */ + const turn = await openWakeTurn(chatId, runId) + await updateRunStatus(turn.id, 'complete', { completedAt: new Date() }) + await releasePendingChatStream(chatId, runId) + + /** The worker saw no run under this ID, so it retries the same wake. */ + expect((await wake(chatId, runId)).status).toBe(404) + + const nextTurn = generateId() + expect(await acquirePendingChatStream(chatId, nextTurn, 0)).toBe(true) + await releasePendingChatStream(chatId, nextTurn) + }) + + it('answers busy, never not-found, while the wake turn under that ID still holds the chat', async () => { + const chatId = await idleChat() + const runId = generateId() + expect((await wake(chatId, runId)).status).toBe(202) + await openWakeTurn(chatId, runId) + + expect((await wake(chatId, runId)).status).toBe(409) + await releasePendingChatStream(chatId, runId) + }, 15_000) +}) diff --git a/apps/sim/lib/mothership/tasks/application/prepare-wake.ts b/apps/sim/lib/mothership/tasks/application/prepare-wake.ts index fd0e7a5ad71..768b12b74b0 100644 --- a/apps/sim/lib/mothership/tasks/application/prepare-wake.ts +++ b/apps/sim/lib/mothership/tasks/application/prepare-wake.ts @@ -1,8 +1,13 @@ import { OrchestrationError } from '@/lib/core/orchestration/types' +import { getLatestRunForStream } from '@/lib/mothership/async-runs/repository' import { defineAuthorizedChatUseCase } from '@/lib/mothership/chat/application/authorized-chat-use-case' import { resolveOwnedChatContext } from '@/lib/mothership/chat/application/context' import type { TaskWakeRequest } from '@/lib/mothership/generated/tasks' -import { acquirePendingChatStream } from '@/lib/mothership/request/session/abort' +import { + acquirePendingChatStream, + getLocalChatStreamLease, + releasePendingChatStream, +} from '@/lib/mothership/request/session/abort' import { taskDelegationPolicy } from '@/lib/mothership/tasks/application/context' import { organizationTaskOperations, @@ -39,6 +44,20 @@ export const prepareTaskWake = defineAuthorizedChatUseCase({ if (!(await acquirePendingChatStream(input.chatId, input.runId))) { throw new OrchestrationError('conflict', 'Another stream holds this chat; retry the wake') } + /** + * The worker retries a wake under the same run ID until its own run appears. A turn sim + * already ran under that ID without reaching the worker (a usage-limit refusal) can never + * open again, so answer not-found: the worker dismisses the notification instead of + * retrying forever. Checked under the chat lock, so an in-flight turn still answers busy. + */ + if (await getLatestRunForStream(input.runId)) { + await releasePendingChatStream( + input.chatId, + input.runId, + getLocalChatStreamLease(input.chatId, input.runId) + ) + throw new OrchestrationError('not_found', 'This wake already ran') + } return { accepted: true } as const }, }) From 0bd25e71a13f9029c1d0d3e6897d42f36d21165d Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 15:18:02 -0700 Subject: [PATCH 2/4] 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. --- .../application/prepare-wake.integration.ts | 26 +++++++++++++++++-- .../tasks/application/prepare-wake.ts | 10 +++++-- 2 files changed, 32 insertions(+), 4 deletions(-) diff --git a/apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts b/apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts index a6cc2f3b832..79a5d3891cf 100644 --- a/apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts +++ b/apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts @@ -2,7 +2,8 @@ * How sim answers the worker's retry of a task wake, against real PostgreSQL and Redis: the * wake route, the chat stream lock, and the run records are production code. Only `after` is * stubbed, so the background wake turn never starts; each test writes the run record that - * turn would have written instead. + * turn would have written instead. The run lookup passes through to PostgreSQL unless a test + * makes it fail. */ import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' @@ -20,13 +21,22 @@ vi.mock('next/server', async (original) => ({ after: () => {}, })) +vi.mock('@/lib/mothership/async-runs/repository', async (original) => { + const actual = await original() + return { ...actual, getLatestRunForStream: vi.fn(actual.getLatestRunForStream) } +}) + import { db } from '@sim/db' import { copilotChats, copilotRuns, permissions, user, workspace } from '@sim/db/schema' import { generateId } from '@sim/utils/id' import { eq, inArray } from 'drizzle-orm' import { NextRequest } from 'next/server' import { closeRedisConnection } from '@/lib/core/config/redis' -import { createRunSegment, updateRunStatus } from '@/lib/mothership/async-runs/repository' +import { + createRunSegment, + getLatestRunForStream, + updateRunStatus, +} from '@/lib/mothership/async-runs/repository' import { chatPubSub } from '@/lib/mothership/chat-status' import { acquirePendingChatStream, @@ -152,4 +162,16 @@ describe.runIf(Boolean(redisUrl))('task wake retries', () => { expect((await wake(chatId, runId)).status).toBe(409) await releasePendingChatStream(chatId, runId) }, 15_000) + + it('frees the chat when the run lookup fails after the wake took it', async () => { + const chatId = await idleChat() + const runId = generateId() + vi.mocked(getLatestRunForStream).mockRejectedValueOnce(new Error('statement timeout')) + + expect((await wake(chatId, runId)).status).toBe(500) + + const nextTurn = generateId() + expect(await acquirePendingChatStream(chatId, nextTurn, 0)).toBe(true) + await releasePendingChatStream(chatId, nextTurn) + }) }) diff --git a/apps/sim/lib/mothership/tasks/application/prepare-wake.ts b/apps/sim/lib/mothership/tasks/application/prepare-wake.ts index 768b12b74b0..79ec93a2639 100644 --- a/apps/sim/lib/mothership/tasks/application/prepare-wake.ts +++ b/apps/sim/lib/mothership/tasks/application/prepare-wake.ts @@ -49,14 +49,20 @@ export const prepareTaskWake = defineAuthorizedChatUseCase({ * already ran under that ID without reaching the worker (a usage-limit refusal) can never * open again, so answer not-found: the worker dismisses the notification instead of * retrying forever. Checked under the chat lock, so an in-flight turn still answers busy. + * Any throw here releases the lock just taken, since the wake turn that would release it + * never starts. */ - if (await getLatestRunForStream(input.runId)) { + try { + if (await getLatestRunForStream(input.runId)) { + throw new OrchestrationError('not_found', 'This wake already ran') + } + } catch (error) { await releasePendingChatStream( input.chatId, input.runId, getLocalChatStreamLease(input.chatId, input.runId) ) - throw new OrchestrationError('not_found', 'This wake already ran') + throw error } return { accepted: true } as const }, From 7cb813b1fc12a2d4739928e6fd33fcb36b5e4ad2 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 15:36:23 -0700 Subject: [PATCH 3/4] fix(mothership): release the wake's own chat lease when the run lookup fails --- .../application/prepare-wake.integration.ts | 39 +++++++++++++++++-- .../tasks/application/prepare-wake.ts | 11 ++---- 2 files changed, 40 insertions(+), 10 deletions(-) diff --git a/apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts b/apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts index 79a5d3891cf..5cf7636e426 100644 --- a/apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts +++ b/apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts @@ -12,7 +12,7 @@ const { redisUrl, inheritedRedisUrl } = await vi.hoisted(async () => { const url = readTestRedisUrl() const inheritedRedisUrl = process.env.REDIS_URL /** The real Redis module reads this at import. */ - process.env.REDIS_URL = url + if (url) process.env.REDIS_URL = url return { redisUrl: url, inheritedRedisUrl } }) @@ -31,7 +31,7 @@ import { copilotChats, copilotRuns, permissions, user, workspace } from '@sim/db import { generateId } from '@sim/utils/id' import { eq, inArray } from 'drizzle-orm' import { NextRequest } from 'next/server' -import { closeRedisConnection } from '@/lib/core/config/redis' +import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis' import { createRunSegment, getLatestRunForStream, @@ -40,6 +40,7 @@ import { import { chatPubSub } from '@/lib/mothership/chat-status' import { acquirePendingChatStream, + getLocalChatStreamLease, releasePendingChatStream, } from '@/lib/mothership/request/session/abort' import { POST as wakeRoute } from '@/app/api/mothership/wake/route' @@ -47,7 +48,7 @@ import { POST as wakeRoute } from '@/app/api/mothership/wake/route' afterAll(async () => { chatPubSub?.dispose() await closeRedisConnection() - if (inheritedRedisUrl === undefined) process.env.REDIS_URL = undefined + if (inheritedRedisUrl === undefined) Reflect.deleteProperty(process.env, 'REDIS_URL') else process.env.REDIS_URL = inheritedRedisUrl }) @@ -174,4 +175,36 @@ describe.runIf(Boolean(redisUrl))('task wake retries', () => { expect(await acquirePendingChatStream(chatId, nextTurn, 0)).toBe(true) await releasePendingChatStream(chatId, nextTurn) }) + + it("keeps a retry's chat lock when an earlier wake's slow lookup fails after its own lock expired", async () => { + const chatId = await idleChat() + const runId = generateId() + let lookupStarted!: () => void + const started = new Promise((resolve) => { + lookupStarted = resolve + }) + let failLookup!: () => void + const failed = new Promise((resolve) => { + failLookup = resolve + }) + vi.mocked(getLatestRunForStream).mockImplementationOnce(async () => { + lookupStarted() + await failed + throw new Error('statement timeout') + }) + + const slowWake = wake(chatId, runId) + await started + /** The first wake's lock outlives its TTL while the lookup hangs. */ + const firstLease = getLocalChatStreamLease(chatId, runId) + await getRedisClient()?.del(firstLease?.key ?? '') + expect((await wake(chatId, runId)).status).toBe(202) + + failLookup() + expect((await slowWake).status).toBe(500) + + const nextTurn = generateId() + expect(await acquirePendingChatStream(chatId, nextTurn, 0)).toBe(false) + await releasePendingChatStream(chatId, runId) + }, 15_000) }) diff --git a/apps/sim/lib/mothership/tasks/application/prepare-wake.ts b/apps/sim/lib/mothership/tasks/application/prepare-wake.ts index 79ec93a2639..b855147b34f 100644 --- a/apps/sim/lib/mothership/tasks/application/prepare-wake.ts +++ b/apps/sim/lib/mothership/tasks/application/prepare-wake.ts @@ -44,24 +44,21 @@ export const prepareTaskWake = defineAuthorizedChatUseCase({ if (!(await acquirePendingChatStream(input.chatId, input.runId))) { throw new OrchestrationError('conflict', 'Another stream holds this chat; retry the wake') } + const lease = getLocalChatStreamLease(input.chatId, input.runId) /** * The worker retries a wake under the same run ID until its own run appears. A turn sim * already ran under that ID without reaching the worker (a usage-limit refusal) can never * open again, so answer not-found: the worker dismisses the notification instead of * retrying forever. Checked under the chat lock, so an in-flight turn still answers busy. - * Any throw here releases the lock just taken, since the wake turn that would release it - * never starts. + * Any throw here releases the lock just taken, by its own lease: a slow lookup can outlive + * the lock, and a retry under the same run ID may hold the chat by then. */ try { if (await getLatestRunForStream(input.runId)) { throw new OrchestrationError('not_found', 'This wake already ran') } } catch (error) { - await releasePendingChatStream( - input.chatId, - input.runId, - getLocalChatStreamLease(input.chatId, input.runId) - ) + await releasePendingChatStream(input.chatId, input.runId, lease) throw error } return { accepted: true } as const From 4ed2f71727bfecf8ab2841318278ed26ae89cb4d Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 15:50:00 -0700 Subject: [PATCH 4/4] test(mothership): stub the chat lease getter in the wake unit test mock --- apps/sim/lib/mothership/tasks/application/tasks.test.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/apps/sim/lib/mothership/tasks/application/tasks.test.ts b/apps/sim/lib/mothership/tasks/application/tasks.test.ts index 9f2e8443aa6..d963414b735 100644 --- a/apps/sim/lib/mothership/tasks/application/tasks.test.ts +++ b/apps/sim/lib/mothership/tasks/application/tasks.test.ts @@ -47,6 +47,7 @@ vi.mock('@/lib/workflows/executor/execution-status', () => ({ })) vi.mock('@/lib/mothership/request/session/abort', () => ({ acquirePendingChatStream: hoisted.acquire, + getLocalChatStreamLease: vi.fn(), })) import type { NextRequest } from 'next/server'