84 lines
2.7 KiB
TypeScript
84 lines
2.7 KiB
TypeScript
|
|
import type { Curve96 } from '@vpp/domain'
|
||
|
|
import { qualityGate } from './quality.js'
|
||
|
|
import type { QualityVerdict } from './quality.js'
|
||
|
|
import type { SnapshotStore } from './snapshot.js'
|
||
|
|
import type { CurveRecord, TimeSeriesStore } from './timeseries.js'
|
||
|
|
|
||
|
|
/**
|
||
|
|
* What a snapshot of an ingested curve contains. Lineage refs point at this,
|
||
|
|
* so a decision can be replayed from exactly the input that was accepted —
|
||
|
|
* or the input that was quarantined and therefore *not* used.
|
||
|
|
*/
|
||
|
|
export interface IngestSnapshot {
|
||
|
|
series_id: string
|
||
|
|
source: string
|
||
|
|
curve: Curve96
|
||
|
|
received_at: string
|
||
|
|
quality: QualityVerdict
|
||
|
|
}
|
||
|
|
|
||
|
|
export type IngestResult =
|
||
|
|
| { accepted: true; record: CurveRecord; snapshot_ref: string }
|
||
|
|
| { accepted: false; issues: string[]; snapshot_ref: string }
|
||
|
|
|
||
|
|
export interface QuarantineEntry {
|
||
|
|
series_id: string
|
||
|
|
date: string
|
||
|
|
source: string
|
||
|
|
issues: string[]
|
||
|
|
snapshot_ref: string
|
||
|
|
received_at: string
|
||
|
|
}
|
||
|
|
|
||
|
|
export interface IngestionPipelineDeps {
|
||
|
|
timeseries: TimeSeriesStore
|
||
|
|
snapshots: SnapshotStore
|
||
|
|
gate?: (curve: Curve96) => QualityVerdict
|
||
|
|
clock?: () => string
|
||
|
|
}
|
||
|
|
|
||
|
|
/**
|
||
|
|
* Ingestion pipeline skeleton (docs/05 §3.1 接入管道): every incoming curve
|
||
|
|
* is snapshotted first (evidence of what arrived), then passed through the
|
||
|
|
* data-quality gate. Passing curves are written to the time-series store;
|
||
|
|
* failing curves are tagged and quarantined (不合格数据打标隔离, §3.2) — they
|
||
|
|
* never reach a store a forecast skill reads from. Cleaning/alignment stages
|
||
|
|
* and the feature-store hand-off are M2 work; this fixes the shape.
|
||
|
|
*/
|
||
|
|
export class IngestionPipeline {
|
||
|
|
private readonly quarantine: QuarantineEntry[] = []
|
||
|
|
private readonly gate: (curve: Curve96) => QualityVerdict
|
||
|
|
private readonly clock: () => string
|
||
|
|
|
||
|
|
constructor(private readonly deps: IngestionPipelineDeps) {
|
||
|
|
this.gate = deps.gate ?? qualityGate
|
||
|
|
this.clock = deps.clock ?? (() => new Date().toISOString())
|
||
|
|
}
|
||
|
|
|
||
|
|
ingestCurve(series_id: string, curve: Curve96, source: string): IngestResult {
|
||
|
|
const received_at = this.clock()
|
||
|
|
const quality = this.gate(curve)
|
||
|
|
const snapshot: IngestSnapshot = { series_id, source, curve, received_at, quality }
|
||
|
|
const snapshot_ref = this.deps.snapshots.put(snapshot)
|
||
|
|
|
||
|
|
if (!quality.ok) {
|
||
|
|
this.quarantine.push({
|
||
|
|
series_id,
|
||
|
|
date: curve.date,
|
||
|
|
source,
|
||
|
|
issues: quality.issues,
|
||
|
|
snapshot_ref,
|
||
|
|
received_at,
|
||
|
|
})
|
||
|
|
return { accepted: false, issues: quality.issues, snapshot_ref }
|
||
|
|
}
|
||
|
|
|
||
|
|
const record = this.deps.timeseries.write(series_id, curve, source)
|
||
|
|
return { accepted: true, record, snapshot_ref }
|
||
|
|
}
|
||
|
|
|
||
|
|
quarantined(): QuarantineEntry[] {
|
||
|
|
return [...this.quarantine]
|
||
|
|
}
|
||
|
|
}
|