1extern crate alloc;
4use alloc::string::String;
5use alloc::vec::Vec;
6use core::fmt;
7
8use crate::{Header, Id128, SubstrateKind};
9
10#[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 pub verb: String,
18 pub substrate: SubstrateKind,
20 pub actor: Option<String>,
22 pub kind: EventKind,
24 pub payload: EventPayload,
26 pub payload_schema_version: u32,
28 pub profile_state_version: Option<u64>,
30 pub aggregate: Option<AggregateRef>,
32}
33
34#[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 #[default]
41 Success,
42 Denied,
44 Error,
46}
47
48impl EventOutcome {
49 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#[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 Audit,
74 RecallExecuted,
76 RerankExecuted,
78 SearchExecuted,
80 LinkCreated,
82 EntityCreated,
84 EntityUpdated,
86 EntityDeleted,
88 EntityMerged,
90 NoteMerged,
92 NoteCreated,
94 NoteUpdated,
96 NoteDeleted,
98 EdgeUpdated,
100 EdgeDeleted,
102 TaskTransitioned,
104 FeedbackExplicit,
106 ProfileResolutionRecommended,
108 ProfileMerged,
110 EmbeddingModelChanged,
112 EmbeddingMigrationCompleted,
114 EmbeddingMigrationFailed,
116 EmbeddingDriftDetected,
118 EmbedderInitialized,
120 ProposalCreated,
122 ProposalReviewed,
124 ProposalApplied,
126 ProposalWithdrawn,
128 ChannelPollStarted,
130 ChannelPollSucceeded,
132 ChannelPollFailed,
134 ChannelBackoffArmed,
136 ChannelBackoffReset,
138 ChannelHeartbeatPersistFailed,
140 ConfigLocked,
142 CheckpointOutcomeRecorded,
144 PhaseStarted,
146 PhaseCompleted,
148 PhaseCancelled,
150}
151
152impl EventKind {
153 pub const ALL: [Self; 39] = [
155 Self::Audit,
156 Self::RecallExecuted,
157 Self::RerankExecuted,
158 Self::SearchExecuted,
159 Self::LinkCreated,
160 Self::EntityCreated,
161 Self::EntityUpdated,
162 Self::EntityDeleted,
163 Self::EntityMerged,
164 Self::NoteMerged,
165 Self::NoteCreated,
166 Self::NoteUpdated,
167 Self::NoteDeleted,
168 Self::EdgeUpdated,
169 Self::EdgeDeleted,
170 Self::TaskTransitioned,
171 Self::FeedbackExplicit,
172 Self::ProfileResolutionRecommended,
173 Self::ProfileMerged,
174 Self::EmbeddingModelChanged,
175 Self::EmbeddingMigrationCompleted,
176 Self::EmbeddingMigrationFailed,
177 Self::EmbeddingDriftDetected,
178 Self::EmbedderInitialized,
179 Self::ProposalCreated,
180 Self::ProposalReviewed,
181 Self::ProposalApplied,
182 Self::ProposalWithdrawn,
183 Self::ChannelPollStarted,
184 Self::ChannelPollSucceeded,
185 Self::ChannelPollFailed,
186 Self::ChannelBackoffArmed,
187 Self::ChannelBackoffReset,
188 Self::ChannelHeartbeatPersistFailed,
189 Self::ConfigLocked,
190 Self::CheckpointOutcomeRecorded,
191 Self::PhaseStarted,
192 Self::PhaseCompleted,
193 Self::PhaseCancelled,
194 ];
195
196 pub const fn name(self) -> &'static str {
198 match self {
199 Self::Audit => "audit",
200 Self::RecallExecuted => "recall_executed",
201 Self::RerankExecuted => "rerank_executed",
202 Self::SearchExecuted => "search_executed",
203 Self::LinkCreated => "link_created",
204 Self::EntityCreated => "entity_created",
205 Self::EntityUpdated => "entity_updated",
206 Self::EntityDeleted => "entity_deleted",
207 Self::EntityMerged => "entity_merged",
208 Self::NoteMerged => "note_merged",
209 Self::NoteCreated => "note_created",
210 Self::NoteUpdated => "note_updated",
211 Self::NoteDeleted => "note_deleted",
212 Self::EdgeUpdated => "edge_updated",
213 Self::EdgeDeleted => "edge_deleted",
214 Self::TaskTransitioned => "task_transitioned",
215 Self::FeedbackExplicit => "feedback_explicit",
216 Self::ProfileResolutionRecommended => "profile_resolution_recommended",
217 Self::ProfileMerged => "profile_merged",
218 Self::EmbeddingModelChanged => "embedding_model_changed",
219 Self::EmbeddingMigrationCompleted => "embedding_migration_completed",
220 Self::EmbeddingMigrationFailed => "embedding_migration_failed",
221 Self::EmbeddingDriftDetected => "embedding_drift_detected",
222 Self::EmbedderInitialized => "embedder_initialized",
223 Self::ProposalCreated => "proposal_created",
224 Self::ProposalReviewed => "proposal_reviewed",
225 Self::ProposalApplied => "proposal_applied",
226 Self::ProposalWithdrawn => "proposal_withdrawn",
227 Self::ChannelPollStarted => "channel_poll_started",
228 Self::ChannelPollSucceeded => "channel_poll_succeeded",
229 Self::ChannelPollFailed => "channel_poll_failed",
230 Self::ChannelBackoffArmed => "channel_backoff_armed",
231 Self::ChannelBackoffReset => "channel_backoff_reset",
232 Self::ChannelHeartbeatPersistFailed => "channel_heartbeat_persist_failed",
233 Self::ConfigLocked => "config_locked",
234 Self::CheckpointOutcomeRecorded => "checkpoint_outcome_recorded",
235 Self::PhaseStarted => "phase_started",
236 Self::PhaseCompleted => "phase_completed",
237 Self::PhaseCancelled => "phase_cancelled",
238 }
239 }
240}
241
242impl fmt::Display for EventKind {
243 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
244 f.write_str(self.name())
245 }
246}
247
248const EVENT_KIND_VALID: &[&str] = &[
249 "audit",
250 "recall_executed",
251 "rerank_executed",
252 "search_executed",
253 "link_created",
254 "entity_created",
255 "entity_updated",
256 "entity_deleted",
257 "entity_merged",
258 "note_merged",
259 "note_created",
260 "note_updated",
261 "note_deleted",
262 "edge_updated",
263 "edge_deleted",
264 "task_transitioned",
265 "feedback_explicit",
266 "profile_resolution_recommended",
267 "profile_merged",
268 "embedding_model_changed",
269 "embedding_migration_completed",
270 "embedding_migration_failed",
271 "embedding_drift_detected",
272 "embedder_initialized",
273 "proposal_created",
274 "proposal_reviewed",
275 "proposal_applied",
276 "proposal_withdrawn",
277 "channel_poll_started",
278 "channel_poll_succeeded",
279 "channel_poll_failed",
280 "channel_backoff_armed",
281 "channel_backoff_reset",
282 "channel_heartbeat_persist_failed",
283 "config_locked",
284 "checkpoint_outcome_recorded",
285 "phase_started",
286 "phase_completed",
287 "phase_cancelled",
288];
289
290impl core::str::FromStr for EventKind {
291 type Err = crate::error::UnknownVariant;
292
293 fn from_str(s: &str) -> Result<Self, Self::Err> {
294 match s.trim().to_ascii_lowercase().as_str() {
295 "audit" => Ok(Self::Audit),
296 "recall_executed" => Ok(Self::RecallExecuted),
297 "rerank_executed" => Ok(Self::RerankExecuted),
298 "search_executed" => Ok(Self::SearchExecuted),
299 "link_created" => Ok(Self::LinkCreated),
300 "entity_created" => Ok(Self::EntityCreated),
301 "entity_updated" => Ok(Self::EntityUpdated),
302 "entity_deleted" => Ok(Self::EntityDeleted),
303 "entity_merged" => Ok(Self::EntityMerged),
304 "note_merged" => Ok(Self::NoteMerged),
305 "note_created" => Ok(Self::NoteCreated),
306 "note_updated" => Ok(Self::NoteUpdated),
307 "note_deleted" => Ok(Self::NoteDeleted),
308 "edge_updated" => Ok(Self::EdgeUpdated),
309 "edge_deleted" => Ok(Self::EdgeDeleted),
310 "task_transitioned" => Ok(Self::TaskTransitioned),
311 "feedback_explicit" => Ok(Self::FeedbackExplicit),
312 "profile_resolution_recommended" => Ok(Self::ProfileResolutionRecommended),
313 "profile_merged" => Ok(Self::ProfileMerged),
314 "embedding_model_changed" => Ok(Self::EmbeddingModelChanged),
315 "embedding_migration_completed" => Ok(Self::EmbeddingMigrationCompleted),
316 "embedding_migration_failed" => Ok(Self::EmbeddingMigrationFailed),
317 "embedding_drift_detected" => Ok(Self::EmbeddingDriftDetected),
318 "embedder_initialized" => Ok(Self::EmbedderInitialized),
319 "proposal_created" => Ok(Self::ProposalCreated),
320 "proposal_reviewed" => Ok(Self::ProposalReviewed),
321 "proposal_applied" => Ok(Self::ProposalApplied),
322 "proposal_withdrawn" => Ok(Self::ProposalWithdrawn),
323 "channel_poll_started" => Ok(Self::ChannelPollStarted),
324 "channel_poll_succeeded" => Ok(Self::ChannelPollSucceeded),
325 "channel_poll_failed" => Ok(Self::ChannelPollFailed),
326 "channel_backoff_armed" => Ok(Self::ChannelBackoffArmed),
327 "channel_backoff_reset" => Ok(Self::ChannelBackoffReset),
328 "channel_heartbeat_persist_failed" => Ok(Self::ChannelHeartbeatPersistFailed),
329 "config_locked" => Ok(Self::ConfigLocked),
330 "checkpoint_outcome_recorded" => Ok(Self::CheckpointOutcomeRecorded),
331 "phase_started" => Ok(Self::PhaseStarted),
332 "phase_completed" => Ok(Self::PhaseCompleted),
333 "phase_cancelled" => Ok(Self::PhaseCancelled),
334 other => Err(crate::error::UnknownVariant::new(
335 "event_kind",
336 other,
337 EVENT_KIND_VALID,
338 )),
339 }
340 }
341}
342
343#[derive(Clone, Debug, PartialEq, Eq)]
348#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
349pub struct AggregateRef {
350 pub kind: String,
352 pub id: Id128,
354}
355
356#[derive(Clone, Debug, PartialEq)]
362#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
363#[cfg_attr(
364 feature = "serde",
365 serde(tag = "kind", content = "payload", rename_all = "snake_case")
366)]
367pub enum EventPayload {
368 Json(String),
370 RerankExecuted(RerankExecutedPayload),
372 #[cfg(feature = "serde")]
374 ProposalCreated(ProposalCreatedPayload),
375 ProposalReviewed(ProposalReviewedPayload),
377 ProposalApplied(ProposalAppliedPayload),
379 ProposalWithdrawn(ProposalWithdrawnPayload),
381}
382
383impl Default for EventPayload {
384 fn default() -> Self {
385 Self::Json("{}".into())
386 }
387}
388
389#[derive(Clone, Debug, PartialEq)]
394#[cfg_attr(feature = "serde", derive(serde::Serialize))]
395pub struct RerankExecutedPayload {
396 pub served_by_profile_id: Option<String>,
398 pub model_id: Id128,
400 pub candidates: Vec<Id128>,
402 pub reranked: Vec<(Id128, Vec<(String, f32)>)>,
404 pub final_scores: Vec<(Id128, f32)>,
406 pub latency_us: u64,
408 pub hook_applied: bool,
410 pub hook_target_match: bool,
412}
413
414impl RerankExecutedPayload {
415 pub fn is_valid(&self) -> bool {
417 let reranked_ok = self
418 .reranked
419 .iter()
420 .all(|(_, scores)| scores.iter().all(|(_, s)| s.is_finite()));
421 let final_ok = self.final_scores.iter().all(|(_, s)| s.is_finite());
422 reranked_ok && final_ok
423 }
424}
425
426#[cfg(feature = "serde")]
427impl<'de> serde::Deserialize<'de> for RerankExecutedPayload {
428 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
429 where
430 D: serde::Deserializer<'de>,
431 {
432 #[derive(serde::Deserialize)]
433 struct Raw {
434 served_by_profile_id: Option<String>,
435 model_id: Id128,
436 candidates: Vec<Id128>,
437 reranked: Vec<(Id128, Vec<(String, f32)>)>,
438 final_scores: Vec<(Id128, f32)>,
439 latency_us: u64,
440 hook_applied: bool,
441 hook_target_match: bool,
442 }
443
444 let raw = Raw::deserialize(deserializer)?;
445
446 for (_, score) in &raw.final_scores {
447 if !score.is_finite() {
448 return Err(serde::de::Error::custom(alloc::format!(
449 "RerankExecutedPayload final_scores must be finite, got {score}"
450 )));
451 }
452 }
453 for (_, sections) in &raw.reranked {
454 for (section_name, score) in sections {
455 if !score.is_finite() {
456 return Err(serde::de::Error::custom(alloc::format!(
457 "RerankExecutedPayload reranked section '{section_name}' score must be finite, got {score}"
458 )));
459 }
460 }
461 }
462
463 Ok(RerankExecutedPayload {
464 served_by_profile_id: raw.served_by_profile_id,
465 model_id: raw.model_id,
466 candidates: raw.candidates,
467 reranked: raw.reranked,
468 final_scores: raw.final_scores,
469 latency_us: raw.latency_us,
470 hook_applied: raw.hook_applied,
471 hook_target_match: raw.hook_target_match,
472 })
473 }
474}
475
476#[cfg(feature = "serde")]
478#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
479pub struct ProposalCreatedPayload {
480 pub proposal_id: Id128,
481 pub proposer: String,
482 pub title: String,
483 pub description: String,
484 pub changeset: ProposalChangeset,
485 pub reviewers: Vec<String>,
486 pub expiry: Option<crate::Timestamp>,
487 pub parent_id: Option<Id128>,
488}
489
490#[cfg(feature = "serde")]
495#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
496pub struct EntityDraft {
497 pub kind: String,
499 pub name: String,
501 #[serde(skip_serializing_if = "Option::is_none")]
503 pub description: Option<String>,
504 #[serde(skip_serializing_if = "Option::is_none")]
506 pub properties: Option<serde_json::Value>,
507 #[serde(default, skip_serializing_if = "Vec::is_empty")]
509 pub tags: Vec<String>,
510}
511
512#[cfg(feature = "serde")]
518#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
519pub struct ProposalEntityPatch {
520 #[serde(skip_serializing_if = "Option::is_none")]
521 pub name: Option<String>,
522 #[serde(
524 default,
525 skip_serializing_if = "Option::is_none",
526 with = "serde_opt_opt"
527 )]
528 pub description: Option<Option<String>>,
529 #[serde(skip_serializing_if = "Option::is_none")]
530 pub properties: Option<serde_json::Value>,
531 #[serde(skip_serializing_if = "Option::is_none")]
532 pub tags: Option<Vec<String>>,
533 #[serde(
537 default,
538 skip_serializing_if = "Option::is_none",
539 with = "serde_opt_opt"
540 )]
541 pub entity_type: Option<Option<String>>,
542}
543
544#[cfg(feature = "serde")]
548#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
549pub struct NoteDraft {
550 pub kind: String,
552 pub content: String,
554 #[serde(skip_serializing_if = "Option::is_none")]
556 pub name: Option<String>,
557 #[serde(skip_serializing_if = "Option::is_none")]
559 pub properties: Option<serde_json::Value>,
560}
561
562#[cfg(feature = "serde")]
564mod serde_opt_opt {
565 use serde::{Deserialize, Deserializer, Serialize, Serializer};
566
567 pub fn serialize<T, S>(val: &Option<Option<T>>, s: S) -> Result<S::Ok, S::Error>
568 where
569 T: Serialize,
570 S: Serializer,
571 {
572 match val {
573 None => unreachable!("skip_serializing_if guards the None case"),
574 Some(inner) => inner.serialize(s),
575 }
576 }
577
578 pub fn deserialize<'de, T, D>(d: D) -> Result<Option<Option<T>>, D::Error>
579 where
580 T: Deserialize<'de>,
581 D: Deserializer<'de>,
582 {
583 let opt: Option<T> = Option::deserialize(d)?;
584 Ok(Some(opt))
585 }
586}
587
588#[cfg(feature = "serde")]
590#[derive(Clone, Debug, PartialEq, serde::Serialize)]
591#[serde(tag = "kind", rename_all = "snake_case")]
592pub enum ProposalChangeset {
593 AddEntity {
595 entity: EntityDraft,
596 },
597 UpdateEntity {
600 id: Id128,
601 patch: ProposalEntityPatch,
602 },
603 AddEdge {
605 source: Id128,
606 target: Id128,
607 relation: crate::EdgeRelation,
608 weight: Option<f32>,
609 },
610 AddNote {
612 note: NoteDraft,
613 },
614 MergeEntities {
615 into: Id128,
616 from: Id128,
617 },
618 SupersedeEntity {
619 old: Id128,
620 new: Id128,
621 },
622 Compound {
623 steps: Vec<ProposalChangeset>,
624 },
625}
626
627#[cfg(feature = "serde")]
628impl ProposalChangeset {
629 fn validate(&self) -> Result<(), alloc::string::String> {
630 match self {
631 Self::AddEdge { weight, .. } => {
632 if let Some(w) = weight {
633 if !w.is_finite() {
634 return Err(alloc::format!(
635 "ProposalChangeset AddEdge weight must be finite, got {w}"
636 ));
637 }
638 if !(*w >= 0.0 && *w <= 1.0) {
639 return Err(alloc::format!(
640 "ProposalChangeset AddEdge weight must be in [0.0, 1.0], got {w}"
641 ));
642 }
643 }
644 Ok(())
645 }
646 Self::Compound { steps } => {
647 for step in steps {
648 step.validate()?;
649 }
650 Ok(())
651 }
652 _ => Ok(()),
653 }
654 }
655}
656
657#[cfg(feature = "serde")]
658impl<'de> serde::Deserialize<'de> for ProposalChangeset {
659 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
660 where
661 D: serde::Deserializer<'de>,
662 {
663 #[derive(serde::Deserialize)]
664 #[serde(tag = "kind", rename_all = "snake_case")]
665 enum ProposalChangesetRaw {
666 AddEntity {
667 entity: EntityDraft,
668 },
669 UpdateEntity {
670 id: Id128,
671 patch: ProposalEntityPatch,
672 },
673 AddEdge {
674 source: Id128,
675 target: Id128,
676 relation: crate::EdgeRelation,
677 weight: Option<f32>,
678 },
679 AddNote {
680 note: NoteDraft,
681 },
682 MergeEntities {
683 into: Id128,
684 from: Id128,
685 },
686 SupersedeEntity {
687 old: Id128,
688 new: Id128,
689 },
690 Compound {
691 steps: Vec<ProposalChangeset>,
692 },
693 }
694
695 let raw = ProposalChangesetRaw::deserialize(deserializer)?;
696 let cs = match raw {
697 ProposalChangesetRaw::AddEntity { entity } => Self::AddEntity { entity },
698 ProposalChangesetRaw::UpdateEntity { id, patch } => Self::UpdateEntity { id, patch },
699 ProposalChangesetRaw::AddEdge {
700 source,
701 target,
702 relation,
703 weight,
704 } => Self::AddEdge {
705 source,
706 target,
707 relation,
708 weight,
709 },
710 ProposalChangesetRaw::AddNote { note } => Self::AddNote { note },
711 ProposalChangesetRaw::MergeEntities { into, from } => {
712 Self::MergeEntities { into, from }
713 }
714 ProposalChangesetRaw::SupersedeEntity { old, new } => {
715 Self::SupersedeEntity { old, new }
716 }
717 ProposalChangesetRaw::Compound { steps } => Self::Compound { steps },
718 };
719 cs.validate().map_err(serde::de::Error::custom)?;
720 Ok(cs)
721 }
722}
723
724#[cfg(not(feature = "serde"))]
725#[derive(Clone, Debug, PartialEq)]
726pub enum ProposalChangeset {
727 AddEdge {
728 source: Id128,
729 target: Id128,
730 relation: crate::EdgeRelation,
731 weight: Option<f32>,
732 },
733 MergeEntities {
734 into: Id128,
735 from: Id128,
736 },
737 SupersedeEntity {
738 old: Id128,
739 new: Id128,
740 },
741 Compound {
742 steps: Vec<ProposalChangeset>,
743 },
744}
745
746#[derive(Clone, Debug, PartialEq)]
748#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
749pub struct ProposalReviewedPayload {
750 pub proposal_id: Id128,
751 pub reviewer: String,
752 pub decision: ProposalDecision,
753 pub comment: Option<String>,
754}
755
756#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
758#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
759#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
760pub enum ProposalDecision {
761 Approve,
763 Reject,
765 Comment,
767 RequestChanges,
769}
770
771impl ProposalDecision {
772 pub fn as_str(self) -> &'static str {
777 match self {
778 Self::Approve => "approve",
779 Self::Reject => "reject",
780 Self::Comment => "comment",
781 Self::RequestChanges => "request_changes",
782 }
783 }
784}
785
786#[derive(Clone, Debug, PartialEq)]
788#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
789pub struct ProposalAppliedPayload {
790 pub proposal_id: Id128,
791 pub applied_at: crate::Timestamp,
792 pub applied_by: String,
793 pub result: ApplyResult,
794}
795
796#[derive(Clone, Debug, PartialEq)]
798#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
799#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
800pub enum ApplyResult {
801 Success {
802 created_records: Vec<Id128>,
803 },
804 Failed {
805 error: String,
806 applied_step_count: u32,
807 },
808}
809
810#[derive(Clone, Debug, PartialEq)]
812#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
813pub struct ProposalWithdrawnPayload {
814 pub proposal_id: Id128,
815 pub by: String,
816 pub reason: Option<String>,
817}
818
819pub struct EventBuilder {
821 verb: String,
822 substrate: SubstrateKind,
823 actor: Option<String>,
824 kind: EventKind,
825 payload: EventPayload,
826 payload_schema_version: u32,
827 profile_state_version: Option<u64>,
828 aggregate: Option<AggregateRef>,
829}
830
831impl EventBuilder {
832 pub fn new(
834 verb: impl Into<String>,
835 substrate: SubstrateKind,
836 actor: impl Into<String>,
837 ) -> Self {
838 Self {
839 verb: verb.into(),
840 substrate,
841 actor: Some(actor.into()),
842 kind: EventKind::Audit,
843 payload: EventPayload::default(),
844 payload_schema_version: 1,
845 profile_state_version: None,
846 aggregate: None,
847 }
848 }
849
850 pub fn kind(mut self, kind: EventKind) -> Self {
852 self.kind = kind;
853 self
854 }
855
856 pub fn payload(mut self, payload: EventPayload) -> Self {
858 self.payload = payload;
859 self
860 }
861
862 pub fn payload_schema_version(mut self, version: u32) -> Self {
864 self.payload_schema_version = version;
865 self
866 }
867
868 pub fn profile_state_version(mut self, version: u64) -> Self {
870 self.profile_state_version = Some(version);
871 self
872 }
873
874 pub fn aggregate(mut self, aggregate: AggregateRef) -> Self {
876 self.aggregate = Some(aggregate);
877 self
878 }
879
880 pub fn build(self, header: Header) -> Event {
882 Event {
883 header,
884 verb: self.verb,
885 substrate: self.substrate,
886 actor: self.actor,
887 kind: self.kind,
888 payload: self.payload,
889 payload_schema_version: self.payload_schema_version,
890 profile_state_version: self.profile_state_version,
891 aggregate: self.aggregate,
892 }
893 }
894}
895
896#[cfg(test)]
897mod tests {
898 extern crate alloc;
899
900 use super::*;
901 use crate::{Namespace, Timestamp};
902 #[cfg(feature = "serde")]
903 use alloc::string::ToString;
904
905 fn header() -> Header {
906 Header::new(
907 Id128::from_u128(1),
908 Namespace::local(),
909 Timestamp::from_secs(1700000000),
910 )
911 }
912
913 #[test]
914 fn event_kind_parse_roundtrip() {
915 for kind in EventKind::ALL {
916 let parsed: EventKind = kind
917 .name()
918 .parse()
919 .expect("EventKind::name must parse back");
920 assert_eq!(parsed, kind);
921 }
922 }
923
924 #[test]
925 fn rerank_payload_records_served_profile() {
926 let payload = EventPayload::RerankExecuted(RerankExecutedPayload {
927 served_by_profile_id: Some("profile-a".into()),
928 model_id: Id128::from_u128(1),
929 candidates: Vec::new(),
930 reranked: Vec::new(),
931 final_scores: Vec::new(),
932 latency_us: 100,
933 hook_applied: false,
934 hook_target_match: false,
935 });
936 let event = EventBuilder::new("rerank", SubstrateKind::Note, "agent:test")
937 .kind(EventKind::RerankExecuted)
938 .payload(payload)
939 .build(header());
940
941 if let EventPayload::RerankExecuted(ref p) = event.payload {
942 assert_eq!(p.served_by_profile_id.as_deref(), Some("profile-a"));
943 } else {
944 panic!("unexpected payload variant");
945 }
946 }
947
948 #[test]
949 fn proposal_payloads_are_typed() {
950 let payload = EventPayload::ProposalReviewed(ProposalReviewedPayload {
951 proposal_id: Id128::from_u128(42),
952 reviewer: "operator".into(),
953 decision: ProposalDecision::Approve,
954 comment: None,
955 });
956 let event = EventBuilder::new("review", SubstrateKind::Entity, "operator")
957 .kind(EventKind::ProposalReviewed)
958 .payload(payload)
959 .build(header());
960 assert_eq!(event.kind.name(), "proposal_reviewed");
961 }
962
963 #[cfg(feature = "serde")]
968 #[test]
969 fn proposal_changeset_id_variants_deserialize_from_value() {
970 let uuid = "7426afd6-0234-4701-9045-83dfd39166e6";
971 let uuid2 = "abcdef01-2345-6789-abcd-ef0123456789";
972
973 let v =
975 serde_json::json!({"kind": "update_entity", "id": uuid, "patch": {"name": "NewName"}});
976 let cs: ProposalChangeset =
977 serde_json::from_value(v).expect("UpdateEntity must deserialize from Value");
978 assert!(
979 matches!(cs, ProposalChangeset::UpdateEntity { .. }),
980 "expected UpdateEntity"
981 );
982
983 for (json, expected) in [
986 (serde_json::json!({}), None),
987 (serde_json::json!({"entity_type": null}), Some(None)),
988 (
989 serde_json::json!({"entity_type": "algorithm"}),
990 Some(Some("algorithm".to_string())),
991 ),
992 ] {
993 let v = serde_json::json!({"kind": "update_entity", "id": uuid, "patch": json});
994 let cs: ProposalChangeset =
995 serde_json::from_value(v).expect("UpdateEntity must deserialize");
996 let ProposalChangeset::UpdateEntity { patch, .. } = cs else {
997 panic!("expected UpdateEntity");
998 };
999 assert_eq!(patch.entity_type, expected, "patch: {json}");
1000 }
1001
1002 let patch = ProposalEntityPatch {
1004 name: None,
1005 description: None,
1006 properties: None,
1007 tags: None,
1008 entity_type: Some(None),
1009 };
1010 let v = serde_json::to_value(&patch).expect("serialize");
1011 assert_eq!(v.get("entity_type"), Some(&serde_json::Value::Null));
1012 assert!(
1013 v.get("name").is_none(),
1014 "absent fields must not be serialized"
1015 );
1016
1017 let v = serde_json::json!({
1019 "kind": "add_edge",
1020 "source": uuid, "target": uuid2,
1021 "relation": "extends", "weight": 1.0
1022 });
1023 let cs: ProposalChangeset =
1024 serde_json::from_value(v).expect("AddEdge must deserialize from Value");
1025 assert!(
1026 matches!(cs, ProposalChangeset::AddEdge { .. }),
1027 "expected AddEdge"
1028 );
1029
1030 let v = serde_json::json!({"kind": "merge_entities", "into": uuid, "from": uuid2});
1032 let cs: ProposalChangeset =
1033 serde_json::from_value(v).expect("MergeEntities must deserialize from Value");
1034 assert!(
1035 matches!(cs, ProposalChangeset::MergeEntities { .. }),
1036 "expected MergeEntities"
1037 );
1038
1039 let v = serde_json::json!({"kind": "supersede_entity", "old": uuid, "new": uuid2});
1041 let cs: ProposalChangeset =
1042 serde_json::from_value(v).expect("SupersedeEntity must deserialize from Value");
1043 assert!(
1044 matches!(cs, ProposalChangeset::SupersedeEntity { .. }),
1045 "expected SupersedeEntity"
1046 );
1047 }
1048
1049 #[cfg(feature = "serde")]
1050 #[test]
1051 fn proposal_changeset_rejects_invalid_edge_weight() {
1052 let uuid = "7426afd6-0234-4701-9045-83dfd39166e6";
1053 let uuid2 = "abcdef01-2345-6789-abcd-ef0123456789";
1054
1055 let v = serde_json::json!({
1056 "kind": "add_edge",
1057 "source": uuid, "target": uuid2,
1058 "relation": "extends", "weight": 2.0
1059 });
1060 let result: Result<ProposalChangeset, _> = serde_json::from_value(v);
1061 assert!(result.is_err());
1062 let err = result.unwrap_err().to_string();
1063 assert!(
1064 err.contains("[0.0, 1.0]"),
1065 "error should mention range: {err}"
1066 );
1067 }
1068
1069 #[cfg(feature = "serde")]
1070 #[test]
1071 fn proposal_changeset_accepts_null_edge_weight() {
1072 let uuid = "7426afd6-0234-4701-9045-83dfd39166e6";
1073 let uuid2 = "abcdef01-2345-6789-abcd-ef0123456789";
1074
1075 let v = serde_json::json!({
1076 "kind": "add_edge",
1077 "source": uuid, "target": uuid2,
1078 "relation": "extends", "weight": null
1079 });
1080 let cs: ProposalChangeset =
1081 serde_json::from_value(v).expect("null weight should be accepted");
1082 assert!(matches!(
1083 cs,
1084 ProposalChangeset::AddEdge { weight: None, .. }
1085 ));
1086 }
1087
1088 #[cfg(feature = "serde")]
1089 #[test]
1090 fn rerank_payload_serde_rejects_non_finite_score() {
1091 let json = serde_json::json!({
1092 "served_by_profile_id": null,
1093 "model_id": "00000000-0000-0000-0000-000000000001",
1094 "candidates": [],
1095 "reranked": [],
1096 "final_scores": [["00000000-0000-0000-0000-000000000001", "Infinity"]],
1097 "latency_us": 100,
1098 "hook_applied": false,
1099 "hook_target_match": false
1100 });
1101 let result: Result<RerankExecutedPayload, _> = serde_json::from_value(json);
1102 assert!(result.is_err());
1103 }
1104
1105 #[test]
1106 fn rerank_payload_is_valid_checks_finite() {
1107 let p = RerankExecutedPayload {
1108 served_by_profile_id: None,
1109 model_id: Id128::from_u128(1),
1110 candidates: Vec::new(),
1111 reranked: Vec::new(),
1112 final_scores: alloc::vec![(Id128::from_u128(1), 0.5)],
1113 latency_us: 100,
1114 hook_applied: false,
1115 hook_target_match: false,
1116 };
1117 assert!(p.is_valid());
1118
1119 let p_inf = RerankExecutedPayload {
1120 served_by_profile_id: None,
1121 model_id: Id128::from_u128(1),
1122 candidates: Vec::new(),
1123 reranked: Vec::new(),
1124 final_scores: alloc::vec![(Id128::from_u128(1), f32::INFINITY)],
1125 latency_us: 100,
1126 hook_applied: false,
1127 hook_target_match: false,
1128 };
1129 assert!(!p_inf.is_valid());
1130 }
1131}