import { existsSync } from 'node:fs' import { describe, expect, it } from 'vitest' import type { Proposal } from '@vpp/domain' import { proposalDigest } from '@vpp/services' import { ScriptedLlm } from '../src/llm.js' import { TriggerService } from '../src/trigger.js' import { MARKET_DATE, envelope, harness } from './helpers.js' const APPROVER = { id: 'user-trader-01', role: 'senior-trader' } async function submitBid(h: Awaited>) { const wf = h.rt.mastra.getWorkflow('dayAheadBid') const run = await wf.createRun() const res = await run.start({ inputData: { market_date: MARKET_DATE } }) if (res.status !== 'success') throw new Error(`bid workflow ${res.status}: ${JSON.stringify((res as { error?: unknown }).error ?? res)}`) return res.result } const statuses = (h: Awaited>, caseId: string) => h.rt.ctx.events.list({ event_type: 'ProposalStatusChanged', correlation_id: caseId }).map((e) => (e.payload as { status: string }).status) describe('docs/07 timeline D-1 06:00 → 08:30 (LLM down: I6)', () => { it('situation → bid → lifecycle end-to-end, auto-approved inside the envelope, released as a bid file', async () => { const h = await harness({ llm: null }) h.rt.ctx.envelopes.register(envelope({ price_deviation_pct: '5.0', max_energy_mwh: '1500.0' })) // 06:00 situation const sit = await h.rt.mastra.getWorkflow('dayAheadSituation').createRun() const sres = await sit.start({ inputData: { market_date: MARKET_DATE } }) expect(sres.status).toBe('success') if (sres.status !== 'success') return expect(sres.result.llm_used).toBe(false) expect(sres.result.report.findings[0]!.summary).toMatch(/template/) expect(sres.result.report.risk_level).toBe('LOW') expect(h.rt.ctx.events.list({ event_type: 'SituationPublished' })).toHaveLength(1) // 08:00 bid → 08:30 chain const out = await submitBid(h) expect(out.llm_used).toBe(false) expect(out.lifecycle_status).toBe('success') expect(out.lifecycle?.outcome).toBe('RELEASED') expect(statuses(h, out.case_id)).toEqual(['DRAFT', 'RULE_CHECK', 'SIMULATION', 'ENVELOPE_CHECK', 'AUTO_APPROVED', 'FRESH_CHECK', 'AUTHORIZED', 'RELEASED']) // P2: every payload number traces to the MILP tool output; P7: the ledger recorded the bid. const lin = h.rt.ctx.caseDesk.expandLineage(out.proposal_id) expect(lin.tool_calls.map((t) => t.tool)).toEqual(['load-forecast', 'pv-forecast', 'price-forecast', 'bid-optimization-milp']) expect((lin.tool_calls[3]!.outputs as { expected_revenue_yuan: string }).expected_revenue_yuan).toBe(out.proposal.payload.expected_revenue_yuan) expect(h.rt.ctx.ledger.viewByTimescale('DAY_AHEAD')).toHaveLength(1) expect(out.proposal.lineage.ledger_version).toBe(1) // I5: receipt + file; case closed with the receipt as outcome. const view = h.rt.ctx.caseDesk.read(out.case_id) expect(view.case.status).toBe('CLOSED_DONE') expect(view.permits).toHaveLength(1) expect(existsSync(h.rt.ctx.gateway.receipts.list()[0]!.value.artifact_ref)).toBe(true) await h.rt.close() }) it('with an LLM, the agents narrate but the numbers are byte-identical to the LLM-down run', async () => { const draftFor = (prompt: string) => { const tc = /tool call id is (tc-[\w-]+)/.exec(prompt)![1]! return { market_date: MARKET_DATE, prices_ref: { tool_call_id: tc, path: 'prices_yuan_per_mwh' }, quantities_ref: { tool_call_id: tc, path: 'quantities_mwh' }, expected_revenue_ref: { tool_call_id: tc, path: 'expected_revenue_yuan' }, rationale: 'Evening peak carries the position; offer as price-taker.' } } const llm = new ScriptedLlm({ 'analysis-agent': { findings: [{ kind: 'TREND', summary: 'Load near seasonal norm; no anomalies.' }] }, 'trading-agent': draftFor, }) const withLlm = await harness({ llm }) const without = await harness({ llm: null }) for (const h of [withLlm, without]) h.rt.ctx.envelopes.register(envelope({ max_energy_mwh: '1500.0' })) const a = await submitBid(withLlm) const b = await submitBid(without) expect(a.llm_used).toBe(true) expect(b.llm_used).toBe(false) expect(a.proposal.payload).toEqual(b.proposal.payload) expect(llm.prompts.map((p) => p.agentId)).toContain('trading-agent') await withLlm.rt.close() await without.rt.close() }) it('an LLM draft that references a tool call outside lineage cannot become a proposal (P2)', async () => { const llm = new ScriptedLlm({ 'trading-agent': { market_date: MARKET_DATE, prices_ref: { tool_call_id: 'tc-forged', path: 'prices_yuan_per_mwh' }, quantities_ref: { tool_call_id: 'tc-forged', path: 'quantities_mwh' }, expected_revenue_ref: { tool_call_id: 'tc-forged', path: 'expected_revenue_yuan' }, rationale: 'x' }, }) const h = await harness({ llm }) const run = await h.rt.mastra.getWorkflow('dayAheadBid').createRun() const res = await run.start({ inputData: { market_date: MARKET_DATE } }) expect(res.status).toBe('failed') if (res.status === 'failed') expect(String(res.error.message)).toMatch(/not in this proposal's lineage/) expect(h.rt.ctx.caseDesk.proposals.list()).toHaveLength(0) await h.rt.close() }) }) describe('human approval path (PENDING_HUMAN → resume)', () => { it('suspends outside the envelope, surfaces in the inbox, and a human approve releases it', async () => { const h = await harness({ llm: null }) // no envelopes → everything manual const out = await submitBid(h) expect(out.lifecycle_status).toBe('suspended') expect(statuses(h, out.case_id).at(-1)).toBe('PENDING_HUMAN') const inbox = h.rt.ctx.caseDesk.inbox() expect(inbox).toHaveLength(1) expect(inbox[0]!.run_id).toBe(out.lifecycle_run_id) expect(inbox[0]!.reasons[0]).toMatch(/no ACTIVE envelope/) const triggers = new TriggerService(h.rt) const res = await triggers.decide(out.lifecycle_run_id, { decision: 'approve', approver: APPROVER, comment: 'ok' }) expect(res.status).toBe('success') if (res.status === 'success') expect(res.result.outcome).toBe('RELEASED') expect(statuses(h, out.case_id).slice(-4)).toEqual(['APPROVED', 'FRESH_CHECK', 'AUTHORIZED', 'RELEASED']) const view = h.rt.ctx.caseDesk.read(out.case_id) expect(view.approvals[0]!.approver).toEqual(APPROVER) expect(view.approvals[0]!.proposal_digest).toBe(out.proposal.digest) await h.rt.close() }) it('I1: the AI (or anyone in the origination chain) cannot approve; the run stays suspended', async () => { const h = await harness({ llm: null }) const out = await submitBid(h) const triggers = new TriggerService(h.rt) for (const approver of [{ id: 'trading-agent', role: 'senior-trader' }, { id: 'user-x', role: 'agent' }, { id: 'bid-optimization-milp', role: 'ops-lead' }]) { const res = await triggers.decide(out.lifecycle_run_id, { decision: 'approve', approver }) expect(res.status).toBe('suspended') } expect(h.rt.ctx.events.list({ event_type: 'ApprovalRefused' })).toHaveLength(3) expect(h.rt.ctx.caseDesk.read(out.case_id).approvals).toHaveLength(0) const still = h.rt.ctx.caseDesk.proposals.get(out.proposal_id)!.value.status expect(still).toBe('PENDING_HUMAN') await h.rt.close() }) it('a human reject ends the run as REJECTED', async () => { const h = await harness({ llm: null }) const out = await submitBid(h) const res = await new TriggerService(h.rt).decide(out.lifecycle_run_id, { decision: 'reject', approver: APPROVER, comment: 'too aggressive' }) expect(res.status).toBe('success') if (res.status === 'success') expect(res.result.outcome).toBe('REJECTED') expect(h.rt.ctx.caseDesk.proposals.get(out.proposal_id)!.value.status).toBe('REJECTED') await h.rt.close() }) it('a simulation alert escalates to a human even inside the envelope (one vote)', async () => { const h = await harness({ llm: null, config: { worstCaseLossBudgetYuan: '0', risk: { risk_aversion: '0.3', commitment_buffer_k: '0.9', min_block_mwh: '0.5', marginal_cost_yuan_per_mwh: '10000' } } }) h.rt.ctx.envelopes.register(envelope({ max_energy_mwh: '1500.0' })) const out = await submitBid(h) expect(out.lifecycle_status).toBe('suspended') expect(h.rt.ctx.caseDesk.inbox()[0]!.reasons[0]).toMatch(/loss budget/) await h.rt.close() }) }) describe('durability: a suspended approval survives process restart', () => { it('resumes on a fresh runtime over the same data dir', async () => { const h1 = await harness({ llm: null }) const out = await submitBid(h1) expect(out.lifecycle_status).toBe('suspended') await h1.rt.close() const h2 = await harness({ llm: null, dataDir: h1.dataDir, seed: false }) expect(h2.rt.ctx.caseDesk.inbox()[0]!.run_id).toBe(out.lifecycle_run_id) expect(h2.rt.ctx.caseDesk.proposals.get(out.proposal_id)!.value.status).toBe('PENDING_HUMAN') // The ledger is in-memory in phase 1: replay its monthly anchor so the fresh check sees the same version. h2.rt.ctx.ledger.append({ id: 'contract-2026-03', timescale: 'MONTHLY', period: '2026-03', kind: 'CONTRACT', energy_mwh: (311.04 * 31).toFixed(3), curve: null, source_ref: 'contract-2026-03-001', expected_version: 0 }) const state = await h2.rt.mastra.getWorkflow('proposalLifecycle').getWorkflowRunById(out.lifecycle_run_id) expect(state?.status).toBe('suspended') const res = await new TriggerService(h2.rt).decide(out.lifecycle_run_id, { decision: 'approve', approver: APPROVER }) expect(res.status).toBe('success') if (res.status === 'success') expect(res.result.outcome).toBe('RELEASED') expect(h2.rt.ctx.caseDesk.read(out.case_id).case.status).toBe('CLOSED_DONE') await h2.rt.close() }) }) describe('staleness and permits (I3/I4)', () => { it('approval goes STALE when the ledger moved between approval and fresh check', async () => { const h = await harness({ llm: null }) const out = await submitBid(h) h.rt.ctx.ledger.append({ id: 'late', timescale: 'MONTHLY', period: '2026-04', kind: 'CONTRACT', energy_mwh: '1.0', curve: null, source_ref: 'c', expected_version: 1 }) // Approval evidence is captured at resume time; the auto path captured version 1 in lineage. const res = await new TriggerService(h.rt).decide(out.lifecycle_run_id, { decision: 'approve', approver: APPROVER }) expect(res.status).toBe('success') if (res.status === 'success') { // The approval itself now references ledger v2 — consistent — so this passes fresh check; // the *auto-approved* path is the one bound to lineage v1. Exercise it below. expect(res.result.outcome).toBe('RELEASED') } await h.rt.close() const g = await harness({ llm: null }) g.rt.ctx.envelopes.register(envelope({ max_energy_mwh: '1500.0' })) // Move the ledger *during* the run: between fetch-context (lineage v1) and fresh check. g.rt.ctx.events.subscribe('Lifecycle.AUTO_APPROVED', () => { g.rt.ctx.ledger.append({ id: 'race', timescale: 'MONTHLY', period: '2026-04', kind: 'CONTRACT', energy_mwh: '1.0', curve: null, source_ref: 'c', expected_version: g.rt.ctx.ledger.read().version }) }) const out2 = await submitBid(g) expect(out2.lifecycle?.outcome).toBe('STALE') expect(out2.lifecycle?.reasons[0]).toMatch(/ledger version 1 at auto-approval, now 2/) expect(g.rt.ctx.gateway.receipts.list()).toHaveLength(0) await g.rt.close() }) it('permit expiry blocks a late release', async () => { const h = await harness({ llm: null }) const out = await submitBid(h) // Resume after the permit TTL would already have passed relative to approval: the gate issues // the permit at `now`, then the gateway is asked at now + TTL + 1s. const triggers = new TriggerService(h.rt) h.rt.ctx.events.subscribe('Lifecycle.AUTHORIZED', () => h.setNow('2026-03-14T06:00:01Z')) // TTL is 6h const res = await triggers.decide(out.lifecycle_run_id, { decision: 'approve', approver: APPROVER }) expect(res.status).toBe('success') if (res.status === 'success') { expect(res.result.outcome).toBe('STALE') expect(res.result.reasons[0]).toMatch(/permit expired/) } expect(h.rt.ctx.gateway.receipts.list()).toHaveLength(0) await h.rt.close() }) it('a revoked permit is refused by the gateway (kill switch)', async () => { const h = await harness({ llm: null }) const out = await submitBid(h) h.rt.ctx.events.subscribe('Lifecycle.AUTHORIZED', (evt) => { h.rt.ctx.authority.revoke((evt.payload as { permit_id: string }).permit_id, 'L1 kill switch', h.rt.ctx.clock()) }) const res = await new TriggerService(h.rt).decide(out.lifecycle_run_id, { decision: 'approve', approver: APPROVER }) if (res.status === 'success') expect(res.result.reasons[0]).toMatch(/revoked/) else throw new Error(res.status) await h.rt.close() }) }) describe('replay (I7) and authoritative store (I8)', () => { it('a released proposal re-checks identically from its snapshots and frozen pack version', async () => { const h = await harness({ llm: null }) h.rt.ctx.envelopes.register(envelope({ max_energy_mwh: '1500.0' })) const out = await submitBid(h) const stored: Proposal = h.rt.ctx.caseDesk.proposals.get(out.proposal_id)!.value expect(proposalDigest(stored)).toBe(out.proposal.digest) const originalCheck = h.rt.ctx.events.list({ event_type: 'RuleCheckCompleted' })[0]!.payload as { result_ref: string } const recorded = h.rt.ctx.snapshots.get(originalCheck.result_ref) as { ok: boolean; rules_evaluated: string[] } // Replay against a ledger reconstructed to the lineage version (the bid append moved it to 2). const replayLedger = h.rt.ctx.ledger const replay = h.rt.ctx.policy.check({ ...stored, status: 'DRAFT' }, { id: 'hubei-spot-bidding', version: stored.lineage.policy_pack_version }, { ledger: { version: stored.lineage.ledger_version, entries: [] }, ledgerService: replayLedger, snapshots: h.rt.ctx.snapshots, now: h.rt.ctx.clock() }) expect(replay.rules_evaluated).toEqual(recorded.rules_evaluated) expect(replay.ok).toBe(recorded.ok) // I8: the ledger (not any agent memory) is the source the proposal cites. expect(stored.lineage.ledger_version).toBe(1) expect(h.rt.ctx.ledger.read().entries.find((e) => e.kind === 'BID_SUBMITTED')?.source_ref).toBe(stored.digest) await h.rt.close() }) }) describe('trigger service', () => { it('scheduled entries fire once per local day and need no LLM', async () => { const h = await harness({ llm: null }) h.rt.ctx.envelopes.register(envelope({ max_energy_mwh: '1500.0' })) const t = new TriggerService(h.rt) expect(await t.tick('2026-03-13T21:00:00Z')).toEqual([]) // 05:00 Shanghai expect(await t.tick('2026-03-13T22:30:00Z')).toEqual(['situation-0600']) expect(await t.tick('2026-03-14T00:05:00Z')).toEqual(['bid-0800']) expect(await t.tick('2026-03-14T00:10:00Z')).toEqual([]) expect(h.rt.ctx.events.list({ event_type: 'WorkflowTriggered' }).map((e) => (e.payload as { workflow: string }).workflow)).toEqual(['dayAheadSituation', 'dayAheadBid']) await h.rt.close() }) it('manual requests route to a registered template; router refuses when the LLM is down', async () => { const llm = new ScriptedLlm({ 'router-agent': { workflow_id: 'day-ahead-situation', params: { market_date: MARKET_DATE }, confidence: 'HIGH' } }) const h = await harness({ llm }) const t = new TriggerService(h.rt) const r = await t.manual('明天的态势报告') expect(r.decision.workflow_id).toBe('day-ahead-situation') expect(r.run_id).not.toBeNull() await h.rt.close() const down = await harness({ llm: null }) await expect(new TriggerService(down.rt).manual('anything')).rejects.toThrow(/pick a template explicitly/) await down.rt.close() }) it('an EXTREME situation opens an abnormal-day case (event trigger)', async () => { const h = await harness({ llm: null }) h.skills.spread = 1.2 // p90/p50 = 2.2 ≥ extreme ratio 2.0 const t = new TriggerService(h.rt) t.startEventConsumers() await t.startTemplate('dayAheadSituation', { market_date: MARKET_DATE }, 'MANUAL', 'test') const cases = h.rt.ctx.caseDesk.cases.list().map((r) => r.value) expect(cases.some((c) => c.kind === 'ADHOC_ANALYSIS' && /Abnormal-day/.test(c.objective))).toBe(true) t.stop() await h.rt.close() }) })