Skip to content

Commit 96fd30d

Browse files
authored
fix(mothership): number recovered turns past an unreadable ring and floor the replay TTL (#8501)
* fix(mothership): number recovered turns past an unreadable ring and floor the replay TTL A recovered run whose replay ring read back empty while its seq counter survived (every retained entry unreadable) resumed numbering at 0, writing over the old range and moving the counter backwards past readers' cursors. Resume from the counter whenever no event was recovered. COPILOT_STREAM_TTL_SECONDS below the 20 s chat-lock heartbeat let an idle live buffer expire between refreshes. Floor it at 60 s. * test(mothership): leave a margin for Redis TTL rounding in the live-buffer floor check
1 parent 2091953 commit 96fd30d

5 files changed

Lines changed: 60 additions & 12 deletions

File tree

‎apps/sim/lib/mothership/request/application/recover-stream.test.ts‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,10 @@ vi.mock('@/lib/mothership/request/session/controller-lease', async (original) =>
4646
vi.mock('@/lib/mothership/request/lifecycle/controller-ownership', () => ({
4747
claimRunController: hoisted.claim,
4848
}))
49-
vi.mock('@/lib/mothership/request/session/buffer', () => ({ readEvents: hoisted.events }))
49+
vi.mock('@/lib/mothership/request/session/buffer', () => ({
50+
readEvents: hoisted.events,
51+
getLatestSeq: async () => null,
52+
}))
5053
vi.mock('@/lib/billing/core/billing-attribution', () => billingAttributionMock)
5154
const mockResolveBillingAttribution = billingAttributionMockFns.mockResolveBillingAttribution
5255
const mockResolveOrganizationBillingAttribution =

‎apps/sim/lib/mothership/request/application/recover-stream.ts‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -125,13 +125,13 @@ export const readChatStream = defineAuthorizedChatUseCase({
125125
* A ring that lost its head is treated like an expired one: the controller starts
126126
* from an empty context and re-attaches with an empty receipt, so the worker re-sends
127127
* the whole response and re-hands its parked calls. Rebuilding from the tail would
128-
* persist a truncated turn.
128+
* persist a truncated turn. Without a recovered event, numbering resumes past the
129+
* stream's counter, which outlives unreadable entries still holding earlier seqs.
129130
*/
130131
const ringIntact = startsAtReplayHead(events[0]?.seq)
131132
const recoveredEvents = ringIntact ? events : []
132-
const resumeSeq = ringIntact
133-
? (events.at(-1)?.seq ?? 0)
134-
: ((await getLatestSeq(run.streamId)) ?? 0)
133+
const lastEvent = recoveredEvents.at(-1)
134+
const resumeSeq = lastEvent ? lastEvent.seq : ((await getLatestSeq(run.streamId)) ?? 0)
135135
const requestId = typeof saved?.requestId === 'string' ? saved.requestId : generateId()
136136
const completion = {
137137
chatId,

‎apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts‎

Lines changed: 21 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ const { redisUrl } = await vi.hoisted(async () => {
99
const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure')
1010
const url = readTestRedisUrl()
1111
process.env.REDIS_URL = url
12-
/** The park below outlasts it more than twice over in real time. */
12+
/** Shorter than the lock heartbeat that refreshes a live buffer. */
1313
process.env.COPILOT_STREAM_TTL_SECONDS = '5'
1414
return { redisUrl: url }
1515
})
@@ -61,9 +61,16 @@ describe.runIf(Boolean(redisUrl))('replay buffer lifetime', () => {
6161
expect(await acquirePendingChatStream(chatId, streamId, 0)).toBe(true)
6262
await appendText(streamId, 'before the park')
6363
const [ownerBudgetKey] = getRedisBudgetKeys({ kind: 'copilot_stream', id: streamId })
64-
const chargedBytes = await getRedisClient()!.get(ownerBudgetKey)
65-
/** The counter's own TTL is an hour; shortening it stands in for a park that long. */
66-
await getRedisClient()!.expire(ownerBudgetKey, 5)
64+
const redis = getRedisClient()!
65+
const chargedBytes = await redis.get(ownerBudgetKey)
66+
/** Shortening every TTL to 5 s stands in for a park longer than each; the park outlasts it twice. */
67+
for (const key of [
68+
ownerBudgetKey,
69+
`mothership_stream:${streamId}:events`,
70+
`mothership_stream:${streamId}:seq`,
71+
]) {
72+
await redis.expire(key, 5)
73+
}
6774

6875
vi.useFakeTimers({ toFake: ['Date'] })
6976
const poller = startAbortPoller(streamId, new AbortController(), { chatId, pollMs: 50 })
@@ -81,10 +88,19 @@ describe.runIf(Boolean(redisUrl))('replay buffer lifetime', () => {
8188
expect(await getLatestSeq(streamId)).toBe(1)
8289
expect((await readEvents(streamId, '0')).map((event) => event.seq)).toEqual([1])
8390
expect(chargedBytes).not.toBeNull()
84-
expect(await getRedisClient()!.get(ownerBudgetKey)).toBe(chargedBytes)
91+
expect(await redis.get(ownerBudgetKey)).toBe(chargedBytes)
8592
expect(await appendText(streamId, 'after the park')).toBe(2)
8693
})
8794

95+
it('keeps a live buffer past the heartbeat that refreshes it when the configured TTL is shorter', async () => {
96+
const streamId = generateId()
97+
await appendText(streamId, 'live')
98+
99+
const redis = getRedisClient()!
100+
expect(await redis.ttl(`mothership_stream:${streamId}:events`)).toBeGreaterThanOrEqual(55)
101+
expect(await redis.ttl(`mothership_stream:${streamId}:seq`)).toBeGreaterThanOrEqual(55)
102+
})
103+
88104
it('never re-extends a finished stream’s buffer after its cleanup was scheduled', async () => {
89105
const streamId = generateId()
90106
await appendText(streamId, 'done')

‎apps/sim/lib/mothership/request/session/buffer.ts‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,11 @@ const logger = createLogger('SessionBuffer')
2121

2222
const STREAM_OUTBOX_PREFIX = 'mothership_stream:'
2323
const DEFAULT_TTL_SECONDS = 60 * 60
24+
/**
25+
* Floor for a configured live TTL: three of the 20 s chat-lock heartbeats that refresh
26+
* an idle live buffer, so a parked run cannot expire between refreshes.
27+
*/
28+
const MIN_TTL_SECONDS = 60
2429
const DEFAULT_COMPLETED_TTL_SECONDS = 5 * 60
2530
const DEFAULT_EVENT_LIMIT = 100_000
2631
const RETRY_DELAYS_MS = [0, 50, 150] as const
@@ -67,7 +72,10 @@ export type StreamConfig = {
6772

6873
export function getStreamConfig(): StreamConfig {
6974
return {
70-
ttlSeconds: envNumber(env.COPILOT_STREAM_TTL_SECONDS, DEFAULT_TTL_SECONDS, { min: 1 }),
75+
ttlSeconds: Math.max(
76+
MIN_TTL_SECONDS,
77+
envNumber(env.COPILOT_STREAM_TTL_SECONDS, DEFAULT_TTL_SECONDS, { min: 1 })
78+
),
7179
eventLimit: envNumber(env.COPILOT_STREAM_EVENT_LIMIT, DEFAULT_EVENT_LIMIT, { min: 1 }),
7280
}
7381
}

‎apps/sim/lib/mothership/request/session/stream-recovery.integration.ts‎

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -69,7 +69,7 @@ import {
6969
createProviderToolCallIdentity,
7070
scopeProviderToolCallId,
7171
} from '@/lib/mothership/request/go/tool-call-identity'
72-
import { appendEvents } from '@/lib/mothership/request/session/buffer'
72+
import { appendEvents, readEvents } from '@/lib/mothership/request/session/buffer'
7373
import { chatStreamLockKey } from '@/lib/mothership/request/session/controller-lease'
7474
import { createEvent } from '@/lib/mothership/request/session/event'
7575
import { GET as streamGET } from '@/app/api/copilot/chat/stream/route'
@@ -346,4 +346,25 @@ describe.runIf(Boolean(redisUrl))('recovering a run whose ring lost its head', (
346346
expect(toolIds).toHaveLength(2)
347347
expect(new Set(toolIds).size).toBe(2)
348348
})
349+
350+
it('numbers a recovered turn past a ring whose events are all unreadable', async () => {
351+
const { streamId, runId, frame } = await orphanedRunWithTrimmedRing()
352+
const redis = getRedisClient()!
353+
const eventsKey = `mothership_stream:${streamId}:events`
354+
await redis.del(eventsKey)
355+
await redis.zadd(eventsKey, 1, 'corrupt-1', 2, 'corrupt-2', 3, 'corrupt-3', 4, 'corrupt-4')
356+
worker.replies = {
357+
'/api/mothership': [
358+
frame(1, 'session', { kind: 'start' }),
359+
frame(2, 'text', { channel: 'assistant', text: FULL_TEXT, textOffset: 0 }),
360+
frame(3, 'complete', { status: 'complete', textLength: FULL_TEXT.length }),
361+
],
362+
}
363+
364+
expect(await recoverAndFinish(streamId, runId)).toBe('complete')
365+
366+
const recovered = await readEvents(streamId, '0')
367+
expect(recovered.length).toBeGreaterThan(0)
368+
expect(Math.min(...recovered.map((event) => event.seq))).toBe(5)
369+
})
349370
})

0 commit comments

Comments
 (0)