1#![forbid(unsafe_code)]
3
4mod journal;
5mod v2;
6
7pub use journal::*;
8pub use v2::*;
9
10use async_trait::async_trait;
11use chrono::{DateTime, Utc};
12use minco_core::{
13 CapabilityProvision, DataClass, Plugin, PluginContext, PluginDescriptor, PluginError, PluginId,
14 PluginStability,
15};
16use semver::{Version, VersionReq};
17use serde::{Deserialize, Serialize};
18use std::{collections::BTreeMap, sync::Arc};
19use tokio::sync::RwLock;
20use uuid::Uuid;
21
22#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
23pub struct AuditEvent {
24 pub id: Uuid,
25 pub action: String,
26 pub resource_type: String,
27 pub resource_id: String,
28 #[serde(default, skip_serializing_if = "Option::is_none")]
29 pub actor_subject: Option<String>,
30 pub correlation_id: Uuid,
31 pub occurred_at: DateTime<Utc>,
32 #[serde(default)]
33 pub metadata: BTreeMap<String, serde_json::Value>,
34}
35
36impl AuditEvent {
37 pub fn new(
38 action: impl Into<String>,
39 resource_type: impl Into<String>,
40 resource_id: impl Into<String>,
41 correlation_id: Uuid,
42 ) -> Self {
43 Self {
44 id: Uuid::now_v7(),
45 action: action.into(),
46 resource_type: resource_type.into(),
47 resource_id: resource_id.into(),
48 actor_subject: None,
49 correlation_id,
50 occurred_at: Utc::now(),
51 metadata: BTreeMap::new(),
52 }
53 }
54}
55
56#[async_trait]
57pub trait AuditSink: Send + Sync + std::fmt::Debug {
58 async fn append(&self, event: AuditEvent) -> Result<(), AuditError>;
59}
60
61#[derive(Clone)]
62pub struct AuditService(pub Arc<dyn AuditSink>);
63
64impl std::fmt::Debug for AuditService {
65 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
66 formatter.debug_tuple("AuditService").finish()
67 }
68}
69
70impl AuditService {
71 pub fn new(sink: Arc<dyn AuditSink>) -> Self {
72 Self(sink)
73 }
74
75 pub async fn append(&self, event: AuditEvent) -> Result<(), AuditError> {
76 self.0.append(event).await
77 }
78}
79
80#[derive(Debug, Default)]
81pub struct MemoryAuditSink {
82 events: RwLock<Vec<AuditEvent>>,
83}
84
85impl MemoryAuditSink {
86 pub async fn all(&self) -> Vec<AuditEvent> {
87 self.events.read().await.clone()
88 }
89}
90
91#[async_trait]
92impl AuditSink for MemoryAuditSink {
93 async fn append(&self, event: AuditEvent) -> Result<(), AuditError> {
94 if event.action.trim().is_empty() || event.resource_id.trim().is_empty() {
95 return Err(AuditError::InvalidEvent);
96 }
97 self.events.write().await.push(event);
98 Ok(())
99 }
100}
101
102#[derive(Debug, Clone)]
103pub struct AuditPlugin {
104 service: AuditService,
105 ledger: Option<AuditLedgerServices>,
106}
107
108impl AuditPlugin {
109 pub fn new(sink: Arc<dyn AuditSink>) -> Self {
110 Self {
111 service: AuditService::new(sink),
112 ledger: None,
113 }
114 }
115
116 #[must_use]
117 pub fn with_ledger<L>(mut self, ledger: Arc<L>) -> Self
118 where
119 L: AuditLedgerWriter + AuditReader + AuditStorageInspector + 'static,
120 {
121 self.ledger = Some(AuditLedgerServices::new(ledger));
122 self
123 }
124
125 #[must_use]
126 pub fn with_ledger_services(mut self, services: AuditLedgerServices) -> Self {
127 self.ledger = Some(services);
128 self
129 }
130
131 pub fn memory() -> (Self, Arc<MemoryAuditSink>) {
132 let sink = Arc::new(MemoryAuditSink::default());
133 let ledger = Arc::new(MemoryAuditLedger::default());
134 (Self::new(sink.clone()).with_ledger(ledger), sink)
135 }
136
137 pub fn memory_v2() -> (Self, Arc<MemoryAuditSink>, Arc<MemoryAuditLedger>) {
138 let sink = Arc::new(MemoryAuditSink::default());
139 let ledger = Arc::new(MemoryAuditLedger::default());
140 (
141 Self::new(sink.clone()).with_ledger(ledger.clone()),
142 sink,
143 ledger,
144 )
145 }
146}
147
148impl Plugin for AuditPlugin {
149 fn descriptor(&self) -> PluginDescriptor {
150 let mut descriptor = PluginDescriptor::new(
151 PluginId::new("audit").expect("static plugin ID"),
152 Version::new(1, 0, 0),
153 "Durable append-only audit history independent of operational logs",
154 );
155 descriptor.documentation = Some("https://docs.rs/minco-plugin-audit".into());
156 descriptor.core_compatibility =
157 VersionReq::parse(concat!("^", env!("CARGO_PKG_VERSION"))).expect("package version");
158 descriptor.stability = PluginStability::Beta;
159 descriptor.data_classes.extend([
160 DataClass::Internal,
161 DataClass::Personal,
162 DataClass::Confidential,
163 ]);
164 descriptor.provides.push(CapabilityProvision {
165 name: "audit.append".into(),
166 version: Version::new(1, 0, 0),
167 });
168 if self.ledger.is_some() {
169 descriptor.provides.extend([
170 CapabilityProvision {
171 name: "audit.ledger".into(),
172 version: Version::new(2, 0, 0),
173 },
174 CapabilityProvision {
175 name: "audit.query".into(),
176 version: Version::new(2, 0, 0),
177 },
178 CapabilityProvision {
179 name: "audit.health".into(),
180 version: Version::new(1, 0, 0),
181 },
182 ]);
183 }
184 descriptor
185 }
186
187 fn install(&self, context: &mut PluginContext<'_>) -> Result<(), PluginError> {
188 context.services().insert(Arc::new(self.service.clone()))?;
189 if let Some(ledger) = &self.ledger {
190 context.services().insert(Arc::new(ledger.clone()))?;
191 }
192 Ok(())
193 }
194}
195
196#[derive(Debug, thiserror::Error)]
197pub enum AuditError {
198 #[error("audit events require a non-empty action and resource ID")]
199 InvalidEvent,
200 #[error("audit append failed: {0}")]
201 Append(String),
202}
203
204#[cfg(test)]
205mod tests {
206 use super::*;
207
208 #[tokio::test]
209 async fn memory_sink_is_append_only_and_ordered() {
210 let sink = MemoryAuditSink::default();
211 let first = AuditEvent::new("feedback.created", "feedback", "one", Uuid::now_v7());
212 let second = AuditEvent::new("feedback.replied", "feedback", "one", Uuid::now_v7());
213 sink.append(first).await.unwrap();
214 sink.append(second).await.unwrap();
215 assert_eq!(
216 sink.all()
217 .await
218 .iter()
219 .map(|event| event.action.as_str())
220 .collect::<Vec<_>>(),
221 ["feedback.created", "feedback.replied"]
222 );
223 }
224
225 #[test]
226 fn memory_plugin_advertises_additive_v2_capabilities() {
227 let descriptor = AuditPlugin::memory().0.descriptor();
228 assert_eq!(descriptor.version, Version::new(1, 0, 0));
229 assert!(
230 descriptor
231 .provides
232 .iter()
233 .any(|capability| capability.name == "audit.ledger")
234 );
235 }
236}