polyc-facts 2026.10.2

Shared semantic-fold library: decode-to-fact functions reused by every consumer that reads the event log, so a payment receipt or a tool call means the same thing everywhere it's read.
//! `observed-routines/v1`'s row type and its canonical payload encoding
//! (7-O).
//!
//! Unlike every other fold in this crate, this one runs in both directions.
//! The control-plane observer (`crates/control-plane/src/
//! routine_observer.rs`) converts each `Routine` custom resource to
//! [`ObservedRoutineRow`] with an explicit `impl From<&Routine>` naming
//! every field — that conversion lives beside `Routine` itself, not here,
//! since this crate keeps zero Kubernetes reach — then calls
//! [`encode_rows`] (re-exported as `encode_observed_routine_rows`) to build
//! the opaque payload `polyc_state::observation::RecordObservation` carries.
//! The projector's observed-source worker
//! (`crates/projector/src/observed_worker.rs`) calls [`decode_rows`]
//! (re-exported as `decode_observed_routine_rows`) on the payload State
//! hands back, and gets the identical rows the observer encoded.
//!
//! This crate's own charter says presentation — truncation, snippet
//! windowing, size budgets — stays with the consumer. [`ObservedRoutineRow`]
//! keeps that split: `prompt` is whatever the observer decided to carry
//! (already bounded there, so the row stays under
//! `polyc_state::observation::MAX_OBSERVED_ROW_BYTES`), and
//! [`ObservedRoutineRow::prompt_truncated`] is the observer's own record of
//! whether it truncated, never a decision this module makes.

use serde::{Deserialize, Serialize};

/// One routine in one observation, as the observer converts it and the
/// projector folds it.
///
/// Field order matches `02-DESIGN.md` §6.2's `observed_routines` table,
/// which is also the Arrow column order the projector publishes
/// (`crates/projector/src/artifact.rs`) — the observation columns first (from the
/// observation record itself, the same on every row of one observation),
/// then the routine's own status/spec surface, in the same order
/// `crates/query/src/routine_catalog.rs::RoutineStatusRecord` already
/// establishes for the fields the two share.
// Five independent flags on a `Routine`'s own status/spec surface
// (`ready`, `prompt_truncated`, `suspended`, `orphaned`,
// `setup_completed`), not a state machine: each names a fact Kubernetes
// or the observer already records independently, and
// `RoutineStatusRecord` (`crates/query/src/routine_catalog.rs`) already
// carries four of the same five.
#[allow(clippy::struct_excessive_bools)]
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ObservedRoutineRow {
    /// The collection's cluster namespace (`ObservationSource`'s own, not a
    /// Polychrome tenant namespace). Known to the observer before it ever
    /// calls State — it is the observer's own configured scope.
    pub namespace: String,
    /// The observation this row belongs to
    /// (`polyc_state::observation::ObservationOrdinal::as_journal_position`).
    ///
    /// Unknown to the observer at encode time — State assigns the ordinal
    /// only once `RecordObservation` commits. The observer encodes `0`
    /// here (never a real ordinal, since [`polyc_state::observation::
    /// ObservationOrdinal::ORIGIN`] means "no observation recorded yet");
    /// the projector overwrites every row's copy with the real ordinal from
    /// the decoded [`polyc_state::observation::ObservationRecord`] before
    /// publishing, the one place this crate's "presentation stays with the
    /// consumer" charter puts that responsibility.
    pub observation_ordinal: u64,
    /// The list call's collection `resourceVersion` — opaque, carried for
    /// display, never compared for ordering. Known at encode time (the
    /// list call already returned it).
    pub observation_resource_version: String,
    /// When the observer's clock made the list call, epoch ms. Known at
    /// encode time.
    pub observed_at_ms: i64,
    /// When State recorded the observation, epoch ms.
    ///
    /// Unknown to the observer at encode time, for the same reason
    /// [`Self::observation_ordinal`] is: the observer encodes `0`, and the
    /// projector overwrites it from the decoded `ObservationRecord`.
    pub recorded_at_ms: i64,
    /// The observing identity (the control plane's own workload identity).
    /// Known at encode time.
    pub observer: String,
    /// The `Routine` resource's own name.
    pub name: String,
    /// The `Routine` resource's own stable Kubernetes `uid`.
    pub uid: String,
    /// The row's own per-object `metadata.resourceVersion` — opaque, distinct
    /// from `observation_resource_version` (the collection's).
    pub resource_version: String,
    /// `metadata.generation`.
    pub generation: i64,
    /// The routine's synthetic prompt-fire conversation id, derived by the
    /// observer the same way `RoutineStatusRecord::fire_conversation_id` is.
    pub fire_conversation_id: String,
    /// `RoutineProvenance::creator_persona`.
    pub creator_persona: String,
    /// `RoutineProvenance::conversation_id`.
    pub provenance_conversation_id: String,
    /// `RoutineSpec::scope`, rendered `"public"`/`"private"`.
    pub scope: String,
    /// `RoutineSpec::display_name`. Empty when unset.
    pub display_name: String,
    /// `RoutineSpec::description`. Empty when unset.
    pub description: String,
    /// The compiled `RoutineSchedule`, as JSON text.
    pub schedule_json: String,
    /// The schedule's own IANA zone name, or `"UTC"` when unset.
    pub schedule_timezone: String,
    /// Up to the next three fire instants, as a JSON array of epoch-ms
    /// numbers. `"[]"` when the schedule yields none.
    pub next_fires_json: String,
    /// `RoutinePayload::prompt`, as the observer bounded it. Never absent.
    pub prompt: String,
    /// Whether the observer truncated [`Self::prompt`] to stay under the
    /// row bound. The observer's own record, not a decision this crate
    /// makes.
    pub prompt_truncated: bool,
    /// `RoutineStatus::ready`.
    pub ready: bool,
    /// `RoutineStatus::phase`. `None` before the first reconcile.
    pub phase: Option<String>,
    /// `RoutineStatus::message`.
    pub message: Option<String>,
    /// `RoutineStatus::last_fire_time`, epoch ms.
    pub last_fire_time_ms: Option<i64>,
    /// `RoutineStatus::next_fire_time`, epoch ms.
    pub next_fire_time_ms: Option<i64>,
    /// `RoutineStatus::conditions`, JSON-array text.
    pub conditions_json: String,
    /// `RoutineSpec::suspend.is_some()`.
    pub suspended: bool,
    /// `RoutineSuspend::paused_by`. `None` while active.
    pub paused_by: Option<String>,
    /// `RoutineSuspend::paused_at`, epoch ms. `None` while active.
    pub paused_at_ms: Option<i64>,
    /// `RoutineSuspend::reason`. `None` while active or unset.
    pub pause_reason: Option<String>,
    /// Whether the scheduler-written `Orphaned` condition is currently
    /// `True`.
    pub orphaned: bool,
    /// Whether the scheduler-written `AwaitingSetup` condition reads
    /// `False` — the setup exchange's marker committed. `True` or absent
    /// reads `false`, the same derivation `RoutineStatusRecord::
    /// setup_completed` (`crates/query/src/routine_catalog.rs`) documents.
    pub setup_completed: bool,
}

