mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 00:02:10 +00:00
fix(agent-session): wait for in-flight session-store writes before teardown returns (#23545)
* fix(agent-session): stop lease renewal before the renewal's write lands Clearing the renewal interval only cancelled the next tick. A tick already past its guard still had a whole-file store transaction to commit, and the store's transaction lock re-creates the store directory before it writes, so that commit could land after host teardown had finished releasing everything it touches. `stop()` now resolves once the tick in flight has finished writing, and host teardown's stop-lease-renewal phase waits for it. The three test harnesses that model a host vanishing without a clean quit shared a copy of the same incomplete shutdown; they now share one helper that waits. The symptom was a CI flake: the refusal-oracle spec removes its temp directory in `afterEach`, and a renewal landing mid-removal put the store directory back, so the removal failed with ENOTEMPTY on the temp root. * fix(agent-session): wait for the delivery loop's restart when abandoning a host The abandon helper disposed the delivery loop and moved on. Disposing only stops the loop's NEXT step: a step already past that check keeps going, and the restart it runs for an accepted send reserves an owner, which is a store commit. The store re-creates its own directory before every commit, so that commit put the directory back under the temp-directory removal the test does next, and the removal failed with ENOTEMPTY. Quit already waits for exactly this work, in its drain-attaches phase — every attach is registered with the task queue from enqueue. The helper now runs the same drain, in quit's order, so it waits for both producers that reach the store after the last awaited call returns. Adds a regression test that holds the loop's restart inside its provider acquisition and asserts abandoning does not return until it lands. * chore: re-trigger PR checks The push to 2ecc9c6e emitted no pull_request event, so the matrix never ran.
This commit is contained in:
+3
-2
@@ -34,8 +34,9 @@ export class StructuredAgentSessionHostRuntimeState {
|
||||
this.leaseRenewer.start()
|
||||
}
|
||||
|
||||
stopLeaseRenewal(): void {
|
||||
this.leaseRenewer.stop()
|
||||
/** Resolves once a renewal tick already in flight has finished writing. */
|
||||
stopLeaseRenewal(): Promise<void> {
|
||||
return this.leaseRenewer.stop()
|
||||
}
|
||||
|
||||
/** The sink the session's current child writes through, created on first use. */
|
||||
|
||||
@@ -49,7 +49,7 @@ async function withPhaseTimeout(run: () => Promise<void>, timeoutMs: number): Pr
|
||||
export function structuredAgentSessionHostTeardownPhases(collaborators: {
|
||||
idleSweep: { dispose: () => Promise<void> | void }
|
||||
runtimeState: {
|
||||
stopLeaseRenewal: () => void
|
||||
stopLeaseRenewal: () => Promise<void> | void
|
||||
flushAllEventSinks: () => Promise<void>
|
||||
}
|
||||
tasks: { drainAttaches: () => Promise<void> }
|
||||
|
||||
+123
@@ -0,0 +1,123 @@
|
||||
// The helper is what the restart specs remove their temp directory behind, so what it waits for is
|
||||
// load-bearing: a store commit that lands afterwards re-creates the directory it just removed.
|
||||
|
||||
import { existsSync } from 'node:fs'
|
||||
import { mkdtemp, rm } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope'
|
||||
import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
|
||||
import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter'
|
||||
import { StructuredAgentSessionHost } from './structured-agent-session-host'
|
||||
import { abandonStructuredAgentSessionHost } from './structured-agent-session-host-test-abandon'
|
||||
import {
|
||||
HOST_TEST_NOW as NOW,
|
||||
HOST_TEST_SESSION as SESSION,
|
||||
HOST_TEST_THREAD as THREAD,
|
||||
hostTestAttachParams,
|
||||
hostTestMessage
|
||||
} from './structured-agent-session-host-test-data'
|
||||
|
||||
const CALLER = { callerKey: 'client-1' }
|
||||
|
||||
describe('abandoning a structured agent-session host', () => {
|
||||
it('waits for the restart the delivery loop woke for an accepted send', async () => {
|
||||
const root = await mkdtemp(join(tmpdir(), 'orca-abandon-'))
|
||||
const store = await AgentSessionRecordStore.open({
|
||||
directory: join(root, 'store'),
|
||||
hostId: 'local'
|
||||
})
|
||||
// Holds the restart inside its provider acquisition, so the point teardown must not run past
|
||||
// is exact rather than a timing window.
|
||||
let gate: Promise<void> | null = null
|
||||
let openGate = (): void => {}
|
||||
let reportEntered = (): void => {}
|
||||
const entered = new Promise<void>((resolve) => {
|
||||
reportEntered = resolve
|
||||
})
|
||||
const adapter: StructuredAgentSessionAdapter = {
|
||||
acquire: async ({ fence }) => {
|
||||
if (gate) {
|
||||
reportEntered()
|
||||
await gate
|
||||
}
|
||||
return {
|
||||
process: {
|
||||
hostId: 'local',
|
||||
pid: 4242,
|
||||
processStartTimeMs: NOW - 1_000,
|
||||
spawnToken: store.getRecord(SESSION)?.lease.reservedSpawnToken ?? 'spawn-a'
|
||||
},
|
||||
link: {
|
||||
linkId: `link-${fence}`,
|
||||
handle: { provider: 'codex', threadId: THREAD },
|
||||
origin: 'created',
|
||||
mintedAtFence: fence,
|
||||
observedAt: NOW
|
||||
}
|
||||
}
|
||||
},
|
||||
dispatch: async () => ({
|
||||
state: 'accepted',
|
||||
providerIdentity: { provider: 'codex', threadId: THREAD, turnId: 'turn-1', ordinal: 1 }
|
||||
}),
|
||||
cancelTurn: async () => ({ cancelled: true }),
|
||||
answerPrompt: async () => undefined,
|
||||
releaseAcquisition: async () => true,
|
||||
setOption: async () => undefined
|
||||
}
|
||||
const host = new StructuredAgentSessionHost({
|
||||
store,
|
||||
adapter,
|
||||
probeOwner: async () => ({
|
||||
outcome: 'identity-matched',
|
||||
matchedOn: ['process-start-time']
|
||||
}),
|
||||
journalRoot: root,
|
||||
claimKeyId: 'key-1',
|
||||
mintSpawnToken: () => 'spawn-a',
|
||||
now: () => NOW
|
||||
})
|
||||
expect(await host.attach(CALLER, hostTestAttachParams(null))).toMatchObject({ ok: true })
|
||||
// The conversation stays and its provider child does not, so the next send makes the delivery
|
||||
// loop start one — the shape the refusal-oracle spec ends on.
|
||||
await host.close(SESSION)
|
||||
gate = new Promise<void>((resolve) => {
|
||||
openGate = resolve
|
||||
})
|
||||
|
||||
const body = hostTestMessage('delivery loop')
|
||||
expect(
|
||||
await host.send(CALLER, {
|
||||
envelope: {
|
||||
sessionId: SESSION,
|
||||
clientOperationId: `${NOW}-000000000000000000000000000003e9`,
|
||||
expectedRuntimeFence: store.getRecord(SESSION)?.lease.runtimeFence ?? 1,
|
||||
payloadFingerprint: computeAgentSessionPayloadFingerprint({
|
||||
method: 'agentSession.send',
|
||||
sessionId: SESSION,
|
||||
fields: { body }
|
||||
})
|
||||
},
|
||||
body
|
||||
})
|
||||
).toMatchObject({ ok: true })
|
||||
await entered
|
||||
|
||||
let abandoned = false
|
||||
const abandoning = abandonStructuredAgentSessionHost(host).then(() => {
|
||||
abandoned = true
|
||||
})
|
||||
await new Promise((resolve) => setTimeout(resolve, 20))
|
||||
expect(abandoned).toBe(false)
|
||||
|
||||
openGate()
|
||||
await abandoning
|
||||
|
||||
// Nothing is left to put the directory back after this returns.
|
||||
await rm(root, { recursive: true })
|
||||
await new Promise((resolve) => setTimeout(resolve, 200))
|
||||
expect(existsSync(root)).toBe(false)
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,25 @@
|
||||
// What a host that VANISHES without a clean quit has to release, for the tests that model one.
|
||||
//
|
||||
// Deliberately not the quit path: quit evicts provider children and releases leases, and that is
|
||||
// the state a restart test is checking gets re-derived from disk. It stops short of ownership.
|
||||
//
|
||||
// What it cannot skip is work already in flight, which is what every copy of this helper used to
|
||||
// get wrong. Two producers reach the session store after the last awaited call has returned — a
|
||||
// lease-renewal tick, and the restart the delivery loop wakes for an accepted send — and the store
|
||||
// re-creates its own directory before every commit. A commit landing after the test removed its
|
||||
// temp directory therefore puts that directory back, and the removal fails with ENOTEMPTY. Quit
|
||||
// waits for both, in its `stop-lease-renewal` and `drain-attaches` phases; so does this.
|
||||
|
||||
import type { StructuredAgentSessionHost } from './structured-agent-session-host'
|
||||
|
||||
export async function abandonStructuredAgentSessionHost(
|
||||
host: StructuredAgentSessionHost
|
||||
): Promise<void> {
|
||||
// First, so the drain below cannot race the loop into enqueueing another step.
|
||||
host['conversationDelivery'].loop.dispose()
|
||||
host['lifetime'].dispose()
|
||||
await host['runtimeState'].stopLeaseRenewal()
|
||||
await host['tasks'].drainAttaches()
|
||||
await Promise.all([...host['sessions'].values()].map((session) => session.journal.close()))
|
||||
host['sessions'].clear()
|
||||
}
|
||||
+26
-1
@@ -190,11 +190,36 @@ describe('structured agent-session lease renewal', () => {
|
||||
{ timeout: 5000 }
|
||||
)
|
||||
} finally {
|
||||
renewer.stop()
|
||||
await renewer.stop()
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
|
||||
it('stops only once a renewal already in flight has finished writing', async () => {
|
||||
const store = await liveStore()
|
||||
let releaseProbe = (): void => {}
|
||||
const probing = new Promise<void>((resolve) => {
|
||||
releaseProbe = resolve
|
||||
})
|
||||
const renewer = new StructuredAgentSessionLeaseRenewer({
|
||||
store,
|
||||
probe: async () => {
|
||||
await probing
|
||||
return { outcome: 'identity-matched', matchedOn: ['process-start-time'] }
|
||||
},
|
||||
now: () => NOW + 10_000
|
||||
})
|
||||
|
||||
void renewer.renewNow()
|
||||
// Nothing has been written yet: the tick is parked in its probe.
|
||||
expect(store.getRecord('session-renewal')?.lease.lastRenewedAt).toBe(NOW)
|
||||
const stopped = renewer.stop()
|
||||
releaseProbe()
|
||||
await stopped
|
||||
|
||||
expect(store.getRecord('session-renewal')?.lease.lastRenewedAt).toBe(NOW + 10_000)
|
||||
})
|
||||
|
||||
it('renews every live owner only after re-proving its child identity', async () => {
|
||||
const store = await liveStore()
|
||||
const probe = vi.fn(async () => ({
|
||||
|
||||
@@ -10,6 +10,9 @@ const RENEW_INTERVAL_MS = Math.floor(AGENT_SESSION_LEASE_TTL_MS / 3)
|
||||
export class StructuredAgentSessionLeaseRenewer {
|
||||
private timer: ReturnType<typeof setInterval> | null = null
|
||||
private running = false
|
||||
/** The tick in flight. Never rejects: the timer path reports renewal failures through `onError`,
|
||||
* and stopping must not turn one into a teardown failure as well. */
|
||||
private inFlight: Promise<void> = Promise.resolve()
|
||||
|
||||
constructor(
|
||||
private readonly input: {
|
||||
@@ -32,70 +35,81 @@ export class StructuredAgentSessionLeaseRenewer {
|
||||
this.timer.unref?.()
|
||||
}
|
||||
|
||||
stop(): void {
|
||||
/** Clearing the interval only stops the NEXT tick. A tick already past its guard still has a
|
||||
* store transaction to commit, and that transaction re-creates the store directory, so a stop
|
||||
* that returned before it landed would let the write outlive whatever tore the host down. */
|
||||
stop(): Promise<void> {
|
||||
if (this.timer) {
|
||||
clearInterval(this.timer)
|
||||
this.timer = null
|
||||
}
|
||||
return this.inFlight
|
||||
}
|
||||
|
||||
async renewNow(): Promise<void> {
|
||||
renewNow(): Promise<void> {
|
||||
if (this.running) {
|
||||
return
|
||||
return this.inFlight
|
||||
}
|
||||
this.running = true
|
||||
try {
|
||||
const records = this.input.store.listRecords().filter(
|
||||
(record) =>
|
||||
!record.lease.unreconciled &&
|
||||
record.lease.claimStatus === 'live' &&
|
||||
record.lease.ownerProcess !== null &&
|
||||
// A record parked in recovery has no transport the host can vouch for; renewing it
|
||||
// keeps an orphan pid's lease reading as a healthy owner.
|
||||
record.lease.handoffStage !== 'recovering'
|
||||
)
|
||||
const probes = await this.probe(records)
|
||||
const renewals: {
|
||||
sessionId: string
|
||||
fence: number
|
||||
childProbe: AgentSessionOwnerProbe
|
||||
now: number
|
||||
}[] = []
|
||||
const now = this.input.now()
|
||||
for (const record of records) {
|
||||
const probe = probes.get(record.sessionId)
|
||||
if (!probe) {
|
||||
continue
|
||||
}
|
||||
renewals.push({
|
||||
sessionId: record.sessionId,
|
||||
fence: record.lease.runtimeFence,
|
||||
childProbe: probe,
|
||||
now
|
||||
})
|
||||
}
|
||||
// The store persists the whole record file per transaction, so keep the healthy path to
|
||||
// one commit. If one renewal is superseded, retrying individually preserves isolation.
|
||||
let results: PromiseSettledResult<AgentSessionRecord>[]
|
||||
try {
|
||||
const renewed = await this.input.store.renewLeases(renewals)
|
||||
results = renewed.map((record) => ({ status: 'fulfilled', value: record }) as const)
|
||||
} catch {
|
||||
results = await Promise.allSettled(
|
||||
renewals.map((renewal) => this.input.store.renewLease(renewal))
|
||||
)
|
||||
}
|
||||
results.forEach((result, index) => {
|
||||
if (result.status === 'rejected') {
|
||||
const renewal = renewals[index]
|
||||
if (renewal) {
|
||||
this.input.onError?.({ sessionId: renewal.sessionId, error: result.reason })
|
||||
}
|
||||
}
|
||||
})
|
||||
} finally {
|
||||
const attempt = this.renewOnce().finally(() => {
|
||||
this.running = false
|
||||
})
|
||||
this.inFlight = attempt.then(
|
||||
() => undefined,
|
||||
() => undefined
|
||||
)
|
||||
return attempt
|
||||
}
|
||||
|
||||
private async renewOnce(): Promise<void> {
|
||||
const records = this.input.store.listRecords().filter(
|
||||
(record) =>
|
||||
!record.lease.unreconciled &&
|
||||
record.lease.claimStatus === 'live' &&
|
||||
record.lease.ownerProcess !== null &&
|
||||
// A record parked in recovery has no transport the host can vouch for; renewing it
|
||||
// keeps an orphan pid's lease reading as a healthy owner.
|
||||
record.lease.handoffStage !== 'recovering'
|
||||
)
|
||||
const probes = await this.probe(records)
|
||||
const renewals: {
|
||||
sessionId: string
|
||||
fence: number
|
||||
childProbe: AgentSessionOwnerProbe
|
||||
now: number
|
||||
}[] = []
|
||||
const now = this.input.now()
|
||||
for (const record of records) {
|
||||
const probe = probes.get(record.sessionId)
|
||||
if (!probe) {
|
||||
continue
|
||||
}
|
||||
renewals.push({
|
||||
sessionId: record.sessionId,
|
||||
fence: record.lease.runtimeFence,
|
||||
childProbe: probe,
|
||||
now
|
||||
})
|
||||
}
|
||||
// The store persists the whole record file per transaction, so keep the healthy path to
|
||||
// one commit. If one renewal is superseded, retrying individually preserves isolation.
|
||||
let results: PromiseSettledResult<AgentSessionRecord>[]
|
||||
try {
|
||||
const renewed = await this.input.store.renewLeases(renewals)
|
||||
results = renewed.map((record) => ({ status: 'fulfilled', value: record }) as const)
|
||||
} catch {
|
||||
results = await Promise.allSettled(
|
||||
renewals.map((renewal) => this.input.store.renewLease(renewal))
|
||||
)
|
||||
}
|
||||
results.forEach((result, index) => {
|
||||
if (result.status === 'rejected') {
|
||||
const renewal = renewals[index]
|
||||
if (renewal) {
|
||||
this.input.onError?.({ sessionId: renewal.sessionId, error: result.reason })
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
private async probe(
|
||||
|
||||
+2
-8
@@ -5,6 +5,7 @@ import { afterEach, describe, expect, it } from 'vitest'
|
||||
import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
|
||||
import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter'
|
||||
import { StructuredAgentSessionHost } from './structured-agent-session-host'
|
||||
import { abandonStructuredAgentSessionHost } from './structured-agent-session-host-test-abandon'
|
||||
import {
|
||||
HOST_TEST_NOW,
|
||||
HOST_TEST_SESSION,
|
||||
@@ -58,15 +59,8 @@ function createHost(
|
||||
return host
|
||||
}
|
||||
|
||||
async function abandonHost(host: StructuredAgentSessionHost): Promise<void> {
|
||||
host['runtimeState'].stopLeaseRenewal()
|
||||
host['lifetime'].dispose()
|
||||
await Promise.all([...host['sessions'].values()].map((session) => session.journal.close()))
|
||||
host['sessions'].clear()
|
||||
}
|
||||
|
||||
afterEach(async () => {
|
||||
await Promise.all(hosts.splice(0).map(abandonHost))
|
||||
await Promise.all(hosts.splice(0).map(abandonStructuredAgentSessionHost))
|
||||
await rm(root, { recursive: true, force: true })
|
||||
root = ''
|
||||
})
|
||||
|
||||
+4
-12
@@ -11,6 +11,7 @@ import { readProcessStartTimeMs } from '../../runtime/agent-session-process-iden
|
||||
import { createStructuredAgentSessionOwnerProbe } from '../../runtime/structured-agent-session-owner-probe'
|
||||
import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter'
|
||||
import { StructuredAgentSessionHost } from './structured-agent-session-host'
|
||||
import { abandonStructuredAgentSessionHost } from './structured-agent-session-host-test-abandon'
|
||||
import type { StructuredAgentSessionHostDeps } from './structured-agent-session-host-types'
|
||||
import {
|
||||
HOST_TEST_NOW as NOW,
|
||||
@@ -88,17 +89,8 @@ function openHost(overrides: Partial<StructuredAgentSessionHostDeps> = {}): void
|
||||
})
|
||||
}
|
||||
|
||||
async function abandonHost(abandonedHost: StructuredAgentSessionHost): Promise<void> {
|
||||
abandonedHost['runtimeState'].stopLeaseRenewal()
|
||||
abandonedHost['lifetime'].dispose()
|
||||
await Promise.all(
|
||||
[...abandonedHost['sessions'].values()].map((session) => session.journal.close())
|
||||
)
|
||||
abandonedHost['sessions'].clear()
|
||||
}
|
||||
|
||||
async function reopenStore(): Promise<void> {
|
||||
await abandonHost(host)
|
||||
await abandonStructuredAgentSessionHost(host)
|
||||
store = await AgentSessionRecordStore.open({ directory: join(root, 'store'), hostId: 'local' })
|
||||
}
|
||||
|
||||
@@ -125,8 +117,8 @@ beforeEach(async () => {
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
await abandonHost(host)
|
||||
await Promise.all([...supersededHosts].map(abandonHost))
|
||||
await abandonStructuredAgentSessionHost(host)
|
||||
await Promise.all([...supersededHosts].map(abandonStructuredAgentSessionHost))
|
||||
supersededHosts.clear()
|
||||
await Promise.all([...spawnedOwners].map((child) => stopOwner(child)))
|
||||
await rm(root, { recursive: true, force: true })
|
||||
|
||||
+2
-11
@@ -19,6 +19,7 @@ import {
|
||||
import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
|
||||
import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter'
|
||||
import { StructuredAgentSessionHost } from './structured-agent-session-host'
|
||||
import { abandonStructuredAgentSessionHost } from './structured-agent-session-host-test-abandon'
|
||||
import {
|
||||
HOST_TEST_NOW as NOW,
|
||||
HOST_TEST_SESSION as SESSION,
|
||||
@@ -103,19 +104,9 @@ async function createHarness(options: { attached?: boolean } = {}) {
|
||||
return harness
|
||||
}
|
||||
|
||||
async function abandonHost(host: StructuredAgentSessionHost): Promise<void> {
|
||||
host['runtimeState'].stopLeaseRenewal()
|
||||
host['lifetime'].dispose()
|
||||
host['conversationDelivery'].loop.dispose()
|
||||
// An accepted send starts the agent in the background; its lease write must land before rm.
|
||||
await host['tasks'].drainAttaches()
|
||||
await Promise.all([...host['sessions'].values()].map((session) => session.journal.close()))
|
||||
host['sessions'].clear()
|
||||
}
|
||||
|
||||
afterEach(async () => {
|
||||
const completed = harnesses.splice(0)
|
||||
await Promise.all(completed.map(async ({ host }) => abandonHost(host)))
|
||||
await Promise.all(completed.map(async ({ host }) => abandonStructuredAgentSessionHost(host)))
|
||||
await Promise.all(completed.map(async ({ root }) => rm(root, { recursive: true })))
|
||||
})
|
||||
|
||||
|
||||
+2
-4
@@ -24,6 +24,7 @@ import {
|
||||
} from './structured-agent-session-adapter'
|
||||
import type { StructuredAgentSessionEventSink } from './structured-agent-session-event-sink'
|
||||
import { StructuredAgentSessionHost } from './structured-agent-session-host'
|
||||
import { abandonStructuredAgentSessionHost } from './structured-agent-session-host-test-abandon'
|
||||
import { unexpectedProviderExitOutcome } from './structured-agent-session-dead-generation-settlement'
|
||||
import { readAgentJournalTurn } from '../../../shared/agent-session-turn-record'
|
||||
import type { StructuredAgentSessionStatusSink } from './structured-agent-session-status-feed'
|
||||
@@ -384,10 +385,7 @@ describe('startup', () => {
|
||||
it('settles an idle absent owner without chat pollution and resumes the same provider identity', async () => {
|
||||
await attach()
|
||||
const beforeRestart = store.getRecord(SESSION)
|
||||
host['runtimeState'].stopLeaseRenewal()
|
||||
host['lifetime'].dispose()
|
||||
await host['sessions'].get(SESSION)?.journal.close()
|
||||
host['sessions'].clear()
|
||||
await abandonStructuredAgentSessionHost(host)
|
||||
|
||||
store = await AgentSessionRecordStore.open({ directory: join(root, 'store'), hostId: 'local' })
|
||||
openHost(async () => ({ outcome: 'pid-absent' }))
|
||||
|
||||
Reference in New Issue
Block a user