M1: close gaps — CI, time-series/relational stores, ingestion skeleton, Python 3.11 pin

- .github/workflows/ci.yml: TS job (typecheck, re-export contracts, drift
  check, vitest) and Python job (3.11, regenerate pydantic models, drift
  check, pytest) — the dual-side contract test now actually runs in CI.
- packages/services: MemoryTimeSeriesStore (append-only daily-curve
  revisions), MemoryRepository + ResourceRegistry (schema-validated,
  optimistic versioning), IngestionPipeline (snapshot → quality gate →
  store or quarantine). In-memory reference semantics; persistent adapters
  arrive with M3 like the ledger.
- skills-py: .python-version + pyproject requires-python >=3.11; codegen
  script asserts interpreter version and runs the generator as a module.
- typecheck scripts per package and in root `check`; README status and dev
  setup; CLAUDE.md current-phase note.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01UoYoGYzHkFyv3ALenkRPhA
This commit is contained in:
Thomas Bayes 2026-09-01 21:58:49 -04:00
parent f681e134cc
commit e7196bc88a
16 changed files with 558 additions and 6 deletions

46
.github/workflows/ci.yml vendored Normal file
View File

@ -0,0 +1,46 @@
name: ci
on:
push:
branches: [master]
pull_request:
jobs:
# Domain schemas + deterministic services (TS side of the contract pipeline).
typescript:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: 24
cache: npm
- run: npm ci
- run: npm run typecheck
# ADR-0002: contracts/ is generated from packages/domain and committed.
# A dirty tree here means someone changed a zod schema without
# re-exporting, or hand-edited a generated file.
- run: npm run export:schemas && npm run make:fixtures
- run: git diff --exit-code -- contracts/
- run: npm test
# Generated pydantic models must accept/reject the same golden fixtures the
# TS side does (docs/11 §3.2 dual-side contract test).
python:
runs-on: ubuntu-latest
defaults:
run:
working-directory: skills-py
steps:
- uses: actions/checkout@v4
- uses: actions/setup-python@v5
with:
python-version-file: skills-py/.python-version
cache: pip
cache-dependency-path: skills-py/requirements.txt
- run: pip install -r requirements.txt
# Generated models are committed but never hand-edited: regenerate and
# require a clean tree.
- run: bash scripts/generate_models.sh
- run: git diff --exit-code -- vpp_contracts/
- run: python -m pytest -q

View File

@ -6,7 +6,10 @@ docs win — or write an ADR changing the doc first.
## Current phase
Docs complete, implementation not started. Build order is ROADMAP.md (M1→M5).
Docs complete; M1 implemented (schemas, contracts pipeline, ledger, snapshot,
time-series, relational and ingestion services, dual-side CI). Storage is
in-memory reference semantics — persistent adapters arrive with M3. Next is
M2. Build order is ROADMAP.md (M1→M5).
Do not start a milestone's work before its predecessor's acceptance criteria are
testable, and do not build phase-2 items (edge control links, federation,
interaction/load-control agents) unless explicitly asked.

View File

@ -6,9 +6,31 @@ a deterministic safety chain (rule check → simulation → envelope/human appro
execution permit) governs everything before any external effect. **The LLM never
computes numbers and never touches the second-level control loop.**
**Status: architecture/design phase.** No implementation code yet. The design is
complete and internally consistent (docs 00–13); implementation follows
[ROADMAP.md](ROADMAP.md), starting with M1.
**Status: M1 (data foundation & contracts) implemented.** Design docs 00–13 are
complete; implementation follows [ROADMAP.md](ROADMAP.md). Present today:
domain schemas (`packages/domain`), the TS↔Python contract pipeline
(`contracts/`, `skills-py/vpp_contracts`), and the ledger, snapshot,
time-series, relational and ingestion services (`packages/services`). M2
(Python skills + eval baseline) is next.
## Development
Requires Node 24 and Python 3.11 (`skills-py/.python-version`; the generated
pydantic models use `StrEnum` and PEP 604 unions).
```sh
npm ci
npm run check # typecheck, re-export contracts, run TS tests
cd skills-py
uv venv --python 3.11 .venv && uv pip install -r requirements.txt # or python3.11 -m venv
source .venv/bin/activate
bash scripts/generate_models.sh # regenerate pydantic models (committed, never hand-edited)
python -m pytest -q
```
CI (`.github/workflows/ci.yml`) runs both sides and fails if `contracts/` or
`skills-py/vpp_contracts` are not regenerated after a schema change.
## Orientation

