curator — full byte-to-byte

Every file. Nothing imported from any prior repo. Compiles with cargo build --workspace --release.

---

Workspace

Cargo.toml

```toml
[workspace]
resolver = "2"
members = [
  "crates/cura-protocol",
  "crates/cura-crypto",
  "crates/cura-sheaf",
  "crates/cura-registry",
  "crates/cura-worker",
  "crates/curator",
  "crates/curator-cli",
]

[workspace.package]
version = "0.1.0"
edition = "2021"
rust-version = "1.83"
license = "UNLICENSED"
authors = ["Comfort Curators Private Limited"]
repository = "https://github.com/yashrajvansh/curator"

[workspace.dependencies]
serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"
thiserror = "1.0"
blake3 = "1.5"
sha2 = "0.10"
hmac = "0.12"
hex = "0.4"
uuid = { version = "1.10", features = ["v7", "serde"] }
anyhow = "1.0"
walkdir = "2.5"
clap = { version = "4.5", features = ["derive", "env"] }
tokio = { version = "1.40", features = ["macros", "rt-multi-thread"] }
reqwest = { version = "0.12", default-features = false, features = ["rustls-tls", "json"] }
toml = "0.8"

[profile.release]
opt-level = "s"
lto = true
codegen-units = 1
strip = true
```

rust-toolchain.toml

```toml
[toolchain]
channel = "1.83.0"
targets = ["wasm32-unknown-unknown"]
components = ["rustfmt", "clippy"]
profile = "minimal"
```

.gitignore

```
target/
_build/
deps/
.wrangler/
.dev.vars
*.wasm
.env
.env.local
.DS_Store
```

---

1. crates/cura-protocol

Cargo.toml

```toml
[package]
name = "cura-protocol"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
publish = false

[dependencies]
serde.workspace = true
blake3.workspace = true
```

src/lib.rs

```rust
//! YA|RA wire types. input · logic · output maps to intent · weave · pattern.

pub mod origin;
pub mod plane;
pub mod stance;
pub mod stalk;
pub mod triple;
pub mod weave;

pub use origin::{Origin, Span};
pub use plane::{ObservationKind, Plane};
pub use stance::Stance;
pub use stalk::{Provenance, Stalk};
pub use triple::{Term, Triple, WeaveId};
pub use weave::{Obstruction, Section, WeaveLink};
```

src/plane.rs

```rust
use serde::{Deserialize, Serialize};

/// The three planes. Each D1, each R2 bucket, one per plane.
///
///   Intent  →  input  — DI  — R2 RI
///   Weave   →  logic  — DIP — R2 RIP
///   Pattern →  output — DP  — R2 RP
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum Plane {
    Intent,
    Weave,
    Pattern,
}

impl Plane {
    pub fn all() -> [Plane; 3] {
        [Plane::Intent, Plane::Weave, Plane::Pattern]
    }

    pub fn as_str(&self) -> &'static str {
        match self {
            Plane::Intent => "intent",
            Plane::Weave => "weave",
            Plane::Pattern => "pattern",
        }
    }

    pub fn r2_bucket(&self) -> &'static str {
        match self {
            Plane::Intent => "cura-ri",
            Plane::Weave => "cura-rip",
            Plane::Pattern => "cura-rp",
        }
    }

    pub fn d1_binding(&self) -> &'static str {
        match self {
            Plane::Intent => "DB_DI",
            Plane::Weave => "DB_DIP",
            Plane::Pattern => "DB_DP",
        }
    }

    pub fn d1_database(&self) -> &'static str {
        match self {
            Plane::Intent => "cura-di",
            Plane::Weave => "cura-dip",
            Plane::Pattern => "cura-dp",
        }
    }
}

#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ObservationKind {
    Lexical,
    Semantic,
    FirmwareSoftware,
    HardwareSoftware,
    Relay,
    Introspective,
    Service,
    Program,
    Sheath,
}

impl ObservationKind {
    pub fn as_str(&self) -> &'static str {
        match self {
            ObservationKind::Lexical => "lexical",
            ObservationKind::Semantic => "semantic",
            ObservationKind::FirmwareSoftware => "firmware_software",
            ObservationKind::HardwareSoftware => "hardware_software",
            ObservationKind::Relay => "relay",
            ObservationKind::Introspective => "introspective",
            ObservationKind::Service => "service",
            ObservationKind::Program => "program",
            ObservationKind::Sheath => "sheath",
        }
    }
}
```

src/triple.rs

```rust
use serde::{Deserialize, Serialize};

#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq, Hash)]
pub struct Term(pub String);

impl Term {
    pub fn new(s: impl Into<String>) -> Self { Term(s.into()) }
    pub fn is_empty(&self) -> bool { self.0.is_empty() }
    pub fn as_bytes(&self) -> &[u8] { self.0.as_bytes() }
}

/// The atomic unit. The mapping is fixed:
///
///   input  → intent
///   logic  → weave
///   output → pattern
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq, Hash)]
pub struct Triple {
    pub input: Term,
    pub logic: Term,
    pub output: Term,
}

impl Triple {
    pub fn new(i: impl Into<String>, l: impl Into<String>, o: impl Into<String>) -> Self {
        Triple { input: Term::new(i), logic: Term::new(l), output: Term::new(o) }
    }

    pub fn is_well_formed(&self) -> bool {
        !(self.input.is_empty() && self.logic.is_empty() && self.output.is_empty())
    }

    /// BLAKE3 over length-prefixed terms. Same value across languages.
    pub fn weave_id(&self) -> WeaveId {
        let mut h = blake3::Hasher::new();
        h.update(b"YA|RA|v1|");
        write_len(&mut h, self.input.as_bytes());
        write_len(&mut h, self.logic.as_bytes());
        write_len(&mut h, self.output.as_bytes());
        WeaveId(*h.finalize().as_bytes())
    }
}

fn write_len(h: &mut blake3::Hasher, b: &[u8]) {
    h.update(&(b.len() as u64).to_be_bytes());
    h.update(b);
}

#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct WeaveId(pub [u8; 32]);

impl WeaveId {
    pub fn to_hex(&self) -> String {
        const H: &[u8; 16] = b"0123456789abcdef";
        let mut s = String::with_capacity(64);
        for b in &self.0 {
            s.push(H[(b >> 4) as usize] as char);
            s.push(H[(b & 0x0f) as usize] as char);
        }
        s
    }

    pub fn short(&self) -> String {
        let h = self.to_hex();
        format!("{}…{}", &h[..12], &h[h.len() - 8..])
    }

    pub fn shard(&self) -> (String, String) {
        let h = self.to_hex();
        (h[0..2].to_string(), h[2..4].to_string())
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn stable() {
        let t = Triple::new("a", "b", "c");
        assert_eq!(t.weave_id(), t.weave_id());
    }

    #[test]
    fn distinct() {
        assert_ne!(
            Triple::new("a", "b", "c").weave_id(),
            Triple::new("a", "b", "d").weave_id()
        );
    }

    #[test]
    fn no_length_collision() {
        assert_ne!(
            Triple::new("ab", "c", "").weave_id(),
            Triple::new("a", "bc", "").weave_id()
        );
    }

    #[test]
    fn hex_is_64() {
        assert_eq!(Triple::new("x", "y", "z").weave_id().to_hex().len(), 64);
    }
}
```

src/stance.rs

```rust
use serde::{Deserialize, Serialize};

/// A program's report about its own attempt to see. Not a verdict.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Stance {
    Wove,
    Refused,
    Malformed,
    OutOfScope,
}

impl Stance {
    pub fn as_str(&self) -> &'static str {
        match self {
            Stance::Wove => "wove",
            Stance::Refused => "refused",
            Stance::Malformed => "malformed",
            Stance::OutOfScope => "out_of_scope",
        }
    }

    pub fn from_i32(n: i32) -> Self {
        match n {
            0 => Stance::Wove,
            -1 => Stance::Refused,
            -2 => Stance::Malformed,
            -3 => Stance::OutOfScope,
            _ => Stance::Malformed,
        }
    }

    pub fn to_i32(self) -> i32 {
        match self {
            Stance::Wove => 0,
            Stance::Refused => -1,
            Stance::Malformed => -2,
            Stance::OutOfScope => -3,
        }
    }
}
```

src/origin.rs

```rust
use serde::{Deserialize, Serialize};

#[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq)]
pub struct Span {
    pub start: u32,
    pub end: u32,
}

#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
pub struct Origin {
    pub repo: String,
    pub path: String,
    pub revision: String,
    pub translator: String,
    pub lang: String,
    #[serde(skip_serializing_if = "Option::is_none", default)]
    pub span: Option<Span>,
}

impl Origin {
    pub fn new(
        repo: impl Into<String>,
        path: impl Into<String>,
        revision: impl Into<String>,
        translator: impl Into<String>,
        lang: impl Into<String>,
    ) -> Self {
        Origin {
            repo: repo.into(),
            path: path.into(),
            revision: revision.into(),
            translator: translator.into(),
            lang: lang.into(),
            span: None,
        }
    }

    pub fn with_span(mut self, start: u32, end: u32) -> Self {
        self.span = Some(Span { start, end });
        self
    }

    pub fn canonical(&self) -> String {
        let span = match self.span {
            Some(s) => format!("{}:{}", s.start, s.end),
            None => String::new(),
        };
        format!(
            "{}|{}|{}|{}|{}|{}",
            self.repo, self.path, self.revision, self.translator, self.lang, span
        )
    }
}
```

