1#[cfg(feature = "replication")]
2use crate::infrastructure::replication::ReplicationMode;
3use crate::{
4 application::{
5 dto::{
6 AckRequest, ConsumerEventDto, ConsumerEventsResponse, ConsumerResponse,
7 DetectDuplicatesRequest, DetectDuplicatesResponse, DuplicateGroup, EntitySummary,
8 EventDto, IngestEventRequest, IngestEventResponse, IngestEventsBatchRequest,
9 IngestEventsBatchResponse, ListEntitiesRequest, ListEntitiesResponse,
10 QueryEventsRequest, QueryEventsResponse, RegisterConsumerRequest,
11 },
12 services::{
13 analytics::{
14 AnalyticsEngine, CorrelationRequest, CorrelationResponse, EventFrequencyRequest,
15 EventFrequencyResponse, StatsSummaryRequest, StatsSummaryResponse,
16 },
17 pipeline::{PipelineConfig, PipelineStats},
18 replay::{ReplayProgress, StartReplayRequest, StartReplayResponse},
19 schema::{
20 CompatibilityMode, RegisterSchemaRequest, RegisterSchemaResponse,
21 ValidateEventRequest, ValidateEventResponse,
22 },
23 webhook::{RegisterWebhookRequest, UpdateWebhookRequest},
24 },
25 },
26 domain::{
27 entities::{Event, SchemaEnforcement},
28 value_objects::TenantId,
29 },
30 error::Result,
31 infrastructure::{
32 persistence::{
33 compaction::CompactionResult,
34 snapshot::{
35 CreateSnapshotRequest, CreateSnapshotResponse, ListSnapshotsRequest,
36 ListSnapshotsResponse, SnapshotInfo,
37 },
38 },
39 query::{
40 geospatial::GeoQueryRequest,
41 graphql::{GraphQLError, GraphQLRequest, GraphQLResponse},
42 },
43 security::middleware::OptionalAuth,
44 web::api_v1::AppState,
45 },
46 store::{EventStore, EventTypeInfo, StreamInfo},
47};
48use axum::{
49 Json, Router,
50 extract::{Path, Query, State, WebSocketUpgrade},
51 response::{IntoResponse, Response},
52 routing::{get, post, put},
53};
54use serde::Deserialize;
55use std::sync::Arc;
56use tower_http::{
57 cors::{Any, CorsLayer},
58 trace::TraceLayer,
59};
60
61type SharedStore = Arc<EventStore>;
62
63#[cfg(feature = "replication")]
69async fn await_replication_ack(state: &AppState) {
70 let shipper_guard = state.wal_shipper.read().await;
71 if let Some(ref shipper) = *shipper_guard {
72 let mode = shipper.replication_mode();
73 if mode == ReplicationMode::Async {
74 return;
75 }
76
77 let target_offset = shipper.current_leader_offset();
78 if target_offset == 0 {
79 return;
80 }
81
82 let shipper = Arc::clone(shipper);
83 drop(shipper_guard);
85
86 let timer = state
87 .store
88 .metrics()
89 .replication_ack_wait_seconds
90 .start_timer();
91 let acked = shipper.wait_for_ack(target_offset).await;
92 timer.observe_duration();
93
94 if !acked {
95 tracing::warn!(
96 "Replication ACK timeout in {} mode (offset {}). \
97 Write succeeded locally but follower confirmation pending.",
98 mode,
99 target_offset,
100 );
101 }
102 }
103}
104#[cfg(not(feature = "replication"))]
105async fn await_replication_ack(_state: &AppState) {
106 }
108
109pub async fn serve(store: SharedStore, addr: &str) -> anyhow::Result<()> {
110 let app = Router::new()
111 .route("/health", get(health))
112 .route("/metrics", get(prometheus_metrics)) .route("/api/v1/events", post(ingest_event))
114 .route("/api/v1/events/batch", post(ingest_events_batch))
115 .route("/api/v1/events/query", get(query_events))
116 .route("/api/v1/events/{event_id}", get(get_event_by_id))
117 .route("/api/v1/events/stream", get(events_websocket)) .route("/api/v1/streams", get(list_streams))
120 .route("/api/v1/event-types", get(list_event_types))
121 .route("/api/v1/entities/duplicates", get(detect_duplicates))
122 .route("/api/v1/entities/{entity_id}/state", get(get_entity_state))
123 .route(
124 "/api/v1/entities/{entity_id}/snapshot",
125 get(get_entity_snapshot),
126 )
127 .route("/api/v1/stats", get(get_stats))
128 .route("/api/v1/analytics/frequency", get(analytics_frequency))
130 .route("/api/v1/analytics/summary", get(analytics_summary))
131 .route("/api/v1/analytics/correlation", get(analytics_correlation))
132 .route("/api/v1/snapshots", post(create_snapshot))
134 .route("/api/v1/snapshots", get(list_snapshots))
135 .route(
136 "/api/v1/snapshots/{entity_id}/latest",
137 get(get_latest_snapshot),
138 )
139 .route("/api/v1/compaction/trigger", post(trigger_compaction))
141 .route("/api/v1/compaction/stats", get(compaction_stats))
142 .route("/api/v1/schemas", post(register_schema))
144 .route("/api/v1/schemas", get(list_subjects))
145 .route("/api/v1/schemas/{subject}", get(get_schema))
146 .route(
147 "/api/v1/schemas/{subject}/versions",
148 get(list_schema_versions),
149 )
150 .route("/api/v1/schemas/validate", post(validate_event_schema))
151 .route(
152 "/api/v1/schemas/{subject}/compatibility",
153 put(set_compatibility_mode),
154 )
155 .route("/api/v1/replay", post(start_replay))
157 .route("/api/v1/replay", get(list_replays))
158 .route("/api/v1/replay/{replay_id}", get(get_replay_progress))
159 .route("/api/v1/replay/{replay_id}/cancel", post(cancel_replay))
160 .route(
161 "/api/v1/replay/{replay_id}",
162 axum::routing::delete(delete_replay),
163 )
164 .route("/api/v1/pipelines", post(register_pipeline))
166 .route("/api/v1/pipelines", get(list_pipelines))
167 .route("/api/v1/pipelines/stats", get(all_pipeline_stats))
168 .route("/api/v1/pipelines/{pipeline_id}", get(get_pipeline))
169 .route(
170 "/api/v1/pipelines/{pipeline_id}",
171 axum::routing::delete(remove_pipeline),
172 )
173 .route(
174 "/api/v1/pipelines/{pipeline_id}/stats",
175 get(get_pipeline_stats),
176 )
177 .route("/api/v1/pipelines/{pipeline_id}/reset", put(reset_pipeline))
178 .route("/api/v1/projections", get(list_projections))
180 .route("/api/v1/projections/{name}", get(get_projection))
181 .route(
182 "/api/v1/projections/{name}",
183 axum::routing::delete(delete_projection),
184 )
185 .route(
186 "/api/v1/projections/{name}/state",
187 get(get_projection_state_summary),
188 )
189 .route("/api/v1/projections/{name}/reset", post(reset_projection))
190 .route(
191 "/api/v1/projections/{name}/{entity_id}/state",
192 get(get_projection_state),
193 )
194 .route(
195 "/api/v1/projections/{name}/{entity_id}/state",
196 post(save_projection_state),
197 )
198 .route(
199 "/api/v1/projections/{name}/{entity_id}/state",
200 put(save_projection_state),
201 )
202 .route(
203 "/api/v1/projections/{name}/bulk",
204 post(bulk_get_projection_states),
205 )
206 .route(
207 "/api/v1/projections/{name}/bulk/save",
208 post(bulk_save_projection_states),
209 )
210 .route("/api/v1/webhooks", post(register_webhook))
212 .route("/api/v1/webhooks", get(list_webhooks))
213 .route("/api/v1/webhooks/{webhook_id}", get(get_webhook))
214 .route("/api/v1/webhooks/{webhook_id}", put(update_webhook))
215 .route(
216 "/api/v1/webhooks/{webhook_id}",
217 axum::routing::delete(delete_webhook),
218 )
219 .route(
220 "/api/v1/webhooks/{webhook_id}/deliveries",
221 get(list_webhook_deliveries),
222 )
223 .route("/api/v1/graphql", post(graphql_query))
225 .route("/api/v1/geospatial/query", post(geo_query))
226 .route("/api/v1/geospatial/stats", get(geo_stats))
227 .route("/api/v1/exactly-once/stats", get(exactly_once_stats))
228 .route(
229 "/api/v1/schema-evolution/history/{event_type}",
230 get(schema_evolution_history),
231 )
232 .route(
233 "/api/v1/schema-evolution/schema/{event_type}",
234 get(schema_evolution_schema),
235 )
236 .route(
237 "/api/v1/schema-evolution/stats",
238 get(schema_evolution_stats),
239 )
240 .layer(
241 CorsLayer::new()
242 .allow_origin(Any)
243 .allow_methods(Any)
244 .allow_headers(Any),
245 )
246 .layer(TraceLayer::new_for_http())
247 .with_state(store);
248
249 let listener = tokio::net::TcpListener::bind(addr).await?;
250 axum::serve(listener, app).await?;
251
252 Ok(())
253}
254
255pub async fn health() -> impl IntoResponse {
256 Json(serde_json::json!({
257 "status": "healthy",
258 "service": "allsource-core",
259 "version": env!("CARGO_PKG_VERSION")
260 }))
261}
262
263pub async fn prometheus_metrics(State(store): State<SharedStore>) -> impl IntoResponse {
265 let metrics = store.metrics();
266
267 match metrics.encode() {
268 Ok(encoded) => Response::builder()
269 .status(200)
270 .header("Content-Type", "text/plain; version=0.0.4")
271 .body(encoded)
272 .unwrap()
273 .into_response(),
274 Err(e) => Response::builder()
275 .status(500)
276 .body(format!("Error encoding metrics: {e}"))
277 .unwrap()
278 .into_response(),
279 }
280}
281
282pub async fn ingest_event(
283 State(store): State<SharedStore>,
284 Json(req): Json<IngestEventRequest>,
285) -> Result<Json<IngestEventResponse>> {
286 let expected_version = req.expected_version;
287
288 let tenant_id = req.tenant_id.unwrap_or_else(|| "default".to_string());
289 let event = Event::from_strings(
290 req.event_type,
291 req.entity_id,
292 tenant_id,
293 req.payload,
294 req.metadata,
295 )?;
296
297 let event_id = event.id;
298 let timestamp = event.timestamp;
299
300 let new_version = store.ingest_with_expected_version(&event, expected_version)?;
301
302 tracing::info!("Event ingested: {}", event_id);
303
304 Ok(Json(IngestEventResponse {
305 event_id,
306 timestamp,
307 version: Some(new_version),
308 }))
309}
310
311async fn enforce_schema_if_configured(
332 state: &AppState,
333 tenant_id: &str,
334 event: &Event,
335) -> Result<()> {
336 let Ok(parsed) = TenantId::new(tenant_id.to_string()) else {
340 return Ok(());
341 };
342 let mode = match state.tenant_repo.find_by_id(&parsed).await {
343 Ok(Some(t)) => t.schema_enforcement(),
344 _ => SchemaEnforcement::Permissive,
347 };
348 if matches!(mode, SchemaEnforcement::Permissive) {
349 return Ok(());
350 }
351
352 let registry = state.store.schema_registry();
356 let Ok(schema) = registry.get_schema(event.event_type.as_str(), None) else {
357 return Ok(());
360 };
361
362 let result = registry
363 .validate(
364 event.event_type.as_str(),
365 Some(schema.version),
366 &event.payload,
367 )
368 .map_err(|e| crate::error::AllSourceError::InternalError(e.to_string()))?;
369
370 if result.valid {
371 return Ok(());
372 }
373
374 match mode {
375 SchemaEnforcement::Strict => Err(crate::error::AllSourceError::SchemaViolation {
376 event_type: event.event_type.as_str().to_string(),
377 schema_version: result.schema_version,
378 errors: result.errors,
379 }),
380 SchemaEnforcement::Warn => {
381 tracing::warn!(
382 tenant = %tenant_id,
383 event_type = %event.event_type.as_str(),
384 schema_version = result.schema_version,
385 errors = ?result.errors,
386 "schema violation (warn mode — write accepted)"
387 );
388 Ok(())
389 }
390 SchemaEnforcement::Permissive => Ok(()),
392 }
393}
394
395pub async fn ingest_event_v1(
396 State(state): State<AppState>,
397 Json(req): Json<IngestEventRequest>,
398) -> Result<Json<IngestEventResponse>> {
399 let expected_version = req.expected_version;
400
401 let tenant_id = req.tenant_id.unwrap_or_else(|| "default".to_string());
402
403 let event = Event::from_strings(
404 req.event_type,
405 req.entity_id,
406 tenant_id.clone(),
407 req.payload,
408 req.metadata,
409 )?;
410
411 enforce_schema_if_configured(&state, &tenant_id, &event).await?;
414
415 let event_id = event.id;
416 let timestamp = event.timestamp;
417
418 let new_version = state
419 .store
420 .ingest_with_expected_version(&event, expected_version)?;
421
422 await_replication_ack(&state).await;
424
425 tracing::info!("Event ingested: {}", event_id);
426
427 Ok(Json(IngestEventResponse {
428 event_id,
429 timestamp,
430 version: Some(new_version),
431 }))
432}
433
434pub async fn ingest_events_batch(
439 State(store): State<SharedStore>,
440 Json(req): Json<IngestEventsBatchRequest>,
441) -> Result<Json<IngestEventsBatchResponse>> {
442 let total = req.events.len();
443 let mut ingested_events = Vec::with_capacity(total);
444
445 for event_req in req.events {
446 let tenant_id = event_req.tenant_id.unwrap_or_else(|| "default".to_string());
447 let expected_version = event_req.expected_version;
448
449 let event = Event::from_strings(
450 event_req.event_type,
451 event_req.entity_id,
452 tenant_id,
453 event_req.payload,
454 event_req.metadata,
455 )?;
456
457 let event_id = event.id;
463 let timestamp = event.timestamp;
464
465 let new_version = store.ingest_with_expected_version(&event, expected_version)?;
466
467 ingested_events.push(IngestEventResponse {
468 event_id,
469 timestamp,
470 version: Some(new_version),
471 });
472 }
473
474 let ingested = ingested_events.len();
475 tracing::info!("Batch ingested {} events", ingested);
476
477 Ok(Json(IngestEventsBatchResponse {
478 total,
479 ingested,
480 events: ingested_events,
481 }))
482}
483
484pub async fn ingest_events_batch_v1(
490 State(state): State<AppState>,
491 Json(req): Json<IngestEventsBatchRequest>,
492) -> Result<Json<IngestEventsBatchResponse>> {
493 let total = req.events.len();
494 let mut ingested_events = Vec::with_capacity(total);
495
496 for event_req in req.events {
497 let tenant_id = event_req.tenant_id.unwrap_or_else(|| "default".to_string());
498 let expected_version = event_req.expected_version;
499
500 let event = Event::from_strings(
501 event_req.event_type,
502 event_req.entity_id,
503 tenant_id.clone(),
504 event_req.payload,
505 event_req.metadata,
506 )?;
507
508 enforce_schema_if_configured(&state, &tenant_id, &event).await?;
509
510 let event_id = event.id;
511 let timestamp = event.timestamp;
512
513 let new_version = state
514 .store
515 .ingest_with_expected_version(&event, expected_version)?;
516
517 ingested_events.push(IngestEventResponse {
518 event_id,
519 timestamp,
520 version: Some(new_version),
521 });
522 }
523
524 await_replication_ack(&state).await;
526
527 let ingested = ingested_events.len();
528 tracing::info!("Batch ingested {} events", ingested);
529
530 Ok(Json(IngestEventsBatchResponse {
531 total,
532 ingested,
533 events: ingested_events,
534 }))
535}
536
537#[derive(Debug, Deserialize)]
541pub struct EventOrderParam {
542 pub order: Option<String>,
544}
545
546#[derive(Debug, Deserialize)]
551pub struct EventOffsetParam {
552 pub offset: Option<usize>,
554}
555
556pub async fn query_events(
557 OptionalAuth(auth): OptionalAuth,
558 Query(req): Query<QueryEventsRequest>,
559 Query(order_param): Query<EventOrderParam>,
560 Query(offset_param): Query<EventOffsetParam>,
561 State(store): State<SharedStore>,
562) -> Result<Json<QueryEventsResponse>> {
563 let offset = offset_param.offset.unwrap_or(0);
564 let queried_entity_id = req.entity_id.clone();
565
566 let descending = match order_param.order.as_deref() {
571 None => false,
572 Some(o) if o.eq_ignore_ascii_case("asc") => false,
573 Some(o) if o.eq_ignore_ascii_case("desc") => true,
574 Some(other) => {
575 return Err(crate::error::AllSourceError::InvalidInput(format!(
576 "invalid 'order' value '{other}': expected 'asc' or 'desc'"
577 )));
578 }
579 };
580
581 let enforced_tenant = req
589 .tenant_id
590 .clone()
591 .or_else(|| auth.as_ref().map(|a| a.tenant_id().to_string()));
592
593 if enforced_tenant.as_deref().unwrap_or("").is_empty() {
600 return Ok(Json(QueryEventsResponse {
601 events: Vec::new(),
602 count: 0,
603 total_count: 0,
604 has_more: false,
605 entity_version: None,
606 }));
607 }
608
609 let scoped_req = QueryEventsRequest {
617 tenant_id: enforced_tenant,
618 ..req
619 };
620 let (limited_events, total_count) = store.query_window(&scoped_req, offset, descending)?;
621
622 let count = limited_events.len();
623 let has_more = offset + count < total_count;
627 let events: Vec<EventDto> = limited_events.iter().map(EventDto::from).collect();
628
629 let entity_version = queried_entity_id
631 .as_deref()
632 .map(|eid| store.get_entity_version(eid));
633
634 tracing::debug!("Query returned {} events (total: {})", count, total_count);
635
636 Ok(Json(QueryEventsResponse {
637 events,
638 count,
639 total_count,
640 has_more,
641 entity_version,
642 }))
643}
644
645pub async fn list_entities(
646 State(store): State<SharedStore>,
647 Query(req): Query<ListEntitiesRequest>,
648) -> Result<Json<ListEntitiesResponse>> {
649 use std::collections::HashMap;
650
651 let query_req = QueryEventsRequest {
653 entity_id: None,
654 event_type: None,
655 tenant_id: None,
656 as_of: None,
657 since: None,
658 until: None,
659 limit: None,
660 event_type_prefix: req.event_type_prefix,
661 exclude_event_type_prefix: None,
662 payload_filter: req.payload_filter,
663 };
664 let events = store.query(&query_req)?;
665
666 let mut entity_map: HashMap<String, Vec<&Event>> = HashMap::new();
668 for event in &events {
669 entity_map
670 .entry(event.entity_id().to_string())
671 .or_default()
672 .push(event);
673 }
674
675 let ascending = match req.order.as_deref() {
679 None => false,
680 Some(o) if o.eq_ignore_ascii_case("desc") => false,
681 Some(o) if o.eq_ignore_ascii_case("asc") => true,
682 Some(other) => {
683 return Err(crate::error::AllSourceError::InvalidInput(format!(
684 "invalid 'order' value '{other}': expected 'asc' or 'desc'"
685 )));
686 }
687 };
688
689 let mut summaries: Vec<EntitySummary> = entity_map
693 .into_iter()
694 .map(|(entity_id, events)| {
695 let last = events.iter().max_by_key(|e| e.timestamp()).unwrap();
696 EntitySummary {
697 entity_id,
698 event_count: events.len(),
699 last_event_type: last.event_type_str().to_string(),
700 last_event_at: last.timestamp(),
701 }
702 })
703 .collect();
704 summaries.sort_by(|a, b| {
705 let by_time = a.last_event_at.cmp(&b.last_event_at);
706 let by_time = if ascending {
707 by_time
708 } else {
709 by_time.reverse()
710 };
711 by_time.then_with(|| a.entity_id.cmp(&b.entity_id))
712 });
713
714 let total = summaries.len();
715
716 let offset = req.offset.unwrap_or(0);
718 let summaries: Vec<EntitySummary> = summaries.into_iter().skip(offset).collect::<Vec<_>>();
719 let summaries = if let Some(limit) = req.limit {
720 let has_more = summaries.len() > limit;
721 let truncated: Vec<EntitySummary> = summaries.into_iter().take(limit).collect();
722 return Ok(Json(ListEntitiesResponse {
723 entities: truncated,
724 total,
725 has_more,
726 }));
727 } else {
728 summaries
729 };
730
731 Ok(Json(ListEntitiesResponse {
732 entities: summaries,
733 total,
734 has_more: false,
735 }))
736}
737
738pub async fn detect_duplicates(
739 State(store): State<SharedStore>,
740 Query(req): Query<DetectDuplicatesRequest>,
741) -> Result<Json<DetectDuplicatesResponse>> {
742 use std::collections::HashMap;
743
744 let group_by_fields: Vec<&str> = req.group_by.split(',').map(str::trim).collect();
745
746 let query_req = QueryEventsRequest {
748 entity_id: None,
749 event_type: None,
750 tenant_id: None,
751 as_of: None,
752 since: None,
753 until: None,
754 limit: None,
755 event_type_prefix: Some(req.event_type_prefix),
756 exclude_event_type_prefix: None,
757 payload_filter: None,
758 };
759 let events = store.query(&query_req)?;
760
761 let mut entity_latest: HashMap<String, &Event> = HashMap::new();
764 for event in &events {
765 let eid = event.entity_id().to_string();
766 entity_latest
767 .entry(eid)
768 .and_modify(|existing| {
769 if event.timestamp() > existing.timestamp() {
770 *existing = event;
771 }
772 })
773 .or_insert(event);
774 }
775
776 let mut groups: HashMap<String, Vec<String>> = HashMap::new();
778 for (entity_id, event) in &entity_latest {
779 let payload = event.payload();
780 let mut key_parts = serde_json::Map::new();
781 for field in &group_by_fields {
782 let value = payload
783 .get(*field)
784 .cloned()
785 .unwrap_or(serde_json::Value::Null);
786 key_parts.insert((*field).to_string(), value);
787 }
788 let key_str = serde_json::to_string(&key_parts).unwrap_or_default();
789 groups.entry(key_str).or_default().push(entity_id.clone());
790 }
791
792 let mut duplicate_groups: Vec<DuplicateGroup> = groups
794 .into_iter()
795 .filter(|(_, ids)| ids.len() > 1)
796 .map(|(key_str, mut ids)| {
797 ids.sort();
798 let key: serde_json::Value =
799 serde_json::from_str(&key_str).unwrap_or(serde_json::Value::Null);
800 let count = ids.len();
801 DuplicateGroup {
802 key,
803 entity_ids: ids,
804 count,
805 }
806 })
807 .collect();
808
809 duplicate_groups.sort_by_key(|a| std::cmp::Reverse(a.count));
811
812 let total = duplicate_groups.len();
813
814 let offset = req.offset.unwrap_or(0);
816 let duplicate_groups: Vec<DuplicateGroup> = duplicate_groups.into_iter().skip(offset).collect();
817
818 if let Some(limit) = req.limit {
819 let has_more = duplicate_groups.len() > limit;
820 let truncated: Vec<DuplicateGroup> = duplicate_groups.into_iter().take(limit).collect();
821 return Ok(Json(DetectDuplicatesResponse {
822 duplicates: truncated,
823 total,
824 has_more,
825 }));
826 }
827
828 Ok(Json(DetectDuplicatesResponse {
829 duplicates: duplicate_groups,
830 total,
831 has_more: false,
832 }))
833}
834
835#[derive(Deserialize)]
836pub struct EntityStateParams {
837 as_of: Option<chrono::DateTime<chrono::Utc>>,
838 tenant_id: Option<String>,
843}
844
845pub async fn get_entity_state(
846 State(store): State<SharedStore>,
847 Path(entity_id): Path<String>,
848 Query(params): Query<EntityStateParams>,
849) -> Result<Json<serde_json::Value>> {
850 let state = match params.tenant_id.as_deref() {
851 Some(tenant_id) => {
852 store.reconstruct_state_for_tenant(&entity_id, params.as_of, tenant_id)?
853 }
854 None => store.reconstruct_state(&entity_id, params.as_of)?,
855 };
856
857 tracing::info!("State reconstructed for entity: {}", entity_id);
858
859 Ok(Json(state))
860}
861
862pub async fn get_entity_snapshot(
863 State(store): State<SharedStore>,
864 Path(entity_id): Path<String>,
865 Query(params): Query<EntityStateParams>,
866) -> Result<Json<serde_json::Value>> {
867 let snapshot = match params.tenant_id.as_deref() {
871 Some(tenant_id) => store.reconstruct_state_for_tenant(&entity_id, None, tenant_id)?,
872 None => store.get_snapshot(&entity_id)?,
873 };
874
875 tracing::debug!("Snapshot retrieved for entity: {}", entity_id);
876
877 Ok(Json(snapshot))
878}
879
880#[derive(Debug, Deserialize)]
882pub struct StatsParams {
883 pub tenant_id: Option<String>,
889}
890
891pub async fn get_stats(
892 State(store): State<SharedStore>,
893 Query(params): Query<StatsParams>,
894) -> impl IntoResponse {
895 match params.tenant_id.as_deref() {
896 Some(tenant_id) => {
897 Json(serde_json::to_value(store.stats_for_tenant(tenant_id)).unwrap_or_default())
898 }
899 None => Json(serde_json::to_value(store.stats()).unwrap_or_default()),
900 }
901}
902
903#[derive(Debug, Deserialize)]
906pub struct ListStreamsParams {
907 pub tenant_id: Option<String>,
909 pub limit: Option<usize>,
911 pub offset: Option<usize>,
913}
914
915#[derive(Debug, serde::Serialize)]
917pub struct ListStreamsResponse {
918 pub streams: Vec<StreamInfo>,
919 pub total: usize,
920}
921
922pub async fn list_streams(
923 OptionalAuth(auth): OptionalAuth,
924 State(store): State<SharedStore>,
925 Query(params): Query<ListStreamsParams>,
926) -> Json<ListStreamsResponse> {
927 let tenant = params
929 .tenant_id
930 .clone()
931 .or_else(|| auth.as_ref().map(|a| a.tenant_id().to_string()))
932 .filter(|t| !t.is_empty());
933 let Some(tenant) = tenant else {
934 return Json(ListStreamsResponse {
935 streams: vec![],
936 total: 0,
937 });
938 };
939 let mut streams = store.list_streams_for_tenant(&tenant);
940 let total = streams.len();
941
942 streams.sort_by_key(|a| std::cmp::Reverse(a.last_event_at));
944
945 if let Some(offset) = params.offset {
947 if offset < streams.len() {
948 streams = streams[offset..].to_vec();
949 } else {
950 streams = vec![];
951 }
952 }
953
954 if let Some(limit) = params.limit {
955 streams.truncate(limit);
956 }
957
958 tracing::debug!("Listed {} streams (total: {})", streams.len(), total);
959
960 Json(ListStreamsResponse { streams, total })
961}
962
963#[derive(Debug, Deserialize)]
966pub struct ListEventTypesParams {
967 pub tenant_id: Option<String>,
969 pub limit: Option<usize>,
971 pub offset: Option<usize>,
973}
974
975#[derive(Debug, serde::Serialize)]
977pub struct ListEventTypesResponse {
978 pub event_types: Vec<EventTypeInfo>,
979 pub total: usize,
980}
981
982pub async fn list_event_types(
983 OptionalAuth(auth): OptionalAuth,
984 State(store): State<SharedStore>,
985 Query(params): Query<ListEventTypesParams>,
986) -> Json<ListEventTypesResponse> {
987 let tenant = params
989 .tenant_id
990 .clone()
991 .or_else(|| auth.as_ref().map(|a| a.tenant_id().to_string()))
992 .filter(|t| !t.is_empty());
993 let Some(tenant) = tenant else {
994 return Json(ListEventTypesResponse {
995 event_types: vec![],
996 total: 0,
997 });
998 };
999 let mut event_types = store.list_event_types_for_tenant(&tenant);
1000 let total = event_types.len();
1001
1002 event_types.sort_by_key(|a| std::cmp::Reverse(a.event_count));
1004
1005 if let Some(offset) = params.offset {
1007 if offset < event_types.len() {
1008 event_types = event_types[offset..].to_vec();
1009 } else {
1010 event_types = vec![];
1011 }
1012 }
1013
1014 if let Some(limit) = params.limit {
1015 event_types.truncate(limit);
1016 }
1017
1018 tracing::debug!(
1019 "Listed {} event types (total: {})",
1020 event_types.len(),
1021 total
1022 );
1023
1024 Json(ListEventTypesResponse { event_types, total })
1025}
1026
1027#[derive(Debug, Deserialize)]
1029pub struct WebSocketParams {
1030 pub consumer_id: Option<String>,
1031}
1032
1033pub async fn events_websocket(
1034 ws: WebSocketUpgrade,
1035 State(store): State<SharedStore>,
1036 Query(params): Query<WebSocketParams>,
1037) -> Response {
1038 let websocket_manager = store.websocket_manager();
1039
1040 ws.on_upgrade(move |socket| async move {
1041 if let Some(consumer_id) = params.consumer_id {
1042 websocket_manager
1043 .handle_socket_with_consumer(socket, consumer_id, store)
1044 .await;
1045 } else {
1046 websocket_manager.handle_socket(socket).await;
1047 }
1048 })
1049}
1050
1051pub async fn analytics_frequency(
1053 State(store): State<SharedStore>,
1054 Query(req): Query<EventFrequencyRequest>,
1055) -> Result<Json<EventFrequencyResponse>> {
1056 let response = AnalyticsEngine::event_frequency(&store, &req)?;
1057
1058 tracing::debug!(
1059 "Frequency analysis returned {} buckets",
1060 response.buckets.len()
1061 );
1062
1063 Ok(Json(response))
1064}
1065
1066pub async fn analytics_summary(
1068 State(store): State<SharedStore>,
1069 Query(req): Query<StatsSummaryRequest>,
1070) -> Result<Json<StatsSummaryResponse>> {
1071 let response = AnalyticsEngine::stats_summary(&store, &req)?;
1072
1073 tracing::debug!(
1074 "Stats summary: {} events across {} entities",
1075 response.total_events,
1076 response.unique_entities
1077 );
1078
1079 Ok(Json(response))
1080}
1081
1082pub async fn analytics_correlation(
1084 State(store): State<SharedStore>,
1085 Query(req): Query<CorrelationRequest>,
1086) -> Result<Json<CorrelationResponse>> {
1087 let response = AnalyticsEngine::analyze_correlation(&store, req)?;
1088
1089 tracing::debug!(
1090 "Correlation analysis: {}/{} correlated pairs ({:.2}%)",
1091 response.correlated_pairs,
1092 response.total_a,
1093 response.correlation_percentage
1094 );
1095
1096 Ok(Json(response))
1097}
1098
1099pub async fn create_snapshot(
1101 State(store): State<SharedStore>,
1102 Json(req): Json<CreateSnapshotRequest>,
1103) -> Result<Json<CreateSnapshotResponse>> {
1104 store.create_snapshot(&req.entity_id)?;
1105
1106 let snapshot_manager = store.snapshot_manager();
1107 let snapshot = snapshot_manager
1108 .get_latest_snapshot(&req.entity_id)
1109 .ok_or_else(|| crate::error::AllSourceError::EntityNotFound(req.entity_id.clone()))?;
1110
1111 tracing::info!("📸 Created snapshot for entity: {}", req.entity_id);
1112
1113 Ok(Json(CreateSnapshotResponse {
1114 snapshot_id: snapshot.id,
1115 entity_id: snapshot.entity_id,
1116 created_at: snapshot.created_at,
1117 event_count: snapshot.event_count,
1118 size_bytes: snapshot.metadata.size_bytes,
1119 }))
1120}
1121
1122pub async fn list_snapshots(
1124 State(store): State<SharedStore>,
1125 Query(req): Query<ListSnapshotsRequest>,
1126) -> Result<Json<ListSnapshotsResponse>> {
1127 let snapshot_manager = store.snapshot_manager();
1128
1129 let snapshots: Vec<SnapshotInfo> = if let Some(entity_id) = req.entity_id {
1130 snapshot_manager
1131 .get_all_snapshots(&entity_id)
1132 .into_iter()
1133 .map(SnapshotInfo::from)
1134 .collect()
1135 } else {
1136 let entities = snapshot_manager.list_entities();
1138 entities
1139 .iter()
1140 .flat_map(|entity_id| {
1141 snapshot_manager
1142 .get_all_snapshots(entity_id)
1143 .into_iter()
1144 .map(SnapshotInfo::from)
1145 })
1146 .collect()
1147 };
1148
1149 let total = snapshots.len();
1150
1151 tracing::debug!("Listed {} snapshots", total);
1152
1153 Ok(Json(ListSnapshotsResponse { snapshots, total }))
1154}
1155
1156pub async fn get_latest_snapshot(
1158 State(store): State<SharedStore>,
1159 Path(entity_id): Path<String>,
1160) -> Result<Json<serde_json::Value>> {
1161 let snapshot_manager = store.snapshot_manager();
1162
1163 let snapshot = snapshot_manager
1164 .get_latest_snapshot(&entity_id)
1165 .ok_or_else(|| crate::error::AllSourceError::EntityNotFound(entity_id.clone()))?;
1166
1167 tracing::debug!("Retrieved latest snapshot for entity: {}", entity_id);
1168
1169 Ok(Json(serde_json::json!({
1170 "snapshot_id": snapshot.id,
1171 "entity_id": snapshot.entity_id,
1172 "created_at": snapshot.created_at,
1173 "as_of": snapshot.as_of,
1174 "event_count": snapshot.event_count,
1175 "size_bytes": snapshot.metadata.size_bytes,
1176 "snapshot_type": snapshot.metadata.snapshot_type,
1177 "state": snapshot.state
1178 })))
1179}
1180
1181pub async fn trigger_compaction(
1183 State(store): State<SharedStore>,
1184) -> Result<Json<CompactionResult>> {
1185 let compaction_manager = store.compaction_manager().ok_or_else(|| {
1186 crate::error::AllSourceError::InternalError(
1187 "Compaction not enabled (no Parquet storage)".to_string(),
1188 )
1189 })?;
1190
1191 tracing::info!("📦 Manual compaction triggered via API");
1192
1193 let result = compaction_manager.compact_now()?;
1194
1195 Ok(Json(result))
1196}
1197
1198pub async fn compaction_stats(State(store): State<SharedStore>) -> Result<Json<serde_json::Value>> {
1200 let compaction_manager = store.compaction_manager().ok_or_else(|| {
1201 crate::error::AllSourceError::InternalError(
1202 "Compaction not enabled (no Parquet storage)".to_string(),
1203 )
1204 })?;
1205
1206 let stats = compaction_manager.stats();
1207 let config = compaction_manager.config();
1208
1209 Ok(Json(serde_json::json!({
1210 "stats": stats,
1211 "config": {
1212 "min_files_to_compact": config.min_files_to_compact,
1213 "target_file_size": config.target_file_size,
1214 "max_file_size": config.max_file_size,
1215 "small_file_threshold": config.small_file_threshold,
1216 "compaction_interval_seconds": config.compaction_interval_seconds,
1217 "auto_compact": config.auto_compact,
1218 "strategy": config.strategy
1219 }
1220 })))
1221}
1222
1223pub async fn register_schema(
1225 State(store): State<SharedStore>,
1226 Json(req): Json<RegisterSchemaRequest>,
1227) -> Result<Json<RegisterSchemaResponse>> {
1228 let schema_registry = store.schema_registry();
1229
1230 let response =
1231 schema_registry.register_schema(req.subject, req.schema, req.description, req.tags)?;
1232
1233 tracing::info!(
1234 "📋 Schema registered: v{} for '{}'",
1235 response.version,
1236 response.subject
1237 );
1238
1239 Ok(Json(response))
1240}
1241
1242#[derive(Deserialize)]
1244pub struct GetSchemaParams {
1245 version: Option<u32>,
1246}
1247
1248pub async fn get_schema(
1249 State(store): State<SharedStore>,
1250 Path(subject): Path<String>,
1251 Query(params): Query<GetSchemaParams>,
1252) -> Result<Json<serde_json::Value>> {
1253 let schema_registry = store.schema_registry();
1254
1255 let schema = schema_registry.get_schema(&subject, params.version)?;
1256
1257 tracing::debug!("Retrieved schema v{} for '{}'", schema.version, subject);
1258
1259 Ok(Json(serde_json::json!({
1260 "id": schema.id,
1261 "subject": schema.subject,
1262 "version": schema.version,
1263 "schema": schema.schema,
1264 "created_at": schema.created_at,
1265 "description": schema.description,
1266 "tags": schema.tags
1267 })))
1268}
1269
1270pub async fn list_schema_versions(
1272 State(store): State<SharedStore>,
1273 Path(subject): Path<String>,
1274) -> Result<Json<serde_json::Value>> {
1275 let schema_registry = store.schema_registry();
1276
1277 let versions = schema_registry.list_versions(&subject)?;
1278
1279 Ok(Json(serde_json::json!({
1280 "subject": subject,
1281 "versions": versions
1282 })))
1283}
1284
1285pub async fn list_subjects(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1287 let schema_registry = store.schema_registry();
1288
1289 let subjects = schema_registry.list_subjects();
1290
1291 Json(serde_json::json!({
1292 "subjects": subjects,
1293 "total": subjects.len()
1294 }))
1295}
1296
1297pub async fn validate_event_schema(
1299 State(store): State<SharedStore>,
1300 Json(req): Json<ValidateEventRequest>,
1301) -> Result<Json<ValidateEventResponse>> {
1302 let schema_registry = store.schema_registry();
1303
1304 let response = schema_registry.validate(&req.subject, req.version, &req.payload)?;
1305
1306 if response.valid {
1307 tracing::debug!(
1308 "✅ Event validated against schema '{}' v{}",
1309 req.subject,
1310 response.schema_version
1311 );
1312 } else {
1313 tracing::warn!(
1314 "❌ Event validation failed for '{}': {:?}",
1315 req.subject,
1316 response.errors
1317 );
1318 }
1319
1320 Ok(Json(response))
1321}
1322
1323#[derive(Deserialize)]
1325pub struct SetCompatibilityRequest {
1326 compatibility: CompatibilityMode,
1327}
1328
1329pub async fn set_compatibility_mode(
1330 State(store): State<SharedStore>,
1331 Path(subject): Path<String>,
1332 Json(req): Json<SetCompatibilityRequest>,
1333) -> Json<serde_json::Value> {
1334 let schema_registry = store.schema_registry();
1335
1336 schema_registry.set_compatibility_mode(subject.clone(), req.compatibility);
1337
1338 tracing::info!(
1339 "🔧 Set compatibility mode for '{}' to {:?}",
1340 subject,
1341 req.compatibility
1342 );
1343
1344 Json(serde_json::json!({
1345 "subject": subject,
1346 "compatibility": req.compatibility
1347 }))
1348}
1349
1350pub async fn start_replay(
1352 State(store): State<SharedStore>,
1353 Json(req): Json<StartReplayRequest>,
1354) -> Result<Json<StartReplayResponse>> {
1355 let replay_manager = store.replay_manager();
1356
1357 let response = replay_manager.start_replay(store, req)?;
1358
1359 tracing::info!(
1360 "🔄 Started replay {} with {} events",
1361 response.replay_id,
1362 response.total_events
1363 );
1364
1365 Ok(Json(response))
1366}
1367
1368pub async fn get_replay_progress(
1370 State(store): State<SharedStore>,
1371 Path(replay_id): Path<uuid::Uuid>,
1372) -> Result<Json<ReplayProgress>> {
1373 let replay_manager = store.replay_manager();
1374
1375 let progress = replay_manager.get_progress(replay_id)?;
1376
1377 Ok(Json(progress))
1378}
1379
1380pub async fn list_replays(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1382 let replay_manager = store.replay_manager();
1383
1384 let replays = replay_manager.list_replays();
1385
1386 Json(serde_json::json!({
1387 "replays": replays,
1388 "total": replays.len()
1389 }))
1390}
1391
1392pub async fn cancel_replay(
1394 State(store): State<SharedStore>,
1395 Path(replay_id): Path<uuid::Uuid>,
1396) -> Result<Json<serde_json::Value>> {
1397 let replay_manager = store.replay_manager();
1398
1399 replay_manager.cancel_replay(replay_id)?;
1400
1401 tracing::info!("🛑 Cancelled replay {}", replay_id);
1402
1403 Ok(Json(serde_json::json!({
1404 "replay_id": replay_id,
1405 "status": "cancelled"
1406 })))
1407}
1408
1409pub async fn delete_replay(
1411 State(store): State<SharedStore>,
1412 Path(replay_id): Path<uuid::Uuid>,
1413) -> Result<Json<serde_json::Value>> {
1414 let replay_manager = store.replay_manager();
1415
1416 let deleted = replay_manager.delete_replay(replay_id)?;
1417
1418 if deleted {
1419 tracing::info!("🗑️ Deleted replay {}", replay_id);
1420 }
1421
1422 Ok(Json(serde_json::json!({
1423 "replay_id": replay_id,
1424 "deleted": deleted
1425 })))
1426}
1427
1428pub async fn register_pipeline(
1430 State(store): State<SharedStore>,
1431 Json(config): Json<PipelineConfig>,
1432) -> Result<Json<serde_json::Value>> {
1433 let pipeline_manager = store.pipeline_manager();
1434
1435 let pipeline_id = pipeline_manager.register(config.clone());
1436
1437 tracing::info!(
1438 "🔀 Pipeline registered: {} (name: {})",
1439 pipeline_id,
1440 config.name
1441 );
1442
1443 Ok(Json(serde_json::json!({
1444 "pipeline_id": pipeline_id,
1445 "name": config.name,
1446 "enabled": config.enabled
1447 })))
1448}
1449
1450pub async fn list_pipelines(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1452 let pipeline_manager = store.pipeline_manager();
1453
1454 let pipelines = pipeline_manager.list();
1455
1456 tracing::debug!("Listed {} pipelines", pipelines.len());
1457
1458 Json(serde_json::json!({
1459 "pipelines": pipelines,
1460 "total": pipelines.len()
1461 }))
1462}
1463
1464pub async fn get_pipeline(
1466 State(store): State<SharedStore>,
1467 Path(pipeline_id): Path<uuid::Uuid>,
1468) -> Result<Json<PipelineConfig>> {
1469 let pipeline_manager = store.pipeline_manager();
1470
1471 let pipeline = pipeline_manager.get(pipeline_id).ok_or_else(|| {
1472 crate::error::AllSourceError::ValidationError(format!("Pipeline not found: {pipeline_id}"))
1473 })?;
1474
1475 Ok(Json(pipeline.config().clone()))
1476}
1477
1478pub async fn remove_pipeline(
1480 State(store): State<SharedStore>,
1481 Path(pipeline_id): Path<uuid::Uuid>,
1482) -> Result<Json<serde_json::Value>> {
1483 let pipeline_manager = store.pipeline_manager();
1484
1485 let removed = pipeline_manager.remove(pipeline_id);
1486
1487 if removed {
1488 tracing::info!("🗑️ Removed pipeline {}", pipeline_id);
1489 }
1490
1491 Ok(Json(serde_json::json!({
1492 "pipeline_id": pipeline_id,
1493 "removed": removed
1494 })))
1495}
1496
1497pub async fn all_pipeline_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1499 let pipeline_manager = store.pipeline_manager();
1500
1501 let stats = pipeline_manager.all_stats();
1502
1503 Json(serde_json::json!({
1504 "stats": stats,
1505 "total": stats.len()
1506 }))
1507}
1508
1509pub async fn get_pipeline_stats(
1511 State(store): State<SharedStore>,
1512 Path(pipeline_id): Path<uuid::Uuid>,
1513) -> Result<Json<PipelineStats>> {
1514 let pipeline_manager = store.pipeline_manager();
1515
1516 let pipeline = pipeline_manager.get(pipeline_id).ok_or_else(|| {
1517 crate::error::AllSourceError::ValidationError(format!("Pipeline not found: {pipeline_id}"))
1518 })?;
1519
1520 Ok(Json(pipeline.stats()))
1521}
1522
1523pub async fn reset_pipeline(
1525 State(store): State<SharedStore>,
1526 Path(pipeline_id): Path<uuid::Uuid>,
1527) -> Result<Json<serde_json::Value>> {
1528 let pipeline_manager = store.pipeline_manager();
1529
1530 let pipeline = pipeline_manager.get(pipeline_id).ok_or_else(|| {
1531 crate::error::AllSourceError::ValidationError(format!("Pipeline not found: {pipeline_id}"))
1532 })?;
1533
1534 pipeline.reset();
1535
1536 tracing::info!("🔄 Reset pipeline {}", pipeline_id);
1537
1538 Ok(Json(serde_json::json!({
1539 "pipeline_id": pipeline_id,
1540 "reset": true
1541 })))
1542}
1543
1544pub async fn get_event_by_id(
1550 State(store): State<SharedStore>,
1551 Path(event_id): Path<uuid::Uuid>,
1552) -> Result<Json<serde_json::Value>> {
1553 let event = store.get_event_by_id(&event_id)?.ok_or_else(|| {
1554 crate::error::AllSourceError::EntityNotFound(format!("Event '{event_id}' not found"))
1555 })?;
1556
1557 let dto = EventDto::from(&event);
1558
1559 tracing::debug!("Event retrieved by ID: {}", event_id);
1560
1561 Ok(Json(serde_json::json!({
1562 "event": dto,
1563 "found": true
1564 })))
1565}
1566
1567pub async fn list_projections(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1573 let projection_manager = store.projection_manager();
1574 let status_map = store.projection_status();
1575
1576 let projections: Vec<serde_json::Value> = projection_manager
1577 .list_projections()
1578 .iter()
1579 .map(|(name, projection)| {
1580 let status = status_map
1581 .get(name)
1582 .map_or_else(|| "running".to_string(), |s| s.value().clone());
1583 serde_json::json!({
1584 "name": name,
1585 "type": format!("{:?}", projection.name()),
1586 "status": status,
1587 })
1588 })
1589 .collect();
1590
1591 tracing::debug!("Listed {} projections", projections.len());
1592
1593 Json(serde_json::json!({
1594 "projections": projections,
1595 "total": projections.len()
1596 }))
1597}
1598
1599pub async fn get_projection(
1601 State(store): State<SharedStore>,
1602 Path(name): Path<String>,
1603) -> Result<Json<serde_json::Value>> {
1604 let projection_manager = store.projection_manager();
1605
1606 let projection = projection_manager.get_projection(&name).ok_or_else(|| {
1607 crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1608 })?;
1609
1610 Ok(Json(serde_json::json!({
1611 "name": projection.name(),
1612 "found": true
1613 })))
1614}
1615
1616pub async fn get_projection_state(
1629 State(store): State<SharedStore>,
1630 Path((name, entity_id)): Path<(String, String)>,
1631) -> Result<Json<serde_json::Value>> {
1632 let state = store
1633 .projection_manager()
1634 .get_projection(&name)
1635 .and_then(|p| p.get_state(&entity_id))
1636 .or_else(|| {
1637 store
1638 .projection_state_cache()
1639 .get(&format!("{name}:{entity_id}"))
1640 .map(|entry| entry.value().clone())
1641 });
1642
1643 tracing::debug!("Projection state retrieved: {} / {}", name, entity_id);
1644
1645 Ok(Json(serde_json::json!({
1646 "projection": name,
1647 "entity_id": entity_id,
1648 "state": state,
1649 "found": state.is_some()
1650 })))
1651}
1652
1653pub async fn delete_projection(
1658 State(store): State<SharedStore>,
1659 Path(name): Path<String>,
1660) -> Result<Json<serde_json::Value>> {
1661 let projection_manager = store.projection_manager();
1662
1663 let projection = projection_manager.get_projection(&name).ok_or_else(|| {
1664 crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1665 })?;
1666
1667 projection.clear();
1668
1669 let cache = store.projection_state_cache();
1671 let prefix = format!("{name}:");
1672 let keys_to_remove: Vec<String> = cache
1673 .iter()
1674 .filter(|entry| entry.key().starts_with(&prefix))
1675 .map(|entry| entry.key().clone())
1676 .collect();
1677 for key in keys_to_remove {
1678 cache.remove(&key);
1679 }
1680
1681 tracing::info!("Projection deleted (cleared): {}", name);
1682
1683 Ok(Json(serde_json::json!({
1684 "projection": name,
1685 "deleted": true
1686 })))
1687}
1688
1689#[derive(Debug, Default, Deserialize)]
1696pub struct ProjectionStateSummaryParams {
1697 pub limit: Option<usize>,
1699 pub offset: Option<usize>,
1701 pub entity_id_prefix: Option<String>,
1704}
1705
1706pub async fn get_projection_state_summary(
1722 State(store): State<SharedStore>,
1723 Path(name): Path<String>,
1724 Query(params): Query<ProjectionStateSummaryParams>,
1725) -> Result<Json<serde_json::Value>> {
1726 let cache = store.projection_state_cache();
1727 let prefix = format!("{name}:");
1728 let offset = params.offset.unwrap_or(0);
1729
1730 let mut entity_ids: Vec<String> = cache
1735 .iter()
1736 .filter_map(|entry| entry.key().strip_prefix(&prefix).map(ToString::to_string))
1737 .filter(|entity_id| {
1738 params
1739 .entity_id_prefix
1740 .as_ref()
1741 .is_none_or(|p| entity_id.starts_with(p))
1742 })
1743 .collect();
1744 entity_ids.sort_unstable();
1745
1746 let total = entity_ids.len();
1747
1748 let page = entity_ids.into_iter().skip(offset);
1749 let page: Vec<String> = match params.limit {
1750 Some(limit) => page.take(limit).collect(),
1751 None => page.collect(),
1752 };
1753
1754 let states: Vec<serde_json::Value> = page
1755 .into_iter()
1756 .filter_map(|entity_id| {
1757 cache.get(&format!("{prefix}{entity_id}")).map(|entry| {
1759 serde_json::json!({
1760 "entity_id": entity_id,
1761 "state": entry.value().clone()
1762 })
1763 })
1764 })
1765 .collect();
1766
1767 let count = states.len();
1768 let has_more = offset + count < total;
1771
1772 tracing::debug!(
1773 "Projection state summary: {} ({} of {} entities, offset {})",
1774 name,
1775 count,
1776 total,
1777 offset
1778 );
1779
1780 Ok(Json(serde_json::json!({
1781 "projection": name,
1782 "states": states,
1783 "count": count,
1784 "total": total,
1785 "has_more": has_more
1786 })))
1787}
1788
1789pub async fn reset_projection(
1793 State(store): State<SharedStore>,
1794 Path(name): Path<String>,
1795) -> Result<Json<serde_json::Value>> {
1796 let reprocessed = store.reset_projection(&name)?;
1797
1798 tracing::info!(
1799 "Projection reset: {} ({} events reprocessed)",
1800 name,
1801 reprocessed
1802 );
1803
1804 Ok(Json(serde_json::json!({
1805 "projection": name,
1806 "reset": true,
1807 "events_reprocessed": reprocessed
1808 })))
1809}
1810
1811pub async fn pause_projection(
1815 State(store): State<SharedStore>,
1816 Path(name): Path<String>,
1817) -> Result<Json<serde_json::Value>> {
1818 let projection_manager = store.projection_manager();
1819
1820 let _projection = projection_manager.get_projection(&name).ok_or_else(|| {
1822 crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1823 })?;
1824
1825 store
1826 .projection_status()
1827 .insert(name.clone(), "paused".to_string());
1828
1829 tracing::info!("Projection paused: {}", name);
1830
1831 Ok(Json(serde_json::json!({
1832 "projection": name,
1833 "status": "paused"
1834 })))
1835}
1836
1837pub async fn start_projection(
1841 State(store): State<SharedStore>,
1842 Path(name): Path<String>,
1843) -> Result<Json<serde_json::Value>> {
1844 let projection_manager = store.projection_manager();
1845
1846 let _projection = projection_manager.get_projection(&name).ok_or_else(|| {
1848 crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1849 })?;
1850
1851 store
1852 .projection_status()
1853 .insert(name.clone(), "running".to_string());
1854
1855 tracing::info!("Projection started: {}", name);
1856
1857 Ok(Json(serde_json::json!({
1858 "projection": name,
1859 "status": "running"
1860 })))
1861}
1862
1863#[derive(Debug, Deserialize)]
1865pub struct SaveProjectionStateRequest {
1866 pub state: serde_json::Value,
1867}
1868
1869pub async fn save_projection_state(
1874 State(store): State<SharedStore>,
1875 Path((name, entity_id)): Path<(String, String)>,
1876 Json(req): Json<SaveProjectionStateRequest>,
1877) -> Result<Json<serde_json::Value>> {
1878 let projection_cache = store.projection_state_cache();
1879
1880 projection_cache.insert(format!("{name}:{entity_id}"), req.state.clone());
1882
1883 tracing::info!("Projection state saved: {} / {}", name, entity_id);
1884
1885 Ok(Json(serde_json::json!({
1886 "projection": name,
1887 "entity_id": entity_id,
1888 "saved": true
1889 })))
1890}
1891
1892#[derive(Debug, Deserialize)]
1896pub struct BulkGetStateRequest {
1897 pub entity_ids: Vec<String>,
1898}
1899
1900#[derive(Debug, Deserialize)]
1904pub struct BulkSaveStateRequest {
1905 pub states: Vec<BulkSaveStateItem>,
1906}
1907
1908#[derive(Debug, Deserialize)]
1909pub struct BulkSaveStateItem {
1910 pub entity_id: String,
1911 pub state: serde_json::Value,
1912}
1913
1914pub async fn bulk_get_projection_states(
1915 State(store): State<SharedStore>,
1916 Path(name): Path<String>,
1917 Json(req): Json<BulkGetStateRequest>,
1918) -> Result<Json<serde_json::Value>> {
1919 let projection = store.projection_manager().get_projection(&name);
1923 let cache = store.projection_state_cache();
1924
1925 let states: Vec<serde_json::Value> = req
1926 .entity_ids
1927 .iter()
1928 .map(|entity_id| {
1929 let state = projection
1930 .as_ref()
1931 .and_then(|p| p.get_state(entity_id))
1932 .or_else(|| {
1933 cache
1934 .get(&format!("{name}:{entity_id}"))
1935 .map(|entry| entry.value().clone())
1936 });
1937 serde_json::json!({
1938 "entity_id": entity_id,
1939 "state": state,
1940 "found": state.is_some()
1941 })
1942 })
1943 .collect();
1944
1945 tracing::debug!(
1946 "Bulk projection state retrieved: {} entities from {}",
1947 states.len(),
1948 name
1949 );
1950
1951 Ok(Json(serde_json::json!({
1952 "projection": name,
1953 "states": states,
1954 "total": states.len()
1955 })))
1956}
1957
1958pub async fn bulk_save_projection_states(
1963 State(store): State<SharedStore>,
1964 Path(name): Path<String>,
1965 Json(req): Json<BulkSaveStateRequest>,
1966) -> Result<Json<serde_json::Value>> {
1967 let projection_cache = store.projection_state_cache();
1968
1969 let mut saved_count = 0;
1970 for item in &req.states {
1971 projection_cache.insert(format!("{name}:{}", item.entity_id), item.state.clone());
1972 saved_count += 1;
1973 }
1974
1975 tracing::info!(
1976 "Bulk projection state saved: {} entities for {}",
1977 saved_count,
1978 name
1979 );
1980
1981 Ok(Json(serde_json::json!({
1982 "projection": name,
1983 "saved": saved_count,
1984 "total": req.states.len()
1985 })))
1986}
1987
1988#[derive(Debug, Deserialize)]
1994pub struct ListWebhooksParams {
1995 pub tenant_id: Option<String>,
1996}
1997
1998pub async fn register_webhook(
2000 State(store): State<SharedStore>,
2001 Json(req): Json<RegisterWebhookRequest>,
2002) -> Json<serde_json::Value> {
2003 let registry = store.webhook_registry();
2004 let webhook = registry.register(req);
2005
2006 tracing::info!("Webhook registered: {} -> {}", webhook.id, webhook.url);
2007
2008 Json(serde_json::json!({
2009 "webhook": webhook,
2010 "created": true
2011 }))
2012}
2013
2014pub async fn list_webhooks(
2016 State(store): State<SharedStore>,
2017 Query(params): Query<ListWebhooksParams>,
2018) -> Json<serde_json::Value> {
2019 let registry = store.webhook_registry();
2020
2021 let webhooks = if let Some(tenant_id) = params.tenant_id {
2022 registry.list_by_tenant(&tenant_id)
2023 } else {
2024 vec![]
2026 };
2027
2028 let total = webhooks.len();
2029
2030 Json(serde_json::json!({
2031 "webhooks": webhooks,
2032 "total": total
2033 }))
2034}
2035
2036pub async fn get_webhook(
2038 State(store): State<SharedStore>,
2039 Path(webhook_id): Path<uuid::Uuid>,
2040) -> Result<Json<serde_json::Value>> {
2041 let registry = store.webhook_registry();
2042
2043 let webhook = registry.get(webhook_id).ok_or_else(|| {
2044 crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2045 })?;
2046
2047 Ok(Json(serde_json::json!({
2048 "webhook": webhook,
2049 "found": true
2050 })))
2051}
2052
2053pub async fn update_webhook(
2055 State(store): State<SharedStore>,
2056 Path(webhook_id): Path<uuid::Uuid>,
2057 Json(req): Json<UpdateWebhookRequest>,
2058) -> Result<Json<serde_json::Value>> {
2059 let registry = store.webhook_registry();
2060
2061 let webhook = registry.update(webhook_id, req).ok_or_else(|| {
2062 crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2063 })?;
2064
2065 tracing::info!("Webhook updated: {}", webhook_id);
2066
2067 Ok(Json(serde_json::json!({
2068 "webhook": webhook,
2069 "updated": true
2070 })))
2071}
2072
2073pub async fn delete_webhook(
2075 State(store): State<SharedStore>,
2076 Path(webhook_id): Path<uuid::Uuid>,
2077) -> Result<Json<serde_json::Value>> {
2078 let registry = store.webhook_registry();
2079
2080 let webhook = registry.delete(webhook_id).ok_or_else(|| {
2081 crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2082 })?;
2083
2084 tracing::info!("Webhook deleted: {} ({})", webhook_id, webhook.url);
2085
2086 Ok(Json(serde_json::json!({
2087 "webhook_id": webhook_id,
2088 "deleted": true
2089 })))
2090}
2091
2092#[derive(Debug, Deserialize)]
2094pub struct ListDeliveriesParams {
2095 pub limit: Option<usize>,
2096}
2097
2098pub async fn list_webhook_deliveries(
2100 State(store): State<SharedStore>,
2101 Path(webhook_id): Path<uuid::Uuid>,
2102 Query(params): Query<ListDeliveriesParams>,
2103) -> Result<Json<serde_json::Value>> {
2104 let registry = store.webhook_registry();
2105
2106 registry.get(webhook_id).ok_or_else(|| {
2108 crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2109 })?;
2110
2111 let limit = params.limit.unwrap_or(50);
2112 let deliveries = registry.get_deliveries(webhook_id, limit);
2113 let total = deliveries.len();
2114
2115 Ok(Json(serde_json::json!({
2116 "webhook_id": webhook_id,
2117 "deliveries": deliveries,
2118 "total": total
2119 })))
2120}
2121
2122#[cfg(feature = "analytics")]
2128pub async fn eventql_query(
2129 State(store): State<SharedStore>,
2130 Json(req): Json<crate::infrastructure::query::eventql::EventQLRequest>,
2131) -> Result<Json<serde_json::Value>> {
2132 let events = store.snapshot_events();
2133 match crate::infrastructure::query::eventql::execute_eventql(&events, &req).await {
2134 Ok(response) => Ok(Json(serde_json::json!({
2135 "columns": response.columns,
2136 "rows": response.rows,
2137 "row_count": response.row_count,
2138 }))),
2139 Err(e) => Err(crate::error::AllSourceError::InvalidQuery(e)),
2140 }
2141}
2142
2143pub async fn graphql_query(
2145 State(store): State<SharedStore>,
2146 Json(req): Json<GraphQLRequest>,
2147) -> Json<serde_json::Value> {
2148 let fields = match crate::infrastructure::query::graphql::parse_query(&req.query) {
2149 Ok(f) => f,
2150 Err(e) => {
2151 return Json(
2152 serde_json::to_value(GraphQLResponse {
2153 data: None,
2154 errors: vec![GraphQLError { message: e }],
2155 })
2156 .unwrap(),
2157 );
2158 }
2159 };
2160
2161 let mut data = serde_json::Map::new();
2162 let mut errors = Vec::new();
2163
2164 for field in &fields {
2165 match field.name.as_str() {
2166 "events" => {
2167 let request = crate::application::dto::QueryEventsRequest {
2168 entity_id: field.arguments.get("entity_id").cloned(),
2169 event_type: field.arguments.get("event_type").cloned(),
2170 tenant_id: field.arguments.get("tenant_id").cloned(),
2171 limit: field.arguments.get("limit").and_then(|l| l.parse().ok()),
2172 as_of: None,
2173 since: None,
2174 until: None,
2175 event_type_prefix: None,
2176 exclude_event_type_prefix: None,
2177 payload_filter: None,
2178 };
2179 match store.query(&request) {
2180 Ok(events) => {
2181 let json_events: Vec<serde_json::Value> = events
2182 .iter()
2183 .map(|e| {
2184 crate::infrastructure::query::graphql::event_to_json(
2185 e,
2186 &field.fields,
2187 )
2188 })
2189 .collect();
2190 data.insert("events".to_string(), serde_json::Value::Array(json_events));
2191 }
2192 Err(e) => errors.push(GraphQLError {
2193 message: format!("events query failed: {e}"),
2194 }),
2195 }
2196 }
2197 "event" => {
2198 if let Some(id_str) = field.arguments.get("id") {
2199 if let Ok(id) = uuid::Uuid::parse_str(id_str) {
2200 match store.get_event_by_id(&id) {
2201 Ok(Some(event)) => {
2202 data.insert(
2203 "event".to_string(),
2204 crate::infrastructure::query::graphql::event_to_json(
2205 &event,
2206 &field.fields,
2207 ),
2208 );
2209 }
2210 Ok(None) => {
2211 data.insert("event".to_string(), serde_json::Value::Null);
2212 }
2213 Err(e) => errors.push(GraphQLError {
2214 message: format!("event lookup failed: {e}"),
2215 }),
2216 }
2217 } else {
2218 errors.push(GraphQLError {
2219 message: format!("Invalid UUID: {id_str}"),
2220 });
2221 }
2222 } else {
2223 errors.push(GraphQLError {
2224 message: "event query requires 'id' argument".to_string(),
2225 });
2226 }
2227 }
2228 "projections" => {
2229 let pm = store.projection_manager();
2230 let names: Vec<serde_json::Value> = pm
2231 .list_projections()
2232 .iter()
2233 .map(|(name, _)| serde_json::Value::String(name.clone()))
2234 .collect();
2235 data.insert("projections".to_string(), serde_json::Value::Array(names));
2236 }
2237 "stats" => {
2238 let stats = store.stats();
2239 data.insert(
2240 "stats".to_string(),
2241 serde_json::json!({
2242 "total_events": stats.total_events,
2243 "total_entities": stats.total_entities,
2244 "total_event_types": stats.total_event_types,
2245 }),
2246 );
2247 }
2248 "__schema" => {
2249 data.insert(
2250 "__schema".to_string(),
2251 crate::infrastructure::query::graphql::introspection_schema(),
2252 );
2253 }
2254 other => {
2255 errors.push(GraphQLError {
2256 message: format!("Unknown field: {other}"),
2257 });
2258 }
2259 }
2260 }
2261
2262 Json(
2263 serde_json::to_value(GraphQLResponse {
2264 data: Some(serde_json::Value::Object(data)),
2265 errors,
2266 })
2267 .unwrap(),
2268 )
2269}
2270
2271pub async fn geo_query(
2273 State(store): State<SharedStore>,
2274 Json(req): Json<GeoQueryRequest>,
2275) -> Json<serde_json::Value> {
2276 let events = store.snapshot_events();
2277 let geo_index = store.geo_index();
2278 let results =
2279 crate::infrastructure::query::geospatial::execute_geo_query(&events, &geo_index, &req);
2280 let total = results.len();
2281 Json(serde_json::json!({
2282 "results": results,
2283 "total": total,
2284 }))
2285}
2286
2287pub async fn geo_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
2289 let stats = store.geo_index().stats();
2290 Json(serde_json::json!(stats))
2291}
2292
2293pub async fn exactly_once_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
2295 let stats = store.exactly_once().stats();
2296 Json(serde_json::json!(stats))
2297}
2298
2299pub async fn schema_evolution_history(
2301 State(store): State<SharedStore>,
2302 Path(event_type): Path<String>,
2303) -> Json<serde_json::Value> {
2304 let mgr = store.schema_evolution();
2305 let history = mgr.get_history(&event_type);
2306 let version = mgr.get_version(&event_type);
2307 Json(serde_json::json!({
2308 "event_type": event_type,
2309 "current_version": version,
2310 "history": history,
2311 }))
2312}
2313
2314pub async fn schema_evolution_schema(
2316 State(store): State<SharedStore>,
2317 Path(event_type): Path<String>,
2318) -> Json<serde_json::Value> {
2319 let mgr = store.schema_evolution();
2320 if let Some(schema) = mgr.get_schema(&event_type) {
2321 let json_schema = crate::application::services::schema_evolution::to_json_schema(&schema);
2322 Json(serde_json::json!({
2323 "event_type": event_type,
2324 "version": mgr.get_version(&event_type),
2325 "inferred_schema": schema,
2326 "json_schema": json_schema,
2327 }))
2328 } else {
2329 Json(serde_json::json!({
2330 "event_type": event_type,
2331 "error": "No schema inferred for this event type"
2332 }))
2333 }
2334}
2335
2336pub async fn schema_evolution_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
2338 let stats = store.schema_evolution().stats();
2339 let event_types = store.schema_evolution().list_event_types();
2340 Json(serde_json::json!({
2341 "stats": stats,
2342 "tracked_event_types": event_types,
2343 }))
2344}
2345
2346#[cfg(feature = "embedded-sync")]
2352pub async fn sync_pull_handler(
2353 State(state): State<AppState>,
2354 Json(request): Json<crate::embedded::sync_types::SyncPullRequest>,
2355) -> Result<Json<crate::embedded::sync_types::SyncPullResponse>> {
2356 use crate::infrastructure::cluster::{crdt::ReplicatedEvent, hlc::HlcTimestamp};
2357
2358 let store = &state.store;
2359
2360 let since = request
2363 .version_vector
2364 .values()
2365 .map(|ts| ts.physical_ms)
2366 .min()
2367 .and_then(|ms| chrono::DateTime::from_timestamp_millis(ms as i64));
2368
2369 let events = store.query(&crate::application::dto::QueryEventsRequest {
2370 entity_id: None,
2371 event_type: None,
2372 tenant_id: None,
2373 as_of: None,
2374 since,
2375 until: None,
2376 limit: None,
2377 event_type_prefix: None,
2378 exclude_event_type_prefix: None,
2379 payload_filter: None,
2380 })?;
2381
2382 let mut replicated = Vec::with_capacity(events.len());
2384 let mut last_ms = 0u64;
2385 let mut logical = 0u32;
2386
2387 for event in &events {
2388 let event_ms = event.timestamp().timestamp_millis() as u64;
2389 if event_ms == last_ms {
2390 logical += 1;
2391 } else {
2392 last_ms = event_ms;
2393 logical = 0;
2394 }
2395
2396 replicated.push(ReplicatedEvent {
2397 event_id: event.id().to_string(),
2398 hlc_timestamp: HlcTimestamp::new(event_ms, logical, 0),
2399 origin_region: "server".to_string(),
2400 event_data: serde_json::json!({
2401 "event_type": event.event_type_str(),
2402 "entity_id": event.entity_id_str(),
2403 "tenant_id": event.tenant_id_str(),
2404 "payload": event.payload,
2405 "metadata": event.metadata,
2406 }),
2407 });
2408 }
2409
2410 Ok(Json(crate::embedded::sync_types::SyncPullResponse {
2411 events: replicated,
2412 version_vector: std::collections::BTreeMap::new(),
2413 }))
2414}
2415
2416#[cfg(feature = "embedded-sync")]
2418pub async fn sync_push_handler(
2419 State(state): State<AppState>,
2420 Json(request): Json<crate::embedded::sync_types::SyncPushRequest>,
2421) -> Result<Json<crate::embedded::sync_types::SyncPushResponse>> {
2422 let store = &state.store;
2423
2424 let mut accepted = 0usize;
2425 let mut skipped = 0usize;
2426
2427 for rep_event in &request.events {
2428 let event_data = &rep_event.event_data;
2429 let event_type = event_data
2430 .get("event_type")
2431 .and_then(|v| v.as_str())
2432 .unwrap_or("unknown")
2433 .to_string();
2434 let entity_id = event_data
2435 .get("entity_id")
2436 .and_then(|v| v.as_str())
2437 .unwrap_or("unknown")
2438 .to_string();
2439 let tenant_id = event_data
2440 .get("tenant_id")
2441 .and_then(|v| v.as_str())
2442 .unwrap_or("default")
2443 .to_string();
2444 let payload = event_data
2445 .get("payload")
2446 .cloned()
2447 .unwrap_or(serde_json::json!({}));
2448 let metadata = event_data.get("metadata").cloned();
2449
2450 match Event::from_strings(event_type, entity_id, tenant_id, payload, metadata) {
2451 Ok(domain_event) => {
2452 store.ingest(&domain_event)?;
2453 accepted += 1;
2454 }
2455 Err(_) => {
2456 skipped += 1;
2457 }
2458 }
2459 }
2460
2461 Ok(Json(crate::embedded::sync_types::SyncPushResponse {
2462 accepted,
2463 skipped,
2464 version_vector: std::collections::BTreeMap::new(),
2465 }))
2466}
2467
2468pub async fn register_consumer(
2474 State(store): State<SharedStore>,
2475 Json(req): Json<RegisterConsumerRequest>,
2476) -> Result<Json<ConsumerResponse>> {
2477 let consumer = store
2478 .consumer_registry()
2479 .register(&req.consumer_id, &req.event_type_filters);
2480
2481 Ok(Json(ConsumerResponse {
2482 consumer_id: consumer.consumer_id,
2483 event_type_filters: consumer.event_type_filters,
2484 cursor_position: consumer.cursor_position,
2485 }))
2486}
2487
2488pub async fn get_consumer(
2490 State(store): State<SharedStore>,
2491 Path(consumer_id): Path<String>,
2492) -> Result<Json<ConsumerResponse>> {
2493 let consumer = store.consumer_registry().get_or_create(&consumer_id);
2494
2495 Ok(Json(ConsumerResponse {
2496 consumer_id: consumer.consumer_id,
2497 event_type_filters: consumer.event_type_filters,
2498 cursor_position: consumer.cursor_position,
2499 }))
2500}
2501
2502#[derive(Debug, Deserialize)]
2504pub struct ConsumerPollQuery {
2505 pub limit: Option<usize>,
2506}
2507
2508pub async fn poll_consumer_events(
2509 State(store): State<SharedStore>,
2510 Path(consumer_id): Path<String>,
2511 Query(query): Query<ConsumerPollQuery>,
2512) -> Result<Json<ConsumerEventsResponse>> {
2513 let consumer = store.consumer_registry().get_or_create(&consumer_id);
2514 let offset = consumer.cursor_position.unwrap_or(0);
2515 let limit = query.limit.unwrap_or(100);
2516
2517 let events = store.events_after_offset(offset, &consumer.event_type_filters, limit);
2518 let count = events.len();
2519
2520 let consumer_events: Vec<ConsumerEventDto> = events
2521 .into_iter()
2522 .map(|(position, event)| ConsumerEventDto {
2523 position,
2524 event: EventDto::from(&event),
2525 })
2526 .collect();
2527
2528 Ok(Json(ConsumerEventsResponse {
2529 events: consumer_events,
2530 count,
2531 }))
2532}
2533
2534pub async fn ack_consumer(
2536 State(store): State<SharedStore>,
2537 Path(consumer_id): Path<String>,
2538 Json(req): Json<AckRequest>,
2539) -> Result<Json<serde_json::Value>> {
2540 let max_offset = store.total_events() as u64;
2541
2542 store
2543 .consumer_registry()
2544 .ack(&consumer_id, req.position, max_offset)
2545 .map_err(crate::error::AllSourceError::InvalidInput)?;
2546
2547 Ok(Json(serde_json::json!({
2548 "status": "ok",
2549 "consumer_id": consumer_id,
2550 "position": req.position,
2551 })))
2552}
2553
2554#[cfg(test)]
2555mod tests {
2556 use super::*;
2557 use crate::{domain::entities::Event, store::EventStore};
2558
2559 fn create_test_store() -> Arc<EventStore> {
2560 Arc::new(EventStore::new())
2561 }
2562
2563 async fn query_page(store: &SharedStore, query: &str) -> QueryEventsResponse {
2572 use axum::extract::{Query, State};
2573
2574 let uri: axum::http::Uri = format!("/api/v1/events/query?tenant_id=test-stream&{query}")
2575 .parse()
2576 .unwrap();
2577 query_events(
2578 OptionalAuth(None),
2579 Query::try_from_uri(&uri).unwrap(),
2580 Query::try_from_uri(&uri).unwrap(),
2581 Query::try_from_uri(&uri).unwrap(),
2582 State(store.clone()),
2583 )
2584 .await
2585 .unwrap()
2586 .0
2587 }
2588
2589 fn create_test_event(entity_id: &str, event_type: &str) -> Event {
2590 Event::from_strings(
2591 event_type.to_string(),
2592 entity_id.to_string(),
2593 "test-stream".to_string(),
2594 serde_json::json!({
2595 "name": "Test",
2596 "value": 42
2597 }),
2598 None,
2599 )
2600 .unwrap()
2601 }
2602
2603 #[tokio::test]
2604 async fn test_query_events_has_more_and_total_count() {
2605 let store = create_test_store();
2606
2607 for i in 0..50 {
2609 store
2610 .ingest(&create_test_event(&format!("entity-{i}"), "user.created"))
2611 .unwrap();
2612 }
2613
2614 let req = QueryEventsRequest {
2616 entity_id: None,
2617 event_type: None,
2618 tenant_id: None,
2619 as_of: None,
2620 since: None,
2621 until: None,
2622 limit: Some(10),
2623 event_type_prefix: None,
2624 exclude_event_type_prefix: None,
2625 payload_filter: None,
2626 };
2627
2628 let requested_limit = req.limit;
2629 let unlimited_req = QueryEventsRequest {
2630 limit: None,
2631 ..QueryEventsRequest {
2632 entity_id: req.entity_id,
2633 event_type: req.event_type,
2634 tenant_id: req.tenant_id,
2635 as_of: req.as_of,
2636 since: req.since,
2637 until: req.until,
2638 limit: None,
2639 event_type_prefix: req.event_type_prefix,
2640 exclude_event_type_prefix: None,
2641 payload_filter: req.payload_filter,
2642 }
2643 };
2644 let all_events = store.query(&unlimited_req).unwrap();
2645 let total_count = all_events.len();
2646 let limited_events: Vec<Event> = if let Some(limit) = requested_limit {
2647 all_events.into_iter().take(limit).collect()
2648 } else {
2649 all_events
2650 };
2651 let count = limited_events.len();
2652 let has_more = count < total_count;
2653
2654 assert_eq!(count, 10);
2655 assert_eq!(total_count, 50);
2656 assert!(has_more);
2657 }
2658
2659 #[tokio::test]
2668 async fn query_events_honours_offset_pagination() {
2669 use axum::extract::{Query, State};
2670
2671 let store = create_test_store();
2672 for i in 0..25 {
2673 store
2674 .ingest(&create_test_event(&format!("e-{i:02}"), "user.created"))
2675 .unwrap();
2676 }
2677
2678 async fn page(store: &SharedStore, limit: usize, offset: usize) -> QueryEventsResponse {
2680 let uri: axum::http::Uri =
2681 format!("/api/v1/events/query?tenant_id=test-stream&limit={limit}&offset={offset}")
2682 .parse()
2683 .unwrap();
2684 let req: Query<QueryEventsRequest> = Query::try_from_uri(&uri).unwrap();
2685 let order: Query<EventOrderParam> = Query::try_from_uri(&uri).unwrap();
2686 let off: Query<EventOffsetParam> = Query::try_from_uri(&uri).unwrap();
2687 assert_eq!(off.0.offset, Some(offset), "offset must deserialize");
2688 query_events(OptionalAuth(None), req, order, off, State(store.clone()))
2689 .await
2690 .unwrap()
2691 .0
2692 }
2693
2694 let p1 = page(&store, 10, 0).await;
2695 let p2 = page(&store, 10, 10).await;
2696 let p3 = page(&store, 10, 20).await;
2697
2698 assert_eq!(p1.count, 10);
2699 assert_eq!(p2.count, 10);
2700 assert_eq!(p3.count, 5, "last page returns the remainder");
2701 assert_eq!(p1.total_count, 25);
2702
2703 let ids = |r: &QueryEventsResponse| -> Vec<String> {
2705 r.events.iter().map(|e| e.entity_id.clone()).collect()
2706 };
2707 assert_ne!(ids(&p1), ids(&p2), "offset=10 must skip the first page");
2708
2709 let mut all = ids(&p1);
2710 all.extend(ids(&p2));
2711 all.extend(ids(&p3));
2712 let unique: std::collections::HashSet<_> = all.iter().cloned().collect();
2713 assert_eq!(
2714 unique.len(),
2715 25,
2716 "paging the whole set must yield 25 distinct entities, not duplicates"
2717 );
2718
2719 assert!(p1.has_more, "25 events, page 1 of 10 → more remain");
2722 assert!(p2.has_more, "25 events, page 2 of 10 → more remain");
2723 assert!(!p3.has_more, "offset=20 + count=5 == total → exhausted");
2724
2725 let past = page(&store, 10, 100).await;
2727 assert_eq!(past.count, 0);
2728 assert!(!past.has_more, "offset beyond the match set is exhausted");
2729 }
2730
2731 #[tokio::test]
2740 async fn projection_state_summary_honours_limit_offset_and_prefix() {
2741 use axum::{
2742 body::{Body, to_bytes},
2743 http::Request,
2744 };
2745 use tower::ServiceExt; let store = create_test_store();
2748 let cache = store.projection_state_cache();
2749 for i in 0..25 {
2750 cache.insert(
2751 format!("demo:tenant-{i:02}"),
2752 serde_json::json!({ "n": i as u64 }),
2753 );
2754 }
2755 cache.insert(
2757 "other:tenant-99".to_string(),
2758 serde_json::json!({ "n": 99 }),
2759 );
2760
2761 let app = Router::new()
2762 .route(
2763 "/api/v1/projections/{name}/state",
2764 get(get_projection_state_summary),
2765 )
2766 .with_state(store.clone());
2767
2768 async fn fetch(app: &Router, uri: &str) -> serde_json::Value {
2769 let resp = app
2770 .clone()
2771 .oneshot(Request::builder().uri(uri).body(Body::empty()).unwrap())
2772 .await
2773 .unwrap();
2774 assert_eq!(resp.status(), axum::http::StatusCode::OK, "GET {uri}");
2775 let bytes = to_bytes(resp.into_body(), usize::MAX).await.unwrap();
2776 serde_json::from_slice(&bytes).unwrap()
2777 }
2778
2779 let ids = |body: &serde_json::Value| -> Vec<String> {
2780 body["states"]
2781 .as_array()
2782 .unwrap()
2783 .iter()
2784 .map(|s| s["entity_id"].as_str().unwrap().to_string())
2785 .collect()
2786 };
2787
2788 let all = fetch(&app, "/api/v1/projections/demo/state").await;
2790 assert_eq!(all["total"], 25);
2791 assert_eq!(ids(&all).len(), 25);
2792
2793 let p1 = fetch(&app, "/api/v1/projections/demo/state?limit=10").await;
2795 assert_eq!(ids(&p1).len(), 10, "limit must bound the response");
2796 assert_eq!(p1["total"], 25);
2797 assert_eq!(p1["count"], 10);
2798 assert_eq!(p1["has_more"], true);
2799
2800 let p2 = fetch(&app, "/api/v1/projections/demo/state?limit=10&offset=10").await;
2802 let p3 = fetch(&app, "/api/v1/projections/demo/state?limit=10&offset=20").await;
2803 assert_eq!(ids(&p3).len(), 5, "last page returns the remainder");
2804 assert_eq!(p3["has_more"], false, "offset + count == total → exhausted");
2805 assert_ne!(ids(&p1), ids(&p2), "offset=10 must skip the first page");
2806
2807 let mut walked = ids(&p1);
2808 walked.extend(ids(&p2));
2809 walked.extend(ids(&p3));
2810 let unique: std::collections::HashSet<_> = walked.iter().cloned().collect();
2811 assert_eq!(
2812 unique.len(),
2813 25,
2814 "paging the whole projection must yield 25 distinct entities"
2815 );
2816
2817 let mut sorted = walked.clone();
2819 sorted.sort();
2820 assert_eq!(walked, sorted, "pages must be ordered by entity_id");
2821
2822 let past = fetch(&app, "/api/v1/projections/demo/state?limit=10&offset=100").await;
2824 assert_eq!(past["count"], 0);
2825 assert_eq!(past["has_more"], false);
2826 assert_eq!(past["total"], 25);
2827
2828 let shard = fetch(
2830 &app,
2831 "/api/v1/projections/demo/state?entity_id_prefix=tenant-1",
2832 )
2833 .await;
2834 assert_eq!(shard["total"], 10, "tenant-10..tenant-19");
2835 assert!(
2836 ids(&shard).iter().all(|id| id.starts_with("tenant-1")),
2837 "entity_id_prefix must filter"
2838 );
2839
2840 let other = fetch(&app, "/api/v1/projections/other/state?limit=10").await;
2842 assert_eq!(ids(&other), vec!["tenant-99".to_string()]);
2843 }
2844
2845 #[test]
2858 fn query_events_limit_does_not_materialize_whole_history() {
2859 const HISTORY: usize = 200;
2860
2861 let store = create_test_store();
2862 for _ in 0..HISTORY {
2863 store
2864 .ingest(&create_test_event("entity-hot", "user.updated"))
2865 .unwrap();
2866 }
2867
2868 let runtime = tokio::runtime::Builder::new_current_thread()
2872 .enable_all()
2873 .build()
2874 .unwrap();
2875 let (resp, materialized) = crate::clone_probe::measure(|| {
2876 runtime.block_on(query_page(
2877 &store,
2878 "entity_id=entity-hot&limit=1&order=desc",
2879 ))
2880 });
2881
2882 assert_eq!(resp.count, 1);
2885 assert_eq!(resp.events.len(), 1);
2886 assert_eq!(resp.total_count, HISTORY);
2887 assert!(resp.has_more);
2888
2889 assert_eq!(
2890 materialized, 1,
2891 "limit=1 materialized {materialized} events out of {HISTORY}: \
2892 `limit` must bound what a request materializes, not just what it \
2893 returns"
2894 );
2895 }
2896
2897 async fn query_page_result(
2900 store: &SharedStore,
2901 query: &str,
2902 ) -> Result<Json<QueryEventsResponse>> {
2903 use axum::extract::{Query, State};
2904
2905 let uri: axum::http::Uri = format!("/api/v1/events/query?tenant_id=test-stream&{query}")
2906 .parse()
2907 .unwrap();
2908 query_events(
2909 OptionalAuth(None),
2910 Query::try_from_uri(&uri).unwrap(),
2911 Query::try_from_uri(&uri).unwrap(),
2912 Query::try_from_uri(&uri).unwrap(),
2913 State(store.clone()),
2914 )
2915 .await
2916 }
2917
2918 #[tokio::test]
2927 async fn query_events_rejects_a_payload_filter_it_cannot_apply() {
2928 let store = create_test_store();
2929 for name in ["alice", "bob", "carol"] {
2930 let mut event = create_test_event(name, "user.created");
2931 event.payload = serde_json::json!({ "user_id": name });
2932 store.ingest(&event).unwrap();
2933 }
2934
2935 let ok = query_page(&store, "payload_filter=%7B%22user_id%22%3A%22alice%22%7D").await;
2938 assert_eq!(ok.count, 1, "a valid payload_filter must still work");
2939 assert_eq!(ok.total_count, 1);
2940
2941 for bad in [
2942 "not-json", "%7B%22user_id%22%3A%22alice", "%5B%22alice%22%5D", "42", ] {
2947 let result = query_page_result(&store, &format!("payload_filter={bad}")).await;
2948 let Err(err) = result else {
2949 let resp = result.unwrap().0;
2950 panic!(
2951 "payload_filter={bad} was silently ignored: returned {} of {} \
2952 events unfiltered instead of rejecting a filter the server \
2953 cannot apply",
2954 resp.count, resp.total_count
2955 );
2956 };
2957 assert!(
2958 matches!(err, crate::error::AllSourceError::InvalidInput(_)),
2959 "payload_filter={bad} must be a 400, got {err:?}"
2960 );
2961 }
2962 }
2963
2964 #[tokio::test]
2970 async fn query_events_rejects_an_unusable_order_value() {
2971 let store = create_test_store();
2972 store
2973 .ingest(&create_test_event("e-1", "user.created"))
2974 .unwrap();
2975
2976 for bad in ["descending", "DESCENDING", "newest", "1", "asc%20"] {
2977 let result = query_page_result(&store, &format!("order={bad}")).await;
2978 let Err(err) = result else {
2979 panic!("order={bad} must be rejected, not silently defaulted");
2980 };
2981 assert!(
2982 matches!(err, crate::error::AllSourceError::InvalidInput(_)),
2983 "order={bad} must be a 400, got {err:?}"
2984 );
2985 }
2986
2987 for good in ["asc", "ASC", "desc", "DeSc"] {
2989 let accepted = query_page_result(&store, &format!("order={good}"))
2990 .await
2991 .unwrap_or_else(|e| panic!("order={good} must be accepted: {e:?}"));
2992 assert_eq!(accepted.0.count, 1, "order={good}");
2993 }
2994 }
2995
2996 #[tokio::test]
3005 async fn query_events_exclude_prefix_applies_before_the_window() {
3006 const TYPES: [&str; 4] = [
3007 "audit.write",
3008 "user.created",
3009 "service.ping",
3010 "user.updated",
3011 ];
3012 let store = create_test_store();
3013 let base = chrono::Utc::now() - chrono::Duration::hours(24);
3014 let mut kept = Vec::new();
3015 for i in 0..12i64 {
3016 let mut event = create_test_event("org-1", TYPES[i as usize % 4]);
3017 event.timestamp = base + chrono::Duration::minutes(i);
3018 event.version = i + 1;
3019 if i % 2 == 1 {
3020 kept.push(event.id);
3021 }
3022 store.ingest(&event).unwrap();
3023 }
3024 assert_eq!(kept.len(), 6, "half the stream is user.*");
3025
3026 let page = query_page(
3030 &store,
3031 "exclude_event_type_prefix=audit.,%20service.&limit=3",
3032 )
3033 .await;
3034 assert_eq!(
3035 page.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3036 kept[..3].to_vec(),
3037 "limit=3 must yield 3 non-excluded events: exclusion runs before \
3038 the window, and the prefix list is comma-separated (whitespace \
3039 trimmed)"
3040 );
3041 assert!(
3042 page.events
3043 .iter()
3044 .all(|e| !e.event_type.starts_with("audit.")
3045 && !e.event_type.starts_with("service.")),
3046 "excluded namespaces must not appear: {:?}",
3047 page.events
3048 .iter()
3049 .map(|e| &e.event_type)
3050 .collect::<Vec<_>>()
3051 );
3052 assert_eq!(
3053 page.total_count, 6,
3054 "total_count is the post-exclusion match set, not the 12 ingested"
3055 );
3056 assert!(page.has_more);
3057
3058 let mut walked = Vec::new();
3060 for offset in [0, 3, 6] {
3061 let p = query_page(
3062 &store,
3063 &format!("exclude_event_type_prefix=audit.,service.&limit=3&offset={offset}"),
3064 )
3065 .await;
3066 assert_eq!(
3067 p.has_more,
3068 offset + p.count < 6,
3069 "has_more must terminate on the excluded view (offset={offset})"
3070 );
3071 walked.extend(p.events.iter().map(|e| e.id));
3072 }
3073 assert_eq!(
3074 walked, kept,
3075 "exclusion + paging must cover the survivors once"
3076 );
3077
3078 let desc = query_page(
3080 &store,
3081 "exclude_event_type_prefix=audit.,service.&limit=1&order=desc",
3082 )
3083 .await;
3084 assert_eq!(
3085 desc.events[0].id,
3086 *kept.last().unwrap(),
3087 "order=desc&limit=1 over an excluded view is the newest SURVIVOR"
3088 );
3089
3090 let audit_only = query_page(&store, "exclude_event_type_prefix=audit.").await;
3093 assert_eq!(audit_only.total_count, 9, "12 minus the 3 audit.* events");
3094 let nothing = query_page(&store, "exclude_event_type_prefix=nosuch.").await;
3095 assert_eq!(nothing.total_count, 12);
3096 }
3097
3098 #[tokio::test]
3112 async fn query_events_honours_time_window_without_an_entity_or_type_filter() {
3113 use chrono::{SecondsFormat, SubsecRound};
3114
3115 let store = create_test_store();
3116 let base = chrono::DateTime::parse_from_rfc3339("2026-01-01T00:00:00.123456789Z")
3130 .unwrap()
3131 .to_utc()
3132 .trunc_subsecs(0);
3133 let mut ids = Vec::new();
3134 for i in 0..5i64 {
3135 let mut event = create_test_event(&format!("e-{i}"), "user.created");
3136 event.timestamp = base + chrono::Duration::hours(i);
3137 event.version = i + 1;
3138 ids.push(event.id);
3139 store.ingest(&event).unwrap();
3140 }
3141 let at = |h: i64| {
3144 let t = base + chrono::Duration::hours(h);
3145 let s = t.to_rfc3339_opts(SecondsFormat::Micros, true);
3146 assert_eq!(
3151 chrono::DateTime::parse_from_rfc3339(&s).unwrap().to_utc(),
3152 t,
3153 "query-string timestamp must round-trip exactly, else the window \
3154 boundary silently excludes the event stamped at it"
3155 );
3156 s
3157 };
3158
3159 for (qs, expected) in [
3160 (format!("since={}", at(2)), vec![ids[2], ids[3], ids[4]]),
3161 (format!("until={}", at(1)), vec![ids[0], ids[1]]),
3162 (format!("as_of={}", at(1)), vec![ids[0], ids[1]]),
3163 (
3164 format!("since={}&until={}", at(1), at(3)),
3165 vec![ids[1], ids[2], ids[3]],
3166 ),
3167 ] {
3168 let resp = query_page(&store, &qs).await;
3169 let got: Vec<_> = resp.events.iter().map(|e| e.id).collect();
3170 assert_eq!(got, expected, "?{qs} must return only the window");
3171 assert_eq!(resp.count, expected.len(), "?{qs}");
3172 assert_eq!(
3173 resp.total_count,
3174 expected.len(),
3175 "?{qs}: total_count must count the window, not the history"
3176 );
3177 assert!(!resp.has_more, "?{qs}: the whole window was served");
3178 }
3179
3180 let page1 = query_page(&store, &format!("since={}&limit=2", at(2))).await;
3183 let page2 = query_page(&store, &format!("since={}&limit=2&offset=2", at(2))).await;
3184 assert_eq!(
3185 page1.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3186 vec![ids[2], ids[3]]
3187 );
3188 assert!(page1.has_more, "3 in the window, 2 served");
3189 assert_eq!(
3190 page2.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3191 vec![ids[4]]
3192 );
3193 assert!(!page2.has_more, "offset 2 + count 1 == the window's 3");
3194 assert_eq!(page2.total_count, 3);
3195
3196 let desc = query_page(&store, &format!("since={}&order=desc", at(2))).await;
3198 assert_eq!(
3199 desc.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3200 vec![ids[4], ids[3], ids[2]]
3201 );
3202
3203 let empty = query_page(&store, &format!("since={}", at(99))).await;
3205 assert_eq!(empty.count, 0);
3206 assert_eq!(empty.total_count, 0);
3207 assert!(!empty.has_more);
3208 }
3209
3210 #[tokio::test]
3214 async fn query_events_fails_closed_without_tenant() {
3215 use axum::extract::{Query, State};
3216
3217 let store = create_test_store();
3218 for i in 0..5 {
3219 store
3220 .ingest(&create_test_event(&format!("e-{i}"), "user.created"))
3221 .unwrap();
3222 }
3223
3224 let resp = query_events(
3225 OptionalAuth(None),
3226 Query(QueryEventsRequest::default()),
3227 Query(EventOrderParam { order: None }),
3228 Query(EventOffsetParam { offset: None }),
3229 State(store.clone()),
3230 )
3231 .await
3232 .unwrap();
3233 assert_eq!(
3234 resp.0.total_count, 0,
3235 "a no-tenant query must NOT return cross-tenant events"
3236 );
3237 assert_eq!(resp.0.count, 0);
3238
3239 let scoped = query_events(
3242 OptionalAuth(None),
3243 Query(QueryEventsRequest {
3244 tenant_id: Some("test-stream".to_string()),
3245 ..QueryEventsRequest::default()
3246 }),
3247 Query(EventOrderParam { order: None }),
3248 Query(EventOffsetParam { offset: None }),
3249 State(store),
3250 )
3251 .await
3252 .unwrap();
3253 assert_eq!(
3254 scoped.0.total_count, 5,
3255 "tenant-scoped query returns its events"
3256 );
3257 }
3258
3259 #[tokio::test]
3264 async fn list_streams_and_types_are_tenant_scoped() {
3265 use crate::domain::entities::Event;
3266 use axum::extract::{Query, State};
3267
3268 let store = create_test_store();
3269 let ev = |entity: &str, etype: &str, tenant: &str| {
3270 Event::from_strings(
3271 etype.to_string(),
3272 entity.to_string(),
3273 tenant.to_string(),
3274 serde_json::json!({}),
3275 None,
3276 )
3277 .unwrap()
3278 };
3279 store.ingest(&ev("e1", "order.placed", "tenant-a")).unwrap();
3281 store.ingest(&ev("e2", "user.created", "tenant-a")).unwrap();
3282 store
3283 .ingest(&ev("e9", "thing.happened", "tenant-b"))
3284 .unwrap();
3285
3286 let streams = |tid: Option<&str>| {
3287 list_streams(
3288 OptionalAuth(None),
3289 State(store.clone()),
3290 Query(ListStreamsParams {
3291 tenant_id: tid.map(String::from),
3292 limit: None,
3293 offset: None,
3294 }),
3295 )
3296 };
3297 assert_eq!(
3298 streams(Some("tenant-a")).await.0.total,
3299 2,
3300 "tenant-a streams"
3301 );
3302 assert_eq!(
3303 streams(Some("tenant-b")).await.0.total,
3304 1,
3305 "tenant-b streams"
3306 );
3307 assert_eq!(
3308 streams(None).await.0.total,
3309 0,
3310 "no tenant -> no streams (fail closed)"
3311 );
3312
3313 let types = |tid: Option<&str>| {
3314 list_event_types(
3315 OptionalAuth(None),
3316 State(store.clone()),
3317 Query(ListEventTypesParams {
3318 tenant_id: tid.map(String::from),
3319 limit: None,
3320 offset: None,
3321 }),
3322 )
3323 };
3324 assert_eq!(
3325 types(Some("tenant-a")).await.0.total,
3326 2,
3327 "tenant-a event types"
3328 );
3329 assert_eq!(
3330 types(Some("tenant-b")).await.0.total,
3331 1,
3332 "tenant-b event types"
3333 );
3334 assert_eq!(
3335 types(None).await.0.total,
3336 0,
3337 "no tenant -> no types (fail closed)"
3338 );
3339 }
3340
3341 #[tokio::test]
3342 async fn test_query_events_no_more_results() {
3343 let store = create_test_store();
3344
3345 for i in 0..5 {
3347 store
3348 .ingest(&create_test_event(&format!("entity-{i}"), "user.created"))
3349 .unwrap();
3350 }
3351
3352 let all_events = store
3354 .query(&QueryEventsRequest {
3355 entity_id: None,
3356 event_type: None,
3357 tenant_id: None,
3358 as_of: None,
3359 since: None,
3360 until: None,
3361 limit: None,
3362 event_type_prefix: None,
3363 exclude_event_type_prefix: None,
3364 payload_filter: None,
3365 })
3366 .unwrap();
3367 let total_count = all_events.len();
3368 let limited_events: Vec<Event> = all_events.into_iter().take(100).collect();
3369 let count = limited_events.len();
3370 let has_more = count < total_count;
3371
3372 assert_eq!(count, 5);
3373 assert_eq!(total_count, 5);
3374 assert!(!has_more);
3375 }
3376
3377 #[tokio::test]
3381 async fn test_query_events_order_desc_returns_latest() {
3382 let store = create_test_store();
3388
3389 let base = chrono::Utc::now();
3392 let mut ascending_ids = Vec::new();
3393 for i in 0..5i64 {
3394 let mut event = create_test_event("org-1", "auth.org.updated");
3395 event.timestamp = base + chrono::Duration::seconds(i);
3396 event.version = i + 1;
3397 ascending_ids.push(event.id);
3398 store.ingest(&event).unwrap();
3399 }
3400 let newest_ts = base + chrono::Duration::seconds(4);
3401
3402 let latest = query_page(&store, "entity_id=org-1&limit=1&order=desc").await;
3404 assert_eq!(latest.count, 1);
3405 assert_eq!(
3406 latest.events[0].id,
3407 ascending_ids[4],
3408 "order=desc&limit=1 must yield the NEWEST event, got the one at \
3409 ascending position {:?}",
3410 ascending_ids
3411 .iter()
3412 .position(|id| *id == latest.events[0].id)
3413 );
3414 assert_eq!(latest.events[0].timestamp, newest_ts);
3415 assert_eq!(latest.total_count, 5, "total is the full match set");
3416 assert!(latest.has_more);
3417
3418 for qs in [
3420 "entity_id=org-1&limit=1",
3421 "entity_id=org-1&limit=1&order=asc",
3422 ] {
3423 let oldest = query_page(&store, qs).await;
3424 assert_eq!(oldest.events[0].id, ascending_ids[0], "{qs}");
3425 assert_eq!(oldest.events[0].timestamp, base);
3426 }
3427
3428 let all_desc = query_page(&store, "entity_id=org-1&order=desc").await;
3431 let got: Vec<_> = all_desc.events.iter().map(|e| e.id).collect();
3432 let expected: Vec<_> = ascending_ids.iter().rev().copied().collect();
3433 assert_eq!(got, expected, "order=desc must return newest-first");
3434 }
3435
3436 #[tokio::test]
3443 async fn query_events_desc_composes_with_offset_and_limit() {
3444 let store = create_test_store();
3445 let base = chrono::Utc::now();
3446 let mut ascending_ids = Vec::new();
3447 for i in 0..5i64 {
3448 let mut event = create_test_event("org-1", "auth.org.updated");
3449 event.timestamp = base + chrono::Duration::seconds(i);
3450 event.version = i + 1;
3451 ascending_ids.push(event.id);
3452 store.ingest(&event).unwrap();
3453 }
3454 let newest_first: Vec<_> = ascending_ids.iter().rev().copied().collect();
3455
3456 for (offset, limit) in [(0, 2), (1, 2), (2, 2), (3, 2), (4, 2), (5, 2), (1, 4)] {
3457 let page = query_page(
3458 &store,
3459 &format!("entity_id=org-1&order=desc&offset={offset}&limit={limit}"),
3460 )
3461 .await;
3462 let got: Vec<_> = page.events.iter().map(|e| e.id).collect();
3463 let expected: Vec<_> = newest_first
3464 .iter()
3465 .skip(offset)
3466 .take(limit)
3467 .copied()
3468 .collect();
3469 assert_eq!(
3470 got, expected,
3471 "order=desc&offset={offset}&limit={limit} must reverse, then \
3472 skip, then take"
3473 );
3474 assert_eq!(page.count, expected.len());
3475 assert_eq!(page.total_count, 5);
3476 assert_eq!(
3477 page.has_more,
3478 offset + expected.len() < 5,
3479 "has_more must account for the offset (offset={offset})"
3480 );
3481 }
3482
3483 let mut walked = Vec::new();
3486 for offset in (0..5).step_by(2) {
3487 let page = query_page(
3488 &store,
3489 &format!("entity_id=org-1&order=desc&offset={offset}&limit=2"),
3490 )
3491 .await;
3492 walked.extend(page.events.iter().map(|e| e.id));
3493 }
3494 assert_eq!(walked, newest_first, "desc paging must cover the set once");
3495 }
3496
3497 #[tokio::test]
3498 async fn test_list_entities_by_type_prefix() {
3499 let store = create_test_store();
3500
3501 store
3503 .ingest(&create_test_event("idx-1", "index.created"))
3504 .unwrap();
3505 store
3506 .ingest(&create_test_event("idx-1", "index.updated"))
3507 .unwrap();
3508 store
3509 .ingest(&create_test_event("idx-2", "index.created"))
3510 .unwrap();
3511 store
3512 .ingest(&create_test_event("idx-3", "index.created"))
3513 .unwrap();
3514 store
3516 .ingest(&create_test_event("trade-1", "trade.created"))
3517 .unwrap();
3518 store
3519 .ingest(&create_test_event("trade-2", "trade.created"))
3520 .unwrap();
3521
3522 let req = ListEntitiesRequest {
3524 event_type_prefix: Some("index.".to_string()),
3525 ..Default::default()
3526 };
3527 let query_req = QueryEventsRequest {
3528 entity_id: None,
3529 event_type: None,
3530 tenant_id: None,
3531 as_of: None,
3532 since: None,
3533 until: None,
3534 limit: None,
3535 event_type_prefix: req.event_type_prefix,
3536 exclude_event_type_prefix: None,
3537 payload_filter: req.payload_filter,
3538 };
3539 let events = store.query(&query_req).unwrap();
3540
3541 let mut entity_map: std::collections::HashMap<String, Vec<&Event>> =
3543 std::collections::HashMap::new();
3544 for event in &events {
3545 entity_map
3546 .entry(event.entity_id().to_string())
3547 .or_default()
3548 .push(event);
3549 }
3550
3551 assert_eq!(entity_map.len(), 3); assert_eq!(entity_map["idx-1"].len(), 2); assert_eq!(entity_map["idx-2"].len(), 1);
3554 assert_eq!(entity_map["idx-3"].len(), 1);
3555 }
3556
3557 #[tokio::test]
3560 async fn test_list_entities_order_and_pagination() {
3561 let store = create_test_store();
3562
3563 let base = chrono::Utc::now();
3565 for (i, eid) in ["org-a", "org-b", "org-c"].iter().enumerate() {
3566 let mut event = create_test_event(eid, "auth.org.created");
3567 event.timestamp = base + chrono::Duration::seconds(i as i64);
3568 store.ingest(&event).unwrap();
3569 }
3570 let prefix = || Some("auth.org.".to_string());
3571
3572 let desc = list_entities(
3574 State(store.clone()),
3575 Query(ListEntitiesRequest {
3576 event_type_prefix: prefix(),
3577 ..Default::default()
3578 }),
3579 )
3580 .await
3581 .unwrap();
3582 let desc_ids: Vec<&str> = desc
3583 .0
3584 .entities
3585 .iter()
3586 .map(|e| e.entity_id.as_str())
3587 .collect();
3588 assert_eq!(desc_ids, ["org-c", "org-b", "org-a"]);
3589
3590 let asc = list_entities(
3592 State(store.clone()),
3593 Query(ListEntitiesRequest {
3594 event_type_prefix: prefix(),
3595 order: Some("asc".to_string()),
3596 ..Default::default()
3597 }),
3598 )
3599 .await
3600 .unwrap();
3601 let asc_ids: Vec<&str> = asc
3602 .0
3603 .entities
3604 .iter()
3605 .map(|e| e.entity_id.as_str())
3606 .collect();
3607 assert_eq!(asc_ids, ["org-a", "org-b", "org-c"]);
3608
3609 let page2 = list_entities(
3611 State(store.clone()),
3612 Query(ListEntitiesRequest {
3613 event_type_prefix: prefix(),
3614 order: Some("asc".to_string()),
3615 limit: Some(1),
3616 offset: Some(1),
3617 ..Default::default()
3618 }),
3619 )
3620 .await
3621 .unwrap();
3622 assert_eq!(page2.0.entities.len(), 1);
3623 assert_eq!(page2.0.entities[0].entity_id, "org-b");
3624 assert_eq!(page2.0.total, 3);
3625 assert!(page2.0.has_more);
3626
3627 let err = list_entities(
3629 State(store.clone()),
3630 Query(ListEntitiesRequest {
3631 event_type_prefix: prefix(),
3632 order: Some("sideways".to_string()),
3633 ..Default::default()
3634 }),
3635 )
3636 .await;
3637 assert!(err.is_err(), "invalid order value must be rejected");
3638 }
3639
3640 fn create_test_event_with_payload(
3641 entity_id: &str,
3642 event_type: &str,
3643 payload: serde_json::Value,
3644 ) -> Event {
3645 Event::from_strings(
3646 event_type.to_string(),
3647 entity_id.to_string(),
3648 "test-stream".to_string(),
3649 payload,
3650 None,
3651 )
3652 .unwrap()
3653 }
3654
3655 #[tokio::test]
3656 async fn test_detect_duplicates_by_payload_fields() {
3657 let store = create_test_store();
3658
3659 store
3661 .ingest(&create_test_event_with_payload(
3662 "idx-1",
3663 "index.created",
3664 serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3665 ))
3666 .unwrap();
3667 store
3668 .ingest(&create_test_event_with_payload(
3669 "idx-2",
3670 "index.created",
3671 serde_json::json!({"name": "S&P 500", "user_id": "bob"}),
3672 ))
3673 .unwrap();
3674 store
3675 .ingest(&create_test_event_with_payload(
3676 "idx-3",
3677 "index.created",
3678 serde_json::json!({"name": "NASDAQ", "user_id": "alice"}),
3679 ))
3680 .unwrap();
3681 store
3682 .ingest(&create_test_event_with_payload(
3683 "idx-4",
3684 "index.created",
3685 serde_json::json!({"name": "NASDAQ", "user_id": "carol"}),
3686 ))
3687 .unwrap();
3688 store
3689 .ingest(&create_test_event_with_payload(
3690 "idx-5",
3691 "index.created",
3692 serde_json::json!({"name": "DAX", "user_id": "dave"}),
3693 ))
3694 .unwrap();
3695
3696 let query_req = QueryEventsRequest {
3698 entity_id: None,
3699 event_type: None,
3700 tenant_id: None,
3701 as_of: None,
3702 since: None,
3703 until: None,
3704 limit: None,
3705 event_type_prefix: Some("index.".to_string()),
3706 exclude_event_type_prefix: None,
3707 payload_filter: None,
3708 };
3709 let events = store.query(&query_req).unwrap();
3710
3711 let group_by_fields = vec!["name"];
3713 let mut entity_latest: std::collections::HashMap<String, &Event> =
3714 std::collections::HashMap::new();
3715 for event in &events {
3716 let eid = event.entity_id().to_string();
3717 entity_latest
3718 .entry(eid)
3719 .and_modify(|existing| {
3720 if event.timestamp() > existing.timestamp() {
3721 *existing = event;
3722 }
3723 })
3724 .or_insert(event);
3725 }
3726
3727 let mut groups: std::collections::HashMap<String, Vec<String>> =
3728 std::collections::HashMap::new();
3729 for (entity_id, event) in &entity_latest {
3730 let payload = event.payload();
3731 let mut key_parts = serde_json::Map::new();
3732 for field in &group_by_fields {
3733 let value = payload
3734 .get(*field)
3735 .cloned()
3736 .unwrap_or(serde_json::Value::Null);
3737 key_parts.insert((*field).to_string(), value);
3738 }
3739 let key_str = serde_json::to_string(&key_parts).unwrap_or_default();
3740 groups.entry(key_str).or_default().push(entity_id.clone());
3741 }
3742
3743 let duplicate_groups: Vec<_> = groups
3744 .into_iter()
3745 .filter(|(_, ids)| ids.len() > 1)
3746 .collect();
3747
3748 assert_eq!(duplicate_groups.len(), 2); for (_, ids) in &duplicate_groups {
3750 assert_eq!(ids.len(), 2);
3751 }
3752 }
3753
3754 #[tokio::test]
3755 async fn test_detect_duplicates_no_duplicates() {
3756 let store = create_test_store();
3757
3758 store
3760 .ingest(&create_test_event_with_payload(
3761 "idx-1",
3762 "index.created",
3763 serde_json::json!({"name": "A"}),
3764 ))
3765 .unwrap();
3766 store
3767 .ingest(&create_test_event_with_payload(
3768 "idx-2",
3769 "index.created",
3770 serde_json::json!({"name": "B"}),
3771 ))
3772 .unwrap();
3773
3774 let query_req = QueryEventsRequest {
3775 entity_id: None,
3776 event_type: None,
3777 tenant_id: None,
3778 as_of: None,
3779 since: None,
3780 until: None,
3781 limit: None,
3782 event_type_prefix: Some("index.".to_string()),
3783 exclude_event_type_prefix: None,
3784 payload_filter: None,
3785 };
3786 let events = store.query(&query_req).unwrap();
3787
3788 let mut entity_latest: std::collections::HashMap<String, &Event> =
3789 std::collections::HashMap::new();
3790 for event in &events {
3791 entity_latest
3792 .entry(event.entity_id().to_string())
3793 .or_insert(event);
3794 }
3795
3796 let mut groups: std::collections::HashMap<String, Vec<String>> =
3797 std::collections::HashMap::new();
3798 for (entity_id, event) in &entity_latest {
3799 let key_str =
3800 serde_json::to_string(&serde_json::json!({"name": event.payload().get("name")}))
3801 .unwrap();
3802 groups.entry(key_str).or_default().push(entity_id.clone());
3803 }
3804
3805 let duplicate_groups: Vec<_> = groups
3806 .into_iter()
3807 .filter(|(_, ids)| ids.len() > 1)
3808 .collect();
3809
3810 assert_eq!(duplicate_groups.len(), 0); }
3812
3813 #[tokio::test]
3814 async fn test_detect_duplicates_multi_field_group_by() {
3815 let store = create_test_store();
3816
3817 store
3819 .ingest(&create_test_event_with_payload(
3820 "idx-1",
3821 "index.created",
3822 serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3823 ))
3824 .unwrap();
3825 store
3826 .ingest(&create_test_event_with_payload(
3827 "idx-2",
3828 "index.created",
3829 serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3830 ))
3831 .unwrap();
3832 store
3834 .ingest(&create_test_event_with_payload(
3835 "idx-3",
3836 "index.created",
3837 serde_json::json!({"name": "S&P 500", "user_id": "bob"}),
3838 ))
3839 .unwrap();
3840
3841 let query_req = QueryEventsRequest {
3842 entity_id: None,
3843 event_type: None,
3844 tenant_id: None,
3845 as_of: None,
3846 since: None,
3847 until: None,
3848 limit: None,
3849 event_type_prefix: Some("index.".to_string()),
3850 exclude_event_type_prefix: None,
3851 payload_filter: None,
3852 };
3853 let events = store.query(&query_req).unwrap();
3854
3855 let group_by_fields = vec!["name", "user_id"];
3856 let mut entity_latest: std::collections::HashMap<String, &Event> =
3857 std::collections::HashMap::new();
3858 for event in &events {
3859 entity_latest
3860 .entry(event.entity_id().to_string())
3861 .and_modify(|existing| {
3862 if event.timestamp() > existing.timestamp() {
3863 *existing = event;
3864 }
3865 })
3866 .or_insert(event);
3867 }
3868
3869 let mut groups: std::collections::HashMap<String, Vec<String>> =
3870 std::collections::HashMap::new();
3871 for (entity_id, event) in &entity_latest {
3872 let payload = event.payload();
3873 let mut key_parts = serde_json::Map::new();
3874 for field in &group_by_fields {
3875 let value = payload
3876 .get(*field)
3877 .cloned()
3878 .unwrap_or(serde_json::Value::Null);
3879 key_parts.insert((*field).to_string(), value);
3880 }
3881 let key_str = serde_json::to_string(&key_parts).unwrap_or_default();
3882 groups.entry(key_str).or_default().push(entity_id.clone());
3883 }
3884
3885 let duplicate_groups: Vec<_> = groups
3886 .into_iter()
3887 .filter(|(_, ids)| ids.len() > 1)
3888 .collect();
3889
3890 assert_eq!(duplicate_groups.len(), 1);
3892 let (_, ref ids) = duplicate_groups[0];
3893 assert_eq!(ids.len(), 2);
3894 let mut sorted_ids = ids.clone();
3895 sorted_ids.sort();
3896 assert_eq!(sorted_ids, vec!["idx-1", "idx-2"]);
3897 }
3898
3899 #[tokio::test]
3900 async fn test_projection_state_cache() {
3901 let store = create_test_store();
3902
3903 let cache = store.projection_state_cache();
3905 cache.insert(
3906 "entity_snapshots:user-123".to_string(),
3907 serde_json::json!({"name": "Test User", "age": 30}),
3908 );
3909
3910 let state = cache.get("entity_snapshots:user-123");
3912 assert!(state.is_some());
3913 let state = state.unwrap();
3914 assert_eq!(state["name"], "Test User");
3915 assert_eq!(state["age"], 30);
3916 }
3917
3918 #[tokio::test]
3919 async fn test_projection_manager_list_projections() {
3920 let store = create_test_store();
3921
3922 let projection_manager = store.projection_manager();
3924 let projections = projection_manager.list_projections();
3925
3926 assert!(projections.len() >= 2);
3928
3929 let names: Vec<&str> = projections.iter().map(|(name, _)| name.as_str()).collect();
3930 assert!(names.contains(&"entity_snapshots"));
3931 assert!(names.contains(&"event_counters"));
3932 }
3933
3934 #[tokio::test]
3935 async fn test_projection_state_after_event_ingestion() {
3936 let store = create_test_store();
3937
3938 let event = create_test_event("user-456", "user.created");
3940 store.ingest(&event).unwrap();
3941
3942 let projection_manager = store.projection_manager();
3944 let snapshot_projection = projection_manager
3945 .get_projection("entity_snapshots")
3946 .unwrap();
3947
3948 let state = snapshot_projection.get_state("user-456");
3949 assert!(state.is_some());
3950 let state = state.unwrap();
3951 assert_eq!(state["name"], "Test");
3952 assert_eq!(state["value"], 42);
3953 }
3954
3955 #[tokio::test]
3956 async fn test_projection_state_cache_multiple_entities() {
3957 let store = create_test_store();
3958 let cache = store.projection_state_cache();
3959
3960 for i in 0..10 {
3962 cache.insert(
3963 format!("entity_snapshots:entity-{i}"),
3964 serde_json::json!({"id": i, "status": "active"}),
3965 );
3966 }
3967
3968 assert_eq!(cache.len(), 10);
3970
3971 for i in 0..10 {
3973 let key = format!("entity_snapshots:entity-{i}");
3974 let state = cache.get(&key);
3975 assert!(state.is_some());
3976 assert_eq!(state.unwrap()["id"], i);
3977 }
3978 }
3979
3980 #[tokio::test]
3981 async fn test_projection_state_update() {
3982 let store = create_test_store();
3983 let cache = store.projection_state_cache();
3984
3985 cache.insert(
3987 "entity_snapshots:user-789".to_string(),
3988 serde_json::json!({"balance": 100}),
3989 );
3990
3991 cache.insert(
3993 "entity_snapshots:user-789".to_string(),
3994 serde_json::json!({"balance": 150}),
3995 );
3996
3997 let state = cache.get("entity_snapshots:user-789").unwrap();
3999 assert_eq!(state["balance"], 150);
4000 }
4001
4002 #[tokio::test]
4003 async fn test_event_counter_projection() {
4004 let store = create_test_store();
4005
4006 store
4008 .ingest(&create_test_event("user-1", "user.created"))
4009 .unwrap();
4010 store
4011 .ingest(&create_test_event("user-2", "user.created"))
4012 .unwrap();
4013 store
4014 .ingest(&create_test_event("user-1", "user.updated"))
4015 .unwrap();
4016
4017 let projection_manager = store.projection_manager();
4019 let counter_projection = projection_manager.get_projection("event_counters").unwrap();
4020
4021 let created_state = counter_projection.get_state("user.created");
4023 assert!(created_state.is_some());
4024 assert_eq!(created_state.unwrap()["count"], 2);
4025
4026 let updated_state = counter_projection.get_state("user.updated");
4027 assert!(updated_state.is_some());
4028 assert_eq!(updated_state.unwrap()["count"], 1);
4029 }
4030
4031 #[tokio::test]
4032 async fn test_projection_state_cache_key_format() {
4033 let store = create_test_store();
4034 let cache = store.projection_state_cache();
4035
4036 let key = "orders:order-12345".to_string();
4038 cache.insert(key.clone(), serde_json::json!({"total": 99.99}));
4039
4040 let state = cache.get(&key).unwrap();
4041 assert_eq!(state["total"], 99.99);
4042 }
4043
4044 #[tokio::test]
4045 async fn test_projection_state_cache_removal() {
4046 let store = create_test_store();
4047 let cache = store.projection_state_cache();
4048
4049 cache.insert(
4051 "test:entity-1".to_string(),
4052 serde_json::json!({"data": "value"}),
4053 );
4054 assert_eq!(cache.len(), 1);
4055
4056 cache.remove("test:entity-1");
4057 assert_eq!(cache.len(), 0);
4058 assert!(cache.get("test:entity-1").is_none());
4059 }
4060
4061 #[tokio::test]
4062 async fn test_get_nonexistent_projection() {
4063 let store = create_test_store();
4064 let projection_manager = store.projection_manager();
4065
4066 let projection = projection_manager.get_projection("nonexistent_projection");
4068 assert!(projection.is_none());
4069 }
4070
4071 #[tokio::test]
4072 async fn test_get_nonexistent_entity_state() {
4073 let store = create_test_store();
4074 let projection_manager = store.projection_manager();
4075
4076 let snapshot_projection = projection_manager
4078 .get_projection("entity_snapshots")
4079 .unwrap();
4080 let state = snapshot_projection.get_state("nonexistent-entity-xyz");
4081 assert!(state.is_none());
4082 }
4083
4084 #[tokio::test]
4085 async fn test_projection_state_cache_concurrent_access() {
4086 let store = create_test_store();
4087 let cache = store.projection_state_cache();
4088
4089 let handles: Vec<_> = (0..10)
4091 .map(|i| {
4092 let cache_clone = cache.clone();
4093 tokio::spawn(async move {
4094 cache_clone.insert(
4095 format!("concurrent:entity-{i}"),
4096 serde_json::json!({"thread": i}),
4097 );
4098 })
4099 })
4100 .collect();
4101
4102 for handle in handles {
4103 handle.await.unwrap();
4104 }
4105
4106 assert_eq!(cache.len(), 10);
4108 }
4109
4110 #[tokio::test]
4111 async fn test_projection_state_large_payload() {
4112 let store = create_test_store();
4113 let cache = store.projection_state_cache();
4114
4115 let large_array: Vec<serde_json::Value> = (0..1000)
4117 .map(|i| serde_json::json!({"item": i, "description": "test item with some padding data to increase size"}))
4118 .collect();
4119
4120 cache.insert(
4121 "large:entity-1".to_string(),
4122 serde_json::json!({"items": large_array}),
4123 );
4124
4125 let state = cache.get("large:entity-1").unwrap();
4126 let items = state["items"].as_array().unwrap();
4127 assert_eq!(items.len(), 1000);
4128 }
4129
4130 #[tokio::test]
4131 async fn test_projection_state_complex_json() {
4132 let store = create_test_store();
4133 let cache = store.projection_state_cache();
4134
4135 let complex_state = serde_json::json!({
4137 "user": {
4138 "id": "user-123",
4139 "profile": {
4140 "name": "John Doe",
4141 "email": "john@example.com",
4142 "settings": {
4143 "theme": "dark",
4144 "notifications": true
4145 }
4146 },
4147 "roles": ["admin", "user"],
4148 "metadata": {
4149 "created_at": "2025-01-01T00:00:00Z",
4150 "last_login": null
4151 }
4152 }
4153 });
4154
4155 cache.insert("complex:user-123".to_string(), complex_state);
4156
4157 let state = cache.get("complex:user-123").unwrap();
4158 assert_eq!(state["user"]["profile"]["name"], "John Doe");
4159 assert_eq!(state["user"]["roles"][0], "admin");
4160 assert!(state["user"]["metadata"]["last_login"].is_null());
4161 }
4162
4163 #[tokio::test]
4164 async fn test_projection_state_cache_iteration() {
4165 let store = create_test_store();
4166 let cache = store.projection_state_cache();
4167
4168 for i in 0..5 {
4170 cache.insert(format!("iter:entity-{i}"), serde_json::json!({"index": i}));
4171 }
4172
4173 let entries: Vec<_> = cache.iter().map(|entry| entry.key().clone()).collect();
4175 assert_eq!(entries.len(), 5);
4176 }
4177
4178 #[tokio::test]
4179 async fn test_projection_manager_get_entity_snapshots() {
4180 let store = create_test_store();
4181 let projection_manager = store.projection_manager();
4182
4183 let projection = projection_manager.get_projection("entity_snapshots");
4185 assert!(projection.is_some());
4186 assert_eq!(projection.unwrap().name(), "entity_snapshots");
4187 }
4188
4189 #[tokio::test]
4190 async fn test_projection_manager_get_event_counters() {
4191 let store = create_test_store();
4192 let projection_manager = store.projection_manager();
4193
4194 let projection = projection_manager.get_projection("event_counters");
4196 assert!(projection.is_some());
4197 assert_eq!(projection.unwrap().name(), "event_counters");
4198 }
4199
4200 #[tokio::test]
4201 async fn test_projection_state_cache_overwrite() {
4202 let store = create_test_store();
4203 let cache = store.projection_state_cache();
4204
4205 cache.insert(
4207 "overwrite:entity-1".to_string(),
4208 serde_json::json!({"version": 1}),
4209 );
4210
4211 cache.insert(
4213 "overwrite:entity-1".to_string(),
4214 serde_json::json!({"version": 2}),
4215 );
4216
4217 cache.insert(
4219 "overwrite:entity-1".to_string(),
4220 serde_json::json!({"version": 3}),
4221 );
4222
4223 let state = cache.get("overwrite:entity-1").unwrap();
4224 assert_eq!(state["version"], 3);
4225
4226 assert_eq!(cache.len(), 1);
4228 }
4229
4230 #[tokio::test]
4231 async fn test_projection_state_multiple_projections() {
4232 let store = create_test_store();
4233 let cache = store.projection_state_cache();
4234
4235 cache.insert(
4237 "entity_snapshots:user-1".to_string(),
4238 serde_json::json!({"name": "Alice"}),
4239 );
4240 cache.insert(
4241 "event_counters:user.created".to_string(),
4242 serde_json::json!({"count": 5}),
4243 );
4244 cache.insert(
4245 "custom_projection:order-1".to_string(),
4246 serde_json::json!({"total": 150.0}),
4247 );
4248
4249 assert_eq!(
4251 cache.get("entity_snapshots:user-1").unwrap()["name"],
4252 "Alice"
4253 );
4254 assert_eq!(
4255 cache.get("event_counters:user.created").unwrap()["count"],
4256 5
4257 );
4258 assert_eq!(
4259 cache.get("custom_projection:order-1").unwrap()["total"],
4260 150.0
4261 );
4262 }
4263
4264 #[tokio::test]
4265 async fn test_bulk_projection_state_access() {
4266 let store = create_test_store();
4267
4268 for i in 0..5 {
4270 let event = create_test_event(&format!("bulk-user-{i}"), "user.created");
4271 store.ingest(&event).unwrap();
4272 }
4273
4274 let projection_manager = store.projection_manager();
4276 let snapshot_projection = projection_manager
4277 .get_projection("entity_snapshots")
4278 .unwrap();
4279
4280 for i in 0..5 {
4282 let state = snapshot_projection.get_state(&format!("bulk-user-{i}"));
4283 assert!(state.is_some(), "Entity bulk-user-{i} should have state");
4284 }
4285 }
4286
4287 #[tokio::test]
4288 async fn test_bulk_save_projection_states() {
4289 let store = create_test_store();
4290 let cache = store.projection_state_cache();
4291
4292 let states = vec![
4294 BulkSaveStateItem {
4295 entity_id: "bulk-entity-1".to_string(),
4296 state: serde_json::json!({"name": "Entity 1", "value": 100}),
4297 },
4298 BulkSaveStateItem {
4299 entity_id: "bulk-entity-2".to_string(),
4300 state: serde_json::json!({"name": "Entity 2", "value": 200}),
4301 },
4302 BulkSaveStateItem {
4303 entity_id: "bulk-entity-3".to_string(),
4304 state: serde_json::json!({"name": "Entity 3", "value": 300}),
4305 },
4306 ];
4307
4308 let projection_name = "test_projection";
4309
4310 for item in &states {
4312 cache.insert(
4313 format!("{projection_name}:{}", item.entity_id),
4314 item.state.clone(),
4315 );
4316 }
4317
4318 assert_eq!(cache.len(), 3);
4320
4321 let state1 = cache.get("test_projection:bulk-entity-1").unwrap();
4322 assert_eq!(state1["name"], "Entity 1");
4323 assert_eq!(state1["value"], 100);
4324
4325 let state2 = cache.get("test_projection:bulk-entity-2").unwrap();
4326 assert_eq!(state2["name"], "Entity 2");
4327 assert_eq!(state2["value"], 200);
4328
4329 let state3 = cache.get("test_projection:bulk-entity-3").unwrap();
4330 assert_eq!(state3["name"], "Entity 3");
4331 assert_eq!(state3["value"], 300);
4332 }
4333
4334 #[tokio::test]
4335 async fn test_bulk_save_empty_states() {
4336 let store = create_test_store();
4337 let cache = store.projection_state_cache();
4338
4339 cache.clear();
4341
4342 let states: Vec<BulkSaveStateItem> = vec![];
4344 assert_eq!(states.len(), 0);
4345
4346 assert_eq!(cache.len(), 0);
4348 }
4349
4350 #[tokio::test]
4351 async fn test_bulk_save_overwrites_existing() {
4352 let store = create_test_store();
4353 let cache = store.projection_state_cache();
4354
4355 cache.insert(
4357 "test:entity-1".to_string(),
4358 serde_json::json!({"version": 1, "data": "initial"}),
4359 );
4360
4361 let new_state = serde_json::json!({"version": 2, "data": "updated"});
4363 cache.insert("test:entity-1".to_string(), new_state);
4364
4365 let state = cache.get("test:entity-1").unwrap();
4367 assert_eq!(state["version"], 2);
4368 assert_eq!(state["data"], "updated");
4369 }
4370
4371 #[tokio::test]
4372 async fn test_bulk_save_high_volume() {
4373 let store = create_test_store();
4374 let cache = store.projection_state_cache();
4375
4376 for i in 0..1000 {
4378 cache.insert(
4379 format!("volume_test:entity-{i}"),
4380 serde_json::json!({"index": i, "status": "active"}),
4381 );
4382 }
4383
4384 assert_eq!(cache.len(), 1000);
4386
4387 assert_eq!(cache.get("volume_test:entity-0").unwrap()["index"], 0);
4389 assert_eq!(cache.get("volume_test:entity-500").unwrap()["index"], 500);
4390 assert_eq!(cache.get("volume_test:entity-999").unwrap()["index"], 999);
4391 }
4392
4393 #[tokio::test]
4394 async fn test_bulk_save_different_projections() {
4395 let store = create_test_store();
4396 let cache = store.projection_state_cache();
4397
4398 let projections = ["entity_snapshots", "event_counters", "custom_analytics"];
4400
4401 for proj in &projections {
4402 for i in 0..5 {
4403 cache.insert(
4404 format!("{proj}:entity-{i}"),
4405 serde_json::json!({"projection": proj, "id": i}),
4406 );
4407 }
4408 }
4409
4410 assert_eq!(cache.len(), 15);
4412
4413 for proj in &projections {
4415 let state = cache.get(&format!("{proj}:entity-0")).unwrap();
4416 assert_eq!(state["projection"], *proj);
4417 }
4418 }
4419
4420 #[tokio::test]
4431 async fn get_projection_state_falls_back_to_cache_when_unregistered() {
4432 let store = create_test_store();
4433 store.projection_state_cache().insert(
4434 "assets:BTC".to_string(),
4435 serde_json::json!({"symbol": "BTC", "altname": "Bitcoin"}),
4436 );
4437
4438 let resp = get_projection_state(
4439 State(Arc::clone(&store)),
4440 Path(("assets".to_string(), "BTC".to_string())),
4441 )
4442 .await
4443 .expect("should not error when projection is not registered");
4444
4445 assert_eq!(resp.0["found"], serde_json::Value::Bool(true));
4446 assert_eq!(resp.0["state"]["symbol"], "BTC");
4447 assert_eq!(resp.0["state"]["altname"], "Bitcoin");
4448 }
4449
4450 #[tokio::test]
4451 async fn get_projection_state_returns_not_found_when_absent_everywhere() {
4452 let store = create_test_store();
4453
4454 let resp = get_projection_state(
4455 State(Arc::clone(&store)),
4456 Path(("assets".to_string(), "UNKNOWN".to_string())),
4457 )
4458 .await
4459 .unwrap();
4460
4461 assert_eq!(resp.0["found"], serde_json::Value::Bool(false));
4462 assert_eq!(resp.0["state"], serde_json::Value::Null);
4463 }
4464
4465 #[tokio::test]
4466 async fn get_projection_state_registered_wins_over_cache() {
4467 let store = create_test_store();
4468
4469 let event = create_test_event("user-777", "user.created");
4471 store.ingest(&event).unwrap();
4472
4473 store.projection_state_cache().insert(
4475 "entity_snapshots:user-777".to_string(),
4476 serde_json::json!({"stolen": "value"}),
4477 );
4478
4479 let resp = get_projection_state(
4480 State(Arc::clone(&store)),
4481 Path(("entity_snapshots".to_string(), "user-777".to_string())),
4482 )
4483 .await
4484 .unwrap();
4485
4486 assert_eq!(resp.0["found"], serde_json::Value::Bool(true));
4489 assert!(
4490 resp.0["state"].get("stolen").is_none(),
4491 "cache entry must not shadow registered projection state: got {:?}",
4492 resp.0["state"]
4493 );
4494 }
4495
4496 #[tokio::test]
4497 async fn get_projection_state_summary_returns_cache_without_registration() {
4498 let store = create_test_store();
4499 let cache = store.projection_state_cache();
4500 cache.insert("assets:BTC".into(), serde_json::json!({"symbol": "BTC"}));
4501 cache.insert("assets:ETH".into(), serde_json::json!({"symbol": "ETH"}));
4502 cache.insert("trades:t-1".into(), serde_json::json!({"x": 1}));
4504
4505 let resp = get_projection_state_summary(
4506 State(Arc::clone(&store)),
4507 Path("assets".to_string()),
4508 Query(ProjectionStateSummaryParams::default()),
4509 )
4510 .await
4511 .unwrap();
4512
4513 assert_eq!(resp.0["total"], 2);
4514 let states = resp.0["states"].as_array().unwrap();
4515 let entity_ids: Vec<&str> = states
4516 .iter()
4517 .map(|s| s["entity_id"].as_str().unwrap())
4518 .collect();
4519 assert!(entity_ids.contains(&"BTC"));
4520 assert!(entity_ids.contains(&"ETH"));
4521 }
4522
4523 #[tokio::test]
4524 async fn bulk_get_projection_states_falls_back_to_cache() {
4525 let store = create_test_store();
4526 let cache = store.projection_state_cache();
4527 cache.insert("assets:BTC".into(), serde_json::json!({"symbol": "BTC"}));
4528 cache.insert("assets:ETH".into(), serde_json::json!({"symbol": "ETH"}));
4529
4530 let req = BulkGetStateRequest {
4531 entity_ids: vec!["BTC".into(), "ETH".into(), "MISSING".into()],
4532 };
4533
4534 let resp = bulk_get_projection_states(
4535 State(Arc::clone(&store)),
4536 Path("assets".to_string()),
4537 Json(req),
4538 )
4539 .await
4540 .unwrap();
4541
4542 assert_eq!(resp.0["total"], 3);
4543 let states = resp.0["states"].as_array().unwrap();
4544 let by_id: std::collections::HashMap<&str, &serde_json::Value> = states
4545 .iter()
4546 .map(|s| (s["entity_id"].as_str().unwrap(), s))
4547 .collect();
4548
4549 assert_eq!(by_id["BTC"]["found"], serde_json::Value::Bool(true));
4550 assert_eq!(by_id["BTC"]["state"]["symbol"], "BTC");
4551 assert_eq!(by_id["ETH"]["found"], serde_json::Value::Bool(true));
4552 assert_eq!(by_id["MISSING"]["found"], serde_json::Value::Bool(false));
4553 }
4554
4555 #[tokio::test]
4559 async fn poll_consumer_events_flattens_event_alongside_position() {
4560 let store = create_test_store();
4561 store
4562 .ingest(&create_test_event("user-1", "user.created"))
4563 .unwrap();
4564 store
4565 .ingest(&create_test_event("user-2", "user.updated"))
4566 .unwrap();
4567 store.consumer_registry().register("w1", &[]);
4568
4569 let resp = poll_consumer_events(
4570 State(Arc::clone(&store)),
4571 Path("w1".to_string()),
4572 Query(ConsumerPollQuery { limit: Some(10) }),
4573 )
4574 .await
4575 .unwrap();
4576
4577 let body = serde_json::to_value(&resp.0).unwrap();
4578 assert_eq!(body["count"], 2);
4579 let first = &body["events"][0];
4580 assert_eq!(first["position"], 1);
4581 assert!(
4582 first.get("event").is_none(),
4583 "event must be flattened, not nested: got {first:?}"
4584 );
4585 assert_eq!(first["event_type"], "user.created");
4586 assert_eq!(first["entity_id"], "user-1");
4587 }
4588}