View File

@ -6,9 +6,10 @@
"packages/*"
],
"scripts": {
"typecheck": "npm run typecheck --workspaces --if-present",
"test": "npm run test --workspaces --if-present",
"export:schemas": "npm run export:schemas -w @vpp/domain",
"make:fixtures": "npm run make:fixtures -w @vpp/domain",
"check": "npm run export:schemas && npm test"
"check": "npm run typecheck && npm run export:schemas && npm test"
}
}

View File

@ -7,6 +7,7 @@
".": "./src/index.ts"
},
"scripts": {
"typecheck": "tsc -p tsconfig.json",
"test": "vitest run",
"export:schemas": "tsx scripts/export-schemas.ts",
"make:fixtures": "tsx scripts/make-fixtures.ts"

View File

@ -7,6 +7,7 @@
".": "./src/index.ts"
},
"scripts": {
"typecheck": "tsc -p tsconfig.json",
"test": "vitest run"
},
"dependencies": {

View File

@ -1,3 +1,6 @@
export * from './snapshot.js'
export * from './ledger.js'
export * from './quality.js'
export * from './timeseries.js'
export * from './relational.js'
export * from './ingest.js'

View File

@ -0,0 +1,83 @@
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]
}
}

View File

@ -0,0 +1,78 @@
import type { z } from 'zod'
import { ResourceProfile as ResourceProfileSchema } from '@vpp/domain'
import type { ResourceProfile } from '@vpp/domain'
export class RepositoryConcurrencyError extends Error {
constructor(key: string, expected: number, actual: number) {
super(`repository row '${key}': expected version ${expected}, current is ${actual}`)
this.name = 'RepositoryConcurrencyError'
}
}
export interface Row<T> {
key: string
/** 0 = never written; increments on every accepted put. */
version: number
value: T
updated_at: string
}
/**
* Relational store port (docs/05 §3, 关系库): master data and business
* records — resource registry (资源台账), contracts, market results. Rows
* are schema-validated on write (schemas come from @vpp/domain only, CLAUDE.md
* rule 6) and versioned with optimistic concurrency, mirroring the ledger.
* Postgres adapter arrives with M3; the ledger keeps its own service because
* its constraint cascade is not generic row storage.
*/
export interface Repository<T> {
get(key: string): Row<T> | undefined
list(): Row<T>[]
/**
* Insert or replace. When expectedVersion is given it must equal the current
* row version (0 for a new row); omit it for unconditional master-data loads.
*/
put(value: T, expectedVersion?: number): Row<T>
}
export class MemoryRepository<T> implements Repository<T> {
private readonly rows = new Map<string, Row<T>>()
constructor(
private readonly schema: z.ZodType<T>,
private readonly keyOf: (value: T) => string,
private readonly clock: () => string = () => new Date().toISOString(),
) {}
get(key: string): Row<T> | undefined {
const row = this.rows.get(key)
return row && { ...row, value: structuredClone(row.value) }
}
list(): Row<T>[] {
return [...this.rows.keys()].sort().map((k) => this.get(k)!)
}
put(value: T, expectedVersion?: number): Row<T> {
const parsed = this.schema.parse(value)
const key = this.keyOf(parsed)
const current = this.rows.get(key)?.version ?? 0
if (expectedVersion !== undefined && expectedVersion !== current) {
throw new RepositoryConcurrencyError(key, expectedVersion, current)
}
const row: Row<T> = {
key,
version: current + 1,
value: structuredClone(parsed),
updated_at: this.clock(),
}
this.rows.set(key, row)
return this.get(key)!
}
}
/** 资源台账: ResourceProfile keyed by resource_id. */
export type ResourceRegistry = Repository<ResourceProfile>
export const createResourceRegistry = (clock?: () => string): ResourceRegistry =>
new MemoryRepository(ResourceProfileSchema, (r) => r.resource_id, clock)

View File

