import { z } from 'zod' import { AwardNotice, Curve96, ExecutionReport, HumanBidRecord, MarketDate, MeteringRecord, RouterDecision } from '@vpp/domain' import type { EventEnvelope, KpiReport, ShadowDayRecord } from '@vpp/domain' import { shadowClearing } from '@vpp/services' import { ENVELOPE_REVIEW_STEP_ID } from './workflows/envelope-review.js' import { LlmUnavailable } from './llm.js' import type { Runtime } from './runtime.js' import { ENVELOPE_GATE_STEP_ID } from './workflows/proposal-lifecycle.js' import type { ResumeDecision } from './workflows/proposal-lifecycle.js' /** * Trigger service (docs/02 §1, docs/09 §1): three sources, one shape — * "instantiate template X with params Y". * scheduled: cron entries (D-1 06:00 situation, 08:00 bid) — no LLM in the path * event: bus consumers with deterministic rules * manual: the router agent turns free text into a RouterDecision * (a template id + params, never a new flow); on LLM outage the * operator picks the template explicitly. * Approval is a resume on the lifecycle run (docs/09 §2). */ export type ScheduledTemplate = 'dayAheadSituation' | 'dayAheadBid' | 'shadowClearing' | 'shadowClose' export interface ScheduleEntry { id: string /** 'HH:MM' Asia/Shanghai */ localTime: string workflow: ScheduledTemplate /** Market date relative to the local day the entry fires on: +1 = tomorrow (D-1 stages), -1 = yesterday (D+1 close). */ marketDateOffsetDays?: 1 | -1 } export const DEFAULT_SCHEDULE: ScheduleEntry[] = [ { id: 'situation-0600', localTime: '06:00', workflow: 'dayAheadSituation' }, // docs/07 D-1 06:00 { id: 'bid-0800', localTime: '08:00', workflow: 'dayAheadBid' }, // docs/07 D-1 08:00; window is OPEN-QUESTION A1 ] /** Shadow mode adds the simulated market answer (docs/07 D-1 16:00) and the D+1 close (execution → metering → review → scoring). */ export const SHADOW_SCHEDULE: ScheduleEntry[] = [ ...DEFAULT_SCHEDULE, { id: 'shadow-clearing-1600', localTime: '16:00', workflow: 'shadowClearing' }, // docs/07 D-1 16:00; publication time is OPEN-QUESTION A1 { id: 'shadow-close-0200', localTime: '02:00', workflow: 'shadowClose', marketDateOffsetDays: -1 }, // D+1 02:00 for D ] export const marketDateFor = (nowIso: string, offsetDays: number): string => { const shanghai = new Date(new Date(nowIso).getTime() + 8 * 3_600_000) shanghai.setUTCDate(shanghai.getUTCDate() + offsetDays) return shanghai.toISOString().slice(0, 10) } export const nextMarketDate = (nowIso: string): string => marketDateFor(nowIso, 1) /** Live-data feed (docs/06 §3): actual curves for a market date; each goes through the quality gate. */ export const MarketDataInput = z.object({ market_date: MarketDate, load_mw: z.array(z.string()).length(96).optional(), pv_mw: z.array(z.string()).length(96).optional(), price_yuan_per_mwh: z.array(z.string()).length(96).optional(), source: z.string().min(1).optional(), }) export type MarketDataInput = z.infer export class TriggerService { private timer: NodeJS.Timeout | null = null private readonly firedToday = new Set() private readonly unsubscribe: Array<() => void> = [] constructor( private readonly rt: Runtime, private readonly schedule: ScheduleEntry[] = rt.ctx.config.mode === 'SHADOW' ? SHADOW_SCHEDULE : DEFAULT_SCHEDULE, ) {} // ---- scheduled --------------------------------------------------------- /** Polls the clock; fires each entry once per local day. Deterministic, LLM-free. */ startScheduler(pollMs = 30_000): void { if (this.timer) return this.timer = setInterval(() => void this.tick(), pollMs) this.timer.unref() } stop(): void { if (this.timer) clearInterval(this.timer) this.timer = null for (const u of this.unsubscribe) u() } async tick(nowIso = this.rt.ctx.clock()): Promise { const local = new Date(new Date(nowIso).getTime() + 8 * 3_600_000) const hhmm = local.toISOString().slice(11, 16) const day = local.toISOString().slice(0, 10) const fired: string[] = [] for (const entry of this.schedule) { const key = `${day}:${entry.id}` if (hhmm >= entry.localTime && !this.firedToday.has(key)) { const marketDate = marketDateFor(nowIso, entry.marketDateOffsetDays ?? 1) try { await this.fire(entry, marketDate) this.firedToday.add(key) fired.push(entry.id) } catch (e) { // Late data (e.g. no clearing price yet) is normal in a live shadow run: record it and retry next tick. this.rt.ctx.events.append({ event_type: 'ScheduledTriggerFailed', payload: { entry: entry.id, workflow: entry.workflow, market_date: marketDate, error: (e as Error).message }, correlation_id: `schedule-${entry.id}` }) } } } return fired } private async fire(entry: ScheduleEntry, marketDate: string): Promise { if (entry.workflow === 'shadowClearing') await this.shadowClearing(marketDate) else if (entry.workflow === 'shadowClose') await this.shadowClose(marketDate) else await this.startTemplate(entry.workflow, { market_date: marketDate }, 'SCHEDULED', entry.id) } // ---- event ------------------------------------------------------------- /** * Deterministic event rules. A SituationPublished with EXTREME risk starts the * abnormal-day protocol (docs/13 §1): envelopes are off for that market date — * every proposal goes to a human — and a case is opened for the review. */ startEventConsumers(): void { this.unsubscribe.push( this.rt.ctx.events.subscribe('SituationPublished', (evt: EventEnvelope) => { const payload = evt.payload as { report: { market_date: string; risk_level: string; id: string } } if (payload.report.risk_level === 'EXTREME') { this.rt.ctx.breakers.declareAbnormalDay(payload.report.market_date, { id: 'runtime', role: 'system' }, `situation ${payload.report.id}: risk EXTREME (OPEN-QUESTION B7 triggers)`, this.rt.ctx.clock()) this.rt.ctx.caseDesk.open({ id: this.rt.ctx.newId('case'), kind: 'ADHOC_ANALYSIS', objective: `Abnormal-day protocol review for ${payload.report.market_date} (situation ${payload.report.id})`, owner: 'user-ops-lead', // OPEN-QUESTION B8 deadline: `${payload.report.market_date}T00:00:00Z`, }) } }), ) } /** Trading-platform adapter stand-in: an award arrives → ledger + decomposition (docs/07 D-1 16:00). */ async onAward(award: AwardNotice) { const a = AwardNotice.parse(award) this.rt.ctx.events.append({ event_type: 'AwardReceived', payload: a, correlation_id: a.id }) const wf = this.rt.mastra.getWorkflow('awardDecomposition') const run = await wf.createRun() const result = await run.start({ inputData: { award: a } }) if (result.status === 'failed') throw result.error return { run_id: run.runId, result: result.status === 'success' ? result.result : null, status: result.status } } /** Execution engine / edge feedback stand-in: execution reports are evidence for review and potential assessment. */ onExecutionReports(reports: ExecutionReport[]) { for (const r of reports) { const parsed = ExecutionReport.parse(r) this.rt.ctx.executionReports.put(parsed) this.rt.ctx.events.append({ event_type: 'ExecutionReport', payload: parsed, correlation_id: parsed.dispatch_proposal_digest }) } } /** Live-data feed: actual curves for a market date through the ingestion pipeline (quality gate + evidence snapshot). */ onMarketData(input: MarketDataInput) { const m = MarketDataInput.parse(input) const source = m.source ?? 'market-data-feed' const series: Array<[string, string[] | undefined]> = [['load:aggregate', m.load_mw], ['pv:aggregate', m.pv_mw], ['price:da', m.price_yuan_per_mwh]] const accepted: string[] = [] const quarantined: Array<{ series: string; issues: string[] }> = [] for (const [id, values] of series) { if (!values) continue const curve = Curve96.parse({ interval_minutes: 15, date: m.market_date, values }) const r = this.rt.ctx.ingestion.ingestCurve(id, curve, source) if (r.accepted) accepted.push(id) else quarantined.push({ series: id, issues: r.issues }) } this.rt.ctx.events.append({ event_type: 'MarketDataIngested', payload: { market_date: m.market_date, source, accepted, quarantined }, correlation_id: `market-data-${m.market_date}` }) return { market_date: m.market_date, accepted, quarantined } } /** The human trader's actual submission for a market date (shadow comparison baseline). */ onHumanBid(record: HumanBidRecord) { const h = HumanBidRecord.parse(record) this.rt.ctx.humanBids.put(h) this.rt.ctx.events.append({ event_type: 'HumanBidRecorded', payload: { id: h.id, market_date: h.market_date, source: h.source }, correlation_id: `shadow-${h.market_date}` }) return h } /** * Shadow market answer (docs/07 D-1 16:00 in shadow mode): the released * shadow bid is cleared against the actual day-ahead price and the resulting * AwardNotice goes down the normal award path. Idempotent per bid digest. */ async shadowClearing(marketDate: string) { const ctx = this.rt.ctx const bid = ctx.caseDesk.proposals .list() .map((r) => r.value) .filter((p) => p.type === 'BID' && p.payload.market_date === marketDate && p.status === 'RELEASED') .sort((a, b) => (a.created_at < b.created_at ? -1 : 1)) .at(-1) if (!bid) { ctx.events.append({ event_type: 'ShadowClearingSkipped', payload: { market_date: marketDate, reason: 'no RELEASED bid for the day' }, correlation_id: `shadow-${marketDate}` }) return null } const existing = ctx.awards.list().find((a) => a.value.bid_proposal_digest === bid.digest) if (existing) return { award: existing.value, decomposition: null } const price = ctx.timeseries.latest('price:da', marketDate)?.curve if (!price) throw new Error(`no clearing price for ${marketDate} yet — ingest market data first`) const award = shadowClearing(bid, price, { id: `award-shadow-${marketDate}`, now: ctx.clock() }) ctx.events.append({ event_type: 'ShadowCleared', payload: { market_date: marketDate, award_id: award.id, bid_proposal_digest: bid.digest, awarded_mwh: award.awarded_mwh.values.reduce((s, v) => s + Number(v), 0).toFixed(3) }, correlation_id: `shadow-${marketDate}` }) const decomposition = await this.onAward(award) return { award, decomposition } } /** D+1 close for a shadow day: simulated execution → metering → review → comparison → KPI (shadow-close workflow). */ async shadowClose(marketDate: string): Promise<{ run_id: string; record: ShadowDayRecord; kpi: KpiReport; review_ran: boolean; execution_simulated: boolean }> { const wf = this.rt.mastra.getWorkflow('shadowClose') const run = await wf.createRun() const result = await run.start({ inputData: { market_date: marketDate } }) if (result.status === 'failed') throw result.error if (result.status !== 'success') throw new Error(`shadow close for ${marketDate} ended ${result.status}`) return { run_id: run.runId, ...result.result } } /** Metering adapter stand-in: D+1 metering arrives → review workflow (docs/07 D+1). */ async onMetering(record: MeteringRecord) { const m = MeteringRecord.parse(record) this.rt.ctx.metering.put(m) this.rt.ctx.events.append({ event_type: 'MeteringArrived', payload: m, correlation_id: m.id }) const wf = this.rt.mastra.getWorkflow('review') const run = await wf.createRun() const result = await run.start({ inputData: { market_date: m.market_date } }) if (result.status === 'failed') throw result.error return { run_id: run.runId, result: result.status === 'success' ? result.result : null } } // ---- manual ------------------------------------------------------------ async route(message: string): Promise { const prompt = `Operator request: "${message}". Today (UTC) is ${this.rt.ctx.clock()}. Default market_date is ${nextMarketDate(this.rt.ctx.clock())}.` try { return await this.rt.ctx.llm.structured('router-agent', prompt, RouterDecision) } catch (e) { if (e instanceof LlmUnavailable) throw new LlmUnavailable('router needs an LLM; pick a template explicitly via startTemplate') throw e } } async manual(message: string): Promise<{ decision: RouterDecision; run_id: string | null }> { const decision = await this.route(message) if (decision.workflow_id === 'adhoc-analysis') return { decision, run_id: null } const run_id = await this.startTemplate(decision.workflow_id === 'day-ahead-bid' ? 'dayAheadBid' : 'dayAheadSituation', decision.params, 'MANUAL', 'router') return { decision, run_id } } async startTemplate(workflow: 'dayAheadSituation' | 'dayAheadBid', params: { market_date: string }, trigger: 'SCHEDULED' | 'EVENT' | 'MANUAL', source: string): Promise { const wf = this.rt.mastra.getWorkflow(workflow) const run = await wf.createRun() this.rt.ctx.events.append({ event_type: 'WorkflowTriggered', payload: { workflow, params, trigger, source, run_id: run.runId }, correlation_id: run.runId }) const result = await run.start({ inputData: params }) this.rt.ctx.events.append({ event_type: 'WorkflowFinished', payload: { workflow, run_id: run.runId, status: result.status }, correlation_id: run.runId }) if (result.status === 'failed') throw result.error return run.runId } // ---- approval = resume ------------------------------------------------- async decide(runId: string, decision: ResumeDecision) { const pending = this.rt.ctx.caseDesk.pending.list().find((p) => p.value.run_id === runId)?.value if (pending?.workflow === 'envelopeReview') { const run = await this.rt.mastra.getWorkflow('envelopeReview').createRun({ runId }) const res = await run.resume({ step: ENVELOPE_REVIEW_STEP_ID, resumeData: decision }) return res.status === 'success' ? { status: 'success' as const, result: { outcome: res.result.applied ? ('RELEASED' as const) : ('REJECTED' as const), ...res.result } } : res } const wf = this.rt.mastra.getWorkflow('proposalLifecycle') const run = await wf.createRun({ runId }) return run.resume({ step: ENVELOPE_GATE_STEP_ID, resumeData: decision }) } }