From 674ba5cfb3ea436e2ca23da3f730430e7ad50e26 Mon Sep 17 00:00:00 2001 From: Christopher Nelson Date: Fri, 2 Oct 2026 19:44:32 -0400 Subject: [PATCH] feat: add batch-capable reflex provider execution --- README.md | 7 +- ROADMAP.md | 13 +- apps/game-api/src/attempt-accounting.test.ts | 61 ++ apps/game-api/src/attempt-accounting.ts | 36 + apps/game-api/src/experiment-export.test.ts | 48 + apps/game-api/src/experiment-export.ts | 26 +- .../src/reflex-execution.batch.test.ts | 412 +++++++++ apps/game-api/src/reflex-execution.ts | 834 +++++++++++++++--- .../src/simulation-service.batch.test.ts | 375 ++++++++ apps/game-api/src/simulation-service.ts | 95 +- apps/game-api/src/swarm-comparison.ts | 5 +- .../src/components/swarm-view.test.tsx | 10 + apps/world-lab/src/components/swarm-view.tsx | 6 + docs/ARCHITECTURE.md | 30 +- docs/EXPERIMENT_ARCHIVE.md | 18 +- docs/SECURITY.md | 16 +- docs/TESTING.md | 15 +- .../adr/0035-batch-capable-reflex-provider.md | 78 ++ packages/agent-runtime/src/reflex-provider.ts | 33 + .../experiment-archive/src/archive.test.ts | 124 ++- .../experiment-archive/src/query-service.ts | 116 ++- packages/shared/src/index.test.ts | 88 ++ packages/shared/src/index.ts | 98 +- tests/e2e/world-lab.spec.ts | 2 +- 24 files changed, 2314 insertions(+), 232 deletions(-) create mode 100644 apps/game-api/src/reflex-execution.batch.test.ts create mode 100644 apps/game-api/src/simulation-service.batch.test.ts create mode 100644 docs/adr/0035-batch-capable-reflex-provider.md diff --git a/README.md b/README.md index dd51ddf..7ae578e 100644 --- a/README.md +++ b/README.md @@ -23,7 +23,8 @@ provider usage. Every agent is visible. or agent IDs. Other ticks reuse the current directives without a planner call. 2. **Workers resolve directives through TypeSafe Jev** using compact observations and opaque, engine-legal action candidates. Workers share a frozen pre-action - world; their calls currently run sequentially under one tick deadline. + world; a bounded adapter runs individual calls under one tick deadline. + Jev remains sequential by default; the seam also accepts keyed native batches. 3. **The deterministic world engine validates and resolves actions** in seeded order. A complete tick commits atomically; cancellation commits no world changes. Provider failures retain safe attempt records and use explicit @@ -89,7 +90,7 @@ are in [the screenshot guide](docs/assets/README.md). - **Inspection:** switch between Live and Agents while the same execution controller stays mounted. Inspect Zero strategy, worker directives, reflex choices, validation outcomes, territory, and safe activity records. -- **Research exports:** generate compact or pretty schema-v13 JSON, download +- **Research exports:** generate compact or pretty schema-v14 JSON, download it, or manually save the exact generated artifact to local SQLite. Exports include bounded safe tick and provider-attempt records, including attempts that did not produce a committed tick. @@ -164,7 +165,7 @@ probe. Neither that probe nor `compare:live` runs in default tests or CI. - [ADR 0033](docs/adr/0033-retire-legacy-multi-agent-architecture.md) — retirement of the previous architecture -Current code reads only schema-v13 exports. Pre-swarm scenarios, snapshots, and +Current code reads only schema-v14 exports. Pre-swarm scenarios, snapshots, and exports require an older Git revision. Historical ADRs and experiment reports remain as decision history. diff --git a/ROADMAP.md b/ROADMAP.md index 9507e75..733bb47 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -56,6 +56,15 @@ call) is retained as the ablation control for the swarm comparisons, not as a second production architecture. Historical milestones below remain as implementation history. +## Reflex-provider follow-up — batch-capable seam + +The local PR 3 change adds optional native-batch worker selection, independent +result validation, and a capped individual-provider adapter. Production Jev +retains cap 1; raising production concurrency is a separate rollout decision. +Schema-v14 exports identify shared batch dispatches and preserve aggregate +billing without per-worker allocations. No native production provider, Laya, +training, planner-cadence changes, or player features are added. See ADR 0035. + ## Agent Zero planner Agent Zero is the sole generative planner. It makes an OpenRouter planning call @@ -171,14 +180,14 @@ paths or SQL, recovery, scheduling, MCP, and archive authority remain deferred. Persistent short- and long-term objectives, compact memories, plan revision, summaries, and longer simulation runs. -_Note: per-agent strategic goals, the compact memory ledger, and the Behavior Trace introduced in this milestone were subsequently removed. The SQLite experiment archive (pre-PR-5 observability slice) remains current, updated to schema version 13. See ADRs 0033 and 0034._ +_Note: per-agent strategic goals, the compact memory ledger, and the Behavior Trace introduced in this milestone were subsequently removed. The SQLite experiment archive (pre-PR-5 observability slice) remains current, updated to schema version 14. See ADRs 0033–0035._ ## PR 6 — Persistent autonomous world Scheduled turns, snapshots, replay, retries, idempotency, durable budget/attempt ledgers, failure recovery, and operation without the World Lab browser being open. The current process-local attempt and credit-admission ceilings are an -operator safety boundary, and their schema-v13 safe ledger can be exported to +operator safety boundary, and their schema-v14 safe ledger can be exported to the analysis archive even when no turn committed. This is not active runtime persistence, restart recovery, or provider-account balance enforcement. diff --git a/apps/game-api/src/attempt-accounting.test.ts b/apps/game-api/src/attempt-accounting.test.ts index ac64195..948add7 100644 --- a/apps/game-api/src/attempt-accounting.test.ts +++ b/apps/game-api/src/attempt-accounting.test.ts @@ -2,6 +2,24 @@ import { describe, expect, it } from 'vitest'; import { AttemptAccounting } from './attempt-accounting'; describe('AttemptAccounting', () => { + it('earmarks only available retry capacity without poisoning admission or recording unstarted work', () => { + const accounting = new AttemptAccounting(10, '0.03', '0.01'); + expect(accounting.reserve(2)).toBe(true); + expect(accounting.reserveUpTo(7)).toBe(1); + expect(accounting.reserveUpTo(1)).toBe(0); + expect(accounting.snapshot()).toMatchObject({ + reservedPermits: 3, + attemptsStarted: 0, + exhausted: false, + }); + accounting.releaseReservations(); + expect(accounting.snapshot()).toMatchObject({ + reservedPermits: 0, + committedCreditExposure: '0', + exhausted: false, + }); + expect(accounting.ledger()).toEqual([]); + }); it('reserves atomically and finalizes known and unknown cost exactly once', () => { const accounting = new AttemptAccounting(3, '0.03', '0.01'); expect(accounting.reserve(3)).toBe(true); @@ -60,6 +78,49 @@ describe('AttemptAccounting', () => { }); }); + it('records one billed attempt for a batch and retains all participant attribution', () => { + const accounting = new AttemptAccounting(1, '0.01', '0.01'); + expect(accounting.reserve(1)).toBe(true); + const first = '11111111-1111-4111-8111-111111111111' as never; + const second = '22222222-2222-4222-8222-222222222222' as never; + const permit = accounting.startReserved({ + agentId: first, + intendedTurnNumber: 1, + intendedTickNumber: 1, + kind: 'initial', + startedAt: '2026-08-13T12:00:00.000Z', + modelId: 'jev-1.13.0' as never, + reasoningProfile: 'provider-default' as never, + batch: { + id: '018f3f38-6b7d-7db7-8e95-751b4ce2681f', + members: [ + { agentId: first, intendedTurnNumber: 1 }, + { agentId: second, intendedTurnNumber: 2 }, + ], + }, + })!; + accounting.finalize(permit, { + outcome: 'completed', + completedAt: '2026-08-13T12:00:01.000Z', + provider: { + provider: 'typesafe', + model: 'jev-1.13.0' as never, + latencyMs: 20, + costCredits: 0.007, + }, + }); + expect(accounting.ledger()).toHaveLength(1); + expect(accounting.ledger()[0]).toMatchObject({ + batch: { members: [{ agentId: first }, { agentId: second }] }, + actualCostCredits: '0.007', + }); + expect(accounting.snapshot()).toMatchObject({ + attemptsStarted: 1, + attemptsFinalized: 1, + knownFinalizedCostCredits: '0.007', + }); + }); + it('reports explicit unlimited capacity', () => { const accounting = new AttemptAccounting(null); expect(accounting.startAdditional()).not.toBeNull(); diff --git a/apps/game-api/src/attempt-accounting.ts b/apps/game-api/src/attempt-accounting.ts index 55faeff..47c6170 100644 --- a/apps/game-api/src/attempt-accounting.ts +++ b/apps/game-api/src/attempt-accounting.ts @@ -4,6 +4,7 @@ import { type ModelAttempt, type ModelId, type ProviderAttemptRecord, + type ProviderAttemptBatch, type ProviderFailure, type ProviderMetadata, type ReflexDecision, @@ -42,6 +43,7 @@ export interface AttemptStart { startedAt: string; modelId: ModelId; reasoningProfile: ReasoningProfile; + batch?: ProviderAttemptBatch; } export interface AttemptCompletion { @@ -200,6 +202,40 @@ export class AttemptAccounting { return permitId; } + /** + * Earmark available retry capacity without exhausting admission on a partial + * allocation. The caller owns a fixed slot per eligible job and must never + * reassign unused slots; all reservations are released at tick completion. + */ + reserveUpTo(count: number): number { + if (!Number.isInteger(count) || count < 0) + throw new Error('Retry reservation count must be a nonnegative integer.'); + if (this.#exhaustionReason !== null) return 0; + let available = 0; + while (available < count) { + const next = available + 1; + if ( + this.limit !== null && + this.#started + this.#reserved + next > this.limit + ) + break; + if ( + this.creditLimit !== null && + compare( + add( + this.#committedExposure, + multiply(this.reservationCreditsPerAttempt, this.#reserved + next), + ), + this.creditLimit, + ) > 0 + ) + break; + available = next; + } + if (available > 0) this.reserve(available); + return available; + } + startAdditional(details?: AttemptStart): number | null { if (!this.reserve(1)) return null; return this.startReserved(details); diff --git a/apps/game-api/src/experiment-export.test.ts b/apps/game-api/src/experiment-export.test.ts index c988e93..4b88e61 100644 --- a/apps/game-api/src/experiment-export.test.ts +++ b/apps/game-api/src/experiment-export.test.ts @@ -4,6 +4,7 @@ import { agentIdSchema, eventIdSchema, h3CellSchema, + providerAttemptRecordSchema, type AgentId, type H3Cell, type WorldAction, @@ -72,6 +73,53 @@ function rejectedMove(tickNumber: number, agentId: AgentId, to: H3Cell) { } describe('movement-pattern metrics', () => { + it('counts shared batch usage once in aggregate and excludes it from per-agent attribution', () => { + const batchAttempt = providerAttemptRecordSchema.parse({ + id: '018f3f38-6b7d-7db7-8e95-751b4ce2681e', + agentId: agentX, + intendedTurnNumber: 1, + intendedTickNumber: 1, + kind: 'initial', + startedAt: '2026-08-13T12:00:00.000Z', + completedAt: '2026-08-13T12:00:01.000Z', + outcome: 'completed', + modelId: 'jev-1.13.0', + reasoningProfile: 'provider-default', + reservedCredits: '0.01', + actualCostCredits: '0.006', + provider: { + provider: 'typesafe', + model: 'jev-1.13.0', + latencyMs: 12, + promptTokens: 30, + completionTokens: 2, + costCredits: 0.006, + }, + batch: { + id: '018f3f38-6b7d-7db7-8e95-751b4ce2681f', + members: [ + { agentId: agentX, intendedTurnNumber: 1 }, + { agentId: agentY, intendedTurnNumber: 2 }, + ], + }, + }); + const metrics = calculateExperimentMetrics( + [], + [agentX, agentY], + [batchAttempt], + ); + expect(metrics.aggregate).toMatchObject({ + modelCalls: 1, + knownCostCredits: 0.006, + tokens: { promptTokens: 30, completionTokens: 2 }, + }); + expect(metrics.byAgent.map(({ metrics: value }) => value)).toEqual( + expect.arrayContaining([ + expect.objectContaining({ modelCalls: 0, knownCostCredits: 0 }), + ]), + ); + }); + it('walks each agent path separately for direction streaks and revisits', () => { const { a, b, direction, opposite } = straightLine(); const yTarget = neighbors(origin).find( diff --git a/apps/game-api/src/experiment-export.ts b/apps/game-api/src/experiment-export.ts index e8d1c55..25433ed 100644 --- a/apps/game-api/src/experiment-export.ts +++ b/apps/game-api/src/experiment-export.ts @@ -32,7 +32,7 @@ import { } from './geographic-direction'; export interface ExperimentSource { - schemaVersion: 13; + schemaVersion: 14; id: ExperimentId; startedAt: string; providerMode: 'openrouter' | 'scripted-test'; @@ -203,7 +203,13 @@ export function createExperimentExport( const initialAgentsById = new Map( source.initialAgents.map((agent) => [agent.id, agent]), ); - const selectedAgents = selectedAgentIds + const selectedProfileIds = new Set([ + ...selectedAgentIds, + ...providerAttempts.flatMap( + (attempt) => attempt.batch?.members.map(({ agentId }) => agentId) ?? [], + ), + ]); + const selectedAgents = [...selectedProfileIds] .map( (agentId) => currentAgentsById.get(agentId) ?? initialAgentsById.get(agentId), @@ -390,8 +396,12 @@ function filterProviderAttempts( selected: Set, tickNumbers: Set | 'all', ): ProviderAttemptRecord[] { - let attempts = source.providerAttempts.filter(({ agentId }) => - selected.has(agentId), + let attempts = source.providerAttempts.filter( + ({ agentId, batch }) => + selected.has(agentId) || + Boolean( + batch?.members.some(({ agentId: memberId }) => selected.has(memberId)), + ), ); if (tickNumbers !== 'all') attempts = attempts.filter( @@ -521,7 +531,9 @@ function attemptMetrics(attempts: readonly ProviderAttemptRecord[]) { ).length; const groups = new Map(); for (const attempt of attempts) { - const key = `${attempt.agentId}:${attempt.intendedTickNumber ?? attempt.intendedTurnNumber}`; + const key = attempt.batch + ? `batch:${attempt.batch.id}` + : `${attempt.agentId}:${attempt.intendedTickNumber ?? attempt.intendedTurnNumber}`; groups.set(key, [...(groups.get(key) ?? []), attempt]); } const retriedGroups = [...groups.values()].filter((group) => @@ -745,7 +757,9 @@ export function calculateExperimentMetrics( agentId, metrics: metricCountsFor( resolvedActions.filter((action) => action.agentId === agentId), - providerAttempts.filter((attempt) => attempt.agentId === agentId), + providerAttempts.filter( + (attempt) => !attempt.batch && attempt.agentId === agentId, + ), controlChanges, [agentId], agentId, diff --git a/apps/game-api/src/reflex-execution.batch.test.ts b/apps/game-api/src/reflex-execution.batch.test.ts new file mode 100644 index 0000000..6dc62ef --- /dev/null +++ b/apps/game-api/src/reflex-execution.batch.test.ts @@ -0,0 +1,412 @@ +import { describe, expect, it, vi } from 'vitest'; +import type { + ReflexBatchAttemptFinalizer, + ReflexBatchDecisionResult, + ReflexAttemptStarter, + ReflexProvider, +} from '@hexzero/agent-runtime'; +import { + reflexDecisionSchema, + type ReflexObservation, + type SwarmDirective, +} from '@hexzero/shared'; +import { createDevelopmentWorld, toWorldState } from '@hexzero/world-engine'; +import { AttemptAccounting } from './attempt-accounting'; +import { + chooseReflexWorldActions, + compileReflexObservation, + ReflexSelectionCancelledError, + ReflexSelectionDeadlineError, +} from './reflex-execution'; + +function setup(workerCount = 3) { + const state = toWorldState(createDevelopmentWorld()); + const agents = [...state.agents.values()].slice(0, workerCount); + const inputs = agents.map((agent, index) => { + const directive: SwarmDirective = { + id: `batch-directive-${index}`, + agentId: agent.id, + mission: 'hold', + targetCell: null, + priority: 'normal', + riskTolerance: 'medium', + issuedAtTick: 0, + expiresAtTick: 5, + }; + return { + compiled: compileReflexObservation(state, directive), + intendedTickNumber: 1, + intendedTurnNumber: index + 1, + initialPermitReserved: true, + }; + }); + return { state, inputs }; +} + +function decision( + observation: ReflexObservation, + choice = observation.candidates[0]!.id, +) { + return reflexDecisionSchema.parse({ + chosenCandidateId: choice, + confidence: 1, + probabilities: Object.fromEntries( + observation.candidates.map(({ id }) => [id, id === choice ? 1 : 0]), + ), + replanProbability: 0, + model: 'jev-1.13.0', + latencyMs: 2, + inputTokens: 12, + outputTokens: 2, + directiveId: observation.directive.id, + cognitionSource: 'jev-reflex', + }); +} + +function provider( + decideBatch: NonNullable, +): ReflexProvider { + return { + mode: 'typesafe-jev', + model: 'jev-1.13.0', + configured: true, + async decide() { + throw new Error('scalar path should not be called'); + }, + decideBatch, + }; +} + +describe('batched reflex execution seam', () => { + it('closes a settled worker attempt gate while another worker remains pending', async () => { + const { inputs } = setup(2); + const accounting = new AttemptAccounting(3); + expect(accounting.reserve(2)).toBe(true); + let lateBegin: ReflexAttemptStarter | undefined; + let release!: () => void; + const scalar: ReflexProvider = { + mode: 'scripted-reflex-test', + configured: true, + async decide(observation, options) { + const finalize = options?.beginAttempt?.('initial'); + if (observation.directive.id === 'batch-directive-0') { + lateBegin = options?.beginAttempt; + } else { + await new Promise((resolve) => { + release = resolve; + }); + } + finalize?.({ + outcome: 'completed', + reflexDecision: decision(observation), + }); + return decision(observation); + }, + }; + const selecting = chooseReflexWorldActions(inputs, scalar, { + accounting, + concurrencyLimit: 2, + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(accounting.snapshot()).toMatchObject({ + attemptsStarted: 2, + attemptsFinalized: 1, + attemptsInFlight: 1, + }); + expect(lateBegin?.('automatic-transport-retry')).toBeNull(); + expect(accounting.ledger()).toHaveLength(2); + release(); + await selecting; + expect(accounting.snapshot()).toMatchObject({ + attemptsStarted: 2, + attemptsFinalized: 2, + attemptsInFlight: 0, + }); + accounting.releaseReservations(); + }); + it('finalizes failed scalar dispatches without retaining an invalid success completion field', async () => { + const { inputs } = setup(1); + const accounting = new AttemptAccounting(1); + expect(accounting.reserve(1)).toBe(true); + const scalar: ReflexProvider = { + mode: 'scripted-reflex-test', + configured: true, + async decide(observation, options) { + options?.beginAttempt?.('initial')?.({ + outcome: 'provider-error', + failure: { + code: 'provider-http', + message: 'Unavailable.', + retryable: false, + }, + reflexDecision: decision(observation), + }); + throw new Error('Unavailable'); + }, + }; + const [selection] = await chooseReflexWorldActions(inputs, scalar, { + accounting, + }); + expect(selection?.cognitionSource).toBe('deterministic-fallback'); + expect(accounting.snapshot()).toMatchObject({ + attemptsStarted: 1, + attemptsFinalized: 1, + attemptsInFlight: 0, + }); + expect(accounting.ledger()[0]).toMatchObject({ outcome: 'provider-error' }); + expect(accounting.ledger()[0]?.reflexDecision).toBeUndefined(); + }); + it('maps reordered native results by worker and records one shared batch attempt', async () => { + const { inputs } = setup(); + const accounting = new AttemptAccounting(1); + expect(accounting.reserve(1)).toBe(true); + const batched = provider(async (observations, options) => { + expect(observations).toHaveLength(inputs.length); + expect(Object.isFrozen(observations[0])).toBe(true); + expect(Object.isFrozen(observations[0]?.candidates)).toBe(true); + const finalize = options?.beginAttempt?.('initial'); + expect(finalize).not.toBeNull(); + finalize?.({ + outcome: 'completed', + provider: { + provider: 'typesafe', + model: 'jev-1.13.0', + latencyMs: 9, + promptTokens: 90, + completionTokens: 12, + totalTokens: 102, + }, + }); + return [...observations].reverse().map((observation) => ({ + agentId: observation.agentId, + status: 'completed' as const, + decision: decision(observation), + })); + }); + + const selected = await chooseReflexWorldActions(inputs, batched, { + accounting, + }); + + expect( + selected.map( + ({ decision: selectedDecision }) => selectedDecision?.directiveId, + ), + ).toEqual(inputs.map(({ compiled }) => compiled.observation.directive.id)); + const [attempt] = accounting.ledger(); + expect(accounting.ledger()).toHaveLength(1); + expect(attempt).toMatchObject({ + agentId: inputs[0]?.compiled.observation.agentId, + batch: { + members: inputs.map(({ compiled }, index) => ({ + agentId: compiled.observation.agentId, + intendedTurnNumber: index + 1, + })), + }, + provider: { promptTokens: 90, completionTokens: 12, totalTokens: 102 }, + }); + expect(attempt).not.toHaveProperty('reflexDecision'); + }); + + it('isolates failed, malformed, missing, duplicate, and extra keyed results', async () => { + const { inputs } = setup(4); + const [first, second, third, fourth] = inputs; + const results: unknown[] = [ + { + agentId: first!.compiled.observation.agentId, + status: 'completed', + decision: decision(first!.compiled.observation), + }, + { + agentId: second!.compiled.observation.agentId, + status: 'failed', + failure: { code: 'invalid-decision', message: 'bad', retryable: false }, + }, + { + agentId: third!.compiled.observation.agentId, + status: 'completed', + decision: { bogus: true }, + }, + { + agentId: first!.compiled.observation.agentId, + status: 'completed', + decision: decision(first!.compiled.observation), + }, + { agentId: '00000000-0000-4000-8000-000000000099', status: 'ignored' }, + ]; + const selected = await chooseReflexWorldActions( + inputs, + provider(async () => results as ReflexBatchDecisionResult[]), + ); + + expect(selected[0]?.cognitionSource).toBe('deterministic-fallback'); + expect(selected[0]?.failure?.code).toBe('unsupported-response'); + expect(selected[1]?.cognitionSource).toBe('deterministic-fallback'); + expect(selected[1]?.failure?.code).toBe('invalid-decision'); + expect(selected[2]?.cognitionSource).toBe('deterministic-fallback'); + expect(selected[2]?.failure?.code).toBe('malformed-response'); + expect(selected[3]?.cognitionSource).toBe('deterministic-fallback'); + expect(selected[3]?.failure?.code).toBe('unsupported-response'); + expect(fourth).toBeDefined(); + }); + + it('uses a bounded scalar adapter and preserves input result order', async () => { + const { inputs } = setup(4); + let active = 0; + let maximumActive = 0; + const pending: Array<() => void> = []; + const scalar: ReflexProvider = { + mode: 'scripted-reflex-test', + model: 'jev-1.13.0', + configured: true, + async decide(observation) { + active += 1; + maximumActive = Math.max(maximumActive, active); + await new Promise((resolve) => pending.push(resolve)); + active -= 1; + return decision(observation); + }, + }; + + const selecting = chooseReflexWorldActions(inputs, scalar, { + concurrencyLimit: 2, + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(pending).toHaveLength(2); + pending.splice(0).forEach((resolve) => resolve()); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(pending).toHaveLength(2); + pending.splice(0).forEach((resolve) => resolve()); + await new Promise((resolve) => setTimeout(resolve, 0)); + pending.splice(0).forEach((resolve) => resolve()); + + const selected = await selecting; + expect(maximumActive).toBe(2); + expect(selected.map(({ decision: item }) => item?.directiveId)).toEqual( + inputs.map(({ compiled }) => compiled.observation.directive.id), + ); + }); + + it.each([true, false])( + 'keeps scarce retry permits assigned to the input prefix (reverse completion: %s)', + async (reverseCompletion) => { + const { inputs } = setup(3); + const accounting = new AttemptAccounting(4); + expect(accounting.reserve(3)).toBe(true); + const gates: Array<() => void> = []; + const scalar: ReflexProvider = { + mode: 'scripted-reflex-test', + model: 'jev-1.13.0', + configured: true, + async decide(observation, options) { + const initial = options?.beginAttempt?.('initial'); + if (!initial) throw new Error('initial permit unavailable'); + initial({ outcome: 'provider-error' }); + expect(options?.beginAttempt?.('initial')).toBeNull(); + if (observation.directive.id !== 'batch-directive-2') + await new Promise((resolve) => gates.push(resolve)); + const retry = options?.beginAttempt?.('automatic-transport-retry'); + if (!retry) throw new Error('retry permit unavailable'); + retry({ + outcome: 'completed', + provider: { + provider: 'typesafe', + model: 'jev-1.13.0', + latencyMs: 1, + }, + }); + expect( + options?.beginAttempt?.('automatic-transport-retry'), + ).toBeNull(); + return decision(observation); + }, + }; + + const selecting = chooseReflexWorldActions(inputs, scalar, { + accounting, + concurrencyLimit: 2, + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(gates).toHaveLength(2); + if (reverseCompletion) { + gates[1]!(); + await new Promise((resolve) => setTimeout(resolve, 0)); + gates[0]!(); + } else { + gates[0]!(); + await new Promise((resolve) => setTimeout(resolve, 0)); + gates[1]!(); + } + const selected = await selecting; + + expect(selected[0]?.cognitionSource).toBe('jev-reflex'); + expect(selected[1]?.cognitionSource).toBe('deterministic-fallback'); + expect(selected[2]?.cognitionSource).toBe('deterministic-fallback'); + expect(accounting.ledger()).toHaveLength(4); + expect( + accounting + .ledger() + .filter(({ kind }) => kind === 'automatic-transport-retry') + .map(({ agentId }) => agentId), + ).toEqual([inputs[0]?.compiled.observation.agentId]); + accounting.releaseReservations(); + }, + ); + + it('cancels promptly and finalizes an open batch attempt once', async () => { + const { inputs } = setup(2); + const accounting = new AttemptAccounting(1); + expect(accounting.reserve(1)).toBe(true); + const controller = new AbortController(); + let release!: (results: readonly ReflexBatchDecisionResult[]) => void; + let lateFinalize: ReflexBatchAttemptFinalizer | null | undefined; + const selecting = chooseReflexWorldActions( + inputs, + provider(async (_observations, options) => { + const finalize = options?.beginAttempt?.('initial'); + lateFinalize = finalize; + return await new Promise((resolve) => { + release = resolve; + }); + }), + { accounting, signal: controller.signal }, + ); + await new Promise((resolve) => setTimeout(resolve, 0)); + controller.abort(); + await expect(selecting).rejects.toBeInstanceOf( + ReflexSelectionCancelledError, + ); + expect(accounting.ledger()).toHaveLength(1); + expect(accounting.ledger()[0]?.outcome).toBe('cancelled'); + lateFinalize?.({ outcome: 'completed' }); + expect(accounting.ledger()[0]?.outcome).toBe('cancelled'); + release([]); + }); + + it('times out the whole batch when the provider ignores abort', async () => { + const { inputs } = setup(2); + vi.useFakeTimers(); + const selecting = chooseReflexWorldActions( + inputs, + provider(async () => await new Promise(() => undefined)), + { deadlineAtMs: Date.now() + 20 }, + ); + const rejected = expect(selecting).rejects.toBeInstanceOf( + ReflexSelectionDeadlineError, + ); + await vi.advanceTimersByTimeAsync(20); + await rejected; + vi.useRealTimers(); + }); + + it('does not dispatch a native batch after its deadline has expired', async () => { + const { inputs } = setup(2); + const decideBatch = vi.fn(async () => []); + await expect( + chooseReflexWorldActions(inputs, provider(decideBatch), { + deadlineAtMs: Date.now() - 1, + }), + ).rejects.toBeInstanceOf(ReflexSelectionDeadlineError); + expect(decideBatch).not.toHaveBeenCalled(); + }); +}); diff --git a/apps/game-api/src/reflex-execution.ts b/apps/game-api/src/reflex-execution.ts index 1ad457e..1b4fa40 100644 --- a/apps/game-api/src/reflex-execution.ts +++ b/apps/game-api/src/reflex-execution.ts @@ -1,10 +1,15 @@ import { gridDistance } from 'h3-js'; import { ReflexProviderError, + type ReflexBatchDecisionResult, + type ReflexAttemptCompletion, type ReflexProvider, } from '@hexzero/agent-runtime'; import { h3CellSchema, + providerFailureSchema, + providerMetadataSchema, + reflexBatchResultEnvelopeSchema, reflexDecisionSchema, reflexObservationSchema, swarmDirectiveSchema, @@ -209,152 +214,731 @@ export interface ReflexSelectionOptions { now?: () => string; } +function failureForCode( + code: ProviderFailure['code'], + model: string, +): ProviderFailure { + const messages: Record = { + configuration: 'The reflex provider is not configured.', + timeout: 'The reflex provider request timed out.', + network: 'The reflex provider could not be reached.', + 'model-unavailable': 'The reflex model is unavailable.', + 'provider-http': 'The reflex provider returned an error.', + cancelled: 'The reflex provider request was cancelled.', + 'malformed-response': 'The reflex provider returned malformed data.', + 'unsupported-response': 'The reflex provider returned unsupported data.', + 'output-length': 'The reflex provider response exceeded its limit.', + 'missing-text-output': 'The reflex provider returned no decision.', + 'invalid-json': 'The reflex provider returned invalid JSON.', + 'missing-tool-call': 'The reflex provider returned no decision.', + 'multiple-tool-calls': 'The reflex provider returned ambiguous data.', + 'wrong-tool': 'The reflex provider returned an unsupported decision.', + 'invalid-tool-arguments': 'The reflex provider returned invalid data.', + 'invalid-decision': 'The reflex provider returned an invalid decision.', + 'simulation-validation': 'The reflex decision failed validation.', + 'budget-exhausted': 'The reflex provider attempt budget is exhausted.', + }; + return { code, message: messages[code], retryable: false, model }; +} + function safeFailure( error: unknown, model: string, -): { - failure: ProviderFailure; - metadata?: ProviderMetadata; -} { - if (error instanceof ReflexProviderError) - return { failure: error.failure, metadata: error.metadata }; +): { failure: ProviderFailure; metadata?: ProviderMetadata } { + if (error instanceof ReflexProviderError) { + const parsedFailure = providerFailureSchema.safeParse(error.failure); + const parsedMetadata = error.metadata + ? providerMetadataSchema.safeParse(error.metadata) + : undefined; + return { + failure: failureForCode( + parsedFailure.success ? parsedFailure.data.code : 'provider-http', + model, + ), + ...(parsedMetadata?.success ? { metadata: parsedMetadata.data } : {}), + }; + } + return { failure: failureForCode('provider-http', model) }; +} + +function normalizeAttemptCompletion( + input: unknown, + model: string, + includeDecision: boolean, +): ReflexAttemptCompletion { + const raw = + input && typeof input === 'object' + ? (input as Record) + : {}; + const outcome = + raw.outcome === 'completed' || + raw.outcome === 'provider-error' || + raw.outcome === 'cancelled' || + raw.outcome === 'timeout' + ? raw.outcome + : 'provider-error'; + const provider = raw.provider + ? providerMetadataSchema.safeParse(raw.provider) + : undefined; + const reportedFailure = raw.failure + ? providerFailureSchema.safeParse(raw.failure) + : undefined; + const failureCode = + outcome === 'cancelled' + ? 'cancelled' + : outcome === 'timeout' + ? 'timeout' + : outcome === 'completed' + ? undefined + : reportedFailure?.success + ? reportedFailure.data.code + : 'provider-http'; + const decision = + includeDecision && outcome === 'completed' + ? reflexDecisionSchema.safeParse(raw.reflexDecision) + : undefined; return { - failure: { - code: 'provider-http', - message: 'The reflex provider failed.', - retryable: false, - model, - }, + outcome, + ...(provider?.success ? { provider: provider.data } : {}), + ...(failureCode ? { failure: failureForCode(failureCode, model) } : {}), + ...(decision?.success ? { reflexDecision: decision.data } : {}), }; } -/** A failed Jev request selects wait; the independent attempt ledger still records it. */ -export async function chooseReflexWorldAction( - state: WorldState, - directive: SwarmDirective, - provider: ReflexProvider, - options: ReflexSelectionOptions = {}, -): Promise { - const compiled = compileReflexObservation(state, directive, options.history); - const fallback = (failure?: ProviderFailure): ReflexSelection => ({ +function fallbackSelection( + compiled: CompiledReflexObservation, + failure?: ProviderFailure, +): ReflexSelection { + return { action: { type: 'wait' }, observation: compiled.observation, decision: null, cognitionSource: 'deterministic-fallback', ...(failure ? { failure } : {}), - }); + }; +} + +function validateReflexDecision( + compiled: CompiledReflexObservation, + rawDecision: unknown, + model: string, +): ReflexSelection { + const parsed = reflexDecisionSchema.safeParse(rawDecision); + if (!parsed.success) + throw new ReflexProviderError(failureForCode('malformed-response', model)); + const decision = parsed.data; + if (decision.directiveId !== compiled.observation.directive.id) + throw new ReflexProviderError( + failureForCode('unsupported-response', model), + ); + if (decision.cognitionSource !== 'jev-reflex') + throw new ReflexProviderError( + failureForCode('unsupported-response', model), + ); + const action = compiled.actions.get(decision.chosenCandidateId); + if (!action) + throw new ReflexProviderError( + failureForCode('unsupported-response', model), + ); + const probabilityIds = Object.keys(decision.probabilities); + const probabilityTotal = Object.values(decision.probabilities).reduce( + (total, value) => total + value, + 0, + ); + if ( + probabilityIds.length !== compiled.actions.size || + probabilityIds.some((id) => !compiled.actions.has(id)) || + Math.abs(probabilityTotal - 1) > 0.001 || + Object.values(decision.probabilities).some( + (value) => value > decision.probabilities[decision.chosenCandidateId]!, + ) + ) + throw new ReflexProviderError( + failureForCode('unsupported-response', model), + ); + return { + action, + observation: compiled.observation, + decision, + cognitionSource: 'jev-reflex', + }; +} + +function deepFreeze(value: T): T { + if (value && typeof value === 'object' && !Object.isFrozen(value)) { + Object.freeze(value); + for (const child of Object.values(value as Record)) + deepFreeze(child); + } + return value; +} + +function freezeObservation(input: ReflexObservation): ReflexObservation { + return deepFreeze(reflexObservationSchema.parse(structuredClone(input))); +} + +export class ReflexSelectionDeadlineError extends Error { + constructor() { + super('The reflex selection exceeded the tick deadline.'); + this.name = 'ReflexSelectionDeadlineError'; + } +} + +export interface ReflexBatchSelectionInput { + compiled: CompiledReflexObservation; + intendedTickNumber?: number; + intendedTurnNumber: number; + /** The initial request consumes a permit reserved before the tick began. */ + initialPermitReserved?: boolean; +} + +export interface ReflexBatchSelectionOptions { + signal?: AbortSignal; + deadlineAtMs?: number; + accounting?: AttemptAccounting; + /** Bounds scalar-provider adapters; native batch providers make one call. */ + concurrencyLimit?: number; + now?: () => string; +} + +const MAX_REFLEX_BATCH_WORKERS = 31; +const MAX_REFLEX_BATCH_RESULTS = 64; + +function validateConcurrency(value: number): number { + if (!Number.isInteger(value) || value < 1 || value > 8) + throw new RangeError( + 'Reflex concurrencyLimit must be an integer from 1 to 8.', + ); + return value; +} + +/** Selects each worker independently while sharing one native batch dispatch. */ +export async function chooseReflexWorldActions( + inputs: readonly ReflexBatchSelectionInput[], + provider: ReflexProvider, + options: ReflexBatchSelectionOptions = {}, +): Promise { + if (inputs.length === 0) return []; + if (inputs.length > MAX_REFLEX_BATCH_WORKERS) + throw new RangeError( + `A reflex batch may contain at most ${MAX_REFLEX_BATCH_WORKERS} workers.`, + ); + const agentIds = inputs.map(({ compiled }) => compiled.observation.agentId); + if (new Set(agentIds).size !== agentIds.length) + throw new Error('A reflex batch cannot contain duplicate worker IDs.'); + const intendedTicks = new Set( + inputs.flatMap(({ intendedTickNumber }) => + intendedTickNumber === undefined ? [] : [intendedTickNumber], + ), + ); + if (intendedTicks.size > 1) + throw new Error('A reflex batch must belong to one intended tick.'); + if ( + provider.decideBatch && + options.accounting && + inputs.some(({ intendedTickNumber }) => intendedTickNumber === undefined) + ) + throw new Error( + 'Accounted native batches require an intended tick number.', + ); if (options.signal?.aborted) throw new ReflexSelectionCancelledError(); + if (options.deadlineAtMs !== undefined && Date.now() >= options.deadlineAtMs) + throw new ReflexSelectionDeadlineError(); + const concurrency = validateConcurrency(options.concurrencyLimit ?? 1); + if ( + concurrency > 1 && + !provider.decideBatch && + options.accounting && + (inputs.some(({ initialPermitReserved }) => !initialPermitReserved) || + options.accounting.snapshot().reservedPermits < inputs.length) + ) + throw new Error( + 'Concurrent scalar reflex selection requires a reserved initial permit for every worker.', + ); + const reservedRetryCount = + concurrency > 1 && !provider.decideBatch && options.accounting + ? options.accounting.reserveUpTo(inputs.length) + : 0; const model = provider.model ?? 'jev-1.13.0'; + const observations = inputs.map(({ compiled }) => + freezeObservation(compiled.observation), + ); + const frozenCompiled = inputs.map((input, index) => ({ + ...input.compiled, + observation: observations[index]!, + })); + const controller = new AbortController(); const now = options.now ?? (() => new Date().toISOString()); - let startedAttempts = 0; - try { - const rawDecision = await provider.decide(compiled.observation, { - signal: options.signal, - deadlineAtMs: options.deadlineAtMs, - beginAttempt: options.accounting - ? (kind) => { - const details = { - agentId: compiled.observation.agentId, - intendedTurnNumber: options.intendedTurnNumber ?? 1, - ...(options.intendedTickNumber - ? { intendedTickNumber: options.intendedTickNumber } - : {}), - kind, - startedAt: now(), - modelId: model, - reasoningProfile: 'provider-default', - } as const; - const permit = - kind === 'initial' && options.initialPermitReserved - ? options.accounting!.startReserved(details) - : options.accounting!.startAdditional(details); - if (permit === null) return null; - startedAttempts += 1; - return (completion) => - options.accounting!.finalize(permit, { - ...completion, - completedAt: now(), - }); - } - : undefined, - }); - if (options.accounting && startedAttempts === 0) - throw new ReflexProviderError({ - code: 'unsupported-response', - message: 'The reflex provider skipped attempt accounting.', - retryable: false, - model, - }); - const parsed = reflexDecisionSchema.safeParse(rawDecision); - if (!parsed.success) - throw new ReflexProviderError({ - code: 'malformed-response', - message: 'The reflex decision failed schema validation.', - retryable: false, - model, - }); - const decision = parsed.data; - if (options.signal?.aborted) throw new ReflexSelectionCancelledError(); - if (decision.directiveId !== compiled.observation.directive.id) - throw new ReflexProviderError({ - code: 'unsupported-response', - message: 'The reflex decision references another directive.', - retryable: false, - model, - }); - if (decision.cognitionSource !== 'jev-reflex') - throw new ReflexProviderError({ - code: 'unsupported-response', - message: 'The reflex decision has incorrect cognition attribution.', - retryable: false, - model, + let closed = false; + let abortReason: 'cancelled' | 'timeout' | undefined; + const openFinalizers = new Map< + (completion: SafeBatchCompletion) => void, + { finalize: (completion: SafeBatchCompletion) => void; agentId?: string } + >(); + type SafeBatchCompletion = { + outcome: 'completed' | 'provider-error' | 'cancelled' | 'timeout'; + provider?: ProviderMetadata; + failure?: ProviderFailure; + reflexDecision?: ReflexDecision; + completedAt: string; + }; + const closeOpenAttempts = ( + outcome: 'cancelled' | 'timeout' | 'provider-error', + ) => { + closed = true; + for (const { finalize } of [...openFinalizers.values()]) + finalize({ + outcome, + failure: + outcome === 'provider-error' + ? failureForCode('unsupported-response', model) + : failureForCode(outcome, model), + completedAt: now(), }); - const action = compiled.actions.get(decision.chosenCandidateId); - if (!action) - throw new ReflexProviderError({ - code: 'unsupported-response', - message: 'The reflex decision selected an unknown candidate.', - retryable: false, - model, + openFinalizers.clear(); + }; + let expireDeadline: () => void = () => undefined; + const finalizeOpenForAgent = (agentId: string, outcome: 'provider-error') => { + for (const [key, entry] of [...openFinalizers]) { + if (entry.agentId !== agentId) continue; + entry.finalize({ + outcome, + failure: failureForCode('unsupported-response', model), + completedAt: now(), }); - const probabilityIds = Object.keys(decision.probabilities); - const probabilityTotal = Object.values(decision.probabilities).reduce( - (total, value) => total + value, - 0, - ); + openFinalizers.delete(key); + } + }; + const startAttempt = ( + kind: 'initial' | 'automatic-transport-retry', + members: readonly ReflexBatchSelectionInput[], + batchId: string, + reserved: boolean, + ) => { + if (!options.accounting) return null; if ( - probabilityIds.length !== compiled.actions.size || - probabilityIds.some((id) => !compiled.actions.has(id)) || - Math.abs(probabilityTotal - 1) > 0.001 || - Object.values(decision.probabilities).some( - (value) => value > decision.probabilities[decision.chosenCandidateId]!, + closed || + controller.signal.aborted || + (options.deadlineAtMs !== undefined && Date.now() >= options.deadlineAtMs) + ) { + if ( + !closed && + options.deadlineAtMs !== undefined && + Date.now() >= options.deadlineAtMs ) - ) - throw new ReflexProviderError({ - code: 'unsupported-response', - message: 'The reflex decision has inconsistent probability telemetry.', - retryable: false, - model, + expireDeadline(); + return null; + } + const first = members[0]!; + const details = { + agentId: first.compiled.observation.agentId, + intendedTurnNumber: first.intendedTurnNumber, + ...(first.intendedTickNumber === undefined + ? {} + : { intendedTickNumber: first.intendedTickNumber }), + kind, + startedAt: now(), + modelId: model, + reasoningProfile: 'provider-default', + batch: { + id: batchId, + members: members.map((member) => ({ + agentId: member.compiled.observation.agentId, + intendedTurnNumber: member.intendedTurnNumber, + })), + }, + } as const; + const permit = + kind === 'initial' && reserved + ? options.accounting.startReserved(details) + : options.accounting.startAdditional(details); + if (permit === null) return null; + let finalized = false; + const finalize = (completion: SafeBatchCompletion) => { + if (finalized) return; + finalized = true; + openFinalizers.delete(finalize); + options.accounting!.finalize(permit, completion); + }; + openFinalizers.set(finalize, { finalize }); + return (completion: unknown) => { + if (closed || finalized) return; + finalize({ + ...normalizeAttemptCompletion(completion, model, false), + completedAt: now(), }); - return { - action, - observation: compiled.observation, - decision, - cognitionSource: 'jev-reflex', }; - } catch (error) { - const cancelled = - error instanceof ReflexSelectionCancelledError || options.signal?.aborted; - const safe = safeFailure(error, model); - const failure: ProviderFailure = cancelled - ? { - code: 'cancelled', - message: 'The reflex request was cancelled.', - retryable: false, - model, + }; + const scalarBeginAttempt = ( + input: ReflexBatchSelectionInput, + inputIndex: number, + onStarted: () => void, + isSettled: () => boolean, + ) => { + const startedKinds = new Set(); + const finalizedKinds = new Set(); + return (kind: 'initial' | 'automatic-transport-retry') => { + if (isSettled() || startedKinds.has(kind)) return null; + if ( + kind === 'automatic-transport-retry' && + !finalizedKinds.has('initial') + ) + return null; + if (!options.accounting) return null; + if ( + closed || + controller.signal.aborted || + (options.deadlineAtMs !== undefined && + Date.now() >= options.deadlineAtMs) + ) { + if ( + !closed && + options.deadlineAtMs !== undefined && + Date.now() >= options.deadlineAtMs + ) + expireDeadline(); + return null; + } + const details = { + agentId: input.compiled.observation.agentId, + intendedTurnNumber: input.intendedTurnNumber, + ...(input.intendedTickNumber === undefined + ? {} + : { intendedTickNumber: input.intendedTickNumber }), + kind, + startedAt: now(), + modelId: model, + reasoningProfile: 'provider-default', + } as const; + const permit = + kind === 'initial' && input.initialPermitReserved + ? options.accounting.startReserved(details) + : kind === 'automatic-transport-retry' && concurrency > 1 + ? inputIndex < reservedRetryCount + ? options.accounting.startReserved(details) + : null + : options.accounting.startAdditional(details); + if (permit === null) return null; + startedKinds.add(kind); + onStarted(); + let finalized = false; + const finalize = (completion: SafeBatchCompletion) => { + if (finalized) return; + finalized = true; + finalizedKinds.add(kind); + openFinalizers.delete(finalize); + options.accounting!.finalize(permit, completion); + }; + openFinalizers.set(finalize, { + finalize, + agentId: input.compiled.observation.agentId, + }); + return (completion: unknown) => { + if (closed || finalized || isSettled()) return; + finalize({ + ...normalizeAttemptCompletion(completion, model, true), + completedAt: now(), + }); + }; + }; + }; + + let deadlineTimer: ReturnType | undefined; + let rejectLifecycle!: (error: Error) => void; + const lifecycle = new Promise((_, reject) => { + rejectLifecycle = reject; + }); + // The deadline may already be expired before the provider operation starts. + // Keep early rejection handled until the racing await is attached below. + void lifecycle.catch(() => undefined); + const onExternalAbort = () => { + if (closed) return; + abortReason = 'cancelled'; + controller.abort(); + closeOpenAttempts('cancelled'); + rejectLifecycle(new ReflexSelectionCancelledError()); + }; + expireDeadline = () => { + if (closed) return; + abortReason = 'timeout'; + controller.abort(); + closeOpenAttempts('timeout'); + rejectLifecycle(new ReflexSelectionDeadlineError()); + }; + options.signal?.addEventListener('abort', onExternalAbort, { once: true }); + if (options.signal?.aborted) onExternalAbort(); + if (options.deadlineAtMs !== undefined) { + const remaining = options.deadlineAtMs - Date.now(); + if (remaining <= 0) { + expireDeadline(); + } else { + deadlineTimer = setTimeout(() => { + expireDeadline(); + }, remaining); + } + } + + const invoke = async (): Promise => { + if (closed || controller.signal.aborted) { + if (abortReason === 'timeout') throw new ReflexSelectionDeadlineError(); + throw new ReflexSelectionCancelledError(); + } + if ( + options.deadlineAtMs !== undefined && + Date.now() >= options.deadlineAtMs + ) { + expireDeadline(); + throw new ReflexSelectionDeadlineError(); + } + if (provider.decideBatch) { + const batchId = crypto.randomUUID(); + let attemptsStarted = 0; + let settled = false; + const startedKinds = new Set(); + const finalizedKinds = new Set(); + const result = await provider + .decideBatch(observations, { + signal: controller.signal, + deadlineAtMs: options.deadlineAtMs, + beginAttempt: options.accounting + ? (kind) => { + if (settled || startedKinds.has(kind)) return null; + if ( + kind === 'automatic-transport-retry' && + !finalizedKinds.has('initial') + ) + return null; + const finalize = startAttempt( + kind, + inputs, + batchId, + inputs[0]?.initialPermitReserved ?? false, + ); + if (finalize === null) return null; + attemptsStarted += 1; + startedKinds.add(kind); + return (completion) => { + if (settled) return; + finalize(completion); + finalizedKinds.add(kind); + }; + } + : undefined, + }) + .finally(() => { + settled = true; + }); + if (options.accounting && attemptsStarted === 0) + return inputs.map(({ compiled }) => ({ + agentId: compiled.observation.agentId, + status: 'failed' as const, + failure: failureForCode('unsupported-response', model), + })); + return result; + } + + const results: ReflexBatchDecisionResult[] = new Array(inputs.length); + const countableInputs = inputs.map((input) => ({ ...input })); + let cursor = 0; + const invokeOne = async (index: number) => { + const input = countableInputs[index]!; + let attemptsStarted = 0; + let settled = false; + try { + const decision = await provider + .decide(observations[index]!, { + signal: controller.signal, + deadlineAtMs: options.deadlineAtMs, + beginAttempt: options.accounting + ? scalarBeginAttempt( + input, + index, + () => attemptsStarted++, + () => settled, + ) + : undefined, + }) + .finally(() => { + settled = true; + }); + results[index] = + attemptsStarted === 0 && options.accounting + ? { + agentId: input.compiled.observation.agentId, + status: 'failed' as const, + failure: failureForCode('unsupported-response', model), + } + : { + agentId: input.compiled.observation.agentId, + status: 'completed' as const, + decision, + }; + finalizeOpenForAgent( + input.compiled.observation.agentId, + 'provider-error', + ); + } catch (error) { + const safe = safeFailure(error, model); + finalizeOpenForAgent( + input.compiled.observation.agentId, + 'provider-error', + ); + results[index] = { + agentId: input.compiled.observation.agentId, + status: 'failed' as const, + failure: safe.failure, + }; + } + }; + cursor = 0; + const boundedWorkers = Array.from( + { length: Math.min(concurrency, inputs.length) }, + async () => { + while (!closed) { + if ( + options.deadlineAtMs !== undefined && + Date.now() >= options.deadlineAtMs + ) { + expireDeadline(); + return; + } + const index = cursor++; + if (index >= inputs.length) return; + await invokeOne(index); } - : safe.failure; - if (cancelled) throw new ReflexSelectionCancelledError(); - return fallback(failure); + }, + ); + await Promise.all(boundedWorkers); + return results.filter( + (result): result is ReflexBatchDecisionResult => result !== undefined, + ); + }; + + try { + const providerOperation = invoke(); + const rawResults = await Promise.race([providerOperation, lifecycle]); + if ( + abortReason === 'timeout' || + (options.deadlineAtMs !== undefined && Date.now() >= options.deadlineAtMs) + ) { + closeOpenAttempts('timeout'); + throw new ReflexSelectionDeadlineError(); + } + if (options.signal?.aborted) throw new ReflexSelectionCancelledError(); + closeOpenAttempts('provider-error'); + return mapBatchResults(frozenCompiled, rawResults, model); + } catch (error) { + if (error instanceof ReflexSelectionCancelledError) throw error; + if (error instanceof ReflexSelectionDeadlineError) throw error; + if (options.signal?.aborted || abortReason === 'cancelled') + throw new ReflexSelectionCancelledError(); + if ( + abortReason === 'timeout' || + (options.deadlineAtMs !== undefined && Date.now() >= options.deadlineAtMs) + ) { + expireDeadline(); + throw new ReflexSelectionDeadlineError(); + } + closeOpenAttempts('provider-error'); + return frozenCompiled.map((compiled) => + fallbackSelection(compiled, safeFailure(error, model).failure), + ); + } finally { + closed = true; + controller.abort(); + if (deadlineTimer) clearTimeout(deadlineTimer); + options.signal?.removeEventListener('abort', onExternalAbort); + } +} + +function mapBatchResults( + compiled: readonly CompiledReflexObservation[], + rawResults: unknown, + model: string, +): ReflexSelection[] { + const byAgent = new Map(); + if ( + Array.isArray(rawResults) && + rawResults.length <= MAX_REFLEX_BATCH_RESULTS + ) { + for (const raw of rawResults) { + if (!raw || typeof raw !== 'object') continue; + const agentId = (raw as { agentId?: unknown }).agentId; + if (typeof agentId !== 'string') continue; + const matches = byAgent.get(agentId) ?? []; + matches.push(raw); + byAgent.set(agentId, matches); + } } + return compiled.map((item) => { + const agentId = item.observation.agentId; + const matches = byAgent.get(agentId) ?? []; + if ( + !Array.isArray(rawResults) || + rawResults.length > MAX_REFLEX_BATCH_RESULTS || + matches.length !== 1 + ) + return fallbackSelection( + item, + failureForCode('unsupported-response', model), + ); + const parsedResult = reflexBatchResultEnvelopeSchema.safeParse(matches[0]); + if (!parsedResult.success) + return fallbackSelection( + item, + failureForCode('unsupported-response', model), + ); + const result = parsedResult.data; + if (result.status === 'failed') { + const parsed = providerFailureSchema.safeParse(result.failure); + return fallbackSelection( + item, + failureForCode( + parsed.success ? parsed.data.code : 'provider-http', + model, + ), + ); + } + try { + return validateReflexDecision(item, result.decision, model); + } catch (error) { + return fallbackSelection(item, safeFailure(error, model).failure); + } + }); +} + +/** A failed Jev request selects wait; the independent attempt ledger records it. */ +export async function chooseReflexWorldAction( + state: WorldState, + directive: SwarmDirective, + provider: ReflexProvider, + options: ReflexSelectionOptions = {}, +): Promise { + const compiled = compileReflexObservation(state, directive, options.history); + const scalarProvider: ReflexProvider = { + mode: provider.mode, + model: provider.model, + configured: provider.configured, + decide: provider.decide.bind(provider), + }; + const [selection] = await chooseReflexWorldActions( + [ + { + compiled, + intendedTurnNumber: options.intendedTurnNumber ?? 1, + ...(options.intendedTickNumber === undefined + ? {} + : { intendedTickNumber: options.intendedTickNumber }), + initialPermitReserved: options.initialPermitReserved, + }, + ], + scalarProvider, + { + signal: options.signal, + deadlineAtMs: options.deadlineAtMs, + accounting: options.accounting, + now: options.now, + }, + ); + return selection!; } diff --git a/apps/game-api/src/simulation-service.batch.test.ts b/apps/game-api/src/simulation-service.batch.test.ts new file mode 100644 index 0000000..e2d7033 --- /dev/null +++ b/apps/game-api/src/simulation-service.batch.test.ts @@ -0,0 +1,375 @@ +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { + DeterministicSwarmPlanner, + DeterministicReflexProvider, + type ReflexProvider, + type ReflexDecisionOptions, + type SwarmPlanner, +} from '@hexzero/agent-runtime'; +import { type ReflexDecision, type ReflexObservation } from '@hexzero/shared'; +import { generateDeterministicRoster } from '@hexzero/world-engine'; +import { + SimulationService, + SimulationTurnCancelledError, +} from './simulation-service'; + +function decisionFor(observation: ReflexObservation): ReflexDecision { + const chosen = + observation.candidates.find(({ description }) => + description.startsWith('Infect'), + ) ?? observation.candidates[0]!; + return { + chosenCandidateId: chosen.id, + confidence: 1, + probabilities: Object.fromEntries( + observation.candidates.map(({ id }) => [id, Number(id === chosen.id)]), + ), + model: 'controlled-reflex', + latencyMs: 0, + directiveId: observation.directive.id, + cognitionSource: 'jev-reflex', + }; +} + +class ControlledReflex implements ReflexProvider { + readonly mode = 'scripted-reflex-test' as const; + readonly model = 'controlled-reflex'; + readonly configured = true; + readonly calls: { observation: ReflexObservation; finish: () => void }[] = []; + async decide( + observation: ReflexObservation, + options: ReflexDecisionOptions = {}, + ) { + const finalize = options.beginAttempt?.('initial'); + if (finalize === null) throw new Error('Attempt denied'); + return await new Promise((resolve) => { + this.calls.push({ + observation, + finish: () => { + const decision = decisionFor(observation); + finalize?.({ outcome: 'completed', reflexDecision: decision }); + resolve(decision); + }, + }); + }); + } +} + +function service( + reflexProvider: ReflexProvider, + concurrency = 3, + planner: SwarmPlanner = new DeterministicSwarmPlanner(), +) { + let event = 0; + const simulation = new SimulationService({ + swarmPlanner: planner, + reflexProvider, + reflexConcurrencyLimit: concurrency, + now: () => '2026-10-02T00:00:00.000Z', + createEventId: () => + `00000000-0000-4000-8000-${String(++event).padStart(12, '0')}`, + createExperimentId: () => '00000000-0000-4000-8000-000000000001', + }); + simulation.applyWorldSetup(simulation.getDefaultWorldSetup()); + return simulation; +} + +// Flush bounded microtask scheduling, without speed or elapsed-time assertions. +async function dispatched(provider: ControlledReflex, count: number) { + for (let pass = 0; pass < 100 && provider.calls.length < count; pass += 1) + await Promise.resolve(); + expect(provider.calls).toHaveLength(count); +} + +afterEach(() => vi.useRealTimers()); + +describe('batch-capable simultaneous worker tick', () => { + it('reuses two eligible directives and makes 17 deterministic decisions after a failed 20-agent replan', async () => { + const planner: SwarmPlanner = { + mode: 'scripted-swarm-test', + configured: true, + async plan(observation, model, options = {}) { + const finalize = options.beginAttempt?.('initial'); + if (observation.tickNumber > 1) { + finalize?.({ + outcome: 'provider-error', + failure: { + code: 'invalid-decision', + message: 'Invalid Zero plan.', + retryable: false, + }, + }); + throw new Error('Invalid Zero plan'); + } + const plan = { + strategySummary: 'Infect local cells; retain two holding directives.', + zeroActionCandidateId: observation.legalZeroActions.find( + ({ action }) => action.type === 'wait', + )!.id, + directives: observation.workerOptions.map((worker, index) => { + const option = worker.options.find((candidate) => + index < 2 + ? candidate.mission === 'hold' + : candidate.mission === 'expand', + )!; + return { + id: `fallback-${index}`, + agentId: worker.agentId, + mission: option.mission, + targetCell: option.targetCell, + priority: 'normal' as const, + riskTolerance: 'low' as const, + issuedAtTick: 1, + expiresAtTick: index < 2 ? 5 : 1, + }; + }), + }; + const metadata = { + provider: 'scripted-test' as const, + model, + latencyMs: 0, + }; + finalize?.({ + outcome: 'completed', + provider: metadata, + swarmPlan: plan, + }); + return { plan, metadata }; + }, + }; + const deterministic = new DeterministicReflexProvider(); + const provider: ReflexProvider = { + mode: deterministic.mode, + model: deterministic.model, + configured: true, + decide: vi.fn((observation, options) => + deterministic.decide(observation, options), + ), + }; + const simulation = service(provider, 3, planner); + const request = simulation.getDefaultWorldSetup(); + const roster = generateDeterministicRoster(20, 'fallback-roster'); + simulation.applyWorldSetup({ + ...request, + radius: 12, + roster, + patientZeroAgentId: roster[0]!.id, + spawnSeed: 'fallback-spawn', + }); + const first = await simulation.executeNextTick(); + expect(first?.planSource).toBe('zero-llm'); + expect(provider.decide).toHaveBeenCalledTimes(19); + const tick = await simulation.executeNextTick(); + expect(tick?.planSource).toBe('deterministic-fallback'); + expect( + tick?.workers.filter(({ source }) => source === 'deterministic-fallback'), + ).toHaveLength(17); + expect( + tick?.workers.filter(({ source }) => source === 'jev-reflex'), + ).toHaveLength(2); + expect(provider.decide).toHaveBeenCalledTimes(21); + expect(simulation.getSnapshot().experiment.attemptAccounting).toMatchObject( + { attemptsStarted: 23, attemptsFinalized: 23, reservedPermits: 0 }, + ); + }); + it('resolves identical frozen observations in seeded order despite reversed completion', async () => { + const run = async (reverse: boolean) => { + const provider = new ControlledReflex(); + const simulation = service(provider, 8); + const before = simulation.getSnapshot(); + const execution = simulation.executeNextTick(); + await dispatched(provider, 7); + expect(simulation.getSnapshot().world).toEqual(before.world); + const calls = reverse ? [...provider.calls].reverse() : provider.calls; + for (const call of calls) call.finish(); + const tick = await execution; + return { + world: simulation.getSnapshot().world, + order: simulation.getSnapshot().resolutionOrder, + tick, + observations: provider.calls.map(({ observation }) => observation), + }; + }; + expect(await run(true)).toEqual(await run(false)); + }); + + it('cancels an ignoring provider, releases queued reservations and rejects late accounting', async () => { + const provider = new ControlledReflex(); + const simulation = service(provider, 2); + const before = simulation.getSnapshot(); + const execution = simulation.executeNextTick(); + await dispatched(provider, 2); + const rejection = expect(execution).rejects.toBeInstanceOf( + SimulationTurnCancelledError, + ); + simulation.cancelCurrentRequest(); + await rejection; + const snapshot = simulation.getSnapshot(); + expect(snapshot.world).toEqual(before.world); + expect(snapshot.tickNumber).toBe(0); + expect(snapshot.swarmTicks).toEqual([]); + expect(snapshot.experiment.attemptAccounting).toMatchObject({ + attemptsStarted: 3, + attemptsFinalized: 3, + attemptsInFlight: 0, + reservedPermits: 0, + }); + const accounting = snapshot.experiment.attemptAccounting; + provider.calls.forEach(({ finish }) => finish()); + await Promise.resolve(); + expect(simulation.getSnapshot().experiment.attemptAccounting).toEqual( + accounting, + ); + expect(provider.calls).toHaveLength(2); + expect(simulation.getSnapshot().world).toEqual(before.world); + }); + + it('rolls back an expired shared deadline, finalizes only dispatched work and ignores late results', async () => { + vi.useFakeTimers(); + const provider = new ControlledReflex(); + const simulation = service(provider, 2); + const before = simulation.getSnapshot(); + const execution = simulation.executeNextTick(); + await dispatched(provider, 2); + const rejection = expect(execution).rejects.toThrow(/deadline/i); + await vi.advanceTimersByTimeAsync(75_001); + await rejection; + expect(simulation.getSnapshot().world).toEqual(before.world); + expect(simulation.getSnapshot().tickNumber).toBe(0); + const accounting = simulation.getSnapshot().experiment.attemptAccounting; + expect(accounting).toMatchObject({ + attemptsStarted: 3, + attemptsFinalized: 3, + attemptsInFlight: 0, + reservedPermits: 0, + }); + provider.calls.forEach(({ finish }) => finish()); + await Promise.resolve(); + expect(simulation.getSnapshot().experiment.attemptAccounting).toEqual( + accounting, + ); + expect(provider.calls).toHaveLength(2); + }); + + it('attributes a native shared dispatch once and retains unrelated valid worker results', async () => { + const provider: ReflexProvider = { + mode: 'scripted-reflex-test', + configured: true, + model: 'controlled-reflex', + decide: vi.fn(async () => { + throw new Error('Scalar path must not run'); + }), + async decideBatch(observations, options) { + options?.beginAttempt?.('initial')?.({ + outcome: 'completed', + provider: { + provider: 'scripted-test', + model: 'controlled-reflex', + latencyMs: 5, + promptTokens: 100, + completionTokens: 50, + totalTokens: 150, + costCredits: 0.02, + }, + }); + return observations + .map((observation, index) => ({ + status: 'completed' as const, + agentId: observation.agentId, + decision: + index === 1 ? { invalid: true } : decisionFor(observation), + })) + .reverse(); + }, + }; + const simulation = service(provider); + const tick = await simulation.executeNextTick(); + expect(provider.decide).not.toHaveBeenCalled(); + expect( + tick?.workers.filter(({ source }) => source === 'jev-reflex'), + ).toHaveLength(6); + expect( + tick?.workers.filter(({ source }) => source === 'deterministic-fallback'), + ).toHaveLength(1); + expect(simulation.getSnapshot().experiment.attemptAccounting).toMatchObject( + { + attemptsStarted: 2, + attemptsFinalized: 2, + reservedPermits: 0, + knownFinalizedCostCredits: '0.02', + }, + ); + const exported = simulation.generateExperimentExport({ + agents: { mode: 'all' }, + turns: { mode: 'entire-retained' }, + outcomes: ['accepted', 'rejected', 'provider-error', 'operator-skipped'], + actions: ['move', 'infect', 'capture', 'wait'], + level: 'full-safe', + serialization: 'compact', + }); + const batch = exported.providerAttempts?.find(({ batch }) => batch); + expect(batch?.batch?.members.map(({ agentId }) => agentId)).toEqual( + tick?.workers.map(({ agentId }) => agentId), + ); + expect(batch?.actualCostCredits).toBe('0.02'); + expect(batch?.reflexDecision).toBeUndefined(); + expect( + tick?.workers.every( + ({ reflexDecision }) => reflexDecision?.inputTokens === undefined, + ), + ).toBe(true); + const memberId = tick!.workers.at(-1)!.agentId; + const filtered = simulation.generateExperimentExport({ + agents: { mode: 'selected', agentIds: [memberId] }, + turns: { mode: 'entire-retained' }, + outcomes: ['accepted', 'rejected', 'provider-error', 'operator-skipped'], + actions: ['move', 'infect', 'capture', 'wait'], + level: 'full-safe', + serialization: 'compact', + }); + expect(filtered.selection.selectedAgentIds).toEqual([memberId]); + expect(filtered.providerAttempts).toHaveLength(1); + expect(filtered.providerAttempts?.[0]?.batch?.members).toEqual( + batch?.batch?.members, + ); + expect(filtered.agents.map(({ id }) => id)).toEqual( + expect.arrayContaining( + batch!.batch!.members.map(({ agentId }) => agentId), + ), + ); + expect(filtered.metrics?.aggregate.knownCostCredits).toBe(0.02); + expect(filtered.metrics?.byAgent[0]?.metrics.knownCostCredits).toBe(0); + }); + + it('preserves deterministic expansion without any reflex dispatch after first planner failure', async () => { + const planner: SwarmPlanner = { + mode: 'scripted-swarm-test', + configured: true, + async plan(_observation, _model, options) { + options?.beginAttempt?.('initial')?.({ + outcome: 'provider-error', + failure: { + code: 'provider-http', + message: 'Planner unavailable.', + retryable: false, + }, + }); + throw new Error('Planner unavailable'); + }, + }; + const provider = new ControlledReflex(); + const simulation = service(provider, 3, planner); + const tick = await simulation.executeNextTick(); + expect(provider.calls).toHaveLength(0); + expect(tick?.planSource).toBe('deterministic-fallback'); + expect( + tick?.workers.every(({ source }) => source === 'deterministic-fallback'), + ).toBe(true); + expect(tick?.workers.some(({ action }) => action?.type === 'infect')).toBe( + true, + ); + expect(simulation.getSnapshot().experiment.attemptAccounting).toMatchObject( + { attemptsStarted: 1, attemptsFinalized: 1, reservedPermits: 0 }, + ); + }); +}); diff --git a/apps/game-api/src/simulation-service.ts b/apps/game-api/src/simulation-service.ts index 6fa0859..6cfc913 100644 --- a/apps/game-api/src/simulation-service.ts +++ b/apps/game-api/src/simulation-service.ts @@ -76,11 +76,13 @@ import { import { compileStrategicOptions } from './strategic-options'; import { AttemptAccounting } from './attempt-accounting'; import { - chooseReflexWorldAction, + chooseReflexWorldActions, compileReflexObservation, ReflexSelectionCancelledError, + ReflexSelectionDeadlineError, type CompiledReflexObservation, type ReflexSelection, + type ReflexBatchSelectionInput, } from './reflex-execution'; import { isSwarmDirectiveComplete, @@ -274,6 +276,8 @@ export class SimulationValidationError extends Error { export interface SimulationServiceOptions { swarmPlanner: SwarmPlanner; reflexProvider: ReflexProvider; + /** Server-owned individual-call pool cap, 1–8. Production Jev remains 1. */ + reflexConcurrencyLimit?: number; /** * Offline comparison seam: choose one opaque, already legal candidate without * calling a reflex provider. Production zero-swarm execution leaves this unset. @@ -290,6 +294,7 @@ export interface SimulationServiceOptions { export class SimulationService { readonly #swarmPlanner: SwarmPlanner; readonly #reflexProvider: ReflexProvider; + readonly #reflexConcurrencyLimit: number; readonly #deterministicWorkerCandidateSelector: ((compiled: CompiledReflexObservation) => string) | undefined; readonly #now: () => string; @@ -329,6 +334,7 @@ export class SimulationService { constructor({ swarmPlanner, reflexProvider, + reflexConcurrencyLimit = 1, deterministicWorkerCandidateSelector, now = () => new Date().toISOString(), createEventId = () => crypto.randomUUID(), @@ -342,6 +348,15 @@ export class SimulationService { throw new Error('Experiment retention limit must be a positive integer.'); this.#swarmPlanner = swarmPlanner; this.#reflexProvider = reflexProvider; + if ( + !Number.isInteger(reflexConcurrencyLimit) || + reflexConcurrencyLimit < 1 || + reflexConcurrencyLimit > 8 + ) + throw new Error( + 'Reflex concurrency limit must be an integer from 1 to 8.', + ); + this.#reflexConcurrencyLimit = reflexConcurrencyLimit; this.#deterministicWorkerCandidateSelector = deterministicWorkerCandidateSelector; this.#now = now; @@ -826,11 +841,12 @@ export class SimulationService { // Planning ticks reserve Zero plus Jev workers. The explicit deterministic // comparison seam only reserves Zero's planner call; it never dispatches a // reflex provider attempt. - const requiredAttempts = this.#deterministicWorkerCandidateSelector - ? replan - ? 1 - : 0 - : agents.length - (replan ? 0 : 1); + const workerAttemptCount = this.#deterministicWorkerCandidateSelector + ? 0 + : this.#reflexProvider.decideBatch + ? Number(agents.length > 1) + : agents.length - 1; + const requiredAttempts = Number(replan) + workerAttemptCount; if ( requiredAttempts > 0 && !this.#attemptAccounting.reserve(requiredAttempts) @@ -946,16 +962,14 @@ export class SimulationService { const zeroAction = observation.legalZeroActions.find( ({ id }) => id === plan.zeroActionCandidateId, )?.action ?? { type: 'wait' as const }; - const selected = new Map< - AgentId, - Awaited> - >(); + if (Date.now() >= deadlineAtMs) throw new ReflexSelectionDeadlineError(); + const selected = new Map(); + const reflexInputs: ReflexBatchSelectionInput[] = []; const workers = agents.filter(({ id }) => id !== zero.id); for (const worker of workers) { const directive = plan.directives.find( ({ agentId }) => agentId === worker.id, )!; - this.#activeAgentId = worker.id; const history = { previousCell: this.#lastSwarmPositions.get(worker.id), pressureEvents, @@ -976,37 +990,49 @@ export class SimulationService { ({ directiveId }) => directiveId === previous.id, ), ); + const compiled = compileReflexObservation( + candidate, + directive, + history, + ); const choice = planSource === 'deterministic-fallback' && !retainedDirective - ? chooseDeterministicWorkerAction( - compileReflexObservation(candidate, directive, history), - (compiled) => - selectNeutralFallbackCandidate(compiled, candidate), + ? chooseDeterministicWorkerAction(compiled, (compiled) => + selectNeutralFallbackCandidate(compiled, candidate), ) : this.#deterministicWorkerCandidateSelector ? chooseDeterministicWorkerAction( - compileReflexObservation(candidate, directive, history), + compiled, this.#deterministicWorkerCandidateSelector, ) - : await chooseReflexWorldAction( - candidate, - directive, - this.#reflexProvider, - { - history, - signal: controller.signal, - deadlineAtMs, - accounting: this.#attemptAccounting, - initialPermitReserved: true, - intendedTickNumber: tickNumber, - intendedTurnNumber: - tickTurnBase + order.indexOf(worker.id) + 1, - now: this.#now, - }, - ); - selected.set(worker.id, choice); + : null; + if (choice) selected.set(worker.id, choice); + else + reflexInputs.push({ + compiled, + initialPermitReserved: true, + intendedTickNumber: tickNumber, + intendedTurnNumber: tickTurnBase + order.indexOf(worker.id) + 1, + }); } + // Completion order never determines observation or resolution order. + this.#activeAgentId = null; + const choices = await chooseReflexWorldActions( + reflexInputs, + this.#reflexProvider, + { + signal: controller.signal, + deadlineAtMs, + accounting: this.#attemptAccounting, + concurrencyLimit: this.#reflexConcurrencyLimit, + now: this.#now, + }, + ); + reflexInputs.forEach(({ compiled }, index) => + selected.set(compiled.observation.agentId, choices[index]!), + ); if (controller.signal.aborted) throw new SimulationTurnCancelledError(); + if (Date.now() >= deadlineAtMs) throw new ReflexSelectionDeadlineError(); let state = candidate; const context = { now: () => virtualTime, @@ -1026,6 +1052,7 @@ export class SimulationService { applied.set(agentId, result.result); } if (controller.signal.aborted) throw new SimulationTurnCancelledError(); + if (Date.now() >= deadlineAtMs) throw new ReflexSelectionDeadlineError(); const signals = workers.flatMap((worker) => { const selection = selected.get(worker.id)!; const probability = selection.decision?.replanProbability; @@ -1725,7 +1752,7 @@ export class SimulationService { currentWorld: this.#worldSnapshot(), modelConfiguration: this.#modelConfiguration, scenario: this.#scenario, - schemaVersion: 13, + schemaVersion: 14, providerAttempts: this.#attemptAccounting.ledger(), attemptRetention: this.#attemptAccounting.retention(), attemptAccounting: this.#attemptAccounting.snapshot(), diff --git a/apps/game-api/src/swarm-comparison.ts b/apps/game-api/src/swarm-comparison.ts index fb67df5..85d3bcc 100644 --- a/apps/game-api/src/swarm-comparison.ts +++ b/apps/game-api/src/swarm-comparison.ts @@ -322,7 +322,10 @@ function sample( latencyMs, promptTokens: inputTokens, completionTokens: outputTokens, - totalTokens: inputTokens + outputTokens, + totalTokens: + inputTokens === undefined || outputTokens === undefined + ? undefined + : inputTokens + outputTokens, })), ]; return { diff --git a/apps/world-lab/src/components/swarm-view.test.tsx b/apps/world-lab/src/components/swarm-view.test.tsx index cd908a1..2e3ea50 100644 --- a/apps/world-lab/src/components/swarm-view.test.tsx +++ b/apps/world-lab/src/components/swarm-view.test.tsx @@ -102,6 +102,16 @@ function snapshot(withTick = true): SimulationSnapshot { } describe('swarm telemetry panels', () => { + it('identifies unknown worker token usage without inventing a batch allocation', () => { + const value = snapshot(); + const decision = value.swarmTicks![0]!.workers[0]!.reflexDecision!; + delete decision.inputTokens; + delete decision.outputTokens; + render(); + expect( + screen.getByText(/usage unknown for 1 worker decisions/), + ).toBeInTheDocument(); + }); it('summarizes inactive player pressure, actions, and provider failures', () => { const value = snapshot(); value.swarmTicks![0]!.workers[0]!.failure = { diff --git a/apps/world-lab/src/components/swarm-view.tsx b/apps/world-lab/src/components/swarm-view.tsx index 265baeb..4000562 100644 --- a/apps/world-lab/src/components/swarm-view.tsx +++ b/apps/world-lab/src/components/swarm-view.tsx @@ -484,6 +484,10 @@ export function SwarmRunPanel({ (item.outputTokens ?? item.completionTokens ?? 0), 0, ); + const workersWithUnknownUsage = jev.filter( + (decision) => + decision.inputTokens === undefined || decision.outputTokens === undefined, + ).length; const providers = snapshot.swarmProviderStatus; const retainedTicks = snapshot.swarmTicks ?? []; const zeroPlans = retainedTicks.filter( @@ -577,6 +581,8 @@ export function SwarmRunPanel({
{total(jev)} tokens ·{' '} {jev.reduce((sum, item) => sum + item.latencyMs, 0)} ms + {workersWithUnknownUsage > 0 && + ` · usage unknown for ${workersWithUnknownUsage} worker decisions`}
diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 7c5d0d8..ef01761 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -128,8 +128,10 @@ tick commits. The Game API owns an operator-triggered tick transaction. It builds Zero's strategic observation from the frozen candidate world and dispatches the single -planning call; worker reflex calls follow sequentially, each receiving an -immutable compiled local observation. All calls share one absolute deadline. +planning call; all worker observations are compiled before dispatch. Individual +providers use an explicit pool cap of 1–8 (production Jev remains 1). Native +batch providers receive one bounded group, at most 31 workers. All calls share +one absolute deadline. See ADR 0035 for attribution and accounting. The simulation derives a reproducible per-tick agent order from the scenario seed. The world engine then resolves all world actions in that order. Only the @@ -146,7 +148,7 @@ The Live workspace is a grid of independently scrolling agent rail, map, context The Game API also owns one process-local experiment record. Each completed safe swarm tick is captured once, independently from the browser snapshot, and server-side export filters apply without affecting provider requests. -Schema-v13 exports may cross a separate offline archive boundary into `packages/experiment-archive`. Node's built-in SQLite stores normalized immutable research records through versioned migrations, foreign keys, prepared statements, and transactional idempotent imports. This downstream observability archive is never consulted by tick execution and cannot recover, resume, or mutate the active world. Its bounded query service is application-independent so a future read-only MCP adapter can reuse it without exposing arbitrary SQL. +Schema-v14 exports may cross a separate offline archive boundary into `packages/experiment-archive`. Node's built-in SQLite stores normalized immutable research records through versioned migrations, foreign keys, prepared statements, and transactional idempotent imports. This downstream observability archive is never consulted by tick execution and cannot recover, resume, or mutate the active world. Its bounded query service is application-independent so a future read-only MCP adapter can reuse it without exposing arbitrary SQL. World Setup uses `world-scenario-v1`. Pure preview computes the actual H3 disk, exact count, summed cell area, deterministic roster/spawns, feasibility, and warnings. Apply recomputes and atomically replaces world and experiment state. Reset reconstructs the current scenario; the Toledo default preserves legacy starts. Explicit location search crosses a replaceable server-owned adapter with no autocomplete, a one-request-per-second Nominatim limit, bounded cache/timeout, normalized results, and OpenStreetMap attribution. Manual coordinates bypass that network boundary. @@ -181,7 +183,7 @@ Equivalent legal moves are ordered reproducibly from world seed, stable agent ID - `POST /api/simulation/experiment/setup/roster/generate` — generate a deterministic roster - `POST /api/simulation/experiment/setup/location-search` — resolve a location query via the Nominatim adapter - `POST /api/simulation/experiment/export/preview` — validate filters and report subset size, retention, and cost -- `POST /api/simulation/experiment/export` — construct one schema-v13 safe JSON document +- `POST /api/simulation/experiment/export` — construct one schema-v14 safe JSON document - `POST /api/simulation/experiment/export/archive` — import the exact generated safe document into the configured local SQLite archive - `GET /api/simulation/models` — return the cached, sanitized compatible model catalog - `POST /api/simulation/models/refresh` — explicitly refresh that catalog @@ -233,7 +235,8 @@ One tick executes as follows: without a provider call (`directive-reuse`). 3. **Provider-attempt reservation.** When replanning, the tick reserves one - Zero planning attempt plus one Jev attempt per worker. Insufficient capacity + Zero planning attempt plus one individual attempt per worker, or one shared + attempt for a native worker batch. Insufficient capacity stops the tick before any provider call. 4. **Zero planning.** The service builds Zero's strategic observation from the @@ -246,10 +249,11 @@ One tick executes as follows: authoritatively. On failure the service falls back to a deterministic plan; the planner attempt is still recorded. -5. **Worker reflex dispatch (sequential).** For each worker in seeded order: - compile a local observation with the assigned directive and history; call - TypeSafe Jev; Jev selects one opaque action candidate ID with a probability - distribution and confidence. If the worker lacks an unexpired directive under +5. **Worker reflex dispatch.** Compile every local observation against the same + candidate before any dispatch. Use one native batch when the provider offers + it, otherwise a bounded individual-call pool; Jev remains supported without + native batching. Each result is attributed by worker ID, validated separately, + and mapped through its own authoritative candidate table. If the worker lacks an unexpired directive under a failed planner, the service substitutes deterministic local expansion without a Jev call. @@ -282,9 +286,9 @@ metrics cannot drift apart. Movement-pattern metrics walk each agent's accepted moves separately, classifying each step with `geographicDirectionBetweenCells`; aggregates sum direction counts and revisits and report the longest single-agent streak. All exports -use schema version 13, which carries `swarmArchitectureVersion: "zero-swarm-v1"` +use schema version 14, which carries `swarmArchitectureVersion: "zero-swarm-v1"` and independent provider-attempt accounting unconditionally. Exports at schema -version 12 and earlier are rejected outright; there is no migration path. The +version 13 and earlier are rejected outright; there is no migration path. The provider-attempt ledger is canonical for attempt counts, latency, token, and cost totals. @@ -292,7 +296,7 @@ The agent runtime follows [OpenRouter's usage-accounting contract](https://openr ## Packages -`packages/shared` owns centralized scenario limits and all public schemas, including model capabilities, swarm directives, metrics, and schema-v13 swarm tick exports. Other-agent observations remain deterministically capped at seven for larger rosters. Types are inferred from Zod. +`packages/shared` owns centralized scenario limits and all public schemas, including model capabilities, swarm directives, metrics, and schema-v14 swarm tick exports. Other-agent observations remain deterministically capped at seven for larger rosters. Types are inferred from Zod. `packages/world-engine` remains deterministic and has no model, HTTP, UI, storage, or credential dependency. It validates world actions independently. Direct proximity is derived from a separately supplied pre-action state. @@ -327,5 +331,5 @@ Structural provider failures retain the broad compatibility code plus bounded de ## Provider-attempt accounting -Provider work has an independent bounded lifecycle ledger. Schema-v13 exports +Provider work has an independent bounded lifecycle ledger. Schema-v14 exports and archive-v4 preserve safe attempt records even when no world tick commits. diff --git a/docs/EXPERIMENT_ARCHIVE.md b/docs/EXPERIMENT_ARCHIVE.md index 063f567..3f31e44 100644 --- a/docs/EXPERIMENT_ARCHIVE.md +++ b/docs/EXPERIMENT_ARCHIVE.md @@ -1,7 +1,15 @@ # Local experiment archive -The archive accepts only schema-v13 exports and rejects any other schema version outright. -Exports at schema version 12 and earlier are rejected with no migration path. Reading a +The archive accepts only schema-v14 exports and rejects any other schema version outright. + +Schema 14 adds explicit native-batch attempt membership. A shared dispatch is +one provider attempt with one aggregate charge; its required agent/turn fields +identify the first member as a storage anchor. Full membership is retained in +source JSON, and agent-filtered attempt queries include participating batches. +Per-agent cost summaries exclude shared charges; aggregate summaries count each +shared dispatch once. No per-worker cost split is inferred. Worker token counts +may be absent when only aggregate batch usage is available. See ADR 0035. +Exports at schema version 13 and earlier are rejected with no migration path. Reading a schema-v12 export requires a Git revision before PR #73 (`swarm-planner-v2`); reading a pre-swarm export (schema versions 9, 10, and 11) requires a revision before PR 1 of the zero-swarm migration. @@ -23,7 +31,7 @@ Migration 6 removes all legacy per-agent-LLM social-system tables: `turns`, `turn_number` column is renamed `tick_number`. Personality and behavior columns are dropped from `agents` and `experiments`. -The experiment archive is a durable, local research surface for completed or partially retained exports. It does not participate in an active simulation: the Game API's in-memory engine remains authoritative, and an archive write cannot change an accepted game outcome. It imports schema-v13 JSON exports only; it is not crash recovery, restartable simulation state, or a scheduler. +The experiment archive is a durable, local research surface for completed or partially retained exports. It does not participate in an active simulation: the Game API's in-memory engine remains authoritative, and an archive write cannot change an accepted game outcome. It imports schema-v14 JSON exports only; it is not crash recovery, restartable simulation state, or a scheduler. ## Storage and configuration @@ -104,7 +112,7 @@ MCP and embeddings are deferred because bounded local retrieval solves the immed Archive schema v4 stores `providerAttempts` independently. Use `pnpm experiment:db provider-attempts ` to inspect committed and uncommitted provider work. Monetary values round-trip as canonical TEXT. This -ledger is canonical for every current (schema-v13) export, which always +ledger is canonical for every current (schema-v14) export, which always carries independent attempt accounting. The SQLite archive is for analysis and is not active runtime recovery. @@ -116,7 +124,7 @@ swarm ticks. ## Zero-swarm comparisons -The archive preserves safe schema-v13 swarm tick records and independent +The archive preserves safe schema-v14 swarm tick records and independent provider attempts, but its `compare` command is not the same-scenario, per-tick swarm harness. Use `pnpm compare:offline` for the reproducible zero-swarm-vs-deterministic-workers fixture report. The runner does not diff --git a/docs/SECURITY.md b/docs/SECURITY.md index 831b493..5d628b2 100644 --- a/docs/SECURITY.md +++ b/docs/SECURITY.md @@ -107,6 +107,20 @@ Explicit scripted mode bypasses repository `.env` loading entirely. This keeps d ## Tick recovery +The optional native-batch reflex boundary attributes results by worker ID and +validates each envelope and decision independently. Duplicate IDs invalidate +only that worker; unknown IDs are ignored and missing results fall back to wait. +Providers receive copies of frozen observations, never authoritative action +maps or mutable world state. The capped individual adapter stops queued work on +abort or deadline expiry. Late results and accounting callbacks are closed out; +they cannot commit a cancelled or expired tick. See ADR 0035. + +Shared batch costs and tokens remain on one dispatch record with explicit +worker membership. The first worker is a storage anchor, not the recipient of +the shared charge. Aggregate metrics count shared usage once; individual usage +remains unknown unless separately reported. No provider secrets, raw bodies, +or private reasoning are added to this contract. + Tick recovery is server-owned and bounded inside the shared deadline. The OpenRouter planner and Jev each allow at most one retry for HTTP 429 or 529. Failed planning and worker choices use deterministic fallbacks. There is no @@ -155,7 +169,7 @@ The Game API captures only schema-validated safe observations, requested world a Export requests, agent IDs, levels, ranges, and Custom dependencies are runtime-validated. Filtering and metrics remain server-owned. The export schema -is exclusively version 13; exports carrying schema version 12 or earlier are +is exclusively version 14; exports carrying schema version 13 or earlier are rejected outright with no migration path. Reset clears swarm tick history and metrics while unlocking preserved roster assignments for the new experiment. diff --git a/docs/TESTING.md b/docs/TESTING.md index b7199a8..bc3bd07 100644 --- a/docs/TESTING.md +++ b/docs/TESTING.md @@ -79,13 +79,24 @@ malformed-entry skipping, cache TTL, stale fallback, and safe failure states. ### Experiment archive (`packages/experiment-archive`) -`archive.test.ts` covers schema-v13-only enforcement: swarm-native provenance +`archive.test.ts` covers schema-v14-only enforcement: swarm-native provenance archival, idempotent import, query service, credential-like-data rejection before -persistence, unknown-architecture-version rejection, and non-v13 schema +persistence, unknown-architecture-version rejection, and non-v14 schema rejection. ### Game API (`apps/game-api`) +Batch reflex tests use offline controlled providers and deferred promises to +assert the individual-call concurrency cap, stable native-batch attribution, +reordered/duplicate/extra/missing/malformed result policy, failure isolation, +and cancellation/deadline closure of started work. Full tick tests compare +frozen observations and committed seeded world results under reversed +completion order, prove unstarted workers create no attempt records, and reject +late callbacks after rollback. Shared batch billing tests count one dispatch +and charge once while omitting fabricated per-worker allocations. Existing +pressure, directive completion, planner-failure fallback and accounting tests +remain required. No speed thresholds or paid calls are used. See ADR 0035. + `app.test.ts` covers the API boundary: repeatable swarm-native scripted providers, health and swarm-setup contracts, atomic swarm tick through the public endpoint, and absence of legacy sequential-turn routes. diff --git a/docs/adr/0035-batch-capable-reflex-provider.md b/docs/adr/0035-batch-capable-reflex-provider.md new file mode 100644 index 0000000..dfc3fb8 --- /dev/null +++ b/docs/adr/0035-batch-capable-reflex-provider.md @@ -0,0 +1,78 @@ +# ADR 0035: Batch-capable reflex execution + +## Status + +Proposed for the local PR 3 change. + +## Context + +Workers already observe one unchanged candidate world and resolve actions +together in seeded order. Awaiting each Jev call serially is not required for +world correctness or accounting. A future native batch endpoint needs stable +worker attribution and honest shared billing before it can replace individual +transport. This change supplies that seam; it adds no production batch provider. + +## Decision + +`ReflexProvider.decide` remains supported. Optional `decideBatch` accepts one +bounded group of observations and returns envelopes keyed by worker `agentId`. +The service compiles all observations and keeps each opaque-choice action map +before dispatch. The individual adapter uses an explicit integer cap from 1 +through 8; `SimulationService.reflexConcurrencyLimit` defaults to 1. Production +Jev remains at that default. Changing its production cap is a separate rollout +decision, not a planner policy change. + +For individual caps above 1, available retry capacity is reserved after Zero +planning and assigned to a fixed prefix of the worker input order. Each job +may consume one initial slot and at most one assigned retry slot. Slots are +never reassigned during the tick, so scarce admission capacity cannot be won +by the fastest response. The anonymous reservation counter holds enough slots +for every unstarted initial even if eligible retries start first. Cap 1 retains +the existing sequential retry admission policy. Unused retry slots are released. + +Results are validated against each worker's own directive and action table. +Out-of-order envelopes are mapped by ID. Unknown workers are ignored; missing, +malformed, or failed items produce an attributed deterministic wait only for +that worker. Duplicate IDs invalidate that worker even if one duplicate is +valid. A transport failure affects the dispatched group. No result may select +another worker's candidates. The engine resolves choices in the existing seeded +order, independently of provider completion order. + +Response arrays are bounded at 64 envelopes. Non-array or oversized responses +fail the dispatched group safely; the per-worker isolation policy applies +within that bound. Concurrent individual selection with accounting requires +initial permits reserved for every input before retry slots are earmarked. + +Cancellation and the absolute tick deadline stop dispatch, settle pending +selection promptly even if a provider ignores abort, and close started attempts +once. Unstarted reservations are released without attempt records. Late +results, retries, or finalization callbacks cannot mutate the world or rewrite +closed accounting. Deadline expiry rolls back the tick; a provider failure +before the deadline retains the existing worker fallback policy. Failed Zero +planning still reuses eligible directives and uses deterministic local expansion +without a reflex call for workers with no reusable directive. + +One native batch HTTP dispatch consumes one permit and records one charge. +`ProviderAttemptRecord.batch` holds a batch ID and ordered worker/turn members; +retries share the batch ID but receive separate dispatch records. Required +`agentId` and `intendedTurnNumber` fields are the first member's storage anchor, +not individual billing attribution. The batch record carries aggregate reported +usage and cost exactly once; individual tick decisions omit token counts when +item usage is unknown. Batch transport completion does not claim that every +worker result was valid. Individual outcomes remain in worker tick records. + +Exports advance to schema 14 because attempt scope and optional worker usage +change meaning. A selected-worker export includes a participating shared batch +with its full membership and full shared charge. Aggregate metrics count it +once; per-agent cost/usage metrics exclude shared batches. Archive records keep +membership in the source JSON and per-agent cost summaries exclude shared +charges. Earlier export versions require an older Git revision. + +## Consequences + +Jev needs no batch endpoint, and existing individual providers remain compatible. +Tests use deferred promises and fake clocks, not latency comparisons or paid +providers. Provider output remains untrusted schema-validated data. Credentials, +raw bodies, prompts, and private reasoning remain outside telemetry. The +deterministic engine, strategic options, planning cadence, and player behavior +are unchanged. diff --git a/packages/agent-runtime/src/reflex-provider.ts b/packages/agent-runtime/src/reflex-provider.ts index 4bfefdb..93044b6 100644 --- a/packages/agent-runtime/src/reflex-provider.ts +++ b/packages/agent-runtime/src/reflex-provider.ts @@ -21,6 +21,12 @@ export interface ReflexProvider { observation: ReflexObservation, options?: ReflexDecisionOptions, ): Promise; + /** Optional native multi-observation request. Results are keyed because a + * provider may return them in a different order from the inputs. */ + decideBatch?( + observations: readonly ReflexObservation[], + options?: ReflexBatchDecisionOptions, + ): Promise; } export interface ReflexDecisionOptions { @@ -29,6 +35,25 @@ export interface ReflexDecisionOptions { beginAttempt?: ReflexAttemptStarter; } +export interface ReflexBatchDecisionOptions { + signal?: AbortSignal; + deadlineAtMs?: number; + /** One callback per dispatched batch HTTP attempt, including retries. */ + beginAttempt?: ReflexBatchAttemptStarter; +} + +export type ReflexBatchDecisionResult = + | { + agentId: string; + status: 'completed'; + decision: unknown; + } + | { + agentId: string; + status: 'failed'; + failure: ProviderFailure; + }; + export interface ReflexAttemptCompletion { outcome: 'completed' | 'provider-error' | 'cancelled' | 'timeout'; provider?: ProviderMetadata; @@ -44,6 +69,14 @@ export type ReflexAttemptStarter = ( kind: 'initial' | 'automatic-transport-retry', ) => ReflexAttemptFinalizer | null; +export type ReflexBatchAttemptFinalizer = ( + completion: Omit, +) => void; + +export type ReflexBatchAttemptStarter = ( + kind: 'initial' | 'automatic-transport-retry', +) => ReflexBatchAttemptFinalizer | null; + export class ReflexProviderError extends Error { constructor( readonly failure: ProviderFailure, diff --git a/packages/experiment-archive/src/archive.test.ts b/packages/experiment-archive/src/archive.test.ts index dfe24a2..376ccf2 100644 --- a/packages/experiment-archive/src/archive.test.ts +++ b/packages/experiment-archive/src/archive.test.ts @@ -36,7 +36,7 @@ async function currentExport(): Promise { describe('experiment archive', () => { it('archives a current swarm export with swarm-native provenance', async () => { const document = await currentExport(); - expect(document.schemaVersion).toBe(13); + expect(document.schemaVersion).toBe(14); expect(document.experiment).toMatchObject({ swarmPlannerContractVersion: 'swarm-planner-v2', scenario: { swarmArchitectureVersion: 'zero-swarm-v1' }, @@ -56,6 +56,124 @@ describe('experiment archive', () => { archive.close(); }); + it('preserves batch membership while charging the call once and excluding anchor per-agent usage', async () => { + const document = await currentExport(); + const [anchor, member] = document.agents; + const attempt = { + id: '018f3f38-6b7d-7db7-8e95-751b4ce2681e', + agentId: anchor!.id, + intendedTurnNumber: 1, + intendedTickNumber: 1, + kind: 'initial', + startedAt: '2026-08-13T12:00:00.000Z', + completedAt: '2026-08-13T12:00:01.000Z', + outcome: 'completed', + modelId: 'jev-1.13.0', + reasoningProfile: 'provider-default', + reservedCredits: '0.01', + actualCostCredits: '0.006', + provider: { + provider: 'typesafe', + model: 'jev-1.13.0', + latencyMs: 12, + promptTokens: 30, + completionTokens: 2, + costCredits: 0.006, + }, + batch: { + id: '018f3f38-6b7d-7db7-8e95-751b4ce2681f', + members: [ + { agentId: anchor!.id, intendedTurnNumber: 1 }, + { agentId: member!.id, intendedTurnNumber: 2 }, + ], + }, + } as const; + const batchDocument = experimentExportDocumentSchema.parse({ + ...document, + providerAttempts: [attempt], + selection: { ...document.selection, matchingProviderAttemptCount: 1 }, + }); + const archive = new ArchiveDatabase({ path: ':memory:' }); + importExperimentExport(archive, batchDocument); + const query = new ExperimentQueryService(archive); + const attempts = query.providerAttempts(batchDocument.experiment.id).rows; + expect(attempts).toHaveLength(1); + expect(attempts[0]).toMatchObject({ + actualCostCredits: '0.006', + batch: { members: [{ agentId: anchor!.id }, { agentId: member!.id }] }, + }); + expect( + query.providerAttempts(batchDocument.experiment.id, { + agent: member!.id, + fromTurn: 2, + toTurn: 2, + }).rows, + ).toHaveLength(1); + expect( + query.providerAttempts(batchDocument.experiment.id, { + agent: member!.id, + fromTurn: 1, + toTurn: 1, + }).rows, + ).toHaveLength(0); + expect( + query.providerAttempts(batchDocument.experiment.id, { + fromTurn: 2, + toTurn: 2, + }).rows, + ).toHaveLength(1); + const report = query.summary(batchDocument.experiment.id) as Record< + string, + unknown + >; + expect(JSON.stringify(report)).toContain(member!.id); + const usage = report.usage as Record; + const perAgent = (usage.byAgent ?? []) as Array>; + expect( + perAgent.every(({ providerAttempts }) => Number(providerAttempts) === 0), + ).toBe(true); + expect(usage.aggregate).toMatchObject({ + providerAttempts: 1, + knownCostCredits: 0.006, + }); + expect( + query.compare(batchDocument.experiment.id, batchDocument.experiment.id), + ).toMatchObject({ + left: { absolute: { activeAgents: 2, providerAttempts: 1 } }, + }); + const failedDocument = experimentExportDocumentSchema.parse({ + ...batchDocument, + providerAttempts: [ + { + ...attempt, + id: '018f3f38-6b7d-7db7-8e95-751b4ce26820', + outcome: 'provider-error', + failure: { + code: 'provider-http', + message: 'Batch unavailable.', + retryable: false, + }, + }, + ], + }); + importExperimentExport(archive, failedDocument); + expect( + query.failures(batchDocument.experiment.id, { + agent: member!.id, + fromTurn: 2, + toTurn: 2, + }).rows[0], + ).toMatchObject({ batch: attempt.batch }); + expect( + query.failures(batchDocument.experiment.id, { + agent: member!.id, + fromTurn: 1, + toTurn: 1, + }).rows, + ).toHaveLength(0); + archive.close(); + }); + it('imports the same swarm export idempotently', async () => { const document = await currentExport(); const archive = new ArchiveDatabase({ path: ':memory:' }); @@ -107,12 +225,12 @@ describe('experiment archive', () => { expect(experimentExportDocumentSchema.safeParse(raw).success).toBe(false); }); - it('rejects an export document with a non-v12 schema version', async () => { + it('rejects an export document with a non-v14 schema version', async () => { const raw = structuredClone(await currentExport()) as unknown as Record< string, unknown >; - raw.schemaVersion = 11; + raw.schemaVersion = 13; expect(experimentExportDocumentSchema.safeParse(raw).success).toBe(false); const archive = new ArchiveDatabase({ path: ':memory:' }); expect(() => diff --git a/packages/experiment-archive/src/query-service.ts b/packages/experiment-archive/src/query-service.ts index 2815943..e5a7032 100644 --- a/packages/experiment-archive/src/query-service.ts +++ b/packages/experiment-archive/src/query-service.ts @@ -46,6 +46,34 @@ function parseJson(value: unknown, fallback: T): T { } } +/** Agent and turn constraints must refer to the same participating worker. */ +function participantFilter(filters: DetailFilters) { + const scalar: string[] = []; + const member: string[] = []; + const values: Array = []; + if (filters.agent !== undefined) { + scalar.push('agent_id = ?'); + member.push("json_extract(member.value, '$.agentId') = ?"); + values.push(filters.agent); + } + if (filters.fromTurn !== undefined) { + scalar.push('intended_turn_number >= ?'); + member.push("json_extract(member.value, '$.intendedTurnNumber') >= ?"); + values.push(filters.fromTurn); + } + if (filters.toTurn !== undefined) { + scalar.push('intended_turn_number <= ?'); + member.push("json_extract(member.value, '$.intendedTurnNumber') <= ?"); + values.push(filters.toTurn); + } + if (!scalar.length) return null; + return { + sql: `((json_extract(source_json, '$.batch') IS NULL AND ${scalar.join(' AND ')}) + OR EXISTS (SELECT 1 FROM json_each(provider_attempts.source_json, '$.batch.members') AS member WHERE ${member.join(' AND ')}))`, + values: [...values, ...values], + }; +} + function addDecimalStrings(left: string, right: string): string { const split = (value: string) => { const [whole, fraction = ''] = value.split('.'); @@ -104,17 +132,10 @@ export class ExperimentQueryService { const limit = boundedLimit(filters.limit); const clauses: string[] = []; const values: Array = []; - if (filters.agent !== undefined) { - clauses.push('agent_id = ?'); - values.push(filters.agent); - } - if (filters.fromTurn !== undefined) { - clauses.push('intended_turn_number >= ?'); - values.push(filters.fromTurn); - } - if (filters.toTurn !== undefined) { - clauses.push('intended_turn_number <= ?'); - values.push(filters.toTurn); + const attribution = participantFilter(filters); + if (attribution) { + clauses.push(attribution.sql); + values.push(...attribution.values); } if (filters.reason !== undefined) { clauses.push('failure_code = ?'); @@ -126,7 +147,7 @@ export class ExperimentQueryService { SELECT id, intended_turn_number AS turn, intended_tick_number AS tick, agent_id AS agent, kind, model_id AS model, failure_code AS code, failure_message AS message, - validation_codes_json AS validationCodes, latency_ms AS latencyMs + validation_codes_json AS validationCodes, latency_ms AS latencyMs, source_json AS source FROM provider_attempts WHERE experiment_id = ? AND failure_code IS NOT NULL ${clauses.map((clause) => `AND ${clause}`).join('\n')} @@ -137,8 +158,12 @@ export class ExperimentQueryService { .all(experimentId, ...values, limit + 1) as Array< Record >; - for (const row of rows) + for (const row of rows) { row.validationCodes = parseJson(row.validationCodes, []); + row.batch = + parseJson>(row.source, {}).batch ?? null; + delete row.source; + } return page(rows, limit); } @@ -149,17 +174,10 @@ export class ExperimentQueryService { const limit = boundedLimit(filters.limit); const clauses: string[] = []; const values: Array = []; - if (filters.agent) { - clauses.push('agent_id = ?'); - values.push(filters.agent); - } - if (filters.fromTurn !== undefined) { - clauses.push('intended_turn_number >= ?'); - values.push(filters.fromTurn); - } - if (filters.toTurn !== undefined) { - clauses.push('intended_turn_number <= ?'); - values.push(filters.toTurn); + const attribution = participantFilter(filters); + if (attribution) { + clauses.push(attribution.sql); + values.push(...attribution.values); } if (filters.outcome) { clauses.push('outcome = ?'); @@ -173,7 +191,7 @@ export class ExperimentQueryService { completed_at AS completedAt, outcome, model_id AS model, reasoning_profile AS reasoning, provider, failure_code AS failureCode, reserved_credits AS reservedCredits, - actual_cost_credits AS actualCostCredits + actual_cost_credits AS actualCostCredits, source_json AS source FROM provider_attempts WHERE experiment_id = ? ${clauses.map((clause) => `AND ${clause}`).join('\n')} ORDER BY started_at ASC, id ASC LIMIT ? @@ -182,6 +200,11 @@ export class ExperimentQueryService { .all(experimentId, ...values, limit + 1) as Array< Record >; + for (const row of rows) { + const source = parseJson>(row.source, {}); + row.batch = source.batch ?? null; + delete row.source; + } return page(rows, limit); } @@ -288,27 +311,49 @@ export class ExperimentQueryService { ROUND(SUM(CAST(actual_cost_credits AS REAL)), 8) AS knownCostCredits, SUM(actual_cost_credits IS NULL) AS attemptsWithUnknownCost, SUM(CASE WHEN actual_cost_credits IS NULL THEN CAST(reserved_credits AS REAL) ELSE 0 END) AS reservedUnknownExposure - FROM provider_attempts WHERE experiment_id = ? + FROM provider_attempts WHERE experiment_id = ? AND json_extract(source_json, '$.batch') IS NULL GROUP BY agent_id ORDER BY agent_id `, ) .all(experimentId) as Array>; - const usageAggregate = aggregateUsage(usageByAgent); + const usageAggregateRows = this.#db + .prepare( + ` + SELECT COUNT(*) AS providerAttempts, + SUM(latency_ms) AS latencyTotalMs, + SUM(latency_ms IS NOT NULL) AS attemptsWithKnownLatency, + SUM(prompt_tokens) AS promptTokens, + SUM(completion_tokens) AS completionTokens, + SUM(total_tokens) AS totalTokens, + ROUND(SUM(CAST(actual_cost_credits AS REAL)), 8) AS knownCostCredits, + SUM(actual_cost_credits IS NULL) AS attemptsWithUnknownCost + FROM provider_attempts WHERE experiment_id = ? + `, + ) + .get(experimentId) as Record; + const usageAggregate = aggregateUsage([usageAggregateRows]); const independentAttemptRows = this.#db .prepare( `SELECT agent_id AS agent, outcome, reserved_credits AS reservedCredits, - actual_cost_credits AS actualCostCredits + actual_cost_credits AS actualCostCredits, + json_extract(source_json, '$.batch') IS NOT NULL AS isBatch FROM provider_attempts WHERE experiment_id = ? ORDER BY started_at, id`, ) .all(experimentId) as Array>; const attemptOutcomes = countBy(independentAttemptRows, 'outcome'); const attemptOutcomesByAgent = [ - ...new Set(independentAttemptRows.map(({ agent }) => String(agent))), + ...new Set( + independentAttemptRows + .filter(({ isBatch }) => !isBatch) + .map(({ agent }) => String(agent)), + ), ].map((agent) => ({ agent, outcomes: countBy( - independentAttemptRows.filter((row) => row.agent === agent), + independentAttemptRows.filter( + (row) => !row.isBatch && row.agent === agent, + ), 'outcome', ), })); @@ -581,14 +626,21 @@ function comparisonMetrics(db: DatabaseSync, experimentId: string) { .prepare( ` SELECT COUNT(*) AS total, - COUNT(DISTINCT agent_id) AS activeAgents, + (SELECT COUNT(DISTINCT participant) FROM ( + SELECT agent_id AS participant FROM provider_attempts + WHERE experiment_id = ? AND json_extract(source_json, '$.batch') IS NULL + UNION ALL + SELECT json_extract(member.value, '$.agentId') AS participant + FROM provider_attempts, json_each(provider_attempts.source_json, '$.batch.members') AS member + WHERE experiment_id = ? + )) AS activeAgents, SUM(outcome = 'accepted') AS accepted, SUM(outcome = 'provider-error') AS failed, SUM(outcome IN ('operator-skipped', 'lost-tick')) AS lost FROM provider_attempts WHERE experiment_id = ? `, ) - .get(experimentId) as Record; + .get(experimentId, experimentId, experimentId) as Record; const simulatedPlayer = db .prepare( ` diff --git a/packages/shared/src/index.test.ts b/packages/shared/src/index.test.ts index c4d294d..03095ff 100644 --- a/packages/shared/src/index.test.ts +++ b/packages/shared/src/index.test.ts @@ -26,6 +26,8 @@ import { swarmPlannerContractVersionSchema, archiveExperimentExportResponseSchema, providerAttemptRecordSchema, + reflexDecisionSchema, + reflexBatchResultEnvelopeSchema, zeroStrategicObservationSchema, } from '.'; @@ -851,6 +853,22 @@ describe('snapshot and export contracts', () => { }); describe('provider and archive contracts', () => { + it('validates batch result envelopes independently by worker', () => { + expect( + reflexBatchResultEnvelopeSchema.safeParse({ + agentId, + status: 'completed', + decision: { arbitrary: 'validated by decision schema later' }, + }).success, + ).toBe(true); + expect( + reflexBatchResultEnvelopeSchema.safeParse({ + agentId, + status: 'failed', + failure: { code: 'timeout', message: 'Timed out', retryable: true }, + }).success, + ).toBe(true); + }); it('accepts complete, partial and tiny-cost provider usage without fabricating unknowns', () => { expect( providerMetadataSchema.parse({ @@ -962,6 +980,24 @@ describe('provider and archive contracts', () => { directiveId: 'directive-1', cognitionSource: 'jev-reflex' as const, }; + expect( + reflexDecisionSchema.safeParse({ + ...reflexDecision, + inputTokens: undefined, + outputTokens: undefined, + }).success, + ).toBe(true); + expect( + reflexDecisionSchema.safeParse({ + chosenCandidateId: 'action_0', + confidence: 0.9, + probabilities: { action_0: 0.9, action_1: 0.1 }, + model: 'jev-1.13.0', + latencyMs: 12, + directiveId: 'directive-1', + cognitionSource: 'jev-reflex', + }).success, + ).toBe(true); expect( providerAttemptRecordSchema.safeParse({ ...base, @@ -992,4 +1028,56 @@ describe('provider and archive contracts', () => { }).success, ).toBe(false); }); + + it('validates batch attempts as one tick-scoped dispatch with honest anchors', () => { + const base = { + id: '018f3f38-6b7d-7db7-8e95-751b4ce2681e', + agentId, + intendedTurnNumber: 1, + intendedTickNumber: 1, + kind: 'initial', + startedAt: '2026-08-13T12:00:00.000Z', + modelId: 'jev-1.13.0', + reasoningProfile: 'provider-default', + reservedCredits: '0.01', + batch: { + id: '018f3f38-6b7d-7db7-8e95-751b4ce2681f', + members: [ + { agentId, intendedTurnNumber: 1 }, + { + agentId: '22222222-2222-4222-8222-222222222222', + intendedTurnNumber: 2, + }, + ], + }, + }; + expect( + providerAttemptRecordSchema.safeParse({ ...base, outcome: 'in-flight' }) + .success, + ).toBe(true); + expect( + providerAttemptRecordSchema.safeParse({ + ...base, + agentId: '22222222-2222-4222-8222-222222222222', + outcome: 'in-flight', + }).success, + ).toBe(false); + expect( + providerAttemptRecordSchema.safeParse({ + ...base, + intendedTickNumber: undefined, + outcome: 'in-flight', + }).success, + ).toBe(false); + expect( + providerAttemptRecordSchema.safeParse({ + ...base, + batch: { + ...base.batch, + members: [base.batch.members[0], base.batch.members[0]], + }, + outcome: 'in-flight', + }).success, + ).toBe(false); + }); }); diff --git a/packages/shared/src/index.ts b/packages/shared/src/index.ts index fbb773a..092df00 100644 --- a/packages/shared/src/index.ts +++ b/packages/shared/src/index.ts @@ -1095,8 +1095,8 @@ export const reflexDecisionSchema = z replanProbability: z.number().finite().min(0).max(1).optional(), model: modelIdSchema, latencyMs: z.number().finite().nonnegative(), - inputTokens: z.number().int().nonnegative(), - outputTokens: z.number().int().nonnegative(), + inputTokens: z.number().int().nonnegative().optional(), + outputTokens: z.number().int().nonnegative().optional(), directiveId: z.string().trim().min(1).max(80), cognitionSource: cognitionSourceSchema, }) @@ -1372,6 +1372,23 @@ export const providerFailureSchema = z.object({ }); export type ProviderFailure = z.infer; +export const reflexBatchResultEnvelopeSchema = z.discriminatedUnion('status', [ + z + .object({ + agentId: agentIdSchema, + status: z.literal('completed'), + decision: z.unknown(), + }) + .strict(), + z + .object({ + agentId: agentIdSchema, + status: z.literal('failed'), + failure: providerFailureSchema, + }) + .strict(), +]); + export const modelAttemptSchema = z.object({ attemptNumber: z.number().int().positive(), kind: z.enum([ @@ -1393,6 +1410,44 @@ export type ModelAttempt = z.infer; export const providerAttemptIdSchema = z.uuid().brand<'ProviderAttemptId'>(); export type ProviderAttemptId = z.infer; +export const providerAttemptBatchSchema = z + .object({ + id: z.uuid(), + members: z + .array( + z + .object({ + agentId: agentIdSchema, + intendedTurnNumber: z.number().int().positive(), + }) + .strict(), + ) + .min(1) + .max(WORLD_SCENARIO_LIMITS.maximumAgents - 1), + }) + .strict() + .superRefine((batch, context) => { + if ( + new Set(batch.members.map(({ agentId }) => agentId)).size !== + batch.members.length + ) + context.addIssue({ + code: 'custom', + path: ['members'], + message: 'Batch members must have unique agent IDs.', + }); + if ( + new Set(batch.members.map(({ intendedTurnNumber }) => intendedTurnNumber)) + .size !== batch.members.length + ) + context.addIssue({ + code: 'custom', + path: ['members'], + message: 'Batch members must have unique intended turns.', + }); + }); +export type ProviderAttemptBatch = z.infer; + export const providerAttemptOutcomeSchema = z.enum([ 'completed', 'provider-error', @@ -1429,6 +1484,8 @@ export const providerAttemptRecordSchema = z agentId: agentIdSchema, intendedTurnNumber: z.number().int().positive(), intendedTickNumber: z.number().int().positive().optional(), + /** One HTTP dispatch shared by these workers; scalar fields anchor member zero. */ + batch: providerAttemptBatchSchema.optional(), kind: modelAttemptSchema.shape.kind, startedAt: z.iso.datetime(), completedAt: z.iso.datetime().optional(), @@ -1448,6 +1505,30 @@ export const providerAttemptRecordSchema = z .strict() .superRefine((attempt, context) => { const finalized = attempt.outcome !== 'in-flight'; + if (attempt.batch) { + const first = attempt.batch.members[0]; + if ( + !first || + attempt.agentId !== first.agentId || + attempt.intendedTurnNumber !== first.intendedTurnNumber + ) + context.addIssue({ + code: 'custom', + path: ['batch'], + message: 'Scalar attribution must match the first batch member.', + }); + if (attempt.intendedTickNumber === undefined) + context.addIssue({ + code: 'custom', + path: ['intendedTickNumber'], + message: 'Batch attempts require a tick number.', + }); + if (attempt.reflexDecision || attempt.swarmPlan) + context.addIssue({ + code: 'custom', + message: 'Batch decisions belong on individual worker tick records.', + }); + } if (finalized !== Boolean(attempt.completedAt)) context.addIssue({ code: 'custom', @@ -2710,7 +2791,7 @@ export type ExperimentExportWorldState = z.infer< const experimentExportDocumentObjectSchema = z .object({ - schemaVersion: z.literal(13), + schemaVersion: z.literal(14), generatedAt: z.iso.datetime(), experiment: experimentManifestSchema, retention: experimentRetentionSchema, @@ -2808,7 +2889,16 @@ const experimentExportDocumentObjectSchema = z }); if ( attempts.some( - ({ agentId }) => !selectedIds.has(agentId) || !exportedIds.has(agentId), + (attempt) => + !exportedIds.has(attempt.agentId) || + (attempt.batch + ? !attempt.batch.members.some(({ agentId }) => + selectedIds.has(agentId), + ) + : !selectedIds.has(attempt.agentId)) || + attempt.batch?.members.some( + ({ agentId }) => !exportedIds.has(agentId), + ), ) ) context.addIssue({ diff --git a/tests/e2e/world-lab.spec.ts b/tests/e2e/world-lab.spec.ts index 0179ece..cc6edbf 100644 --- a/tests/e2e/world-lab.spec.ts +++ b/tests/e2e/world-lab.spec.ts @@ -109,7 +109,7 @@ test('runs a deterministic swarm tick and exports safe telemetry', async ({ const exported = experimentExportDocumentSchema.parse( JSON.parse(await readFile(downloadedPath!, 'utf8')), ); - expect(exported.schemaVersion).toBe(13); + expect(exported.schemaVersion).toBe(14); expect(exported.experiment.scenario?.swarmArchitectureVersion).toBe( 'zero-swarm-v1', );