@ -0,0 +1,96 @@
import { Curve96 as Curve96Schema, MarketDate } from '@vpp/domain'
import type { Curve96 } from '@vpp/domain'
import { canonicalJson } from './snapshot.js'
/**
* One stored revision of a daily curve. Telemetry is corrected after the fact
* (late meter reads, re-transmissions), so a (series, date) has a revision
* history — never an in-place overwrite. Lineage that referenced revision N
* still resolves to revision N (docs/05 §3.2 快照与引用).
*/
export interface CurveRecord {
series_id: string
date: string
/** 1-based, monotonically increasing per (series_id, date). */
version: number
curve: Curve96
/** Where the curve came from: adapter name, file, edge node id. */
source: string
recorded_at: string
}
/**
* Time-series store port (docs/05 §3, 时序库). Curves are stored as whole
* market days (Curve96) — the granularity every forecast and bid consumes;
* sub-day streaming telemetry lives in the edge protocol domain, not here
* (docs/11 §4). TimescaleDB adapter arrives with M3; the in-memory
* implementation is the reference semantics.
*/
export interface TimeSeriesStore {
/** Appends a revision. Identical content to the latest revision is a no-op returning it. */
write(series_id: string, curve: Curve96, source: string): CurveRecord
latest(series_id: string, date: string): CurveRecord | undefined
revisions(series_id: string, date: string): CurveRecord[]
/** Latest revision per date, inclusive of both bounds, ascending by date. */
range(series_id: string, from: string, to: string): CurveRecord[]
seriesIds(): string[]
}
const keyOf = (series_id: string, date: string) => `${series_id} ${date}`
export class MemoryTimeSeriesStore implements TimeSeriesStore {
private readonly byKey = new Map<string, CurveRecord[]>()
constructor(private readonly clock: () => string = () => new Date().toISOString()) {}
write(series_id: string, curve: Curve96, source: string): CurveRecord {
if (!series_id) throw new Error('series_id must be non-empty')
Curve96Schema.parse(curve)
const key = keyOf(series_id, curve.date)
const history = this.byKey.get(key) ?? []
const last = history[history.length - 1]
if (last && canonicalJson(last.curve) === canonicalJson(curve)) return last
const record: CurveRecord = {
series_id,
date: curve.date,
version: history.length + 1,
curve: structuredClone(curve),
source,
recorded_at: this.clock(),
}
this.byKey.set(key, [...history, record])
return record
}
latest(series_id: string, date: string): CurveRecord | undefined {
const history = this.byKey.get(keyOf(series_id, date))
return history?.[history.length - 1]
}
revisions(series_id: string, date: string): CurveRecord[] {
return [...(this.byKey.get(keyOf(series_id, date)) ?? [])]
}
range(series_id: string, from: string, to: string): CurveRecord[] {
MarketDate.parse(from)
MarketDate.parse(to)
const out: CurveRecord[] = []
for (const history of this.byKey.values()) {
const last = history[history.length - 1]
if (!last || last.series_id !== series_id) continue
// ISO dates compare lexicographically.
if (last.date >= from && last.date <= to) out.push(last)
}
return out.sort((a, b) => (a.date < b.date ? -1 : a.date > b.date ? 1 : 0))
}
seriesIds(): string[] {
const ids = new Set<string>()
for (const history of this.byKey.values()) {
const first = history[0]
if (first) ids.add(first.series_id)
}
return [...ids].sort()
}
}

View File

