Skip to main content

khive_types/
event.rs

1//! Event substrate — append-only log produced by every verb execution.
2
3extern crate alloc;
4use alloc::string::String;
5use alloc::vec::Vec;
6use core::fmt;
7
8use crate::{Header, Id128, SubstrateKind};
9
10/// A system event. Append-only, never mutated or deleted.
11#[derive(Clone, Debug)]
12#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
13pub struct Event {
14    #[cfg_attr(feature = "serde", serde(flatten))]
15    pub header: Header,
16    /// The verb that produced the event.
17    pub verb: String,
18    /// Which substrate type was acted upon.
19    pub substrate: SubstrateKind,
20    /// Who performed the action. Profile- or system-produced events may omit it.
21    pub actor: Option<String>,
22    /// Typed event discriminant used by replay, projections, and workers.
23    pub kind: EventKind,
24    /// Typed payload surface for known event families; raw JSON is still allowed.
25    pub payload: EventPayload,
26    /// Payload schema version interpreted per `kind`.
27    pub payload_schema_version: u32,
28    /// Brain profile state version observed when the event was emitted.
29    pub profile_state_version: Option<u64>,
30    /// Logical aggregate threaded across related event ids.
31    pub aggregate: Option<AggregateRef>,
32}
33
34/// Outcome of a verb execution recorded in an event log entry.
35#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Default)]
36#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
37#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
38pub enum EventOutcome {
39    /// The verb executed successfully.
40    #[default]
41    Success,
42    /// The verb was denied by a policy check.
43    Denied,
44    /// The verb encountered a runtime error.
45    Error,
46}
47
48impl EventOutcome {
49    /// Return the canonical lowercase string for this outcome.
50    pub const fn name(self) -> &'static str {
51        match self {
52            Self::Success => "success",
53            Self::Denied => "denied",
54            Self::Error => "error",
55        }
56    }
57}
58
59impl fmt::Display for EventOutcome {
60    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
61        f.write_str(self.name())
62    }
63}
64
65/// Discriminant for the 41 typed event variants produced by the verb dispatch path
66/// and by lifecycle telemetry producers (channel polling/backoff, config-lock,
67/// checkpoint outcome, background phase spans).
68#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
69#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
70#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
71pub enum EventKind {
72    /// Generic audit event with no structured payload.
73    Audit,
74    /// A `recall` verb was executed and results were returned.
75    RecallExecuted,
76    /// A rerank pass was applied to search candidates.
77    RerankExecuted,
78    /// A `search` verb was executed.
79    SearchExecuted,
80    /// A tool-use policy decision was computed by tool.check or exec.run.
81    ToolCheckDecided,
82    /// A new directed edge was created between two nodes.
83    LinkCreated,
84    /// A new entity was created.
85    EntityCreated,
86    /// An existing entity was patched.
87    EntityUpdated,
88    /// An entity was soft- or hard-deleted.
89    EntityDeleted,
90    /// Two entities were merged (deduplication).
91    EntityMerged,
92    /// Two notes were merged (deduplication).
93    NoteMerged,
94    /// A new note was created.
95    NoteCreated,
96    /// An existing note was patched.
97    NoteUpdated,
98    /// A note was soft- or hard-deleted.
99    NoteDeleted,
100    /// An edge's relation or weight was updated.
101    EdgeUpdated,
102    /// An edge was removed.
103    EdgeDeleted,
104    /// A GTD task moved between lifecycle states.
105    TaskTransitioned,
106    /// An explicit user feedback signal was recorded.
107    FeedbackExplicit,
108    /// The brain recommended a profile resolution update.
109    ProfileResolutionRecommended,
110    /// Two brain profiles were merged.
111    ProfileMerged,
112    /// The active embedding model was changed.
113    EmbeddingModelChanged,
114    /// An embedding migration batch completed successfully.
115    EmbeddingMigrationCompleted,
116    /// An embedding migration batch failed.
117    EmbeddingMigrationFailed,
118    /// Drift was detected between stored and live embeddings.
119    EmbeddingDriftDetected,
120    /// A lazily loaded embedder finished initialization.
121    EmbedderInitialized,
122    /// A proposal was submitted for review.
123    ProposalCreated,
124    /// A reviewer accepted, rejected, or commented on a proposal.
125    ProposalReviewed,
126    /// A proposal was applied to the graph.
127    ProposalApplied,
128    /// A proposal was withdrawn before it was applied.
129    ProposalWithdrawn,
130    /// A channel poll cycle started for one `(kind, slug)` credential.
131    ChannelPollStarted,
132    /// A channel poll cycle returned envelopes after a prior failure.
133    ChannelPollSucceeded,
134    /// A channel poll cycle failed.
135    ChannelPollFailed,
136    /// A channel's backoff escalated to a new step after a failure.
137    ChannelBackoffArmed,
138    /// A channel's backoff reset to base after a success.
139    ChannelBackoffReset,
140    /// Persisting a channel heartbeat row failed.
141    ChannelHeartbeatPersistFailed,
142    /// A process-lifetime `OnceLock` configuration value was locked in.
143    ConfigLocked,
144    /// A WAL checkpoint tick's outcome was recorded (ADR-091 elevated/drain edge).
145    CheckpointOutcomeRecorded,
146    /// A background phase (ANN warm, index rebuild/backfill, ...) started (ADR-103 Stage 1).
147    PhaseStarted,
148    /// A background phase completed (ADR-103 Stage 1).
149    PhaseCompleted,
150    /// A background phase was cancelled before completion (ADR-103 Stage 1).
151    PhaseCancelled,
152    /// An attempted operation was refused while retaining its existing subject.
153    Refusal,
154}
155
156impl EventKind {
157    /// All 41 event kind variants in declaration order.
158    pub const ALL: [Self; 41] = [
159        Self::Audit,
160        Self::RecallExecuted,
161        Self::RerankExecuted,
162        Self::SearchExecuted,
163        Self::ToolCheckDecided,
164        Self::LinkCreated,
165        Self::EntityCreated,
166        Self::EntityUpdated,
167        Self::EntityDeleted,
168        Self::EntityMerged,
169        Self::NoteMerged,
170        Self::NoteCreated,
171        Self::NoteUpdated,
172        Self::NoteDeleted,
173        Self::EdgeUpdated,
174        Self::EdgeDeleted,
175        Self::TaskTransitioned,
176        Self::FeedbackExplicit,
177        Self::ProfileResolutionRecommended,
178        Self::ProfileMerged,
179        Self::EmbeddingModelChanged,
180        Self::EmbeddingMigrationCompleted,
181        Self::EmbeddingMigrationFailed,
182        Self::EmbeddingDriftDetected,
183        Self::EmbedderInitialized,
184        Self::ProposalCreated,
185        Self::ProposalReviewed,
186        Self::ProposalApplied,
187        Self::ProposalWithdrawn,
188        Self::ChannelPollStarted,
189        Self::ChannelPollSucceeded,
190        Self::ChannelPollFailed,
191        Self::ChannelBackoffArmed,
192        Self::ChannelBackoffReset,
193        Self::ChannelHeartbeatPersistFailed,
194        Self::ConfigLocked,
195        Self::CheckpointOutcomeRecorded,
196        Self::PhaseStarted,
197        Self::PhaseCompleted,
198        Self::PhaseCancelled,
199        Self::Refusal,
200    ];
201
202    /// Return the canonical snake_case string for this event kind.
203    pub const fn name(self) -> &'static str {
204        match self {
205            Self::Audit => "audit",
206            Self::RecallExecuted => "recall_executed",
207            Self::RerankExecuted => "rerank_executed",
208            Self::SearchExecuted => "search_executed",
209            Self::ToolCheckDecided => "tool_check_decided",
210            Self::LinkCreated => "link_created",
211            Self::EntityCreated => "entity_created",
212            Self::EntityUpdated => "entity_updated",
213            Self::EntityDeleted => "entity_deleted",
214            Self::EntityMerged => "entity_merged",
215            Self::NoteMerged => "note_merged",
216            Self::NoteCreated => "note_created",
217            Self::NoteUpdated => "note_updated",
218            Self::NoteDeleted => "note_deleted",
219            Self::EdgeUpdated => "edge_updated",
220            Self::EdgeDeleted => "edge_deleted",
221            Self::TaskTransitioned => "task_transitioned",
222            Self::FeedbackExplicit => "feedback_explicit",
223            Self::ProfileResolutionRecommended => "profile_resolution_recommended",
224            Self::ProfileMerged => "profile_merged",
225            Self::EmbeddingModelChanged => "embedding_model_changed",
226            Self::EmbeddingMigrationCompleted => "embedding_migration_completed",
227            Self::EmbeddingMigrationFailed => "embedding_migration_failed",
228            Self::EmbeddingDriftDetected => "embedding_drift_detected",
229            Self::EmbedderInitialized => "embedder_initialized",
230            Self::ProposalCreated => "proposal_created",
231            Self::ProposalReviewed => "proposal_reviewed",
232            Self::ProposalApplied => "proposal_applied",
233            Self::ProposalWithdrawn => "proposal_withdrawn",
234            Self::ChannelPollStarted => "channel_poll_started",
235            Self::ChannelPollSucceeded => "channel_poll_succeeded",
236            Self::ChannelPollFailed => "channel_poll_failed",
237            Self::ChannelBackoffArmed => "channel_backoff_armed",
238            Self::ChannelBackoffReset => "channel_backoff_reset",
239            Self::ChannelHeartbeatPersistFailed => "channel_heartbeat_persist_failed",
240            Self::ConfigLocked => "config_locked",
241            Self::CheckpointOutcomeRecorded => "checkpoint_outcome_recorded",
242            Self::PhaseStarted => "phase_started",
243            Self::PhaseCompleted => "phase_completed",
244            Self::PhaseCancelled => "phase_cancelled",
245            Self::Refusal => "refusal",
246        }
247    }
248}
249
250impl fmt::Display for EventKind {
251    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
252        f.write_str(self.name())
253    }
254}
255
256const EVENT_KIND_VALID: &[&str] = &[
257    "audit",
258    "recall_executed",
259    "rerank_executed",
260    "search_executed",
261    "tool_check_decided",
262    "link_created",
263    "entity_created",
264    "entity_updated",
265    "entity_deleted",
266    "entity_merged",
267    "note_merged",
268    "note_created",
269    "note_updated",
270    "note_deleted",
271    "edge_updated",
272    "edge_deleted",
273    "task_transitioned",
274    "feedback_explicit",
275    "profile_resolution_recommended",
276    "profile_merged",
277    "embedding_model_changed",
278    "embedding_migration_completed",
279    "embedding_migration_failed",
280    "embedding_drift_detected",
281    "embedder_initialized",
282    "proposal_created",
283    "proposal_reviewed",
284    "proposal_applied",
285    "proposal_withdrawn",
286    "channel_poll_started",
287    "channel_poll_succeeded",
288    "channel_poll_failed",
289    "channel_backoff_armed",
290    "channel_backoff_reset",
291    "channel_heartbeat_persist_failed",
292    "config_locked",
293    "checkpoint_outcome_recorded",
294    "phase_started",
295    "phase_completed",
296    "phase_cancelled",
297    "refusal",
298];
299
300impl core::str::FromStr for EventKind {
301    type Err = crate::error::UnknownVariant;
302
303    fn from_str(s: &str) -> Result<Self, Self::Err> {
304        match s.trim().to_ascii_lowercase().as_str() {
305            "audit" => Ok(Self::Audit),
306            "recall_executed" => Ok(Self::RecallExecuted),
307            "rerank_executed" => Ok(Self::RerankExecuted),
308            "search_executed" => Ok(Self::SearchExecuted),
309            "tool_check_decided" => Ok(Self::ToolCheckDecided),
310            "link_created" => Ok(Self::LinkCreated),
311            "entity_created" => Ok(Self::EntityCreated),
312            "entity_updated" => Ok(Self::EntityUpdated),
313            "entity_deleted" => Ok(Self::EntityDeleted),
314            "entity_merged" => Ok(Self::EntityMerged),
315            "note_merged" => Ok(Self::NoteMerged),
316            "note_created" => Ok(Self::NoteCreated),
317            "note_updated" => Ok(Self::NoteUpdated),
318            "note_deleted" => Ok(Self::NoteDeleted),
319            "edge_updated" => Ok(Self::EdgeUpdated),
320            "edge_deleted" => Ok(Self::EdgeDeleted),
321            "task_transitioned" => Ok(Self::TaskTransitioned),
322            "feedback_explicit" => Ok(Self::FeedbackExplicit),
323            "profile_resolution_recommended" => Ok(Self::ProfileResolutionRecommended),
324            "profile_merged" => Ok(Self::ProfileMerged),
325            "embedding_model_changed" => Ok(Self::EmbeddingModelChanged),
326            "embedding_migration_completed" => Ok(Self::EmbeddingMigrationCompleted),
327            "embedding_migration_failed" => Ok(Self::EmbeddingMigrationFailed),
328            "embedding_drift_detected" => Ok(Self::EmbeddingDriftDetected),
329            "embedder_initialized" => Ok(Self::EmbedderInitialized),
330            "proposal_created" => Ok(Self::ProposalCreated),
331            "proposal_reviewed" => Ok(Self::ProposalReviewed),
332            "proposal_applied" => Ok(Self::ProposalApplied),
333            "proposal_withdrawn" => Ok(Self::ProposalWithdrawn),
334            "channel_poll_started" => Ok(Self::ChannelPollStarted),
335            "channel_poll_succeeded" => Ok(Self::ChannelPollSucceeded),
336            "channel_poll_failed" => Ok(Self::ChannelPollFailed),
337            "channel_backoff_armed" => Ok(Self::ChannelBackoffArmed),
338            "channel_backoff_reset" => Ok(Self::ChannelBackoffReset),
339            "channel_heartbeat_persist_failed" => Ok(Self::ChannelHeartbeatPersistFailed),
340            "config_locked" => Ok(Self::ConfigLocked),
341            "checkpoint_outcome_recorded" => Ok(Self::CheckpointOutcomeRecorded),
342            "phase_started" => Ok(Self::PhaseStarted),
343            "phase_completed" => Ok(Self::PhaseCompleted),
344            "phase_cancelled" => Ok(Self::PhaseCancelled),
345            "refusal" => Ok(Self::Refusal),
346            other => Err(crate::error::UnknownVariant::new(
347                "event_kind",
348                other,
349                EVENT_KIND_VALID,
350            )),
351        }
352    }
353}
354
355/// A reference to the logical aggregate that an event belongs to.
356///
357/// Used to thread related events (e.g. proposal lifecycle events) into a
358/// single auditable chain identified by `kind` and `id`.
359#[derive(Clone, Debug, PartialEq, Eq)]
360#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
361pub struct AggregateRef {
362    /// The aggregate type string (e.g. `"proposal"`).
363    pub kind: String,
364    /// The aggregate instance identifier.
365    pub id: Id128,
366}
367
368/// Typed payload for an [`Event`], dispatched by [`EventKind`].
369///
370/// The `Json` variant is a catch-all for events whose payload has not yet
371/// been promoted to a structured type. All other variants carry a concrete
372/// typed struct that can be pattern-matched without round-tripping through JSON.
373#[derive(Clone, Debug, PartialEq)]
374#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
375#[cfg_attr(
376    feature = "serde",
377    serde(tag = "kind", content = "payload", rename_all = "snake_case")
378)]
379pub enum EventPayload {
380    /// Raw JSON payload for untyped events.
381    Json(String),
382    /// Structured payload for a rerank pass event.
383    RerankExecuted(RerankExecutedPayload),
384    /// Structured payload for a tool-use policy decision.
385    ToolCheckDecided(ToolCheckDecidedPayload),
386    /// Structured payload for a proposal-created event (requires `serde` feature).
387    #[cfg(feature = "serde")]
388    ProposalCreated(ProposalCreatedPayload),
389    /// Structured payload for a proposal-reviewed event.
390    ProposalReviewed(ProposalReviewedPayload),
391    /// Structured payload for a proposal-applied event.
392    ProposalApplied(ProposalAppliedPayload),
393    /// Structured payload for a proposal-withdrawn event.
394    ProposalWithdrawn(ProposalWithdrawnPayload),
395}
396
397impl Default for EventPayload {
398    fn default() -> Self {
399        Self::Json("{}".into())
400    }
401}
402
403/// The decision computed for one tool-use policy check.
404#[derive(Clone, Debug, PartialEq, Eq)]
405#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
406pub struct ToolCheckDecidedPayload {
407    pub actor: String,
408    pub tool: String,
409    pub registered: bool,
410    pub decision: String,
411    pub source: String,
412    pub id: Option<String>,
413    pub scope: Option<String>,
414    pub caller_verb: String,
415}
416
417/// Payload for a rerank pass event, recording per-candidate scores.
418///
419/// All score values (`reranked` section scores, `final_scores`) must be finite.
420/// When the `serde` feature is enabled, deserialization rejects non-finite scores.
421#[derive(Clone, Debug, PartialEq)]
422#[cfg_attr(feature = "serde", derive(serde::Serialize))]
423pub struct RerankExecutedPayload {
424    /// Brain profile that served this rerank, if any.
425    pub served_by_profile_id: Option<String>,
426    /// Model used for reranking.
427    pub model_id: Id128,
428    /// Candidate IDs in input order.
429    pub candidates: Vec<Id128>,
430    /// Per-candidate named sub-scores from the reranker.
431    pub reranked: Vec<(Id128, Vec<(String, f32)>)>,
432    /// Final aggregated score per candidate.
433    pub final_scores: Vec<(Id128, f32)>,
434    /// Wall-clock latency of the rerank operation in microseconds.
435    pub latency_us: u64,
436    /// Whether a brain hook was applied during this rerank.
437    pub hook_applied: bool,
438    /// Whether the hook matched the intended target.
439    pub hook_target_match: bool,
440}
441
442impl RerankExecutedPayload {
443    /// Return `true` if all score values are finite.
444    pub fn is_valid(&self) -> bool {
445        let reranked_ok = self
446            .reranked
447            .iter()
448            .all(|(_, scores)| scores.iter().all(|(_, s)| s.is_finite()));
449        let final_ok = self.final_scores.iter().all(|(_, s)| s.is_finite());
450        reranked_ok && final_ok
451    }
452}
453
454#[cfg(feature = "serde")]
455impl<'de> serde::Deserialize<'de> for RerankExecutedPayload {
456    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
457    where
458        D: serde::Deserializer<'de>,
459    {
460        #[derive(serde::Deserialize)]
461        struct Raw {
462            served_by_profile_id: Option<String>,
463            model_id: Id128,
464            candidates: Vec<Id128>,
465            reranked: Vec<(Id128, Vec<(String, f32)>)>,
466            final_scores: Vec<(Id128, f32)>,
467            latency_us: u64,
468            hook_applied: bool,
469            hook_target_match: bool,
470        }
471
472        let raw = Raw::deserialize(deserializer)?;
473
474        for (_, score) in &raw.final_scores {
475            if !score.is_finite() {
476                return Err(serde::de::Error::custom(alloc::format!(
477                    "RerankExecutedPayload final_scores must be finite, got {score}"
478                )));
479            }
480        }
481        for (_, sections) in &raw.reranked {
482            for (section_name, score) in sections {
483                if !score.is_finite() {
484                    return Err(serde::de::Error::custom(alloc::format!(
485                        "RerankExecutedPayload reranked section '{section_name}' score must be finite, got {score}"
486                    )));
487                }
488            }
489        }
490
491        Ok(RerankExecutedPayload {
492            served_by_profile_id: raw.served_by_profile_id,
493            model_id: raw.model_id,
494            candidates: raw.candidates,
495            reranked: raw.reranked,
496            final_scores: raw.final_scores,
497            latency_us: raw.latency_us,
498            hook_applied: raw.hook_applied,
499            hook_target_match: raw.hook_target_match,
500        })
501    }
502}
503
504/// Payload for the `ProposalCreated` event — captures the full initial proposal state.
505#[cfg(feature = "serde")]
506#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
507pub struct ProposalCreatedPayload {
508    pub proposal_id: Id128,
509    pub proposer: String,
510    pub title: String,
511    pub description: String,
512    pub changeset: ProposalChangeset,
513    pub reviewers: Vec<String>,
514    pub expiry: Option<crate::Timestamp>,
515    pub parent_id: Option<Id128>,
516}
517
518/// Structured draft for adding a new entity via a proposal.
519///
520/// Fields mirror the `create(kind=<entity kind>)` verb surface; `kind` is
521/// validated against the closed 8-kind entity taxonomy at apply time.
522#[cfg(feature = "serde")]
523#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
524pub struct EntityDraft {
525    /// Entity kind — must be one of the 8 closed entity kind values.
526    pub kind: String,
527    /// Human-readable name (required).
528    pub name: String,
529    /// Optional long-form description.
530    #[serde(skip_serializing_if = "Option::is_none")]
531    pub description: Option<String>,
532    /// Arbitrary structured metadata.
533    #[serde(skip_serializing_if = "Option::is_none")]
534    pub properties: Option<serde_json::Value>,
535    /// Classification tags.
536    #[serde(default, skip_serializing_if = "Vec::is_empty")]
537    pub tags: Vec<String>,
538}
539
540/// Structured patch for modifying an existing entity via a proposal.
541///
542/// Absent fields mean "leave unchanged". Setting `description` or
543/// `entity_type` to `null` explicitly clears it; a string `entity_type` sets
544/// it, validated against the kind's registered vocabulary at apply time.
545#[cfg(feature = "serde")]
546#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
547pub struct ProposalEntityPatch {
548    #[serde(skip_serializing_if = "Option::is_none")]
549    pub name: Option<String>,
550    /// `null` clears the description; absent leaves it unchanged.
551    #[serde(
552        default,
553        skip_serializing_if = "Option::is_none",
554        with = "serde_opt_opt"
555    )]
556    pub description: Option<Option<String>>,
557    #[serde(skip_serializing_if = "Option::is_none")]
558    pub properties: Option<serde_json::Value>,
559    #[serde(skip_serializing_if = "Option::is_none")]
560    pub tags: Option<Vec<String>>,
561    /// ADR-014 tri-state: absent leaves the type unchanged, `null` explicitly
562    /// clears it, a string sets it (validated against the kind's vocabulary
563    /// at apply time, per ADR-046 parity with the runtime entity update).
564    #[serde(
565        default,
566        skip_serializing_if = "Option::is_none",
567        with = "serde_opt_opt"
568    )]
569    pub entity_type: Option<Option<String>>,
570}
571
572/// Structured draft for adding a new note via a proposal.
573///
574/// Fields mirror the `create(kind=<note kind>)` verb surface.
575#[cfg(feature = "serde")]
576#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
577pub struct NoteDraft {
578    /// Note kind string (validated by the loaded pack at apply time).
579    pub kind: String,
580    /// Note body / content (required).
581    pub content: String,
582    /// Optional short name.
583    #[serde(skip_serializing_if = "Option::is_none")]
584    pub name: Option<String>,
585    /// Arbitrary structured metadata.
586    #[serde(skip_serializing_if = "Option::is_none")]
587    pub properties: Option<serde_json::Value>,
588}
589
590/// Serde helper for `Option<Option<T>>` — distinguishes absent vs. explicit null.
591#[cfg(feature = "serde")]
592mod serde_opt_opt {
593    use serde::{Deserialize, Deserializer, Serialize, Serializer};
594
595    pub fn serialize<T, S>(val: &Option<Option<T>>, s: S) -> Result<S::Ok, S::Error>
596    where
597        T: Serialize,
598        S: Serializer,
599    {
600        match val {
601            None => unreachable!("skip_serializing_if guards the None case"),
602            Some(inner) => inner.serialize(s),
603        }
604    }
605
606    pub fn deserialize<'de, T, D>(d: D) -> Result<Option<Option<T>>, D::Error>
607    where
608        T: Deserialize<'de>,
609        D: Deserializer<'de>,
610    {
611        let opt: Option<T> = Option::deserialize(d)?;
612        Ok(Some(opt))
613    }
614}
615
616/// The set of KG mutations a proposal intends to apply as a proposal changeset.
617#[cfg(feature = "serde")]
618#[derive(Clone, Debug, PartialEq, serde::Serialize)]
619#[serde(tag = "kind", rename_all = "snake_case")]
620pub enum ProposalChangeset {
621    /// Add a new entity. `entity.kind` validated at apply time.
622    AddEntity {
623        entity: EntityDraft,
624    },
625    /// Modify an existing entity's properties / tags / description / entity
626    /// type (absent = unchanged, null = clear, string = set-and-validate).
627    UpdateEntity {
628        id: Id128,
629        patch: ProposalEntityPatch,
630    },
631    /// Add a typed edge. `weight` must be finite and in `[0.0, 1.0]` if present.
632    AddEdge {
633        source: Id128,
634        target: Id128,
635        relation: crate::EdgeRelation,
636        weight: Option<f32>,
637    },
638    /// Add a note (entity-annotating or stand-alone).
639    AddNote {
640        note: NoteDraft,
641    },
642    MergeEntities {
643        into: Id128,
644        from: Id128,
645    },
646    SupersedeEntity {
647        old: Id128,
648        new: Id128,
649    },
650    Compound {
651        steps: Vec<ProposalChangeset>,
652    },
653}
654
655#[cfg(feature = "serde")]
656impl ProposalChangeset {
657    fn validate(&self) -> Result<(), alloc::string::String> {
658        match self {
659            Self::AddEdge { weight, .. } => {
660                if let Some(w) = weight {
661                    if !w.is_finite() {
662                        return Err(alloc::format!(
663                            "ProposalChangeset AddEdge weight must be finite, got {w}"
664                        ));
665                    }
666                    if !(*w >= 0.0 && *w <= 1.0) {
667                        return Err(alloc::format!(
668                            "ProposalChangeset AddEdge weight must be in [0.0, 1.0], got {w}"
669                        ));
670                    }
671                }
672                Ok(())
673            }
674            Self::Compound { steps } => {
675                for step in steps {
676                    step.validate()?;
677                }
678                Ok(())
679            }
680            _ => Ok(()),
681        }
682    }
683}
684
685#[cfg(feature = "serde")]
686impl<'de> serde::Deserialize<'de> for ProposalChangeset {
687    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
688    where
689        D: serde::Deserializer<'de>,
690    {
691        #[derive(serde::Deserialize)]
692        #[serde(tag = "kind", rename_all = "snake_case")]
693        enum ProposalChangesetRaw {
694            AddEntity {
695                entity: EntityDraft,
696            },
697            UpdateEntity {
698                id: Id128,
699                patch: ProposalEntityPatch,
700            },
701            AddEdge {
702                source: Id128,
703                target: Id128,
704                relation: crate::EdgeRelation,
705                weight: Option<f32>,
706            },
707            AddNote {
708                note: NoteDraft,
709            },
710            MergeEntities {
711                into: Id128,
712                from: Id128,
713            },
714            SupersedeEntity {
715                old: Id128,
716                new: Id128,
717            },
718            Compound {
719                steps: Vec<ProposalChangeset>,
720            },
721        }
722
723        let raw = ProposalChangesetRaw::deserialize(deserializer)?;
724        let cs = match raw {
725            ProposalChangesetRaw::AddEntity { entity } => Self::AddEntity { entity },
726            ProposalChangesetRaw::UpdateEntity { id, patch } => Self::UpdateEntity { id, patch },
727            ProposalChangesetRaw::AddEdge {
728                source,
729                target,
730                relation,
731                weight,
732            } => Self::AddEdge {
733                source,
734                target,
735                relation,
736                weight,
737            },
738            ProposalChangesetRaw::AddNote { note } => Self::AddNote { note },
739            ProposalChangesetRaw::MergeEntities { into, from } => {
740                Self::MergeEntities { into, from }
741            }
742            ProposalChangesetRaw::SupersedeEntity { old, new } => {
743                Self::SupersedeEntity { old, new }
744            }
745            ProposalChangesetRaw::Compound { steps } => Self::Compound { steps },
746        };
747        cs.validate().map_err(serde::de::Error::custom)?;
748        Ok(cs)
749    }
750}
751
752#[cfg(not(feature = "serde"))]
753#[derive(Clone, Debug, PartialEq)]
754pub enum ProposalChangeset {
755    AddEdge {
756        source: Id128,
757        target: Id128,
758        relation: crate::EdgeRelation,
759        weight: Option<f32>,
760    },
761    MergeEntities {
762        into: Id128,
763        from: Id128,
764    },
765    SupersedeEntity {
766        old: Id128,
767        new: Id128,
768    },
769    Compound {
770        steps: Vec<ProposalChangeset>,
771    },
772}
773
774/// Payload for the `ProposalReviewed` event — records a single reviewer's decision.
775#[derive(Clone, Debug, PartialEq)]
776#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
777pub struct ProposalReviewedPayload {
778    pub proposal_id: Id128,
779    pub reviewer: String,
780    pub decision: ProposalDecision,
781    pub comment: Option<String>,
782}
783
784/// A reviewer's decision on a proposal.
785#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
786#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
787#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
788pub enum ProposalDecision {
789    /// The reviewer approved the proposal for application.
790    Approve,
791    /// The reviewer rejected the proposal; it will not be applied.
792    Reject,
793    /// The reviewer left a comment without blocking the proposal.
794    Comment,
795    /// The reviewer requested changes before the proposal can proceed.
796    RequestChanges,
797}
798
799impl ProposalDecision {
800    /// Returns the bare variant name as a lowercase string, matching the serde
801    /// `rename_all = "snake_case"` representation.  Use this when storing the
802    /// decision as a plain TEXT column — **not** `serde_json::to_string`, which
803    /// would produce a JSON-quoted string (`"\"approve\""` instead of `"approve"`).
804    pub fn as_str(self) -> &'static str {
805        match self {
806            Self::Approve => "approve",
807            Self::Reject => "reject",
808            Self::Comment => "comment",
809            Self::RequestChanges => "request_changes",
810        }
811    }
812}
813
814/// Payload for the `ProposalApplied` event — records the outcome of the apply attempt.
815#[derive(Clone, Debug, PartialEq)]
816#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
817pub struct ProposalAppliedPayload {
818    pub proposal_id: Id128,
819    pub applied_at: crate::Timestamp,
820    pub applied_by: String,
821    pub result: ApplyResult,
822}
823
824/// Outcome of applying a proposal: either all steps succeeded or the apply failed with an error.
825#[derive(Clone, Debug, PartialEq)]
826#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
827#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
828pub enum ApplyResult {
829    Success {
830        created_records: Vec<Id128>,
831    },
832    Failed {
833        error: String,
834        applied_step_count: u32,
835    },
836}
837
838/// Payload for the `ProposalWithdrawn` event — records who withdrew and an optional reason.
839#[derive(Clone, Debug, PartialEq)]
840#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
841pub struct ProposalWithdrawnPayload {
842    pub proposal_id: Id128,
843    pub by: String,
844    pub reason: Option<String>,
845}
846
847/// Builder for events. Used by the verb dispatch path.
848pub struct EventBuilder {
849    verb: String,
850    substrate: SubstrateKind,
851    actor: Option<String>,
852    kind: EventKind,
853    payload: EventPayload,
854    payload_schema_version: u32,
855    profile_state_version: Option<u64>,
856    aggregate: Option<AggregateRef>,
857}
858
859impl EventBuilder {
860    /// Create a new builder for an event produced by `verb` acting on `substrate` as `actor`.
861    pub fn new(
862        verb: impl Into<String>,
863        substrate: SubstrateKind,
864        actor: impl Into<String>,
865    ) -> Self {
866        Self {
867            verb: verb.into(),
868            substrate,
869            actor: Some(actor.into()),
870            kind: EventKind::Audit,
871            payload: EventPayload::default(),
872            payload_schema_version: 1,
873            profile_state_version: None,
874            aggregate: None,
875        }
876    }
877
878    /// Override the event kind discriminant.
879    pub fn kind(mut self, kind: EventKind) -> Self {
880        self.kind = kind;
881        self
882    }
883
884    /// Set the typed payload for this event.
885    pub fn payload(mut self, payload: EventPayload) -> Self {
886        self.payload = payload;
887        self
888    }
889
890    /// Set the payload schema version (defaults to 1).
891    pub fn payload_schema_version(mut self, version: u32) -> Self {
892        self.payload_schema_version = version;
893        self
894    }
895
896    /// Record the brain profile state version observed at emit time.
897    pub fn profile_state_version(mut self, version: u64) -> Self {
898        self.profile_state_version = Some(version);
899        self
900    }
901
902    /// Thread this event into an aggregate chain.
903    pub fn aggregate(mut self, aggregate: AggregateRef) -> Self {
904        self.aggregate = Some(aggregate);
905        self
906    }
907
908    /// Consume the builder and produce an [`Event`] with the given `header`.
909    pub fn build(self, header: Header) -> Event {
910        Event {
911            header,
912            verb: self.verb,
913            substrate: self.substrate,
914            actor: self.actor,
915            kind: self.kind,
916            payload: self.payload,
917            payload_schema_version: self.payload_schema_version,
918            profile_state_version: self.profile_state_version,
919            aggregate: self.aggregate,
920        }
921    }
922}
923
924#[cfg(test)]
925mod tests {
926    extern crate alloc;
927
928    use super::*;
929    use crate::{Namespace, Timestamp};
930    #[cfg(feature = "serde")]
931    use alloc::string::ToString;
932
933    fn header() -> Header {
934        Header::new(
935            Id128::from_u128(1),
936            Namespace::local(),
937            Timestamp::from_secs(1700000000),
938        )
939    }
940
941    #[test]
942    fn event_kind_parse_roundtrip() {
943        for kind in EventKind::ALL {
944            let parsed: EventKind = kind
945                .name()
946                .parse()
947                .expect("EventKind::name must parse back");
948            assert_eq!(parsed, kind);
949        }
950    }
951
952    #[cfg(feature = "serde")]
953    #[test]
954    fn refusal_kind_has_canonical_string_and_serde_roundtrip() {
955        let kind = EventKind::Refusal;
956        assert!(EventKind::ALL.contains(&kind));
957        assert_eq!(kind.name(), "refusal");
958        assert_eq!(kind.to_string(), "refusal");
959        assert_eq!("refusal".parse::<EventKind>().unwrap(), kind);
960        assert_eq!(serde_json::to_string(&kind).unwrap(), "\"refusal\"");
961        assert_eq!(
962            serde_json::from_str::<EventKind>("\"refusal\"").unwrap(),
963            kind
964        );
965    }
966
967    #[test]
968    fn rerank_payload_records_served_profile() {
969        let payload = EventPayload::RerankExecuted(RerankExecutedPayload {
970            served_by_profile_id: Some("profile-a".into()),
971            model_id: Id128::from_u128(1),
972            candidates: Vec::new(),
973            reranked: Vec::new(),
974            final_scores: Vec::new(),
975            latency_us: 100,
976            hook_applied: false,
977            hook_target_match: false,
978        });
979        let event = EventBuilder::new("rerank", SubstrateKind::Note, "agent:test")
980            .kind(EventKind::RerankExecuted)
981            .payload(payload)
982            .build(header());
983
984        if let EventPayload::RerankExecuted(ref p) = event.payload {
985            assert_eq!(p.served_by_profile_id.as_deref(), Some("profile-a"));
986        } else {
987            panic!("unexpected payload variant");
988        }
989    }
990
991    #[test]
992    fn proposal_payloads_are_typed() {
993        let payload = EventPayload::ProposalReviewed(ProposalReviewedPayload {
994            proposal_id: Id128::from_u128(42),
995            reviewer: "operator".into(),
996            decision: ProposalDecision::Approve,
997            comment: None,
998        });
999        let event = EventBuilder::new("review", SubstrateKind::Entity, "operator")
1000            .kind(EventKind::ProposalReviewed)
1001            .payload(payload)
1002            .build(header());
1003        assert_eq!(event.kind.name(), "proposal_reviewed");
1004    }
1005
1006    /// C1 regression: all ProposalChangeset variants that carry Id128 fields must
1007    /// round-trip through serde_json::Value.  Previously `Id128::deserialize` used
1008    /// `<&str>::deserialize` which fails when the deserializer holds owned data
1009    /// (the Value-backed path used by the MCP DSL parser).
1010    #[cfg(feature = "serde")]
1011    #[test]
1012    fn proposal_changeset_id_variants_deserialize_from_value() {
1013        let uuid = "7426afd6-0234-4701-9045-83dfd39166e6";
1014        let uuid2 = "abcdef01-2345-6789-abcd-ef0123456789";
1015
1016        // UpdateEntity — patch is now a structured ProposalEntityPatch object
1017        let v =
1018            serde_json::json!({"kind": "update_entity", "id": uuid, "patch": {"name": "NewName"}});
1019        let cs: ProposalChangeset =
1020            serde_json::from_value(v).expect("UpdateEntity must deserialize from Value");
1021        assert!(
1022            matches!(cs, ProposalChangeset::UpdateEntity { .. }),
1023            "expected UpdateEntity"
1024        );
1025
1026        // Tri-state entity_type (ADR-014): absent vs explicit null vs set
1027        // must all survive the proposal wire boundary distinctly.
1028        for (json, expected) in [
1029            (serde_json::json!({}), None),
1030            (serde_json::json!({"entity_type": null}), Some(None)),
1031            (
1032                serde_json::json!({"entity_type": "algorithm"}),
1033                Some(Some("algorithm".to_string())),
1034            ),
1035        ] {
1036            let v = serde_json::json!({"kind": "update_entity", "id": uuid, "patch": json});
1037            let cs: ProposalChangeset =
1038                serde_json::from_value(v).expect("UpdateEntity must deserialize");
1039            let ProposalChangeset::UpdateEntity { patch, .. } = cs else {
1040                panic!("expected UpdateEntity");
1041            };
1042            assert_eq!(patch.entity_type, expected, "patch: {json}");
1043        }
1044
1045        // The clear must serialize back as an explicit null.
1046        let patch = ProposalEntityPatch {
1047            name: None,
1048            description: None,
1049            properties: None,
1050            tags: None,
1051            entity_type: Some(None),
1052        };
1053        let v = serde_json::to_value(&patch).expect("serialize");
1054        assert_eq!(v.get("entity_type"), Some(&serde_json::Value::Null));
1055        assert!(
1056            v.get("name").is_none(),
1057            "absent fields must not be serialized"
1058        );
1059
1060        // AddEdge
1061        let v = serde_json::json!({
1062            "kind": "add_edge",
1063            "source": uuid, "target": uuid2,
1064            "relation": "extends", "weight": 1.0
1065        });
1066        let cs: ProposalChangeset =
1067            serde_json::from_value(v).expect("AddEdge must deserialize from Value");
1068        assert!(
1069            matches!(cs, ProposalChangeset::AddEdge { .. }),
1070            "expected AddEdge"
1071        );
1072
1073        // MergeEntities
1074        let v = serde_json::json!({"kind": "merge_entities", "into": uuid, "from": uuid2});
1075        let cs: ProposalChangeset =
1076            serde_json::from_value(v).expect("MergeEntities must deserialize from Value");
1077        assert!(
1078            matches!(cs, ProposalChangeset::MergeEntities { .. }),
1079            "expected MergeEntities"
1080        );
1081
1082        // SupersedeEntity
1083        let v = serde_json::json!({"kind": "supersede_entity", "old": uuid, "new": uuid2});
1084        let cs: ProposalChangeset =
1085            serde_json::from_value(v).expect("SupersedeEntity must deserialize from Value");
1086        assert!(
1087            matches!(cs, ProposalChangeset::SupersedeEntity { .. }),
1088            "expected SupersedeEntity"
1089        );
1090    }
1091
1092    #[cfg(feature = "serde")]
1093    #[test]
1094    fn proposal_changeset_rejects_invalid_edge_weight() {
1095        let uuid = "7426afd6-0234-4701-9045-83dfd39166e6";
1096        let uuid2 = "abcdef01-2345-6789-abcd-ef0123456789";
1097
1098        let v = serde_json::json!({
1099            "kind": "add_edge",
1100            "source": uuid, "target": uuid2,
1101            "relation": "extends", "weight": 2.0
1102        });
1103        let result: Result<ProposalChangeset, _> = serde_json::from_value(v);
1104        assert!(result.is_err());
1105        let err = result.unwrap_err().to_string();
1106        assert!(
1107            err.contains("[0.0, 1.0]"),
1108            "error should mention range: {err}"
1109        );
1110    }
1111
1112    #[cfg(feature = "serde")]
1113    #[test]
1114    fn proposal_changeset_accepts_null_edge_weight() {
1115        let uuid = "7426afd6-0234-4701-9045-83dfd39166e6";
1116        let uuid2 = "abcdef01-2345-6789-abcd-ef0123456789";
1117
1118        let v = serde_json::json!({
1119            "kind": "add_edge",
1120            "source": uuid, "target": uuid2,
1121            "relation": "extends", "weight": null
1122        });
1123        let cs: ProposalChangeset =
1124            serde_json::from_value(v).expect("null weight should be accepted");
1125        assert!(matches!(
1126            cs,
1127            ProposalChangeset::AddEdge { weight: None, .. }
1128        ));
1129    }
1130
1131    #[cfg(feature = "serde")]
1132    #[test]
1133    fn rerank_payload_serde_rejects_non_finite_score() {
1134        let json = serde_json::json!({
1135            "served_by_profile_id": null,
1136            "model_id": "00000000-0000-0000-0000-000000000001",
1137            "candidates": [],
1138            "reranked": [],
1139            "final_scores": [["00000000-0000-0000-0000-000000000001", "Infinity"]],
1140            "latency_us": 100,
1141            "hook_applied": false,
1142            "hook_target_match": false
1143        });
1144        let result: Result<RerankExecutedPayload, _> = serde_json::from_value(json);
1145        assert!(result.is_err());
1146    }
1147
1148    #[test]
1149    fn rerank_payload_is_valid_checks_finite() {
1150        let p = RerankExecutedPayload {
1151            served_by_profile_id: None,
1152            model_id: Id128::from_u128(1),
1153            candidates: Vec::new(),
1154            reranked: Vec::new(),
1155            final_scores: alloc::vec![(Id128::from_u128(1), 0.5)],
1156            latency_us: 100,
1157            hook_applied: false,
1158            hook_target_match: false,
1159        };
1160        assert!(p.is_valid());
1161
1162        let p_inf = RerankExecutedPayload {
1163            served_by_profile_id: None,
1164            model_id: Id128::from_u128(1),
1165            candidates: Vec::new(),
1166            reranked: Vec::new(),
1167            final_scores: alloc::vec![(Id128::from_u128(1), f32::INFINITY)],
1168            latency_us: 100,
1169            hook_applied: false,
1170            hook_target_match: false,
1171        };
1172        assert!(!p_inf.is_valid());
1173    }
1174}