diff --git a/Cargo.toml b/Cargo.toml index 7cf3cba..4db6085 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,6 +6,10 @@ edition = "2021" [lib] crate-type = ["cdylib", "rlib"] +[[bin]] +name = "proofstell_service" +path = "src/main.rs" + [dependencies] soroban-sdk = "23.4.1" stellar-strkey = "0.0.9" diff --git a/src/cache.rs b/src/cache.rs index b0961ae..a4df679 100644 --- a/src/cache.rs +++ b/src/cache.rs @@ -16,6 +16,7 @@ use crate::metrics::MetricsRegistry; pub enum CacheKey { Verification(String), Config(String), + Events(String), } impl CacheKey { @@ -23,6 +24,7 @@ impl CacheKey { match self { CacheKey::Verification(hash) => format!("verification:{}", hash), CacheKey::Config(key) => format!("config:{}", key), + CacheKey::Events(hash) => format!("events:{}", hash), } } } @@ -407,4 +409,30 @@ mod tests { let output = metrics.render(); assert!(output.contains("cache_serialization_failures_total")); } + + #[tokio::test] + async fn event_cache_stores_and_retrieves_events() { + let cache = CacheBackend::InMemory(InMemoryCache::new()); + let key = CacheKey::Events("doc-hash-1".to_string()); + let events = vec!["{\"seq\":1}", "{\"seq\":2}"]; + let serialized = serde_json::to_string(&events).unwrap(); + + cache.set_raw(&key, &serialized, 60).await.unwrap(); + let retrieved: Option> = cache.get(&key).await.unwrap(); + + assert!(retrieved.is_some()); + assert_eq!(retrieved.unwrap().len(), 2); + } + + #[tokio::test] + async fn event_cache_events_namespace_does_not_collide() { + let cache = CacheBackend::InMemory(InMemoryCache::new()); + let v_key = CacheKey::Verification("x".to_string()); + let e_key = CacheKey::Events("x".to_string()); + + cache.set_raw(&v_key, "verification_val", 60).await.unwrap(); + cache.set_raw(&e_key, "events_val", 60).await.unwrap(); + assert_eq!(cache.get_raw(&v_key).await.unwrap(), Some("verification_val".to_string())); + assert_eq!(cache.get_raw(&e_key).await.unwrap(), Some("events_val".to_string())); + } } diff --git a/src/event.rs b/src/event.rs index 0b75c09..47273a2 100644 --- a/src/event.rs +++ b/src/event.rs @@ -12,6 +12,13 @@ use chrono::{DateTime, Utc}; use serde::{Deserialize, Serialize}; use uuid::Uuid; +/// Canonical event type identifiers for off-chain consumers. +pub const EVENT_DOCUMENT_REGISTERED: &str = "DocumentRegistered"; +pub const EVENT_DOCUMENT_REVOKED: &str = "DocumentRevoked"; +pub const EVENT_DOCUMENT_VERIFIED: &str = "DocumentVerified"; +pub const EVENT_DOCUMENT_AUTHORIZATION_FAILED: &str = "DocumentAuthorizationFailed"; +pub const EVENT_DOCUMENT_OWNER_CHANGED: &str = "DocumentOwnerChanged"; + /// Source of an audit event. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub enum EventSource { @@ -288,6 +295,66 @@ impl Default for EventIngestor { } } +/// Persistent event store that maintains append-only logs per aggregate. +/// +/// This is a minimal in-memory implementation suitable for testing and local development. +/// Production deployments should replace this with a durable store (e.g. Redis Streams, +/// PostgreSQL, or event-sourcing infrastructure). +/// +/// **Note on memory growth:** The `events` map grows unboundedly without eviction. +/// A production implementation should cap retention per aggregate. +pub struct EventStore { + events: HashMap>, +} + +impl EventStore { + pub fn new() -> Self { + Self { + events: HashMap::new(), + } + } + + /// Append an event to the store for the given aggregate. + /// + /// Returns the event with its final sequence number populated. + pub fn append(&mut self, aggregate_id: impl Into, event: Event) -> Event { + let aggregate_id = aggregate_id.into(); + let events = self.events.entry(aggregate_id.clone()).or_default(); + + let sequence = events.len() as u64 + 1; + let mut finalized = event.with_sequence(sequence); + finalized.aggregate_id = aggregate_id; + + events.push(finalized.clone()); + + finalized + } + + /// Retrieve the full event history for an aggregate. + pub fn get_history(&self, aggregate_id: &str) -> Option<&Vec> { + self.events.get(aggregate_id) + } + + /// Retrieve the latest sequence number for an aggregate. + pub fn get_latest_sequence(&self, aggregate_id: &str) -> Option { + self.events + .get(aggregate_id) + .and_then(|events| events.last()) + .map(|e| e.sequence) + } + + /// Return the number of events stored for an aggregate. + pub fn count(&self, aggregate_id: &str) -> usize { + self.events.get(aggregate_id).map(|v| v.len()).unwrap_or(0) + } +} + +impl Default for EventStore { + fn default() -> Self { + Self::new() + } +} + #[cfg(test)] mod tests { use super::*; @@ -464,3 +531,121 @@ mod tests { assert!(output.contains("event_backlog_size")); } } + +#[cfg(test)] +mod store_tests { + use super::*; + + #[test] + fn event_store_appends_events_in_sequence() { + let mut store = EventStore::new(); + + let e1 = store.append( + "doc-1", + Event::new( + "doc-1".to_string(), + EVENT_DOCUMENT_REGISTERED.to_string(), + serde_json::json!({"issuer": "addr1"}), + "issuer".to_string(), + ), + ); + + let e2 = store.append( + "doc-1", + Event::new( + "doc-1".to_string(), + EVENT_DOCUMENT_REVOKED.to_string(), + serde_json::json!({}), + "issuer".to_string(), + ), + ); + + assert_eq!(e1.sequence, 1); + assert_eq!(e2.sequence, 2); + assert_eq!(store.count("doc-1"), 2); + } + + #[test] + fn event_store_retrieves_history() { + let mut store = EventStore::new(); + + store.append( + "doc-1", + Event::new( + "doc-1".to_string(), + EVENT_DOCUMENT_REGISTERED.to_string(), + serde_json::json!({}), + "issuer".to_string(), + ), + ); + + let history = store.get_history("doc-1").unwrap(); + assert_eq!(history.len(), 1); + assert_eq!(history[0].event_type, EVENT_DOCUMENT_REGISTERED); + } + + #[test] + fn event_store_empty_history_returns_none() { + let store = EventStore::new(); + assert!(store.get_history("missing").is_none()); + assert_eq!(store.get_latest_sequence("missing"), None); + } + + #[test] + fn event_store_get_latest_sequence() { + let mut store = EventStore::new(); + + store.append( + "doc-1", + Event::new( + "doc-1".to_string(), + EVENT_DOCUMENT_REGISTERED.to_string(), + serde_json::json!({}), + "issuer".to_string(), + ), + ); + store.append( + "doc-1", + Event::new( + "doc-1".to_string(), + EVENT_DOCUMENT_VERIFIED.to_string(), + serde_json::json!({}), + "caller".to_string(), + ), + ); + + assert_eq!(store.get_latest_sequence("doc-1"), Some(2)); + } + + #[test] + fn event_store_separates_aggregates() { + let mut store = EventStore::new(); + + store.append( + "doc-1", + Event::new( + "doc-1".to_string(), + EVENT_DOCUMENT_REGISTERED.to_string(), + serde_json::json!({}), + "issuer".to_string(), + ), + ); + store.append( + "doc-2", + Event::new( + "doc-2".to_string(), + EVENT_DOCUMENT_OWNER_CHANGED.to_string(), + serde_json::json!({}), + "owner".to_string(), + ), + ); + + assert_eq!(store.count("doc-1"), 1); + assert_eq!(store.count("doc-2"), 1); + assert_eq!(store.get_history("doc-1").unwrap()[0].event_type, EVENT_DOCUMENT_REGISTERED); + assert_eq!( + store.get_history("doc-2").unwrap()[0].event_type, + EVENT_DOCUMENT_OWNER_CHANGED + ); + } +}