Skip to main content

minco_plugin_audit/
lib.rs

1//! Append-only audit events and a deterministic memory reference sink.
2#![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}