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 Refusal,
152}
153
154impl EventKind {
155 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 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#[derive(Clone, Debug, PartialEq, Eq)]
354#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
355pub struct AggregateRef {
356 pub kind: String,
358 pub id: Id128,
360}
361
362#[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 Json(String),
376 RerankExecuted(RerankExecutedPayload),
378 #[cfg(feature = "serde")]
380 ProposalCreated(ProposalCreatedPayload),
381 ProposalReviewed(ProposalReviewedPayload),
383 ProposalApplied(ProposalAppliedPayload),
385 ProposalWithdrawn(ProposalWithdrawnPayload),
387}
388
389impl Default for EventPayload {
390 fn default() -> Self {
391 Self::Json("{}".into())
392 }
393}
394
395#[derive(Clone, Debug, PartialEq)]
400#[cfg_attr(feature = "serde", derive(serde::Serialize))]
401pub struct RerankExecutedPayload {
402 pub served_by_profile_id: Option<String>,
404 pub model_id: Id128,
406 pub candidates: Vec<Id128>,
408 pub reranked: Vec<(Id128, Vec<(String, f32)>)>,
410 pub final_scores: Vec<(Id128, f32)>,
412 pub latency_us: u64,
414 pub hook_applied: bool,
416 pub hook_target_match: bool,
418}
419
420impl RerankExecutedPayload {
421 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#[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#[cfg(feature = "serde")]
501#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
502pub struct EntityDraft {
503 pub kind: String,
505 pub name: String,
507 #[serde(skip_serializing_if = "Option::is_none")]
509 pub description: Option<String>,
510 #[serde(skip_serializing_if = "Option::is_none")]
512 pub properties: Option<serde_json::Value>,
513 #[serde(default, skip_serializing_if = "Vec::is_empty")]
515 pub tags: Vec<String>,
516}
517
518#[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 #[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 #[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#[cfg(feature = "serde")]
554#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
555pub struct NoteDraft {
556 pub kind: String,
558 pub content: String,
560 #[serde(skip_serializing_if = "Option::is_none")]
562 pub name: Option<String>,
563 #[serde(skip_serializing_if = "Option::is_none")]
565 pub properties: Option<serde_json::Value>,
566}
567
568#[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#[cfg(feature = "serde")]
596#[derive(Clone, Debug, PartialEq, serde::Serialize)]
597#[serde(tag = "kind", rename_all = "snake_case")]
598pub enum ProposalChangeset {
599 AddEntity {
601 entity: EntityDraft,
602 },
603 UpdateEntity {
606 id: Id128,
607 patch: ProposalEntityPatch,
608 },
609 AddEdge {
611 source: Id128,
612 target: Id128,
613 relation: crate::EdgeRelation,
614 weight: Option<f32>,
615 },
616 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#[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#[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 Approve,
769 Reject,
771 Comment,
773 RequestChanges,
775}
776
777impl ProposalDecision {
778 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#[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#[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#[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
825pub 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 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 pub fn kind(mut self, kind: EventKind) -> Self {
858 self.kind = kind;
859 self
860 }
861
862 pub fn payload(mut self, payload: EventPayload) -> Self {
864 self.payload = payload;
865 self
866 }
867
868 pub fn payload_schema_version(mut self, version: u32) -> Self {
870 self.payload_schema_version = version;
871 self
872 }
873
874 pub fn profile_state_version(mut self, version: u64) -> Self {
876 self.profile_state_version = Some(version);
877 self
878 }
879
880 pub fn aggregate(mut self, aggregate: AggregateRef) -> Self {
882 self.aggregate = Some(aggregate);
883 self
884 }
885
886 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 #[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 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 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 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 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 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 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}