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
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down
8 changes: 4 additions & 4 deletions apps/sim/lib/mothership/request/application/recover-stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
26 changes: 21 additions & 5 deletions apps/sim/lib/mothership/request/session/buffer-ttl.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
})
Expand Down Expand Up @@ -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 })
Expand All @@ -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')
Expand Down
10 changes: 9 additions & 1 deletion apps/sim/lib/mothership/request/session/buffer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 }),
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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)
})
})
Loading