/// The wire shape [`encode_rows`] writes and [`decode_rows`] reads.
///
/// A one-field envelope rather than a bare array, so a future field this
/// module adds beside the row list (a payload-wide flag, say) is additive
/// rather than a breaking reshape of the top-level JSON value.
#[derive(Debug, Clone, Serialize, Deserialize)]
struct ObservedRoutinesPayload {
    /// The payload format this module writes and reads. A build that widens
    /// the shape bumps this and refuses an older one, the same "no
    /// compatibility reader" rule every durable format in this workspace
    /// follows — the wipe is at the `ObservationOrdinal` layer (a new
    /// observation), never a silent reinterpretation of old bytes.
    version: u8,
    rows: Vec<ObservedRoutineRow>,
}

// Version 2 widened `ObservedRoutineRow` by one flag, `setup_completed` —
// the first widening this format has taken, since the version-1 shape is
// only hours older on `main` and never shipped.
const PAYLOAD_VERSION: u8 = 2;

/// Refuses to decode a payload this build did not write.
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum DecodeError {
    /// The bytes are not this module's JSON shape at all.
    #[error("the observed-routines payload is not readable JSON: {0}")]
    Malformed(String),
    /// The bytes declare a payload version this build does not read.
    #[error("the observed-routines payload is version {found}; this build reads {expected}")]
    UnknownVersion {
        /// The version the payload declared.
        found: u8,
        /// The version this build writes and reads.
        expected: u8,
    },
}