src/stalk.rs

```rust
use crate::origin::Origin;
use crate::stance::Stance;
use serde::{Deserialize, Serialize};

#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Provenance {
    pub source: String,
    #[serde(skip_serializing_if = "Option::is_none", default)]
    pub body_key: Option<String>,
    #[serde(skip_serializing_if = "Option::is_none", default)]
    pub stance: Option<Stance>,
    #[serde(skip_serializing_if = "Option::is_none", default)]
    pub program_id: Option<String>,
    #[serde(skip_serializing_if = "Option::is_none", default)]
    pub fuel_consumed: Option<u64>,
    #[serde(skip_serializing_if = "Option::is_none", default)]
    pub peak_memory_bytes: Option<usize>,
    #[serde(skip_serializing_if = "Option::is_none", default)]
    pub origin: Option<Origin>,
}

impl Provenance {
    pub fn new(source: impl Into<String>) -> Self {
        Provenance {
            source: source.into(),
            body_key: None,
            stance: None,
            program_id: None,
            fuel_consumed: None,
            peak_memory_bytes: None,
            origin: None,
        }
    }
}

#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Stalk {
    pub vector: Vec<f32>,
    pub provenance: Provenance,
    pub observed_at_ms: i64,
    /// The cage: how long Curator looked without looking elsewhere.
    pub duration_ms: u64,
}

impl Stalk {
    pub fn byte_size(&self) -> usize { self.vector.len() * 4 }
}
```

src/weave.rs

```rust
use crate::triple::WeaveId;
use serde::{Deserialize, Serialize};

#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct WeaveLink {
    pub from: WeaveId,
    pub to: WeaveId,
    pub strength: f64,
    pub confirmations: u64,
    pub denials: u64,
    pub first_seen_ms: i64,
    pub last_confirmed_ms: i64,
}

#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Section {
    pub support: Vec<WeaveId>,
    pub vector: Vec<f32>,
    pub consistent: bool,
    pub obstruction: Option<Obstruction>,
}

#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Obstruction {
    DurationExceeded,
    IncompatibleGlue,
    Empty,
}
```

---

2. crates/cura-crypto

Cargo.toml

```toml
[package]
name = "cura-crypto"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
publish = false

[dependencies]
hmac.workspace = true
sha2.workspace = true
blake3.workspace = true
hex.workspace = true
```

src/lib.rs

```rust
//! HMAC-SHA-256 and SHA-256, matching the Elixir orchestrator byte-for-byte.

use hmac::{Hmac, Mac};
use sha2::{Digest, Sha256};

type HmacSha256 = Hmac<Sha256>;

pub fn sign(payload: &[u8], secret: &[u8]) -> String {
    let mut mac = HmacSha256::new_from_slice(secret).expect("hmac accepts any key length");
    mac.update(payload);
    hex::encode(mac.finalize().into_bytes())
}

pub fn verify(payload: &[u8], signature_hex: &str, secret: &[u8]) -> bool {
    let Ok(expected) = hex::decode(signature_hex) else { return false };
    let mut mac = HmacSha256::new_from_slice(secret).expect("hmac accepts any key length");
    mac.update(payload);
    mac.verify_slice(&expected).is_ok()
}

pub fn sha256_hex(bytes: &[u8]) -> String {
    let mut h = Sha256::new();
    h.update(bytes);
    hex::encode(h.finalize())
}

pub fn blake3_hex(bytes: &[u8]) -> String {
    blake3::hash(bytes).to_hex().to_string()
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn sign_shape() {
        let s = sign(b"hello", b"secret");
        assert_eq!(s.len(), 64);
        assert!(s.chars().all(|c| c.is_ascii_hexdigit() && !c.is_ascii_uppercase()));
    }

    #[test]
    fn verify_ok() {
        let s = sign(b"p", b"k");
        assert!(verify(b"p", &s, b"k"));
    }

    #[test]
    fn verify_rejects() {
        let s = sign(b"p", b"k");
        assert!(!verify(b"q", &s, b"k"));
    }

    #[test]
    fn rfc_4231_case_2() {
        assert_eq!(
            sign(b"what do ya want for nothing?", b"Jefe"),
            "5bdcc146bf60754e6a042426089575c75a003f089d2739839dec58b964ec3843"
        );
    }

    #[test]
    fn sha256_empty() {
        assert_eq!(
            sha256_hex(b""),
            "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
        );
    }
}
```

---

3. crates/cura-sheaf

Cargo.toml

```toml
[package]
name = "cura-sheaf"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
publish = false

[dependencies]
cura-protocol = { path = "../cura-protocol" }
thiserror.workspace = true
```

src/lib.rs

```rust
//! A sheaf is a switchable encoder. Every kind of seeing returns the
//! same shape. The default is the sheath: a lexical encoder.

pub mod lexical;
pub mod sheath;

pub use lexical::LexicalSheaf;
pub use sheath::SheathSheaf;

use cura_protocol::*;
use std::time::Duration;

pub type SheafResult<T> = Result<T, SheafError>;

#[derive(Debug, thiserror::Error)]
pub enum SheafError {
    #[error("encoder unavailable: {0}")] EncoderUnavailable(String),
    #[error("restriction failed: {0}")]  RestrictionFailed(String),
    #[error("glue failed: {0}")]         GlueFailed(String),
    #[error("duration exceeded")]        DurationExceeded,
    #[error("empty input")]              Empty,
}

pub trait Sheaf: Send + Sync {
    fn kind(&self) -> ObservationKind;

    /// Local data at a point. `budget` is the cage.
    fn stalk(&self, t: &Triple, budget: Duration, origin: Option<Origin>) -> SheafResult<Stalk>;

    /// How a stalk transforms along a weave.
    fn restrict(&self, from: &Stalk, to: &Stalk) -> SheafResult<Stalk>;

    /// `new = old + inertia * (current - old)`.
    fn remeasure(&self, link: &WeaveLink, current: f64) -> SheafResult<WeaveLink>;

    /// Glue stalks into a section.
    fn glue(&self, points: &[(WeaveId, Stalk)]) -> SheafResult<Section>;

    /// When gluing cannot succeed.
    fn obstruct(&self, points: &[(WeaveId, Stalk)]) -> Option<Obstruction>;
}
```

src/lexical.rs

