From 4db0af1a66516364f9607ba787f57854db6f33c4 Mon Sep 17 00:00:00 2001 From: Adil Date: Sun, 20 Sep 2026 02:18:34 +0500 Subject: [PATCH] feat(thrum): views.rs as protocol source with generated mirrors Canonical per-chi body structs live in thrum-core/src/views.rs; codegen parses them with syn and emits protocol.{ts,py,go} so the three clients cannot drift. thrum-core's build fails if any chi lacks a body view. Golden fixture (15 tones) locks the wire; go/python parse tests and Crockford-52 rid() parity (hum_identity vectors) run in CI via wire.yml. --- .github/workflows/wire.yml | 51 ++ Cargo.lock | 1 + codegen/Cargo.toml | 1 + codegen/src/lib.rs | 194 ++++-- codegen/src/main.rs | 52 +- codegen/src/paths.rs | 53 +- codegen/src/protocol.rs | 620 ++++++++++++++++++++ thrum-clients/go/thrum/helpers.go | 58 +- thrum-clients/go/thrum/helpers_test.go | 51 ++ thrum-clients/go/thrum/protocol.go | 262 +++++++++ thrum-clients/go/thrum/protocol_test.go | 126 ++++ thrum-clients/python/tests/test_helpers.py | 37 ++ thrum-clients/python/tests/test_protocol.py | 80 +++ thrum-clients/python/thrum/helpers.py | 37 +- thrum-clients/python/thrum/protocol.py | 521 ++++++++++++++++ thrum-clients/ts/helpers.ts | 43 +- thrum-clients/ts/index.ts | 24 +- thrum-clients/ts/package.json | 4 +- thrum-clients/ts/protocol.ts | 263 +++++++++ thrum-core/build.rs | 54 +- thrum-core/src/lib.rs | 2 + thrum-core/src/views.rs | 419 +++++++++++++ thrum-core/tests/fixtures/golden.ndjson | 15 + thrum-core/tests/golden.rs | 278 +++++++++ 24 files changed, 3123 insertions(+), 123 deletions(-) create mode 100644 .github/workflows/wire.yml create mode 100644 codegen/src/protocol.rs create mode 100644 thrum-clients/go/thrum/helpers_test.go create mode 100644 thrum-clients/go/thrum/protocol.go create mode 100644 thrum-clients/go/thrum/protocol_test.go create mode 100644 thrum-clients/python/tests/test_helpers.py create mode 100644 thrum-clients/python/tests/test_protocol.py create mode 100644 thrum-clients/python/thrum/protocol.py create mode 100644 thrum-clients/ts/protocol.ts create mode 100644 thrum-core/src/views.rs create mode 100644 thrum-core/tests/fixtures/golden.ndjson create mode 100644 thrum-core/tests/golden.rs diff --git a/.github/workflows/wire.yml b/.github/workflows/wire.yml new file mode 100644 index 00000000..32d29388 --- /dev/null +++ b/.github/workflows/wire.yml @@ -0,0 +1,51 @@ +name: Wire protocol + +# views.rs is the single source of truth for every chi body and the +# envelope. These checks keep the generated client mirrors (TS/Python/Go) +# and the golden fixture byte-locked to it. + +on: + push: + pull_request: + +permissions: + contents: read + +jobs: + rust: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@stable + - name: Every chi has a body view (fails the build otherwise) + run: cargo build -p thrum-core + - name: Generated mirrors are byte-identical to views.rs + run: cargo run -p codegen -- --check + - name: Codegen unit tests + run: cargo test -p codegen --lib + - name: Golden fixture round-trips through views + run: cargo test -p thrum-core + + go: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-go@v5 + with: + go-version: stable + - name: Golden-wire parse + rid parity + working-directory: thrum-clients/go + run: | + go build ./... + go vet ./... + go test ./... + + python: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-python@v5 + with: + python-version: "3.x" + - name: Golden-wire parse + rid parity (stdlib only) + run: python -m unittest discover -s thrum-clients/python/tests diff --git a/Cargo.lock b/Cargo.lock index 6b13e024..0c98d6cb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -517,6 +517,7 @@ version = "0.32.0" dependencies = [ "anyhow", "regex", + "syn", ] [[package]] diff --git a/codegen/Cargo.toml b/codegen/Cargo.toml index 713f8073..fc26c96b 100644 --- a/codegen/Cargo.toml +++ b/codegen/Cargo.toml @@ -15,3 +15,4 @@ path = "src/main.rs" [dependencies] anyhow = { workspace = true } regex = { workspace = true } +syn = { version = "2", features = ["full", "parsing"] } diff --git a/codegen/src/lib.rs b/codegen/src/lib.rs index 0179a4d8..64de877d 100644 --- a/codegen/src/lib.rs +++ b/codegen/src/lib.rs @@ -13,6 +13,7 @@ use anyhow::{anyhow, bail, Context}; use regex::Regex; pub mod paths; +pub mod protocol; pub use anyhow::Result; @@ -32,7 +33,6 @@ fn chi_at_generated_header(comment: &str) -> String { ) } - #[derive(Debug, Clone)] pub struct Variant { /// PascalCase name as it appears in the Rust enum. @@ -53,10 +53,10 @@ pub struct ChiSpec { /// Parse a chi.rs + lib.rs pair into a [`ChiSpec`]. pub fn parse(chi_rs: &Path, lib_rs: &Path) -> Result { - let chi_src = fs::read_to_string(chi_rs) - .with_context(|| format!("read {}", chi_rs.display()))?; - let lib_src = fs::read_to_string(lib_rs) - .with_context(|| format!("read {}", lib_rs.display()))?; + let chi_src = + fs::read_to_string(chi_rs).with_context(|| format!("read {}", chi_rs.display()))?; + let lib_src = + fs::read_to_string(lib_rs).with_context(|| format!("read {}", lib_rs.display()))?; Ok(ChiSpec { version: extract_version(&lib_src)?, chi: extract_enum(&chi_src, "Chi")?, @@ -133,7 +133,9 @@ pub fn emit_go_helpers(output: &Path) -> Result<()> { fn extract_version(lib: &str) -> Result { let re = Regex::new(r#"pub\s+const\s+THRUM_VERSION\s*:\s*&\s*str\s*=\s*"([^"]+)""#)?; - let caps = re.captures(lib).ok_or_else(|| anyhow!("THRUM_VERSION not found in lib.rs"))?; + let caps = re + .captures(lib) + .ok_or_else(|| anyhow!("THRUM_VERSION not found in lib.rs"))?; Ok(caps[1].to_string()) } @@ -142,7 +144,9 @@ fn extract_version(lib: &str) -> Result { /// doc comments that immediately precede it. fn extract_enum(src: &str, name: &str) -> Result> { let opener = Regex::new(&format!(r"pub\s+enum\s+{}\s*\{{", regex::escape(name)))?; - let m = opener.find(src).ok_or_else(|| anyhow!("`pub enum {}` not found", name))?; + let m = opener + .find(src) + .ok_or_else(|| anyhow!("`pub enum {}` not found", name))?; let body_start = m.end(); let bytes = src.as_bytes(); let mut depth = 1usize; @@ -152,13 +156,17 @@ fn extract_enum(src: &str, name: &str) -> Result> { b'{' => depth += 1, b'}' => { depth -= 1; - if depth == 0 { break; } + if depth == 0 { + break; + } } _ => {} } i += 1; } - if depth != 0 { bail!("unterminated `{}` block", name); } + if depth != 0 { + bail!("unterminated `{}` block", name); + } let body = &src[body_start..i]; let mut out = Vec::new(); @@ -191,7 +199,9 @@ fn extract_enum(src: &str, name: &str) -> Result> { pending_doc.clear(); out.push(Variant { pascal, wire, doc }); } - if out.is_empty() { bail!("`{}` block parsed but yielded no variants", name); } + if out.is_empty() { + bail!("`{}` block parsed but yielded no variants", name); + } Ok(out) } @@ -199,7 +209,9 @@ fn pascal_to_kebab(s: &str) -> String { let mut out = String::with_capacity(s.len() + 4); for (i, c) in s.chars().enumerate() { if c.is_ascii_uppercase() { - if i > 0 { out.push('-'); } + if i > 0 { + out.push('-'); + } out.push(c.to_ascii_lowercase()); } else { out.push(c); @@ -245,7 +257,7 @@ fn render_helpers() -> String { // helper, add it in Rust first and extend codegen's render_helpers.\n\n", paths::HELPERS_SOURCE_REF, ); - const BODY: &str = r#"import { createHash } from "crypto"; + const BODY: &str = r#"import { createHash, randomBytes } from "crypto"; /** * Deterministic identity for a (nest, session) pair. @@ -263,10 +275,45 @@ export function sigil(sid: string, nest: string): string { .slice(0, 12); } -/** Monotonic request id — base36 timestamp + counter. */ -let __ridCounter = 0; +const CROCKFORD = "0123456789ABCDEFGHJKMNPQRSTVWXYZ"; + +/** + * Encode 32 bytes (48-bit ms prefix + 208 random bits, big-endian) as a + * 52-char Crockford base32 id. Mirrors hum_identity's encode bit-for-bit. + */ +export function humIdEncode(buf: Uint8Array): string { + let v = 0n; + for (const b of buf) v = (v << 8n) | BigInt(b); + v <<= 4n; + let out = ""; + for (let i = 0; i < 52; i++) { + out += CROCKFORD[Number((v >> BigInt(255 - 5 * i)) & 0x1fn)]; + } + return out; +} + +/** + * Correlation id in canonical HumId form — 52-char Crockford base32, + * ts-prefixed, matching hum_identity::HumId::mint(). + */ export function rid(): string { - return `${Date.now().toString(36)}-${(__ridCounter++).toString(36)}`; + const buf = new Uint8Array(32); + let ms = BigInt(Date.now()); + for (let i = 0; i < 6; i++) { + buf[5 - i] = Number(ms & 0xffn); + ms >>= 8n; + } + buf.set(randomBytes(26), 6); + return humIdEncode(buf); +} + +/** True iff `value` is a 52-char Crockford base32 id. */ +export function isValidRid(value: unknown): value is string { + return ( + typeof value === "string" && + value.length === 52 && + [...value].every((c) => CROCKFORD.includes(c)) + ); } /** Absolute ms timestamp ms in the future. */ @@ -305,7 +352,10 @@ fn render_ts(spec: &ChiSpec) -> String { let mut s = String::new(); s.push_str(&chi_at_generated_header("//")); - s.push_str(&format!("export const THRUM_VERSION = \"{}\" as const;\n\n", spec.version)); + s.push_str(&format!( + "export const THRUM_VERSION = \"{}\" as const;\n\n", + spec.version + )); s.push_str("// Every wire-known chi value. Adding a new variant bumps the\n"); s.push_str("// protocol minor; renaming/removing bumps major.\n"); @@ -314,7 +364,11 @@ fn render_ts(spec: &ChiSpec) -> String { if !v.doc.is_empty() { s.push_str(&format!(" /** {} */\n", v.doc)); } - s.push_str(&format!(" {}: \"{}\",\n", pascal_to_camel(&v.pascal), v.wire)); + s.push_str(&format!( + " {}: \"{}\",\n", + pascal_to_camel(&v.pascal), + v.wire + )); } s.push_str("} as const;\n"); s.push_str("export type ChiKind = typeof Chi[keyof typeof Chi];\n\n"); @@ -327,7 +381,11 @@ fn render_ts(spec: &ChiSpec) -> String { if !v.doc.is_empty() { s.push_str(&format!(" /** {} */\n", v.doc)); } - s.push_str(&format!(" {}: \"{}\",\n", pascal_to_camel(&v.pascal), v.wire)); + s.push_str(&format!( + " {}: \"{}\",\n", + pascal_to_camel(&v.pascal), + v.wire + )); } s.push_str("} as const;\n"); s.push_str("export type PulseKindT = typeof PulseKind[keyof typeof PulseKind];\n\n"); @@ -427,6 +485,7 @@ fn render_py_helpers() -> String { import hashlib import json import os +import secrets import threading import time from typing import Any, Mapping @@ -448,28 +507,30 @@ def now_ms() -> int: return int(time.time() * 1000) -_RID_LOCK = threading.Lock() -_RID_COUNTER = 0 +_CROCKFORD = "0123456789ABCDEFGHJKMNPQRSTVWXYZ" -def _base36(n: int) -> str: - if n == 0: - return "0" - alphabet = "0123456789abcdefghijklmnopqrstuvwxyz" - out = [] - while n > 0: - out.append(alphabet[n % 36]) - n //= 36 - return "".join(reversed(out)) +def _hum_id_from_bytes(buf: bytes) -> str: + """Encode 32 bytes (48-bit ms prefix + 208 random bits) as a 52-char + Crockford base32 id. Mirrors hum_identity's encode bit-for-bit.""" + v = int.from_bytes(buf, "big") << 4 + return "".join(_CROCKFORD[(v >> (255 - 5 * i)) & 0x1F] for i in range(52)) def rid() -> str: - """Monotonic correlation id: '{base36-ms}-{base36-counter}'.""" - global _RID_COUNTER - with _RID_LOCK: - n = _RID_COUNTER - _RID_COUNTER += 1 - return f"{_base36(now_ms())}-{_base36(n)}" + """Correlation id in canonical HumId form: 52-char Crockford base32, + ts-prefixed, matching hum_identity::HumId::mint().""" + ts = int(time.time() * 1000) + buf = ts.to_bytes(6, "big") + secrets.token_bytes(26) + return _hum_id_from_bytes(buf) + + +def is_valid_rid(value: str) -> bool: + """True iff `value` is a 52-char Crockford base32 id.""" + return ( + len(value) == 52 + and all(c in _CROCKFORD for c in value) + ) def dusk_in(ms: int) -> int: @@ -572,7 +633,10 @@ fn render_go(spec: &ChiSpec) -> String { if !v.doc.is_empty() { s.push_str(&format!(" // {}\n", v.doc)); } - s.push_str(&format!(" PulseKind{} PulseKind = \"{}\"\n", v.pascal, v.wire)); + s.push_str(&format!( + " PulseKind{} PulseKind = \"{}\"\n", + v.pascal, v.wire + )); } s.push_str(")\n"); @@ -590,15 +654,15 @@ fn render_go_helpers() -> String { const BODY: &str = r#"package thrum import ( + "crypto/rand" "crypto/sha256" "encoding/hex" "encoding/json" - "fmt" + "math/big" "os" "path/filepath" - "strconv" + "strings" "sync" - "sync/atomic" "time" ) @@ -613,19 +677,53 @@ func Sigil(sid, nest string) string { // NowMs returns wall-clock milliseconds since the Unix epoch. func NowMs() int64 { return time.Now().UnixMilli() } -var ridCounter uint64 - -func base36(n uint64) string { - if n == 0 { - return "0" +const crockford = "0123456789ABCDEFGHJKMNPQRSTVWXYZ" + +// HumIdEncode encodes 32 bytes (48-bit ms prefix + 208 random bits, +// big-endian) as a 52-char Crockford base32 id. Mirrors hum_identity's +// encode bit-for-bit. +func HumIdEncode(buf [32]byte) string { + var v big.Int + v.SetBytes(buf[:]) + v.Lsh(&v, 4) + var out strings.Builder + out.Grow(52) + for i := 0; i < 52; i++ { + var idx big.Int + idx.Rsh(&v, uint(255-5*i)) + idx.And(&idx, big.NewInt(0x1f)) + out.WriteByte(crockford[idx.Int64()]) } - return strconv.FormatUint(n, 36) + return out.String() } -// Rid mints a monotonic correlation id: "{base36-ms}-{base36-counter}". +// Rid mints a correlation id in canonical HumId form — 52-char +// Crockford base32, ts-prefixed, matching hum_identity::HumId::mint(). func Rid() string { - n := atomic.AddUint64(&ridCounter, 1) - 1 - return fmt.Sprintf("%s-%s", base36(uint64(NowMs())), base36(n)) + var buf [32]byte + ms := uint64(NowMs()) + for i := 0; i < 6; i++ { + buf[5-i] = byte(ms >> (8 * uint(i))) + } + if _, err := rand.Read(buf[6:]); err != nil { + // 26 bytes from the CSPRNG is practically infallible; fall + // back to a time-seeded value rather than mint nothing. + buf[6] = byte(ms >> 40) + } + return HumIdEncode(buf) +} + +// IsValidRid reports whether value is a 52-char Crockford base32 id. +func IsValidRid(value string) bool { + if len(value) != 52 { + return false + } + for i := 0; i < len(value); i++ { + if !strings.ContainsRune(crockford, rune(value[i])) { + return false + } + } + return true } // DuskIn returns the absolute ms timestamp at which a tone with this diff --git a/codegen/src/main.rs b/codegen/src/main.rs index 1ff4ea82..f3b12411 100644 --- a/codegen/src/main.rs +++ b/codegen/src/main.rs @@ -38,39 +38,57 @@ fn run() -> Result<()> { positional }; - let spec = codegen::parse(&paths::chi_rs(), &paths::lib_rs()) - .context("parse chi spec")?; + let spec = codegen::parse(&paths::chi_rs(), &paths::lib_rs()).context("parse chi spec")?; + let proto = codegen::protocol::protocol_spec( + &paths::views_rs(), + &paths::envelope_rs(), + &paths::lib_rs(), + ) + .context("parse protocol spec")?; for t in &targets { - run_target(t, &spec, check)?; + run_target(t, &spec, &proto, check)?; } Ok(()) } -fn run_target(target: &str, spec: &ChiSpec, check: bool) -> Result<()> { - let (chi_out, helpers_out, emit_chi, emit_helpers): ( +fn run_target( + target: &str, + spec: &ChiSpec, + proto: &codegen::protocol::ProtocolSpec, + check: bool, +) -> Result<()> { + let (chi_out, helpers_out, protocol_out, emit_chi, emit_helpers, emit_protocol): ( std::path::PathBuf, std::path::PathBuf, + std::path::PathBuf, + Box Result<()>>, Box Result<()>>, Box Result<()>>, ) = match target { "ts" => ( paths::ts_chi(), paths::ts_helpers(), + paths::ts_protocol(), Box::new(|p| codegen::emit_ts(spec, p)), Box::new(codegen::emit_helpers), + Box::new(|p| codegen::protocol::emit_protocol_ts(proto, p)), ), "python" | "py" => ( paths::py_chi(), paths::py_helpers(), + paths::py_protocol(), Box::new(|p| codegen::emit_py(spec, p)), Box::new(codegen::emit_py_helpers), + Box::new(|p| codegen::protocol::emit_protocol_py(proto, p)), ), "go" => ( paths::go_chi(), paths::go_helpers(), + paths::go_protocol(), Box::new(|p| codegen::emit_go(spec, p)), Box::new(codegen::emit_go_helpers), + Box::new(|p| codegen::protocol::emit_protocol_go(proto, p)), ), other => anyhow::bail!("unknown target {other}; valid: ts, python, go"), }; @@ -78,20 +96,34 @@ fn run_target(target: &str, spec: &ChiSpec, check: bool) -> Result<()> { if check { check_against(&chi_out, &emit_chi)?; check_against(&helpers_out, &emit_helpers)?; - eprintln!("codegen: {} + {} up to date", chi_out.display(), helpers_out.display()); + check_against(&protocol_out, &emit_protocol)?; + eprintln!( + "codegen: {} + {} + {} up to date", + chi_out.display(), + helpers_out.display(), + protocol_out.display() + ); } else { emit_chi(&chi_out)?; emit_helpers(&helpers_out)?; + emit_protocol(&protocol_out)?; eprintln!( - "codegen {target}: {} ({} chi, {} pulse) -> {} + {}", - spec.version, spec.chi.len(), spec.pulse.len(), - chi_out.display(), helpers_out.display(), + "codegen {target}: {} ({} chi, {} pulse) -> {} + {} + {}", + spec.version, + spec.chi.len(), + spec.pulse.len(), + chi_out.display(), + helpers_out.display(), + protocol_out.display(), ); } Ok(()) } -fn check_against(output: &std::path::Path, emit: &dyn Fn(&std::path::Path) -> Result<()>) -> Result<()> { +fn check_against( + output: &std::path::Path, + emit: &dyn Fn(&std::path::Path) -> Result<()>, +) -> Result<()> { let tmp = tempfile_path(output); emit(&tmp)?; let generated = std::fs::read(&tmp).context("read tmp")?; diff --git a/codegen/src/paths.rs b/codegen/src/paths.rs index 81f21130..649340c6 100644 --- a/codegen/src/paths.rs +++ b/codegen/src/paths.rs @@ -21,6 +21,9 @@ use std::path::PathBuf; pub const CHI_RS_REL: &str = "thrum-core/src/chi.rs"; pub const LIB_RS_REL: &str = "thrum-core/src/lib.rs"; +/// Wire-view sources — parsed into typed protocol mirrors for each client. +pub const VIEWS_RS_REL: &str = "thrum-core/src/views.rs"; +pub const ENVELOPE_RS_REL: &str = "thrum-core/src/envelope.rs"; /// Compact ref for the two runtime-helper sources, used in the /// "Runtime helpers that mirror …" header emitted into every helpers /// file. They live next to chi.rs / lib.rs in `thrum-core/src/`. @@ -28,10 +31,13 @@ pub const HELPERS_SOURCE_REF: &str = "thrum-core/src/{prims,wane}.rs"; pub const TS_CHI_REL: &str = "thrum-clients/ts/chi.ts"; pub const TS_HELPERS_REL: &str = "thrum-clients/ts/helpers.ts"; +pub const TS_PROTOCOL_REL: &str = "thrum-clients/ts/protocol.ts"; pub const PY_CHI_REL: &str = "thrum-clients/python/thrum/chi.py"; pub const PY_HELPERS_REL: &str = "thrum-clients/python/thrum/helpers.py"; +pub const PY_PROTOCOL_REL: &str = "thrum-clients/python/thrum/protocol.py"; pub const GO_CHI_REL: &str = "thrum-clients/go/thrum/chi.go"; pub const GO_HELPERS_REL: &str = "thrum-clients/go/thrum/helpers.go"; +pub const GO_PROTOCOL_REL: &str = "thrum-clients/go/thrum/protocol.go"; // ── absolute resolution ─────────────────────────────────────────────────── @@ -39,11 +45,42 @@ fn repo_root() -> PathBuf { PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("..") } -pub fn chi_rs() -> PathBuf { repo_root().join(CHI_RS_REL) } -pub fn lib_rs() -> PathBuf { repo_root().join(LIB_RS_REL) } -pub fn ts_chi() -> PathBuf { repo_root().join(TS_CHI_REL) } -pub fn ts_helpers() -> PathBuf { repo_root().join(TS_HELPERS_REL) } -pub fn py_chi() -> PathBuf { repo_root().join(PY_CHI_REL) } -pub fn py_helpers() -> PathBuf { repo_root().join(PY_HELPERS_REL) } -pub fn go_chi() -> PathBuf { repo_root().join(GO_CHI_REL) } -pub fn go_helpers() -> PathBuf { repo_root().join(GO_HELPERS_REL) } +pub fn chi_rs() -> PathBuf { + repo_root().join(CHI_RS_REL) +} +pub fn lib_rs() -> PathBuf { + repo_root().join(LIB_RS_REL) +} +pub fn views_rs() -> PathBuf { + repo_root().join(VIEWS_RS_REL) +} +pub fn envelope_rs() -> PathBuf { + repo_root().join(ENVELOPE_RS_REL) +} +pub fn ts_chi() -> PathBuf { + repo_root().join(TS_CHI_REL) +} +pub fn ts_helpers() -> PathBuf { + repo_root().join(TS_HELPERS_REL) +} +pub fn ts_protocol() -> PathBuf { + repo_root().join(TS_PROTOCOL_REL) +} +pub fn py_chi() -> PathBuf { + repo_root().join(PY_CHI_REL) +} +pub fn py_helpers() -> PathBuf { + repo_root().join(PY_HELPERS_REL) +} +pub fn py_protocol() -> PathBuf { + repo_root().join(PY_PROTOCOL_REL) +} +pub fn go_chi() -> PathBuf { + repo_root().join(GO_CHI_REL) +} +pub fn go_helpers() -> PathBuf { + repo_root().join(GO_HELPERS_REL) +} +pub fn go_protocol() -> PathBuf { + repo_root().join(GO_PROTOCOL_REL) +} diff --git a/codegen/src/protocol.rs b/codegen/src/protocol.rs new file mode 100644 index 00000000..a9137a64 --- /dev/null +++ b/codegen/src/protocol.rs @@ -0,0 +1,620 @@ +//! Protocol-view codegen: parse `thrum-core/src/views.rs` + `envelope.rs` +//! with `syn` and emit typed per-chi bodies for TS / Python / Go. +//! +//! Same constraint as the rest of codegen: no dependency on `thrum-core`. +//! We parse source text, so this crate stays an ordinary build dependency +//! of thrum-core (no cycle). + +use std::fs; +use std::path::Path; + +use anyhow::{Context, Result, anyhow}; +use syn::{Field, Fields, Item, ItemStruct, Meta, Type}; + +use crate::paths; + +/// Simplified wire type — everything `views.rs` / `envelope.rs` can +/// express that the client languages need to render. +#[derive(Debug, Clone)] +pub enum Repr { + Str, + Int { + unsigned: bool, + }, + Bool, + /// `serde_json::Value` — opaque JSON. + Json, + /// `BTreeMap`. + Map(Box), + /// `Vec`. + Arr(Box), + /// Reference to another generated type (`DroneLoad`, …). + Ref(String), +} + +#[derive(Debug, Clone)] +pub struct FieldSpec { + /// Rust (snake_case) name, also the Python field name. + pub rust: String, + /// Wire key — from `#[serde(rename = "...")]` else the rust name. + pub json: String, + pub repr: Repr, + pub optional: bool, +} + +#[derive(Debug, Clone)] +pub struct TypeSpec { + /// `HelloBody`, `Envelope`, … + pub name: String, + pub fields: Vec, +} + +#[derive(Debug, Clone)] +pub struct ProtocolSpec { + pub version: String, + pub types: Vec, +} + +/// Parse a views.rs + envelope.rs pair (plus lib.rs for the version). +pub fn protocol_spec(views_rs: &Path, envelope_rs: &Path, lib_rs: &Path) -> Result { + let lib_src = + fs::read_to_string(lib_rs).with_context(|| format!("read {}", lib_rs.display()))?; + let version = extract_version(&lib_src)? + .ok_or_else(|| anyhow!("THRUM_VERSION not found in {}", lib_rs.display()))?; + + let mut types = Vec::new(); + for path in [views_rs, envelope_rs] { + let src = fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?; + for st in structs_in(&src)? { + types.push(st); + } + } + Ok(ProtocolSpec { version, types }) +} + +fn structs_in(src: &str) -> Result> { + let file = syn::parse_file(src).map_err(|e| anyhow::anyhow!("syn parse: {e}"))?; + let mut out = Vec::new(); + for item in file.items { + let Item::Struct(s) = item else { continue }; + let name = s.ident.to_string(); + // Only wire-facing types: the envelope, per-chi bodies, and the + // nested body-support structs (DroneLoad). Test-only types never + // appear here — unit-test structs are filtered by these names. + if name != "Envelope" && name != "DroneLoad" && !name.ends_with("Body") { + continue; + } + out.push(type_spec(s)); + } + Ok(out) +} + +fn type_spec(s: ItemStruct) -> TypeSpec { + let mut fields = Vec::new(); + if let Fields::Named(named) = &s.fields { + for f in &named.named { + let Some(field) = field_spec(f) else { continue }; + fields.push(field); + } + } + TypeSpec { + name: s.ident.to_string(), + fields, + } +} + +fn field_spec(f: &Field) -> Option { + let rust = f.ident.as_ref()?.to_string(); + let json = serde_rename(f).unwrap_or_else(|| rust.clone()); + let (optional, repr) = simplify(&f.ty); + Some(FieldSpec { + rust, + json, + repr, + optional, + }) +} + +/// `Option` → (true, T); anything else → (false, itself). +fn simplify(ty: &Type) -> (bool, Repr) { + if let Type::Path(p) = ty { + // Check the leading segment name directly: syn's `is_ident` rejects + // paths with generic arguments, so `Option` would slip by. + if let Some(seg) = p.path.segments.first() { + if seg.ident == "Option" { + if let syn::PathArguments::AngleBracketed(a) = &seg.arguments { + if let Some(syn::GenericArgument::Type(inner)) = a.args.first() { + return (true, simplify(inner).1); + } + } + } + } + } + (false, repr_of(ty)) +} + +fn repr_of(ty: &Type) -> Repr { + let Type::Path(p) = ty else { + return Repr::Json; // opaque fallback (bare trait object, etc.) + }; + let seg = &p.path.segments[0]; + match seg.ident.to_string().as_str() { + "String" => Repr::Str, + "i64" | "i32" | "i16" => Repr::Int { unsigned: false }, + "u64" | "u32" | "u16" => Repr::Int { unsigned: true }, + "bool" => Repr::Bool, + "Value" => Repr::Json, + "Vec" => { + let inner = angle_ty(&seg.arguments, 0).unwrap_or_else(|| Repr::Json); + Repr::Arr(Box::new(inner)) + } + "BTreeMap" => { + // The value type is the SECOND generic arg (first is the key). + let inner = angle_ty(&seg.arguments, 1).unwrap_or_else(|| Repr::Json); + Repr::Map(Box::new(inner)) + } + other => { + // Multi-segment paths (serde_json::Value) land on the last + // segment; otherwise treat as a reference to another type. + if p.path.segments.len() > 1 { + let last = &p.path.segments[p.path.segments.len() - 1]; + return repr_of(&Type::Path(syn::TypePath { + qself: None, + path: last.clone().into(), + })); + } + Repr::Ref(other.to_string()) + } + } +} + +fn angle_ty(args: &syn::PathArguments, nth: usize) -> Option { + if let syn::PathArguments::AngleBracketed(a) = args { + if let Some(syn::GenericArgument::Type(t)) = a.args.iter().nth(nth) { + return Some(simplify(t).1); + } + } + None +} + +/// Verify every Chi variant has a matching `{Pascal}Body` view, and every +/// body view maps back to a real chi. Called on every thrum-core build so +/// a new chi without a view (or an orphan view) fails fast instead of +/// silently generating a partial client. +pub fn check_coverage(chis: &[crate::Variant], proto: &ProtocolSpec) -> Result<()> { + for v in chis { + let expect = format!("{}Body", v.pascal); + if !proto.types.iter().any(|t| t.name == expect) { + anyhow::bail!( + "chi `{}` ({}) has no body view `{expect}` in {}", + v.wire, + v.pascal, + paths::VIEWS_RS_REL + ); + } + } + let mut wire_names: Vec = proto + .types + .iter() + .filter(|t| names_body_view(&t.name)) + .map(|t| body_chi(&t.name)) + .collect(); + let mut chi_wires: Vec = chis.iter().map(|v| v.wire.clone()).collect(); + wire_names.sort(); + chi_wires.sort(); + if wire_names != chi_wires { + anyhow::bail!( + "view <-> chi mismatch in {}: views {wire_names:?} vs chis {chi_wires:?}", + paths::VIEWS_RS_REL + ); + } + Ok(()) +} + +/// First `#[serde(rename = "...")]` value on the field, if any. +/// +/// Parsed by regex over the raw attribute tokens rather than +/// `parse_nested_meta`: that API's handling of unconsumed `key = value` +/// pairs changed across syn 2.x patch releases, and the token text of +/// `rename = "…"` is verbatim in both, so regex is version-proof. +fn serde_rename(f: &Field) -> Option { + let re = regex::Regex::new(r#"rename\s*=\s*"([^"]*)""#).expect("static regex"); + for attr in &f.attrs { + if !attr.path().is_ident("serde") { + continue; + } + let Meta::List(list) = &attr.meta else { + continue; + }; + if let Some(cap) = re.captures(&list.tokens.to_string()) { + return Some(cap[1].to_string()); + } + } + None +} + +/// THRUM_VERSION from lib.rs (same regex as the chi registry parser). +fn extract_version(lib: &str) -> Result> { + let re = regex::Regex::new(r#"pub\s+const\s+THRUM_VERSION\s*:\s*&\s*str\s*=\s*"([^"]+)""#)?; + Ok(re.captures(lib).map(|c| c[1].to_string())) +} + +// ── emitters ─────────────────────────────────────────────────────────── + +fn regen_header(comment: &str, refs: &[&str]) -> String { + let refs = refs.join(", "); + format!( + "{c} @generated by `cargo run -p codegen` from {refs} — DO NOT EDIT.\n\ + {c}\n\ + {c} Canonical wire views: every chi body is a Rust struct in\n\ + {c} thrum-core/src/views.rs, the envelope in envelope.rs. These\n\ + {c} client mirrors come out of the same source — no hand-sync.\n\ + {c} Manual regen: `cargo run -p codegen`.\n\n", + c = comment, + ) +} + +/// Emit `protocol.ts` to `output`. +pub fn emit_protocol_ts(spec: &ProtocolSpec, output: &Path) -> Result<()> { + let mut s = String::new(); + s.push_str(®en_header( + "//", + &[paths::VIEWS_RS_REL, paths::ENVELOPE_RS_REL], + )); + s.push_str("export type { ChiKind } from \"./chi.ts\";\n"); + s.push_str("import type { Chi } from \"./chi.ts\";\n\n"); + + for t in &spec.types { + s.push_str(&format!("export interface {} {{\n", t.name)); + for f in &t.fields { + let opt = if f.optional { "?" } else { "" }; + s.push_str(&format!(" {}{}: {};\n", f.json, opt, ts_repr(&f.repr))); + } + s.push_str("}\n\n"); + } + + // chi → view type map for discriminated tones. + s.push_str("// Body view per chi wire value.\n"); + s.push_str("export interface ToneViews {\n"); + for t in &spec.types { + if !names_body_view(&t.name) { + continue; + } + let wire = body_chi(&t.name); + s.push_str(&format!(" \"{}\": {};\n", wire, t.name)); + } + s.push_str("}\n\n"); + s.push_str("export type ToneView = Envelope & ToneViews[C];\n"); + write_out(output, &s) +} + +/// True for structs that are per-chi bodies (skips Envelope and the +/// nested support structs in chi → view maps). +fn names_body_view(name: &str) -> bool { + name.ends_with("Body") +} + +/// `HelloBody` → `hello`; `KadFindNodeRespBody` → `kad-find-node-resp`. +fn body_chi(name: &str) -> String { + let pascal = name.strip_suffix("Body").unwrap_or(name); + let mut out = String::with_capacity(pascal.len() + 4); + for (i, c) in pascal.chars().enumerate() { + if c.is_ascii_uppercase() { + if i > 0 { + out.push('-'); + } + out.push(c.to_ascii_lowercase()); + } else { + out.push(c); + } + } + out +} + +fn ts_repr(r: &Repr) -> String { + match r { + Repr::Str => "string".into(), + Repr::Int { .. } => "number".into(), + Repr::Bool => "boolean".into(), + Repr::Json => "unknown".into(), + Repr::Map(inner) => format!("Record", ts_repr(inner)), + Repr::Arr(inner) => format!("{}[]", ts_repr(inner)), + Repr::Ref(name) => name.clone(), + } +} + +/// Emit `protocol.py` to `output`. +pub fn emit_protocol_py(spec: &ProtocolSpec, output: &Path) -> Result<()> { + let mut s = String::new(); + s.push_str(®en_header( + "#", + &[paths::VIEWS_RS_REL, paths::ENVELOPE_RS_REL], + )); + s.push_str("from __future__ import annotations\n\n"); + s.push_str("from dataclasses import dataclass\n"); + s.push_str("from typing import Any, ClassVar, Optional\n\n"); + + for t in &spec.types { + s.push_str(&format!("@dataclass\nclass {}:\n", t.name)); + if t.fields.is_empty() { + s.push_str(" pass\n\n"); + continue; + } + // Dataclasses require defaulted (optional) fields after required + // ones — Rust order is wire order and mixes them freely. + let ordered: Vec<&FieldSpec> = t + .fields + .iter() + .filter(|f| !f.optional) + .chain(t.fields.iter().filter(|f| f.optional)) + .collect(); + let mut pairs = Vec::new(); + for f in ordered.iter() { + let python = py_field_name(&f.rust); + let default = if f.optional { " = None" } else { "" }; + s.push_str(&format!( + " {}: {}{}\n", + python, + py_repr(&f.repr, f.optional), + default + )); + pairs.push(format!(" {:?}: {:?},", python, f.json)); + } + s.push('\n'); + s.push_str(" # rust field name -> wire key (single source in views.rs)\n"); + s.push_str(" _wire: ClassVar[dict[str, str]] = {\n"); + for p in &pairs { + s.push_str(&format!("{p}\n")); + } + s.push_str(" }\n\n"); + } + + s.push_str( + "# Serialize a view to the flat wire body (envelope keys not included).\n\ + # Absent optionals, like skip_serializing_if on the Rust side, are omitted.\n\ + def to_wire(view: Any) -> dict[str, Any]:\n\ + \x20\x20\x20\x20payload: dict[str, Any] = {}\n\ + \x20\x20\x20\x20for rust_name, wire_key in type(view)._wire.items():\n\ + \x20\x20\x20\x20\x20\x20\x20\x20value = getattr(view, rust_name)\n\ + \x20\x20\x20\x20\x20\x20\x20\x20if value is not None:\n\ + \x20\x20\x20\x20\x20\x20\x20\x20\x20\x20\x20\x20payload[wire_key] = value\n\ + \x20\x20\x20\x20return payload\n\n", + ); + s.push_str( + "# Parse a flat wire body into the view for `cls`, ignoring envelope\n\ + # keys (already stripped) and any unknown keys. Missing required\n\ + # fields raise TypeError, surfacing spec drift loudly.\n\ + def from_wire(cls: type, payload: dict[str, Any]) -> Any:\n\ + \x20\x20\x20\x20wire_to_rust = {j: f for f, j in cls._wire.items()}\n\ + \x20\x20\x20\x20kwargs = {\n\ + \x20\x20\x20\x20\x20\x20\x20\x20wire_to_rust[k]: v for k, v in payload.items() if k in wire_to_rust\n\ + \x20\x20\x20\x20}\n\ + \x20\x20\x20\x20return cls(**kwargs)\n\n", + ); + + s.push_str("# chi wire value -> body view class.\n"); + s.push_str("TONE_VIEWS: dict[str, type] = {\n"); + for t in &spec.types { + if !names_body_view(&t.name) { + continue; + } + s.push_str(&format!(" \"{}\": {},\n", body_chi(&t.name), t.name)); + } + s.push_str("}\n"); + write_out(output, &s) +} + +fn py_repr(r: &Repr, optional: bool) -> String { + let base = match r { + Repr::Str => "str", + Repr::Int { .. } => "int", + Repr::Bool => "bool", + Repr::Json => "Any", + Repr::Map(inner) => { + return map_opt(optional, &format!("dict[str, {}]", py_repr(inner, false))); + } + Repr::Arr(inner) => return map_opt(optional, &format!("list[{}]", py_repr(inner, false))), + Repr::Ref(name) => name, + }; + if optional { + format!("Optional[{base}]") + } else { + base.to_string() + } +} + +/// Rust snake field name → valid Python identifier. `from` collides with +/// the `from` keyword (gossip-publish body). +fn py_field_name(rust: &str) -> String { + match rust { + "from" => "from_".to_string(), + other => other.to_string(), + } +} + +/// Optional-list/map annotations still need `Optional` even though `None` +/// is a legal value: `list[str] = None` fails type checkers. +fn map_opt(optional: bool, inner: &str) -> String { + if optional { + format!("Optional[{inner}]") + } else { + inner.to_string() + } +} + +/// Emit `protocol.go` to `output`. +pub fn emit_protocol_go(spec: &ProtocolSpec, output: &Path) -> Result<()> { + let mut s = String::new(); + s.push_str(®en_header( + "//", + &[paths::VIEWS_RS_REL, paths::ENVELOPE_RS_REL], + )); + s.push_str("package thrum\n\n"); + s.push_str("import \"encoding/json\"\n\n"); + + for t in &spec.types { + s.push_str(&format!("type {} struct {{\n", t.name)); + if t.fields.is_empty() { + s.push_str("}\n\n"); + continue; + } + for f in &t.fields { + let tag = if f.optional { + format!("json:\"{},omitempty\"", f.json) + } else { + format!("json:\"{}\"", f.json) + }; + s.push_str(&format!( + "\t{} {} `{}`\n", + snake_to_pascal(&f.rust), + go_repr(&f.repr, f.optional), + tag + )); + } + s.push_str("}\n\n"); + } + + s.push_str("// ToneViews maps each chi wire value to its body view type.\n"); + s.push_str("var ToneViews = map[Chi]any{\n"); + for t in &spec.types { + if !names_body_view(&t.name) { + continue; + } + s.push_str(&format!("\tChi{}: {}{{}},\n", pascal_of(&t.name), t.name)); + } + s.push_str("}\n"); + write_out(output, &s) +} + +fn pascal_of(name: &str) -> String { + // HelloBody → HelloBody (already Pascal); strip impossible suffix. + name.strip_suffix("Body") + .map(|p| p.to_string()) + .unwrap_or_else(|| name.to_string()) +} + +fn go_repr(r: &Repr, optional: bool) -> String { + let base = match r { + Repr::Str => "string".to_string(), + Repr::Int { unsigned } => { + if *unsigned { + "uint64".into() + } else { + "int64".into() + } + } + Repr::Bool => "bool".to_string(), + Repr::Json => "json.RawMessage".to_string(), + Repr::Map(inner) => format!("map[string]{}", go_repr(inner, false)), + Repr::Arr(inner) => format!("[]{}", go_repr(inner, false)), + Repr::Ref(name) => name.clone(), + }; + if optional && matches!(r, Repr::Str | Repr::Int { .. } | Repr::Bool | Repr::Ref(_)) { + format!("*{base}") + } else { + base + } +} + +fn snake_to_pascal(s: &str) -> String { + let mut out = String::with_capacity(s.len()); + for part in s.split('_') { + let mut cs = part.chars(); + if let Some(f) = cs.next() { + out.push(f.to_ascii_uppercase()); + } + out.push_str(cs.as_str()); + } + out +} + +fn write_out(output: &Path, s: &str) -> Result<()> { + if let Some(parent) = output.parent() { + fs::create_dir_all(parent).ok(); + } + fs::write(output, s).with_context(|| format!("write {}", output.display()))?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn kebab_of_body() { + assert_eq!(body_chi("HelloBody"), "hello"); + assert_eq!(body_chi("KadFindNodeRespBody"), "kad-find-node-resp"); + assert_eq!(body_chi("ToolCallBody"), "tool-call"); + } + + #[test] + fn pascal_field_names() { + assert_eq!(snake_to_pascal("proto_version"), "ProtoVersion"); + assert_eq!(snake_to_pascal("humd_id"), "HumdId"); + assert_eq!(snake_to_pascal("block_idx"), "BlockIdx"); + } + + #[test] + fn parses_views_slice() { + // A tiny slice of the wire spec enough to exercise Option / Vec / + // BTreeMap / serde rename. + let src = r#" + #[derive(Serialize, Deserialize)] + pub struct ChunkBody { + #[serde(rename = "chunkType")] + pub chunk_type: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub delta: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub tools: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub load: Option>, + } + "#; + let ts = structs_in(src).unwrap(); + assert_eq!(ts.len(), 1); + let f = &ts[0].fields; + assert_eq!(f.len(), 4); + assert_eq!(f[0].json, "chunkType"); + assert_eq!(f[0].repr.matches(&Repr::Str), true); + assert_eq!(f[1].optional, true); + assert!(matches!(f[1].repr, Repr::Json)); + assert!(matches!(f[2].repr, Repr::Arr(_))); + assert!(matches!(f[3].repr, Repr::Map(_))); + } + + #[test] + fn parses_real_views_renames() { + let spec = + protocol_spec(&paths::views_rs(), &paths::envelope_rs(), &paths::lib_rs()).unwrap(); + let chunk = spec + .types + .iter() + .find(|t| t.name == "ChunkBody") + .expect("ChunkBody in spec"); + for f in &chunk.fields { + match f.rust.as_str() { + "chunk_type" => assert_eq!(f.json, "chunkType"), + "block_idx" => assert_eq!(f.json, "blockIdx"), + "partial_json" => assert_eq!(f.json, "partialJson"), + "delta" => assert_eq!(f.json, "delta"), + other => panic!("unexpected ChunkBody field {other}"), + } + } + let prompt = spec.types.iter().find(|t| t.name == "PromptBody").unwrap(); + for f in &prompt.fields { + if f.rust == "model_id" { + assert_eq!(f.json, "modelId"); + } + if f.rust == "system_prompt" { + assert_eq!(f.json, "systemPrompt"); + } + } + } + + impl Repr { + fn matches(&self, other: &Repr) -> bool { + std::mem::discriminant(self) == std::mem::discriminant(other) + } + } +} diff --git a/thrum-clients/go/thrum/helpers.go b/thrum-clients/go/thrum/helpers.go index 455d268c..7ff8e6bc 100644 --- a/thrum-clients/go/thrum/helpers.go +++ b/thrum-clients/go/thrum/helpers.go @@ -6,15 +6,15 @@ package thrum import ( + "crypto/rand" "crypto/sha256" "encoding/hex" "encoding/json" - "fmt" + "math/big" "os" "path/filepath" - "strconv" + "strings" "sync" - "sync/atomic" "time" ) @@ -29,19 +29,53 @@ func Sigil(sid, nest string) string { // NowMs returns wall-clock milliseconds since the Unix epoch. func NowMs() int64 { return time.Now().UnixMilli() } -var ridCounter uint64 - -func base36(n uint64) string { - if n == 0 { - return "0" +const crockford = "0123456789ABCDEFGHJKMNPQRSTVWXYZ" + +// HumIdEncode encodes 32 bytes (48-bit ms prefix + 208 random bits, +// big-endian) as a 52-char Crockford base32 id. Mirrors hum_identity's +// encode bit-for-bit. +func HumIdEncode(buf [32]byte) string { + var v big.Int + v.SetBytes(buf[:]) + v.Lsh(&v, 4) + var out strings.Builder + out.Grow(52) + for i := 0; i < 52; i++ { + var idx big.Int + idx.Rsh(&v, uint(255-5*i)) + idx.And(&idx, big.NewInt(0x1f)) + out.WriteByte(crockford[idx.Int64()]) } - return strconv.FormatUint(n, 36) + return out.String() } -// Rid mints a monotonic correlation id: "{base36-ms}-{base36-counter}". +// Rid mints a correlation id in canonical HumId form — 52-char +// Crockford base32, ts-prefixed, matching hum_identity::HumId::mint(). func Rid() string { - n := atomic.AddUint64(&ridCounter, 1) - 1 - return fmt.Sprintf("%s-%s", base36(uint64(NowMs())), base36(n)) + var buf [32]byte + ms := uint64(NowMs()) + for i := 0; i < 6; i++ { + buf[5-i] = byte(ms >> (8 * uint(i))) + } + if _, err := rand.Read(buf[6:]); err != nil { + // 26 bytes from the CSPRNG is practically infallible; fall + // back to a time-seeded value rather than mint nothing. + buf[6] = byte(ms >> 40) + } + return HumIdEncode(buf) +} + +// IsValidRid reports whether value is a 52-char Crockford base32 id. +func IsValidRid(value string) bool { + if len(value) != 52 { + return false + } + for i := 0; i < len(value); i++ { + if !strings.ContainsRune(crockford, rune(value[i])) { + return false + } + } + return true } // DuskIn returns the absolute ms timestamp at which a tone with this diff --git a/thrum-clients/go/thrum/helpers_test.go b/thrum-clients/go/thrum/helpers_test.go new file mode 100644 index 00000000..e36f8b1f --- /dev/null +++ b/thrum-clients/go/thrum/helpers_test.go @@ -0,0 +1,51 @@ +package thrum + +import ( + "strings" + "testing" +) + +// Byte-for-byte parity with hum_identity's encoder. These vectors come +// straight from ids.rs ts_parity_vectors — if hum_identity's alphabet or +// shift layout ever changes, this test names the drift. +func TestHumIdParityVectors(t *testing.T) { + var zero [32]byte + if got := HumIdEncode(zero); got != strings.Repeat("0", 52) { + t.Fatalf("zeros: got %q", got) + } + + // 1_700_000_000_000 ms in the top 6 bytes, 208 zero bits after. + ts := uint64(1700000000000) + var tsOnly [32]byte + for i := 0; i < 6; i++ { + tsOnly[5-i] = byte(ts >> (8 * uint(i))) + } + want := "065WZSB8" + strings.Repeat("0", 44) + if got := HumIdEncode(tsOnly); got != want { + t.Fatalf("ts-only: got %q want %q", got, want) + } + + var ones [32]byte + for i := range ones { + ones[i] = 0xff + } + want = strings.Repeat("Z", 51) + "G" + if got := HumIdEncode(ones); got != want { + t.Fatalf("ones: got %q want %q", got, want) + } +} + +func TestRidIsValidCrockford52(t *testing.T) { + for i := 0; i < 100; i++ { + id := Rid() + if !IsValidRid(id) { + t.Fatalf("Rid()=%q invalid", id) + } + if len(id) != 52 { + t.Fatalf("Rid()=%q len %d", id, len(id)) + } + } + if IsValidRid("") || IsValidRid(strings.Repeat("I", 52)) { + t.Fatal("IsValidRid accepted garbage") + } +} diff --git a/thrum-clients/go/thrum/protocol.go b/thrum-clients/go/thrum/protocol.go new file mode 100644 index 00000000..c4414453 --- /dev/null +++ b/thrum-clients/go/thrum/protocol.go @@ -0,0 +1,262 @@ +// @generated by `cargo run -p codegen` from thrum-core/src/views.rs, thrum-core/src/envelope.rs — DO NOT EDIT. +// +// Canonical wire views: every chi body is a Rust struct in +// thrum-core/src/views.rs, the envelope in envelope.rs. These +// client mirrors come out of the same source — no hand-sync. +// Manual regen: `cargo run -p codegen`. + +package thrum + +import "encoding/json" + +type HelloBody struct { + ProtoVersion string `json:"protoVersion"` + Bee string `json:"bee"` + Version string `json:"version"` + Hid *string `json:"hid,omitempty"` + ToolNames []string `json:"tool_names,omitempty"` + Provides []string `json:"provides,omitempty"` + Source *string `json:"source,omitempty"` +} + +type PromptBody struct { + ModelId *string `json:"modelId,omitempty"` + Cwd *string `json:"cwd,omitempty"` + SystemPrompt *string `json:"systemPrompt,omitempty"` + Text *string `json:"text,omitempty"` + Content []json.RawMessage `json:"content,omitempty"` + Tools []json.RawMessage `json:"tools,omitempty"` + ForagerTools []json.RawMessage `json:"foragerTools,omitempty"` + Provided []string `json:"provided,omitempty"` + DisallowedTools []string `json:"disallowedTools,omitempty"` +} + +type BreathBody struct { + Sessions json.RawMessage `json:"sessions"` + ProtoVersion *string `json:"protoVersion,omitempty"` +} + +type ChunkBody struct { + ChunkType string `json:"chunkType"` + BlockIdx *uint64 `json:"blockIdx,omitempty"` + Delta json.RawMessage `json:"delta,omitempty"` + PartialJson *string `json:"partialJson,omitempty"` +} + +type FinishBody struct { + FinishReason string `json:"finishReason"` + Usage json.RawMessage `json:"usage,omitempty"` + ExitCode *int64 `json:"exitCode,omitempty"` + Subtype *string `json:"subtype,omitempty"` +} + +type ErrorBody struct { + Message string `json:"message"` + Code *string `json:"code,omitempty"` + Subtype *string `json:"subtype,omitempty"` + Usage json.RawMessage `json:"usage,omitempty"` +} + +type SessionReadyBody struct { + NestId string `json:"nestId"` + Model string `json:"model"` + Tools []json.RawMessage `json:"tools,omitempty"` +} + +type ToolCallBody struct { + CallId string `json:"callId"` + ToolName string `json:"toolName"` + Name *string `json:"name,omitempty"` + Args json.RawMessage `json:"args,omitempty"` +} + +type ToolResultBody struct { + CallId string `json:"callId"` + Output *string `json:"output,omitempty"` + Result *string `json:"result,omitempty"` + IsError *bool `json:"isError,omitempty"` + Title *string `json:"title,omitempty"` + Metadata json.RawMessage `json:"metadata,omitempty"` +} + +type ToolInfoBody struct { + CallId string `json:"callId"` + Name *string `json:"name,omitempty"` + Args json.RawMessage `json:"args,omitempty"` + Result json.RawMessage `json:"result,omitempty"` +} + +type ToolMetaBody struct { + CallId *string `json:"callId,omitempty"` + ToolName *string `json:"toolName,omitempty"` + Metadata json.RawMessage `json:"metadata,omitempty"` +} + +type PulseBody struct { + Kind string `json:"kind"` + Pid *int64 `json:"pid,omitempty"` +} + +type PermissionAskBody struct { + CallId *string `json:"callId,omitempty"` + ToolName *string `json:"toolName,omitempty"` + Message *string `json:"message,omitempty"` + Arg json.RawMessage `json:"arg,omitempty"` +} + +type ReleasePermitBody struct { + CallId *string `json:"callId,omitempty"` + Ok *bool `json:"ok,omitempty"` + Error *string `json:"error,omitempty"` +} + +type EchoBody struct { + Ok bool `json:"ok"` + Error *string `json:"error,omitempty"` +} + +type LogBody struct { + Level *string `json:"level,omitempty"` + Message *string `json:"message,omitempty"` + Fields json.RawMessage `json:"fields,omitempty"` +} + +type DroneBody struct { + Health string `json:"health"` + RhythmMs *uint64 `json:"rhythm_ms,omitempty"` + PendingEchoes []json.RawMessage `json:"pending_echoes,omitempty"` + Load *DroneLoad `json:"load,omitempty"` +} + +type DroneLoad struct { + ActiveSigils *uint64 `json:"active_sigils,omitempty"` + PendingPermissions *uint64 `json:"pending_permissions,omitempty"` + InflightTools *uint64 `json:"inflight_tools,omitempty"` + TokensBurned *uint64 `json:"tokens_burned,omitempty"` +} + +type DroneRetrofitBody struct { + Sigil *string `json:"sigil,omitempty"` + Reason *string `json:"reason,omitempty"` +} + +type PeerAddBody struct { + HumdId string `json:"humd_id"` + Hints []json.RawMessage `json:"hints,omitempty"` +} + +type PeerRemoveBody struct { + HumdId string `json:"humd_id"` +} + +type AttachBody struct { + HearOnly *bool `json:"hearOnly,omitempty"` +} + +type DetachBody struct { +} + +type WaneSyncBody struct { + Snapshot map[string]uint64 `json:"snapshot"` +} + +type GossipPublishBody struct { + Topic string `json:"topic"` + Payload json.RawMessage `json:"payload"` + From *string `json:"from,omitempty"` + MsgId string `json:"msg_id"` +} + +type KadFindNodeBody struct { + QueryId string `json:"query_id"` + Target string `json:"target"` +} + +type KadFindNodeRespBody struct { + QueryId string `json:"query_id"` + Closest []json.RawMessage `json:"closest"` +} + +type BackfillBody struct { + Author string `json:"author"` + From *int64 `json:"from,omitempty"` +} + +type PerfMarkBody struct { + At *int64 `json:"at,omitempty"` + Note *string `json:"note,omitempty"` +} + +type CancelBody struct { +} + +type CleanupBody struct { +} + +type CurateBody struct { +} + +type TendrilResultBody struct { + CallId *string `json:"callId,omitempty"` + Output json.RawMessage `json:"output,omitempty"` +} + +type TendrilReachBody struct { + Task *string `json:"task,omitempty"` + Tools []json.RawMessage `json:"tools,omitempty"` +} + +type PetalCellBody struct { + Cell json.RawMessage `json:"cell,omitempty"` +} + +type Envelope struct { + Chi Chi `json:"chi"` + Rid string `json:"rid"` + From *string `json:"from,omitempty"` + To *string `json:"to,omitempty"` + Sigil *string `json:"sigil,omitempty"` + Sid *string `json:"sid,omitempty"` + Wane *uint64 `json:"wane,omitempty"` + SentAt *int64 `json:"sentAt,omitempty"` + Dusk *int64 `json:"dusk,omitempty"` + Ext map[string]map[string]json.RawMessage `json:"ext,omitempty"` +} + +// ToneViews maps each chi wire value to its body view type. +var ToneViews = map[Chi]any{ + ChiHello: HelloBody{}, + ChiPrompt: PromptBody{}, + ChiBreath: BreathBody{}, + ChiChunk: ChunkBody{}, + ChiFinish: FinishBody{}, + ChiError: ErrorBody{}, + ChiSessionReady: SessionReadyBody{}, + ChiToolCall: ToolCallBody{}, + ChiToolResult: ToolResultBody{}, + ChiToolInfo: ToolInfoBody{}, + ChiToolMeta: ToolMetaBody{}, + ChiPulse: PulseBody{}, + ChiPermissionAsk: PermissionAskBody{}, + ChiReleasePermit: ReleasePermitBody{}, + ChiEcho: EchoBody{}, + ChiLog: LogBody{}, + ChiDrone: DroneBody{}, + ChiDroneRetrofit: DroneRetrofitBody{}, + ChiPeerAdd: PeerAddBody{}, + ChiPeerRemove: PeerRemoveBody{}, + ChiAttach: AttachBody{}, + ChiDetach: DetachBody{}, + ChiWaneSync: WaneSyncBody{}, + ChiGossipPublish: GossipPublishBody{}, + ChiKadFindNode: KadFindNodeBody{}, + ChiKadFindNodeResp: KadFindNodeRespBody{}, + ChiBackfill: BackfillBody{}, + ChiPerfMark: PerfMarkBody{}, + ChiCancel: CancelBody{}, + ChiCleanup: CleanupBody{}, + ChiCurate: CurateBody{}, + ChiTendrilResult: TendrilResultBody{}, + ChiTendrilReach: TendrilReachBody{}, + ChiPetalCell: PetalCellBody{}, +} diff --git a/thrum-clients/go/thrum/protocol_test.go b/thrum-clients/go/thrum/protocol_test.go new file mode 100644 index 00000000..a5709652 --- /dev/null +++ b/thrum-clients/go/thrum/protocol_test.go @@ -0,0 +1,126 @@ +package thrum + +// Golden-wire parse test: every committed tone byte-decodes into its +// generated view struct without losing or inventing keys. When views.rs +// changes a wire key or shape, regen + rerun me. + +import ( + "bufio" + "encoding/json" + "os" + "path/filepath" + "reflect" + "strings" + "testing" +) + +var envKeys = map[string]bool{ + "chi": true, "rid": true, "from": true, "to": true, "sigil": true, + "sid": true, "wane": true, "sentAt": true, "dusk": true, "ext": true, +} + +func goldenLines(t *testing.T) []map[string]json.RawMessage { + t.Helper() + fixture := filepath.Join("..", "..", "..", "thrum-core", "tests", "fixtures", "golden.ndjson") + f, err := os.Open(fixture) + if err != nil { + t.Fatalf("open fixture: %v", err) + } + defer f.Close() + var out []map[string]json.RawMessage + sc := bufio.NewScanner(f) + for sc.Scan() { + line := strings.TrimSpace(sc.Text()) + if line == "" { + continue + } + var m map[string]json.RawMessage + if err := json.Unmarshal([]byte(line), &m); err != nil { + t.Fatalf("fixture line: %v", err) + } + out = append(out, m) + } + return out +} + +func TestGoldenDecodesIntoViews(t *testing.T) { + tone := goldenLines(t) + for i, m := range tone { + var chi Chi + if err := json.Unmarshal(m["chi"], &chi); err != nil { + t.Fatalf("line %d: chi: %v", i, err) + } + if !IsValidChi(string(chi)) { + t.Fatalf("line %d: unknown chi %q", i, chi) + } + body := make(map[string]json.RawMessage) + for k, v := range m { + if !envKeys[k] { + body[k] = v + } + } + pick, ok := ToneViews[chi] + if !ok { + t.Fatalf("line %d: chi %q has no view", i, chi) + } + // Allocate a fresh instance of the body view type and unmarshal + // into it — this is the decode path real balks take. + blob, _ := json.Marshal(body) + inst := reflect.New(reflect.TypeOf(pick)) + if err := json.Unmarshal(blob, inst.Interface()); err != nil { + t.Fatalf("line %d (%s): decode into %T: %v", i, chi, pick, err) + } + t.Logf("line %d: %s -> %+v", i, chi, inst.Elem().Interface()) + } +} + +func TestGoldenSpotChecks(t *testing.T) { + for _, m := range goldenLines(t) { + var chi Chi + json.Unmarshal(m["chi"], &chi) + body := make(map[string]json.RawMessage) + for k, v := range m { + if !envKeys[k] { + body[k] = v + } + } + switch chi { + case ChiHello: + var h HelloBody + blob, _ := json.Marshal(body) + if err := json.Unmarshal(blob, &h); err != nil { + t.Fatal(err) + } + if h.ProtoVersion != "0.7.0" || h.Bee != "claude-cli" { + t.Fatalf("hello mismatch: %+v", h) + } + case ChiToolResult: + var r ToolResultBody + blob, _ := json.Marshal(body) + if err := json.Unmarshal(blob, &r); err != nil { + t.Fatal(err) + } + if r.CallId != "call-1" || r.Output == nil || *r.Output != "file contents" { + t.Fatalf("tool-result mismatch: %+v", r) + } + case ChiWaneSync: + var w WaneSyncBody + blob, _ := json.Marshal(body) + if err := json.Unmarshal(blob, &w); err != nil { + t.Fatal(err) + } + if w.Snapshot["a1b2c3d4e5f6"] != 12 { + t.Fatalf("wane-sync mismatch: %+v", w) + } + case ChiEcho: + var e EchoBody + blob, _ := json.Marshal(body) + if err := json.Unmarshal(blob, &e); err != nil { + t.Fatal(err) + } + if !e.Ok { + t.Fatalf("echo mismatch: %+v", e) + } + } + } +} diff --git a/thrum-clients/python/tests/test_helpers.py b/thrum-clients/python/tests/test_helpers.py new file mode 100644 index 00000000..e9f0831b --- /dev/null +++ b/thrum-clients/python/tests/test_helpers.py @@ -0,0 +1,37 @@ +#!/usr/bin/env python3 +"""rid() parity vectors against hum_identity's encoder, taken verbatim +from thrum-core/../hum-identity/src/ids.rs ts_parity_vectors.""" + +import os +import sys +import unittest + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + +from thrum.helpers import _hum_id_from_bytes, is_valid_rid, rid # noqa: E402 + + +class RidParityTest(unittest.TestCase): + def test_vectors_match_rust(self): + self.assertEqual(_hum_id_from_bytes(bytes(32)), "0" * 52) + self.assertEqual( + _hum_id_from_bytes((1700000000000).to_bytes(6, "big") + bytes(26)), + "065WZSB8" + "0" * 44, + ) + self.assertEqual(_hum_id_from_bytes(bytes([0xFF] * 32)), "Z" * 51 + "G") + + def test_rid_shape(self): + for _ in range(100): + id_ = rid() + self.assertTrue(is_valid_rid(id_), id_) + self.assertEqual(len(id_), 52) + + def test_is_valid_rid(self): + self.assertFalse(is_valid_rid("")) + self.assertFalse(is_valid_rid("I" * 52)) + self.assertFalse(is_valid_rid("0" * 51)) + self.assertTrue(is_valid_rid(rid())) + + +if __name__ == "__main__": + unittest.main() diff --git a/thrum-clients/python/tests/test_protocol.py b/thrum-clients/python/tests/test_protocol.py new file mode 100644 index 00000000..504b1ac7 --- /dev/null +++ b/thrum-clients/python/tests/test_protocol.py @@ -0,0 +1,80 @@ +#!/usr/bin/env python3 +"""Golden-wire parse test: every committed tone decodes into its typed +view with the exact wire keys views.rs declares. Drift in views.rs or in +this client's generated mirrors fails here.""" + +import json +import os +import sys +import unittest + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + +from thrum import protocol as p # noqa: E402 + +FIXTURE = os.path.join( + os.path.dirname(os.path.abspath(__file__)), + "..", + "..", + "..", + "thrum-core", + "tests", + "fixtures", + "golden.ndjson", +) + +ENV_KEYS = {"chi", "rid", "from", "to", "sigil", "sid", "wane", "sentAt", "dusk", "ext"} + + +def lines(): + with open(FIXTURE) as f: + return f.read().splitlines() + + +class GoldenWireTest(unittest.TestCase): + def test_every_line_is_a_known_tone(self): + for raw in lines(): + tone = json.loads(raw) + self.assertIn(tone["chi"], p.TONE_VIEWS, raw) + self.assertIn("rid", tone) + self.assertIn("sentAt", tone) + + def test_bodies_decode_through_generated_views(self): + decoded = 0 + for raw in lines(): + tone = json.loads(raw) + chi = tone["chi"] + body = {k: v for k, v in tone.items() if k not in ENV_KEYS} + view = p.from_wire(p.TONE_VIEWS[chi], body) + # every wire key must be a declared field of the view + declared = set(p.TONE_VIEWS[chi]._wire.values()) + unknown = set(body) - declared + self.assertEqual(unknown, set(), f"{raw}: {unknown}") + # and the view must re-serialize to exactly this body + self.assertEqual(p.to_wire(view), body, raw) + decoded += 1 + self.assertEqual(decoded, 15) + + def test_spot_checks(self): + tones = {json.loads(r)["chi"]: json.loads(r) for r in lines()} + hello = p.from_wire( + p.HelloBody, {k: v for k, v in tones["hello"].items() if k not in ENV_KEYS} + ) + self.assertEqual(hello.proto_version, "0.7.0") + self.assertEqual(hello.tool_names, ["Read", "Bash"]) + chunk = p.from_wire( + p.ChunkBody, {k: v for k, v in tones["chunk"].items() if k not in ENV_KEYS} + ) + self.assertEqual(chunk.chunk_type, "tool_input_delta") + self.assertEqual(chunk.block_idx, 2) + self.assertEqual(chunk.partial_json, '{"a":') + result = p.from_wire( + p.ToolResultBody, + {k: v for k, v in tones["tool-result"].items() if k not in ENV_KEYS}, + ) + self.assertEqual(result.output, "file contents") + self.assertEqual(result.is_error, False) + + +if __name__ == "__main__": + unittest.main() diff --git a/thrum-clients/python/thrum/helpers.py b/thrum-clients/python/thrum/helpers.py index 99322894..f61bb049 100644 --- a/thrum-clients/python/thrum/helpers.py +++ b/thrum-clients/python/thrum/helpers.py @@ -8,6 +8,7 @@ import hashlib import json import os +import secrets import threading import time from typing import Any, Mapping @@ -29,28 +30,30 @@ def now_ms() -> int: return int(time.time() * 1000) -_RID_LOCK = threading.Lock() -_RID_COUNTER = 0 +_CROCKFORD = "0123456789ABCDEFGHJKMNPQRSTVWXYZ" -def _base36(n: int) -> str: - if n == 0: - return "0" - alphabet = "0123456789abcdefghijklmnopqrstuvwxyz" - out = [] - while n > 0: - out.append(alphabet[n % 36]) - n //= 36 - return "".join(reversed(out)) +def _hum_id_from_bytes(buf: bytes) -> str: + """Encode 32 bytes (48-bit ms prefix + 208 random bits) as a 52-char + Crockford base32 id. Mirrors hum_identity's encode bit-for-bit.""" + v = int.from_bytes(buf, "big") << 4 + return "".join(_CROCKFORD[(v >> (255 - 5 * i)) & 0x1F] for i in range(52)) def rid() -> str: - """Monotonic correlation id: '{base36-ms}-{base36-counter}'.""" - global _RID_COUNTER - with _RID_LOCK: - n = _RID_COUNTER - _RID_COUNTER += 1 - return f"{_base36(now_ms())}-{_base36(n)}" + """Correlation id in canonical HumId form: 52-char Crockford base32, + ts-prefixed, matching hum_identity::HumId::mint().""" + ts = int(time.time() * 1000) + buf = ts.to_bytes(6, "big") + secrets.token_bytes(26) + return _hum_id_from_bytes(buf) + + +def is_valid_rid(value: str) -> bool: + """True iff `value` is a 52-char Crockford base32 id.""" + return ( + len(value) == 52 + and all(c in _CROCKFORD for c in value) + ) def dusk_in(ms: int) -> int: diff --git a/thrum-clients/python/thrum/protocol.py b/thrum-clients/python/thrum/protocol.py new file mode 100644 index 00000000..ab7c08d5 --- /dev/null +++ b/thrum-clients/python/thrum/protocol.py @@ -0,0 +1,521 @@ +# @generated by `cargo run -p codegen` from thrum-core/src/views.rs, thrum-core/src/envelope.rs — DO NOT EDIT. +# +# Canonical wire views: every chi body is a Rust struct in +# thrum-core/src/views.rs, the envelope in envelope.rs. These +# client mirrors come out of the same source — no hand-sync. +# Manual regen: `cargo run -p codegen`. + +from __future__ import annotations + +from dataclasses import dataclass +from typing import Any, ClassVar, Optional + +@dataclass +class HelloBody: + proto_version: str + bee: str + version: str + hid: Optional[str] = None + tool_names: Optional[list[str]] = None + provides: Optional[list[str]] = None + source: Optional[str] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "proto_version": "protoVersion", + "bee": "bee", + "version": "version", + "hid": "hid", + "tool_names": "tool_names", + "provides": "provides", + "source": "source", + } + +@dataclass +class PromptBody: + model_id: Optional[str] = None + cwd: Optional[str] = None + system_prompt: Optional[str] = None + text: Optional[str] = None + content: Optional[list[Any]] = None + tools: Optional[list[Any]] = None + forager_tools: Optional[list[Any]] = None + provided: Optional[list[str]] = None + disallowed_tools: Optional[list[str]] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "model_id": "modelId", + "cwd": "cwd", + "system_prompt": "systemPrompt", + "text": "text", + "content": "content", + "tools": "tools", + "forager_tools": "foragerTools", + "provided": "provided", + "disallowed_tools": "disallowedTools", + } + +@dataclass +class BreathBody: + sessions: Any + proto_version: Optional[str] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "sessions": "sessions", + "proto_version": "protoVersion", + } + +@dataclass +class ChunkBody: + chunk_type: str + block_idx: Optional[int] = None + delta: Optional[Any] = None + partial_json: Optional[str] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "chunk_type": "chunkType", + "block_idx": "blockIdx", + "delta": "delta", + "partial_json": "partialJson", + } + +@dataclass +class FinishBody: + finish_reason: str + usage: Optional[Any] = None + exit_code: Optional[int] = None + subtype: Optional[str] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "finish_reason": "finishReason", + "usage": "usage", + "exit_code": "exitCode", + "subtype": "subtype", + } + +@dataclass +class ErrorBody: + message: str + code: Optional[str] = None + subtype: Optional[str] = None + usage: Optional[Any] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "message": "message", + "code": "code", + "subtype": "subtype", + "usage": "usage", + } + +@dataclass +class SessionReadyBody: + nest_id: str + model: str + tools: Optional[list[Any]] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "nest_id": "nestId", + "model": "model", + "tools": "tools", + } + +@dataclass +class ToolCallBody: + call_id: str + tool_name: str + name: Optional[str] = None + args: Optional[Any] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "call_id": "callId", + "tool_name": "toolName", + "name": "name", + "args": "args", + } + +@dataclass +class ToolResultBody: + call_id: str + output: Optional[str] = None + result: Optional[str] = None + is_error: Optional[bool] = None + title: Optional[str] = None + metadata: Optional[Any] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "call_id": "callId", + "output": "output", + "result": "result", + "is_error": "isError", + "title": "title", + "metadata": "metadata", + } + +@dataclass +class ToolInfoBody: + call_id: str + name: Optional[str] = None + args: Optional[Any] = None + result: Optional[Any] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "call_id": "callId", + "name": "name", + "args": "args", + "result": "result", + } + +@dataclass +class ToolMetaBody: + call_id: Optional[str] = None + tool_name: Optional[str] = None + metadata: Optional[Any] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "call_id": "callId", + "tool_name": "toolName", + "metadata": "metadata", + } + +@dataclass +class PulseBody: + kind: str + pid: Optional[int] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "kind": "kind", + "pid": "pid", + } + +@dataclass +class PermissionAskBody: + call_id: Optional[str] = None + tool_name: Optional[str] = None + message: Optional[str] = None + arg: Optional[Any] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "call_id": "callId", + "tool_name": "toolName", + "message": "message", + "arg": "arg", + } + +@dataclass +class ReleasePermitBody: + call_id: Optional[str] = None + ok: Optional[bool] = None + error: Optional[str] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "call_id": "callId", + "ok": "ok", + "error": "error", + } + +@dataclass +class EchoBody: + ok: bool + error: Optional[str] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "ok": "ok", + "error": "error", + } + +@dataclass +class LogBody: + level: Optional[str] = None + message: Optional[str] = None + fields: Optional[Any] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "level": "level", + "message": "message", + "fields": "fields", + } + +@dataclass +class DroneBody: + health: str + rhythm_ms: Optional[int] = None + pending_echoes: Optional[list[Any]] = None + load: Optional[DroneLoad] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "health": "health", + "rhythm_ms": "rhythm_ms", + "pending_echoes": "pending_echoes", + "load": "load", + } + +@dataclass +class DroneLoad: + active_sigils: Optional[int] = None + pending_permissions: Optional[int] = None + inflight_tools: Optional[int] = None + tokens_burned: Optional[int] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "active_sigils": "active_sigils", + "pending_permissions": "pending_permissions", + "inflight_tools": "inflight_tools", + "tokens_burned": "tokens_burned", + } + +@dataclass +class DroneRetrofitBody: + sigil: Optional[str] = None + reason: Optional[str] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "sigil": "sigil", + "reason": "reason", + } + +@dataclass +class PeerAddBody: + humd_id: str + hints: Optional[list[Any]] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "humd_id": "humd_id", + "hints": "hints", + } + +@dataclass +class PeerRemoveBody: + humd_id: str + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "humd_id": "humd_id", + } + +@dataclass +class AttachBody: + hear_only: Optional[bool] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "hear_only": "hearOnly", + } + +@dataclass +class DetachBody: + pass + +@dataclass +class WaneSyncBody: + snapshot: dict[str, int] + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "snapshot": "snapshot", + } + +@dataclass +class GossipPublishBody: + topic: str + payload: Any + msg_id: str + from_: Optional[str] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "topic": "topic", + "payload": "payload", + "msg_id": "msg_id", + "from_": "from", + } + +@dataclass +class KadFindNodeBody: + query_id: str + target: str + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "query_id": "query_id", + "target": "target", + } + +@dataclass +class KadFindNodeRespBody: + query_id: str + closest: list[Any] + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "query_id": "query_id", + "closest": "closest", + } + +@dataclass +class BackfillBody: + author: str + from_: Optional[int] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "author": "author", + "from_": "from", + } + +@dataclass +class PerfMarkBody: + at: Optional[int] = None + note: Optional[str] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "at": "at", + "note": "note", + } + +@dataclass +class CancelBody: + pass + +@dataclass +class CleanupBody: + pass + +@dataclass +class CurateBody: + pass + +@dataclass +class TendrilResultBody: + call_id: Optional[str] = None + output: Optional[Any] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "call_id": "callId", + "output": "output", + } + +@dataclass +class TendrilReachBody: + task: Optional[str] = None + tools: Optional[list[Any]] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "task": "task", + "tools": "tools", + } + +@dataclass +class PetalCellBody: + cell: Optional[Any] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "cell": "cell", + } + +@dataclass +class Envelope: + chi: Chi + rid: str + from_: Optional[str] = None + to: Optional[str] = None + sigil: Optional[str] = None + sid: Optional[str] = None + wane: Optional[int] = None + sent_at: Optional[int] = None + dusk: Optional[int] = None + ext: Optional[dict[str, dict[str, Any]]] = None + + # rust field name -> wire key (single source in views.rs) + _wire: ClassVar[dict[str, str]] = { + "chi": "chi", + "rid": "rid", + "from_": "from", + "to": "to", + "sigil": "sigil", + "sid": "sid", + "wane": "wane", + "sent_at": "sentAt", + "dusk": "dusk", + "ext": "ext", + } + +# Serialize a view to the flat wire body (envelope keys not included). +# Absent optionals, like skip_serializing_if on the Rust side, are omitted. +def to_wire(view: Any) -> dict[str, Any]: + payload: dict[str, Any] = {} + for rust_name, wire_key in type(view)._wire.items(): + value = getattr(view, rust_name) + if value is not None: + payload[wire_key] = value + return payload + +# Parse a flat wire body into the view for `cls`, ignoring envelope +# keys (already stripped) and any unknown keys. Missing required +# fields raise TypeError, surfacing spec drift loudly. +def from_wire(cls: type, payload: dict[str, Any]) -> Any: + wire_to_rust = {j: f for f, j in cls._wire.items()} + kwargs = { + wire_to_rust[k]: v for k, v in payload.items() if k in wire_to_rust + } + return cls(**kwargs) + +# chi wire value -> body view class. +TONE_VIEWS: dict[str, type] = { + "hello": HelloBody, + "prompt": PromptBody, + "breath": BreathBody, + "chunk": ChunkBody, + "finish": FinishBody, + "error": ErrorBody, + "session-ready": SessionReadyBody, + "tool-call": ToolCallBody, + "tool-result": ToolResultBody, + "tool-info": ToolInfoBody, + "tool-meta": ToolMetaBody, + "pulse": PulseBody, + "permission-ask": PermissionAskBody, + "release-permit": ReleasePermitBody, + "echo": EchoBody, + "log": LogBody, + "drone": DroneBody, + "drone-retrofit": DroneRetrofitBody, + "peer-add": PeerAddBody, + "peer-remove": PeerRemoveBody, + "attach": AttachBody, + "detach": DetachBody, + "wane-sync": WaneSyncBody, + "gossip-publish": GossipPublishBody, + "kad-find-node": KadFindNodeBody, + "kad-find-node-resp": KadFindNodeRespBody, + "backfill": BackfillBody, + "perf-mark": PerfMarkBody, + "cancel": CancelBody, + "cleanup": CleanupBody, + "curate": CurateBody, + "tendril-result": TendrilResultBody, + "tendril-reach": TendrilReachBody, + "petal-cell": PetalCellBody, +} diff --git a/thrum-clients/ts/helpers.ts b/thrum-clients/ts/helpers.ts index 029a7c67..546d0c0c 100644 --- a/thrum-clients/ts/helpers.ts +++ b/thrum-clients/ts/helpers.ts @@ -4,7 +4,7 @@ // so the TS side cannot drift from the Rust side. If you need a new // helper, add it in Rust first and extend codegen's render_helpers. -import { createHash } from "crypto"; +import { createHash, randomBytes } from "crypto"; /** * Deterministic identity for a (nest, session) pair. @@ -22,10 +22,45 @@ export function sigil(sid: string, nest: string): string { .slice(0, 12); } -/** Monotonic request id — base36 timestamp + counter. */ -let __ridCounter = 0; +const CROCKFORD = "0123456789ABCDEFGHJKMNPQRSTVWXYZ"; + +/** + * Encode 32 bytes (48-bit ms prefix + 208 random bits, big-endian) as a + * 52-char Crockford base32 id. Mirrors hum_identity's encode bit-for-bit. + */ +export function humIdEncode(buf: Uint8Array): string { + let v = 0n; + for (const b of buf) v = (v << 8n) | BigInt(b); + v <<= 4n; + let out = ""; + for (let i = 0; i < 52; i++) { + out += CROCKFORD[Number((v >> BigInt(255 - 5 * i)) & 0x1fn)]; + } + return out; +} + +/** + * Correlation id in canonical HumId form — 52-char Crockford base32, + * ts-prefixed, matching hum_identity::HumId::mint(). + */ export function rid(): string { - return `${Date.now().toString(36)}-${(__ridCounter++).toString(36)}`; + const buf = new Uint8Array(32); + let ms = BigInt(Date.now()); + for (let i = 0; i < 6; i++) { + buf[5 - i] = Number(ms & 0xffn); + ms >>= 8n; + } + buf.set(randomBytes(26), 6); + return humIdEncode(buf); +} + +/** True iff `value` is a 52-char Crockford base32 id. */ +export function isValidRid(value: unknown): value is string { + return ( + typeof value === "string" && + value.length === 52 && + [...value].every((c) => CROCKFORD.includes(c)) + ); } /** Absolute ms timestamp ms in the future. */ diff --git a/thrum-clients/ts/index.ts b/thrum-clients/ts/index.ts index 811bd212..4d6b6e69 100644 --- a/thrum-clients/ts/index.ts +++ b/thrum-clients/ts/index.ts @@ -1,16 +1,13 @@ // Thrum — TS surface of the hum protocol. // -// Both files in this directory are generated from thrum-core (the Rust -// source of truth): -// - chi.ts — Chi enum, PulseKind, Envelope, validators +// All three generated files come from thrum-core (the Rust source): +// - chi.ts — Chi enum, PulseKind, validators // - helpers.ts — sigil, rid, duskIn, isDusk, WaneTracker +// - protocol.ts — per-chi body views + the definitive Envelope // -// Hand-edit chi.rs (or extend codegen for new helpers). cargo build -// regenerates both files via thrum-core/build.rs. -// +// Hand-edit the Rust source, then regen (`cargo run -p codegen`). // This index.ts is a thin barrel — every export here flows through -// from one of the two generated files. New protocol primitives should -// be added in Rust first, then surfaced via codegen. +// from one of the three generated files. export { // Registry + version @@ -24,12 +21,7 @@ export { isKnownTone, } from "./chi.ts"; -export type { - ChiKind, - PulseKindT, - Envelope, - Tone, -} from "./chi.ts"; +export type { ChiKind, PulseKindT, Tone } from "./chi.ts"; export { sigil, @@ -38,3 +30,7 @@ export { isDusk, WaneTracker, } from "./helpers.ts"; + +// Protocol views — per-chi body types, ToneViews, plus the definitive +// wire Envelope. Prefer this Envelope over chi.ts's validator-local copy. +export * from "./protocol.ts"; diff --git a/thrum-clients/ts/package.json b/thrum-clients/ts/package.json index 0dfd0b83..056e9698 100644 --- a/thrum-clients/ts/package.json +++ b/thrum-clients/ts/package.json @@ -8,12 +8,14 @@ "exports": { ".": "./index.ts", "./chi": "./chi.ts", - "./helpers": "./helpers.ts" + "./helpers": "./helpers.ts", + "./protocol": "./protocol.ts" }, "files": [ "index.ts", "chi.ts", "helpers.ts", + "protocol.ts", "README.md" ], "keywords": [ diff --git a/thrum-clients/ts/protocol.ts b/thrum-clients/ts/protocol.ts new file mode 100644 index 00000000..4bad6a27 --- /dev/null +++ b/thrum-clients/ts/protocol.ts @@ -0,0 +1,263 @@ +// @generated by `cargo run -p codegen` from thrum-core/src/views.rs, thrum-core/src/envelope.rs — DO NOT EDIT. +// +// Canonical wire views: every chi body is a Rust struct in +// thrum-core/src/views.rs, the envelope in envelope.rs. These +// client mirrors come out of the same source — no hand-sync. +// Manual regen: `cargo run -p codegen`. + +export type { ChiKind } from "./chi.ts"; +import type { Chi } from "./chi.ts"; + +export interface HelloBody { + protoVersion: string; + bee: string; + version: string; + hid?: string; + tool_names?: string[]; + provides?: string[]; + source?: string; +} + +export interface PromptBody { + modelId?: string; + cwd?: string; + systemPrompt?: string; + text?: string; + content?: unknown[]; + tools?: unknown[]; + foragerTools?: unknown[]; + provided?: string[]; + disallowedTools?: string[]; +} + +export interface BreathBody { + sessions: unknown; + protoVersion?: string; +} + +export interface ChunkBody { + chunkType: string; + blockIdx?: number; + delta?: unknown; + partialJson?: string; +} + +export interface FinishBody { + finishReason: string; + usage?: unknown; + exitCode?: number; + subtype?: string; +} + +export interface ErrorBody { + message: string; + code?: string; + subtype?: string; + usage?: unknown; +} + +export interface SessionReadyBody { + nestId: string; + model: string; + tools?: unknown[]; +} + +export interface ToolCallBody { + callId: string; + toolName: string; + name?: string; + args?: unknown; +} + +export interface ToolResultBody { + callId: string; + output?: string; + result?: string; + isError?: boolean; + title?: string; + metadata?: unknown; +} + +export interface ToolInfoBody { + callId: string; + name?: string; + args?: unknown; + result?: unknown; +} + +export interface ToolMetaBody { + callId?: string; + toolName?: string; + metadata?: unknown; +} + +export interface PulseBody { + kind: string; + pid?: number; +} + +export interface PermissionAskBody { + callId?: string; + toolName?: string; + message?: string; + arg?: unknown; +} + +export interface ReleasePermitBody { + callId?: string; + ok?: boolean; + error?: string; +} + +export interface EchoBody { + ok: boolean; + error?: string; +} + +export interface LogBody { + level?: string; + message?: string; + fields?: unknown; +} + +export interface DroneBody { + health: string; + rhythm_ms?: number; + pending_echoes?: unknown[]; + load?: DroneLoad; +} + +export interface DroneLoad { + active_sigils?: number; + pending_permissions?: number; + inflight_tools?: number; + tokens_burned?: number; +} + +export interface DroneRetrofitBody { + sigil?: string; + reason?: string; +} + +export interface PeerAddBody { + humd_id: string; + hints?: unknown[]; +} + +export interface PeerRemoveBody { + humd_id: string; +} + +export interface AttachBody { + hearOnly?: boolean; +} + +export interface DetachBody { +} + +export interface WaneSyncBody { + snapshot: Record; +} + +export interface GossipPublishBody { + topic: string; + payload: unknown; + from?: string; + msg_id: string; +} + +export interface KadFindNodeBody { + query_id: string; + target: string; +} + +export interface KadFindNodeRespBody { + query_id: string; + closest: unknown[]; +} + +export interface BackfillBody { + author: string; + from?: number; +} + +export interface PerfMarkBody { + at?: number; + note?: string; +} + +export interface CancelBody { +} + +export interface CleanupBody { +} + +export interface CurateBody { +} + +export interface TendrilResultBody { + callId?: string; + output?: unknown; +} + +export interface TendrilReachBody { + task?: string; + tools?: unknown[]; +} + +export interface PetalCellBody { + cell?: unknown; +} + +export interface Envelope { + chi: Chi; + rid: string; + from?: string; + to?: string; + sigil?: string; + sid?: string; + wane?: number; + sentAt?: number; + dusk?: number; + ext?: Record>; +} + +// Body view per chi wire value. +export interface ToneViews { + "hello": HelloBody; + "prompt": PromptBody; + "breath": BreathBody; + "chunk": ChunkBody; + "finish": FinishBody; + "error": ErrorBody; + "session-ready": SessionReadyBody; + "tool-call": ToolCallBody; + "tool-result": ToolResultBody; + "tool-info": ToolInfoBody; + "tool-meta": ToolMetaBody; + "pulse": PulseBody; + "permission-ask": PermissionAskBody; + "release-permit": ReleasePermitBody; + "echo": EchoBody; + "log": LogBody; + "drone": DroneBody; + "drone-retrofit": DroneRetrofitBody; + "peer-add": PeerAddBody; + "peer-remove": PeerRemoveBody; + "attach": AttachBody; + "detach": DetachBody; + "wane-sync": WaneSyncBody; + "gossip-publish": GossipPublishBody; + "kad-find-node": KadFindNodeBody; + "kad-find-node-resp": KadFindNodeRespBody; + "backfill": BackfillBody; + "perf-mark": PerfMarkBody; + "cancel": CancelBody; + "cleanup": CleanupBody; + "curate": CurateBody; + "tendril-result": TendrilResultBody; + "tendril-reach": TendrilReachBody; + "petal-cell": PetalCellBody; +} + +export type ToneView = Envelope & ToneViews[C]; diff --git a/thrum-core/build.rs b/thrum-core/build.rs index afc3e537..0ad3e026 100644 --- a/thrum-core/build.rs +++ b/thrum-core/build.rs @@ -11,9 +11,13 @@ use codegen::paths; fn main() { let chi_rs = paths::chi_rs(); let lib_rs = paths::lib_rs(); + let views_rs = paths::views_rs(); + let envelope_rs = paths::envelope_rs(); println!("cargo:rerun-if-changed={}", chi_rs.display()); println!("cargo:rerun-if-changed={}", lib_rs.display()); + println!("cargo:rerun-if-changed={}", views_rs.display()); + println!("cargo:rerun-if-changed={}", envelope_rs.display()); println!("cargo:rerun-if-changed=build.rs"); let spec = match codegen::parse(&chi_rs, &lib_rs) { @@ -23,17 +27,49 @@ fn main() { return; } }; + let proto = match codegen::protocol::protocol_spec(&views_rs, &envelope_rs, &lib_rs) { + Ok(p) => p, + Err(e) => { + println!("cargo:warning=thrum-core build.rs: protocol parse failed: {e}"); + return; + } + }; + if let Err(e) = codegen::protocol::check_coverage(&spec.chi, &proto) { + println!("cargo:warning=thrum-core build.rs: coverage: {e}"); + return; + } - let emits: [(&str, std::path::PathBuf, &dyn Fn(&std::path::Path) -> codegen::Result<()>); 6] = [ - ("emit_ts", paths::ts_chi(), &|p| codegen::emit_ts(&spec, p)), - ("emit_helpers", paths::ts_helpers(), &codegen::emit_helpers), - ("emit_py", paths::py_chi(), &|p| codegen::emit_py(&spec, p)), - ("emit_py_helpers", paths::py_helpers(), &codegen::emit_py_helpers), - ("emit_go", paths::go_chi(), &|p| codegen::emit_go(&spec, p)), - ("emit_go_helpers", paths::go_helpers(), &codegen::emit_go_helpers), + let emits: [(&str, std::path::PathBuf, &dyn Fn() -> codegen::Result<()>); 9] = [ + ("emit_ts", paths::ts_chi(), &|| { + codegen::emit_ts(&spec, &paths::ts_chi()) + }), + ("emit_helpers", paths::ts_helpers(), &|| { + codegen::emit_helpers(&paths::ts_helpers()) + }), + ("emit_protocol_ts", paths::ts_protocol(), &|| { + codegen::protocol::emit_protocol_ts(&proto, &paths::ts_protocol()) + }), + ("emit_py", paths::py_chi(), &|| { + codegen::emit_py(&spec, &paths::py_chi()) + }), + ("emit_py_helpers", paths::py_helpers(), &|| { + codegen::emit_py_helpers(&paths::py_helpers()) + }), + ("emit_protocol_py", paths::py_protocol(), &|| { + codegen::protocol::emit_protocol_py(&proto, &paths::py_protocol()) + }), + ("emit_go", paths::go_chi(), &|| { + codegen::emit_go(&spec, &paths::go_chi()) + }), + ("emit_go_helpers", paths::go_helpers(), &|| { + codegen::emit_go_helpers(&paths::go_helpers()) + }), + ("emit_protocol_go", paths::go_protocol(), &|| { + codegen::protocol::emit_protocol_go(&proto, &paths::go_protocol()) + }), ]; - for (label, out, emit) in &emits { - if let Err(e) = emit(out) { + for (label, _out, emit) in &emits { + if let Err(e) = emit() { println!("cargo:warning=thrum-core build.rs: {label}: {e}"); } } diff --git a/thrum-core/src/lib.rs b/thrum-core/src/lib.rs index 90281ace..79e91f8d 100644 --- a/thrum-core/src/lib.rs +++ b/thrum-core/src/lib.rs @@ -9,11 +9,13 @@ mod chi; mod envelope; mod prims; +mod views; mod wane; pub use chi::{Chi, PulseKind}; pub use envelope::{Envelope, Tone}; pub use prims::{dusk_in, echo_for, is_dusk, now_ms, rid, sigil}; +pub use views::*; pub use wane::WaneTracker; /// Protocol semver — independent of any package version. diff --git a/thrum-core/src/views.rs b/thrum-core/src/views.rs new file mode 100644 index 00000000..d1247912 --- /dev/null +++ b/thrum-core/src/views.rs @@ -0,0 +1,419 @@ +//! Canonical per-chi tone bodies. +//! +//! One struct per chi, field names renamed to EXACTLY the wire keys the +//! runtime produces. This module is the single source of truth for body +//! shapes; codegen emits the TS/Python/Go mirrors from it, and the +//! golden-message test locks the bytes. +//! +//! Optional fields are also `skip_serializing_if` so absent fields are +//! absent from the wire, mirroring how builders construct tones. + +use std::collections::BTreeMap; + +use serde::{Deserialize, Serialize}; +use serde_json::Value; + +/// `hello` — nest → daemon bootstrap. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct HelloBody { + #[serde(rename = "protoVersion")] + pub proto_version: String, + pub bee: String, + pub version: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub hid: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub tool_names: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub provides: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub source: Option, +} + +/// `prompt` — start a turn. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct PromptBody { + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "modelId")] + pub model_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub cwd: Option, + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "systemPrompt")] + pub system_prompt: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub text: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub content: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub tools: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "foragerTools")] + pub forager_tools: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub provided: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "disallowedTools")] + pub disallowed_tools: Option>, +} + +/// `breath` — full state sync on connect. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct BreathBody { + pub sessions: Value, + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "protoVersion")] + pub proto_version: Option, +} + +/// `chunk` — model output partwise. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ChunkBody { + #[serde(rename = "chunkType")] + pub chunk_type: String, + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "blockIdx")] + pub block_idx: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub delta: Option, + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "partialJson")] + pub partial_json: Option, +} + +/// `finish` — turn complete. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct FinishBody { + #[serde(rename = "finishReason")] + pub finish_reason: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub usage: Option, + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "exitCode")] + pub exit_code: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub subtype: Option, +} + +/// `error` — turn aborted. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ErrorBody { + pub message: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub code: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub subtype: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub usage: Option, +} + +/// `session-ready` — nest spawned, claude session id known. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SessionReadyBody { + #[serde(rename = "nestId")] + pub nest_id: String, + pub model: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub tools: Option>, +} + +/// `tool-call` — nestler-declared tool dispatch. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ToolCallBody { + #[serde(rename = "callId")] + pub call_id: String, + #[serde(rename = "toolName")] + pub tool_name: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub name: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub args: Option, +} + +/// `tool-result` — nestler-declared tool answered. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ToolResultBody { + #[serde(rename = "callId")] + pub call_id: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub output: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub result: Option, + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "isError")] + pub is_error: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub title: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub metadata: Option, +} + +/// `tool-info` — completed-in-server tool run, informational. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ToolInfoBody { + #[serde(rename = "callId")] + pub call_id: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub name: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub args: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub result: Option, +} + +/// `tool-meta` — out-of-band metadata for a tool result. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ToolMetaBody { + #[serde(rename = "callId")] + pub call_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "toolName")] + pub tool_name: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub metadata: Option, +} + +/// `pulse` — process lifecycle event. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct PulseBody { + pub kind: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub pid: Option, +} + +/// `permission-ask` — mid-stream permission needed. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct PermissionAskBody { + #[serde(rename = "callId")] + pub call_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "toolName")] + pub tool_name: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub message: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub arg: Option, +} + +/// `release-permit` — resolve an earlier permission-ask. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ReleasePermitBody { + #[serde(rename = "callId")] + pub call_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub ok: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +/// `echo` — delivery ack for a rid. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct EchoBody { + pub ok: bool, + #[serde(skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +/// `log` — structured log forwarding. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogBody { + #[serde(skip_serializing_if = "Option::is_none")] + pub level: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub message: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub fields: Option, +} + +/// `drone` — drone heartbeat. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DroneBody { + pub health: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub rhythm_ms: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub pending_echoes: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub load: Option, +} + +/// `drone.load` +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DroneLoad { + #[serde(skip_serializing_if = "Option::is_none")] + pub active_sigils: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub pending_permissions: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub inflight_tools: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub tokens_burned: Option, +} + +/// `drone-retrofit` — swallow + retry signal. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DroneRetrofitBody { + #[serde(skip_serializing_if = "Option::is_none")] + pub sigil: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub reason: Option, +} + +/// `peer-add` — register a peer humd. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct PeerAddBody { + pub humd_id: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub hints: Option>, +} + +/// `peer-remove` — drop a peer humd. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct PeerRemoveBody { + pub humd_id: String, +} + +/// `attach` — peer humd observes a sid. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct AttachBody { + #[serde(rename = "hearOnly")] + pub hear_only: Option, +} + +/// `detach` — peer humd stops observing. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DetachBody {} + +/// `wane-sync` — reconcile WaneTracker after partition heal. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WaneSyncBody { + pub snapshot: BTreeMap, +} + +/// `gossip-publish` — ensemble-wide mesh message. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct GossipPublishBody { + pub topic: String, + pub payload: Value, + #[serde(skip_serializing_if = "Option::is_none")] + pub from: Option, + #[serde(rename = "msg_id")] + pub msg_id: String, +} + +/// `kad-find-node` — DHT FIND_NODE query. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct KadFindNodeBody { + #[serde(rename = "query_id")] + pub query_id: String, + pub target: String, +} + +/// `kad-find-node-resp` — DHT FIND_NODE response. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct KadFindNodeRespBody { + #[serde(rename = "query_id")] + pub query_id: String, + pub closest: Vec, +} + +/// `backfill` — thehum chi-log replay request. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct BackfillBody { + pub author: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub from: Option, +} + +/// `perf-mark` — drift timing. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct PerfMarkBody { + #[serde(skip_serializing_if = "Option::is_none")] + pub at: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub note: Option, +} + +/// `cancel` — interrupt mid-turn. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CancelBody {} + +/// `cleanup` — session deleted. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CleanupBody {} + +/// `curate` — manual compaction request. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CurateBody {} + +/// `tendril-result` — task subagent answered. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct TendrilResultBody { + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(rename = "callId")] + pub call_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub output: Option, +} + +/// `tendril-reach` — task subagent dispatch. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct TendrilReachBody { + #[serde(skip_serializing_if = "Option::is_none")] + pub task: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub tools: Option>, +} + +/// `petal-cell` — OC message-graph update. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct PetalCellBody { + #[serde(skip_serializing_if = "Option::is_none")] + pub cell: Option, +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + #[test] + fn echo_shapes_canonically() { + let b = EchoBody { + ok: true, + error: None, + }; + assert_eq!(serde_json::to_value(&b).unwrap(), json!({ "ok": true })); + } + + #[test] + fn tool_result_wire_names() { + let b = ToolResultBody { + call_id: "c1".into(), + output: Some("done".into()), + result: None, + is_error: Some(false), + title: None, + metadata: None, + }; + let v = serde_json::to_value(&b).unwrap(); + assert_eq!(v["callId"], "c1"); + assert_eq!(v["output"], "done"); + assert_eq!(v["isError"], false); + assert!(!v.as_object().unwrap().contains_key("result")); + } + + #[test] + fn chunk_wire_names() { + let b = ChunkBody { + chunk_type: "text_delta".into(), + block_idx: None, + delta: Some(json!("hi")), + partial_json: None, + }; + let v = serde_json::to_value(&b).unwrap(); + assert_eq!(v["chunkType"], "text_delta"); + assert_eq!(v["delta"], "hi"); + assert!(v.as_object().unwrap().contains_key("delta")); + } +} diff --git a/thrum-core/tests/fixtures/golden.ndjson b/thrum-core/tests/fixtures/golden.ndjson new file mode 100644 index 00000000..f231c65a --- /dev/null +++ b/thrum-core/tests/fixtures/golden.ndjson @@ -0,0 +1,15 @@ +{"bee":"claude-cli","chi":"hello","from":"humd","hid":"h-9","protoVersion":"0.7.0","provides":["env","token"],"rid":"0000000000000000000000000000000000000000000000000000","sentAt":1700000000000,"sid":"sess-1","sigil":"a1b2c3d4e5f6","source":"cli","tool_names":["Read","Bash"],"version":"0.32.0","wane":0} +{"chi":"prompt","cwd":"/tmp/w","disallowedTools":[],"foragerTools":[{"name":"Brave"}],"from":"humd","modelId":"claude-3-5-sonnet","provided":["env"],"rid":"0000000000000000000000000000000000000000000000000001","sentAt":1700000000001,"sid":"sess-1","sigil":"a1b2c3d4e5f6","systemPrompt":"be brief","text":"hi","tools":["builtin"],"wane":1} +{"blockIdx":0,"chi":"chunk","chunkType":"text_delta","delta":"Hello","from":"humd","rid":"0000000000000000000000000000000000000000000000000002","sentAt":1700000000002,"sid":"sess-1","sigil":"a1b2c3d4e5f6","wane":2} +{"blockIdx":2,"chi":"chunk","chunkType":"tool_input_delta","from":"humd","partialJson":"{\"a\":","rid":"0000000000000000000000000000000000000000000000000003","sentAt":1700000000003,"sid":"sess-1","sigil":"a1b2c3d4e5f6","wane":3} +{"chi":"finish","finishReason":"stop","from":"humd","rid":"0000000000000000000000000000000000000000000000000004","sentAt":1700000000004,"sid":"sess-1","sigil":"a1b2c3d4e5f6","usage":{"input_tokens":12,"output_tokens":34},"wane":4} +{"chi":"error","code":"worker_error","from":"humd","message":"nest crashed","rid":"0000000000000000000000000000000000000000000000000005","sentAt":1700000000005,"sid":"sess-1","sigil":"a1b2c3d4e5f6","subtype":"nest_crash","wane":5} +{"args":{"path":"/x"},"callId":"call-1","chi":"tool-call","from":"humd","name":"Read","rid":"0000000000000000000000000000000000000000000000000006","sentAt":1700000000006,"sid":"sess-1","sigil":"a1b2c3d4e5f6","toolName":"Read","wane":6} +{"callId":"call-1","chi":"tool-result","from":"humd","isError":false,"output":"file contents","rid":"0000000000000000000000000000000000000000000000000007","sentAt":1700000000007,"sid":"sess-1","sigil":"a1b2c3d4e5f6","title":"Read ok","wane":7} +{"chi":"session-ready","from":"humd","model":"claude-3-5-sonnet","nestId":"n-1","rid":"0000000000000000000000000000000000000000000000000008","sentAt":1700000000008,"sid":"sess-1","sigil":"a1b2c3d4e5f6","tools":[{"name":"Read"}],"wane":8} +{"chi":"breath","from":"humd","protoVersion":"0.7.0","rid":"0000000000000000000000000000000000000000000000000009","sentAt":1700000000009,"sessions":[{"model":"claude","sid":"s1"}],"sid":"sess-1","sigil":"a1b2c3d4e5f6","wane":9} +{"chi":"attach","from":"humd","hearOnly":true,"rid":"0000000000000000000000000000000000000000000000000010","sentAt":1700000000010,"sid":"sess-1","sigil":"a1b2c3d4e5f6","wane":10} +{"chi":"wane-sync","from":"humd","rid":"0000000000000000000000000000000000000000000000000011","sentAt":1700000000011,"sid":"sess-1","sigil":"a1b2c3d4e5f6","snapshot":{"a1b2c3d4e5f6":12},"wane":11} +{"chi":"gossip-publish","from":"humd","msg_id":"m-1","payload":{"k":"v"},"rid":"0000000000000000000000000000000000000000000000000012","sentAt":1700000000012,"sid":"sess-1","sigil":"a1b2c3d4e5f6","topic":"gossip","wane":12} +{"chi":"peer-add","from":"humd","hints":[{"addr":"ws://x"}],"humd_id":"humd-1","rid":"0000000000000000000000000000000000000000000000000013","sentAt":1700000000013,"sid":"sess-1","sigil":"a1b2c3d4e5f6","wane":13} +{"chi":"echo","from":"humd","ok":true,"rid":"0000000000000000000000000000000000000000000000000014","sentAt":1700000000014,"sid":"sess-1","sigil":"a1b2c3d4e5f6","wane":14} diff --git a/thrum-core/tests/golden.rs b/thrum-core/tests/golden.rs new file mode 100644 index 00000000..d0d37b29 --- /dev/null +++ b/thrum-core/tests/golden.rs @@ -0,0 +1,278 @@ +//! Golden wire fixture — locks the serialized shape of canonical tones. +//! +//! Every line of `fixtures/golden.ndjson` is a full tone parsed straight +//! off the wire. The test rebuilds each tone from the typed views and +//! deep-compares bytes with what's committed. If `views.rs` changes a +//! key, adds/removes a field, or alters optionality, this test fails +//! until the fixture — and thereby every generated client — catches up. +//! +//! The fixture is written on first run; after that it's compare-only. + +use serde_json::{Map, Value}; + +use thrum_core::{Chi, EchoBody, Envelope, Tone}; + +/// Body view → flat body map (serializes absent optionals off the wire). +fn body_map(b: &T) -> Map { + serde_json::to_value(b) + .expect("body serializes") + .as_object() + .expect("body is an object") + .clone() +} + +/// Fixed valid 52-char Crockford rid: 50 zeroes + the line index. +fn rid(n: usize) -> String { + format!("{:0>50}{:02}", "", n) +} + +fn tone(env: Envelope, b: Map) -> Value { + serde_json::to_value(Tone::with_body(env, b)).expect("tone serializes") +} + +fn env(chi: Chi, n: usize) -> Envelope { + Envelope { + rid: rid(n), + from: Some("humd".into()), + sigil: Some("a1b2c3d4e5f6".into()), + sid: Some("sess-1".into()), + wane: Some(n as u64), + sent_at: Some(1_700_000_000_000 + n as i64), + ..Envelope::new(chi, rid(n)) + } +} + +fn golden() -> Vec { + let mut out = Vec::new(); + let mut n = 0usize; + macro_rules! push { + ($chi:expr, $body:expr) => {{ + out.push(tone(env($chi, n), body_map(&$body))); + n += 1; + }}; + } + + push!( + Chi::Hello, + thrum_core::HelloBody { + proto_version: "0.7.0".into(), + bee: "claude-cli".into(), + version: "0.32.0".into(), + hid: Some("h-9".into()), + tool_names: Some(vec!["Read".into(), "Bash".into()]), + provides: Some(vec!["env".into(), "token".into()]), + source: Some("cli".into()), + } + ); + push!( + Chi::Prompt, + thrum_core::PromptBody { + model_id: Some("claude-3-5-sonnet".into()), + cwd: Some("/tmp/w".into()), + system_prompt: Some("be brief".into()), + text: Some("hi".into()), + content: None, + tools: Some(vec![Value::from("builtin")]), + forager_tools: Some(vec![serde_json::json!({ "name": "Brave" })]), + provided: Some(vec!["env".into()]), + disallowed_tools: Some(vec![]), + } + ); + push!( + Chi::Chunk, + thrum_core::ChunkBody { + chunk_type: "text_delta".into(), + block_idx: Some(0), + delta: Some(Value::from("Hello")), + partial_json: None, + } + ); + push!( + Chi::Chunk, + thrum_core::ChunkBody { + chunk_type: "tool_input_delta".into(), + block_idx: Some(2), + delta: None, + partial_json: Some("{\"a\":".into()), + } + ); + push!( + Chi::Finish, + thrum_core::FinishBody { + finish_reason: "stop".into(), + usage: Some(serde_json::json!({ + "input_tokens": 12, + "output_tokens": 34 + })), + exit_code: None, + subtype: None, + } + ); + push!( + Chi::Error, + thrum_core::ErrorBody { + message: "nest crashed".into(), + code: Some("worker_error".into()), + subtype: Some("nest_crash".into()), + usage: None, + } + ); + push!( + Chi::ToolCall, + thrum_core::ToolCallBody { + call_id: "call-1".into(), + tool_name: "Read".into(), + name: Some("Read".into()), + args: Some(serde_json::json!({ "path": "/x" })), + } + ); + push!( + Chi::ToolResult, + thrum_core::ToolResultBody { + call_id: "call-1".into(), + output: Some("file contents".into()), + result: None, + is_error: Some(false), + title: Some("Read ok".into()), + metadata: None, + } + ); + push!( + Chi::SessionReady, + thrum_core::SessionReadyBody { + nest_id: "n-1".into(), + model: "claude-3-5-sonnet".into(), + tools: Some(vec![serde_json::json!({ "name": "Read" })]), + } + ); + push!( + Chi::Breath, + thrum_core::BreathBody { + sessions: serde_json::json!([{ "sid": "s1", "model": "claude" }]), + proto_version: Some("0.7.0".into()), + } + ); + push!( + Chi::Attach, + thrum_core::AttachBody { + hear_only: Some(true), + } + ); + push!( + Chi::WaneSync, + thrum_core::WaneSyncBody { + snapshot: [("a1b2c3d4e5f6".to_string(), 12u64)].into_iter().collect(), + } + ); + push!( + Chi::GossipPublish, + thrum_core::GossipPublishBody { + topic: "gossip".into(), + payload: serde_json::json!({ "k": "v" }), + from: Some("gossip-0".into()), + msg_id: "m-1".into(), + } + ); + push!( + Chi::PeerAdd, + thrum_core::PeerAddBody { + humd_id: "humd-1".into(), + hints: Some(vec![serde_json::json!({ "addr": "ws://x" })]), + } + ); + push!( + Chi::Echo, + EchoBody { + ok: true, + error: None, + } + ); + out +} + +#[test] +fn golden_wire_locked() { + let fixture = + std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/golden.ndjson"); + let want = golden(); + if !fixture.exists() { + // First run: materialize the fixture, then pass. + std::fs::create_dir_all(fixture.parent().unwrap()).unwrap(); + let mut s = String::new(); + for t in &want { + s.push_str(&serde_json::to_string(t).unwrap()); + s.push('\n'); + } + std::fs::write(&fixture, s).unwrap(); + eprintln!("golden fixture written: {}", fixture.display()); + return; + } + let have: Vec = std::fs::read_to_string(&fixture) + .expect("read fixture") + .lines() + .map(|l| serde_json::from_str(l).expect("fixture line parses")) + .collect(); + assert_eq!(have.len(), want.len(), "fixture/code line count drift"); + for (i, (h, w)) in have.iter().zip(want.iter()).enumerate() { + assert_eq!( + h, w, + "\nline {i} drifted — regen with `rm tests/fixtures/golden.ndjson && cargo test -p thrum-core`\n fixture: {}\n current: {}", + h, w + ); + } +} + +/// Every tone passes through the wire serde round-trip: serialize to a +/// flat frame, parse back into Envelope + raw body, decode the body into +/// its typed view, and confirm envelope fields survived. +#[test] +fn golden_round_trips_through_views() { + for (i, wire) in golden().iter().enumerate() { + let tone: Tone = serde_json::from_value(wire.clone()) + .unwrap_or_else(|e| panic!("line {i} fails as Tone: {e}")); + check_round_trip(&tone); + } +} + +fn check_round_trip(tone: &Tone) { + use thrum_core::*; + // Re-serialize body map from the parsed tone and decode into the view. + let body = &tone.body; + // spot-check a representative set of shapes with the envelope merged + match tone.chi() { + Chi::Chunk => { + let v: ChunkBody = + serde_json::from_value(Value::clone(&Value::Object(body.clone()))).unwrap(); + assert!( + v.chunk_type == "text_delta" || v.chunk_type == "tool_input_delta", + "unexpected chunk_type {}", + v.chunk_type + ); + if let Some(b) = v.block_idx { + assert!(b == 0 || b == 2); + } + } + Chi::Hello => { + let v: HelloBody = serde_json::from_value(Value::Object(body.clone())).unwrap(); + assert_eq!(v.proto_version, "0.7.0"); + assert_eq!(v.bee, "claude-cli"); + } + Chi::ToolResult => { + let v: ToolResultBody = serde_json::from_value(Value::Object(body.clone())).unwrap(); + assert_eq!(v.call_id, "call-1"); + assert_eq!(v.output.as_deref(), Some("file contents")); + assert_eq!(v.is_error, Some(false)); + } + Chi::WaneSync => { + let v: WaneSyncBody = serde_json::from_value(Value::Object(body.clone())).unwrap(); + assert_eq!(v.snapshot.get("a1b2c3d4e5f6"), Some(&12)); + } + Chi::GossipPublish => { + let v: GossipPublishBody = serde_json::from_value(Value::Object(body.clone())).unwrap(); + assert_eq!(v.msg_id, "m-1"); + } + _ => {} + } + // Envelope fields came back intact. + assert!(tone.envelope.sent_at.is_some()); +}