vpp-ai-platform/packages/services/test/chain.test.ts

294 lines
17 KiB
TypeScript
Raw Normal View History

M3: Mastra runtime, safety chain, two agents, Case Desk v1 - packages/domain: safety-chain objects (ValidationResult, SimulationResult, EnvelopeMatch, StaleDenial, ExecutionReceipt, LineageRef/BidProposalDraft, RouterDecision, BidExportFile) + fixtures on both sides. - packages/services: proposal digest; PolicyEngine + hubei-spot-bidding pack (digest-valid, bid-format, price-limits, quantity-non-negative, ledger-consistency, lineage-integrity, originator-permission — each with pass/fail tests); EnvelopeService; AuthorityService (fresh check, permits, revoke, gateway validate); FileExportGateway (idempotent receipts); Memory/File EventBus; LineageRecorder + P2 assembler; RevenueScenario simulator; CaseDeskService; FsRepository; skill HTTP client moved here. - packages/runtime: createRuntime (LibSQL storage, per-runtime workflow factories), proposal-lifecycle (rule check → simulation → envelope gate with suspend/resume → fresh check + permit → release), day-ahead-situation, day-ahead-bid, TriggerService (scheduled/event/manual), LlmPort (Mastra/Scripted/Null), Case Desk HTTP API, dev entry point. - Tests: all eight docs/01 invariants, docs/07 06:00→08:30 end to end with LLM down, restart survival of a suspended approval, permit expiry and revocation, replay of a released proposal, trigger scheduling. 141 TS + 60 Python tests. - Known gaps: ledger not yet persisted (replayed on restart); STALE ends the run instead of looping to rule check; synthetic data stands in for historical replay. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UoYoGYzHkFyv3ALenkRPhA
2026-09-02 06:55:55 -04:00
import { mkdtempSync, readFileSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { describe, expect, it } from 'vitest'
import { BidProposalDraft, Proposal as ProposalSchema } from '@vpp/domain'
import { AuthorityService, assertApproverIndependent, isStaleDenial } from '../src/authority.js'
import { CaseDeskService } from '../src/casedesk.js'
import { proposalDigest } from '../src/digest.js'
import { EnvelopeService } from '../src/envelope.js'
import { FileEventBus, MemoryEventBus } from '../src/events.js'
import { FileExportGateway, GatewayRejected } from '../src/gateway.js'
import { LineageRefError, assembleBidProposal } from '../src/lineage.js'
import { FsRepository } from '../src/relational.js'
import { RevenueScenarioSimulator } from '../src/simulation.js'
import { MemorySnapshotStore } from '../src/snapshot.js'
import { DATE, NOW, PACK_VERSION, approvalFor, curve, envelope, priceForecast, world } from './helpers.js'
const HUMAN_ROLES = ['senior-trader', 'ops-lead', 'risk-officer']
const freshCtx = (w: ReturnType<typeof world>, over: Record<string, unknown> = {}) => ({
now: NOW, ledger: w.ledger, currentPolicyPackVersion: PACK_VERSION, offlineResources: [], permitTtlMs: 3_600_000, ...over,
})
describe('policy engine: hubei-spot-bidding pack (each rule has a pass and a fail case)', () => {
it('passes a well-formed, traceable, in-band bid', () => {
const w = world()
const r = w.engine.check(w.proposal, { id: 'hubei-spot-bidding', version: PACK_VERSION }, w.ruleCtx())
expect(r.ok, JSON.stringify(r.violations)).toBe(true)
expect(r.rules_evaluated).toContain('lineage-integrity')
expect(r.policy_pack.version).toBe(PACK_VERSION)
})
it('digest-valid: a tampered payload no longer matches its digest', () => {
const w = world()
const tampered = { ...w.proposal, payload: { ...w.proposal.payload, expected_revenue_yuan: '999.00' } }
const r = w.engine.check(tampered, { id: 'hubei-spot-bidding', version: PACK_VERSION }, w.ruleCtx())
expect(r.violations.map((v) => v.rule_id)).toContain('digest-valid')
})
it('price-limits: offer above the cap is rejected', () => {
const w = world({ price: '1500.01' })
const r = w.engine.check(w.proposal, { id: 'hubei-spot-bidding', version: PACK_VERSION }, w.ruleCtx())
expect(r.ok).toBe(false)
expect(r.violations[0]!.rule_id).toBe('price-limits')
})
it('ledger-consistency: energy outside the cascade band, or a stale ledger version, is rejected', () => {
const over = world({ qty: '14.0' }) // 1344 MWh > 1260
let r = over.engine.check(over.proposal, { id: 'hubei-spot-bidding', version: PACK_VERSION }, over.ruleCtx())
expect(r.violations.some((v) => v.rule_id === 'ledger-consistency' && /outside/.test(v.message))).toBe(true)
const w = world()
w.ledger.append({ id: 'x', timescale: 'MONTHLY', period: '2026-04', kind: 'CONTRACT', energy_mwh: '1.0', curve: null, source_ref: 'c', expected_version: 1 })
r = w.engine.check(w.proposal, { id: 'hubei-spot-bidding', version: PACK_VERSION }, w.ruleCtx())
expect(r.violations.some((v) => v.rule_id === 'ledger-consistency' && /ledger_version/.test(v.message))).toBe(true)
})
it('lineage-integrity: a number not present in any referenced tool output is rejected (P2)', () => {
const w = world()
const forged = { ...w.proposal, payload: { ...w.proposal.payload, quantities_mwh: curve('12.6') } }
forged.digest = proposalDigest(forged) // even a re-digested forgery fails: the number has no provenance
const r = w.engine.check(forged, { id: 'hubei-spot-bidding', version: PACK_VERSION }, w.ruleCtx())
expect(r.violations.map((v) => v.rule_id)).toContain('lineage-integrity')
expect(r.violations.map((v) => v.rule_id)).not.toContain('digest-valid')
})
it('originator-permission: an agent outside the allowlist may not originate a BID (I2)', () => {
const w = world({ agent: 'load-control-agent' })
const r = w.engine.check(w.proposal, { id: 'hubei-spot-bidding', version: PACK_VERSION }, w.ruleCtx())
expect(r.violations.map((v) => v.rule_id)).toContain('originator-permission')
})
it('old pack versions stay registered for replay (I7)', () => {
const w = world()
expect(w.engine.current('hubei-spot-bidding').version).toBe(PACK_VERSION)
expect(() => w.engine.pack('hubei-spot-bidding', '2025.01')).toThrow(/not registered/)
})
})
describe('P2 assembler: numbers enter a proposal only by lineage reference', () => {
it('rejects a draft carrying a literal number at the schema boundary', () => {
const draft = { market_date: DATE, prices_ref: { tool_call_id: 'tc-1', path: 'p' }, quantities_ref: { tool_call_id: 'tc-1', path: 'q' }, expected_revenue_ref: '510600.00', rationale: 'x' }
expect(BidProposalDraft.safeParse(draft).success).toBe(false)
})
it('rejects a reference to a tool call that is not in lineage', () => {
const w = world()
expect(() =>
assembleBidProposal(
{ market_date: DATE, prices_ref: { tool_call_id: 'tc-99', path: 'prices_yuan_per_mwh' }, quantities_ref: { tool_call_id: 'tc-99', path: 'quantities_mwh' }, expected_revenue_ref: { tool_call_id: 'tc-99', path: 'expected_revenue_yuan' }, rationale: 'x' },
{ id: 'p', originator: w.proposal.originator, lineage: w.proposal.lineage, ledger_version: 1, policy_pack_version: PACK_VERSION, snapshots: w.snapshots, now: NOW },
),
).toThrow(LineageRefError)
})
it('assembled proposal validates against the domain schema and its digest is reproducible', () => {
const w = world()
expect(ProposalSchema.safeParse(w.proposal).success).toBe(true)
expect(proposalDigest(w.proposal)).toBe(w.proposal.digest)
expect(w.proposal.lineage.tool_calls.map((t) => t.tool)).toEqual(['price-forecast', 'bid-optimization-milp'])
})
})
describe('envelope gate (P3)', () => {
it('empty envelope set → everything goes to a human (autonomy dial at zero)', () => {
const w = world()
const m = new EnvelopeService().match(w.proposal, { now: NOW, priceBaseline: priceForecast().quantiles.p50 })
expect(m.within).toBe(false)
expect(m.envelope_id).toBeNull()
expect(m.required_level).toBe('L2')
})
it('auto-approves inside bounds and reports which envelope', () => {
const w = world()
const svc = new EnvelopeService()
svc.register(envelope({ price_deviation_pct: '5.0', max_energy_mwh: '1500.0' }))
const m = svc.match(w.proposal, { now: NOW, priceBaseline: priceForecast().quantiles.p50 })
expect(m.within).toBe(true)
expect(m.envelope_id).toBe('env-bid-001')
})
it('exceeding a bound, an expired window, or SUSPENDED status all fall out of the envelope', () => {
const w = world()
const outOfEnergy = new EnvelopeService(); outOfEnergy.register(envelope({ max_energy_mwh: '1000.0' }))
expect(outOfEnergy.match(w.proposal, { now: NOW, priceBaseline: null }).reasons[0]).toMatch(/max_energy_mwh/)
const expired = new EnvelopeService(); expired.register(envelope({ max_energy_mwh: '1500.0' }, { validity: { from: '2026-01-01T00:00:00Z', to: '2026-02-01T00:00:00Z' } }))
expect(expired.match(w.proposal, { now: NOW, priceBaseline: null }).within).toBe(false)
const suspended = new EnvelopeService(); suspended.register(envelope({ max_energy_mwh: '1500.0' }, { status: 'SUSPENDED' }))
expect(suspended.match(w.proposal, { now: NOW, priceBaseline: null }).within).toBe(false)
const priced = world({ price: '450.00' }) // 5.8% above 425.50 baseline
const tight = new EnvelopeService(); tight.register(envelope({ price_deviation_pct: '5.0' }))
expect(tight.match(priced.proposal, { now: NOW, priceBaseline: priceForecast().quantiles.p50 }).reasons[0]).toMatch(/deviates/)
})
})
describe('approval identity (I1/I2)', () => {
it('rejects the originating agent, any tool in the chain, and non-human roles', () => {
const w = world()
expect(() => assertApproverIndependent(w.proposal, { id: 'trading-agent', role: 'senior-trader' }, HUMAN_ROLES)).toThrow(/origination chain/)
expect(() => assertApproverIndependent(w.proposal, { id: 'bid-optimization-milp', role: 'senior-trader' }, HUMAN_ROLES)).toThrow(/origination chain/)
expect(() => assertApproverIndependent(w.proposal, { id: 'user-x', role: 'agent' }, HUMAN_ROLES)).toThrow(/AI cannot approve/)
expect(() => assertApproverIndependent(w.proposal, { id: 'user-x', role: 'senior-trader' }, HUMAN_ROLES)).not.toThrow()
})
})
describe('authority: fresh check, permits, revocation (I3/I4)', () => {
it('issues a permit bound to the digest with the narrowest effect limits', () => {
const w = world()
const permit = w.authority.authorize(w.proposal, { approvals: [approvalFor(w.proposal, 1, { scope: { effect_type: 'BID', limits: { max_energy_mwh: '1250.0' } } })], envelopeMatch: null }, freshCtx(w))
expect(isStaleDenial(permit)).toBe(false)
if (isStaleDenial(permit)) throw new Error()
expect(permit.proposal_digest).toBe(w.proposal.digest)
expect(permit.effect_limits['max_energy_mwh']).toBe('1250.0')
expect(permit.expires_at).toBe('2026-03-14T09:00:00.000Z')
})
it('goes STALE when the digest, window, ledger, policy pack or resources changed since approval (I3)', () => {
const w = world()
const a = approvalFor(w.proposal, 1)
const stale = (basis: Parameters<AuthorityService['authorize']>[1], ctx = freshCtx(w)) => {
const r = w.authority.authorize(w.proposal, basis, ctx)
if (!isStaleDenial(r)) throw new Error('expected stale')
return r.reasons.join(' | ')
}
expect(stale({ approvals: [{ ...a, proposal_digest: 'e'.repeat(64) }], envelopeMatch: null })).toMatch(/digest/)
expect(stale({ approvals: [a], envelopeMatch: null }, freshCtx(w, { now: '2026-03-15T10:00:00Z' }))).toMatch(/validity/)
expect(stale({ approvals: [{ ...a, evidence_versions: { ...a.evidence_versions, ledger_version: 0 } }], envelopeMatch: null })).toMatch(/ledger version/)
expect(stale({ approvals: [a], envelopeMatch: null }, freshCtx(w, { currentPolicyPackVersion: '2026.04' }))).toMatch(/policy pack/)
expect(stale({ approvals: [a], envelopeMatch: null }, freshCtx(w, { offlineResources: ['res-storage-01'] }))).toMatch(/offline/)
expect(stale({ approvals: [{ ...a, decision: 'REJECT' }], envelopeMatch: null })).toMatch(/REJECT/)
expect(stale({ approvals: [], envelopeMatch: null })).toMatch(/no approval basis/)
})
it('an auto-approval is a valid basis only while its evidence is current', () => {
const w = world()
const match = { proposal_digest: w.proposal.digest, within: true, envelope_id: 'env-bid-001', required_level: 'L1' as const, reasons: [], checked_at: NOW }
expect(isStaleDenial(w.authority.authorize(w.proposal, { approvals: [], envelopeMatch: match }, freshCtx(w)))).toBe(false)
w.ledger.append({ id: 'y', timescale: 'MONTHLY', period: '2026-04', kind: 'CONTRACT', energy_mwh: '1.0', curve: null, source_ref: 'c', expected_version: 1 })
expect(isStaleDenial(w.authority.authorize(w.proposal, { approvals: [], envelopeMatch: match }, freshCtx(w)))).toBe(true)
})
it('validate: expiry, revocation, digest mismatch and unknown permits are all rejected (I4)', () => {
const w = world()
const permit = w.authority.authorize(w.proposal, { approvals: [approvalFor(w.proposal, 1)], envelopeMatch: null }, freshCtx(w))
if (isStaleDenial(permit)) throw new Error()
expect(w.authority.validate(permit, w.proposal, NOW)).toEqual([])
expect(w.authority.validate(permit, w.proposal, '2026-03-14T09:00:00.000Z')[0]).toMatch(/expired/)
expect(w.authority.validate({ ...permit, proposal_digest: 'e'.repeat(64) }, w.proposal, NOW)[0]).toMatch(/digest/)
expect(w.authority.validate({ ...permit, id: 'permit-forged' }, w.proposal, NOW)[0]).toMatch(/not on record/)
w.authority.revoke(permit.id, 'storage offline', '2026-03-14T08:30:00Z')
expect(w.authority.validate(permit, w.proposal, NOW)[0]).toMatch(/revoked/)
})
})
describe('gateway: only (Proposal, Permit) pairs; idempotent receipts (I4/I5)', () => {
it('exports a bid file once and returns the same receipt on re-dispatch', () => {
const w = world()
const dir = mkdtempSync(join(tmpdir(), 'vpp-export-'))
const gw = new FileExportGateway(w.authority, dir)
const permit = w.authority.authorize(w.proposal, { approvals: [approvalFor(w.proposal, 1)], envelopeMatch: null }, freshCtx(w))
if (isStaleDenial(permit)) throw new Error()
const r1 = gw.dispatch(w.proposal, permit, NOW)
const r2 = gw.dispatch(w.proposal, permit, '2026-03-14T08:05:00Z')
expect(r2).toEqual(r1)
expect(gw.receipts.list()).toHaveLength(1)
const file = JSON.parse(readFileSync(r1.artifact_ref, 'utf8'))
expect(file.proposal_digest).toBe(w.proposal.digest)
expect(file.quantities_mwh.values[0]).toBe('12.5')
})
it('refuses a late release, a revoked permit, and a bare proposal with a forged permit', () => {
const w = world()
const gw = new FileExportGateway(w.authority, mkdtempSync(join(tmpdir(), 'vpp-export-')))
const permit = w.authority.authorize(w.proposal, { approvals: [approvalFor(w.proposal, 1)], envelopeMatch: null }, freshCtx(w))
if (isStaleDenial(permit)) throw new Error()
expect(() => gw.dispatch(w.proposal, permit, '2026-03-15T09:00:01Z')).toThrow(GatewayRejected)
expect(() => gw.dispatch(w.proposal, { ...permit, id: 'permit-forged' }, NOW)).toThrow(/not on record/)
w.authority.revoke(permit.id, 'kill switch', NOW)
expect(() => gw.dispatch(w.proposal, permit, NOW)).toThrow(/revoked/)
expect(gw.receipts.list()).toHaveLength(0)
})
})
describe('revenue scenario simulation', () => {
it('is deterministic for a seed and does not alert on a sane price-taker bid', () => {
const w = world()
const sim = new RevenueScenarioSimulator()
const ctx = { now: NOW, priceForecast: priceForecast(), marginalCostYuanPerMwh: '0', worstCaseLossBudgetYuan: '100000' }
const a = sim.simulate(w.proposal, ctx)
const b = sim.simulate(w.proposal, ctx)
expect(a).toEqual(b)
expect(a.alert).toBe(false)
expect(a.scenario_count).toBe(1000)
expect(Number(a.metrics['revenue_p05_yuan'])).toBeLessThan(Number(a.metrics['revenue_p95_yuan']))
})
it('alerts when the P05 scenario breaches the loss budget, and when there is no forecast', () => {
const w = world()
const sim = new RevenueScenarioSimulator()
const costly = sim.simulate(w.proposal, { now: NOW, priceForecast: priceForecast(), marginalCostYuanPerMwh: '600', worstCaseLossBudgetYuan: '1000' })
expect(costly.alert).toBe(true)
expect(costly.alerts[0]).toMatch(/loss budget/)
expect(sim.simulate(w.proposal, { now: NOW, priceForecast: null, marginalCostYuanPerMwh: '0', worstCaseLossBudgetYuan: '1' }).alerts[0]).toMatch(/no price forecast/)
})
})
describe('case desk projections', () => {
it('opens a case, tracks proposal/approval/permit refs, and expands lineage to tool IO', () => {
const w = world()
const desk = new CaseDeskService({ authority: w.authority, snapshots: w.snapshots, events: w.events, clock: () => NOW })
const c = desk.open({ id: 'case-1', kind: 'DAY_AHEAD_BID', objective: 'D bid', owner: 'user-trader-01', deadline: '2026-03-14T09:00:00Z' })
desk.recordProposal(c.id, w.proposal)
desk.requestApproval({ proposal_id: w.proposal.id, proposal_digest: w.proposal.digest, case_id: c.id, run_id: 'run-1', required_level: 'L2', reasons: ['no envelope'], since: NOW })
expect(desk.inbox().map((p) => p.run_id)).toEqual(['run-1'])
expect(desk.read(c.id).case.status).toBe('AWAITING_APPROVAL')
desk.recordApproval(c.id, approvalFor(w.proposal, 1))
const permit = w.authority.authorize(w.proposal, { approvals: [approvalFor(w.proposal, 1)], envelopeMatch: null }, freshCtx(w))
if (isStaleDenial(permit)) throw new Error()
desk.recordPermit(c.id, permit)
const view = desk.read(c.id)
expect(view.approvals[0]!.approver.id).toBe('user-trader-01')
expect(view.permits[0]!.id).toBe('permit-001')
expect(view.events.map((e) => e.event_type)).toEqual(['CaseOpened', 'ProposalStatusChanged', 'ApprovalRequested', 'HumanDecision'])
const lin = desk.expandLineage(w.proposal.id)
expect(lin.tool_calls[1]!.tool).toBe('bid-optimization-milp')
expect((lin.tool_calls[1]!.outputs as { expected_revenue_yuan: string }).expected_revenue_yuan).toBe('510600.00')
})
})
describe('persistence for restart survival', () => {
it('FsRepository and FileEventBus reload their state in a new instance', () => {
const dir = mkdtempSync(join(tmpdir(), 'vpp-persist-'))
const w = world()
const repo = new FsRepository(join(dir, 'proposals.json'), ProposalSchema, (p) => p.id)
repo.put(w.proposal)
repo.put({ ...w.proposal, status: 'PENDING_HUMAN' }, 1)
const reloaded = new FsRepository(join(dir, 'proposals.json'), ProposalSchema, (p) => p.id)
expect(reloaded.get('prop-001')?.version).toBe(2)
expect(reloaded.get('prop-001')?.value.status).toBe('PENDING_HUMAN')
const bus = new FileEventBus(join(dir, 'events.jsonl'), () => NOW)
bus.append({ event_type: 'A', payload: { n: 1 }, correlation_id: 'c1' })
bus.append({ event_type: 'B', payload: { n: 2 }, correlation_id: 'c1', causation_id: 'evt-000001' })
const again = new FileEventBus(join(dir, 'events.jsonl'), () => NOW)
expect(again.list().map((e) => e.event_type)).toEqual(['A', 'B'])
expect(again.append({ event_type: 'C', payload: {}, correlation_id: 'c2' }).event_id).toBe('evt-000003')
expect(new MemoryEventBus().list()).toEqual([])
})
})