134 lines
5.7 KiB
TypeScript
134 lines
5.7 KiB
TypeScript
|
|
import { RouterDecision } from '@vpp/domain'
|
||
|
|
import type { EventEnvelope } from '@vpp/domain'
|
||
|
|
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 interface ScheduleEntry {
|
||
|
|
id: string
|
||
|
|
/** 'HH:MM' Asia/Shanghai on D-1 */
|
||
|
|
localTime: string
|
||
|
|
workflow: 'dayAheadSituation' | 'dayAheadBid'
|
||
|
|
}
|
||
|
|
|
||
|
|
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
|
||
|
|
]
|
||
|
|
|
||
|
|
export const nextMarketDate = (nowIso: string): string => {
|
||
|
|
const shanghai = new Date(new Date(nowIso).getTime() + 8 * 3_600_000)
|
||
|
|
shanghai.setUTCDate(shanghai.getUTCDate() + 1)
|
||
|
|
return shanghai.toISOString().slice(0, 10)
|
||
|
|
}
|
||
|
|
|
||
|
|
export class TriggerService {
|
||
|
|
private timer: NodeJS.Timeout | null = null
|
||
|
|
private readonly firedToday = new Set<string>()
|
||
|
|
private readonly unsubscribe: Array<() => void> = []
|
||
|
|
|
||
|
|
constructor(
|
||
|
|
private readonly rt: Runtime,
|
||
|
|
private readonly schedule: ScheduleEntry[] = 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<string[]> {
|
||
|
|
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)) {
|
||
|
|
this.firedToday.add(key)
|
||
|
|
await this.startTemplate(entry.workflow, { market_date: nextMarketDate(nowIso) }, 'SCHEDULED', entry.id)
|
||
|
|
fired.push(entry.id)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return fired
|
||
|
|
}
|
||
|
|
|
||
|
|
// ---- event -------------------------------------------------------------
|
||
|
|
|
||
|
|
/** Deterministic event rules. M3: a SituationPublished with EXTREME risk opens an abnormal-day case (docs/13 §1). */
|
||
|
|
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.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`,
|
||
|
|
})
|
||
|
|
}
|
||
|
|
}),
|
||
|
|
)
|
||
|
|
}
|
||
|
|
|
||
|
|
// ---- manual ------------------------------------------------------------
|
||
|
|
|
||
|
|
async route(message: string): Promise<RouterDecision> {
|
||
|
|
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<string> {
|
||
|
|
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 wf = this.rt.mastra.getWorkflow('proposalLifecycle')
|
||
|
|
const run = await wf.createRun({ runId })
|
||
|
|
return run.resume({ step: ENVELOPE_GATE_STEP_ID, resumeData: decision })
|
||
|
|
}
|
||
|
|
}
|