From 5b55bd3aa3d904cdc62e8ca70fcc9307915c3ef3 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 18:17:24 -0700 Subject: [PATCH 1/2] 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. --- .../application/recover-stream.test.ts | 5 +++- .../request/application/recover-stream.ts | 8 +++--- .../request/session/buffer-ttl.integration.ts | 26 +++++++++++++++---- .../lib/mothership/request/session/buffer.ts | 10 ++++++- .../session/stream-recovery.integration.ts | 23 +++++++++++++++- 5 files changed, 60 insertions(+), 12 deletions(-) diff --git a/apps/sim/lib/mothership/request/application/recover-stream.test.ts b/apps/sim/lib/mothership/request/application/recover-stream.test.ts index 2a1eedf766b..67009e610e1 100644 --- a/apps/sim/lib/mothership/request/application/recover-stream.test.ts +++ b/apps/sim/lib/mothership/request/application/recover-stream.test.ts @@ -46,7 +46,10 @@ vi.mock('@/lib/mothership/request/session/controller-lease', async (original) => vi.mock('@/lib/mothership/request/lifecycle/controller-ownership', () => ({ claimRunController: hoisted.claim, })) -vi.mock('@/lib/mothership/request/session/buffer', () => ({ readEvents: hoisted.events })) +vi.mock('@/lib/mothership/request/session/buffer', () => ({ + readEvents: hoisted.events, + getLatestSeq: async () => null, +})) vi.mock('@/lib/billing/core/billing-attribution', () => billingAttributionMock) const mockResolveBillingAttribution = billingAttributionMockFns.mockResolveBillingAttribution const mockResolveOrganizationBillingAttribution = diff --git a/apps/sim/lib/mothership/request/application/recover-stream.ts b/apps/sim/lib/mothership/request/application/recover-stream.ts index 61eff64d315..10b1c9ad390 100644 --- a/apps/sim/lib/mothership/request/application/recover-stream.ts +++ b/apps/sim/lib/mothership/request/application/recover-stream.ts @@ -125,13 +125,13 @@ export const readChatStream = defineAuthorizedChatUseCase({ * A ring that lost its head is treated like an expired one: the controller starts * from an empty context and re-attaches with an empty receipt, so the worker re-sends * the whole response and re-hands its parked calls. Rebuilding from the tail would - * persist a truncated turn. + * persist a truncated turn. Without a recovered event, numbering resumes past the + * stream's counter, which outlives unreadable entries still holding earlier seqs. */ const ringIntact = startsAtReplayHead(events[0]?.seq) const recoveredEvents = ringIntact ? events : [] - const resumeSeq = ringIntact - ? (events.at(-1)?.seq ?? 0) - : ((await getLatestSeq(run.streamId)) ?? 0) + const lastEvent = recoveredEvents.at(-1) + const resumeSeq = lastEvent ? lastEvent.seq : ((await getLatestSeq(run.streamId)) ?? 0) const requestId = typeof saved?.requestId === 'string' ? saved.requestId : generateId() const completion = { chatId, diff --git a/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts b/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts index 59d9b49d6a5..85bf013e8f9 100644 --- a/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts +++ b/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts @@ -9,7 +9,7 @@ const { redisUrl } = await vi.hoisted(async () => { const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') const url = readTestRedisUrl() process.env.REDIS_URL = url - /** The park below outlasts it more than twice over in real time. */ + /** Shorter than the lock heartbeat that refreshes a live buffer. */ process.env.COPILOT_STREAM_TTL_SECONDS = '5' return { redisUrl: url } }) @@ -61,9 +61,16 @@ describe.runIf(Boolean(redisUrl))('replay buffer lifetime', () => { expect(await acquirePendingChatStream(chatId, streamId, 0)).toBe(true) await appendText(streamId, 'before the park') const [ownerBudgetKey] = getRedisBudgetKeys({ kind: 'copilot_stream', id: streamId }) - const chargedBytes = await getRedisClient()!.get(ownerBudgetKey) - /** The counter's own TTL is an hour; shortening it stands in for a park that long. */ - await getRedisClient()!.expire(ownerBudgetKey, 5) + const redis = getRedisClient()! + const chargedBytes = await redis.get(ownerBudgetKey) + /** Shortening every TTL to 5 s stands in for a park longer than each; the park outlasts it twice. */ + for (const key of [ + ownerBudgetKey, + `mothership_stream:${streamId}:events`, + `mothership_stream:${streamId}:seq`, + ]) { + await redis.expire(key, 5) + } vi.useFakeTimers({ toFake: ['Date'] }) const poller = startAbortPoller(streamId, new AbortController(), { chatId, pollMs: 50 }) @@ -81,10 +88,19 @@ describe.runIf(Boolean(redisUrl))('replay buffer lifetime', () => { expect(await getLatestSeq(streamId)).toBe(1) expect((await readEvents(streamId, '0')).map((event) => event.seq)).toEqual([1]) expect(chargedBytes).not.toBeNull() - expect(await getRedisClient()!.get(ownerBudgetKey)).toBe(chargedBytes) + expect(await redis.get(ownerBudgetKey)).toBe(chargedBytes) expect(await appendText(streamId, 'after the park')).toBe(2) }) + it('keeps a live buffer past the heartbeat that refreshes it when the configured TTL is shorter', async () => { + const streamId = generateId() + await appendText(streamId, 'live') + + const redis = getRedisClient()! + expect(await redis.ttl(`mothership_stream:${streamId}:events`)).toBeGreaterThanOrEqual(60) + expect(await redis.ttl(`mothership_stream:${streamId}:seq`)).toBeGreaterThanOrEqual(60) + }) + it('never re-extends a finished stream’s buffer after its cleanup was scheduled', async () => { const streamId = generateId() await appendText(streamId, 'done') diff --git a/apps/sim/lib/mothership/request/session/buffer.ts b/apps/sim/lib/mothership/request/session/buffer.ts index b990f543262..883f7ba3d5c 100644 --- a/apps/sim/lib/mothership/request/session/buffer.ts +++ b/apps/sim/lib/mothership/request/session/buffer.ts @@ -21,6 +21,11 @@ const logger = createLogger('SessionBuffer') const STREAM_OUTBOX_PREFIX = 'mothership_stream:' const DEFAULT_TTL_SECONDS = 60 * 60 +/** + * Floor for a configured live TTL: three of the 20 s chat-lock heartbeats that refresh + * an idle live buffer, so a parked run cannot expire between refreshes. + */ +const MIN_TTL_SECONDS = 60 const DEFAULT_COMPLETED_TTL_SECONDS = 5 * 60 const DEFAULT_EVENT_LIMIT = 100_000 const RETRY_DELAYS_MS = [0, 50, 150] as const @@ -67,7 +72,10 @@ export type StreamConfig = { export function getStreamConfig(): StreamConfig { return { - ttlSeconds: envNumber(env.COPILOT_STREAM_TTL_SECONDS, DEFAULT_TTL_SECONDS, { min: 1 }), + ttlSeconds: Math.max( + MIN_TTL_SECONDS, + envNumber(env.COPILOT_STREAM_TTL_SECONDS, DEFAULT_TTL_SECONDS, { min: 1 }) + ), eventLimit: envNumber(env.COPILOT_STREAM_EVENT_LIMIT, DEFAULT_EVENT_LIMIT, { min: 1 }), } } diff --git a/apps/sim/lib/mothership/request/session/stream-recovery.integration.ts b/apps/sim/lib/mothership/request/session/stream-recovery.integration.ts index 6b714c19cb9..1665d56a39e 100644 --- a/apps/sim/lib/mothership/request/session/stream-recovery.integration.ts +++ b/apps/sim/lib/mothership/request/session/stream-recovery.integration.ts @@ -69,7 +69,7 @@ import { createProviderToolCallIdentity, scopeProviderToolCallId, } from '@/lib/mothership/request/go/tool-call-identity' -import { appendEvents } from '@/lib/mothership/request/session/buffer' +import { appendEvents, readEvents } from '@/lib/mothership/request/session/buffer' import { chatStreamLockKey } from '@/lib/mothership/request/session/controller-lease' import { createEvent } from '@/lib/mothership/request/session/event' 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', ( expect(toolIds).toHaveLength(2) expect(new Set(toolIds).size).toBe(2) }) + + it('numbers a recovered turn past a ring whose events are all unreadable', async () => { + const { streamId, runId, frame } = await orphanedRunWithTrimmedRing() + const redis = getRedisClient()! + const eventsKey = `mothership_stream:${streamId}:events` + await redis.del(eventsKey) + await redis.zadd(eventsKey, 1, 'corrupt-1', 2, 'corrupt-2', 3, 'corrupt-3', 4, 'corrupt-4') + worker.replies = { + '/api/mothership': [ + frame(1, 'session', { kind: 'start' }), + frame(2, 'text', { channel: 'assistant', text: FULL_TEXT, textOffset: 0 }), + frame(3, 'complete', { status: 'complete', textLength: FULL_TEXT.length }), + ], + } + + expect(await recoverAndFinish(streamId, runId)).toBe('complete') + + const recovered = await readEvents(streamId, '0') + expect(recovered.length).toBeGreaterThan(0) + expect(Math.min(...recovered.map((event) => event.seq))).toBe(5) + }) }) From 52927955a7f9d1ed0d658d01b19064048b7282e9 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 18:43:30 -0700 Subject: [PATCH 2/2] test(mothership): leave a margin for Redis TTL rounding in the live-buffer floor check --- .../lib/mothership/request/session/buffer-ttl.integration.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts b/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts index 85bf013e8f9..a84bc2bf7c0 100644 --- a/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts +++ b/apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts @@ -97,8 +97,8 @@ describe.runIf(Boolean(redisUrl))('replay buffer lifetime', () => { await appendText(streamId, 'live') const redis = getRedisClient()! - expect(await redis.ttl(`mothership_stream:${streamId}:events`)).toBeGreaterThanOrEqual(60) - expect(await redis.ttl(`mothership_stream:${streamId}:seq`)).toBeGreaterThanOrEqual(60) + expect(await redis.ttl(`mothership_stream:${streamId}:events`)).toBeGreaterThanOrEqual(55) + expect(await redis.ttl(`mothership_stream:${streamId}:seq`)).toBeGreaterThanOrEqual(55) }) it('never re-extends a finished stream’s buffer after its cleanup was scheduled', async () => {