Skip to content

Commit 3752dfe

Browse files
authored
fix(typecheck): recreate missing CLI Transport interface (Gitlawb#1581)
* fix(typecheck): recreate missing CLI Transport interface * fix(transports): implement async close in CLI transports * fix(transports): harden async close cleanup * fix: address async transport close review feedback * test: isolate environment-sensitive suites * fix(transports): drain uploader after close failure * fix(transports): type guard CCR stream events * test: remove unused auto-compact fixture helper
1 parent 94d2a6a commit 3752dfe

17 files changed

Lines changed: 983 additions & 106 deletions

src/bridge/remoteBridgeCore.ts

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -490,7 +490,7 @@ export async function initEnvLessBridgeCore(
490490
flushGate.start()
491491
try {
492492
const seq = transport.getLastSequenceNum()
493-
transport.close()
493+
await transport.close()
494494
transport = await createV2ReplTransport({
495495
sessionUrl: buildCCRv2SdkUrl(fresh.api_base_url, sessionId),
496496
ingressToken: fresh.worker_jwt,
@@ -506,7 +506,7 @@ export async function initEnvLessBridgeCore(
506506
// Teardown fired during the async createV2ReplTransport window.
507507
// Don't wire/connect/schedule — we'd re-arm timers after cancelAll()
508508
// and fire onInboundMessage into a torn-down bridge.
509-
transport.close()
509+
await transport.close()
510510
return
511511
}
512512
wireTransportCallbacks()
@@ -717,7 +717,14 @@ export async function initEnvLessBridgeCore(
717717
}
718718
}
719719

720-
transport.close()
720+
try {
721+
await transport.close()
722+
} catch (err) {
723+
logForDebugging(
724+
`[remote-bridge] Transport close threw during teardown: ${errorMessage(err)}`,
725+
{ level: 'error' },
726+
)
727+
}
721728

722729
const archiveStatus: ArchiveTelemetryStatus =
723730
status === 'no_token'

src/bridge/replBridge.ts

Lines changed: 39 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -626,6 +626,20 @@ export async function initBridgeCore(
626626
}
627627
}
628628

629+
async function closeTransportBestEffort(
630+
transportToClose: ReplBridgeTransport,
631+
reason: string,
632+
): Promise<void> {
633+
try {
634+
await transportToClose.close()
635+
} catch (err) {
636+
logForDebugging(
637+
`[bridge:repl] Transport close threw ${reason}: ${errorMessage(err)}`,
638+
{ level: 'error' },
639+
)
640+
}
641+
}
642+
629643
async function doReconnect(): Promise<boolean> {
630644
environmentRecreations++
631645
// Invalidate any in-flight v2 handshake — the environment is being
@@ -652,7 +666,7 @@ export async function initBridgeCore(
652666
if (seq > lastTransportSequenceNum) {
653667
lastTransportSequenceNum = seq
654668
}
655-
transport.close()
669+
await closeTransportBestEffort(transport, 'during reconnect')
656670
transport = null
657671
}
658672
// Transport is gone — wake the poll loop out of its at-capacity
@@ -1055,7 +1069,7 @@ export async function initBridgeCore(
10551069
// forwarding prompts → ~25-min dead window observed in daemon logs.
10561070
// Kill the transport + work state so isAtCapacity()=false; the loop
10571071
// fast-polls and picks up the server's re-dispatched work in seconds.
1058-
onHeartbeatFatal: (err: BridgeFatalError) => {
1072+
onHeartbeatFatal: async (err: BridgeFatalError) => {
10591073
logForDebugging(
10601074
`[bridge:repl] heartbeatWork fatal (status=${err.status}) — tearing down work item for fast re-dispatch`,
10611075
)
@@ -1064,7 +1078,7 @@ export async function initBridgeCore(
10641078
if (seq > lastTransportSequenceNum) {
10651079
lastTransportSequenceNum = seq
10661080
}
1067-
transport.close()
1081+
await closeTransportBestEffort(transport, 'after heartbeat fatal')
10681082
transport = null
10691083
}
10701084
flushGate.drop()
@@ -1094,7 +1108,7 @@ export async function initBridgeCore(
10941108
}
10951109
return { environmentId, environmentSecret }
10961110
},
1097-
onWorkReceived: (
1111+
onWorkReceived: async (
10981112
workSessionId: string,
10991113
ingressToken: string,
11001114
workId: string,
@@ -1199,7 +1213,10 @@ export async function initBridgeCore(
11991213
if (oldSeq > lastTransportSequenceNum) {
12001214
lastTransportSequenceNum = oldSeq
12011215
}
1202-
oldTransport.close()
1216+
await closeTransportBestEffort(
1217+
oldTransport,
1218+
'while replacing work transport',
1219+
)
12031220
}
12041221
// Reset flush state — the old flush (if any) is no longer relevant.
12051222
// Preserve pending messages so they're drained after the new
@@ -1410,13 +1427,16 @@ export async function initBridgeCore(
14101427
sessionId: workSessionId,
14111428
initialSequenceNum: lastTransportSequenceNum,
14121429
}).then(
1413-
t => {
1430+
async t => {
14141431
// Teardown started while registerWorker was in flight. Teardown
14151432
// saw transport === null and skipped close(); installing now
14161433
// would leak CCRClient heartbeat timers and reset
14171434
// teardownStarted via wireTransport's side effects.
14181435
if (pollController.signal.aborted) {
1419-
t.close()
1436+
await closeTransportBestEffort(
1437+
t,
1438+
'while discarding aborted CCR v2 transport',
1439+
)
14201440
return
14211441
}
14221442
// onWorkReceived may have fired again while registerWorker()
@@ -1429,7 +1449,10 @@ export async function initBridgeCore(
14291449
logForDebugging(
14301450
`[bridge:repl] CCR v2: discarding stale handshake gen=${thisGen} current=${v2Generation}`,
14311451
)
1432-
t.close()
1452+
await closeTransportBestEffort(
1453+
t,
1454+
'while discarding stale CCR v2 transport',
1455+
)
14331456
return
14341457
}
14351458
wireTransport(t)
@@ -1668,8 +1691,10 @@ export async function initBridgeCore(
16681691
// log their own success/failure internally.
16691692
await Promise.all([stopWorkP, archiveSession(currentSessionId)])
16701693

1671-
teardownTransport?.close()
1672-
logForDebugging('[bridge:repl] Teardown: transport closed')
1694+
if (teardownTransport) {
1695+
await closeTransportBestEffort(teardownTransport, 'during teardown')
1696+
logForDebugging('[bridge:repl] Teardown: transport closed')
1697+
}
16731698

16741699
await api.deregisterEnvironment(environmentId).catch((err: unknown) => {
16751700
logForDebugging(
@@ -1892,7 +1917,7 @@ async function startWorkPollLoop({
18921917
ingressToken: string,
18931918
workId: string,
18941919
useCodeSessions: boolean,
1895-
) => void
1920+
) => Promise<void>
18961921
/** Called when the environment has been deleted. Returns new credentials or null. */
18971922
onEnvironmentLost?: () => Promise<{
18981923
environmentId: string
@@ -1935,7 +1960,7 @@ async function startWorkPollLoop({
19351960
* ~10-minute dead window before recovery). When omitted, falls back to
19361961
* the backoff sleep to avoid a tight poll+heartbeat loop.
19371962
*/
1938-
onHeartbeatFatal?: (err: BridgeFatalError) => void
1963+
onHeartbeatFatal?: (err: BridgeFatalError) => Promise<void>
19391964
}): Promise<void> {
19401965
const MAX_ENVIRONMENT_RECREATIONS = 3
19411966

@@ -2063,7 +2088,7 @@ async function startWorkPollLoop({
20632088
// for the server's re-dispatched work item. Without
20642089
// the hook, backoff to avoid tight poll+heartbeat loop.
20652090
if (onHeartbeatFatal) {
2066-
onHeartbeatFatal(err)
2091+
await onHeartbeatFatal(err)
20672092
logForDebugging(
20682093
`[bridge:repl:heartbeat] Fatal (status=${err.status}), work state cleared — fast-polling for re-dispatch`,
20692094
)
@@ -2197,7 +2222,7 @@ async function startWorkPollLoop({
21972222
continue
21982223
}
21992224

2200-
onWorkReceived(
2225+
await onWorkReceived(
22012226
workSessionId,
22022227
secret.session_ingress_token,
22032228
work.id,

src/bridge/replBridgeTransport.ts

Lines changed: 30 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ import { registerWorker } from './workSecret.js'
2323
export type ReplBridgeTransport = {
2424
write(message: StdoutMessage): Promise<void>
2525
writeBatch(messages: StdoutMessage[]): Promise<void>
26-
close(): void
26+
close(): Promise<void>
2727
isConnectedStatus(): boolean
2828
getStateLabel(): string
2929
setOnData(callback: (data: string) => void): void
@@ -210,20 +210,8 @@ export async function createV2ReplTransport(opts: {
210210
logForDebugging(
211211
'[bridge:repl] CCR v2: epoch superseded (409) — closing for poll-loop recovery',
212212
)
213-
// Close resources in a try block so the throw always executes.
214-
// If ccr.close() or sse.close() throw, we still need to unwind
215-
// the caller (request()) — otherwise handleEpochMismatch's `never`
216-
// return type is violated at runtime and control falls through.
217-
try {
218-
ccr.close()
219-
sse.close()
220-
onCloseCb?.(4090)
221-
} catch (closeErr: unknown) {
222-
logForDebugging(
223-
`[bridge:repl] CCR v2: error during epoch-mismatch cleanup: ${errorMessage(closeErr)}`,
224-
{ level: 'error' },
225-
)
226-
}
213+
closeResourcesBestEffort('epoch-mismatch cleanup')
214+
onCloseCb?.(4090)
227215
// Don't return — the calling request() code continues after the 409
228216
// branch, so callers see the logged warning and a false return. We
229217
// throw to unwind; the uploaders catch it as a send failure.
@@ -267,6 +255,30 @@ export async function createV2ReplTransport(opts: {
267255
let ccrInitialized = false
268256
let closed = false
269257

258+
function closeResources(): Promise<void> {
259+
closed = true
260+
return Promise.all([ccr.close(), sse.close()]).then(() => {})
261+
}
262+
263+
function closeResourcesBestEffort(reason: string): void {
264+
void closeResources().catch((closeErr: unknown) => {
265+
logForDebugging(
266+
`[bridge:repl] CCR v2: error during ${reason}: ${errorMessage(closeErr)}`,
267+
{ level: 'error' },
268+
)
269+
})
270+
}
271+
272+
function closeCcrBestEffort(reason: string): void {
273+
closed = true
274+
void ccr.close().catch((closeErr: unknown) => {
275+
logForDebugging(
276+
`[bridge:repl] CCR v2: error during ${reason}: ${errorMessage(closeErr)}`,
277+
{ level: 'error' },
278+
)
279+
})
280+
}
281+
270282
return {
271283
write(msg) {
272284
return ccr.writeEvent(msg)
@@ -281,11 +293,7 @@ export async function createV2ReplTransport(opts: {
281293
await ccr.writeEvent(m)
282294
}
283295
},
284-
close() {
285-
closed = true
286-
ccr.close()
287-
sse.close()
288-
},
296+
close: closeResources,
289297
isConnectedStatus() {
290298
// Write-readiness, not read-readiness — replBridge checks this
291299
// before calling writeBatch. SSE open state is orthogonal.
@@ -309,7 +317,7 @@ export async function createV2ReplTransport(opts: {
309317
// heartbeat timer before notifying replBridge. (sse.close() doesn't
310318
// invoke this, so the epoch-mismatch path above isn't double-firing.)
311319
sse.setOnClose(code => {
312-
ccr.close()
320+
closeCcrBestEffort('SSE close cleanup')
313321
cb(code ?? 4092)
314322
})
315323
},
@@ -360,8 +368,7 @@ export async function createV2ReplTransport(opts: {
360368
// so the poll loop can retry on the next work dispatch.
361369
// Without this callback, replBridge never learns the transport
362370
// failed to initialize and sits with transport === null forever.
363-
ccr.close()
364-
sse.close()
371+
closeResourcesBestEffort('initialize failure cleanup')
365372
onCloseCb?.(4091) // 4091 = init failure, distinguishable from 4090 epoch mismatch
366373
},
367374
)
Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,99 @@
1+
import { expect, test } from 'bun:test'
2+
import { BridgeFatalError } from './bridgeApi.js'
3+
import { DEFAULT_POLL_CONFIG } from './pollConfigDefaults.js'
4+
import { _startWorkPollLoopForTesting } from './replBridge.js'
5+
import type { ReplBridgeTransport } from './replBridgeTransport.js'
6+
import type { BridgeApiClient } from './types.js'
7+
8+
type AssertPromiseVoid<T extends Promise<void>> = T
9+
10+
type _CloseReturnsPromise = AssertPromiseVoid<
11+
ReturnType<ReplBridgeTransport['close']>
12+
>
13+
14+
test('work poll loop waits for heartbeat-fatal cleanup before fast-polling again', async () => {
15+
const abort = new AbortController()
16+
let pollCount = 0
17+
let resolveCleanup: (() => void) | undefined
18+
const cleanupStarted = deferred<void>()
19+
const secondPollStarted = deferred<void>()
20+
const cleanupReleased = new Promise<void>(resolve => {
21+
resolveCleanup = resolve
22+
})
23+
24+
const api = {
25+
pollForWork: async () => {
26+
pollCount += 1
27+
if (pollCount === 2) {
28+
secondPollStarted.resolve()
29+
abort.abort()
30+
}
31+
return null
32+
},
33+
heartbeatWork: async () => {
34+
throw new BridgeFatalError('work item gone', 404)
35+
},
36+
} as unknown as BridgeApiClient
37+
38+
let atCapacity = true
39+
const loop = _startWorkPollLoopForTesting({
40+
api,
41+
getCredentials: () => ({
42+
environmentId: 'env-1',
43+
environmentSecret: 'secret-1',
44+
}),
45+
signal: abort.signal,
46+
isAtCapacity: () => atCapacity,
47+
capacitySignal: createCapacitySignal,
48+
getHeartbeatInfo: () => ({
49+
environmentId: 'env-1',
50+
workId: 'work-1',
51+
sessionToken: 'token-1',
52+
}),
53+
getPollIntervalConfig: () => ({
54+
...DEFAULT_POLL_CONFIG,
55+
poll_interval_ms_not_at_capacity: 1,
56+
poll_interval_ms_at_capacity: 10_000,
57+
non_exclusive_heartbeat_interval_ms: 1,
58+
reclaim_older_than_ms: 0,
59+
}),
60+
onWorkReceived: async () => {},
61+
onHeartbeatFatal: async () => {
62+
atCapacity = false
63+
cleanupStarted.resolve()
64+
await cleanupReleased
65+
},
66+
})
67+
68+
await cleanupStarted.promise
69+
await Promise.resolve()
70+
await Promise.resolve()
71+
expect(pollCount).toBe(1)
72+
73+
resolveCleanup?.()
74+
await secondPollStarted.promise
75+
await loop
76+
expect(pollCount).toBe(2)
77+
})
78+
79+
function createCapacitySignal(): {
80+
signal: AbortSignal
81+
cleanup: () => void
82+
} {
83+
const controller = new AbortController()
84+
return {
85+
signal: controller.signal,
86+
cleanup: () => {},
87+
}
88+
}
89+
90+
function deferred<T>(): {
91+
promise: Promise<T>
92+
resolve: (value: T | PromiseLike<T>) => void
93+
} {
94+
let resolve!: (value: T | PromiseLike<T>) => void
95+
const promise = new Promise<T>(res => {
96+
resolve = res
97+
})
98+
return { promise, resolve }
99+
}

src/cli/remoteIO.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -244,12 +244,12 @@ export class RemoteIO extends StructuredIO {
244244
/**
245245
* Clean up connections gracefully
246246
*/
247-
close(): void {
247+
async close(): Promise<void> {
248248
if (this.keepAliveTimer) {
249249
clearInterval(this.keepAliveTimer)
250250
this.keepAliveTimer = null
251251
}
252-
this.transport.close()
252+
await this.transport.close()
253253
this.inputStream.end()
254254
}
255255
}

0 commit comments

Comments
 (0)