Skip to content

Commit 25a7d2d

Browse files
committed
fix(mothership): refresh the stream byte counter with its buffer, and never revive a closed buffer
- The chat-lock heartbeat now slides the owner byte counter's TTL along with the replay buffer's. Otherwise, after a park longer than an hour, the counter expired while the ring survived: the ring could grow to about twice its target, and its refunds drained the user counter. - Scheduling a finished stream's cleanup now marks it closed, and the refresh leaves a closed stream alone. A heartbeat still in flight at teardown can no longer re-extend a buffer whose cleanup was already scheduled. - Describe the worker idle timeout in terms of network intermediaries, and state that the delegation TTL reuses the long-running tool watchdog's cap.
1 parent be26984 commit 25a7d2d

5 files changed

Lines changed: 70 additions & 18 deletions

File tree

‎apps/sim/lib/mothership/auth/application-delegation.ts‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,9 @@ import type { DelegatedPrincipal, OrganizationDelegatedPrincipal } from '@sim/au
22
import { TOOL_WATCHDOG_LONG_RUNNING_MS } from '@/lib/mothership/constants'
33

44
/**
5-
* Delegated authority is minted per operation, so it must outlive the longest
6-
* single tool call it authorizes, never the whole run.
5+
* Delegated authority is minted per operation and must outlive the longest single
6+
* tool call it authorizes, never the whole run, so it reuses the long-running tool
7+
* watchdog's cap rather than a lifetime of its own.
78
*/
89
export const COPILOT_APPLICATION_DELEGATION_TTL_MS = TOOL_WATCHDOG_LONG_RUNNING_MS
910

‎apps/sim/lib/mothership/constants.ts‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -14,9 +14,9 @@ export const SIM_AGENT_API_URL =
1414
* How long a worker SSE leg may stay silent before Sim treats the connection as
1515
* lost and re-attaches. The worker writes a keepalive comment whenever a leg has
1616
* been quiet for 10 s (checked every 15 s), independent of model or tool progress,
17-
* so a healthy leg is never silent for more than about 25 s. This must stay under
18-
* the NAT's fixed 350 s idle timeout, which drops a silent connection without
19-
* closing it.
17+
* so a healthy leg is never silent for more than about 25 s. It stays well under
18+
* the idle timeouts of network intermediaries, which can drop a silent connection
19+
* without closing it.
2020
*/
2121
export const WORKER_STREAM_IDLE_TIMEOUT_MS = 120_000
2222

‎apps/sim/lib/mothership/request/go/stream.test.ts‎

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1027,7 +1027,8 @@ describe('copilot go stream helpers', () => {
10271027
})
10281028

