mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 16:02:29 +00:00
Cut CI time in store, Git contention and readiness tests (#25147)
* Check store retention boundaries with a faster independent oracle * Replace CI diagnostic sleeps with gated contention and scoped transcript clocks
This commit is contained in:
@@ -1959,3 +1959,60 @@ same ordering assertion. The gate was released, both controls drained the captur
|
||||
job, and the instrumentation was removed. Changed-code quality passed. This proves
|
||||
the teardown ordering mechanism, not a measured avoided-retry saving. Final-head
|
||||
hosted qualification remains required.
|
||||
|
||||
## October 4 store oracle and retention fixtures
|
||||
|
||||
The randomized in-place-store test validated the copying oracle twice after
|
||||
accepted mutations and compared snapshots through the same production parser.
|
||||
Its 5,000-step retention fixture generated enough tombstones to hit the count
|
||||
limit, but never reached the 4,096-revision age boundary.
|
||||
|
||||
The test retains all four seeds and 1,500 mutations per seed, removes the duplicate
|
||||
validation, and projects snapshots directly from the copying oracle's validated
|
||||
maps. Separate fixtures now check the revision before, at and after expiry and
|
||||
count overflow. Production code is unchanged.
|
||||
|
||||
Three alternating one-worker pairs on `ubuntu-24.04-arm` in
|
||||
[37180517143](https://github.com/stablyai/orca/actions/runs/37180517143)
|
||||
measured baseline invocation times 33.551 / 33.304 / 33.529 seconds and candidate
|
||||
13.848 / 13.816 / 13.875 seconds: median 33.529 to 13.848 seconds, saving 19.681
|
||||
seconds (58.7%). Baseline passed seven tests; candidate passed eight. This is a
|
||||
focused test saving, not a measured whole-shard or queue-delay change.
|
||||
|
||||
Hosted Node typecheck passed. Separate fault controls failed the intended
|
||||
assertion for early, late and disabled age expiry, disabled count compaction,
|
||||
and a snapshot that drops child descriptions. The description fault passes with
|
||||
the original parser-sharing oracle and fails with the independent projection.
|
||||
|
||||
## October 4 Git contention and remaining readiness waits
|
||||
|
||||
The full Git admission benchmark compared a disabled arm with no correctness
|
||||
assertions to an enabled arm with structural ledger checks. Its default CI test
|
||||
now saturates the real base and headroom budgets with FIFO-gated child processes,
|
||||
queues older background and newer interactive work, releases base slots, and
|
||||
requires interactive priority, matching outputs and complete permit release.
|
||||
The full original diagnostic remains opt-in through
|
||||
`ORCA_GIT_ADMISSION_STORM_MEASUREMENT=1`; both opt-in tests passed locally.
|
||||
The existing Windows real-Git parity tests remain unchanged; this fixture retains
|
||||
its existing POSIX platform scope.
|
||||
|
||||
Two remaining Antigravity transcript tests used real 5,000ms refusal windows.
|
||||
They now use the existing scoped `waitForTranscriptIdle` timer harness after the
|
||||
emulator drains. All 60 tests, original captured transcripts, deadlines and
|
||||
readiness assertions remain.
|
||||
|
||||
Three alternating one-worker hosted ARM pairs in
|
||||
[37180614492](https://github.com/stablyai/orca/actions/runs/37180614492)
|
||||
measured these complete focused invocations:
|
||||
|
||||
| Suite | Baseline seconds | Candidate seconds | Median saving |
|
||||
| --- | --- | --- | --- |
|
||||
| Git admission storm | 26.619 / 26.635 / 26.582 | 1.017 / 1.018 / 1.016 | 25.602s (96.2%) |
|
||||
| Antigravity readiness | 27.347 / 27.910 / 27.550 | 13.855 / 13.894 / 13.800 | 13.695s (49.7%) |
|
||||
|
||||
Each candidate passed its original meaningful checks. Hosted Node typecheck
|
||||
passed. Separate scheduler faults for bypassed admission, withheld release and
|
||||
FIFO-only priority failed the queued-contention or interactive-start assertion.
|
||||
Two additional local transcript faults failed the original picker-rejection and
|
||||
repaint-readiness assertions. These are focused suite savings; whole-shard time
|
||||
and queue delay were not measured by this experiment.
|
||||
|
||||
@@ -1,250 +1,183 @@
|
||||
/**
|
||||
* POSIX-only measurement: a PATH-injected git fixture exercises the real spawn path.
|
||||
* Timing distributions are reported for field comparison; CI assertions stay structural.
|
||||
*/
|
||||
import { chmod, mkdtemp, mkdir, readdir, rm, writeFile } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { existsSync, watch } from 'node:fs'
|
||||
import { chmod, mkdir, readdir, writeFile } from 'node:fs/promises'
|
||||
import path from 'node:path'
|
||||
import { afterEach, describe, expect, it } from 'vitest'
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { runProcess } from '../../../shared/child-process/run-process'
|
||||
import { gitExecFileAsync } from './git-exec-file'
|
||||
import {
|
||||
GIT_ADMISSION_AGING_MS,
|
||||
GENERAL_CAP,
|
||||
GENERAL_HEADROOM,
|
||||
GitAdmissionScheduler,
|
||||
MAX_GIT_CHILDREN,
|
||||
_gitAdmissionSnapshotForTests,
|
||||
_resetGitAdmissionForTests,
|
||||
type GitAdmissionEvent
|
||||
} from './git-subprocess-admission'
|
||||
import {
|
||||
assertAdmissionLedger,
|
||||
createStormRoot,
|
||||
formatMeasurementTable,
|
||||
measureStorm,
|
||||
type InteractiveQueueSnapshot
|
||||
} from './git-admission-storm-test-fixture'
|
||||
|
||||
type StormMeasurement = {
|
||||
mode: 'disabled' | 'enabled'
|
||||
maxConcurrentChildren: number
|
||||
eventLoopMaxDriftMs: number
|
||||
eventLoopP99DriftMs: number
|
||||
interactiveP50Ms: number
|
||||
interactiveP95Ms: number
|
||||
totalWallMs: number
|
||||
admissionEvents: GitAdmissionEvent[]
|
||||
interactiveQueueSnapshots: InteractiveQueueSnapshot[]
|
||||
function waitForLiveChild(stateDir: string, id: string, signal: AbortSignal): Promise<void> {
|
||||
return new Promise((resolve, reject) => {
|
||||
const watcher = watch(stateDir)
|
||||
const finish = (error?: Error): void => {
|
||||
clearTimeout(deadline)
|
||||
watcher.close()
|
||||
signal.removeEventListener('abort', aborted)
|
||||
if (error) {
|
||||
reject(error)
|
||||
} else {
|
||||
resolve()
|
||||
}
|
||||
}
|
||||
const aborted = (): void => finish(new Error(`Aborted waiting for ${id}`))
|
||||
const deadline = setTimeout(() => finish(new Error(`Child ${id} did not start`)), 5_000)
|
||||
const check = (): void => {
|
||||
if (existsSync(path.join(stateDir, `${id}.live`))) {
|
||||
finish()
|
||||
}
|
||||
}
|
||||
watcher.on('change', check)
|
||||
watcher.once('error', finish)
|
||||
signal.addEventListener('abort', aborted, { once: true })
|
||||
if (signal.aborted) {
|
||||
aborted()
|
||||
} else {
|
||||
check()
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
type InteractiveQueueSnapshot = {
|
||||
commandLabel: string
|
||||
backgroundWaiterIds: number[]
|
||||
}
|
||||
|
||||
const tempRoots: string[] = []
|
||||
const originalAdmissionDisabled = process.env.ORCA_GIT_ADMISSION_DISABLED
|
||||
|
||||
afterEach(async () => {
|
||||
if (originalAdmissionDisabled === undefined) {
|
||||
delete process.env.ORCA_GIT_ADMISSION_DISABLED
|
||||
} else {
|
||||
process.env.ORCA_GIT_ADMISSION_DISABLED = originalAdmissionDisabled
|
||||
}
|
||||
_resetGitAdmissionForTests()
|
||||
await Promise.all(tempRoots.splice(0).map((root) => rm(root, { recursive: true, force: true })))
|
||||
})
|
||||
|
||||
function percentile(values: readonly number[], percentileValue: number): number {
|
||||
if (values.length === 0) {
|
||||
return 0
|
||||
}
|
||||
const sorted = [...values].sort((a, b) => a - b)
|
||||
return sorted[Math.min(sorted.length - 1, Math.ceil(sorted.length * percentileValue) - 1)]
|
||||
}
|
||||
|
||||
async function liveChildCount(stateDir: string): Promise<number> {
|
||||
return (await readdir(stateDir)).filter((name) => name.endsWith('.live')).length
|
||||
}
|
||||
|
||||
async function createStubGit(root: string): Promise<string> {
|
||||
async function verifySpawnContention(): Promise<void> {
|
||||
delete process.env.ORCA_GIT_ADMISSION_DISABLED
|
||||
const admissionEvents: GitAdmissionEvent[] = []
|
||||
_resetGitAdmissionForTests(
|
||||
new GitAdmissionScheduler({
|
||||
now: () => 0,
|
||||
onAdmissionEvent: (event) => admissionEvents.push(event)
|
||||
})
|
||||
)
|
||||
const root = await createStormRoot('contention')
|
||||
const stateDir = path.join(root, 'state')
|
||||
const gateDir = path.join(root, 'gates')
|
||||
const binDir = path.join(root, 'bin')
|
||||
await mkdir(binDir)
|
||||
const stubPath = path.join(binDir, 'git')
|
||||
await Promise.all([stateDir, gateDir, binDir].map((directory) => mkdir(directory)))
|
||||
const backgroundIds = Array.from({ length: GENERAL_CAP + 4 }, (_, index) => `background-${index}`)
|
||||
const headroomIds = Array.from({ length: GENERAL_HEADROOM }, (_, index) => `interactive-${index}`)
|
||||
const queuedIds = ['interactive-queued-0', 'interactive-queued-1']
|
||||
const ids = [...backgroundIds, ...headroomIds, ...queuedIds]
|
||||
const gates = await runProcess({
|
||||
program: 'mkfifo',
|
||||
args: ids.map((id) => path.join(gateDir, id))
|
||||
})
|
||||
expect(gates.code, gates.stderr).toBe(0)
|
||||
const stub = path.join(binDir, 'git')
|
||||
await writeFile(
|
||||
stubPath,
|
||||
stub,
|
||||
`#!/bin/sh
|
||||
set -eu
|
||||
live="$ORCA_STUB_STATE_DIR/$ORCA_STUB_ID.live"
|
||||
: > "$live"
|
||||
trap 'rm -f "$live"' EXIT HUP INT TERM
|
||||
sleep "$(awk "BEGIN { print $ORCA_STUB_SLEEP_MS / 1000 }")"
|
||||
IFS= read -r release < "$ORCA_STUB_GATE_DIR/$ORCA_STUB_ID"
|
||||
printf 'stub:%s\\n' "$*"
|
||||
`
|
||||
)
|
||||
await chmod(stubPath, 0o755)
|
||||
return binDir
|
||||
}
|
||||
|
||||
function setAdmissionMode(mode: StormMeasurement['mode']): void {
|
||||
if (mode === 'disabled') {
|
||||
process.env.ORCA_GIT_ADMISSION_DISABLED = '1'
|
||||
} else {
|
||||
delete process.env.ORCA_GIT_ADMISSION_DISABLED
|
||||
}
|
||||
}
|
||||
|
||||
async function measureStorm(mode: StormMeasurement['mode']): Promise<StormMeasurement> {
|
||||
setAdmissionMode(mode)
|
||||
const admissionEvents: GitAdmissionEvent[] = []
|
||||
const admissionClockStartedAt = performance.now()
|
||||
_resetGitAdmissionForTests(
|
||||
new GitAdmissionScheduler({
|
||||
now: () => Math.min(performance.now() - admissionClockStartedAt, GIT_ADMISSION_AGING_MS - 1),
|
||||
onAdmissionEvent: (event) => admissionEvents.push(event)
|
||||
})
|
||||
)
|
||||
const root = await mkdtemp(path.join(tmpdir(), `orca-git-storm-${mode}-`))
|
||||
tempRoots.push(root)
|
||||
const stateDir = path.join(root, 'state')
|
||||
await mkdir(stateDir)
|
||||
const binDir = await createStubGit(root)
|
||||
const repoDirs = await Promise.all(
|
||||
Array.from({ length: 6 }, async (_, index) => {
|
||||
const repoDir = path.join(root, `repo-${index}`)
|
||||
await mkdir(repoDir)
|
||||
return repoDir
|
||||
})
|
||||
)
|
||||
const baseEnv = { ...process.env, PATH: `${binDir}${path.delimiter}${process.env.PATH ?? ''}` }
|
||||
const startedAt = performance.now()
|
||||
const drifts: number[] = []
|
||||
let nextJankSample = performance.now() + 50
|
||||
const jankTimer = setInterval(() => {
|
||||
const now = performance.now()
|
||||
drifts.push(Math.max(0, now - nextJankSample))
|
||||
nextJankSample = now + 50
|
||||
}, 50)
|
||||
let maxConcurrentChildren = 0
|
||||
const censusTimer = setInterval(() => {
|
||||
void liveChildCount(stateDir).then((count) => {
|
||||
maxConcurrentChildren = Math.max(maxConcurrentChildren, count)
|
||||
})
|
||||
}, 5)
|
||||
|
||||
const background = Array.from({ length: 60 }, (_, index) =>
|
||||
gitExecFileAsync(['status', '--porcelain=v2', `storm-${index}`], {
|
||||
cwd: repoDirs[index % repoDirs.length],
|
||||
env: {
|
||||
...baseEnv,
|
||||
ORCA_STUB_ID: `background-${index}`,
|
||||
ORCA_STUB_SLEEP_MS: index % 10 === 0 ? '5000' : '200',
|
||||
ORCA_STUB_STATE_DIR: stateDir
|
||||
},
|
||||
admissionTier: 'background'
|
||||
})
|
||||
)
|
||||
await chmod(stub, 0o755)
|
||||
const controller = new AbortController()
|
||||
const commands = new Map<string, ReturnType<typeof gitExecFileAsync>>()
|
||||
const interactiveQueueSnapshots: InteractiveQueueSnapshot[] = []
|
||||
const interactiveLatencies = await Promise.all(
|
||||
Array.from(
|
||||
{ length: 10 },
|
||||
(_, index) =>
|
||||
new Promise<number>((resolve, reject) => {
|
||||
setTimeout(
|
||||
() => {
|
||||
void (async () => {
|
||||
const concurrentAtInjection = await liveChildCount(stateDir)
|
||||
const commandStartedAt = performance.now()
|
||||
const commandLabel = `interactive-${index}`
|
||||
interactiveQueueSnapshots.push({
|
||||
commandLabel,
|
||||
backgroundWaiterIds: _gitAdmissionSnapshotForTests()
|
||||
.queuedWaiters.filter((waiter) => waiter.tier === 'background')
|
||||
.map((waiter) => waiter.id)
|
||||
})
|
||||
await gitExecFileAsync(['rev-parse', commandLabel], {
|
||||
cwd: repoDirs[index % repoDirs.length],
|
||||
env: {
|
||||
...baseEnv,
|
||||
ORCA_STUB_ID: `interactive-${index}`,
|
||||
ORCA_STUB_SLEEP_MS: String(Math.max(10, concurrentAtInjection * 12)),
|
||||
ORCA_STUB_STATE_DIR: stateDir
|
||||
},
|
||||
admissionTier: 'interactive'
|
||||
})
|
||||
resolve(performance.now() - commandStartedAt)
|
||||
})().catch(reject)
|
||||
},
|
||||
50 * (index + 1)
|
||||
)
|
||||
})
|
||||
)
|
||||
)
|
||||
await Promise.all(background)
|
||||
clearInterval(censusTimer)
|
||||
clearInterval(jankTimer)
|
||||
maxConcurrentChildren = Math.max(maxConcurrentChildren, await liveChildCount(stateDir))
|
||||
const totalWallMs = performance.now() - startedAt
|
||||
return {
|
||||
mode,
|
||||
maxConcurrentChildren,
|
||||
eventLoopMaxDriftMs: Math.max(0, ...drifts),
|
||||
eventLoopP99DriftMs: percentile(drifts, 0.99),
|
||||
interactiveP50Ms: percentile(interactiveLatencies, 0.5),
|
||||
interactiveP95Ms: percentile(interactiveLatencies, 0.95),
|
||||
totalWallMs,
|
||||
admissionEvents,
|
||||
interactiveQueueSnapshots
|
||||
}
|
||||
}
|
||||
|
||||
function formatMeasurementTable(rows: readonly StormMeasurement[]): string {
|
||||
const rounded = rows.map((row) => ({
|
||||
mode: row.mode,
|
||||
maxConcurrentChildren: row.maxConcurrentChildren,
|
||||
eventLoopMaxDriftMs: row.eventLoopMaxDriftMs.toFixed(1),
|
||||
eventLoopP99DriftMs: row.eventLoopP99DriftMs.toFixed(1),
|
||||
interactiveP50Ms: row.interactiveP50Ms.toFixed(1),
|
||||
interactiveP95Ms: row.interactiveP95Ms.toFixed(1),
|
||||
totalWallMs: row.totalWallMs.toFixed(1)
|
||||
}))
|
||||
return JSON.stringify(rounded)
|
||||
}
|
||||
|
||||
function assertAdmissionLedger(measurement: StormMeasurement): void {
|
||||
const activeWaiters = new Set<number>()
|
||||
const grantSequenceByWaiter = new Map<number, number>()
|
||||
const grantByLabel = new Map<string, GitAdmissionEvent>()
|
||||
|
||||
measurement.admissionEvents.forEach((event, index) => {
|
||||
expect(event.sequence).toBe(index)
|
||||
if (event.phase === 'grant') {
|
||||
activeWaiters.add(event.waiterId)
|
||||
grantSequenceByWaiter.set(event.waiterId, event.sequence)
|
||||
const label = event.args.find((arg) => arg.startsWith('interactive-'))
|
||||
if (label) {
|
||||
grantByLabel.set(label, event)
|
||||
const start = (id: string, tier: 'background' | 'interactive'): void => {
|
||||
if (tier === 'interactive') {
|
||||
interactiveQueueSnapshots.push({
|
||||
commandLabel: id,
|
||||
backgroundWaiterIds: _gitAdmissionSnapshotForTests()
|
||||
.queuedWaiters.filter((waiter) => waiter.tier === 'background')
|
||||
.map((waiter) => waiter.id)
|
||||
})
|
||||
}
|
||||
const command = gitExecFileAsync([tier === 'background' ? 'status' : 'rev-parse', id], {
|
||||
cwd: root,
|
||||
admissionTier: tier,
|
||||
signal: controller.signal,
|
||||
env: {
|
||||
...process.env,
|
||||
PATH: `${binDir}${path.delimiter}${process.env.PATH ?? ''}`,
|
||||
ORCA_STUB_ID: id,
|
||||
ORCA_STUB_STATE_DIR: stateDir,
|
||||
ORCA_STUB_GATE_DIR: gateDir
|
||||
}
|
||||
} else {
|
||||
expect(activeWaiters.delete(event.waiterId)).toBe(true)
|
||||
})
|
||||
void command.catch(() => {})
|
||||
commands.set(id, command)
|
||||
}
|
||||
const release = async (id: string): Promise<void> => {
|
||||
await waitForLiveChild(stateDir, id, controller.signal)
|
||||
await writeFile(path.join(gateDir, id), 'release\n')
|
||||
await expect(commands.get(id)).resolves.toEqual({
|
||||
stdout: `stub:${id.startsWith('background-') ? 'status' : 'rev-parse'} ${id}\n`,
|
||||
stderr: ''
|
||||
})
|
||||
}
|
||||
try {
|
||||
backgroundIds.forEach((id) => start(id, 'background'))
|
||||
expect(_gitAdmissionSnapshotForTests().queued).toBe(4)
|
||||
headroomIds.forEach((id) => start(id, 'interactive'))
|
||||
queuedIds.forEach((id) => start(id, 'interactive'))
|
||||
expect(_gitAdmissionSnapshotForTests()).toMatchObject({
|
||||
queued: 6,
|
||||
budgets: { general: { baseUsed: GENERAL_CAP, headroomUsed: GENERAL_HEADROOM } }
|
||||
})
|
||||
await Promise.all(
|
||||
[...backgroundIds.slice(0, GENERAL_CAP), ...headroomIds].map((id) =>
|
||||
waitForLiveChild(stateDir, id, controller.signal)
|
||||
)
|
||||
)
|
||||
expect(await readdir(stateDir)).toHaveLength(GENERAL_CAP + GENERAL_HEADROOM)
|
||||
for (const [index, id] of queuedIds.entries()) {
|
||||
await release(backgroundIds[index])
|
||||
await waitForLiveChild(stateDir, id, controller.signal)
|
||||
}
|
||||
expect(activeWaiters.size).toBeLessThanOrEqual(MAX_GIT_CHILDREN)
|
||||
for (const budget of event.budgets) {
|
||||
expect(budget.baseUsed).toBeLessThanOrEqual(budget.baseCapacity)
|
||||
expect(budget.headroomUsed).toBeLessThanOrEqual(budget.headroomCapacity)
|
||||
expect(_gitAdmissionSnapshotForTests().queuedWaiters.map((waiter) => waiter.tier)).toEqual([
|
||||
'background',
|
||||
'background',
|
||||
'background',
|
||||
'background'
|
||||
])
|
||||
for (const id of [...headroomIds, ...queuedIds, ...backgroundIds.slice(queuedIds.length)]) {
|
||||
await release(id)
|
||||
}
|
||||
})
|
||||
expect(activeWaiters.size).toBe(0)
|
||||
|
||||
for (const snapshot of measurement.interactiveQueueSnapshots) {
|
||||
const interactiveGrant = grantByLabel.get(snapshot.commandLabel)
|
||||
expect(interactiveGrant).toBeDefined()
|
||||
const preceded = snapshot.backgroundWaiterIds.filter(
|
||||
(waiterId) => interactiveGrant!.sequence < (grantSequenceByWaiter.get(waiterId) ?? Infinity)
|
||||
).length
|
||||
expect(preceded).toBeGreaterThanOrEqual(Math.ceil(snapshot.backgroundWaiterIds.length * 0.9))
|
||||
expect(admissionEvents.filter((event) => event.phase === 'grant')).toHaveLength(ids.length)
|
||||
expect(admissionEvents.filter((event) => event.phase === 'release')).toHaveLength(ids.length)
|
||||
assertAdmissionLedger({ admissionEvents, interactiveQueueSnapshots })
|
||||
expect(_gitAdmissionSnapshotForTests()).toMatchObject({
|
||||
queued: 0,
|
||||
budgets: { general: { baseUsed: 0, headroomUsed: 0 } }
|
||||
})
|
||||
expect(await readdir(stateDir)).toEqual([])
|
||||
} finally {
|
||||
controller.abort()
|
||||
await Promise.allSettled(commands.values())
|
||||
}
|
||||
}
|
||||
|
||||
describe.skipIf(process.platform === 'win32')('git admission storm measurement', () => {
|
||||
it('reports bounded-concurrency before and after measurements', async () => {
|
||||
const disabled = await measureStorm('disabled')
|
||||
const enabled = await measureStorm('enabled')
|
||||
console.info(`GIT_ADMISSION_STORM_MEASUREMENT=${formatMeasurementTable([disabled, enabled])}`)
|
||||
assertAdmissionLedger(enabled)
|
||||
})
|
||||
it(
|
||||
'bounds real children and drains interactive work before the background backlog',
|
||||
verifySpawnContention
|
||||
)
|
||||
|
||||
// Output parity lives in git-admission-output-parity.test.ts: it needs real git,
|
||||
// not the PATH stub, so it runs on win32 too and cannot share this file's gate.
|
||||
it.runIf(process.env.ORCA_GIT_ADMISSION_STORM_MEASUREMENT === '1')(
|
||||
'reports bounded-concurrency before and after measurements',
|
||||
async () => {
|
||||
const disabled = await measureStorm('disabled')
|
||||
const enabled = await measureStorm('enabled')
|
||||
console.info(`GIT_ADMISSION_STORM_MEASUREMENT=${formatMeasurementTable([disabled, enabled])}`)
|
||||
assertAdmissionLedger(enabled)
|
||||
}
|
||||
)
|
||||
// Real-git output parity remains in git-admission-output-parity.test.ts, including Windows.
|
||||
})
|
||||
|
||||
@@ -0,0 +1,245 @@
|
||||
/**
|
||||
* POSIX-only measurement: a PATH-injected git fixture exercises the real spawn path.
|
||||
* Timing distributions are reported for field comparison; CI assertions stay structural.
|
||||
*/
|
||||
import { chmod, mkdtemp, mkdir, readdir, rm, writeFile } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import path from 'node:path'
|
||||
import { afterEach, expect } from 'vitest'
|
||||
import { gitExecFileAsync } from './git-exec-file'
|
||||
import {
|
||||
GIT_ADMISSION_AGING_MS,
|
||||
GitAdmissionScheduler,
|
||||
MAX_GIT_CHILDREN,
|
||||
_gitAdmissionSnapshotForTests,
|
||||
_resetGitAdmissionForTests,
|
||||
type GitAdmissionEvent
|
||||
} from './git-subprocess-admission'
|
||||
|
||||
type StormMeasurement = {
|
||||
mode: 'disabled' | 'enabled'
|
||||
maxConcurrentChildren: number
|
||||
eventLoopMaxDriftMs: number
|
||||
eventLoopP99DriftMs: number
|
||||
interactiveP50Ms: number
|
||||
interactiveP95Ms: number
|
||||
totalWallMs: number
|
||||
admissionEvents: GitAdmissionEvent[]
|
||||
interactiveQueueSnapshots: InteractiveQueueSnapshot[]
|
||||
}
|
||||
|
||||
export type InteractiveQueueSnapshot = {
|
||||
commandLabel: string
|
||||
backgroundWaiterIds: number[]
|
||||
}
|
||||
|
||||
const tempRoots: string[] = []
|
||||
const originalAdmissionDisabled = process.env.ORCA_GIT_ADMISSION_DISABLED
|
||||
|
||||
afterEach(async () => {
|
||||
if (originalAdmissionDisabled === undefined) {
|
||||
delete process.env.ORCA_GIT_ADMISSION_DISABLED
|
||||
} else {
|
||||
process.env.ORCA_GIT_ADMISSION_DISABLED = originalAdmissionDisabled
|
||||
}
|
||||
_resetGitAdmissionForTests()
|
||||
await Promise.all(tempRoots.splice(0).map((root) => rm(root, { recursive: true, force: true })))
|
||||
})
|
||||
|
||||
export async function createStormRoot(label: string): Promise<string> {
|
||||
const root = await mkdtemp(path.join(tmpdir(), `orca-git-storm-${label}-`))
|
||||
tempRoots.push(root)
|
||||
return root
|
||||
}
|
||||
|
||||
function percentile(values: readonly number[], percentileValue: number): number {
|
||||
if (values.length === 0) {
|
||||
return 0
|
||||
}
|
||||
const sorted = [...values].sort((a, b) => a - b)
|
||||
return sorted[Math.min(sorted.length - 1, Math.ceil(sorted.length * percentileValue) - 1)]
|
||||
}
|
||||
|
||||
async function liveChildCount(stateDir: string): Promise<number> {
|
||||
return (await readdir(stateDir)).filter((name) => name.endsWith('.live')).length
|
||||
}
|
||||
|
||||
async function createStubGit(root: string): Promise<string> {
|
||||
const binDir = path.join(root, 'bin')
|
||||
await mkdir(binDir)
|
||||
const stubPath = path.join(binDir, 'git')
|
||||
await writeFile(
|
||||
stubPath,
|
||||
`#!/bin/sh
|
||||
set -eu
|
||||
live="$ORCA_STUB_STATE_DIR/$ORCA_STUB_ID.live"
|
||||
: > "$live"
|
||||
trap 'rm -f "$live"' EXIT HUP INT TERM
|
||||
sleep "$(awk "BEGIN { print $ORCA_STUB_SLEEP_MS / 1000 }")"
|
||||
printf 'stub:%s\\n' "$*"
|
||||
`
|
||||
)
|
||||
await chmod(stubPath, 0o755)
|
||||
return binDir
|
||||
}
|
||||
|
||||
function setAdmissionMode(mode: StormMeasurement['mode']): void {
|
||||
if (mode === 'disabled') {
|
||||
process.env.ORCA_GIT_ADMISSION_DISABLED = '1'
|
||||
} else {
|
||||
delete process.env.ORCA_GIT_ADMISSION_DISABLED
|
||||
}
|
||||
}
|
||||
|
||||
export async function measureStorm(mode: StormMeasurement['mode']): Promise<StormMeasurement> {
|
||||
setAdmissionMode(mode)
|
||||
const admissionEvents: GitAdmissionEvent[] = []
|
||||
const admissionClockStartedAt = performance.now()
|
||||
_resetGitAdmissionForTests(
|
||||
new GitAdmissionScheduler({
|
||||
now: () => Math.min(performance.now() - admissionClockStartedAt, GIT_ADMISSION_AGING_MS - 1),
|
||||
onAdmissionEvent: (event) => admissionEvents.push(event)
|
||||
})
|
||||
)
|
||||
const root = await createStormRoot(mode)
|
||||
const stateDir = path.join(root, 'state')
|
||||
await mkdir(stateDir)
|
||||
const binDir = await createStubGit(root)
|
||||
const repoDirs = await Promise.all(
|
||||
Array.from({ length: 6 }, async (_, index) => {
|
||||
const repoDir = path.join(root, `repo-${index}`)
|
||||
await mkdir(repoDir)
|
||||
return repoDir
|
||||
})
|
||||
)
|
||||
const baseEnv = { ...process.env, PATH: `${binDir}${path.delimiter}${process.env.PATH ?? ''}` }
|
||||
const startedAt = performance.now()
|
||||
const drifts: number[] = []
|
||||
let nextJankSample = performance.now() + 50
|
||||
const jankTimer = setInterval(() => {
|
||||
const now = performance.now()
|
||||
drifts.push(Math.max(0, now - nextJankSample))
|
||||
nextJankSample = now + 50
|
||||
}, 50)
|
||||
let maxConcurrentChildren = 0
|
||||
const censusTimer = setInterval(() => {
|
||||
void liveChildCount(stateDir).then((count) => {
|
||||
maxConcurrentChildren = Math.max(maxConcurrentChildren, count)
|
||||
})
|
||||
}, 5)
|
||||
|
||||
const background = Array.from({ length: 60 }, (_, index) =>
|
||||
gitExecFileAsync(['status', '--porcelain=v2', `storm-${index}`], {
|
||||
cwd: repoDirs[index % repoDirs.length],
|
||||
env: {
|
||||
...baseEnv,
|
||||
ORCA_STUB_ID: `background-${index}`,
|
||||
ORCA_STUB_SLEEP_MS: index % 10 === 0 ? '5000' : '200',
|
||||
ORCA_STUB_STATE_DIR: stateDir
|
||||
},
|
||||
admissionTier: 'background'
|
||||
})
|
||||
)
|
||||
const interactiveQueueSnapshots: InteractiveQueueSnapshot[] = []
|
||||
const interactiveLatencies = await Promise.all(
|
||||
Array.from(
|
||||
{ length: 10 },
|
||||
(_, index) =>
|
||||
new Promise<number>((resolve, reject) => {
|
||||
setTimeout(
|
||||
() => {
|
||||
void (async () => {
|
||||
const concurrentAtInjection = await liveChildCount(stateDir)
|
||||
const commandStartedAt = performance.now()
|
||||
const commandLabel = `interactive-${index}`
|
||||
interactiveQueueSnapshots.push({
|
||||
commandLabel,
|
||||
backgroundWaiterIds: _gitAdmissionSnapshotForTests()
|
||||
.queuedWaiters.filter((waiter) => waiter.tier === 'background')
|
||||
.map((waiter) => waiter.id)
|
||||
})
|
||||
await gitExecFileAsync(['rev-parse', commandLabel], {
|
||||
cwd: repoDirs[index % repoDirs.length],
|
||||
env: {
|
||||
...baseEnv,
|
||||
ORCA_STUB_ID: `interactive-${index}`,
|
||||
ORCA_STUB_SLEEP_MS: String(Math.max(10, concurrentAtInjection * 12)),
|
||||
ORCA_STUB_STATE_DIR: stateDir
|
||||
},
|
||||
admissionTier: 'interactive'
|
||||
})
|
||||
resolve(performance.now() - commandStartedAt)
|
||||
})().catch(reject)
|
||||
},
|
||||
50 * (index + 1)
|
||||
)
|
||||
})
|
||||
)
|
||||
)
|
||||
await Promise.all(background)
|
||||
clearInterval(censusTimer)
|
||||
clearInterval(jankTimer)
|
||||
maxConcurrentChildren = Math.max(maxConcurrentChildren, await liveChildCount(stateDir))
|
||||
const totalWallMs = performance.now() - startedAt
|
||||
return {
|
||||
mode,
|
||||
maxConcurrentChildren,
|
||||
eventLoopMaxDriftMs: Math.max(0, ...drifts),
|
||||
eventLoopP99DriftMs: percentile(drifts, 0.99),
|
||||
interactiveP50Ms: percentile(interactiveLatencies, 0.5),
|
||||
interactiveP95Ms: percentile(interactiveLatencies, 0.95),
|
||||
totalWallMs,
|
||||
admissionEvents,
|
||||
interactiveQueueSnapshots
|
||||
}
|
||||
}
|
||||
|
||||
export function formatMeasurementTable(rows: readonly StormMeasurement[]): string {
|
||||
const rounded = rows.map((row) => ({
|
||||
mode: row.mode,
|
||||
maxConcurrentChildren: row.maxConcurrentChildren,
|
||||
eventLoopMaxDriftMs: row.eventLoopMaxDriftMs.toFixed(1),
|
||||
eventLoopP99DriftMs: row.eventLoopP99DriftMs.toFixed(1),
|
||||
interactiveP50Ms: row.interactiveP50Ms.toFixed(1),
|
||||
interactiveP95Ms: row.interactiveP95Ms.toFixed(1),
|
||||
totalWallMs: row.totalWallMs.toFixed(1)
|
||||
}))
|
||||
return JSON.stringify(rounded)
|
||||
}
|
||||
|
||||
export function assertAdmissionLedger(
|
||||
measurement: Pick<StormMeasurement, 'admissionEvents' | 'interactiveQueueSnapshots'>
|
||||
): void {
|
||||
const activeWaiters = new Set<number>()
|
||||
const grantSequenceByWaiter = new Map<number, number>()
|
||||
const grantByLabel = new Map<string, GitAdmissionEvent>()
|
||||
|
||||
measurement.admissionEvents.forEach((event, index) => {
|
||||
expect(event.sequence).toBe(index)
|
||||
if (event.phase === 'grant') {
|
||||
activeWaiters.add(event.waiterId)
|
||||
grantSequenceByWaiter.set(event.waiterId, event.sequence)
|
||||
const label = event.args.find((arg) => arg.startsWith('interactive-'))
|
||||
if (label) {
|
||||
grantByLabel.set(label, event)
|
||||
}
|
||||
} else {
|
||||
expect(activeWaiters.delete(event.waiterId)).toBe(true)
|
||||
}
|
||||
expect(activeWaiters.size).toBeLessThanOrEqual(MAX_GIT_CHILDREN)
|
||||
for (const budget of event.budgets) {
|
||||
expect(budget.baseUsed).toBeLessThanOrEqual(budget.baseCapacity)
|
||||
expect(budget.headroomUsed).toBeLessThanOrEqual(budget.headroomCapacity)
|
||||
}
|
||||
})
|
||||
expect(activeWaiters.size).toBe(0)
|
||||
|
||||
for (const snapshot of measurement.interactiveQueueSnapshots) {
|
||||
const interactiveGrant = grantByLabel.get(snapshot.commandLabel)
|
||||
expect(interactiveGrant).toBeDefined()
|
||||
const preceded = snapshot.backgroundWaiterIds.filter(
|
||||
(waiterId) => interactiveGrant!.sequence < (grantSequenceByWaiter.get(waiterId) ?? Infinity)
|
||||
).length
|
||||
expect(preceded).toBeGreaterThanOrEqual(Math.ceil(snapshot.backgroundWaiterIds.length * 0.9))
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,9 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { createTranscriptPane, TRANSCRIPT_PANE_PTY_ID } from './agent-transcript-pane-test-harness'
|
||||
import {
|
||||
createTranscriptPane,
|
||||
TRANSCRIPT_PANE_PTY_ID,
|
||||
waitForTranscriptIdle
|
||||
} from './agent-transcript-pane-test-harness'
|
||||
import {
|
||||
finalReplayFrame,
|
||||
readRuntimeFixture,
|
||||
@@ -112,9 +116,8 @@ describe('Antigravity 1.2.14 readiness from captured bytes', () => {
|
||||
data: `${String.fromCharCode(27)}]0;agy${String.fromCharCode(7)}${readRuntimeFixture('antigravity-1-2-14-model-picker')}`,
|
||||
size: { cols: 120, rows: 40 }
|
||||
})
|
||||
await expect(
|
||||
runtime.waitForTerminal(handle, { condition: 'tui-idle', timeoutMs: 5_000 })
|
||||
).rejects.toThrow(/timeout/)
|
||||
await runtime.readTerminal(handle, { screen: true })
|
||||
await expect(waitForTranscriptIdle({ runtime, handle }, 5_000)).rejects.toThrow(/timeout/)
|
||||
}, 15_000)
|
||||
|
||||
// Why this recording: only the screen reads it ready, so settling proves the grid is trusted.
|
||||
@@ -133,11 +136,8 @@ describe('Antigravity 1.2.14 readiness from captured bytes', () => {
|
||||
runtime.onExternalPtyResize(TRANSCRIPT_PANE_PTY_ID, cols, rows)
|
||||
}
|
||||
const settles = async () =>
|
||||
(
|
||||
await runtime
|
||||
.waitForTerminal(handle, { condition: 'tui-idle', timeoutMs: 5_000 })
|
||||
.catch(() => ({ satisfied: false }))
|
||||
).satisfied
|
||||
(await waitForTranscriptIdle({ runtime, handle }, 5_000).catch(() => ({ satisfied: false })))
|
||||
.satisfied
|
||||
options.size = { cols: 120, rows: 40 }
|
||||
runtime.reflowHeadlessTerminalToPtyGrid(TRANSCRIPT_PANE_PTY_ID, 120, 40)
|
||||
await runtime.readTerminal(handle, { screen: true })
|
||||
|
||||
@@ -10,12 +10,15 @@ import {
|
||||
import { deserializeAgentChildWorkBindingKey } from './agent-status-child-work-binding'
|
||||
import { agentChildWorkBelongsTo, type AgentChildWorkRecord } from './agent-status-child-work'
|
||||
import { parseAgentStatusStoreMutation } from './agent-status-store-codec'
|
||||
import type { AgentStatusStoreSnapshot } from './agent-status-store-contract'
|
||||
import {
|
||||
AGENT_STATUS_STORE_SNAPSHOT_VERSION,
|
||||
type AgentStatusStoreSnapshot
|
||||
} from './agent-status-store-contract'
|
||||
import { applyAgentStatusStoreMutationSteps } from './agent-status-store-mutation'
|
||||
import {
|
||||
cloneAgentStatusStoreState,
|
||||
createEmptyAgentStatusStoreState,
|
||||
snapshotFromAgentStatusStoreState,
|
||||
deepFreezeAgentStatusStoreValue,
|
||||
validateAgentStatusStoreState,
|
||||
type AgentStatusStoreState
|
||||
} from './agent-status-store-state'
|
||||
@@ -59,7 +62,7 @@ function applyByCopy(current: AgentStatusStoreState, value: unknown): AgentStatu
|
||||
export type CopyingAgentStatusStoreOracle = {
|
||||
applyMutation(mutation: unknown): boolean
|
||||
getSnapshot(): AgentStatusStoreSnapshot
|
||||
/** The oracle's own state, for checking every invariant after each step. */
|
||||
/** The oracle's own maps, for fixture coverage and diagnostics. */
|
||||
state(): AgentStatusStoreState
|
||||
getChildren(subject: AgentStatusSubject): AgentChildWorkRecord[]
|
||||
getAliasesForChild(childWorkId: string): AgentChildWorkAliasRecord[]
|
||||
@@ -76,7 +79,18 @@ export function createCopyingAgentStatusStoreOracle(epoch: string): CopyingAgent
|
||||
}
|
||||
return next !== null
|
||||
},
|
||||
getSnapshot: () => snapshotFromAgentStatusStoreState(current),
|
||||
// Project the validated maps independently of the production snapshot parser.
|
||||
getSnapshot: () =>
|
||||
deepFreezeAgentStatusStoreValue({
|
||||
version: AGENT_STATUS_STORE_SNAPSHOT_VERSION,
|
||||
epoch: current.epoch,
|
||||
revision: current.revision,
|
||||
parents: [...current.parents.values()],
|
||||
children: [...current.children.values()],
|
||||
aliases: [...current.aliases.values()],
|
||||
facts: [...current.facts.values()],
|
||||
tombstones: [...current.tombstones.values()]
|
||||
}),
|
||||
state: () => current,
|
||||
getChildren: (subject) =>
|
||||
[...current.children.values()].filter((child) => agentChildWorkBelongsTo(child, subject)),
|
||||
|
||||
@@ -8,16 +8,17 @@ import type { AgentChildWorkInput } from './agent-status-child-work'
|
||||
import { createAgentStatusStore } from './agent-status-store'
|
||||
import { parseAgentStatusStoreMutation } from './agent-status-store-codec'
|
||||
import { commitAgentStatusStoreMutation } from './agent-status-store-commit'
|
||||
import type { AgentStatusStoreMutation } from './agent-status-store-contract'
|
||||
import {
|
||||
AGENT_STATUS_STORE_LIMITS,
|
||||
AGENT_STATUS_STORE_TOMBSTONE_RETENTION_REVISIONS,
|
||||
type AgentStatusStoreMutation
|
||||
} from './agent-status-store-contract'
|
||||
import { createCopyingAgentStatusStoreOracle } from './agent-status-store-copying-oracle.test-fixture'
|
||||
import {
|
||||
indexAgentStatusStoreState,
|
||||
type AgentStatusStoreIndexes
|
||||
} from './agent-status-store-indexes'
|
||||
import {
|
||||
createEmptyAgentStatusStoreState,
|
||||
validateAgentStatusStoreState
|
||||
} from './agent-status-store-state'
|
||||
import { createEmptyAgentStatusStoreState } from './agent-status-store-state'
|
||||
import {
|
||||
makeStructuredAgentStatusSubject,
|
||||
serializeAgentStatusSubject,
|
||||
@@ -163,7 +164,6 @@ describe('AgentStatusStore applied in place', () => {
|
||||
accepted += 1
|
||||
expect(replica.applyTransportEnvelope(envelope)).toBe(true)
|
||||
}
|
||||
expect(validateAgentStatusStoreState(oracle.state())).toBe(true)
|
||||
expect(JSON.stringify(store.getSnapshot())).toBe(JSON.stringify(oracle.getSnapshot()))
|
||||
for (const parent of parents) {
|
||||
expect(store.getChildren(parent)).toEqual(oracle.getChildren(parent))
|
||||
@@ -206,18 +206,51 @@ describe('AgentStatusStore applied in place', () => {
|
||||
60_000
|
||||
)
|
||||
|
||||
it('compacts tombstones exactly as the copying store does past the retention window', () => {
|
||||
it('retires tombstones exactly at the retention revision boundary', () => {
|
||||
const store = createAgentStatusStore({ epoch: 'epoch-a', mode: 'authority' })
|
||||
const oracle = createCopyingAgentStatusStoreOracle('epoch-a')
|
||||
const parent = parents[0]!
|
||||
for (let step = 0; step < 5_000; step += 1) {
|
||||
const mutation: AgentStatusStoreMutation =
|
||||
step % 2 === 0
|
||||
? { parent: { subject: parent }, facts: [{ subject: parent, key: `k${step}`, value: 1 }] }
|
||||
: { removeFacts: [{ subject: parent, key: `k${step - 1}` }] }
|
||||
for (const childWorkId of ['oldest', 'newer']) {
|
||||
const mutation = { removeChildren: [childWorkId] }
|
||||
expect(store.applyMutation(mutation) !== null).toBe(oracle.applyMutation(mutation))
|
||||
}
|
||||
expect(oracle.getSnapshot().tombstones.length).toBeGreaterThan(1_000)
|
||||
const advance = { parent: { subject: parent } }
|
||||
for (
|
||||
let revision = 3;
|
||||
revision <= AGENT_STATUS_STORE_TOMBSTONE_RETENTION_REVISIONS;
|
||||
revision += 1
|
||||
) {
|
||||
expect(store.applyMutation(advance) !== null).toBe(oracle.applyMutation(advance))
|
||||
}
|
||||
expect(store.getSnapshot().revision).toBe(AGENT_STATUS_STORE_TOMBSTONE_RETENTION_REVISIONS)
|
||||
expect(store.getSnapshot().tombstones).toEqual([
|
||||
{ entity: 'child', key: 'oldest', revision: 1 },
|
||||
{ entity: 'child', key: 'newer', revision: 2 }
|
||||
])
|
||||
expect(JSON.stringify(store.getSnapshot())).toBe(JSON.stringify(oracle.getSnapshot()))
|
||||
|
||||
expect(store.applyMutation(advance) !== null).toBe(oracle.applyMutation(advance))
|
||||
expect(store.getSnapshot().tombstones).toEqual([{ entity: 'child', key: 'newer', revision: 2 }])
|
||||
expect(JSON.stringify(store.getSnapshot())).toBe(JSON.stringify(oracle.getSnapshot()))
|
||||
|
||||
expect(store.applyMutation(advance) !== null).toBe(oracle.applyMutation(advance))
|
||||
expect(store.getSnapshot().tombstones).toEqual([])
|
||||
expect(JSON.stringify(store.getSnapshot())).toBe(JSON.stringify(oracle.getSnapshot()))
|
||||
}, 60_000)
|
||||
|
||||
it('keeps the newest tombstones in order when their count overflows', () => {
|
||||
const store = createAgentStatusStore({ epoch: 'epoch-a', mode: 'authority' })
|
||||
const oracle = createCopyingAgentStatusStoreOracle('epoch-a')
|
||||
const childWorkIds = Array.from(
|
||||
{ length: AGENT_STATUS_STORE_LIMITS.tombstones },
|
||||
(_, index) => `removed-${index}`
|
||||
)
|
||||
for (const mutation of [{ removeChildren: ['oldest'] }, { removeChildren: childWorkIds }]) {
|
||||
expect(store.applyMutation(mutation) !== null).toBe(oracle.applyMutation(mutation))
|
||||
}
|
||||
expect(store.getSnapshot().tombstones).toEqual(
|
||||
childWorkIds.map((key) => ({ entity: 'child', key, revision: 2 }))
|
||||
)
|
||||
expect(JSON.stringify(store.getSnapshot())).toBe(JSON.stringify(oracle.getSnapshot()))
|
||||
}, 60_000)
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user