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