```rust
use crate::*;
use cura_protocol::*;
use std::time::Duration;

const DIMS: usize = 64;
const FNV_OFFSET: u32 = 0x811c9dc5;
const FNV_PRIME: u32 = 0x01000193;
const INERTIA: f64 = 0.35;
const DISSOLVE: f64 = 0.24;

pub struct LexicalSheaf {
    dims: usize,
    default_budget: Duration,
}

impl Default for LexicalSheaf {
    fn default() -> Self {
        Self { dims: DIMS, default_budget: Duration::from_secs(88) }
    }
}

impl LexicalSheaf {
    fn fnv1a(bytes: &[u8], seed: u32) -> u32 {
        let mut h = seed;
        for b in bytes {
            h ^= *b as u32;
            h = h.wrapping_mul(FNV_PRIME);
        }
        h
    }

    fn features(text: &str) -> Vec<(String, f64)> {
        let lower = text.to_lowercase();
        let words: Vec<&str> = lower
            .split(|c: char| !c.is_ascii_alphanumeric())
            .filter(|s| !s.is_empty())
            .collect();

        let mut c: std::collections::HashMap<String, u32> = Default::default();
        for w in &words {
            *c.entry((*w).to_string()).or_insert(0) += 1;
        }
        for p in words.windows(2) {
            *c.entry(format!("{} {}", p[0], p[1])).or_insert(0) += 1;
        }
        for w in &words {
            let ch: Vec<char> = w.chars().collect();
            if ch.len() > 4 {
                for i in 0..ch.len() - 2 {
                    *c.entry(ch[i..i + 3].iter().collect()).or_insert(0) += 1;
                }
            }
        }

        c.into_iter().map(|(k, n)| (k, 1.0 + (n as f64).ln())).collect()
    }

    fn embed(&self, text: &str) -> Vec<f32> {
        let mut v = vec![0.0f32; self.dims];
        for (f, w) in Self::features(text) {
            let h = Self::fnv1a(f.as_bytes(), FNV_OFFSET);
            let idx = (h as usize) % self.dims;
            let sign = if (h >> 31) & 1 == 0 { 1.0 } else { -1.0 };
            v[idx] += sign * w as f32;
        }
        let norm: f32 = v.iter().map(|x| x * x).sum::<f32>().sqrt();
        if norm > 0.0 {
            for x in v.iter_mut() {
                *x /= norm;
            }
        }
        v
    }
}

impl Sheaf for LexicalSheaf {
    fn kind(&self) -> ObservationKind {
        ObservationKind::Lexical
    }

    fn stalk(&self, t: &Triple, budget: Duration, origin: Option<Origin>) -> SheafResult<Stalk> {
        if !t.is_well_formed() {
            return Err(SheafError::Empty);
        }
        let mut v = Vec::with_capacity(self.dims * 3);
        v.extend(self.embed(&t.input.0));
        v.extend(self.embed(&t.logic.0));
        v.extend(self.embed(&t.output.0));
        let n: f32 = v.iter().map(|x| x * x).sum::<f32>().sqrt();
        if n > 0.0 {
            for x in v.iter_mut() {
                *x /= n;
            }
        }

        let mut prov = Provenance::new("lexical");
        prov.origin = origin;

        Ok(Stalk {
            vector: v,
            provenance: prov,
            observed_at_ms: now_ms(),
            duration_ms: budget.as_millis() as u64,
        })
    }

    fn restrict(&self, a: &Stalk, b: &Stalk) -> SheafResult<Stalk> {
        if a.vector.len() != b.vector.len() {
            return Err(SheafError::RestrictionFailed("dim mismatch".into()));
        }
        let v: Vec<f32> = a
            .vector
            .iter()
            .zip(&b.vector)
            .map(|(x, y)| 0.5 * x + 0.5 * y)
            .collect();
        Ok(Stalk {
            vector: v,
            provenance: a.provenance.clone(),
            observed_at_ms: now_ms(),
            duration_ms: a.duration_ms.min(b.duration_ms),
        })
    }

    fn remeasure(&self, l: &WeaveLink, current: f64) -> SheafResult<WeaveLink> {
        let new = l.strength + INERTIA * (current - l.strength);
        Ok(WeaveLink {
            from: l.from,
            to: l.to,
            strength: new,
            confirmations: if current > l.strength { l.confirmations + 1 } else { l.confirmations },
            denials: if new < DISSOLVE { l.denials + 1 } else { l.denials },
            first_seen_ms: l.first_seen_ms,
            last_confirmed_ms: now_ms(),
        })
    }

    fn glue(&self, points: &[(WeaveId, Stalk)]) -> SheafResult<Section> {
        if points.is_empty() {
            return Err(SheafError::GlueFailed("empty".into()));
        }
        let n = points[0].1.vector.len();
        let mut mean = vec![0.0f32; n];
        for (_, s) in points {
            if s.vector.len() != n {
                return Err(SheafError::GlueFailed("dim".into()));
            }
            for i in 0..n {
                mean[i] += s.vector[i];
            }
        }
        for x in mean.iter_mut() {
            *x /= points.len() as f32;
        }
        Ok(Section {
            support: points.iter().map(|(id, _)| *id).collect(),
            vector: mean,
            consistent: true,
            obstruction: None,
        })
    }

    fn obstruct(&self, _: &[(WeaveId, Stalk)]) -> Option<Obstruction> {
        None
    }
}

pub fn now_ms() -> i64 {
    std::time::SystemTime::now()
        .duration_since(std::time::UNIX_EPOCH)
        .map(|d| d.as_millis() as i64)
        .unwrap_or(0)
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn stalk_shape() {
        let s = LexicalSheaf::default();
        let t = Triple::new("a b c", "d e", "f g h");
        let st = s.stalk(&t, Duration::from_secs(88), None).unwrap();
        assert_eq!(st.vector.len(), 192);
        let n: f32 = st.vector.iter().map(|x| x * x).sum::<f32>().sqrt();
        assert!((n - 1.0).abs() < 1e-4);
    }

    #[test]
    fn deterministic() {
        let s = LexicalSheaf::default();
        let t = Triple::new("a", "b", "c");
        let a = s.stalk(&t, Duration::from_secs(1), None).unwrap();
        let b = s.stalk(&t, Duration::from_secs(1), None).unwrap();
        assert_eq!(a.vector, b.vector);
    }

    #[test]
    fn empty_refused() {
        let s = LexicalSheaf::default();
        let t = Triple::new("", "", "");
        assert!(matches!(s.stalk(&t, Duration::from_secs(1), None), Err(SheafError::Empty)));
    }
}
```

src/sheath.rs

```rust
//! The sheath is the default sheaf. Same lexical geometry, named
//! for the fact that it wraps every observation without changing it.

use crate::lexical::LexicalSheaf;
use crate::*;
use cura_protocol::*;
use std::time::Duration;

pub struct SheathSheaf {
    inner: LexicalSheaf,
}

impl Default for SheathSheaf {
    fn default() -> Self {
        SheathSheaf { inner: LexicalSheaf::default() }
    }
}

impl Sheaf for SheathSheaf {
    fn kind(&self) -> ObservationKind {
        ObservationKind::Sheath
    }

    fn stalk(&self, t: &Triple, budget: Duration, origin: Option<Origin>) -> SheafResult<Stalk> {
        let mut st = self.inner.stalk(t, budget, origin)?;
        st.provenance.source = "sheath".to_string();
        Ok(st)
    }

    fn restrict(&self, a: &Stalk, b: &Stalk) -> SheafResult<Stalk> {
        self.inner.restrict(a, b)
    }

    fn remeasure(&self, l: &WeaveLink, current: f64) -> SheafResult<WeaveLink> {
        self.inner.remeasure(l, current)
    }

    fn glue(&self, points: &[(WeaveId, Stalk)]) -> SheafResult<Section> {
        self.inner.glue(points)
    }

    fn obstruct(&self, points: &[(WeaveId, Stalk)]) -> Option<Obstruction> {
        self.inner.obstruct(points)
    }
}
```

---

4. crates/cura-registry

Cargo.toml

```toml
[package]
name = "cura-registry"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
publish = false

[dependencies]
cura-protocol = { path = "../cura-protocol" }
cura-sheaf = { path = "../cura-sheaf" }
serde.workspace = true
serde_json.workspace = true
thiserror.workspace = true
```

src/lib.rs

```rust
//! Cura — the registry. Stores, indexes, and serves YA|RA.

pub mod body;
pub mod classify;
pub mod record;
pub mod store;

pub use body::body_key_for;
pub use classify::classify_from_source;
pub use record::{Record, RecordError, RecordStatus};
pub use store::{MemoryStore, Store, StoreError};
```

src/body.rs

```rust
use cura_protocol::{Plane, WeaveId};

/// R2 key layout:
///
///     <plane-bucket-prefix>/<hex[0..2]>/<hex[2..4]>/<hex>.bin
///
/// The prefix is the plane name so a single bucket can hold all planes
/// when the deployment is small; in production three buckets are bound
/// separately and the prefix is redundant but harmless.
pub fn body_key_for(plane: &Plane, id: &WeaveId) -> String {
    let hex = id.to_hex();
    format!("{}/{}/{}/{}.bin", plane.as_str(), &hex[0..2], &hex[2..4], hex)
}

#[cfg(test)]
mod tests {
    use super::*;
    use cura_protocol::Triple;

    #[test]
    fn key_layout() {
        let id = Triple::new("a", "b", "c").weave_id();
        let k = body_key_for(&Plane::Weave, &id);
        let parts: Vec<&str> = k.split('/').collect();
        assert_eq!(parts[0], "weave");
        assert_eq!(parts[1].len(), 2);
        assert_eq!(parts[2].len(), 2);
        assert_eq!(parts[3].len(), 68);
        assert!(parts[3].ends_with(".bin"));
    }
}
```

src/classify.rs

```rust
use cura_protocol::ObservationKind;

/// The default sheaf is the sheath. Every source maps to Sheath unless
/// a specific translator has already declared a kind.
pub fn classify_from_source(source: &str) -> ObservationKind {
    let s = source.to_lowercase();
    if s.starts_with("program:") || s.starts_with("wasm:") {
        return ObservationKind::Program;
    }
    if s.starts_with("self-report") {
        return ObservationKind::Introspective;
    }
    if s.starts_with("founder-message") {
        return ObservationKind::Relay;
    }
    if s.contains("firmware") {
        return ObservationKind::FirmwareSoftware;
    }
    if s.contains("hardware") {
        return ObservationKind::HardwareSoftware;
    }
    if s.starts_with("door:")
        || s.starts_with("evidence:")
        || s.starts_with("synthetic:")
        || s.starts_with("service:")
    {
        return ObservationKind::Service;
    }
    ObservationKind::Sheath
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn plain_is_sheath() {
        assert_eq!(classify_from_source("ordinary"), ObservationKind::Sheath);
    }

    #[test]
    fn program_is_program() {
        assert_eq!(classify_from_source("program:abc"), ObservationKind::Program);
    }

    #[test]
    fn self_report_is_introspective() {
        assert_eq!(classify_from_source("self-report:x"), ObservationKind::Introspective);
    }
}
```

src/record.rs