@ -0,0 +1,68 @@
import { describe, expect, it } from 'vitest'
import type { Curve96 } from '@vpp/domain'
import { IngestionPipeline } from '../src/ingest.js'
import type { IngestSnapshot } from '../src/ingest.js'
import { MemorySnapshotStore } from '../src/snapshot.js'
import { MemoryTimeSeriesStore } from '../src/timeseries.js'
const curve = (fill: string): Curve96 => ({
interval_minutes: 15,
date: '2026-03-14',
values: Array.from({ length: 96 }, () => fill),
})
const clock = () => '2026-03-14T08:00:00Z'
const make = () => {
const timeseries = new MemoryTimeSeriesStore(clock)
const snapshots = new MemorySnapshotStore()
return { timeseries, snapshots, pipeline: new IngestionPipeline({ timeseries, snapshots, clock }) }
}
describe('ingestion pipeline skeleton', () => {
it('snapshots then stores a curve that passes the quality gate', () => {
const { timeseries, snapshots, pipeline } = make()
const result = pipeline.ingestCurve('load:agg-01', curve('12.5'), 'meter-file')
expect(result.accepted).toBe(true)
if (!result.accepted) throw new Error('unreachable')
expect(result.record.version).toBe(1)
expect(timeseries.latest('load:agg-01', '2026-03-14')).toEqual(result.record)
const snap = snapshots.get(result.snapshot_ref) as IngestSnapshot
expect(snap.series_id).toBe('load:agg-01')
expect(snap.source).toBe('meter-file')
expect(snap.quality.ok).toBe(true)
expect(snap.curve).toEqual(curve('12.5'))
expect(pipeline.quarantined()).toHaveLength(0)
})
it('quarantines a curve that fails the gate and keeps it out of the store', () => {
const { timeseries, snapshots, pipeline } = make()
const result = pipeline.ingestCurve('load:agg-01', curve('0'), 'edge')
expect(result.accepted).toBe(false)
if (result.accepted) throw new Error('unreachable')
expect(result.issues[0]).toMatch(/flat-zero/)
expect(timeseries.latest('load:agg-01', '2026-03-14')).toBeUndefined()
// Quarantined input is still evidence: tagged with its issues and snapshotted.
const q = pipeline.quarantined()
expect(q).toHaveLength(1)
expect(q[0]).toMatchObject({ series_id: 'load:agg-01', date: '2026-03-14', source: 'edge' })
const snap = snapshots.get(q[0]!.snapshot_ref) as IngestSnapshot
expect(snap.quality.ok).toBe(false)
})
it('accepts a custom gate so stricter checks can be plugged in later', () => {
const timeseries = new MemoryTimeSeriesStore(clock)
const snapshots = new MemorySnapshotStore()
const pipeline = new IngestionPipeline({
timeseries,
snapshots,
clock,
gate: () => ({ ok: false, issues: ['always rejected'] }),
})
const result = pipeline.ingestCurve('load:x', curve('1.0'), 'edge')
expect(result.accepted).toBe(false)
expect(timeseries.seriesIds()).toEqual([])
})
})

View File

@ -0,0 +1,60 @@
import { readFileSync } from 'node:fs'
import { fileURLToPath } from 'node:url'
import { describe, expect, it } from 'vitest'
import type { ResourceProfile } from '@vpp/domain'
import { RepositoryConcurrencyError, createResourceRegistry } from '../src/relational.js'
const fixture: ResourceProfile = JSON.parse(
readFileSync(
fileURLToPath(new URL('../../../contracts/fixtures/resource_profile/storage.json', import.meta.url)),
'utf8',
),
)
const clock = () => '2026-03-14T08:00:00Z'
describe('relational store: resource registry (资源台账)', () => {
it('puts and gets a schema-valid row with version 1', () => {
const reg = createResourceRegistry(clock)
const row = reg.put(fixture)
expect(row.key).toBe('res-storage-01')
expect(row.version).toBe(1)
expect(reg.get('res-storage-01')?.value).toEqual(fixture)
expect(reg.get('missing')).toBeUndefined()
})
it('rejects rows that fail the domain schema (no unvalidated master data)', () => {
const reg = createResourceRegistry(clock)
expect(() => reg.put({ ...fixture, rated_power_mw: 10 as unknown as string })).toThrow()
expect(() => reg.put({ ...fixture, type: 'WINDMILL' as never })).toThrow()
expect(reg.list()).toHaveLength(0)
})
it('enforces optimistic concurrency when expectedVersion is supplied', () => {
const reg = createResourceRegistry(clock)
reg.put(fixture, 0)
expect(() => reg.put({ ...fixture, reliability_score: '0.90' }, 0)).toThrow(
RepositoryConcurrencyError,
)
const updated = reg.put({ ...fixture, reliability_score: '0.90' }, 1)
expect(updated.version).toBe(2)
expect(updated.value.reliability_score).toBe('0.90')
})
it('allows unconditional replace for master-data loads', () => {
const reg = createResourceRegistry(clock)
reg.put(fixture)
const row = reg.put({ ...fixture, name: 'renamed' })
expect(row.version).toBe(2)
})
it('lists rows sorted by key and isolates stored values from mutation', () => {
const reg = createResourceRegistry(clock)
reg.put({ ...fixture, resource_id: 'res-b' })
reg.put({ ...fixture, resource_id: 'res-a' })
expect(reg.list().map((r) => r.key)).toEqual(['res-a', 'res-b'])
const row = reg.get('res-a')!
row.value.name = 'mutated'
expect(reg.get('res-a')!.value.name).toBe(fixture.name)
})
})

View File

