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..a84bc2bf7c0 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(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 () => { 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) + }) })