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