Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 10 additions & 2 deletions apps/sim/lib/mothership/request/lifecycle/run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Comment thread
waleedlatif1 marked this conversation as resolved.
})
throw new Error('Chat could not start because its execution record is unavailable', {
cause: error,
})
}
}

Expand Down
210 changes: 210 additions & 0 deletions apps/sim/lib/mothership/tasks/application/prepare-wake.integration.ts
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'
Comment thread
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)
})
24 changes: 23 additions & 1 deletion apps/sim/lib/mothership/tasks/application/prepare-wake.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -39,6 +44,23 @@ 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, 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, lease)
throw error
}
return { accepted: true } as const
},
})
Expand Down
1 change: 1 addition & 0 deletions apps/sim/lib/mothership/tasks/application/tasks.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down
Loading