```rust
use cura_protocol::{Origin, Plane, Stalk, Triple, WeaveId};
use serde::{Deserialize, Serialize};

#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Record {
    pub weave_id: String,
    pub plane: Plane,
    pub triple: Triple,
    pub sheaf_kind: String,
    pub stalk: Option<Stalk>,
    pub body_key: Option<String>,
    pub source: String,
    pub origin: Option<Origin>,
    pub observed_at_ms: i64,
    pub duration_ms: u64,
}

#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum RecordStatus {
    Stored,
    Duplicate,
    Obstruction,
}

#[derive(Debug, thiserror::Error)]
pub enum RecordError {
    #[error("serialize: {0}")]
    Serialize(#[from] serde_json::Error),
    #[error("missing weave_id")]
    MissingId,
}

impl Record {
    pub fn new(
        plane: Plane,
        triple: Triple,
        sheaf_kind: impl Into<String>,
        source: impl Into<String>,
        duration_ms: u64,
    ) -> Self {
        let id = triple.weave_id();
        Record {
            weave_id: id.to_hex(),
            plane,
            triple,
            sheaf_kind: sheaf_kind.into(),
            stalk: None,
            body_key: None,
            source: source.into(),
            origin: None,
            observed_at_ms: 0,
            duration_ms,
        }
    }

    pub fn with_origin(mut self, origin: Origin) -> Self {
        self.origin = Some(origin);
        self
    }

    pub fn with_stalk(mut self, stalk: Stalk) -> Self {
        self.stalk = Some(stalk);
        self
    }

    pub fn with_body_key(mut self, key: impl Into<String>) -> Self {
        self.body_key = Some(key.into());
        self
    }

    pub fn id(&self) -> WeaveId {
        self.triple.weave_id()
    }
}
```

src/store.rs

```rust
use crate::record::{Record, RecordError};
use cura_protocol::Plane;
use std::collections::HashMap;
use std::sync::{Arc, Mutex};

#[derive(Debug, thiserror::Error)]
pub enum StoreError {
    #[error("not found")] NotFound,
    #[error("duplicate")] Duplicate,
    #[error("io: {0}")] Io(String),
    #[error(transparent)] Record(#[from] RecordError),
}

/// Cura's storage surface.
pub trait Store: Send + Sync {
    fn put(&self, record: Record) -> Result<(), StoreError>;
    fn get(&self, plane: Plane, weave_id: &str) -> Result<Record, StoreError>;
    fn list(&self, plane: Plane, limit: usize, before_ms: Option<i64>) -> Result<Vec<Record>, StoreError>;
    fn has(&self, plane: Plane, weave_id: &str) -> Result<bool, StoreError>;
    fn put_body(&self, key: &str, bytes: &[u8]) -> Result<(), StoreError>;
    fn get_body(&self, key: &str) -> Result<Vec<u8>, StoreError>;
}

#[derive(Default, Clone)]
pub struct MemoryStore {
    inner: Arc<Mutex<MemoryInner>>,
}

#[derive(Default)]
struct MemoryInner {
    intent: HashMap<String, Record>,
    weave: HashMap<String, Record>,
    pattern: HashMap<String, Record>,
    bodies: HashMap<String, Vec<u8>>,
}

impl MemoryStore {
    pub fn new() -> Self { Self::default() }

    pub fn len(&self, plane: Plane) -> usize {
        let i = self.inner.lock().unwrap();
        match plane {
            Plane::Intent => i.intent.len(),
            Plane::Weave => i.weave.len(),
            Plane::Pattern => i.pattern.len(),
        }
    }

    pub fn is_empty(&self, plane: Plane) -> bool { self.len(plane) == 0 }

    fn map_mut<'a>(i: &'a mut MemoryInner, p: Plane) -> &'a mut HashMap<String, Record> {
        match p {
            Plane::Intent => &mut i.intent,
            Plane::Weave => &mut i.weave,
            Plane::Pattern => &mut i.pattern,
        }
    }

    fn map<'a>(i: &'a MemoryInner, p: Plane) -> &'a HashMap<String, Record> {
        match p {
            Plane::Intent => &i.intent,
            Plane::Weave => &i.weave,
            Plane::Pattern => &i.pattern,
        }
    }
}

impl Store for MemoryStore {
    fn put(&self, record: Record) -> Result<(), StoreError> {
        let mut i = self.inner.lock().unwrap();
        let m = Self::map_mut(&mut i, record.plane);
        if m.contains_key(&record.weave_id) {
            return Err(StoreError::Duplicate);
        }
        m.insert(record.weave_id.clone(), record);
        Ok(())
    }

    fn get(&self, plane: Plane, weave_id: &str) -> Result<Record, StoreError> {
        let i = self.inner.lock().unwrap();
        Self::map(&i, plane).get(weave_id).cloned().ok_or(StoreError::NotFound)
    }

    fn list(&self, plane: Plane, limit: usize, before_ms: Option<i64>) -> Result<Vec<Record>, StoreError> {
        let i = self.inner.lock().unwrap();
        let mut rows: Vec<Record> = Self::map(&i, plane).values().cloned().collect();
        rows.sort_by(|a, b| b.observed_at_ms.cmp(&a.observed_at_ms));
        let filtered = match before_ms {
            None => rows,
            Some(t) => rows.into_iter().filter(|r| r.observed_at_ms < t).collect(),
        };
        Ok(filtered.into_iter().take(limit).collect())
    }

    fn has(&self, plane: Plane, weave_id: &str) -> Result<bool, StoreError> {
        let i = self.inner.lock().unwrap();
        Ok(Self::map(&i, plane).contains_key(weave_id))
    }

    fn put_body(&self, key: &str, bytes: &[u8]) -> Result<(), StoreError> {
        self.inner.lock().unwrap().bodies.insert(key.to_string(), bytes.to_vec());
        Ok(())
    }

    fn get_body(&self, key: &str) -> Result<Vec<u8>, StoreError> {
        self.inner.lock().unwrap().bodies.get(key).cloned().ok_or(StoreError::NotFound)
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::record::Record;

    #[test]
    fn put_get() {
        let s = MemoryStore::new();
        let r = Record::new(Plane::Weave, cura_protocol::Triple::new("a", "b", "c"), "sheath", "t", 1000);
        let id = r.weave_id.clone();
        s.put(r).unwrap();
        assert!(s.has(Plane::Weave, &id).unwrap());
        assert_eq!(s.get(Plane::Weave, &id).unwrap().weave_id, id);
    }

    #[test]
    fn duplicate_rejected() {
        let s = MemoryStore::new();
        let r = Record::new(Plane::Weave, cura_protocol::Triple::new("a", "b", "c"), "sheath", "t", 1000);
        s.put(r.clone()).unwrap();
        assert!(matches!(s.put(r), Err(StoreError::Duplicate)));
    }

    #[test]
    fn planes_isolated() {
        let s = MemoryStore::new();
        let r = Record::new(Plane::Weave, cura_protocol::Triple::new("a", "b", "c"), "sheath", "t", 1000);
        let id = r.weave_id.clone();
        s.put(r).unwrap();
        assert!(matches!(s.get(Plane::Intent, &id), Err(StoreError::NotFound)));
    }
}
```

---

5. crates/cura-worker

Cargo.toml

```toml
[package]
name = "cura-worker"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
publish = false

[lib]
crate-type = ["cdylib", "rlib"]

[dependencies]
cura-protocol = { path = "../cura-protocol" }
cura-crypto = { path = "../cura-crypto" }
cura-registry = { path = "../cura-registry" }
serde.workspace = true
serde_json.workspace = true
worker = "0.5"
js-sys = "0.3"
console_error_panic_hook = "0.1"
```

wrangler.toml

```toml
name = "cura"
main = "build/worker/shim.mjs"
compatibility_date = "2024-12-01"
compatibility_flags = ["nodejs_compat"]

[build]
command = "cargo install -q worker-build && worker-build --release"

[[r2_buckets]]
binding = "BODIES_RI"
bucket_name = "cura-ri"

[[r2_buckets]]
binding = "BODIES_RIP"
bucket_name = "cura-rip"

[[r2_buckets]]
binding = "BODIES_RP"
bucket_name = "cura-rp"

[[d1_databases]]
binding = "DB_DI"
database_name = "cura-di"
database_id = "REPLACE_DI_ID"

[[d1_databases]]
binding = "DB_DIP"
database_name = "cura-dip"
database_id = "REPLACE_DIP_ID"

[[d1_databases]]
binding = "DB_DP"
database_name = "cura-dp"
database_id = "REPLACE_DP_ID"

[vars]
CURA_SERVICE_ID = "cura"
LOG_LEVEL = "info"

[observability]
enabled = true
```

migrations/intent/0001_init.sql

```sql
CREATE TABLE IF NOT EXISTS intent (
  weave_id    TEXT PRIMARY KEY,
  input       TEXT NOT NULL,
  logic       TEXT NOT NULL,
  output      TEXT NOT NULL DEFAULT '',
  sheaf_kind  TEXT NOT NULL,
  stalk       TEXT NOT NULL DEFAULT '',
  observed_at INTEGER NOT NULL,
  duration_ms INTEGER NOT NULL,
  body_key    TEXT,
  source      TEXT NOT NULL,
  origin      TEXT,
  created_at  TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS intent_observed_idx ON intent (observed_at DESC);
```

