textmachine/platform/internal/ingest/events.go

141 lines
5.6 KiB
Go

// Package ingest is the platform side of the engine seam (D39.85): it supervises a tmctl process,
// reads its NDJSON event stream and hands each event to a sink that materializes it into the
// reporting database. It never opens the engine's SQLite and never parses human output.
//
// ⚠ The emitter does not exist yet — it is row 103 of the engine backlog. The vocabulary below is
// therefore the platform's PROPOSAL, derived from research/23 §7 (which call sites already carry
// the data) and from the API contract §6 (what a reader must be told). It is written as code
// rather than prose so the engine zone can answer it with a diff.
package ingest
import (
"encoding/json"
"fmt"
"strconv"
"strings"
"time"
)
// StreamVersion is the version this decoder implements. The rule is terraform's, ratified by
// D39.85: a MINOR bump adds fields and event types — unknown ones are ignored; a MAJOR bump is
// refused, because a stream whose meaning changed must not be materialized as if it had not.
const StreamVersion = "1.0"
// Type is the event name. Unknown values are legal under the minor rule and are dropped by the
// sink, not by the decoder.
type Type string
const (
// TypeHello is always the first line: the version handshake.
TypeHello Type = "hello"
// TypeProgress carries the per-phase counters. Per phase because a unit is done only once its
// edit resolved, so one end-to-end counter reads zero for the whole draft wave.
TypeProgress Type = "progress"
// TypeUnitDone is one shipped (or withheld) edit unit.
TypeUnitDone Type = "unit_done"
// TypeBankStop is the book-wide signing stop before the edit wave.
TypeBankStop Type = "bank_stop"
// TypeCeiling is the resumable halt on the spend ceiling. It carries NO figures.
TypeCeiling Type = "ceiling"
// TypeFinished is the last line of a clean stream.
TypeFinished Type = "finished"
)
// Envelope is one line of the stream.
type Envelope struct {
// Seq is per PROCESS and starts at 1. It is the second half of the ratified idempotency key
// (run_id, seq) — where run_id is Hello.EngineRunID, not the platform's run: a resumed run is
// a new process whose seq restarts, so keying on the platform run would drop its whole stream.
Seq int64 `json:"seq"`
Type Type `json:"type"`
Time time.Time `json:"time"`
Data json.RawMessage `json:"data"`
}
// Hello is the handshake payload.
type Hello struct {
StreamVersion string `json:"stream_version"`
// EngineRunID is the id the engine mints per invocation (today: the trace id, main.go:67-70).
// Once the engine accepts an external trace context (row 102) the platform can supply it, and
// the idempotency namespace becomes ours by construction.
EngineRunID string `json:"engine_run_id"`
BookID string `json:"book_id"`
// ChunkerVersion lets the platform notice that the chapter manifest it persisted was produced
// by a different chunker — the case that silently re-numbers chapters (row 100).
ChunkerVersion string `json:"chunker_version"`
}
// Counter is a phase counter, in units.
type Counter struct {
Done int `json:"done"`
Total int `json:"total"`
}
// Progress is the run's counters. ETASeconds mirrors what status already computes
// (pipeline/status.go:112) — the owner asked for a visible ETA (K-5).
type Progress struct {
Draft Counter `json:"draft"`
Edit Counter `json:"edit"`
ETASeconds int `json:"eta_seconds,omitempty"`
}
// UnitDone is one resolved edit unit. Chapter and Unit are the engine's own numbering; mapping
// them onto the platform's opaque ids is the materializer's job.
//
// Shipped and Flagged are BOTH needed and neither implies the other: a flagged unit legally ships
// with text (sanitizer cleanup, c-lite member drop), and the pair of them is exactly the contract's
// derivation of unit state — shipped → translated, flagged without text → withheld.
type UnitDone struct {
Chapter int `json:"chapter"`
Unit int `json:"unit"`
Wave string `json:"wave"` // draft | edit
Shipped bool `json:"shipped"`
Flagged bool `json:"flagged"`
// Reason is the ENGINE's flag reason (glossary_miss, sanitizer_stripped, …). It is stored and
// never projected: the product phrase is applied at read time from the contract's map, so a
// reason this platform has never heard of still gets a neutral phrase instead of a hole.
Reason string `json:"reason,omitempty"`
}
// BankStop is the signing stop. The full table travels as an artifact (engine backlog row 101),
// not through the stream: a thousand rows are not an event.
type BankStop struct {
TermsProposed int `json:"terms_proposed"`
}
// Ceiling is the ceiling halt. It carries the FACT and nothing else: money never reaches the
// platform's wire or its INFO logs (D39.84), and the stop is resumable, so it is not a failure.
type Ceiling struct {
Halted bool `json:"halted"`
}
// Finished is the terminal line. Outcome mirrors the engine's exit contract so a stream that ends
// cleanly needs no exit-code archaeology: clean | flagged | bank_stop | failed.
type Finished struct {
Outcome string `json:"outcome"`
}
// checkVersion applies the semver rule to a handshake.
func checkVersion(got string) error {
gotMajor, err := major(got)
if err != nil {
return err
}
wantMajor, err := major(StreamVersion)
if err != nil {
return err
}
if gotMajor != wantMajor {
return fmt.Errorf("%w: stream is %s, this build speaks %s", ErrUnsupportedVersion, got, StreamVersion)
}
return nil
}
func major(v string) (int, error) {
head, _, _ := strings.Cut(v, ".")
n, err := strconv.Atoi(head)
if err != nil {
return 0, fmt.Errorf("ingest: malformed stream version %q", v)
}
return n, nil
}