/// Encodes `rows` as the canonical bytes an observation payload carries.
///
/// Deterministic: the same rows in the same order always encode to the same
/// bytes, which is what lets the observer derive a stable payload digest
/// before ever calling State, and the projector re-derive the identical
/// digest after decoding.
///
/// # Panics
///
/// Never in practice: every field of [`ObservedRoutineRow`] is a plain
/// scalar or `String`, none of which `serde_json` can refuse to serialize.
#[must_use]
pub fn encode_rows(rows: &[ObservedRoutineRow]) -> Vec<u8> {
    let payload = ObservedRoutinesPayload {
        version: PAYLOAD_VERSION,
        rows: rows.to_vec(),
    };
    // A `Vec<u8>` sink cannot fail to write to; the only way `to_writer`
    // returns an error is a type that refuses to serialize, and every field
    // here is a plain, always-serializable scalar or `String`.
    serde_json::to_vec(&payload).expect("an observed-routines payload always serializes")
}

/// Decodes the bytes [`encode_rows`] wrote.
///
/// # Errors
///
/// Returns [`DecodeError::Malformed`] for bytes that are not this module's
/// JSON shape, and [`DecodeError::UnknownVersion`] for a payload this build
/// does not read — refused rather than reinterpreted, the observation's own
/// ordinal is what a caller relists and re-records under to recover.
pub fn decode_rows(bytes: &[u8]) -> Result<Vec<ObservedRoutineRow>, DecodeError> {
    let payload: ObservedRoutinesPayload =
        serde_json::from_slice(bytes).map_err(|error| DecodeError::Malformed(error.to_string()))?;
    if payload.version != PAYLOAD_VERSION {
        return Err(DecodeError::UnknownVersion {
            found: payload.version,
            expected: PAYLOAD_VERSION,
        });
    }
    Ok(payload.rows)
}

#[cfg(test)]
mod tests {
    #![allow(clippy::pedantic, clippy::nursery, missing_docs)]

    use super::{DecodeError, ObservedRoutineRow, decode_rows, encode_rows};

    fn row(name: &str) -> ObservedRoutineRow {
        ObservedRoutineRow {
            namespace: "polychrome".to_owned(),
            observation_ordinal: 1,
            observation_resource_version: "999".to_owned(),
            observed_at_ms: 1_700_000_000_000,
            recorded_at_ms: 1_700_000_000_001,
            observer: "control-plane/test".to_owned(),
            name: name.to_owned(),
            uid: "uid-1".to_owned(),
            resource_version: "42".to_owned(),
            generation: 3,
            fire_conversation_id: "fire-1".to_owned(),
            creator_persona: "persona-1".to_owned(),
            provenance_conversation_id: "conv-1".to_owned(),
            scope: "private".to_owned(),
            display_name: String::new(),
            description: String::new(),
            schedule_json: "{}".to_owned(),
            schedule_timezone: "UTC".to_owned(),
            next_fires_json: "[]".to_owned(),
            prompt: "do the thing".to_owned(),
            prompt_truncated: false,
            ready: true,
            phase: Some("Ready".to_owned()),
            message: None,
            last_fire_time_ms: None,
            next_fire_time_ms: None,
            conditions_json: "[]".to_owned(),
            suspended: false,
            paused_by: None,
            paused_at_ms: None,
            pause_reason: None,
            orphaned: false,
            setup_completed: false,
        }
    }

    #[test]
    fn rows_round_trip_through_the_canonical_encoding() {
        let rows = vec![row("a"), row("b")];
        let bytes = encode_rows(&rows);
        assert_eq!(decode_rows(&bytes).expect("decodes"), rows);
    }

    #[test]
    fn encoding_is_deterministic() {
        let rows = vec![row("a"), row("b")];
        assert_eq!(encode_rows(&rows), encode_rows(&rows));
    }

    #[test]
    fn an_empty_collection_round_trips() {
        let bytes = encode_rows(&[]);
        assert_eq!(decode_rows(&bytes).expect("decodes"), Vec::new());
    }

    #[test]
    fn a_foreign_version_is_refused() {
        let bytes = br#"{"version":99,"rows":[]}"#;
        match decode_rows(bytes) {
            Err(DecodeError::UnknownVersion {
                found: 99,
                expected: 2,
            }) => {}
            other => panic!("expected UnknownVersion, got {other:?}"),
        }
    }

    #[test]
    fn garbage_bytes_are_refused_rather_than_panicking() {
        assert!(matches!(
            decode_rows(b"not json"),
            Err(DecodeError::Malformed(_))
        ));
    }
}