migrations/weave/0001_init.sql

```sql
CREATE TABLE IF NOT EXISTS weave (
  weave_id     TEXT PRIMARY KEY,
  input        TEXT NOT NULL,
  logic        TEXT NOT NULL,
  output       TEXT NOT NULL,
  sheaf_kind   TEXT NOT NULL,
  stalk        TEXT NOT NULL DEFAULT '',
  section      TEXT,
  obstruction  TEXT,
  observed_at  INTEGER NOT NULL,
  duration_ms  INTEGER NOT NULL,
  body_key     TEXT,
  source       TEXT NOT NULL,
  origin       TEXT,
  created_at   TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS weave_observed_idx ON weave (observed_at DESC);
CREATE INDEX IF NOT EXISTS weave_obstruction_idx
  ON weave (observed_at DESC) WHERE obstruction IS NOT NULL;

CREATE TABLE IF NOT EXISTS weave_link (
  from_id           TEXT NOT NULL,
  to_id             TEXT NOT NULL,
  strength          REAL NOT NULL DEFAULT 0.0,
  confirmations     INTEGER NOT NULL DEFAULT 0,
  denials           INTEGER NOT NULL DEFAULT 0,
  first_seen_at     INTEGER NOT NULL,
  last_confirmed_at INTEGER NOT NULL,
  PRIMARY KEY (from_id, to_id)
);
CREATE INDEX IF NOT EXISTS weave_link_strength_idx ON weave_link (strength DESC);
```

migrations/pattern/0001_init.sql

```sql
CREATE TABLE IF NOT EXISTS pattern (
  weave_id    TEXT PRIMARY KEY,
  input       TEXT NOT NULL,
  logic       TEXT NOT NULL,
  output      TEXT NOT NULL,
  sheaf_kind  TEXT NOT NULL,
  stalk       TEXT NOT NULL DEFAULT '',
  observed_at INTEGER NOT NULL,
  duration_ms INTEGER NOT NULL,
  body_key    TEXT,
  source      TEXT NOT NULL,
  origin      TEXT,
  created_at  TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS pattern_kind_idx ON pattern (sheaf_kind, observed_at DESC);
```

src/lib.rs

```rust
//! Cura — the registry. Serves YA|RA over HTTP.

use cura_protocol::Plane;
use cura_registry::{body_key_for, classify_from_source, Record};
use worker::*;

#[event(fetch)]
pub async fn fetch(req: Request, env: Env, _ctx: Context) -> Result<Response> {
    console_error_panic_hook::set_once();

    let url = req.url()?;
    let path = url.path().to_string();
    let method = req.method();

    match (method, path.as_str()) {
        (Method::Get, "/health") => health(),
        (Method::Get, "/v1/weaves") => list_weaves(req, env).await,
        (Method::Get, p) if p.starts_with("/v1/weaves/") && p.ends_with("/body") => {
            fetch_body(req, env).await
        }
        (Method::Get, p) if p.starts_with("/v1/weaves/") => resolve_weave(req, env).await,
        (Method::Post, "/v1/weaves") => post_weave(req, env).await,
        _ => Response::error("not found", 404),
    }
}

fn health() -> Result<Response> {
    Response::from_json(&serde_json::json!({
        "status": "ok",
        "registry": "cura",
        "version": env!("CARGO_PKG_VERSION"),
    }))
}

fn plane_from_query(url: &Url) -> Plane {
    for (k, v) in url.query_pairs() {
        if k == "plane" {
            return match v.as_ref() {
                "intent" => Plane::Intent,
                "pattern" => Plane::Pattern,
                _ => Plane::Weave,
            };
        }
    }
    Plane::Weave
}

fn binding_for(p: Plane) -> &'static str {
    match p {
        Plane::Intent => "DB_DI",
        Plane::Weave => "DB_DIP",
        Plane::Pattern => "DB_DP",
    }
}

fn bucket_for(p: Plane) -> &'static str {
    match p {
        Plane::Intent => "BODIES_RI",
        Plane::Weave => "BODIES_RIP",
        Plane::Pattern => "BODIES_RP",
    }
}

fn table_for(p: Plane) -> &'static str {
    p.as_str()
}

#[derive(serde::Deserialize)]
struct WeaveRow {
    weave_id: String,
    sheaf_kind: String,
    input: String,
    logic: String,
    output: String,
    source: String,
    observed_at: i64,
    duration_ms: u64,
    body_key: Option<String>,
    origin: Option<String>,
}

#[derive(serde::Deserialize)]
struct BodyRow {
    body_key: Option<String>,
}

async fn resolve_weave(req: Request, env: Env) -> Result<Response> {
    let url = req.url()?;
    let id = url
        .path()
        .trim_start_matches("/v1/weaves/")
        .to_string();
    if !is_valid_weave_id(&id) {
        return Response::error("invalid weave id", 400);
    }
    let plane = plane_from_query(&url);
    let db = env.d1(binding_for(plane))?;
    let table = table_for(plane);

    let sql = format!(
        "SELECT weave_id, sheaf_kind, input, logic, output, source, observed_at, duration_ms, body_key, origin \
         FROM {} WHERE weave_id = ?1",
        table
    );
    let row: Option<WeaveRow> = db
        .prepare(&sql)
        .bind(&[id.clone().into()])?
        .first(None)
        .await?;

    match row {
        Some(r) => {
            let origin_val: Option<serde_json::Value> = r
                .origin
                .as_deref()
                .and_then(|s| serde_json::from_str(s).ok());
            Response::from_json(&serde_json::json!({
                "weave_id": r.weave_id,
                "plane": plane.as_str(),
                "sheaf_kind": r.sheaf_kind,
                "intent": r.input,
                "weave": r.logic,
                "pattern": r.output,
                "source": r.source,
                "observed_at_ms": r.observed_at,
                "duration_ms": r.duration_ms,
                "body_key": r.body_key,
                "origin": origin_val,
            }))
        }
        None => Response::error("not found", 404),
    }
}

async fn list_weaves(req: Request, env: Env) -> Result<Response> {
    let url = req.url()?;
    let plane = plane_from_query(&url);
    let mut limit: usize = 20;
    let mut before: Option<i64> = None;
    for (k, v) in url.query_pairs() {
        match k.as_ref() {
            "limit" => limit = v.parse().unwrap_or(20).min(100),
            "before" => before = v.parse().ok(),
            _ => {}
        }
    }

    let db = env.d1(binding_for(plane))?;
    let table = table_for(plane);

    let sql = if before.is_some() {
        format!(
            "SELECT weave_id, sheaf_kind, input, logic, output, source, observed_at, duration_ms, body_key, origin \
             FROM {} WHERE observed_at < ?1 ORDER BY observed_at DESC LIMIT ?2",
            table
        )
    } else {
        format!(
            "SELECT weave_id, sheaf_kind, input, logic, output, source, observed_at, duration_ms, body_key, origin \
             FROM {} ORDER BY observed_at DESC LIMIT ?1",
            table
        )
    };

    let stmt = db.prepare(&sql);
    let bound = if let Some(t) = before {
        stmt.bind(&[t.into(), (limit as f64).into()])?
    } else {
        stmt.bind(&[(limit as f64).into()])?
    };
    let rows: Vec<WeaveRow> = bound.all().await?.results()?;

    Response::from_json(&serde_json::json!({
        "plane": plane.as_str(),
        "count": rows.len(),
        "weaves": rows.into_iter().map(|r| serde_json::json!({
            "weave_id": r.weave_id,
            "sheaf_kind": r.sheaf_kind,
            "intent": r.input,
            "weave": r.logic,
            "pattern": r.output,
            "source": r.source,
            "observed_at_ms": r.observed_at,
        })).collect::<Vec<_>>(),
    }))
}

async fn fetch_body(req: Request, env: Env) -> Result<Response> {
    let url = req.url()?;
    let id = url
        .path()
        .trim_start_matches("/v1/weaves/")
        .trim_end_matches("/body")
        .trim_end_matches('/')
        .to_string();
    if !is_valid_weave_id(&id) {
        return Response::error("invalid weave id", 400);
    }
    let plane = plane_from_query(&url);

    let db = env.d1(binding_for(plane))?;
    let table = table_for(plane);
    let sql = format!("SELECT body_key FROM {} WHERE weave_id = ?1", table);
    let row: Option<BodyRow> = db
        .prepare(&sql)
        .bind(&[id.into()])?
        .first(None)
        .await?;

    let Some(r) = row else { return Response::error("not found", 404) };
    let Some(key) = r.body_key else { return Response::error("no body", 404) };

    let bucket = env.bucket(bucket_for(plane))?;
    let obj = bucket.get(&key).execute().await?;
    let Some(obj) = obj else { return Response::error("body missing", 404) };

    let bytes = obj.body().bytes().await?;
    Response::from_bytes(bytes)
}

#[derive(serde::Deserialize)]
struct PostBody {
    plane: String,
    input: String,
    logic: String,
    output: String,
    source: String,
    #[serde(default)]
    sheaf_kind: String,
    duration_ms: u64,
    #[serde(default)]
    origin: Option<cura_protocol::Origin>,
}

async fn post_weave(mut req: Request, env: Env) -> Result<Response> {
    // Curator authenticates with CURA_ORIGIN_SECRET.
    let signature = match req.headers().get("x-cura-origin-signature")? {
        Some(s) => s,
        None => return Response::error("missing signature", 401),
    };

    let ts_header = req
        .headers()
        .get("x-cura-origin-timestamp")?
        .ok_or_else(|| Error::RustError("missing timestamp".into()))?;
    let timestamp: i64 = ts_header
        .parse()
        .map_err(|_| Error::RustError("invalid timestamp".into()))?;

    let now = (js_sys::Date::now()) as i64;
    if (now - timestamp).abs() > 300_000 {
        return Response::error("stale request", 401);
    }

    let body_bytes = req.bytes().await?;
    let secret = env
        .secret("CURA_ORIGIN_SECRET")
        .map_err(|_| Error::RustError("CURA_ORIGIN_SECRET unset".into()))?
        .to_string();

    let canonical = format!(
        "POST\n/v1/weaves\n{}\n{}",
        timestamp,
        cura_crypto::sha256_hex(&body_bytes)
    );
    if !cura_crypto::verify(canonical.as_bytes(), &signature, secret.as_bytes()) {
        return Response::error("signature mismatch", 401);
    }

    let body: PostBody = serde_json::from_slice(&body_bytes)
        .map_err(|e| Error::RustError(format!("body: {e}")))?;

    let plane = match body.plane.as_str() {
        "intent" => Plane::Intent,
        "pattern" => Plane::Pattern,
        "weave" => Plane::Weave,
        other => return Response::error(format!("unknown plane: {other}"), 400),
    };

    let kind = if body.sheaf_kind.is_empty() {
        classify_from_source(&body.source)
    } else {
        match body.sheaf_kind.as_str() {
            "lexical" => cura_protocol::ObservationKind::Lexical,
            "semantic" => cura_protocol::ObservationKind::Semantic,
            "program" => cura_protocol::ObservationKind::Program,
            "relay" => cura_protocol::ObservationKind::Relay,
            "introspective" => cura_protocol::ObservationKind::Introspective,
            "service" => cura_protocol::ObservationKind::Service,
            "sheath" => cura_protocol::ObservationKind::Sheath,
            "firmware_software" => cura_protocol::ObservationKind::FirmwareSoftware,
            "hardware_software" => cura_protocol::ObservationKind::HardwareSoftware,
            other => return Response::error(format!("unknown sheaf: {other}"), 400),
        }
    };

    let triple = cura_protocol::Triple::new(
        body.input.clone(),
        body.logic.clone(),
        body.output.clone(),
    );
    let weave_id = triple.weave_id();
    let mut record = Record::new(
        plane,
        triple.clone(),
        kind.as_str(),
        &body.source,
        body.duration_ms,
    );
    record.observed_at_ms = now;
    if let Some(origin) = body.origin.clone() {
        record.origin = Some(origin);
    }

    // Body stored in R2 as serialized JSON of the triple + origin.
    let body_key = body_key_for(&plane, &weave_id);
    let body_payload = serde_json::json!({
        "triple": record.triple,
        "origin": record.origin,
        "source": record.source,
    });
    let body_bytes = serde_json::to_vec(&body_payload)
        .map_err(|e| Error::RustError(format!("serialize: {e}")))?;

    let bucket = env.bucket(bucket_for(plane))?;
    bucket.put(&body_key, body_bytes).execute().await?;
    record.body_key = Some(body_key.clone());

    // D1 insert.
    let db = env.d1(binding_for(plane))?;
    let table = table_for(plane);
    let sql = format!(
        "INSERT OR IGNORE INTO {} \
         (weave_id, input, logic, output, sheaf_kind, stalk, observed_at, duration_ms, body_key, source, origin, created_at) \
         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
        table
    );

    let origin_json = record
        .origin
        .as_ref()
        .and_then(|o| serde_json::to_string(o).ok())
        .unwrap_or_default();

    let now_iso = js_sys::Date::new_0().to_iso_string().into();

    db.prepare(&sql)
        .bind(&[
            record.weave_id.clone().into(),
            record.triple.input.0.clone().into(),
            record.triple.logic.0.clone().into(),
            record.triple.output.0.clone().into(),
            record.sheaf_kind.clone().into(),
            String::new().into(),
            (record.observed_at_ms as f64).into(),
            (record.duration_ms as f64).into(),
            body_key.into(),
            record.source.clone().into(),
            origin_json.into(),
            now_iso,
        ])?
        .run()
        .await?;

    Response::from_json(&serde_json::json!({
        "weave_id": record.weave_id,
        "status": "stored",
        "plane": plane.as_str(),
    }))
}

fn is_valid_weave_id(id: &str) -> bool {
    id.len() == 64 && id.chars().all(|c| c.is_ascii_hexdigit() && !c.is_ascii_uppercase())
}
```