10291029
describe('worker stream liveness without a caller deadline', () => {
1030-
const NAT_IDLE_MS = 350_000
1030+
/** Well past the idle timeout, and under common intermediary idle cuts. */
1031+
const INTERMEDIARY_IDLE_MS = 300_000
10311032
const encoder = new TextEncoder()
10321033
const frame = (event: unknown) => encoder.encode(`data: ${JSON.stringify(event)}\n\n`)
10331034
const firstText = createEvent({
@@ -1112,7 +1113,7 @@ describe('copilot go stream helpers', () => {
11121113
expect(context.streamComplete).toBe(true)
11131114
})
11141115

1115-
it('fails a silent leg as a retryable interruption before the NAT drops it', async () => {
1116+
it('fails a silent leg as a retryable interruption before an intermediary drops it', async () => {
11161117
vi.useFakeTimers()
11171118
vi.mocked(fetch).mockResolvedValueOnce(
11181119
new Response(
@@ -1136,13 +1137,13 @@ describe('copilot go stream helpers', () => {
11361137

11371138
await vi.advanceTimersByTimeAsync(60_000)
11381139
expect(state.done).toBe(false)
1139-
await vi.advanceTimersByTimeAsync(NAT_IDLE_MS - 60_000)
1140+
await vi.advanceTimersByTimeAsync(INTERMEDIARY_IDLE_MS - 60_000)
11401141

11411142
expect(state.done).toBe(true)
11421143
expect(state.error).toBeInstanceOf(WorkerStreamInterruptedError)
11431144
})
11441145

1145-
it('fails a worker that never answers as unreachable before the NAT drops it', async () => {
1146+
it('fails a worker that never answers as unreachable before an intermediary drops it', async () => {
11461147
vi.useFakeTimers()
11471148
vi.mocked(fetch).mockImplementationOnce(
11481149
(_url, options) =>
@@ -1163,7 +1164,7 @@ describe('copilot go stream helpers', () => {
11631164
)
11641165
)
11651166

1166-
await vi.advanceTimersByTimeAsync(NAT_IDLE_MS)
1167+
await vi.advanceTimersByTimeAsync(INTERMEDIARY_IDLE_MS)
11671168

11681169
expect(state.done).toBe(true)
11691170
expect(state.error).toBeInstanceOf(WorkerUnreachableError)

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

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ const { redisUrl } = await vi.hoisted(async () => {
1717
import { sleep } from '@sim/utils/helpers'
1818
import { generateId } from '@sim/utils/id'
1919
import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis'
20+
import { getRedisBudgetKeys } from '@/lib/core/redis/byte-budget.server'
2021
import { MothershipStreamV1EventType } from '@/lib/mothership/generated/mothership-stream-v1'
2122
import {
2223
acquirePendingChatStream,
@@ -28,6 +29,8 @@ import {
2829
appendEvents,
2930
getLatestSeq,
3031
readEvents,
32+
refreshBufferTtl,
33+
scheduleBufferCleanup,
3134
} from '@/lib/mothership/request/session/buffer'
3235
import { createEvent } from '@/lib/mothership/request/session/event'
3336
import { checkForReplayGap } from '@/lib/mothership/request/session/recovery'
@@ -52,11 +55,15 @@ describe.runIf(Boolean(redisUrl))('replay buffer lifetime', () => {
5255
await closeRedisConnection()
5356
})
5457

55-
it('keeps a live run’s buffer through a park longer than its TTL', async () => {
58+
it('keeps a live run’s buffer and its byte counter through a park longer than their TTLs', async () => {
5659
const chatId = generateId()
5760
const streamId = generateId()
5861
expect(await acquirePendingChatStream(chatId, streamId, 0)).toBe(true)
5962
await appendText(streamId, 'before the park')
63+
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, 2)
6067

6168
vi.useFakeTimers({ toFake: ['Date'] })
6269
const poller = startAbortPoller(streamId, new AbortController(), { chatId, pollMs: 50 })
@@ -73,9 +80,23 @@ describe.runIf(Boolean(redisUrl))('replay buffer lifetime', () => {
7380

7481
expect(await getLatestSeq(streamId)).toBe(1)
7582
expect((await readEvents(streamId, '0')).map((event) => event.seq)).toEqual([1])
83+
expect(chargedBytes).not.toBeNull()
84+
expect(await getRedisClient()!.get(ownerBudgetKey)).toBe(chargedBytes)
7685
expect(await appendText(streamId, 'after the park')).toBe(2)
7786
})
7887

88+
it('never re-extends a finished stream’s buffer after its cleanup was scheduled', async () => {
89+
const streamId = generateId()
90+
await appendText(streamId, 'done')
91+
await scheduleBufferCleanup(streamId)
92+
93+
await refreshBufferTtl(streamId)
94+
95+
const redis = getRedisClient()!
96+
expect(await redis.ttl(`mothership_stream:${streamId}:events`)).toBeGreaterThan(250)
97+
expect(await redis.ttl(`mothership_stream:${streamId}:seq`)).toBeGreaterThan(250)
98+
})
99+
79100
it('reports a gap to a cursor ahead of a buffer whose numbering restarted', async () => {
80101
const streamId = generateId()
81102
for (let index = 0; index < 5; index++) await appendText(streamId, `part ${index}`)

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

Lines changed: 36 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,11 @@ function getAbortKey(streamId: string) {
4242
return `${STREAM_OUTBOX_PREFIX}${streamId}:abort`
4343
}
4444

45+
/** Marks a stream whose cleanup is scheduled, so a late heartbeat cannot revive it. */
46+
function getClosedKey(streamId: string) {
47+
return `${STREAM_OUTBOX_PREFIX}${streamId}:closed`
48+
}
49+
4550
export type StreamConfig = {
4651
ttlSeconds: number
4752
eventLimit: number
@@ -127,23 +132,46 @@ export async function clearBuffer(streamId: string, operation = 'clear_outbox'):
127132
getEventsKey(streamId),
128133
getSeqKey(streamId),
129134
getAbortKey(streamId),
135+
getClosedKey(streamId),
130136
ownerBudgetKey
131137
)
132138
})
133139
}
134140

135141
/**
136-
* Slides a live stream's replay TTLs without an append. The TTLs otherwise move only
137-
* when an event lands, so a run parked on a long tool call or approval would lose its
138-
* replay history and restart its numbering while it is still running.
142+
* KEYS: [events, seq, ownerBudget, closed]
143+
* ARGV: [ttlSeconds, budgetTtlSeconds]
144+
*/
145+
const REFRESH_BUFFER_TTL_SCRIPT = `
146+
if redis.call('EXISTS', KEYS[4]) == 1 then return 0 end
147+
redis.call('EXPIRE', KEYS[1], ARGV[1])
148+
redis.call('EXPIRE', KEYS[2], ARGV[1])
149+
redis.call('EXPIRE', KEYS[3], ARGV[2])
150+
return 1
151+
`
152+
153+
/**
154+
* Slides a live stream's replay TTLs, and its byte counter's, without an append. They
155+
* otherwise move only when an event lands, so a run parked on a long tool call or
156+
* approval would lose its replay history and restart its numbering while it still
157+
* runs, or keep its history after the counter that accounts for it expired. A stream
158+
* whose cleanup is already scheduled is left to expire.
139159
*/
140160
export async function refreshBufferTtl(streamId: string): Promise<void> {
141161
const { ttlSeconds } = getStreamConfig()
162+
const [ownerBudgetKey] = getRedisBudgetKeys({ kind: 'copilot_stream', id: streamId })
163+
const budgetTtlSeconds = Math.max(getRedisBudgetLimits('copilot_stream').ttlSeconds, ttlSeconds)
142164
await withRedisRetry({ operation: 'refresh_outbox_ttl', streamId }, async (redis) => {
143-
const pipeline = redis.pipeline()
144-
pipeline.expire(getEventsKey(streamId), ttlSeconds)
145-
pipeline.expire(getSeqKey(streamId), ttlSeconds)
146-
await pipeline.exec()
165+
await redis.eval(
166+
REFRESH_BUFFER_TTL_SCRIPT,
167+
4,
168+
getEventsKey(streamId),
169+
getSeqKey(streamId),
170+
ownerBudgetKey,
171+
getClosedKey(streamId),
172+
ttlSeconds,
173+
budgetTtlSeconds
174+
)
147175
})
148176
}
149177

@@ -157,6 +185,7 @@ export async function scheduleBufferCleanup(
157185
pipeline.expire(getEventsKey(streamId), ttlSeconds)
158186
pipeline.expire(getSeqKey(streamId), ttlSeconds)
159187
pipeline.expire(getAbortKey(streamId), ttlSeconds)
188+
pipeline.set(getClosedKey(streamId), '1', 'EX', ttlSeconds)
160189
await pipeline.exec()
161190
})
162191
} catch (error) {

0 commit comments

Comments
 (0)