-
Notifications
You must be signed in to change notification settings - Fork 3.8k
fix(mothership): answer a retried task wake whose turn already ran #8484
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
+244
−3
Merged
Changes from all commits
Commits
Show all changes
4 commits
Select commit
Hold shift + click to select a range
85986a7
fix(mothership): answer a retried task wake whose turn already ran
waleedlatif1 0bd25e7
fix(mothership): release the wake's chat lock when the run lookup fails
waleedlatif1 7cb813b
fix(mothership): release the wake's own chat lease when the run looku…
waleedlatif1 4ed2f71
test(mothership): stub the chat lease getter in the wake unit test mock
waleedlatif1 File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
210 changes: 210 additions & 0 deletions
210
apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,210 @@ | ||
| /** | ||
| * 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. The run lookup passes through to PostgreSQL unless a test | ||
| * makes it fail. | ||
| */ | ||
| 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. */ | ||
| if (url) process.env.REDIS_URL = url | ||
| return { redisUrl: url, inheritedRedisUrl } | ||
| }) | ||
|
|
||
| vi.mock('next/server', async (original) => ({ | ||
| ...(await original<typeof import('next/server')>()), | ||
| after: () => {}, | ||
| })) | ||
|
|
||
| vi.mock('@/lib/mothership/async-runs/repository', async (original) => { | ||
| const actual = await original<typeof import('@/lib/mothership/async-runs/repository')>() | ||
| 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' | ||
|
waleedlatif1 marked this conversation as resolved.
|
||
| import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis' | ||
| import { | ||
| createRunSegment, | ||
| getLatestRunForStream, | ||
| updateRunStatus, | ||
| } from '@/lib/mothership/async-runs/repository' | ||
| 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' | ||
|
|
||
| afterAll(async () => { | ||
| chatPubSub?.dispose() | ||
| await closeRedisConnection() | ||
| if (inheritedRedisUrl === undefined) Reflect.deleteProperty(process.env, 'REDIS_URL') | ||
| 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) | ||
|
|
||
| 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) | ||
| }) | ||
|
|
||
| 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<void>((resolve) => { | ||
| lookupStarted = resolve | ||
| }) | ||
| let failLookup!: () => void | ||
| const failed = new Promise<void>((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) | ||
| }) | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.