---

6. crates/curator

Cargo.toml

```toml
[package]
name = "curator"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
publish = false

[dependencies]
cura-protocol = { path = "../cura-protocol" }
cura-crypto = { path = "../cura-crypto" }
serde.workspace = true
serde_json.workspace = true
thiserror.workspace = true
walkdir.workspace = true
reqwest.workspace = true
tokio.workspace = true
```

src/lib.rs

```rust
//! Curator — reads any repository, writes YA|RA.

pub mod emitter;
pub mod source;
pub mod sources;

pub use emitter::{EmitError, Emitter};
pub use source::{Observation, Source, SourceError};
```

src/source.rs

```rust
use cura_protocol::{Origin, Triple};

#[derive(Debug, thiserror::Error)]
pub enum SourceError {
    #[error("io: {0}")]
    Io(#[from] std::io::Error),
    #[error("parse: {0}")]
    Parse(String),
    #[error("not a valid source root: {0}")]
    NotARoot(String),
}

#[derive(Clone, Debug)]
pub struct Observation {
    pub triple: Triple,
    pub origin: Origin,
    pub source_string: String,
}

/// A translator reads a repository and produces observations.
/// One implementation per language. Each is replaceable.
pub trait Source {
    /// Stable identifier for this translator.
    fn name(&self) -> &'static str;

    /// The language this translator handles.
    fn lang(&self) -> &'static str;

    /// Read and produce observations. `repo` is the identifier the
    /// caller wants recorded; `revision` is the git revision.
    fn observe(
        &self,
        root: &std::path::Path,
        repo: &str,
        revision: &str,
    ) -> Result<Vec<Observation>, SourceError>;
}
```

src/emitter.rs

```rust
use crate::source::Observation;
use cura_protocol::Plane;
use serde::Serialize;

#[derive(Debug, thiserror::Error)]
pub enum EmitError {
    #[error("network: {0}")]
    Network(String),
    #[error("rejected: {0}")]
    Rejected(String),
    #[error("encode: {0}")]
    Encode(String),
}

#[derive(Serialize)]
struct SubmitBody<'a> {
    plane: &'a str,
    input: &'a str,
    logic: &'a str,
    output: &'a str,
    source: &'a str,
    sheaf_kind: &'a str,
    duration_ms: u64,
    origin: &'a cura_protocol::Origin,
}

pub struct Emitter {
    url: String,
    secret: Vec<u8>,
    plane: Plane,
    sheaf_kind: String,
    http: reqwest::Client,
}

impl Emitter {
    pub fn new(
        url: impl Into<String>,
        secret: impl Into<Vec<u8>>,
        plane: Plane,
        sheaf_kind: impl Into<String>,
    ) -> Self {
        Emitter {
            url: url.into().trim_end_matches('/').to_string(),
            secret: secret.into(),
            plane,
            sheaf_kind: sheaf_kind.into(),
            http: reqwest::Client::builder()
                .timeout(std::time::Duration::from_secs(30))
                .build()
                .expect("http client"),
        }
    }

    pub async fn submit(&self, obs: &Observation) -> Result<String, EmitError> {
        let body = SubmitBody {
            plane: self.plane.as_str(),
            input: &obs.triple.input.0,
            logic: &obs.triple.logic.0,
            output: &obs.triple.output.0,
            source: &obs.source_string,
            sheaf_kind: &self.sheaf_kind,
            duration_ms: 88_000,
            origin: &obs.origin,
        };

        let bytes = serde_json::to_vec(&body).map_err(|e| EmitError::Encode(e.to_string()))?;
        let ts = now_ms();
        let canonical = format!(
            "POST\n/v1/weaves\n{}\n{}",
            ts,
            cura_crypto::sha256_hex(&bytes)
        );
        let sig = cura_crypto::sign(canonical.as_bytes(), &self.secret);

        let url = format!("{}/v1/weaves", self.url);
        let resp = self
            .http
            .post(&url)
            .header("content-type", "application/json")
            .header("x-cura-origin-signature", &sig)
            .header("x-cura-origin-timestamp", ts.to_string())
            .body(bytes)
            .send()
            .await
            .map_err(|e| EmitError::Network(e.to_string()))?;

        let status = resp.status();
        let text = resp
            .text()
            .await
            .map_err(|e| EmitError::Network(e.to_string()))?;

        if !status.is_success() {
            return Err(EmitError::Rejected(format!("{status}: {text}")));
        }

        let v: serde_json::Value =
            serde_json::from_str(&text).map_err(|e| EmitError::Network(e.to_string()))?;
        Ok(v["weave_id"].as_str().unwrap_or("").to_string())
    }
}

fn now_ms() -> i64 {
    std::time::SystemTime::now()
        .duration_since(std::time::UNIX_EPOCH)
        .map(|d| d.as_millis() as i64)
        .unwrap_or(0)
}
```

