- packages/domain: AwardNotice, DISPATCH_PLAN proposal payload, ExecutionReport, MeteringRecord, potential-assessment and dispatch-optimization contracts, ReviewFinding (+ typed writebacks), SemanticMemoryEntry, EnvelopeChangeRequest, InsightCard — exported with fixtures on both sides. - skills-py: potential-assessment (certified × rolling fulfilment, evidence days) and dispatch-optimization (per-interval LP on HiGHS, shortfall reported) skills + routes + tests. - packages/services: dispatch rules in the policy pack (over-allocation, award anchor, lineage integrity for allocations); PowerBalanceSimulator; SimulationGateway (permit-only, idempotent, seeded execute → ExecutionReports); envelope deviation-streak suspension + apply(); ReviewService (attribution, reliability EWMA writeback, semantic memory, envelope recommendations as change requests); dispatch assembler; skill client methods. - packages/runtime: resource agent; award-decomposition, review and envelope-review workflows; lifecycle selects simulator/gateway by proposal type; trigger hooks for awards, execution reports, metering; decide() resumes either lifecycle or envelope-review runs; insight cards API. - Tests: docs/07 D-1 16:00 and D+1 end to end; reliability score 0.9 → 0.880 and the next assessment de-rates capacity; envelope suspension on a seeded 3-day streak; WIDEN request applied only by a human. 184 TS + 80 Python. - docs/open-questions: B10 (reliability/potential parameters). README and CLAUDE.md status → M4 done, M5 next. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UoYoGYzHkFyv3ALenkRPhA
250 lines
12 KiB
TypeScript
250 lines
12 KiB
TypeScript
import { createStep, createWorkflow } from '@mastra/core/workflows'
|
|
import { z } from 'zod'
|
|
import type { Approval, ForecastBundle, Proposal, ProposalStatus } from '@vpp/domain'
|
|
import { Approval as ApprovalSchema, ForecastBundle as ForecastBundleSchema, Proposal as ProposalSchema } from '@vpp/domain'
|
|
import { GatewayRejected, assertApproverIndependent, isStaleDenial } from '@vpp/services'
|
|
import type { RuntimeContext } from '../context.js'
|
|
|
|
/**
|
|
* The safety chain as ONE workflow for every proposal type (ADR-0005,
|
|
* docs/03 §2, docs/09 §2): rule check → simulation → envelope gate
|
|
* (suspend for humans) → fresh check + permit → release. Every transition is
|
|
* recorded through the case desk (event-sourced) so the run is replayable.
|
|
*/
|
|
export const LifecycleInput = z.object({ proposal: ProposalSchema, case_id: z.string().min(1) })
|
|
|
|
export const LifecycleOutcome = z.enum(['RELEASED', 'REJECTED', 'STALE', 'SIMULATION_ALERT'])
|
|
|
|
export const LifecycleOutput = z.object({
|
|
outcome: LifecycleOutcome,
|
|
proposal: ProposalSchema,
|
|
permit_id: z.string().nullable(),
|
|
receipt_id: z.string().nullable(),
|
|
reasons: z.array(z.string()),
|
|
})
|
|
|
|
const Carry = z.object({
|
|
proposal: ProposalSchema,
|
|
case_id: z.string(),
|
|
simulation_alert: z.boolean(),
|
|
simulation_alerts: z.array(z.string()),
|
|
approvals: z.array(ApprovalSchema),
|
|
envelope_match: z.unknown().nullable(),
|
|
done: LifecycleOutput.nullable(),
|
|
})
|
|
type Carry = z.infer<typeof Carry>
|
|
|
|
export const ResumeDecision = z.object({
|
|
decision: z.enum(['approve', 'reject']),
|
|
approver: z.object({ id: z.string().min(1), role: z.string().min(1) }),
|
|
comment: z.string().optional(),
|
|
})
|
|
export type ResumeDecision = z.infer<typeof ResumeDecision>
|
|
|
|
const SuspendInfo = z.object({
|
|
reason: z.string(),
|
|
proposal_id: z.string(),
|
|
required_level: z.string(),
|
|
reasons: z.array(z.string()),
|
|
})
|
|
|
|
function transition(ctx: RuntimeContext, caseId: string, p: Proposal, status: ProposalStatus, extra: Record<string, unknown> = {}): Proposal {
|
|
const next: Proposal = { ...p, status }
|
|
ctx.caseDesk.recordProposal(caseId, next)
|
|
if (Object.keys(extra).length) {
|
|
ctx.events.append({ event_type: `Lifecycle.${status}`, payload: { proposal_id: p.id, digest: p.digest, ...extra }, correlation_id: caseId, causation_id: p.id })
|
|
}
|
|
return next
|
|
}
|
|
|
|
function priceForecastFromLineage(ctx: RuntimeContext, p: Proposal): ForecastBundle | null {
|
|
for (const tc of [...p.lineage.tool_calls].reverse()) {
|
|
if (tc.tool !== 'price-forecast') continue
|
|
const parsed = ForecastBundleSchema.safeParse(ctx.snapshots.get(tc.outputs_ref))
|
|
if (parsed.success) return parsed.data
|
|
}
|
|
return null
|
|
}
|
|
|
|
const finish = (c: Carry, outcome: z.infer<typeof LifecycleOutcome>, reasons: string[], permit_id: string | null = null, receipt_id: string | null = null): Carry => ({
|
|
...c,
|
|
done: { outcome, proposal: c.proposal, permit_id, receipt_id, reasons },
|
|
})
|
|
|
|
/** One lifecycle workflow per runtime: steps close over that runtime's services. */
|
|
export function createProposalLifecycle(ctx: RuntimeContext) {
|
|
const ruleCheck = createStep({
|
|
id: 'rule-check',
|
|
inputSchema: LifecycleInput,
|
|
outputSchema: Carry,
|
|
execute: async ({ inputData }) => {
|
|
const pack = ctx.policy.pack(ctx.config.policyPackId, inputData.proposal.lineage.policy_pack_version)
|
|
let p = transition(ctx, inputData.case_id, inputData.proposal, 'RULE_CHECK')
|
|
const result = ctx.policy.check(p, { id: pack.id, version: pack.version }, { ledger: ctx.ledger.read(), ledgerService: ctx.ledger, snapshots: ctx.snapshots, now: ctx.clock() })
|
|
const ref = ctx.snapshots.put(result)
|
|
ctx.events.append({ event_type: 'RuleCheckCompleted', payload: { proposal_id: p.id, ok: result.ok, result_ref: ref, violations: result.violations }, correlation_id: inputData.case_id, causation_id: p.id })
|
|
const carry: Carry = { proposal: p, case_id: inputData.case_id, simulation_alert: false, simulation_alerts: [], approvals: [], envelope_match: null, done: null }
|
|
if (!result.ok) {
|
|
p = transition(ctx, inputData.case_id, p, 'REJECTED', { by: 'rule-check', violations: result.violations })
|
|
return finish({ ...carry, proposal: p }, 'REJECTED', result.violations.map((v) => `${v.rule_id}: ${v.message}`))
|
|
}
|
|
return carry
|
|
},
|
|
})
|
|
|
|
const simulation = createStep({
|
|
id: 'simulation',
|
|
inputSchema: Carry,
|
|
outputSchema: Carry,
|
|
execute: async ({ inputData }) => {
|
|
if (inputData.done) return inputData
|
|
const p = transition(ctx, inputData.case_id, inputData.proposal, 'SIMULATION')
|
|
const result = ctx.simulationFor(p.type).simulate(p, {
|
|
now: ctx.clock(),
|
|
priceForecast: priceForecastFromLineage(ctx, p),
|
|
marginalCostYuanPerMwh: ctx.config.risk.marginal_cost_yuan_per_mwh,
|
|
worstCaseLossBudgetYuan: ctx.config.worstCaseLossBudgetYuan,
|
|
})
|
|
const ref = ctx.snapshots.put(result)
|
|
ctx.events.append({ event_type: 'SimulationCompleted', payload: { proposal_id: p.id, alert: result.alert, alerts: result.alerts, metrics: result.metrics, result_ref: ref }, correlation_id: inputData.case_id, causation_id: p.id })
|
|
return { ...inputData, proposal: p, simulation_alert: result.alert, simulation_alerts: result.alerts }
|
|
},
|
|
})
|
|
|
|
const envelopeGate = createStep({
|
|
id: 'envelope-gate',
|
|
inputSchema: Carry,
|
|
outputSchema: Carry,
|
|
suspendSchema: SuspendInfo,
|
|
resumeSchema: ResumeDecision,
|
|
execute: async ({ inputData, resumeData, suspend }) => {
|
|
if (inputData.done) return inputData
|
|
const now = ctx.clock()
|
|
const p0 = inputData.proposal
|
|
|
|
if (resumeData) {
|
|
// —— Human decision path (I1/I2 enforced here; the AI has no route to APPROVED).
|
|
try {
|
|
assertApproverIndependent(p0, resumeData.approver, ctx.config.approverRoles)
|
|
} catch (e) {
|
|
ctx.events.append({ event_type: 'ApprovalRefused', payload: { proposal_id: p0.id, approver: resumeData.approver, reason: (e as Error).message }, correlation_id: inputData.case_id, causation_id: p0.id })
|
|
return await suspend({ reason: `approval refused: ${(e as Error).message}`, proposal_id: p0.id, required_level: 'L2', reasons: [(e as Error).message] })
|
|
}
|
|
const approval: Approval = {
|
|
id: ctx.newId('appr'),
|
|
proposal_digest: p0.digest,
|
|
decision: resumeData.decision === 'approve' ? 'APPROVE' : 'REJECT',
|
|
approver: resumeData.approver,
|
|
scope: { effect_type: p0.type, limits: {} },
|
|
validity: { from: now, to: new Date(new Date(now).getTime() + ctx.config.bidPermitTtlMs).toISOString() },
|
|
evidence_versions: { policy_pack_version: p0.lineage.policy_pack_version, ledger_version: ctx.ledger.read().version, data_snapshot_refs: p0.lineage.data_refs },
|
|
...(resumeData.comment ? { comment: resumeData.comment } : {}),
|
|
decided_at: now,
|
|
}
|
|
ctx.caseDesk.recordApproval(inputData.case_id, approval)
|
|
ctx.caseDesk.clearPending(p0.id)
|
|
if (approval.decision === 'REJECT') {
|
|
const p = transition(ctx, inputData.case_id, p0, 'REJECTED', { by: approval.approver.id, comment: resumeData.comment ?? null })
|
|
return finish({ ...inputData, proposal: p, approvals: [approval] }, 'REJECTED', [`rejected by ${approval.approver.id}`])
|
|
}
|
|
const p = transition(ctx, inputData.case_id, p0, 'APPROVED', { by: approval.approver.id })
|
|
return { ...inputData, proposal: p, approvals: [approval] }
|
|
}
|
|
|
|
// —— Automatic path.
|
|
let p = transition(ctx, inputData.case_id, p0, 'ENVELOPE_CHECK')
|
|
if (inputData.simulation_alert) {
|
|
// One vote: any simulation alert escalates to a human, no exceptions.
|
|
p = transition(ctx, inputData.case_id, p, 'PENDING_HUMAN', { because: 'simulation alert', alerts: inputData.simulation_alerts })
|
|
ctx.caseDesk.requestApproval({ proposal_id: p.id, proposal_digest: p.digest, case_id: inputData.case_id, run_id: '', required_level: 'L2', reasons: inputData.simulation_alerts, since: now })
|
|
return await suspend({ reason: 'simulation alert', proposal_id: p.id, required_level: 'L2', reasons: inputData.simulation_alerts })
|
|
}
|
|
const match = ctx.envelopes.match(p, { now, priceBaseline: priceForecastFromLineage(ctx, p)?.quantiles.p50 ?? null })
|
|
ctx.snapshots.put(match)
|
|
if (match.within) {
|
|
p = transition(ctx, inputData.case_id, { ...p, envelope_ref: match.envelope_id }, 'AUTO_APPROVED', { envelope_id: match.envelope_id })
|
|
return { ...inputData, proposal: p, envelope_match: match }
|
|
}
|
|
p = transition(ctx, inputData.case_id, p, 'PENDING_HUMAN', { because: 'outside envelope', reasons: match.reasons })
|
|
ctx.caseDesk.requestApproval({ proposal_id: p.id, proposal_digest: p.digest, case_id: inputData.case_id, run_id: '', required_level: match.required_level, reasons: match.reasons, since: now })
|
|
return await suspend({ reason: 'outside envelope', proposal_id: p.id, required_level: match.required_level, reasons: match.reasons })
|
|
},
|
|
})
|
|
|
|
const freshCheckAndPermit = createStep({
|
|
id: 'fresh-check',
|
|
inputSchema: Carry,
|
|
outputSchema: Carry,
|
|
execute: async ({ inputData }) => {
|
|
if (inputData.done) return inputData
|
|
const now = ctx.clock()
|
|
let p = transition(ctx, inputData.case_id, inputData.proposal, 'FRESH_CHECK')
|
|
const result = ctx.authority.authorize(
|
|
p,
|
|
{ approvals: inputData.approvals, envelopeMatch: (inputData.envelope_match as Parameters<typeof ctx.authority.authorize>[1]['envelopeMatch']) ?? null },
|
|
{ now, ledger: ctx.ledger, currentPolicyPackVersion: ctx.policy.current(ctx.config.policyPackId).version, offlineResources: [], permitTtlMs: ctx.config.bidPermitTtlMs },
|
|
)
|
|
if (isStaleDenial(result)) {
|
|
p = transition(ctx, inputData.case_id, p, 'STALE', { reasons: result.reasons })
|
|
return finish({ ...inputData, proposal: p }, 'STALE', result.reasons)
|
|
}
|
|
ctx.caseDesk.recordPermit(inputData.case_id, result)
|
|
p = transition(ctx, inputData.case_id, p, 'AUTHORIZED', { permit_id: result.id, expires_at: result.expires_at })
|
|
return { ...inputData, proposal: p, done: null, envelope_match: inputData.envelope_match, approvals: inputData.approvals, permit: result } as Carry & { permit: unknown }
|
|
},
|
|
})
|
|
|
|
const release = createStep({
|
|
id: 'release',
|
|
inputSchema: Carry.extend({ permit: z.unknown().optional() }),
|
|
outputSchema: LifecycleOutput,
|
|
execute: async ({ inputData }) => {
|
|
if (inputData.done) return inputData.done
|
|
const now = ctx.clock()
|
|
const permitId = (inputData.permit as { id: string } | undefined)?.id
|
|
const permit = permitId ? ctx.authority.permits.get(permitId)?.value : undefined
|
|
if (!permit) throw new Error('release reached without a permit — invariant I4 violated in workflow wiring')
|
|
try {
|
|
const receipt = ctx.gatewayFor(inputData.proposal.type).dispatch(inputData.proposal, permit, now)
|
|
const p = transition(ctx, inputData.case_id, inputData.proposal, 'RELEASED', { receipt_id: receipt.receipt_id, artifact_ref: receipt.artifact_ref })
|
|
ctx.events.append({ event_type: 'ExecutionReceipt', payload: receipt, correlation_id: inputData.case_id, causation_id: p.id })
|
|
if (p.type === 'BID') {
|
|
ctx.ledger.append({
|
|
id: ctx.newId('pos'),
|
|
timescale: 'DAY_AHEAD',
|
|
period: p.payload.market_date,
|
|
kind: 'BID_SUBMITTED',
|
|
energy_mwh: p.payload.quantities_mwh.values.reduce((s, v) => (Number(s) + Number(v)).toFixed(3), '0'),
|
|
curve: p.payload.quantities_mwh,
|
|
source_ref: p.digest,
|
|
expected_version: ctx.ledger.read().version,
|
|
})
|
|
}
|
|
ctx.caseDesk.close(inputData.case_id, receipt.receipt_id, true)
|
|
return { outcome: 'RELEASED' as const, proposal: p, permit_id: permit.id, receipt_id: receipt.receipt_id, reasons: [] }
|
|
} catch (e) {
|
|
if (e instanceof GatewayRejected) {
|
|
const p = transition(ctx, inputData.case_id, inputData.proposal, 'STALE', { gateway_reasons: e.reasons })
|
|
return { outcome: 'STALE' as const, proposal: p, permit_id: permit.id, receipt_id: null, reasons: e.reasons }
|
|
}
|
|
throw e
|
|
}
|
|
},
|
|
})
|
|
|
|
return createWorkflow({
|
|
id: 'proposal-lifecycle',
|
|
inputSchema: LifecycleInput,
|
|
outputSchema: LifecycleOutput,
|
|
})
|
|
.then(ruleCheck)
|
|
.then(simulation)
|
|
.then(envelopeGate)
|
|
.then(freshCheckAndPermit)
|
|
.then(release)
|
|
.commit()
|
|
|
|
}
|
|
|
|
export const ENVELOPE_GATE_STEP_ID = 'envelope-gate'
|