mirror of
https://github.com/stablyai/orca.git
synced 2026-10-03 00:02:19 +00:00
Dormant seam for the Phase 3 port of #16741 (design D9.3): the shared
contracts later T1/T2 slices build on. No production module imports any of
it yet; only tests do. Existing behavior is unchanged: the call-queue hold
check is a no-op until something calls holdIdleSelectors.
Adds:
- RelayWorkAdmission + relay-work-drain-contract: closes relay handler
admission and waits for in-flight work, still admitting the cleanup
requests/notifications (cancel, acks, unwatch) that let that work finish.
- TransportPublicationDrain: counts local write completions against one
transport and fails sticky once that transport is replaced.
- pty-ownership-transfer-release-gate: exact-match opt-in
(ORCA_ENABLE_PTY_OWNERSHIP_TRANSFER_MUTATION=1) that SSH core uses to keep
live transfer off (D9.2).
- RuntimeRpcCallQueuePool.holdIdleSelectors + RuntimeRpcCallQueueBusyError:
atomically fence idle selectors while routing changes.
Porting note (source: #16741 head a68b6f3531, merge-base 277c289bd4; T1 per
the #24137 track manifest):
- Taken verbatim: pty-ownership-transfer-release-gate(.test),
runtime-rpc-call-queue delta (+ retirement test), relay-work-admission
API and behavior, transport-publication-drain API and behavior.
- Adapted: relay-work-drain-contract omits relay.reset,
relay.recoverPreparedReset (T3) and relay.networkTunnel.close/.frame
(T4); those contracts do not exist on this base and join the sets when
their tracks land. Promise.withResolvers replaced with plain promises in
relay-reachable code, because the legacy relay bundle still targets
node18 hosts. RequestContext.transportGeneration is not added here (it
belongs with the dispatcher slice), so the admission test drops it.
Added transport-publication-drain.test.ts, since #16741 covered it only
through relay-producer-publication-drain.
- Dropped: the pty-source-credit-contract/-validation/-record/ledger-test
deltas. They only carry the ownershipTransfer output envelope, which is
live-transfer engine (T7, deferred per D9.2), so the credit contract on
this base is already the seam T2 needs. The post-slice
assertPtySourceSpan re-validation came with that envelope and is
dropped with it.
- Not in this slice (next T1 slices, which are also the first consumers):
dispatcher wiring (dispatcher-client-state / rpc-routing /
producer-transport / notification-publication +
RequestContext.transportGeneration) that constructs RelayWorkAdmission;
relay-producer-publication-drain on TransportPublicationDrain; fs/git
stream shutdown; ws-transport + node-websocket-lifecycle +
websocket-transport-limits; renderer flow-controller deltas. Later
consumers: T2 ssh-pty-preparation-admission (drain contract),
ssh-pty-source-preparation-drain (release gate); T4 tunnels
(publication drain); T5 runtime-environment-call-queue (selector hold).
No wire change: no new RPC fields, methods, or opcodes.
Co-authored-by: m4air <m4air@m4airs-Air.localdomain>
This commit is contained in:
@@ -0,0 +1,223 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import type { RequestContext } from './dispatcher-contract'
|
||||
import { RelayWorkAdmission } from './relay-work-admission'
|
||||
import { SKILL_SSH_RELAY_CANCEL_UPLOAD_METHOD } from '../shared/skill-ssh-relay-contract'
|
||||
|
||||
function context(): RequestContext {
|
||||
return { clientId: 1, isStale: () => false }
|
||||
}
|
||||
|
||||
function pendingOperation() {
|
||||
let resolve!: () => void
|
||||
const promise = new Promise<void>((settle) => {
|
||||
resolve = settle
|
||||
})
|
||||
return { promise, resolve }
|
||||
}
|
||||
|
||||
async function nextTurn(): Promise<void> {
|
||||
await new Promise<void>((resolve) => setImmediate(resolve))
|
||||
}
|
||||
|
||||
describe('RelayWorkAdmission', () => {
|
||||
it('starts admitted operations synchronously and preserves their result', async () => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
const operation = vi.fn(() => 'started')
|
||||
const result = admission.run('fs.writeFile', context(), operation)
|
||||
expect(operation).toHaveBeenCalledOnce()
|
||||
await expect(result).resolves.toBe('started')
|
||||
})
|
||||
|
||||
it('closes admission synchronously and drains existing work despite abort and disconnect', async () => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
const controller = new AbortController()
|
||||
let disconnected = false
|
||||
const mutation = pendingOperation()
|
||||
const running = admission.run(
|
||||
'fs.writeFile',
|
||||
{ ...context(), signal: controller.signal, isStale: () => disconnected },
|
||||
() => mutation.promise
|
||||
)
|
||||
const drained = vi.fn()
|
||||
const drain = admission.beginDrain().then(drained)
|
||||
const rejectedOperation = vi.fn()
|
||||
await expect(admission.run('fs.writeFile', context(), rejectedOperation)).rejects.toThrow(
|
||||
'relay_work_admission_closed'
|
||||
)
|
||||
expect(rejectedOperation).not.toHaveBeenCalled()
|
||||
controller.abort()
|
||||
disconnected = true
|
||||
await nextTurn()
|
||||
expect(drained).not.toHaveBeenCalled()
|
||||
mutation.resolve()
|
||||
await running
|
||||
await drain
|
||||
expect(drained).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it.each([
|
||||
'relay.status',
|
||||
'fs.unwatchAndWait',
|
||||
'agent.cancelExec',
|
||||
SKILL_SSH_RELAY_CANCEL_UPLOAD_METHOD
|
||||
])('admits cleanup request %s while draining', async (method) => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
const mutation = pendingOperation()
|
||||
const running = admission.run('git.diff', context(), () => mutation.promise)
|
||||
const drain = admission.beginDrain()
|
||||
await expect(
|
||||
admission.run(method, context(), () => {
|
||||
mutation.resolve()
|
||||
return 'acknowledged'
|
||||
})
|
||||
).resolves.toBe('acknowledged')
|
||||
await running
|
||||
await drain
|
||||
})
|
||||
|
||||
it.each([
|
||||
'rpc.cancel',
|
||||
'git.responseAck',
|
||||
'git.cancelResponseStream',
|
||||
'fs.streamAck',
|
||||
'fs.cancelStream',
|
||||
'fs.unwatch',
|
||||
'pty.ackData',
|
||||
'pty.setDeliveryPaused'
|
||||
])('admits cleanup notification %s while draining', async (method) => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
const mutation = pendingOperation()
|
||||
const running = admission.run('git.diff', context(), () => mutation.promise)
|
||||
const drain = admission.beginDrain()
|
||||
const acknowledge = vi.fn(() => mutation.resolve())
|
||||
admission.runNotification(method, context(), acknowledge)
|
||||
expect(acknowledge).toHaveBeenCalledOnce()
|
||||
await running
|
||||
await drain
|
||||
await expect(admission.run('git.diff', context(), () => {})).rejects.toThrow(
|
||||
'relay_work_admission_closed'
|
||||
)
|
||||
})
|
||||
|
||||
it('excludes the exact reset request without excluding another request from the same client', async () => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
const mutation = pendingOperation()
|
||||
const running = admission.run('fs.writeFile', context(), () => mutation.promise)
|
||||
const resetContext = context()
|
||||
const drained = vi.fn()
|
||||
const reset = admission.run('relay.reset', resetContext, async () => {
|
||||
await admission.beginDrain(resetContext)
|
||||
drained()
|
||||
})
|
||||
await nextTurn()
|
||||
expect(drained).not.toHaveBeenCalled()
|
||||
mutation.resolve()
|
||||
await running
|
||||
await reset
|
||||
expect(drained).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('rejects a forged context without closing admission', async () => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
const actualContext = context()
|
||||
await admission.run('relay.reset', actualContext, async () => {
|
||||
await expect(admission.beginDrain({ ...actualContext })).rejects.toThrow(
|
||||
'relay_work_drain_context_not_active'
|
||||
)
|
||||
await expect(admission.run('fs.writeFile', context(), () => 'open')).resolves.toBe('open')
|
||||
})
|
||||
})
|
||||
|
||||
it('rejects an already-settled request context without closing admission', async () => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
const oldContext = context()
|
||||
await admission.run('relay.reset', oldContext, () => {})
|
||||
await expect(admission.beginDrain(oldContext)).rejects.toThrow(
|
||||
'relay_work_drain_context_not_active'
|
||||
)
|
||||
await expect(admission.run('fs.writeFile', context(), () => 'open')).resolves.toBe('open')
|
||||
})
|
||||
|
||||
it('waits for reset work when the host lifecycle supplies no exclusion', async () => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
const mutation = pendingOperation()
|
||||
const reset = admission.run('relay.reset', context(), () => mutation.promise)
|
||||
const drained = vi.fn()
|
||||
const drain = admission.beginDrain().then(drained)
|
||||
await nextTurn()
|
||||
expect(drained).not.toHaveBeenCalled()
|
||||
mutation.resolve()
|
||||
await reset
|
||||
await drain
|
||||
expect(drained).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('drops new mutation notifications without invoking their handlers during drain', async () => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
await admission.beginDrain()
|
||||
const mutate = vi.fn()
|
||||
admission.runNotification('pty.write', context(), mutate)
|
||||
expect(mutate).not.toHaveBeenCalled()
|
||||
await expect(admission.run('pty.ackData', context(), mutate)).rejects.toThrow(
|
||||
'relay_work_admission_closed'
|
||||
)
|
||||
})
|
||||
|
||||
it('tracks asynchronous notifications through settlement while draining', async () => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
const mutation = pendingOperation()
|
||||
const operation = vi.fn(async () => mutation.promise)
|
||||
admission.runNotification('pty.write', context(), operation)
|
||||
expect(operation).toHaveBeenCalledOnce()
|
||||
const drained = vi.fn()
|
||||
const drain = admission.beginDrain().then(drained)
|
||||
await nextTurn()
|
||||
expect(drained).not.toHaveBeenCalled()
|
||||
mutation.resolve()
|
||||
await drain
|
||||
expect(drained).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('releases tracking after synchronous notification failure', async () => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
expect(() =>
|
||||
admission.runNotification('pty.write', context(), () => {
|
||||
throw new Error('notification_failed')
|
||||
})
|
||||
).toThrow('notification_failed')
|
||||
await admission.beginDrain()
|
||||
const mutate = vi.fn()
|
||||
admission.runNotification('pty.write', context(), mutate)
|
||||
expect(mutate).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('settles concurrent drains after an admitted operation fails without reopening', async () => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
const mutation = pendingOperation()
|
||||
const running = admission.run('fs.writeFile', context(), async () => {
|
||||
await mutation.promise
|
||||
throw new Error('write_failed')
|
||||
})
|
||||
const failure = expect(running).rejects.toThrow('write_failed')
|
||||
const firstDrain = admission.beginDrain()
|
||||
const secondDrain = admission.beginDrain()
|
||||
mutation.resolve()
|
||||
await failure
|
||||
await Promise.all([firstDrain, secondDrain])
|
||||
await expect(admission.run('fs.writeFile', context(), () => {})).rejects.toThrow(
|
||||
'relay_work_admission_closed'
|
||||
)
|
||||
})
|
||||
|
||||
it('turns synchronous operation failures into rejected promises and releases tracking', async () => {
|
||||
const admission = new RelayWorkAdmission()
|
||||
const result = admission.run('fs.writeFile', context(), () => {
|
||||
throw new Error('write_failed')
|
||||
})
|
||||
await expect(result).rejects.toThrow('write_failed')
|
||||
await admission.beginDrain()
|
||||
await expect(admission.run('fs.writeFile', context(), () => {})).rejects.toThrow(
|
||||
'relay_work_admission_closed'
|
||||
)
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,84 @@
|
||||
import type { RequestContext } from './dispatcher-contract'
|
||||
import { allowsRelayWorkDuringDrain } from '../shared/relay-work-drain-contract'
|
||||
|
||||
/** Handler drain is not proof of physical process exit or detached producer cleanup. */
|
||||
export class RelayWorkAdmission {
|
||||
private draining = false
|
||||
private readonly active = new Map<object, { context: RequestContext; done: Promise<void> }>()
|
||||
|
||||
allows(method: string, notification = false): boolean {
|
||||
return !this.draining || allowsRelayWorkDuringDrain(method, notification)
|
||||
}
|
||||
|
||||
async run<T>(
|
||||
method: string,
|
||||
context: RequestContext,
|
||||
operation: () => T | Promise<T>
|
||||
): Promise<T> {
|
||||
if (!this.allows(method)) {
|
||||
throw new Error('relay_work_admission_closed')
|
||||
}
|
||||
return this.track(context, operation)
|
||||
}
|
||||
|
||||
runNotification(method: string, context: RequestContext, operation: () => void): void {
|
||||
if (!this.allows(method, true)) {
|
||||
return
|
||||
}
|
||||
const finish = this.trackEntry(context)
|
||||
try {
|
||||
const result = operation()
|
||||
void Promise.resolve(result).then(finish, (error) => {
|
||||
finish()
|
||||
process.stderr.write(`[relay] Notification handler failed: ${String(error)}\n`)
|
||||
})
|
||||
} catch (error) {
|
||||
finish()
|
||||
throw error
|
||||
}
|
||||
}
|
||||
|
||||
assertActiveContext(context: RequestContext): void {
|
||||
if (![...this.active.values()].some((entry) => entry.context === context)) {
|
||||
throw new Error('relay_work_drain_context_not_active')
|
||||
}
|
||||
}
|
||||
|
||||
/** Only the actual initiating request context may be excluded from its own drain. */
|
||||
async beginDrain(exclude?: RequestContext): Promise<void> {
|
||||
if (exclude) {
|
||||
this.assertActiveContext(exclude)
|
||||
}
|
||||
this.draining = true
|
||||
for (;;) {
|
||||
const pending = [...this.active.values()].filter((entry) => entry.context !== exclude)
|
||||
if (pending.length === 0) {
|
||||
return
|
||||
}
|
||||
await Promise.all(pending.map((entry) => entry.done))
|
||||
}
|
||||
}
|
||||
|
||||
private async track<T>(context: RequestContext, operation: () => T | Promise<T>): Promise<T> {
|
||||
const finish = this.trackEntry(context)
|
||||
try {
|
||||
return await operation()
|
||||
} finally {
|
||||
finish()
|
||||
}
|
||||
}
|
||||
|
||||
// Why: no Promise.withResolvers — the legacy relay bundle still targets Node 18 hosts.
|
||||
private trackEntry(context: RequestContext): () => void {
|
||||
const key = {}
|
||||
let resolve!: () => void
|
||||
const done = new Promise<void>((settle) => {
|
||||
resolve = settle
|
||||
})
|
||||
this.active.set(key, { context, done })
|
||||
return () => {
|
||||
this.active.delete(key)
|
||||
resolve()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import {
|
||||
isPtyOwnershipTransferMutationEnabled,
|
||||
PTY_OWNERSHIP_TRANSFER_CANARY_ENV
|
||||
} from './pty-ownership-transfer-release-gate'
|
||||
|
||||
describe('PTY ownership-transfer release gate', () => {
|
||||
it('is closed unless the exact canary value is present', () => {
|
||||
expect(isPtyOwnershipTransferMutationEnabled({})).toBe(false)
|
||||
expect(
|
||||
isPtyOwnershipTransferMutationEnabled({
|
||||
[PTY_OWNERSHIP_TRANSFER_CANARY_ENV]: 'true'
|
||||
})
|
||||
).toBe(false)
|
||||
expect(
|
||||
isPtyOwnershipTransferMutationEnabled({
|
||||
[PTY_OWNERSHIP_TRANSFER_CANARY_ENV]: '1 '
|
||||
})
|
||||
).toBe(false)
|
||||
})
|
||||
|
||||
it('opens only for an explicit canary opt-in', () => {
|
||||
expect(
|
||||
isPtyOwnershipTransferMutationEnabled({
|
||||
[PTY_OWNERSHIP_TRANSFER_CANARY_ENV]: '1'
|
||||
})
|
||||
).toBe(true)
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,12 @@
|
||||
/**
|
||||
* The live PTY ownership-transfer path is a release canary until routed,
|
||||
* cross-platform recovery evidence is complete. Keep the opt-in exact and
|
||||
* process-scoped so ordinary profiles cannot inherit it accidentally.
|
||||
*/
|
||||
export const PTY_OWNERSHIP_TRANSFER_CANARY_ENV = 'ORCA_ENABLE_PTY_OWNERSHIP_TRANSFER_MUTATION'
|
||||
|
||||
export function isPtyOwnershipTransferMutationEnabled(
|
||||
env: Record<string, string | undefined> = process.env
|
||||
): boolean {
|
||||
return env[PTY_OWNERSHIP_TRANSFER_CANARY_ENV] === '1'
|
||||
}
|
||||
@@ -0,0 +1,24 @@
|
||||
import { SKILL_SSH_RELAY_CANCEL_UPLOAD_METHOD } from './skill-ssh-relay-contract'
|
||||
|
||||
// Why: a drain must still admit the requests and notifications that let in-flight work finish or cancel.
|
||||
// Owner-reset (T3) and network-tunnel (T4) methods join these sets when their contracts land.
|
||||
const drainRequests = new Set<string>([
|
||||
'relay.status',
|
||||
'fs.unwatchAndWait',
|
||||
'agent.cancelExec',
|
||||
SKILL_SSH_RELAY_CANCEL_UPLOAD_METHOD
|
||||
])
|
||||
const drainNotifications = new Set<string>([
|
||||
'rpc.cancel',
|
||||
'git.responseAck',
|
||||
'git.cancelResponseStream',
|
||||
'fs.streamAck',
|
||||
'fs.cancelStream',
|
||||
'fs.unwatch',
|
||||
'pty.ackData',
|
||||
'pty.setDeliveryPaused'
|
||||
])
|
||||
|
||||
export function allowsRelayWorkDuringDrain(method: string, notification = false): boolean {
|
||||
return (notification ? drainNotifications : drainRequests).has(method)
|
||||
}
|
||||
@@ -0,0 +1,65 @@
|
||||
import { expect, it, vi } from 'vitest'
|
||||
import { RuntimeRpcCallQueueBusyError, RuntimeRpcCallQueuePool } from './runtime-rpc-call-queue'
|
||||
|
||||
it('refuses an active or queued selector without interrupting or replaying its calls', async () => {
|
||||
const queue = new RuntimeRpcCallQueuePool(1, 1)
|
||||
const pending = Promise.withResolvers<string>()
|
||||
const first = queue.enqueue('host', 'terminal.send', () => pending.promise)
|
||||
const secondRun = vi.fn(async () => 'second-ack')
|
||||
const second = queue.enqueue('host', 'terminal.send', secondRun)
|
||||
expect(() => queue.holdIdleSelectors(['host'])).toThrow(RuntimeRpcCallQueueBusyError)
|
||||
pending.resolve('first-ack')
|
||||
await expect(first).resolves.toBe('first-ack')
|
||||
await expect(second).resolves.toBe('second-ack')
|
||||
expect(secondRun).toHaveBeenCalledOnce()
|
||||
await vi.waitFor(() => queue.holdIdleSelectors(['host'])())
|
||||
})
|
||||
|
||||
it('holds both identities, refuses new calls before dispatch, and leaves other hosts running', async () => {
|
||||
const queue = new RuntimeRpcCallQueuePool()
|
||||
const release = queue.holdIdleSelectors(['canonical', 'historical', 'canonical'])
|
||||
const run = vi.fn(async () => 'ack')
|
||||
for (const id of ['canonical', 'historical']) {
|
||||
await expect(queue.enqueue(id, 'terminal.send', run)).rejects.toMatchObject({
|
||||
code: 'runtime_rpc_queue_busy'
|
||||
})
|
||||
}
|
||||
expect(run).not.toHaveBeenCalled()
|
||||
await expect(queue.enqueue('other', 'terminal.send', run)).resolves.toBe('ack')
|
||||
release()
|
||||
await expect(queue.enqueue('historical', 'terminal.send', run)).resolves.toBe('ack')
|
||||
expect(run).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
it('acquires a group atomically without retaining a partial hold on failure', async () => {
|
||||
const queue = new RuntimeRpcCallQueuePool()
|
||||
const release = queue.holdIdleSelectors(['historical'])
|
||||
expect(() => queue.holdIdleSelectors(['canonical', 'historical'])).toThrow(
|
||||
RuntimeRpcCallQueueBusyError
|
||||
)
|
||||
await expect(queue.enqueue('canonical', 'repo.list', async () => 'ok')).resolves.toBe('ok')
|
||||
release()
|
||||
})
|
||||
|
||||
it('does not let an old release clear a newer hold', async () => {
|
||||
const queue = new RuntimeRpcCallQueuePool()
|
||||
const first = queue.holdIdleSelectors(['host'])
|
||||
first()
|
||||
const second = queue.holdIdleSelectors(['host'])
|
||||
first()
|
||||
await expect(queue.enqueue('host', 'repo.list', async () => 'ok')).rejects.toThrow(
|
||||
RuntimeRpcCallQueueBusyError
|
||||
)
|
||||
second()
|
||||
await expect(queue.enqueue('host', 'repo.list', async () => 'ok')).resolves.toBe('ok')
|
||||
})
|
||||
|
||||
it('counts a selector busy while only its long-wait lane has a call', async () => {
|
||||
const queue = new RuntimeRpcCallQueuePool()
|
||||
const pending = Promise.withResolvers<string>()
|
||||
const removal = queue.enqueue('host', 'worktree.rm', () => pending.promise)
|
||||
expect(() => queue.holdIdleSelectors(['host'])).toThrow(RuntimeRpcCallQueueBusyError)
|
||||
pending.resolve('removed')
|
||||
await expect(removal).resolves.toBe('removed')
|
||||
await vi.waitFor(() => queue.holdIdleSelectors(['host'])())
|
||||
})
|
||||
@@ -16,6 +16,15 @@ export class RuntimeRpcCallQueueOverloadError extends Error {
|
||||
}
|
||||
}
|
||||
|
||||
export class RuntimeRpcCallQueueBusyError extends Error {
|
||||
readonly code = 'runtime_rpc_queue_busy'
|
||||
|
||||
constructor() {
|
||||
super('Runtime calls are active or their routing is changing; retry after they settle.')
|
||||
this.name = 'RuntimeRpcCallQueueBusyError'
|
||||
}
|
||||
}
|
||||
|
||||
type QueuedRuntimeCall<T> = {
|
||||
background: boolean
|
||||
retainedBytes: number
|
||||
@@ -57,8 +66,13 @@ function isLongWaitRuntimeMethod(method: string): boolean {
|
||||
return method === 'worktree.rm'
|
||||
}
|
||||
|
||||
function longWaitQueueKey(selector: string): string {
|
||||
return `${selector}\u0000long-wait`
|
||||
}
|
||||
|
||||
export class RuntimeRpcCallQueuePool {
|
||||
private readonly queues = new Map<string, RuntimeCallQueue>()
|
||||
private readonly heldSelectors = new Set<string>()
|
||||
private queuedCallCount = 0
|
||||
private retainedCallBytes = 0
|
||||
|
||||
@@ -70,6 +84,32 @@ export class RuntimeRpcCallQueuePool {
|
||||
private readonly maxRetainedBytes = REMOTE_RUNTIME_MAX_PREPARED_RPC_BYTES
|
||||
) {}
|
||||
|
||||
/** Acquires all idle selectors atomically; release never replays refused calls. */
|
||||
holdIdleSelectors(selectors: readonly string[]): () => void {
|
||||
const unique = [...new Set(selectors)]
|
||||
// A selector's long-wait lane counts too: holding it must see every call still in flight.
|
||||
const busy = (selector: string): boolean =>
|
||||
this.heldSelectors.has(selector) ||
|
||||
this.queues.has(selector) ||
|
||||
this.queues.has(longWaitQueueKey(selector))
|
||||
if (unique.some(busy)) {
|
||||
throw new RuntimeRpcCallQueueBusyError()
|
||||
}
|
||||
for (const selector of unique) {
|
||||
this.heldSelectors.add(selector)
|
||||
}
|
||||
let released = false
|
||||
return () => {
|
||||
if (released) {
|
||||
return
|
||||
}
|
||||
released = true
|
||||
for (const selector of unique) {
|
||||
this.heldSelectors.delete(selector)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
enqueue<T>(
|
||||
selector: string,
|
||||
method: string,
|
||||
@@ -80,8 +120,11 @@ export class RuntimeRpcCallQueuePool {
|
||||
if (signal?.aborted) {
|
||||
return Promise.reject(abortSignalReason(signal))
|
||||
}
|
||||
if (this.heldSelectors.has(selector)) {
|
||||
return Promise.reject(new RuntimeRpcCallQueueBusyError())
|
||||
}
|
||||
// Same concurrency bound, counted apart from the selector's other calls; global caps still apply.
|
||||
const queueKey = isLongWaitRuntimeMethod(method) ? `${selector}\u0000long-wait` : selector
|
||||
const queueKey = isLongWaitRuntimeMethod(method) ? longWaitQueueKey(selector) : selector
|
||||
if (this.queuedCallCount >= this.maxQueuedTotal) {
|
||||
return Promise.reject(new RuntimeRpcCallQueueOverloadError('global'))
|
||||
}
|
||||
|
||||
@@ -0,0 +1,70 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { TransportPublicationDrain } from './transport-publication-drain'
|
||||
|
||||
describe('TransportPublicationDrain', () => {
|
||||
it('waits for every tracked write, including writes tracked while draining', async () => {
|
||||
const publication = new TransportPublicationDrain(() => {})
|
||||
const first = publication.trackWrite()
|
||||
const done = vi.fn()
|
||||
const drain = publication.drain(new AbortController().signal).then(done)
|
||||
const second = publication.trackWrite()
|
||||
first({ ok: true })
|
||||
await Promise.resolve()
|
||||
expect(done).not.toHaveBeenCalled()
|
||||
expect(() => publication.assertDrained()).toThrow('transport_publication_not_drained')
|
||||
second({ ok: true })
|
||||
await drain
|
||||
expect(done).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('ignores a repeated settlement of the same write', async () => {
|
||||
const publication = new TransportPublicationDrain(() => {})
|
||||
const settle = publication.trackWrite()
|
||||
const other = publication.trackWrite()
|
||||
settle({ ok: true })
|
||||
settle({ ok: true })
|
||||
expect(() => publication.assertDrained()).toThrow('transport_publication_not_drained')
|
||||
other({ ok: true })
|
||||
await publication.drain(new AbortController().signal)
|
||||
})
|
||||
|
||||
it('aborts only the observer and keeps a later write failure sticky', async () => {
|
||||
const onFailure = vi.fn()
|
||||
const publication = new TransportPublicationDrain(() => {}, onFailure)
|
||||
const settle = publication.trackWrite()
|
||||
const controller = new AbortController()
|
||||
const drain = publication.drain(controller.signal)
|
||||
controller.abort(new Error('observer cancelled'))
|
||||
await expect(drain).rejects.toThrow('observer cancelled')
|
||||
settle({ ok: false, error: new Error('lost write') })
|
||||
expect(onFailure).toHaveBeenCalledOnce()
|
||||
await expect(publication.drain(new AbortController().signal)).rejects.toThrow('lost write')
|
||||
expect(() => publication.assertCurrent()).toThrow('lost write')
|
||||
})
|
||||
|
||||
it('fails once when the transport it was bound to is replaced', async () => {
|
||||
let current = true
|
||||
const onFailure = vi.fn()
|
||||
const publication = new TransportPublicationDrain(() => {
|
||||
if (!current) {
|
||||
throw new Error('transport_replaced')
|
||||
}
|
||||
}, onFailure)
|
||||
const settle = publication.trackWrite()
|
||||
current = false
|
||||
const drain = publication.drain(new AbortController().signal)
|
||||
await expect(drain).rejects.toThrow('transport_replaced')
|
||||
settle({ ok: true })
|
||||
expect(() => publication.assertDrained()).toThrow('transport_replaced')
|
||||
expect(onFailure).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('refuses construction against a transport that is already stale', () => {
|
||||
expect(
|
||||
() =>
|
||||
new TransportPublicationDrain(() => {
|
||||
throw new Error('transport_replaced')
|
||||
})
|
||||
).toThrow('transport_replaced')
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,87 @@
|
||||
import { waitForPromiseWithSignal } from './abort-signal-reason'
|
||||
|
||||
export type TransportPublicationSettlement = { ok: true } | { ok: false; error: Error }
|
||||
|
||||
/** Local write completion only; downstream consumption requires separate proof. */
|
||||
export class TransportPublicationDrain {
|
||||
private pending = 0
|
||||
protected failure: Error | undefined
|
||||
private readonly observers = new Set<() => void>()
|
||||
|
||||
constructor(
|
||||
private readonly assertTransport: () => void,
|
||||
private readonly onFailure: (error: Error) => void = () => {}
|
||||
) {
|
||||
this.assertCurrent()
|
||||
}
|
||||
|
||||
trackWrite(): (result: TransportPublicationSettlement) => void {
|
||||
this.pending++
|
||||
let settled = false
|
||||
return (result) => {
|
||||
if (settled) {
|
||||
return
|
||||
}
|
||||
settled = true
|
||||
this.pending--
|
||||
if (!result.ok) {
|
||||
this.fail(result.error)
|
||||
}
|
||||
this.changed()
|
||||
}
|
||||
}
|
||||
|
||||
fail(error: Error): void {
|
||||
if (this.failure) {
|
||||
return
|
||||
}
|
||||
this.failure = error
|
||||
this.changed()
|
||||
this.onFailure(error)
|
||||
}
|
||||
|
||||
assertDrained(): void {
|
||||
this.assertCurrent()
|
||||
if (this.pending !== 0) {
|
||||
throw new Error('transport_publication_not_drained')
|
||||
}
|
||||
}
|
||||
|
||||
async drain(signal: AbortSignal): Promise<void> {
|
||||
signal.throwIfAborted()
|
||||
while (this.pending !== 0) {
|
||||
this.assertCurrent()
|
||||
// Why: no Promise.withResolvers — the legacy relay bundle still targets Node 18 hosts.
|
||||
let notify!: () => void
|
||||
const changed = new Promise<void>((resolve) => {
|
||||
notify = resolve
|
||||
})
|
||||
this.observers.add(notify)
|
||||
try {
|
||||
await waitForPromiseWithSignal(changed, signal)
|
||||
} finally {
|
||||
this.observers.delete(notify)
|
||||
}
|
||||
}
|
||||
signal.throwIfAborted()
|
||||
this.assertDrained()
|
||||
}
|
||||
|
||||
assertCurrent(): void {
|
||||
if (this.failure) {
|
||||
throw this.failure
|
||||
}
|
||||
try {
|
||||
this.assertTransport()
|
||||
} catch (error) {
|
||||
this.fail(error instanceof Error ? error : new Error(String(error)))
|
||||
throw this.failure
|
||||
}
|
||||
}
|
||||
|
||||
private changed(): void {
|
||||
for (const notify of this.observers) {
|
||||
notify()
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user