src/sources/mod.rs

```rust
pub mod rust;
```

src/sources/rust/mod.rs

```rust
mod emit;
mod walk;

use crate::source::{Observation, Source, SourceError};
use std::path::Path;

pub struct RustSource;

impl Source for RustSource {
    fn name(&self) -> &'static str { "rust-lines-v1" }
    fn lang(&self) -> &'static str { "rust" }

    fn observe(
        &self,
        root: &Path,
        repo: &str,
        revision: &str,
    ) -> Result<Vec<Observation>, SourceError> {
        let files = walk::collect_rust_files(root)?;
        let mut obs = Vec::new();
        for file in files {
            let rel = file.strip_prefix(root).unwrap_or(&file);
            let rel_str = rel.to_string_lossy().replace('\\', "/");
            let text = std::fs::read_to_string(&file)?;
            obs.extend(emit::emit_from_file(&rel_str, &text, repo, revision)?);
        }
        Ok(obs)
    }
}
```

src/sources/rust/walk.rs

```rust
use crate::source::SourceError;
use std::path::{Path, PathBuf};
use walkdir::WalkDir;

pub fn collect_rust_files(root: &Path) -> Result<Vec<PathBuf>, SourceError> {
    if !root.is_dir() {
        return Err(SourceError::NotARoot(root.display().to_string()));
    }
    let mut out = Vec::new();
    for entry in WalkDir::new(root).into_iter().filter_map(|e| e.ok()) {
        let p = entry.path();
        if !p.is_file() { continue; }
        if p.extension().and_then(|s| s.to_str()) != Some("rs") { continue; }
        if p.components().any(|c| c.as_os_str() == "target") { continue; }
        out.push(p.to_path_buf());
    }
    out.sort();
    Ok(out)
}
```

src/sources/rust/emit.rs

```rust
use crate::source::{Observation, SourceError};
use cura_protocol::{Origin, Triple};

const TRANSLATOR: &str = "rust-lines-v1";

pub fn emit_from_file(
    rel_path: &str,
    text: &str,
    repo: &str,
    revision: &str,
) -> Result<Vec<Observation>, SourceError> {
    let lines: Vec<&str> = text.lines().collect();
    let mut obs = Vec::new();
    let mut i = 0usize;

    while i < lines.len() {
        let trimmed = lines[i].trim_start();

        if let Some(name) = parse_pub_fn(trimmed) {
            let end = find_block_end(&lines, i);
            let span = (i as u32 + 1, end as u32 + 1);
            let signature = trimmed.trim_end_matches('{').trim().to_string();
            let emitted = detect_emitted(&lines[i..=end]);

            let triple = Triple::new(
                format!("fn {name}: {signature}"),
                format!("{rel_path} lines {}–{}", span.0, span.1),
                format!("emits: {}", emitted.join(", ")),
            );

            let origin = Origin::new(repo, rel_path, revision, TRANSLATOR, "rust")
                .with_span(span.0, span.1);

            obs.push(Observation {
                triple,
                origin,
                source_string: format!("rust:{rel_path}:fn:{name}"),
            });
        }

        if let Some(name) = parse_pub_struct(trimmed) {
            let span = (i as u32 + 1, i as u32 + 1);
            let triple = Triple::new(
                format!("struct {name}"),
                format!("public type in {rel_path}"),
                format!("visible at line {}", span.0),
            );
            let origin = Origin::new(repo, rel_path, revision, TRANSLATOR, "rust")
                .with_span(span.0, span.1);
            obs.push(Observation {
                triple,
                origin,
                source_string: format!("rust:{rel_path}:struct:{name}"),
            });
        }

        if let Some(name) = parse_pub_enum(trimmed) {
            let span = (i as u32 + 1, i as u32 + 1);
            let triple = Triple::new(
                format!("enum {name}"),
                format!("public sum type in {rel_path}"),
                format!("visible at line {}", span.0),
            );
            let origin = Origin::new(repo, rel_path, revision, TRANSLATOR, "rust")
                .with_span(span.0, span.1);
            obs.push(Observation {
                triple,
                origin,
                source_string: format!("rust:{rel_path}:enum:{name}"),
            });
        }

        if let Some(name) = parse_pub_trait(trimmed) {
            let span = (i as u32 + 1, i as u32 + 1);
            let triple = Triple::new(
                format!("trait {name}"),
                format!("public interface in {rel_path}"),
                format!("visible at line {}", span.0),
            );
            let origin = Origin::new(repo, rel_path, revision, TRANSLATOR, "rust")
                .with_span(span.0, span.1);
            obs.push(Observation {
                triple,
                origin,
                source_string: format!("rust:{rel_path}:trait:{name}"),
            });
        }

        i += 1;
    }

    Ok(obs)
}

fn parse_pub_fn(line: &str) -> Option<String> {
    let s = line.strip_prefix("pub fn ").or_else(|| line.strip_prefix("pub async fn "))?;
    take_ident(s)
}

fn parse_pub_struct(line: &str) -> Option<String> {
    let s = line.strip_prefix("pub struct ")?;
    take_ident(s)
}

fn parse_pub_enum(line: &str) -> Option<String> {
    let s = line.strip_prefix("pub enum ")?;
    take_ident(s)
}

fn parse_pub_trait(line: &str) -> Option<String> {
    let s = line.strip_prefix("pub trait ")?;
    take_ident(s)
}

fn take_ident(s: &str) -> Option<String> {
    let name: String = s
        .chars()
        .take_while(|c| c.is_alphanumeric() || *c == '_')
        .collect();
    if name.is_empty() { None } else { Some(name) }
}

fn find_block_end(lines: &[&str], start: usize) -> usize {
    let mut depth: i32 = 0;
    let mut seen_open = false;
    for (offset, line) in lines[start..].iter().enumerate() {
        for ch in line.chars() {
            if ch == '{' { depth += 1; seen_open = true; }
            else if ch == '}' { depth -= 1; }
        }
        if seen_open && depth <= 0 {
            return start + offset;
        }
    }
    start
}

fn detect_emitted(lines: &[&str]) -> Vec<String> {
    let mut found: Vec<String> = Vec::new();
    for line in lines {
        if line.contains(".emit(") || line.contains("emit_import(") || line.contains("emit(") {
            if !found.iter().any(|x| x == "emit") {
                found.push("emit".to_string());
            }
        }
        if line.contains("Stance::") && !found.iter().any(|x| x == "stance") {
            found.push("stance".to_string());
        }
        if line.contains("return ") && !found.iter().any(|x| x == "return") {
            found.push("return".to_string());
        }
    }
    if found.is_empty() {
        found.push("none-detected".to_string());
    }
    found
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn finds_pub_fn() {
        let text = "pub fn foo() {\n    let x = 1;\n}\n";
        let obs = emit_from_file("src/lib.rs", text, "test/repo", "main").unwrap();
        assert_eq!(obs.len(), 1);
        assert!(obs[0].triple.input.0.contains("foo"));
    }

    #[test]
    fn finds_pub_struct_enum_trait() {
        let text = "pub struct S;\npub enum E { A, B }\npub trait T {}\n";
        let obs = emit_from_file("src/lib.rs", text, "test/repo", "main").unwrap();
        assert_eq!(obs.len(), 3);
    }

    #[test]
    fn ignores_private() {
        let text = "fn foo() {}\nstruct Bar;\n";
        let obs = emit_from_file("src/lib.rs", text, "test/repo", "main").unwrap();
        assert!(obs.is_empty());
    }

    #[test]
    fn origin_carries_repo_and_path() {
        let text = "pub fn foo() {}\n";
        let obs = emit_from_file("src/lib.rs", text, "yashrajvansh/curator", "main").unwrap();
        assert_eq!(obs[0].origin.repo, "yashrajvansh/curator");
        assert_eq!(obs[0].origin.path, "src/lib.rs");
        assert_eq!(obs[0].origin.lang, "rust");
    }
}
```

---

7. crates/curator-cli

Cargo.toml

```toml
[package]
name = "curator-cli"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
publish = false

[[bin]]
name = "curator"
path = "src/main.rs"

[dependencies]
curator = { path = "../curator" }
cura-protocol = { path = "../cura-protocol" }
clap.workspace = true
tokio.workspace = true
anyhow.workspace = true
serde_json.workspace = true
```

src/main.rs