@ -0,0 +1,78 @@
import { describe, expect, it } from 'vitest'
import type { Curve96 } from '@vpp/domain'
import { MemoryTimeSeriesStore } from '../src/timeseries.js'
const curve = (date: string, fill = '1.0'): Curve96 => ({
interval_minutes: 15,
date,
values: Array.from({ length: 96 }, () => fill),
})
let tick = 0
const clock = () => `2026-03-14T08:00:${String(tick++).padStart(2, '0')}Z`
describe('time-series store', () => {
it('writes a revision and reads it back as latest', () => {
const store = new MemoryTimeSeriesStore(clock)
const rec = store.write('load:agg-01', curve('2026-03-14'), 'meter-file')
expect(rec.version).toBe(1)
expect(store.latest('load:agg-01', '2026-03-14')).toEqual(rec)
expect(store.latest('load:agg-01', '2026-03-15')).toBeUndefined()
})
it('keeps revision history instead of overwriting (late corrections)', () => {
const store = new MemoryTimeSeriesStore(clock)
store.write('load:agg-01', curve('2026-03-14', '1.0'), 'edge')
const corrected = store.write('load:agg-01', curve('2026-03-14', '1.1'), 'settlement-meter')
expect(corrected.version).toBe(2)
const history = store.revisions('load:agg-01', '2026-03-14')
expect(history.map((r) => r.version)).toEqual([1, 2])
expect(history[0]!.curve.values[0]).toBe('1.0')
expect(store.latest('load:agg-01', '2026-03-14')!.curve.values[0]).toBe('1.1')
})
it('is idempotent for identical content', () => {
const store = new MemoryTimeSeriesStore(clock)
const a = store.write('pv:site-7', curve('2026-03-14'), 'edge')
const b = store.write('pv:site-7', curve('2026-03-14'), 'edge-retry')
expect(b).toBe(a)
expect(store.revisions('pv:site-7', '2026-03-14')).toHaveLength(1)
})
it('returns an inclusive date range with latest revision per date, sorted', () => {
const store = new MemoryTimeSeriesStore(clock)
store.write('load:agg-01', curve('2026-03-16'), 'edge')
store.write('load:agg-01', curve('2026-03-14', '1.0'), 'edge')
store.write('load:agg-01', curve('2026-03-14', '2.0'), 'edge')
store.write('load:agg-01', curve('2026-03-13'), 'edge')
store.write('load:agg-01', curve('2026-03-17'), 'edge')
store.write('pv:site-7', curve('2026-03-15'), 'edge')
const out = store.range('load:agg-01', '2026-03-14', '2026-03-16')
expect(out.map((r) => r.date)).toEqual(['2026-03-14', '2026-03-16'])
expect(out[0]!.version).toBe(2)
})
it('rejects malformed curves and empty series ids', () => {
const store = new MemoryTimeSeriesStore(clock)
expect(() => store.write('', curve('2026-03-14'), 'edge')).toThrow(/series_id/)
expect(() =>
store.write('load:x', { ...curve('2026-03-14'), values: ['1.0'] }, 'edge'),
).toThrow()
expect(() => store.range('load:x', '2026/03/14', '2026-03-15')).toThrow()
})
it('stored curves are isolated from caller mutation', () => {
const store = new MemoryTimeSeriesStore(clock)
const input = curve('2026-03-14')
store.write('load:x', input, 'edge')
input.values[0] = '999'
expect(store.latest('load:x', '2026-03-14')!.curve.values[0]).toBe('1.0')
})
it('lists series ids', () => {
const store = new MemoryTimeSeriesStore(clock)
store.write('pv:site-7', curve('2026-03-14'), 'edge')
store.write('load:agg-01', curve('2026-03-14'), 'edge')
expect(store.seriesIds()).toEqual(['load:agg-01', 'pv:site-7'])
})
})

View File

@ -0,0 +1 @@
3.11

9
skills-py/pyproject.toml Normal file
View File

@ -0,0 +1,9 @@
[project]
name = "vpp-skills"
version = "0.1.0"
description = "Python skill services and generated contract models (see ../ROADMAP.md M2)"
# Generated models use StrEnum and PEP 604 unions (generate_models.sh targets 3.11).
requires-python = ">=3.11"
[tool.pytest.ini_options]
testpaths = ["tests"]

View File

@ -4,7 +4,9 @@
set -euo pipefail
cd "$(dirname "$0")/.."
datamodel-codegen \
python3 -c 'import sys; assert sys.version_info >= (3, 11), f"need Python >=3.11, got {sys.version}"'
python3 -m datamodel_code_generator \
--input ../contracts/schema \
--input-file-type jsonschema \
--output vpp_contracts \