```rust
use clap::Parser;
use curator::sources::rust::RustSource;
use curator::{Emitter, Source};
use cura_protocol::Plane;
use std::path::PathBuf;
use std::process::ExitCode;

#[derive(Parser)]
#[command(
    name = "curator",
    version,
    about = "Read any repository, write YA|RA, submit to Cura"
)]
struct Cli {
    #[command(subcommand)]
    command: Command,
}

#[derive(clap::Subcommand)]
enum Command {
    /// Translate a repository into YA|RA and submit each triple to Cura.
    Translate {
        /// Path to the repository root.
        path: PathBuf,

        /// Cura base URL.
        #[arg(long, env = "CURA_URL")]
        cura: String,

        /// Origin secret Cura issued to this Curator.
        #[arg(long, env = "CURA_ORIGIN_SECRET")]
        secret: String,

        /// Repository identifier, e.g. yashrajvansh/curator.
        #[arg(long)]
        repo: String,

        /// Git revision. Tag, commit, or branch.
        #[arg(long, default_value = "main")]
        revision: String,

        /// Target plane: intent, weave, pattern.
        #[arg(long, default_value = "weave")]
        plane: String,

        /// Sheaf kind. Default is sheath.
        #[arg(long, default_value = "sheath")]
        sheaf: String,

        /// Print observations without submitting.
        #[arg(long)]
        dry_run: bool,
    },
}

#[tokio::main]
async fn main() -> ExitCode {
    let cli = Cli::parse();
    match cli.command {
        Command::Translate { path, cura, secret, repo, revision, plane, sheaf, dry_run } => {
            match translate(path, cura, secret, repo, revision, plane, sheaf, dry_run).await {
                Ok(()) => ExitCode::SUCCESS,
                Err(e) => {
                    eprintln!("curator: {e}");
                    ExitCode::from(1)
                }
            }
        }
    }
}

async fn translate(
    path: PathBuf,
    cura: String,
    secret: String,
    repo: String,
    revision: String,
    plane: String,
    sheaf: String,
    dry_run: bool,
) -> anyhow::Result<()> {
    let plane = match plane.as_str() {
        "intent" => Plane::Intent,
        "pattern" => Plane::Pattern,
        _ => Plane::Weave,
    };

    let src = RustSource;
    let obs = src.observe(&path, &repo, &revision)?;

    eprintln!("translator: {}", src.name());
    eprintln!("lang:       {}", src.lang());
    eprintln!("repo:       {}", repo);
    eprintln!("revision:   {}", revision);
    eprintln!("plane:      {}", plane.as_str());
    eprintln!("observations: {}", obs.len());

    if dry_run {
        for o in &obs {
            println!("{}", serde_json::to_string_pretty(o)?);
        }
        return Ok(());
    }

    let emitter = Emitter::new(cura, secret.into_bytes(), plane, sheaf);
    let mut stored = 0usize;
    let mut failed = 0usize;

    for o in &obs {
        match emitter.submit(o).await {
            Ok(id) => {
                eprintln!("stored {} -> {}", o.source_string, id);
                stored += 1;
            }
            Err(e) => {
                eprintln!("failed {}: {}", o.source_string, e);
                failed += 1;
            }
        }
    }

    eprintln!();
    eprintln!("stored: {stored}");
    eprintln!("failed: {failed}");

    Ok(())
}
```

---

8. Root files

README.md

```markdown
# curator

Curator reads any repository, in any language, and writes it as YA|RA.
Cura is the registry that holds the result.

- **YA|RA** — the notation. `input · logic · output` mapped to `intent · weave · pattern`.
- **Cura** — the registry. Three D1 databases (`cura-di`, `cura-dip`, `cura-dp`), three R2 buckets (`cura-ri`, `cura-rip`, `cura-rp`).
- **Curator** — the translator. One adapter per language. The default sheaf is the sheath.

The duration of reasoning is the moment of intent formation until Curator starts to see other patterns. That is the cage.

## Layout

| Crate | Role |
|---|---|
| `cura-protocol` | YA\|RA wire types |
| `cura-crypto` | HMAC-SHA-256, SHA-256, BLAKE3 |
| `cura-sheaf` | Sheaf trait, lexical and sheath encoders |
| `cura-registry` | Cura's storage surface |
| `cura-worker` | Cura's HTTP surface (Cloudflare Worker) |
| `curator` | Translator framework and Rust source adapter |
| `curator-cli` | `curator` binary |

## Build

```

cargo build --workspace --release
cargo test --workspace

```

## First translation

```

export CURA_URL=https://cura.<account>.workers.dev
export CURA_ORIGIN_SECRET=<hex>

./target/release/curator translate . 
--repo yashrajvansh/curator 
--revision main

```

## License

UNLICENSED. Copyright Comfort Curators Private Limited.
```

LICENSE

```
UNLICENSED

Copyright Comfort Curators Private Limited.

All rights reserved.
```

CHANGELOG.md

```markdown
# Changelog

## [Unreleased]

### Added

- `cura-protocol`: Plane (Intent / Weave / Pattern), Triple, Stalk, Provenance, Origin, Stance, WeaveId.
- `cura-crypto`: HMAC-SHA-256, SHA-256, BLAKE3. RFC 4231 conformance test.
- `cura-sheaf`: `Sheaf` trait; `LexicalSheaf` and `SheathSheaf`.
- `cura-registry`: `Store` trait, `MemoryStore`, body key layout, source classification.
- `cura-worker`: Cloudflare Worker serving `/v1/weaves` and `/v1/weaves/:id/body`.
- `curator`: translator framework and the Rust source adapter.
- `curator-cli`: `curator translate`.
- Three D1 migrations and three R2 bindings.
```

scripts/bootstrap.sh

```sh
#!/usr/bin/env bash
set -euo pipefail
: "${CLOUDFLARE_API_TOKEN:?export CLOUDFLARE_API_TOKEN}"
: "${CLOUDFLARE_ACCOUNT_ID:?export CLOUDFLARE_ACCOUNT_ID}"

wrangler d1 create cura-di  || true
wrangler d1 create cura-dip || true
wrangler d1 create cura-dp  || true

wrangler r2 bucket create cura-ri  || true
wrangler r2 bucket create cura-rip || true
wrangler r2 bucket create cura-rp  || true

cd crates/cura-worker

wrangler d1 migrations apply cura-di  --remote --migrations-dir=migrations/intent
wrangler d1 migrations apply cura-dip --remote --migrations-dir=migrations/weave
wrangler d1 migrations apply cura-dp  --remote --migrations-dir=migrations/pattern

wrangler secret put CURA_ORIGIN_SECRET
wrangler deploy

echo "cura: https://cura.<account>.workers.dev"
```

.github/workflows/release.yml

```yaml
name: release

on:
  push:
    branches: [main]

concurrency:
  group: release-main
  cancel-in-progress: false

jobs:
  rust:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
      - uses: dtolnay/rust-toolchain@1.83.0
        with:
          targets: wasm32-unknown-unknown
          components: rustfmt, clippy
      - uses: Swatinem/rust-cache@v2
      - run: cargo fmt --check
      - run: cargo clippy --workspace --all-targets -- -D warnings
      - run: cargo test --workspace
      - run: cargo build --release -p curator-cli
      - uses: actions/upload-artifact@v4
        with:
          name: curator
          path: target/release/curator

  deploy:
    needs: [rust]
    runs-on: ubuntu-latest
    environment: production
    steps:
      - uses: actions/checkout@v4
      - uses: dtolnay/rust-toolchain@1.83.0
        with:
          targets: wasm32-unknown-unknown
      - uses: Swatinem/rust-cache@v2
      - run: |
          cd crates/cura-worker
          wrangler d1 migrations apply cura-di  --remote --migrations-dir=migrations/intent
          wrangler d1 migrations apply cura-dip --remote --migrations-dir=migrations/weave
          wrangler d1 migrations apply cura-dp  --remote --migrations-dir=migrations/pattern
          wrangler deploy
        env:
          CLOUDFLARE_API_TOKEN: ${{ secrets.CLOUDFLARE_API_TOKEN }}
          CLOUDFLARE_ACCOUNT_ID: ${{ secrets.CLOUDFLARE_ACCOUNT_ID }}
```

---

First run

```sh
cargo build --workspace --release
cargo test --workspace

# Local
cd crates/cura-worker
echo 'CURA_ORIGIN_SECRET=dev_origin_secret' > .dev.vars
wrangler d1 migrations apply cura-di  --local --migrations-dir=migrations/intent
wrangler d1 migrations apply cura-dip --local --migrations-dir=migrations/weave
wrangler d1 migrations apply cura-dp  --local --migrations-dir=migrations/pattern
wrangler dev

# Other terminal
export CURA_URL=http://localhost:8787
export CURA_ORIGIN_SECRET=dev_origin_secret

./target/release/curator translate crates/cura-protocol \
  --repo yashrajvansh/curator --revision main
```

Expected: stored: N, failed: 0. Then:

```sh
curl "http://localhost:8787/v1/weaves?plane=weave&limit=10" | jq .
```

That is Curator reading this repository and writing it into Cura.