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 =
301 super::archive_work::append(Arc::clone(&store), event, expected_version).await?;
302
303 tracing::info!("Event ingested: {}", event_id);
304
305 Ok(Json(IngestEventResponse {
306 event_id,
307 timestamp,
308 version: Some(new_version),
309 }))
310}
311
312async fn enforce_schema_if_configured(
333 state: &AppState,
334 tenant_id: &str,
335 event: &Event,
336) -> Result<()> {
337 let Ok(parsed) = TenantId::new(tenant_id.to_string()) else {
341 return Ok(());
342 };
343 let mode = match state.tenant_repo.find_by_id(&parsed).await {
344 Ok(Some(t)) => t.schema_enforcement(),
345 _ => SchemaEnforcement::Permissive,
348 };
349 if matches!(mode, SchemaEnforcement::Permissive) {
350 return Ok(());
351 }
352
353 let registry = state.store.schema_registry();
357 let Ok(schema) = registry.get_schema(event.event_type.as_str(), None) else {
358 return Ok(());
361 };
362
363 let result = registry
364 .validate(
365 event.event_type.as_str(),
366 Some(schema.version),
367 &event.payload,
368 )
369 .map_err(|e| crate::error::AllSourceError::InternalError(e.to_string()))?;
370
371 if result.valid {
372 return Ok(());
373 }
374
375 match mode {
376 SchemaEnforcement::Strict => Err(crate::error::AllSourceError::SchemaViolation {
377 event_type: event.event_type.as_str().to_string(),
378 schema_version: result.schema_version,
379 errors: result.errors,
380 }),
381 SchemaEnforcement::Warn => {
382 tracing::warn!(
383 tenant = %tenant_id,
384 event_type = %event.event_type.as_str(),
385 schema_version = result.schema_version,
386 errors = ?result.errors,
387 "schema violation (warn mode — write accepted)"
388 );
389 Ok(())
390 }
391 SchemaEnforcement::Permissive => Ok(()),
393 }
394}
395
396pub async fn ingest_event_v1(
397 State(state): State<AppState>,
398 Json(req): Json<IngestEventRequest>,
399) -> Result<Json<IngestEventResponse>> {
400 let expected_version = req.expected_version;
401
402 let tenant_id = req.tenant_id.unwrap_or_else(|| "default".to_string());
403
404 let event = Event::from_strings(
405 req.event_type,
406 req.entity_id,
407 tenant_id.clone(),
408 req.payload,
409 req.metadata,
410 )?;
411
412 enforce_schema_if_configured(&state, &tenant_id, &event).await?;
415
416 let event_id = event.id;
417 let timestamp = event.timestamp;
418
419 let new_version =
420 super::archive_work::append(Arc::clone(&state.store), event, expected_version).await?;
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 =
466 super::archive_work::append(Arc::clone(&store), event, expected_version).await?;
467
468 ingested_events.push(IngestEventResponse {
469 event_id,
470 timestamp,
471 version: Some(new_version),
472 });
473 }
474
475 let ingested = ingested_events.len();
476 tracing::info!("Batch ingested {} events", ingested);
477
478 Ok(Json(IngestEventsBatchResponse {
479 total,
480 ingested,
481 events: ingested_events,
482 }))
483}
484
485pub async fn ingest_events_batch_v1(
491 State(state): State<AppState>,
492 Json(req): Json<IngestEventsBatchRequest>,
493) -> Result<Json<IngestEventsBatchResponse>> {
494 let total = req.events.len();
495 let mut ingested_events = Vec::with_capacity(total);
496
497 for event_req in req.events {
498 let tenant_id = event_req.tenant_id.unwrap_or_else(|| "default".to_string());
499 let expected_version = event_req.expected_version;
500
501 let event = Event::from_strings(
502 event_req.event_type,
503 event_req.entity_id,
504 tenant_id.clone(),
505 event_req.payload,
506 event_req.metadata,
507 )?;
508
509 enforce_schema_if_configured(&state, &tenant_id, &event).await?;
510
511 let event_id = event.id;
512 let timestamp = event.timestamp;
513
514 let new_version =
515 super::archive_work::append(Arc::clone(&state.store), event, expected_version).await?;
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
556#[derive(Debug, Default, Deserialize)]
558pub struct EventIntegrityParam {
559 pub integrity: Option<String>,
560}
561
562pub async fn query_events(
563 OptionalAuth(auth): OptionalAuth,
564 Query(req): Query<QueryEventsRequest>,
565 Query(order_param): Query<EventOrderParam>,
566 Query(offset_param): Query<EventOffsetParam>,
567 Query(integrity_param): Query<EventIntegrityParam>,
568 State(store): State<SharedStore>,
569) -> Result<Json<QueryEventsResponse>> {
570 let offset = offset_param.offset.unwrap_or(0);
571 let queried_entity_id = req.entity_id.clone();
572
573 let descending = match order_param.order.as_deref() {
578 None => false,
579 Some(o) if o.eq_ignore_ascii_case("asc") => false,
580 Some(o) if o.eq_ignore_ascii_case("desc") => true,
581 Some(other) => {
582 return Err(crate::error::AllSourceError::InvalidInput(format!(
583 "invalid 'order' value '{other}': expected 'asc' or 'desc'"
584 )));
585 }
586 };
587
588 let enforced_tenant = req
596 .tenant_id
597 .clone()
598 .or_else(|| auth.as_ref().map(|a| a.tenant_id().to_string()));
599
600 if let Some(protocol) = integrity_param.integrity {
601 return super::retained_query::query(
602 store,
603 QueryEventsRequest {
604 tenant_id: enforced_tenant,
605 ..req
606 },
607 offset,
608 descending,
609 &protocol,
610 )
611 .await
612 .map(Json);
613 }
614
615 if enforced_tenant.as_deref().unwrap_or("").is_empty() {
622 return Ok(Json(QueryEventsResponse {
623 events: Vec::new(),
624 count: 0,
625 total_count: 0,
626 has_more: false,
627 entity_version: None,
628 archive_integrity: None,
629 }));
630 }
631
632 let scoped_req = QueryEventsRequest {
640 tenant_id: enforced_tenant,
641 ..req
642 };
643 let (limited_events, total_count) = store.query_window(&scoped_req, offset, descending)?;
644
645 let count = limited_events.len();
646 let has_more = offset + count < total_count;
650 let events: Vec<EventDto> = limited_events.iter().map(EventDto::from).collect();
651
652 let entity_version = queried_entity_id
654 .as_deref()
655 .map(|eid| store.get_entity_version(eid));
656
657 tracing::debug!("Query returned {} events (total: {})", count, total_count);
658
659 Ok(Json(QueryEventsResponse {
660 events,
661 count,
662 total_count,
663 has_more,
664 entity_version,
665 archive_integrity: None,
666 }))
667}
668
669pub async fn list_entities(
670 State(store): State<SharedStore>,
671 Query(req): Query<ListEntitiesRequest>,
672) -> Result<Json<ListEntitiesResponse>> {
673 use std::collections::HashMap;
674
675 let query_req = QueryEventsRequest {
677 entity_id: None,
678 event_type: None,
679 tenant_id: None,
680 as_of: None,
681 since: None,
682 until: None,
683 limit: None,
684 event_type_prefix: req.event_type_prefix,
685 exclude_event_type_prefix: None,
686 payload_filter: req.payload_filter,
687 };
688 let events = store.query(&query_req)?;
689
690 let mut entity_map: HashMap<String, Vec<&Event>> = HashMap::new();
692 for event in &events {
693 entity_map
694 .entry(event.entity_id().to_string())
695 .or_default()
696 .push(event);
697 }
698
699 let ascending = match req.order.as_deref() {
703 None => false,
704 Some(o) if o.eq_ignore_ascii_case("desc") => false,
705 Some(o) if o.eq_ignore_ascii_case("asc") => true,
706 Some(other) => {
707 return Err(crate::error::AllSourceError::InvalidInput(format!(
708 "invalid 'order' value '{other}': expected 'asc' or 'desc'"
709 )));
710 }
711 };
712
713 let mut summaries: Vec<EntitySummary> = entity_map
717 .into_iter()
718 .map(|(entity_id, events)| {
719 let last = events.iter().max_by_key(|e| e.timestamp()).unwrap();
720 EntitySummary {
721 entity_id,
722 event_count: events.len(),
723 last_event_type: last.event_type_str().to_string(),
724 last_event_at: last.timestamp(),
725 }
726 })
727 .collect();
728 summaries.sort_by(|a, b| {
729 let by_time = a.last_event_at.cmp(&b.last_event_at);
730 let by_time = if ascending {
731 by_time
732 } else {
733 by_time.reverse()
734 };
735 by_time.then_with(|| a.entity_id.cmp(&b.entity_id))
736 });
737
738 let total = summaries.len();
739
740 let offset = req.offset.unwrap_or(0);
742 let summaries: Vec<EntitySummary> = summaries.into_iter().skip(offset).collect::<Vec<_>>();
743 let summaries = if let Some(limit) = req.limit {
744 let has_more = summaries.len() > limit;
745 let truncated: Vec<EntitySummary> = summaries.into_iter().take(limit).collect();
746 return Ok(Json(ListEntitiesResponse {
747 entities: truncated,
748 total,
749 has_more,
750 }));
751 } else {
752 summaries
753 };
754
755 Ok(Json(ListEntitiesResponse {
756 entities: summaries,
757 total,
758 has_more: false,
759 }))
760}
761
762pub async fn detect_duplicates(
763 State(store): State<SharedStore>,
764 Query(req): Query<DetectDuplicatesRequest>,
765) -> Result<Json<DetectDuplicatesResponse>> {
766 use std::collections::HashMap;
767
768 let group_by_fields: Vec<&str> = req.group_by.split(',').map(str::trim).collect();
769
770 let query_req = QueryEventsRequest {
772 entity_id: None,
773 event_type: None,
774 tenant_id: None,
775 as_of: None,
776 since: None,
777 until: None,
778 limit: None,
779 event_type_prefix: Some(req.event_type_prefix),
780 exclude_event_type_prefix: None,
781 payload_filter: None,
782 };
783 let events = store.query(&query_req)?;
784
785 let mut entity_latest: HashMap<String, &Event> = HashMap::new();
788 for event in &events {
789 let eid = event.entity_id().to_string();
790 entity_latest
791 .entry(eid)
792 .and_modify(|existing| {
793 if event.timestamp() > existing.timestamp() {
794 *existing = event;
795 }
796 })
797 .or_insert(event);
798 }
799
800 let mut groups: HashMap<String, Vec<String>> = HashMap::new();
802 for (entity_id, event) in &entity_latest {
803 let payload = event.payload();
804 let mut key_parts = serde_json::Map::new();
805 for field in &group_by_fields {
806 let value = payload
807 .get(*field)
808 .cloned()
809 .unwrap_or(serde_json::Value::Null);
810 key_parts.insert((*field).to_string(), value);
811 }
812 let key_str = serde_json::to_string(&key_parts).unwrap_or_default();
813 groups.entry(key_str).or_default().push(entity_id.clone());
814 }
815
816 let mut duplicate_groups: Vec<DuplicateGroup> = groups
818 .into_iter()
819 .filter(|(_, ids)| ids.len() > 1)
820 .map(|(key_str, mut ids)| {
821 ids.sort();
822 let key: serde_json::Value =
823 serde_json::from_str(&key_str).unwrap_or(serde_json::Value::Null);
824 let count = ids.len();
825 DuplicateGroup {
826 key,
827 entity_ids: ids,
828 count,
829 }
830 })
831 .collect();
832
833 duplicate_groups.sort_by_key(|a| std::cmp::Reverse(a.count));
835
836 let total = duplicate_groups.len();
837
838 let offset = req.offset.unwrap_or(0);
840 let duplicate_groups: Vec<DuplicateGroup> = duplicate_groups.into_iter().skip(offset).collect();
841
842 if let Some(limit) = req.limit {
843 let has_more = duplicate_groups.len() > limit;
844 let truncated: Vec<DuplicateGroup> = duplicate_groups.into_iter().take(limit).collect();
845 return Ok(Json(DetectDuplicatesResponse {
846 duplicates: truncated,
847 total,
848 has_more,
849 }));
850 }
851
852 Ok(Json(DetectDuplicatesResponse {
853 duplicates: duplicate_groups,
854 total,
855 has_more: false,
856 }))
857}
858
859#[derive(Deserialize)]
860pub struct EntityStateParams {
861 as_of: Option<chrono::DateTime<chrono::Utc>>,
862 tenant_id: Option<String>,
867}
868
869pub async fn get_entity_state(
870 State(store): State<SharedStore>,
871 Path(entity_id): Path<String>,
872 Query(params): Query<EntityStateParams>,
873) -> Result<Json<serde_json::Value>> {
874 let state = match params.tenant_id.as_deref() {
875 Some(tenant_id) => {
876 store.reconstruct_state_for_tenant(&entity_id, params.as_of, tenant_id)?
877 }
878 None => store.reconstruct_state(&entity_id, params.as_of)?,
879 };
880
881 tracing::info!("State reconstructed for entity: {}", entity_id);
882
883 Ok(Json(state))
884}
885
886pub async fn get_entity_snapshot(
887 State(store): State<SharedStore>,
888 Path(entity_id): Path<String>,
889 Query(params): Query<EntityStateParams>,
890) -> Result<Json<serde_json::Value>> {
891 let snapshot = match params.tenant_id.as_deref() {
895 Some(tenant_id) => store.reconstruct_state_for_tenant(&entity_id, None, tenant_id)?,
896 None => store.get_snapshot(&entity_id)?,
897 };
898
899 tracing::debug!("Snapshot retrieved for entity: {}", entity_id);
900
901 Ok(Json(snapshot))
902}
903
904#[derive(Debug, Deserialize)]
906pub struct StatsParams {
907 pub tenant_id: Option<String>,
913}
914
915pub async fn get_stats(
916 State(store): State<SharedStore>,
917 Query(params): Query<StatsParams>,
918) -> impl IntoResponse {
919 match params.tenant_id.as_deref() {
920 Some(tenant_id) => {
921 Json(serde_json::to_value(store.stats_for_tenant(tenant_id)).unwrap_or_default())
922 }
923 None => Json(serde_json::to_value(store.stats()).unwrap_or_default()),
924 }
925}
926
927#[derive(Debug, Deserialize)]
930pub struct ListStreamsParams {
931 pub tenant_id: Option<String>,
933 pub limit: Option<usize>,
935 pub offset: Option<usize>,
937}
938
939#[derive(Debug, serde::Serialize)]
941pub struct ListStreamsResponse {
942 pub streams: Vec<StreamInfo>,
943 pub total: usize,
944}
945
946pub async fn list_streams(
947 OptionalAuth(auth): OptionalAuth,
948 State(store): State<SharedStore>,
949 Query(params): Query<ListStreamsParams>,
950) -> Json<ListStreamsResponse> {
951 let tenant = params
953 .tenant_id
954 .clone()
955 .or_else(|| auth.as_ref().map(|a| a.tenant_id().to_string()))
956 .filter(|t| !t.is_empty());
957 let Some(tenant) = tenant else {
958 return Json(ListStreamsResponse {
959 streams: vec![],
960 total: 0,
961 });
962 };
963 let mut streams = store.list_streams_for_tenant(&tenant);
964 let total = streams.len();
965
966 streams.sort_by_key(|a| std::cmp::Reverse(a.last_event_at));
968
969 if let Some(offset) = params.offset {
971 if offset < streams.len() {
972 streams = streams[offset..].to_vec();
973 } else {
974 streams = vec![];
975 }
976 }
977
978 if let Some(limit) = params.limit {
979 streams.truncate(limit);
980 }
981
982 tracing::debug!("Listed {} streams (total: {})", streams.len(), total);
983
984 Json(ListStreamsResponse { streams, total })
985}
986
987#[derive(Debug, Deserialize)]
990pub struct ListEventTypesParams {
991 pub tenant_id: Option<String>,
993 pub limit: Option<usize>,
995 pub offset: Option<usize>,
997}
998
999#[derive(Debug, serde::Serialize)]
1001pub struct ListEventTypesResponse {
1002 pub event_types: Vec<EventTypeInfo>,
1003 pub total: usize,
1004}
1005
1006pub async fn list_event_types(
1007 OptionalAuth(auth): OptionalAuth,
1008 State(store): State<SharedStore>,
1009 Query(params): Query<ListEventTypesParams>,
1010) -> Json<ListEventTypesResponse> {
1011 let tenant = params
1013 .tenant_id
1014 .clone()
1015 .or_else(|| auth.as_ref().map(|a| a.tenant_id().to_string()))
1016 .filter(|t| !t.is_empty());
1017 let Some(tenant) = tenant else {
1018 return Json(ListEventTypesResponse {
1019 event_types: vec![],
1020 total: 0,
1021 });
1022 };
1023 let mut event_types = store.list_event_types_for_tenant(&tenant);
1024 let total = event_types.len();
1025
1026 event_types.sort_by_key(|a| std::cmp::Reverse(a.event_count));
1028
1029 if let Some(offset) = params.offset {
1031 if offset < event_types.len() {
1032 event_types = event_types[offset..].to_vec();
1033 } else {
1034 event_types = vec![];
1035 }
1036 }
1037
1038 if let Some(limit) = params.limit {
1039 event_types.truncate(limit);
1040 }
1041
1042 tracing::debug!(
1043 "Listed {} event types (total: {})",
1044 event_types.len(),
1045 total
1046 );
1047
1048 Json(ListEventTypesResponse { event_types, total })
1049}
1050
1051#[derive(Debug, Deserialize)]
1053pub struct WebSocketParams {
1054 pub consumer_id: Option<String>,
1055}
1056
1057pub async fn events_websocket(
1058 ws: WebSocketUpgrade,
1059 State(store): State<SharedStore>,
1060 Query(params): Query<WebSocketParams>,
1061) -> Response {
1062 let websocket_manager = store.websocket_manager();
1063
1064 ws.on_upgrade(move |socket| async move {
1065 if let Some(consumer_id) = params.consumer_id {
1066 websocket_manager
1067 .handle_socket_with_consumer(socket, consumer_id, store)
1068 .await;
1069 } else {
1070 websocket_manager.handle_socket(socket).await;
1071 }
1072 })
1073}
1074
1075pub async fn analytics_frequency(
1077 State(store): State<SharedStore>,
1078 Query(req): Query<EventFrequencyRequest>,
1079) -> Result<Json<EventFrequencyResponse>> {
1080 let response = AnalyticsEngine::event_frequency(&store, &req)?;
1081
1082 tracing::debug!(
1083 "Frequency analysis returned {} buckets",
1084 response.buckets.len()
1085 );
1086
1087 Ok(Json(response))
1088}
1089
1090pub async fn analytics_summary(
1092 State(store): State<SharedStore>,
1093 Query(req): Query<StatsSummaryRequest>,
1094) -> Result<Json<StatsSummaryResponse>> {
1095 let response = AnalyticsEngine::stats_summary(&store, &req)?;
1096
1097 tracing::debug!(
1098 "Stats summary: {} events across {} entities",
1099 response.total_events,
1100 response.unique_entities
1101 );
1102
1103 Ok(Json(response))
1104}
1105
1106pub async fn analytics_correlation(
1108 State(store): State<SharedStore>,
1109 Query(req): Query<CorrelationRequest>,
1110) -> Result<Json<CorrelationResponse>> {
1111 let response = AnalyticsEngine::analyze_correlation(&store, req)?;
1112
1113 tracing::debug!(
1114 "Correlation analysis: {}/{} correlated pairs ({:.2}%)",
1115 response.correlated_pairs,
1116 response.total_a,
1117 response.correlation_percentage
1118 );
1119
1120 Ok(Json(response))
1121}
1122
1123pub async fn create_snapshot(
1125 State(store): State<SharedStore>,
1126 Json(req): Json<CreateSnapshotRequest>,
1127) -> Result<Json<CreateSnapshotResponse>> {
1128 store.create_snapshot(&req.entity_id)?;
1129
1130 let snapshot_manager = store.snapshot_manager();
1131 let snapshot = snapshot_manager
1132 .get_latest_snapshot(&req.entity_id)
1133 .ok_or_else(|| crate::error::AllSourceError::EntityNotFound(req.entity_id.clone()))?;
1134
1135 tracing::info!("📸 Created snapshot for entity: {}", req.entity_id);
1136
1137 Ok(Json(CreateSnapshotResponse {
1138 snapshot_id: snapshot.id,
1139 entity_id: snapshot.entity_id,
1140 created_at: snapshot.created_at,
1141 event_count: snapshot.event_count,
1142 size_bytes: snapshot.metadata.size_bytes,
1143 }))
1144}
1145
1146pub async fn list_snapshots(
1148 State(store): State<SharedStore>,
1149 Query(req): Query<ListSnapshotsRequest>,
1150) -> Result<Json<ListSnapshotsResponse>> {
1151 let snapshot_manager = store.snapshot_manager();
1152
1153 let snapshots: Vec<SnapshotInfo> = if let Some(entity_id) = req.entity_id {
1154 snapshot_manager
1155 .get_all_snapshots(&entity_id)
1156 .into_iter()
1157 .map(SnapshotInfo::from)
1158 .collect()
1159 } else {
1160 let entities = snapshot_manager.list_entities();
1162 entities
1163 .iter()
1164 .flat_map(|entity_id| {
1165 snapshot_manager
1166 .get_all_snapshots(entity_id)
1167 .into_iter()
1168 .map(SnapshotInfo::from)
1169 })
1170 .collect()
1171 };
1172
1173 let total = snapshots.len();
1174
1175 tracing::debug!("Listed {} snapshots", total);
1176
1177 Ok(Json(ListSnapshotsResponse { snapshots, total }))
1178}
1179
1180pub async fn get_latest_snapshot(
1182 State(store): State<SharedStore>,
1183 Path(entity_id): Path<String>,
1184) -> Result<Json<serde_json::Value>> {
1185 let snapshot_manager = store.snapshot_manager();
1186
1187 let snapshot = snapshot_manager
1188 .get_latest_snapshot(&entity_id)
1189 .ok_or_else(|| crate::error::AllSourceError::EntityNotFound(entity_id.clone()))?;
1190
1191 tracing::debug!("Retrieved latest snapshot for entity: {}", entity_id);
1192
1193 Ok(Json(serde_json::json!({
1194 "snapshot_id": snapshot.id,
1195 "entity_id": snapshot.entity_id,
1196 "created_at": snapshot.created_at,
1197 "as_of": snapshot.as_of,
1198 "event_count": snapshot.event_count,
1199 "size_bytes": snapshot.metadata.size_bytes,
1200 "snapshot_type": snapshot.metadata.snapshot_type,
1201 "state": snapshot.state
1202 })))
1203}
1204
1205pub async fn trigger_compaction(
1207 State(store): State<SharedStore>,
1208) -> Result<Json<CompactionResult>> {
1209 let compaction_manager = store.compaction_manager().ok_or_else(|| {
1210 crate::error::AllSourceError::InternalError(
1211 "Compaction not enabled (no Parquet storage)".to_string(),
1212 )
1213 })?;
1214
1215 tracing::info!("📦 Manual compaction triggered via API");
1216
1217 let result = compaction_manager.compact_now()?;
1218
1219 Ok(Json(result))
1220}
1221
1222pub async fn compaction_stats(State(store): State<SharedStore>) -> Result<Json<serde_json::Value>> {
1224 let compaction_manager = store.compaction_manager().ok_or_else(|| {
1225 crate::error::AllSourceError::InternalError(
1226 "Compaction not enabled (no Parquet storage)".to_string(),
1227 )
1228 })?;
1229
1230 let stats = compaction_manager.stats();
1231 let config = compaction_manager.config();
1232
1233 Ok(Json(serde_json::json!({
1234 "stats": stats,
1235 "config": {
1236 "min_files_to_compact": config.min_files_to_compact,
1237 "target_file_size": config.target_file_size,
1238 "max_file_size": config.max_file_size,
1239 "small_file_threshold": config.small_file_threshold,
1240 "compaction_interval_seconds": config.compaction_interval_seconds,
1241 "auto_compact": config.auto_compact,
1242 "strategy": config.strategy
1243 }
1244 })))
1245}
1246
1247pub async fn register_schema(
1249 State(store): State<SharedStore>,
1250 Json(req): Json<RegisterSchemaRequest>,
1251) -> Result<Json<RegisterSchemaResponse>> {
1252 let schema_registry = store.schema_registry();
1253
1254 let response =
1255 schema_registry.register_schema(req.subject, req.schema, req.description, req.tags)?;
1256
1257 tracing::info!(
1258 "📋 Schema registered: v{} for '{}'",
1259 response.version,
1260 response.subject
1261 );
1262
1263 Ok(Json(response))
1264}
1265
1266#[derive(Deserialize)]
1268pub struct GetSchemaParams {
1269 version: Option<u32>,
1270}
1271
1272pub async fn get_schema(
1273 State(store): State<SharedStore>,
1274 Path(subject): Path<String>,
1275 Query(params): Query<GetSchemaParams>,
1276) -> Result<Json<serde_json::Value>> {
1277 let schema_registry = store.schema_registry();
1278
1279 let schema = schema_registry.get_schema(&subject, params.version)?;
1280
1281 tracing::debug!("Retrieved schema v{} for '{}'", schema.version, subject);
1282
1283 Ok(Json(serde_json::json!({
1284 "id": schema.id,
1285 "subject": schema.subject,
1286 "version": schema.version,
1287 "schema": schema.schema,
1288 "created_at": schema.created_at,
1289 "description": schema.description,
1290 "tags": schema.tags
1291 })))
1292}
1293
1294pub async fn list_schema_versions(
1296 State(store): State<SharedStore>,
1297 Path(subject): Path<String>,
1298) -> Result<Json<serde_json::Value>> {
1299 let schema_registry = store.schema_registry();
1300
1301 let versions = schema_registry.list_versions(&subject)?;
1302
1303 Ok(Json(serde_json::json!({
1304 "subject": subject,
1305 "versions": versions
1306 })))
1307}
1308
1309pub async fn list_subjects(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1311 let schema_registry = store.schema_registry();
1312
1313 let subjects = schema_registry.list_subjects();
1314
1315 Json(serde_json::json!({
1316 "subjects": subjects,
1317 "total": subjects.len()
1318 }))
1319}
1320
1321pub async fn validate_event_schema(
1323 State(store): State<SharedStore>,
1324 Json(req): Json<ValidateEventRequest>,
1325) -> Result<Json<ValidateEventResponse>> {
1326 let schema_registry = store.schema_registry();
1327
1328 let response = schema_registry.validate(&req.subject, req.version, &req.payload)?;
1329
1330 if response.valid {
1331 tracing::debug!(
1332 "✅ Event validated against schema '{}' v{}",
1333 req.subject,
1334 response.schema_version
1335 );
1336 } else {
1337 tracing::warn!(
1338 "❌ Event validation failed for '{}': {:?}",
1339 req.subject,
1340 response.errors
1341 );
1342 }
1343
1344 Ok(Json(response))
1345}
1346
1347#[derive(Deserialize)]
1349pub struct SetCompatibilityRequest {
1350 compatibility: CompatibilityMode,
1351}
1352
1353pub async fn set_compatibility_mode(
1354 State(store): State<SharedStore>,
1355 Path(subject): Path<String>,
1356 Json(req): Json<SetCompatibilityRequest>,
1357) -> Json<serde_json::Value> {
1358 let schema_registry = store.schema_registry();
1359
1360 schema_registry.set_compatibility_mode(subject.clone(), req.compatibility);
1361
1362 tracing::info!(
1363 "🔧 Set compatibility mode for '{}' to {:?}",
1364 subject,
1365 req.compatibility
1366 );
1367
1368 Json(serde_json::json!({
1369 "subject": subject,
1370 "compatibility": req.compatibility
1371 }))
1372}
1373
1374pub async fn start_replay(
1376 State(store): State<SharedStore>,
1377 Json(req): Json<StartReplayRequest>,
1378) -> Result<Json<StartReplayResponse>> {
1379 let replay_manager = store.replay_manager();
1380
1381 let response = replay_manager.start_replay(store, req)?;
1382
1383 tracing::info!(
1384 "🔄 Started replay {} with {} events",
1385 response.replay_id,
1386 response.total_events
1387 );
1388
1389 Ok(Json(response))
1390}
1391
1392pub async fn get_replay_progress(
1394 State(store): State<SharedStore>,
1395 Path(replay_id): Path<uuid::Uuid>,
1396) -> Result<Json<ReplayProgress>> {
1397 let replay_manager = store.replay_manager();
1398
1399 let progress = replay_manager.get_progress(replay_id)?;
1400
1401 Ok(Json(progress))
1402}
1403
1404pub async fn list_replays(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1406 let replay_manager = store.replay_manager();
1407
1408 let replays = replay_manager.list_replays();
1409
1410 Json(serde_json::json!({
1411 "replays": replays,
1412 "total": replays.len()
1413 }))
1414}
1415
1416pub async fn cancel_replay(
1418 State(store): State<SharedStore>,
1419 Path(replay_id): Path<uuid::Uuid>,
1420) -> Result<Json<serde_json::Value>> {
1421 let replay_manager = store.replay_manager();
1422
1423 replay_manager.cancel_replay(replay_id)?;
1424
1425 tracing::info!("🛑 Cancelled replay {}", replay_id);
1426
1427 Ok(Json(serde_json::json!({
1428 "replay_id": replay_id,
1429 "status": "cancelled"
1430 })))
1431}
1432
1433pub async fn delete_replay(
1435 State(store): State<SharedStore>,
1436 Path(replay_id): Path<uuid::Uuid>,
1437) -> Result<Json<serde_json::Value>> {
1438 let replay_manager = store.replay_manager();
1439
1440 let deleted = replay_manager.delete_replay(replay_id)?;
1441
1442 if deleted {
1443 tracing::info!("🗑️ Deleted replay {}", replay_id);
1444 }
1445
1446 Ok(Json(serde_json::json!({
1447 "replay_id": replay_id,
1448 "deleted": deleted
1449 })))
1450}
1451
1452pub async fn register_pipeline(
1454 State(store): State<SharedStore>,
1455 Json(config): Json<PipelineConfig>,
1456) -> Result<Json<serde_json::Value>> {
1457 let pipeline_manager = store.pipeline_manager();
1458
1459 let pipeline_id = pipeline_manager.register(config.clone());
1460
1461 tracing::info!(
1462 "🔀 Pipeline registered: {} (name: {})",
1463 pipeline_id,
1464 config.name
1465 );
1466
1467 Ok(Json(serde_json::json!({
1468 "pipeline_id": pipeline_id,
1469 "name": config.name,
1470 "enabled": config.enabled
1471 })))
1472}
1473
1474pub async fn list_pipelines(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1476 let pipeline_manager = store.pipeline_manager();
1477
1478 let pipelines = pipeline_manager.list();
1479
1480 tracing::debug!("Listed {} pipelines", pipelines.len());
1481
1482 Json(serde_json::json!({
1483 "pipelines": pipelines,
1484 "total": pipelines.len()
1485 }))
1486}
1487
1488pub async fn get_pipeline(
1490 State(store): State<SharedStore>,
1491 Path(pipeline_id): Path<uuid::Uuid>,
1492) -> Result<Json<PipelineConfig>> {
1493 let pipeline_manager = store.pipeline_manager();
1494
1495 let pipeline = pipeline_manager.get(pipeline_id).ok_or_else(|| {
1496 crate::error::AllSourceError::ValidationError(format!("Pipeline not found: {pipeline_id}"))
1497 })?;
1498
1499 Ok(Json(pipeline.config().clone()))
1500}
1501
1502pub async fn remove_pipeline(
1504 State(store): State<SharedStore>,
1505 Path(pipeline_id): Path<uuid::Uuid>,
1506) -> Result<Json<serde_json::Value>> {
1507 let pipeline_manager = store.pipeline_manager();
1508
1509 let removed = pipeline_manager.remove(pipeline_id);
1510
1511 if removed {
1512 tracing::info!("🗑️ Removed pipeline {}", pipeline_id);
1513 }
1514
1515 Ok(Json(serde_json::json!({
1516 "pipeline_id": pipeline_id,
1517 "removed": removed
1518 })))
1519}
1520
1521pub async fn all_pipeline_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1523 let pipeline_manager = store.pipeline_manager();
1524
1525 let stats = pipeline_manager.all_stats();
1526
1527 Json(serde_json::json!({
1528 "stats": stats,
1529 "total": stats.len()
1530 }))
1531}
1532
1533pub async fn get_pipeline_stats(
1535 State(store): State<SharedStore>,
1536 Path(pipeline_id): Path<uuid::Uuid>,
1537) -> Result<Json<PipelineStats>> {
1538 let pipeline_manager = store.pipeline_manager();
1539
1540 let pipeline = pipeline_manager.get(pipeline_id).ok_or_else(|| {
1541 crate::error::AllSourceError::ValidationError(format!("Pipeline not found: {pipeline_id}"))
1542 })?;
1543
1544 Ok(Json(pipeline.stats()))
1545}
1546
1547pub async fn reset_pipeline(
1549 State(store): State<SharedStore>,
1550 Path(pipeline_id): Path<uuid::Uuid>,
1551) -> Result<Json<serde_json::Value>> {
1552 let pipeline_manager = store.pipeline_manager();
1553
1554 let pipeline = pipeline_manager.get(pipeline_id).ok_or_else(|| {
1555 crate::error::AllSourceError::ValidationError(format!("Pipeline not found: {pipeline_id}"))
1556 })?;
1557
1558 pipeline.reset();
1559
1560 tracing::info!("🔄 Reset pipeline {}", pipeline_id);
1561
1562 Ok(Json(serde_json::json!({
1563 "pipeline_id": pipeline_id,
1564 "reset": true
1565 })))
1566}
1567
1568pub async fn get_event_by_id(
1574 State(store): State<SharedStore>,
1575 Path(event_id): Path<uuid::Uuid>,
1576) -> Result<Json<serde_json::Value>> {
1577 let event = store.get_event_by_id(&event_id)?.ok_or_else(|| {
1578 crate::error::AllSourceError::EntityNotFound(format!("Event '{event_id}' not found"))
1579 })?;
1580
1581 let dto = EventDto::from(&event);
1582
1583 tracing::debug!("Event retrieved by ID: {}", event_id);
1584
1585 Ok(Json(serde_json::json!({
1586 "event": dto,
1587 "found": true
1588 })))
1589}
1590
1591pub async fn list_projections(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1597 let projection_manager = store.projection_manager();
1598 let status_map = store.projection_status();
1599
1600 let projections: Vec<serde_json::Value> = projection_manager
1601 .list_projections()
1602 .iter()
1603 .map(|(name, projection)| {
1604 let status = status_map
1605 .get(name)
1606 .map_or_else(|| "running".to_string(), |s| s.value().clone());
1607 serde_json::json!({
1608 "name": name,
1609 "type": format!("{:?}", projection.name()),
1610 "status": status,
1611 })
1612 })
1613 .collect();
1614
1615 tracing::debug!("Listed {} projections", projections.len());
1616
1617 Json(serde_json::json!({
1618 "projections": projections,
1619 "total": projections.len()
1620 }))
1621}
1622
1623pub async fn get_projection(
1625 State(store): State<SharedStore>,
1626 Path(name): Path<String>,
1627) -> Result<Json<serde_json::Value>> {
1628 let projection_manager = store.projection_manager();
1629
1630 let projection = projection_manager.get_projection(&name).ok_or_else(|| {
1631 crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1632 })?;
1633
1634 Ok(Json(serde_json::json!({
1635 "name": projection.name(),
1636 "found": true
1637 })))
1638}
1639
1640pub async fn get_projection_state(
1653 State(store): State<SharedStore>,
1654 Path((name, entity_id)): Path<(String, String)>,
1655) -> Result<Json<serde_json::Value>> {
1656 let state = store
1657 .projection_manager()
1658 .get_projection(&name)
1659 .and_then(|p| p.get_state(&entity_id))
1660 .or_else(|| {
1661 store
1662 .projection_state_cache()
1663 .get(&format!("{name}:{entity_id}"))
1664 .map(|entry| entry.value().clone())
1665 });
1666
1667 tracing::debug!("Projection state retrieved: {} / {}", name, entity_id);
1668
1669 Ok(Json(serde_json::json!({
1670 "projection": name,
1671 "entity_id": entity_id,
1672 "state": state,
1673 "found": state.is_some()
1674 })))
1675}
1676
1677pub async fn delete_projection(
1682 State(store): State<SharedStore>,
1683 Path(name): Path<String>,
1684) -> Result<Json<serde_json::Value>> {
1685 let projection_manager = store.projection_manager();
1686
1687 let projection = projection_manager.get_projection(&name).ok_or_else(|| {
1688 crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1689 })?;
1690
1691 projection.clear();
1692
1693 let cache = store.projection_state_cache();
1695 let prefix = format!("{name}:");
1696 let keys_to_remove: Vec<String> = cache
1697 .iter()
1698 .filter(|entry| entry.key().starts_with(&prefix))
1699 .map(|entry| entry.key().clone())
1700 .collect();
1701 for key in keys_to_remove {
1702 cache.remove(&key);
1703 }
1704
1705 tracing::info!("Projection deleted (cleared): {}", name);
1706
1707 Ok(Json(serde_json::json!({
1708 "projection": name,
1709 "deleted": true
1710 })))
1711}
1712
1713#[derive(Debug, Default, Deserialize)]
1720pub struct ProjectionStateSummaryParams {
1721 pub limit: Option<usize>,
1723 pub offset: Option<usize>,
1725 pub entity_id_prefix: Option<String>,
1728}
1729
1730pub async fn get_projection_state_summary(
1746 State(store): State<SharedStore>,
1747 Path(name): Path<String>,
1748 Query(params): Query<ProjectionStateSummaryParams>,
1749) -> Result<Json<serde_json::Value>> {
1750 let cache = store.projection_state_cache();
1751 let prefix = format!("{name}:");
1752 let offset = params.offset.unwrap_or(0);
1753
1754 let mut entity_ids: Vec<String> = cache
1759 .iter()
1760 .filter_map(|entry| entry.key().strip_prefix(&prefix).map(ToString::to_string))
1761 .filter(|entity_id| {
1762 params
1763 .entity_id_prefix
1764 .as_ref()
1765 .is_none_or(|p| entity_id.starts_with(p))
1766 })
1767 .collect();
1768 entity_ids.sort_unstable();
1769
1770 let total = entity_ids.len();
1771
1772 let page = entity_ids.into_iter().skip(offset);
1773 let page: Vec<String> = match params.limit {
1774 Some(limit) => page.take(limit).collect(),
1775 None => page.collect(),
1776 };
1777
1778 let states: Vec<serde_json::Value> = page
1779 .into_iter()
1780 .filter_map(|entity_id| {
1781 cache.get(&format!("{prefix}{entity_id}")).map(|entry| {
1783 serde_json::json!({
1784 "entity_id": entity_id,
1785 "state": entry.value().clone()
1786 })
1787 })
1788 })
1789 .collect();
1790
1791 let count = states.len();
1792 let has_more = offset + count < total;
1795
1796 tracing::debug!(
1797 "Projection state summary: {} ({} of {} entities, offset {})",
1798 name,
1799 count,
1800 total,
1801 offset
1802 );
1803
1804 Ok(Json(serde_json::json!({
1805 "projection": name,
1806 "states": states,
1807 "count": count,
1808 "total": total,
1809 "has_more": has_more
1810 })))
1811}
1812
1813pub async fn reset_projection(
1817 State(store): State<SharedStore>,
1818 Path(name): Path<String>,
1819) -> Result<Json<serde_json::Value>> {
1820 let reprocessed = store.reset_projection(&name)?;
1821
1822 tracing::info!(
1823 "Projection reset: {} ({} events reprocessed)",
1824 name,
1825 reprocessed
1826 );
1827
1828 Ok(Json(serde_json::json!({
1829 "projection": name,
1830 "reset": true,
1831 "events_reprocessed": reprocessed
1832 })))
1833}
1834
1835pub async fn pause_projection(
1839 State(store): State<SharedStore>,
1840 Path(name): Path<String>,
1841) -> Result<Json<serde_json::Value>> {
1842 let projection_manager = store.projection_manager();
1843
1844 let _projection = projection_manager.get_projection(&name).ok_or_else(|| {
1846 crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1847 })?;
1848
1849 store
1850 .projection_status()
1851 .insert(name.clone(), "paused".to_string());
1852
1853 tracing::info!("Projection paused: {}", name);
1854
1855 Ok(Json(serde_json::json!({
1856 "projection": name,
1857 "status": "paused"
1858 })))
1859}
1860
1861pub async fn start_projection(
1865 State(store): State<SharedStore>,
1866 Path(name): Path<String>,
1867) -> Result<Json<serde_json::Value>> {
1868 let projection_manager = store.projection_manager();
1869
1870 let _projection = projection_manager.get_projection(&name).ok_or_else(|| {
1872 crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1873 })?;
1874
1875 store
1876 .projection_status()
1877 .insert(name.clone(), "running".to_string());
1878
1879 tracing::info!("Projection started: {}", name);
1880
1881 Ok(Json(serde_json::json!({
1882 "projection": name,
1883 "status": "running"
1884 })))
1885}
1886
1887#[derive(Debug, Deserialize)]
1889pub struct SaveProjectionStateRequest {
1890 pub state: serde_json::Value,
1891}
1892
1893pub async fn save_projection_state(
1898 State(store): State<SharedStore>,
1899 Path((name, entity_id)): Path<(String, String)>,
1900 Json(req): Json<SaveProjectionStateRequest>,
1901) -> Result<Json<serde_json::Value>> {
1902 let projection_cache = store.projection_state_cache();
1903
1904 projection_cache.insert(format!("{name}:{entity_id}"), req.state.clone());
1906
1907 tracing::info!("Projection state saved: {} / {}", name, entity_id);
1908
1909 Ok(Json(serde_json::json!({
1910 "projection": name,
1911 "entity_id": entity_id,
1912 "saved": true
1913 })))
1914}
1915
1916#[derive(Debug, Deserialize)]
1920pub struct BulkGetStateRequest {
1921 pub entity_ids: Vec<String>,
1922}
1923
1924#[derive(Debug, Deserialize)]
1928pub struct BulkSaveStateRequest {
1929 pub states: Vec<BulkSaveStateItem>,
1930}
1931
1932#[derive(Debug, Deserialize)]
1933pub struct BulkSaveStateItem {
1934 pub entity_id: String,
1935 pub state: serde_json::Value,
1936}
1937
1938pub async fn bulk_get_projection_states(
1939 State(store): State<SharedStore>,
1940 Path(name): Path<String>,
1941 Json(req): Json<BulkGetStateRequest>,
1942) -> Result<Json<serde_json::Value>> {
1943 let projection = store.projection_manager().get_projection(&name);
1947 let cache = store.projection_state_cache();
1948
1949 let states: Vec<serde_json::Value> = req
1950 .entity_ids
1951 .iter()
1952 .map(|entity_id| {
1953 let state = projection
1954 .as_ref()
1955 .and_then(|p| p.get_state(entity_id))
1956 .or_else(|| {
1957 cache
1958 .get(&format!("{name}:{entity_id}"))
1959 .map(|entry| entry.value().clone())
1960 });
1961 serde_json::json!({
1962 "entity_id": entity_id,
1963 "state": state,
1964 "found": state.is_some()
1965 })
1966 })
1967 .collect();
1968
1969 tracing::debug!(
1970 "Bulk projection state retrieved: {} entities from {}",
1971 states.len(),
1972 name
1973 );
1974
1975 Ok(Json(serde_json::json!({
1976 "projection": name,
1977 "states": states,
1978 "total": states.len()
1979 })))
1980}
1981
1982pub async fn bulk_save_projection_states(
1987 State(store): State<SharedStore>,
1988 Path(name): Path<String>,
1989 Json(req): Json<BulkSaveStateRequest>,
1990) -> Result<Json<serde_json::Value>> {
1991 let projection_cache = store.projection_state_cache();
1992
1993 let mut saved_count = 0;
1994 for item in &req.states {
1995 projection_cache.insert(format!("{name}:{}", item.entity_id), item.state.clone());
1996 saved_count += 1;
1997 }
1998
1999 tracing::info!(
2000 "Bulk projection state saved: {} entities for {}",
2001 saved_count,
2002 name
2003 );
2004
2005 Ok(Json(serde_json::json!({
2006 "projection": name,
2007 "saved": saved_count,
2008 "total": req.states.len()
2009 })))
2010}
2011
2012#[derive(Debug, Deserialize)]
2018pub struct ListWebhooksParams {
2019 pub tenant_id: Option<String>,
2020}
2021
2022pub async fn register_webhook(
2024 State(store): State<SharedStore>,
2025 Json(req): Json<RegisterWebhookRequest>,
2026) -> Json<serde_json::Value> {
2027 let registry = store.webhook_registry();
2028 let webhook = registry.register(req);
2029
2030 tracing::info!("Webhook registered: {} -> {}", webhook.id, webhook.url);
2031
2032 Json(serde_json::json!({
2033 "webhook": webhook,
2034 "created": true
2035 }))
2036}
2037
2038pub async fn list_webhooks(
2040 State(store): State<SharedStore>,
2041 Query(params): Query<ListWebhooksParams>,
2042) -> Json<serde_json::Value> {
2043 let registry = store.webhook_registry();
2044
2045 let webhooks = if let Some(tenant_id) = params.tenant_id {
2046 registry.list_by_tenant(&tenant_id)
2047 } else {
2048 vec![]
2050 };
2051
2052 let total = webhooks.len();
2053
2054 Json(serde_json::json!({
2055 "webhooks": webhooks,
2056 "total": total
2057 }))
2058}
2059
2060pub async fn get_webhook(
2062 State(store): State<SharedStore>,
2063 Path(webhook_id): Path<uuid::Uuid>,
2064) -> Result<Json<serde_json::Value>> {
2065 let registry = store.webhook_registry();
2066
2067 let webhook = registry.get(webhook_id).ok_or_else(|| {
2068 crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2069 })?;
2070
2071 Ok(Json(serde_json::json!({
2072 "webhook": webhook,
2073 "found": true
2074 })))
2075}
2076
2077pub async fn update_webhook(
2079 State(store): State<SharedStore>,
2080 Path(webhook_id): Path<uuid::Uuid>,
2081 Json(req): Json<UpdateWebhookRequest>,
2082) -> Result<Json<serde_json::Value>> {
2083 let registry = store.webhook_registry();
2084
2085 let webhook = registry.update(webhook_id, req).ok_or_else(|| {
2086 crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2087 })?;
2088
2089 tracing::info!("Webhook updated: {}", webhook_id);
2090
2091 Ok(Json(serde_json::json!({
2092 "webhook": webhook,
2093 "updated": true
2094 })))
2095}
2096
2097pub async fn delete_webhook(
2099 State(store): State<SharedStore>,
2100 Path(webhook_id): Path<uuid::Uuid>,
2101) -> Result<Json<serde_json::Value>> {
2102 let registry = store.webhook_registry();
2103
2104 let webhook = registry.delete(webhook_id).ok_or_else(|| {
2105 crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2106 })?;
2107
2108 tracing::info!("Webhook deleted: {} ({})", webhook_id, webhook.url);
2109
2110 Ok(Json(serde_json::json!({
2111 "webhook_id": webhook_id,
2112 "deleted": true
2113 })))
2114}
2115
2116#[derive(Debug, Deserialize)]
2118pub struct ListDeliveriesParams {
2119 pub limit: Option<usize>,
2120}
2121
2122pub async fn list_webhook_deliveries(
2124 State(store): State<SharedStore>,
2125 Path(webhook_id): Path<uuid::Uuid>,
2126 Query(params): Query<ListDeliveriesParams>,
2127) -> Result<Json<serde_json::Value>> {
2128 let registry = store.webhook_registry();
2129
2130 registry.get(webhook_id).ok_or_else(|| {
2132 crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2133 })?;
2134
2135 let limit = params.limit.unwrap_or(50);
2136 let deliveries = registry.get_deliveries(webhook_id, limit);
2137 let total = deliveries.len();
2138
2139 Ok(Json(serde_json::json!({
2140 "webhook_id": webhook_id,
2141 "deliveries": deliveries,
2142 "total": total
2143 })))
2144}
2145
2146#[cfg(feature = "analytics")]
2152pub async fn eventql_query(
2153 State(store): State<SharedStore>,
2154 Json(req): Json<crate::infrastructure::query::eventql::EventQLRequest>,
2155) -> Result<Json<serde_json::Value>> {
2156 let events = store.snapshot_events();
2157 match crate::infrastructure::query::eventql::execute_eventql(&events, &req).await {
2158 Ok(response) => Ok(Json(serde_json::json!({
2159 "columns": response.columns,
2160 "rows": response.rows,
2161 "row_count": response.row_count,
2162 }))),
2163 Err(e) => Err(crate::error::AllSourceError::InvalidQuery(e)),
2164 }
2165}
2166
2167pub async fn graphql_query(
2169 State(store): State<SharedStore>,
2170 Json(req): Json<GraphQLRequest>,
2171) -> Json<serde_json::Value> {
2172 let fields = match crate::infrastructure::query::graphql::parse_query(&req.query) {
2173 Ok(f) => f,
2174 Err(e) => {
2175 return Json(
2176 serde_json::to_value(GraphQLResponse {
2177 data: None,
2178 errors: vec![GraphQLError { message: e }],
2179 })
2180 .unwrap(),
2181 );
2182 }
2183 };
2184
2185 let mut data = serde_json::Map::new();
2186 let mut errors = Vec::new();
2187
2188 for field in &fields {
2189 match field.name.as_str() {
2190 "events" => {
2191 let request = crate::application::dto::QueryEventsRequest {
2192 entity_id: field.arguments.get("entity_id").cloned(),
2193 event_type: field.arguments.get("event_type").cloned(),
2194 tenant_id: field.arguments.get("tenant_id").cloned(),
2195 limit: field.arguments.get("limit").and_then(|l| l.parse().ok()),
2196 as_of: None,
2197 since: None,
2198 until: None,
2199 event_type_prefix: None,
2200 exclude_event_type_prefix: None,
2201 payload_filter: None,
2202 };
2203 match store.query(&request) {
2204 Ok(events) => {
2205 let json_events: Vec<serde_json::Value> = events
2206 .iter()
2207 .map(|e| {
2208 crate::infrastructure::query::graphql::event_to_json(
2209 e,
2210 &field.fields,
2211 )
2212 })
2213 .collect();
2214 data.insert("events".to_string(), serde_json::Value::Array(json_events));
2215 }
2216 Err(e) => errors.push(GraphQLError {
2217 message: format!("events query failed: {e}"),
2218 }),
2219 }
2220 }
2221 "event" => {
2222 if let Some(id_str) = field.arguments.get("id") {
2223 if let Ok(id) = uuid::Uuid::parse_str(id_str) {
2224 match store.get_event_by_id(&id) {
2225 Ok(Some(event)) => {
2226 data.insert(
2227 "event".to_string(),
2228 crate::infrastructure::query::graphql::event_to_json(
2229 &event,
2230 &field.fields,
2231 ),
2232 );
2233 }
2234 Ok(None) => {
2235 data.insert("event".to_string(), serde_json::Value::Null);
2236 }
2237 Err(e) => errors.push(GraphQLError {
2238 message: format!("event lookup failed: {e}"),
2239 }),
2240 }
2241 } else {
2242 errors.push(GraphQLError {
2243 message: format!("Invalid UUID: {id_str}"),
2244 });
2245 }
2246 } else {
2247 errors.push(GraphQLError {
2248 message: "event query requires 'id' argument".to_string(),
2249 });
2250 }
2251 }
2252 "projections" => {
2253 let pm = store.projection_manager();
2254 let names: Vec<serde_json::Value> = pm
2255 .list_projections()
2256 .iter()
2257 .map(|(name, _)| serde_json::Value::String(name.clone()))
2258 .collect();
2259 data.insert("projections".to_string(), serde_json::Value::Array(names));
2260 }
2261 "stats" => {
2262 let stats = store.stats();
2263 data.insert(
2264 "stats".to_string(),
2265 serde_json::json!({
2266 "total_events": stats.total_events,
2267 "total_entities": stats.total_entities,
2268 "total_event_types": stats.total_event_types,
2269 }),
2270 );
2271 }
2272 "__schema" => {
2273 data.insert(
2274 "__schema".to_string(),
2275 crate::infrastructure::query::graphql::introspection_schema(),
2276 );
2277 }
2278 other => {
2279 errors.push(GraphQLError {
2280 message: format!("Unknown field: {other}"),
2281 });
2282 }
2283 }
2284 }
2285
2286 Json(
2287 serde_json::to_value(GraphQLResponse {
2288 data: Some(serde_json::Value::Object(data)),
2289 errors,
2290 })
2291 .unwrap(),
2292 )
2293}
2294
2295pub async fn geo_query(
2297 State(store): State<SharedStore>,
2298 Json(req): Json<GeoQueryRequest>,
2299) -> Json<serde_json::Value> {
2300 let events = store.snapshot_events();
2301 let geo_index = store.geo_index();
2302 let results =
2303 crate::infrastructure::query::geospatial::execute_geo_query(&events, &geo_index, &req);
2304 let total = results.len();
2305 Json(serde_json::json!({
2306 "results": results,
2307 "total": total,
2308 }))
2309}
2310
2311pub async fn geo_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
2313 let stats = store.geo_index().stats();
2314 Json(serde_json::json!(stats))
2315}
2316
2317pub async fn exactly_once_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
2319 let stats = store.exactly_once().stats();
2320 Json(serde_json::json!(stats))
2321}
2322
2323pub async fn schema_evolution_history(
2325 State(store): State<SharedStore>,
2326 Path(event_type): Path<String>,
2327) -> Json<serde_json::Value> {
2328 let mgr = store.schema_evolution();
2329 let history = mgr.get_history(&event_type);
2330 let version = mgr.get_version(&event_type);
2331 Json(serde_json::json!({
2332 "event_type": event_type,
2333 "current_version": version,
2334 "history": history,
2335 }))
2336}
2337
2338pub async fn schema_evolution_schema(
2340 State(store): State<SharedStore>,
2341 Path(event_type): Path<String>,
2342) -> Json<serde_json::Value> {
2343 let mgr = store.schema_evolution();
2344 if let Some(schema) = mgr.get_schema(&event_type) {
2345 let json_schema = crate::application::services::schema_evolution::to_json_schema(&schema);
2346 Json(serde_json::json!({
2347 "event_type": event_type,
2348 "version": mgr.get_version(&event_type),
2349 "inferred_schema": schema,
2350 "json_schema": json_schema,
2351 }))
2352 } else {
2353 Json(serde_json::json!({
2354 "event_type": event_type,
2355 "error": "No schema inferred for this event type"
2356 }))
2357 }
2358}
2359
2360pub async fn schema_evolution_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
2362 let stats = store.schema_evolution().stats();
2363 let event_types = store.schema_evolution().list_event_types();
2364 Json(serde_json::json!({
2365 "stats": stats,
2366 "tracked_event_types": event_types,
2367 }))
2368}
2369
2370#[cfg(feature = "embedded-sync")]
2376pub async fn sync_pull_handler(
2377 State(state): State<AppState>,
2378 Json(request): Json<crate::embedded::sync_types::SyncPullRequest>,
2379) -> Result<Json<crate::embedded::sync_types::SyncPullResponse>> {
2380 use crate::infrastructure::cluster::{crdt::ReplicatedEvent, hlc::HlcTimestamp};
2381
2382 let store = &state.store;
2383
2384 let since = request
2387 .version_vector
2388 .values()
2389 .map(|ts| ts.physical_ms)
2390 .min()
2391 .and_then(|ms| chrono::DateTime::from_timestamp_millis(ms as i64));
2392
2393 let events = store.query(&crate::application::dto::QueryEventsRequest {
2394 entity_id: None,
2395 event_type: None,
2396 tenant_id: None,
2397 as_of: None,
2398 since,
2399 until: None,
2400 limit: None,
2401 event_type_prefix: None,
2402 exclude_event_type_prefix: None,
2403 payload_filter: None,
2404 })?;
2405
2406 let mut replicated = Vec::with_capacity(events.len());
2408 let mut last_ms = 0u64;
2409 let mut logical = 0u32;
2410
2411 for event in &events {
2412 let event_ms = event.timestamp().timestamp_millis() as u64;
2413 if event_ms == last_ms {
2414 logical += 1;
2415 } else {
2416 last_ms = event_ms;
2417 logical = 0;
2418 }
2419
2420 replicated.push(ReplicatedEvent {
2421 event_id: event.id().to_string(),
2422 hlc_timestamp: HlcTimestamp::new(event_ms, logical, 0),
2423 origin_region: "server".to_string(),
2424 event_data: serde_json::json!({
2425 "event_type": event.event_type_str(),
2426 "entity_id": event.entity_id_str(),
2427 "tenant_id": event.tenant_id_str(),
2428 "payload": event.payload,
2429 "metadata": event.metadata,
2430 }),
2431 });
2432 }
2433
2434 Ok(Json(crate::embedded::sync_types::SyncPullResponse {
2435 events: replicated,
2436 version_vector: std::collections::BTreeMap::new(),
2437 }))
2438}
2439
2440#[cfg(feature = "embedded-sync")]
2442pub async fn sync_push_handler(
2443 State(state): State<AppState>,
2444 Json(request): Json<crate::embedded::sync_types::SyncPushRequest>,
2445) -> Result<Json<crate::embedded::sync_types::SyncPushResponse>> {
2446 let store = &state.store;
2447
2448 let mut accepted = 0usize;
2449 let mut skipped = 0usize;
2450
2451 for rep_event in &request.events {
2452 let event_data = &rep_event.event_data;
2453 let event_type = event_data
2454 .get("event_type")
2455 .and_then(|v| v.as_str())
2456 .unwrap_or("unknown")
2457 .to_string();
2458 let entity_id = event_data
2459 .get("entity_id")
2460 .and_then(|v| v.as_str())
2461 .unwrap_or("unknown")
2462 .to_string();
2463 let tenant_id = event_data
2464 .get("tenant_id")
2465 .and_then(|v| v.as_str())
2466 .unwrap_or("default")
2467 .to_string();
2468 let payload = event_data
2469 .get("payload")
2470 .cloned()
2471 .unwrap_or(serde_json::json!({}));
2472 let metadata = event_data.get("metadata").cloned();
2473
2474 match Event::from_strings(event_type, entity_id, tenant_id, payload, metadata) {
2475 Ok(domain_event) => {
2476 store.ingest(&domain_event)?;
2477 accepted += 1;
2478 }
2479 Err(_) => {
2480 skipped += 1;
2481 }
2482 }
2483 }
2484
2485 Ok(Json(crate::embedded::sync_types::SyncPushResponse {
2486 accepted,
2487 skipped,
2488 version_vector: std::collections::BTreeMap::new(),
2489 }))
2490}
2491
2492pub async fn register_consumer(
2498 State(store): State<SharedStore>,
2499 Json(req): Json<RegisterConsumerRequest>,
2500) -> Result<Json<ConsumerResponse>> {
2501 let consumer = store
2502 .consumer_registry()
2503 .register(&req.consumer_id, &req.event_type_filters);
2504
2505 Ok(Json(ConsumerResponse {
2506 consumer_id: consumer.consumer_id,
2507 event_type_filters: consumer.event_type_filters,
2508 cursor_position: consumer.cursor_position,
2509 }))
2510}
2511
2512pub async fn get_consumer(
2514 State(store): State<SharedStore>,
2515 Path(consumer_id): Path<String>,
2516) -> Result<Json<ConsumerResponse>> {
2517 let consumer = store.consumer_registry().get_or_create(&consumer_id);
2518
2519 Ok(Json(ConsumerResponse {
2520 consumer_id: consumer.consumer_id,
2521 event_type_filters: consumer.event_type_filters,
2522 cursor_position: consumer.cursor_position,
2523 }))
2524}
2525
2526#[derive(Debug, Deserialize)]
2528pub struct ConsumerPollQuery {
2529 pub limit: Option<usize>,
2530}
2531
2532pub async fn poll_consumer_events(
2533 State(store): State<SharedStore>,
2534 Path(consumer_id): Path<String>,
2535 Query(query): Query<ConsumerPollQuery>,
2536) -> Result<Json<ConsumerEventsResponse>> {
2537 let consumer = store.consumer_registry().get_or_create(&consumer_id);
2538 let offset = consumer.cursor_position.unwrap_or(0);
2539 let limit = query.limit.unwrap_or(100);
2540
2541 let events = store.events_after_offset(offset, &consumer.event_type_filters, limit);
2542 let count = events.len();
2543
2544 let consumer_events: Vec<ConsumerEventDto> = events
2545 .into_iter()
2546 .map(|(position, event)| ConsumerEventDto {
2547 position,
2548 event: EventDto::from(&event),
2549 })
2550 .collect();
2551
2552 Ok(Json(ConsumerEventsResponse {
2553 events: consumer_events,
2554 count,
2555 }))
2556}
2557
2558pub async fn ack_consumer(
2560 State(store): State<SharedStore>,
2561 Path(consumer_id): Path<String>,
2562 Json(req): Json<AckRequest>,
2563) -> Result<Json<serde_json::Value>> {
2564 let max_offset = store.total_events() as u64;
2565
2566 store
2567 .consumer_registry()
2568 .ack(&consumer_id, req.position, max_offset)
2569 .map_err(crate::error::AllSourceError::InvalidInput)?;
2570
2571 Ok(Json(serde_json::json!({
2572 "status": "ok",
2573 "consumer_id": consumer_id,
2574 "position": req.position,
2575 })))
2576}
2577
2578#[cfg(test)]
2579mod tests {
2580 use super::*;
2581 use crate::{domain::entities::Event, store::EventStore};
2582
2583 fn create_test_store() -> Arc<EventStore> {
2584 Arc::new(EventStore::new())
2585 }
2586
2587 async fn query_page(store: &SharedStore, query: &str) -> QueryEventsResponse {
2596 use axum::extract::{Query, State};
2597
2598 let uri: axum::http::Uri = format!("/api/v1/events/query?tenant_id=test-stream&{query}")
2599 .parse()
2600 .unwrap();
2601 query_events(
2602 OptionalAuth(None),
2603 Query::try_from_uri(&uri).unwrap(),
2604 Query::try_from_uri(&uri).unwrap(),
2605 Query::try_from_uri(&uri).unwrap(),
2606 Query::try_from_uri(&uri).unwrap(),
2607 State(store.clone()),
2608 )
2609 .await
2610 .unwrap()
2611 .0
2612 }
2613
2614 fn create_test_event(entity_id: &str, event_type: &str) -> Event {
2615 Event::from_strings(
2616 event_type.to_string(),
2617 entity_id.to_string(),
2618 "test-stream".to_string(),
2619 serde_json::json!({
2620 "name": "Test",
2621 "value": 42
2622 }),
2623 None,
2624 )
2625 .unwrap()
2626 }
2627
2628 #[tokio::test]
2629 async fn test_query_events_has_more_and_total_count() {
2630 let store = create_test_store();
2631
2632 for i in 0..50 {
2634 store
2635 .ingest(&create_test_event(&format!("entity-{i}"), "user.created"))
2636 .unwrap();
2637 }
2638
2639 let req = QueryEventsRequest {
2641 entity_id: None,
2642 event_type: None,
2643 tenant_id: None,
2644 as_of: None,
2645 since: None,
2646 until: None,
2647 limit: Some(10),
2648 event_type_prefix: None,
2649 exclude_event_type_prefix: None,
2650 payload_filter: None,
2651 };
2652
2653 let requested_limit = req.limit;
2654 let unlimited_req = QueryEventsRequest {
2655 limit: None,
2656 ..QueryEventsRequest {
2657 entity_id: req.entity_id,
2658 event_type: req.event_type,
2659 tenant_id: req.tenant_id,
2660 as_of: req.as_of,
2661 since: req.since,
2662 until: req.until,
2663 limit: None,
2664 event_type_prefix: req.event_type_prefix,
2665 exclude_event_type_prefix: None,
2666 payload_filter: req.payload_filter,
2667 }
2668 };
2669 let all_events = store.query(&unlimited_req).unwrap();
2670 let total_count = all_events.len();
2671 let limited_events: Vec<Event> = if let Some(limit) = requested_limit {
2672 all_events.into_iter().take(limit).collect()
2673 } else {
2674 all_events
2675 };
2676 let count = limited_events.len();
2677 let has_more = count < total_count;
2678
2679 assert_eq!(count, 10);
2680 assert_eq!(total_count, 50);
2681 assert!(has_more);
2682 }
2683
2684 #[tokio::test]
2693 async fn query_events_honours_offset_pagination() {
2694 use axum::extract::{Query, State};
2695
2696 let store = create_test_store();
2697 for i in 0..25 {
2698 store
2699 .ingest(&create_test_event(&format!("e-{i:02}"), "user.created"))
2700 .unwrap();
2701 }
2702
2703 async fn page(store: &SharedStore, limit: usize, offset: usize) -> QueryEventsResponse {
2705 let uri: axum::http::Uri =
2706 format!("/api/v1/events/query?tenant_id=test-stream&limit={limit}&offset={offset}")
2707 .parse()
2708 .unwrap();
2709 let req: Query<QueryEventsRequest> = Query::try_from_uri(&uri).unwrap();
2710 let order: Query<EventOrderParam> = Query::try_from_uri(&uri).unwrap();
2711 let off: Query<EventOffsetParam> = Query::try_from_uri(&uri).unwrap();
2712 assert_eq!(off.0.offset, Some(offset), "offset must deserialize");
2713 query_events(
2714 OptionalAuth(None),
2715 req,
2716 order,
2717 off,
2718 Query::try_from_uri(&uri).unwrap(),
2719 State(store.clone()),
2720 )
2721 .await
2722 .unwrap()
2723 .0
2724 }
2725
2726 let p1 = page(&store, 10, 0).await;
2727 let p2 = page(&store, 10, 10).await;
2728 let p3 = page(&store, 10, 20).await;
2729
2730 assert_eq!(p1.count, 10);
2731 assert_eq!(p2.count, 10);
2732 assert_eq!(p3.count, 5, "last page returns the remainder");
2733 assert_eq!(p1.total_count, 25);
2734
2735 let ids = |r: &QueryEventsResponse| -> Vec<String> {
2737 r.events.iter().map(|e| e.entity_id.clone()).collect()
2738 };
2739 assert_ne!(ids(&p1), ids(&p2), "offset=10 must skip the first page");
2740
2741 let mut all = ids(&p1);
2742 all.extend(ids(&p2));
2743 all.extend(ids(&p3));
2744 let unique: std::collections::HashSet<_> = all.iter().cloned().collect();
2745 assert_eq!(
2746 unique.len(),
2747 25,
2748 "paging the whole set must yield 25 distinct entities, not duplicates"
2749 );
2750
2751 assert!(p1.has_more, "25 events, page 1 of 10 → more remain");
2754 assert!(p2.has_more, "25 events, page 2 of 10 → more remain");
2755 assert!(!p3.has_more, "offset=20 + count=5 == total → exhausted");
2756
2757 let past = page(&store, 10, 100).await;
2759 assert_eq!(past.count, 0);
2760 assert!(!past.has_more, "offset beyond the match set is exhausted");
2761 }
2762
2763 #[tokio::test]
2772 async fn projection_state_summary_honours_limit_offset_and_prefix() {
2773 use axum::{
2774 body::{Body, to_bytes},
2775 http::Request,
2776 };
2777 use tower::ServiceExt; let store = create_test_store();
2780 let cache = store.projection_state_cache();
2781 for i in 0..25 {
2782 cache.insert(
2783 format!("demo:tenant-{i:02}"),
2784 serde_json::json!({ "n": i as u64 }),
2785 );
2786 }
2787 cache.insert(
2789 "other:tenant-99".to_string(),
2790 serde_json::json!({ "n": 99 }),
2791 );
2792
2793 let app = Router::new()
2794 .route(
2795 "/api/v1/projections/{name}/state",
2796 get(get_projection_state_summary),
2797 )
2798 .with_state(store.clone());
2799
2800 async fn fetch(app: &Router, uri: &str) -> serde_json::Value {
2801 let resp = app
2802 .clone()
2803 .oneshot(Request::builder().uri(uri).body(Body::empty()).unwrap())
2804 .await
2805 .unwrap();
2806 assert_eq!(resp.status(), axum::http::StatusCode::OK, "GET {uri}");
2807 let bytes = to_bytes(resp.into_body(), usize::MAX).await.unwrap();
2808 serde_json::from_slice(&bytes).unwrap()
2809 }
2810
2811 let ids = |body: &serde_json::Value| -> Vec<String> {
2812 body["states"]
2813 .as_array()
2814 .unwrap()
2815 .iter()
2816 .map(|s| s["entity_id"].as_str().unwrap().to_string())
2817 .collect()
2818 };
2819
2820 let all = fetch(&app, "/api/v1/projections/demo/state").await;
2822 assert_eq!(all["total"], 25);
2823 assert_eq!(ids(&all).len(), 25);
2824
2825 let p1 = fetch(&app, "/api/v1/projections/demo/state?limit=10").await;
2827 assert_eq!(ids(&p1).len(), 10, "limit must bound the response");
2828 assert_eq!(p1["total"], 25);
2829 assert_eq!(p1["count"], 10);
2830 assert_eq!(p1["has_more"], true);
2831
2832 let p2 = fetch(&app, "/api/v1/projections/demo/state?limit=10&offset=10").await;
2834 let p3 = fetch(&app, "/api/v1/projections/demo/state?limit=10&offset=20").await;
2835 assert_eq!(ids(&p3).len(), 5, "last page returns the remainder");
2836 assert_eq!(p3["has_more"], false, "offset + count == total → exhausted");
2837 assert_ne!(ids(&p1), ids(&p2), "offset=10 must skip the first page");
2838
2839 let mut walked = ids(&p1);
2840 walked.extend(ids(&p2));
2841 walked.extend(ids(&p3));
2842 let unique: std::collections::HashSet<_> = walked.iter().cloned().collect();
2843 assert_eq!(
2844 unique.len(),
2845 25,
2846 "paging the whole projection must yield 25 distinct entities"
2847 );
2848
2849 let mut sorted = walked.clone();
2851 sorted.sort();
2852 assert_eq!(walked, sorted, "pages must be ordered by entity_id");
2853
2854 let past = fetch(&app, "/api/v1/projections/demo/state?limit=10&offset=100").await;
2856 assert_eq!(past["count"], 0);
2857 assert_eq!(past["has_more"], false);
2858 assert_eq!(past["total"], 25);
2859
2860 let shard = fetch(
2862 &app,
2863 "/api/v1/projections/demo/state?entity_id_prefix=tenant-1",
2864 )
2865 .await;
2866 assert_eq!(shard["total"], 10, "tenant-10..tenant-19");
2867 assert!(
2868 ids(&shard).iter().all(|id| id.starts_with("tenant-1")),
2869 "entity_id_prefix must filter"
2870 );
2871
2872 let other = fetch(&app, "/api/v1/projections/other/state?limit=10").await;
2874 assert_eq!(ids(&other), vec!["tenant-99".to_string()]);
2875 }
2876
2877 #[test]
2890 fn query_events_limit_does_not_materialize_whole_history() {
2891 const HISTORY: usize = 200;
2892
2893 let store = create_test_store();
2894 for _ in 0..HISTORY {
2895 store
2896 .ingest(&create_test_event("entity-hot", "user.updated"))
2897 .unwrap();
2898 }
2899
2900 let runtime = tokio::runtime::Builder::new_current_thread()
2904 .enable_all()
2905 .build()
2906 .unwrap();
2907 let (resp, materialized) = crate::clone_probe::measure(|| {
2908 runtime.block_on(query_page(
2909 &store,
2910 "entity_id=entity-hot&limit=1&order=desc",
2911 ))
2912 });
2913
2914 assert_eq!(resp.count, 1);
2917 assert_eq!(resp.events.len(), 1);
2918 assert_eq!(resp.total_count, HISTORY);
2919 assert!(resp.has_more);
2920
2921 assert_eq!(
2922 materialized, 1,
2923 "limit=1 materialized {materialized} events out of {HISTORY}: \
2924 `limit` must bound what a request materializes, not just what it \
2925 returns"
2926 );
2927 }
2928
2929 async fn query_page_result(
2932 store: &SharedStore,
2933 query: &str,
2934 ) -> Result<Json<QueryEventsResponse>> {
2935 use axum::extract::{Query, State};
2936
2937 let uri: axum::http::Uri = format!("/api/v1/events/query?tenant_id=test-stream&{query}")
2938 .parse()
2939 .unwrap();
2940 query_events(
2941 OptionalAuth(None),
2942 Query::try_from_uri(&uri).unwrap(),
2943 Query::try_from_uri(&uri).unwrap(),
2944 Query::try_from_uri(&uri).unwrap(),
2945 Query::try_from_uri(&uri).unwrap(),
2946 State(store.clone()),
2947 )
2948 .await
2949 }
2950
2951 #[tokio::test]
2960 async fn query_events_rejects_a_payload_filter_it_cannot_apply() {
2961 let store = create_test_store();
2962 for name in ["alice", "bob", "carol"] {
2963 let mut event = create_test_event(name, "user.created");
2964 event.payload = serde_json::json!({ "user_id": name });
2965 store.ingest(&event).unwrap();
2966 }
2967
2968 let ok = query_page(&store, "payload_filter=%7B%22user_id%22%3A%22alice%22%7D").await;
2971 assert_eq!(ok.count, 1, "a valid payload_filter must still work");
2972 assert_eq!(ok.total_count, 1);
2973
2974 for bad in [
2975 "not-json", "%7B%22user_id%22%3A%22alice", "%5B%22alice%22%5D", "42", ] {
2980 let result = query_page_result(&store, &format!("payload_filter={bad}")).await;
2981 let Err(err) = result else {
2982 let resp = result.unwrap().0;
2983 panic!(
2984 "payload_filter={bad} was silently ignored: returned {} of {} \
2985 events unfiltered instead of rejecting a filter the server \
2986 cannot apply",
2987 resp.count, resp.total_count
2988 );
2989 };
2990 assert!(
2991 matches!(err, crate::error::AllSourceError::InvalidInput(_)),
2992 "payload_filter={bad} must be a 400, got {err:?}"
2993 );
2994 }
2995 }
2996
2997 #[tokio::test]
3003 async fn query_events_rejects_an_unusable_order_value() {
3004 let store = create_test_store();
3005 store
3006 .ingest(&create_test_event("e-1", "user.created"))
3007 .unwrap();
3008
3009 for bad in ["descending", "DESCENDING", "newest", "1", "asc%20"] {
3010 let result = query_page_result(&store, &format!("order={bad}")).await;
3011 let Err(err) = result else {
3012 panic!("order={bad} must be rejected, not silently defaulted");
3013 };
3014 assert!(
3015 matches!(err, crate::error::AllSourceError::InvalidInput(_)),
3016 "order={bad} must be a 400, got {err:?}"
3017 );
3018 }
3019
3020 for good in ["asc", "ASC", "desc", "DeSc"] {
3022 let accepted = query_page_result(&store, &format!("order={good}"))
3023 .await
3024 .unwrap_or_else(|e| panic!("order={good} must be accepted: {e:?}"));
3025 assert_eq!(accepted.0.count, 1, "order={good}");
3026 }
3027 }
3028
3029 #[tokio::test]
3038 async fn query_events_exclude_prefix_applies_before_the_window() {
3039 const TYPES: [&str; 4] = [
3040 "audit.write",
3041 "user.created",
3042 "service.ping",
3043 "user.updated",
3044 ];
3045 let store = create_test_store();
3046 let base = chrono::Utc::now() - chrono::Duration::hours(24);
3047 let mut kept = Vec::new();
3048 for i in 0..12i64 {
3049 let mut event = create_test_event("org-1", TYPES[i as usize % 4]);
3050 event.timestamp = base + chrono::Duration::minutes(i);
3051 event.version = i + 1;
3052 if i % 2 == 1 {
3053 kept.push(event.id);
3054 }
3055 store.ingest(&event).unwrap();
3056 }
3057 assert_eq!(kept.len(), 6, "half the stream is user.*");
3058
3059 let page = query_page(
3063 &store,
3064 "exclude_event_type_prefix=audit.,%20service.&limit=3",
3065 )
3066 .await;
3067 assert_eq!(
3068 page.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3069 kept[..3].to_vec(),
3070 "limit=3 must yield 3 non-excluded events: exclusion runs before \
3071 the window, and the prefix list is comma-separated (whitespace \
3072 trimmed)"
3073 );
3074 assert!(
3075 page.events
3076 .iter()
3077 .all(|e| !e.event_type.starts_with("audit.")
3078 && !e.event_type.starts_with("service.")),
3079 "excluded namespaces must not appear: {:?}",
3080 page.events
3081 .iter()
3082 .map(|e| &e.event_type)
3083 .collect::<Vec<_>>()
3084 );
3085 assert_eq!(
3086 page.total_count, 6,
3087 "total_count is the post-exclusion match set, not the 12 ingested"
3088 );
3089 assert!(page.has_more);
3090
3091 let mut walked = Vec::new();
3093 for offset in [0, 3, 6] {
3094 let p = query_page(
3095 &store,
3096 &format!("exclude_event_type_prefix=audit.,service.&limit=3&offset={offset}"),
3097 )
3098 .await;
3099 assert_eq!(
3100 p.has_more,
3101 offset + p.count < 6,
3102 "has_more must terminate on the excluded view (offset={offset})"
3103 );
3104 walked.extend(p.events.iter().map(|e| e.id));
3105 }
3106 assert_eq!(
3107 walked, kept,
3108 "exclusion + paging must cover the survivors once"
3109 );
3110
3111 let desc = query_page(
3113 &store,
3114 "exclude_event_type_prefix=audit.,service.&limit=1&order=desc",
3115 )
3116 .await;
3117 assert_eq!(
3118 desc.events[0].id,
3119 *kept.last().unwrap(),
3120 "order=desc&limit=1 over an excluded view is the newest SURVIVOR"
3121 );
3122
3123 let audit_only = query_page(&store, "exclude_event_type_prefix=audit.").await;
3126 assert_eq!(audit_only.total_count, 9, "12 minus the 3 audit.* events");
3127 let nothing = query_page(&store, "exclude_event_type_prefix=nosuch.").await;
3128 assert_eq!(nothing.total_count, 12);
3129 }
3130
3131 #[tokio::test]
3145 async fn query_events_honours_time_window_without_an_entity_or_type_filter() {
3146 use chrono::{SecondsFormat, SubsecRound};
3147
3148 let store = create_test_store();
3149 let base = chrono::DateTime::parse_from_rfc3339("2026-01-01T00:00:00.123456789Z")
3163 .unwrap()
3164 .to_utc()
3165 .trunc_subsecs(0);
3166 let mut ids = Vec::new();
3167 for i in 0..5i64 {
3168 let mut event = create_test_event(&format!("e-{i}"), "user.created");
3169 event.timestamp = base + chrono::Duration::hours(i);
3170 event.version = i + 1;
3171 ids.push(event.id);
3172 store.ingest(&event).unwrap();
3173 }
3174 let at = |h: i64| {
3177 let t = base + chrono::Duration::hours(h);
3178 let s = t.to_rfc3339_opts(SecondsFormat::Micros, true);
3179 assert_eq!(
3184 chrono::DateTime::parse_from_rfc3339(&s).unwrap().to_utc(),
3185 t,
3186 "query-string timestamp must round-trip exactly, else the window \
3187 boundary silently excludes the event stamped at it"
3188 );
3189 s
3190 };
3191
3192 for (qs, expected) in [
3193 (format!("since={}", at(2)), vec![ids[2], ids[3], ids[4]]),
3194 (format!("until={}", at(1)), vec![ids[0], ids[1]]),
3195 (format!("as_of={}", at(1)), vec![ids[0], ids[1]]),
3196 (
3197 format!("since={}&until={}", at(1), at(3)),
3198 vec![ids[1], ids[2], ids[3]],
3199 ),
3200 ] {
3201 let resp = query_page(&store, &qs).await;
3202 let got: Vec<_> = resp.events.iter().map(|e| e.id).collect();
3203 assert_eq!(got, expected, "?{qs} must return only the window");
3204 assert_eq!(resp.count, expected.len(), "?{qs}");
3205 assert_eq!(
3206 resp.total_count,
3207 expected.len(),
3208 "?{qs}: total_count must count the window, not the history"
3209 );
3210 assert!(!resp.has_more, "?{qs}: the whole window was served");
3211 }
3212
3213 let page1 = query_page(&store, &format!("since={}&limit=2", at(2))).await;
3216 let page2 = query_page(&store, &format!("since={}&limit=2&offset=2", at(2))).await;
3217 assert_eq!(
3218 page1.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3219 vec![ids[2], ids[3]]
3220 );
3221 assert!(page1.has_more, "3 in the window, 2 served");
3222 assert_eq!(
3223 page2.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3224 vec![ids[4]]
3225 );
3226 assert!(!page2.has_more, "offset 2 + count 1 == the window's 3");
3227 assert_eq!(page2.total_count, 3);
3228
3229 let desc = query_page(&store, &format!("since={}&order=desc", at(2))).await;
3231 assert_eq!(
3232 desc.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3233 vec![ids[4], ids[3], ids[2]]
3234 );
3235
3236 let empty = query_page(&store, &format!("since={}", at(99))).await;
3238 assert_eq!(empty.count, 0);
3239 assert_eq!(empty.total_count, 0);
3240 assert!(!empty.has_more);
3241 }
3242
3243 #[tokio::test]
3247 async fn query_events_fails_closed_without_tenant() {
3248 use axum::extract::{Query, State};
3249
3250 let store = create_test_store();
3251 for i in 0..5 {
3252 store
3253 .ingest(&create_test_event(&format!("e-{i}"), "user.created"))
3254 .unwrap();
3255 }
3256
3257 let resp = query_events(
3258 OptionalAuth(None),
3259 Query(QueryEventsRequest::default()),
3260 Query(EventOrderParam { order: None }),
3261 Query(EventOffsetParam { offset: None }),
3262 Query(EventIntegrityParam::default()),
3263 State(store.clone()),
3264 )
3265 .await
3266 .unwrap();
3267 assert_eq!(
3268 resp.0.total_count, 0,
3269 "a no-tenant query must NOT return cross-tenant events"
3270 );
3271 assert_eq!(resp.0.count, 0);
3272
3273 let scoped = query_events(
3276 OptionalAuth(None),
3277 Query(QueryEventsRequest {
3278 tenant_id: Some("test-stream".to_string()),
3279 ..QueryEventsRequest::default()
3280 }),
3281 Query(EventOrderParam { order: None }),
3282 Query(EventOffsetParam { offset: None }),
3283 Query(EventIntegrityParam::default()),
3284 State(store),
3285 )
3286 .await
3287 .unwrap();
3288 assert_eq!(
3289 scoped.0.total_count, 5,
3290 "tenant-scoped query returns its events"
3291 );
3292 }
3293
3294 #[tokio::test]
3299 async fn list_streams_and_types_are_tenant_scoped() {
3300 use crate::domain::entities::Event;
3301 use axum::extract::{Query, State};
3302
3303 let store = create_test_store();
3304 let ev = |entity: &str, etype: &str, tenant: &str| {
3305 Event::from_strings(
3306 etype.to_string(),
3307 entity.to_string(),
3308 tenant.to_string(),
3309 serde_json::json!({}),
3310 None,
3311 )
3312 .unwrap()
3313 };
3314 store.ingest(&ev("e1", "order.placed", "tenant-a")).unwrap();
3316 store.ingest(&ev("e2", "user.created", "tenant-a")).unwrap();
3317 store
3318 .ingest(&ev("e9", "thing.happened", "tenant-b"))
3319 .unwrap();
3320
3321 let streams = |tid: Option<&str>| {
3322 list_streams(
3323 OptionalAuth(None),
3324 State(store.clone()),
3325 Query(ListStreamsParams {
3326 tenant_id: tid.map(String::from),
3327 limit: None,
3328 offset: None,
3329 }),
3330 )
3331 };
3332 assert_eq!(
3333 streams(Some("tenant-a")).await.0.total,
3334 2,
3335 "tenant-a streams"
3336 );
3337 assert_eq!(
3338 streams(Some("tenant-b")).await.0.total,
3339 1,
3340 "tenant-b streams"
3341 );
3342 assert_eq!(
3343 streams(None).await.0.total,
3344 0,
3345 "no tenant -> no streams (fail closed)"
3346 );
3347
3348 let types = |tid: Option<&str>| {
3349 list_event_types(
3350 OptionalAuth(None),
3351 State(store.clone()),
3352 Query(ListEventTypesParams {
3353 tenant_id: tid.map(String::from),
3354 limit: None,
3355 offset: None,
3356 }),
3357 )
3358 };
3359 assert_eq!(
3360 types(Some("tenant-a")).await.0.total,
3361 2,
3362 "tenant-a event types"
3363 );
3364 assert_eq!(
3365 types(Some("tenant-b")).await.0.total,
3366 1,
3367 "tenant-b event types"
3368 );
3369 assert_eq!(
3370 types(None).await.0.total,
3371 0,
3372 "no tenant -> no types (fail closed)"
3373 );
3374 }
3375
3376 #[tokio::test]
3377 async fn test_query_events_no_more_results() {
3378 let store = create_test_store();
3379
3380 for i in 0..5 {
3382 store
3383 .ingest(&create_test_event(&format!("entity-{i}"), "user.created"))
3384 .unwrap();
3385 }
3386
3387 let all_events = store
3389 .query(&QueryEventsRequest {
3390 entity_id: None,
3391 event_type: None,
3392 tenant_id: None,
3393 as_of: None,
3394 since: None,
3395 until: None,
3396 limit: None,
3397 event_type_prefix: None,
3398 exclude_event_type_prefix: None,
3399 payload_filter: None,
3400 })
3401 .unwrap();
3402 let total_count = all_events.len();
3403 let limited_events: Vec<Event> = all_events.into_iter().take(100).collect();
3404 let count = limited_events.len();
3405 let has_more = count < total_count;
3406
3407 assert_eq!(count, 5);
3408 assert_eq!(total_count, 5);
3409 assert!(!has_more);
3410 }
3411
3412 #[tokio::test]
3416 async fn test_query_events_order_desc_returns_latest() {
3417 let store = create_test_store();
3423
3424 let base = chrono::Utc::now();
3427 let mut ascending_ids = Vec::new();
3428 for i in 0..5i64 {
3429 let mut event = create_test_event("org-1", "auth.org.updated");
3430 event.timestamp = base + chrono::Duration::seconds(i);
3431 event.version = i + 1;
3432 ascending_ids.push(event.id);
3433 store.ingest(&event).unwrap();
3434 }
3435 let newest_ts = base + chrono::Duration::seconds(4);
3436
3437 let latest = query_page(&store, "entity_id=org-1&limit=1&order=desc").await;
3439 assert_eq!(latest.count, 1);
3440 assert_eq!(
3441 latest.events[0].id,
3442 ascending_ids[4],
3443 "order=desc&limit=1 must yield the NEWEST event, got the one at \
3444 ascending position {:?}",
3445 ascending_ids
3446 .iter()
3447 .position(|id| *id == latest.events[0].id)
3448 );
3449 assert_eq!(latest.events[0].timestamp, newest_ts);
3450 assert_eq!(latest.total_count, 5, "total is the full match set");
3451 assert!(latest.has_more);
3452
3453 for qs in [
3455 "entity_id=org-1&limit=1",
3456 "entity_id=org-1&limit=1&order=asc",
3457 ] {
3458 let oldest = query_page(&store, qs).await;
3459 assert_eq!(oldest.events[0].id, ascending_ids[0], "{qs}");
3460 assert_eq!(oldest.events[0].timestamp, base);
3461 }
3462
3463 let all_desc = query_page(&store, "entity_id=org-1&order=desc").await;
3466 let got: Vec<_> = all_desc.events.iter().map(|e| e.id).collect();
3467 let expected: Vec<_> = ascending_ids.iter().rev().copied().collect();
3468 assert_eq!(got, expected, "order=desc must return newest-first");
3469 }
3470
3471 #[tokio::test]
3478 async fn query_events_desc_composes_with_offset_and_limit() {
3479 let store = create_test_store();
3480 let base = chrono::Utc::now();
3481 let mut ascending_ids = Vec::new();
3482 for i in 0..5i64 {
3483 let mut event = create_test_event("org-1", "auth.org.updated");
3484 event.timestamp = base + chrono::Duration::seconds(i);
3485 event.version = i + 1;
3486 ascending_ids.push(event.id);
3487 store.ingest(&event).unwrap();
3488 }
3489 let newest_first: Vec<_> = ascending_ids.iter().rev().copied().collect();
3490
3491 for (offset, limit) in [(0, 2), (1, 2), (2, 2), (3, 2), (4, 2), (5, 2), (1, 4)] {
3492 let page = query_page(
3493 &store,
3494 &format!("entity_id=org-1&order=desc&offset={offset}&limit={limit}"),
3495 )
3496 .await;
3497 let got: Vec<_> = page.events.iter().map(|e| e.id).collect();
3498 let expected: Vec<_> = newest_first
3499 .iter()
3500 .skip(offset)
3501 .take(limit)
3502 .copied()
3503 .collect();
3504 assert_eq!(
3505 got, expected,
3506 "order=desc&offset={offset}&limit={limit} must reverse, then \
3507 skip, then take"
3508 );
3509 assert_eq!(page.count, expected.len());
3510 assert_eq!(page.total_count, 5);
3511 assert_eq!(
3512 page.has_more,
3513 offset + expected.len() < 5,
3514 "has_more must account for the offset (offset={offset})"
3515 );
3516 }
3517
3518 let mut walked = Vec::new();
3521 for offset in (0..5).step_by(2) {
3522 let page = query_page(
3523 &store,
3524 &format!("entity_id=org-1&order=desc&offset={offset}&limit=2"),
3525 )
3526 .await;
3527 walked.extend(page.events.iter().map(|e| e.id));
3528 }
3529 assert_eq!(walked, newest_first, "desc paging must cover the set once");
3530 }
3531
3532 #[tokio::test]
3533 async fn test_list_entities_by_type_prefix() {
3534 let store = create_test_store();
3535
3536 store
3538 .ingest(&create_test_event("idx-1", "index.created"))
3539 .unwrap();
3540 store
3541 .ingest(&create_test_event("idx-1", "index.updated"))
3542 .unwrap();
3543 store
3544 .ingest(&create_test_event("idx-2", "index.created"))
3545 .unwrap();
3546 store
3547 .ingest(&create_test_event("idx-3", "index.created"))
3548 .unwrap();
3549 store
3551 .ingest(&create_test_event("trade-1", "trade.created"))
3552 .unwrap();
3553 store
3554 .ingest(&create_test_event("trade-2", "trade.created"))
3555 .unwrap();
3556
3557 let req = ListEntitiesRequest {
3559 event_type_prefix: Some("index.".to_string()),
3560 ..Default::default()
3561 };
3562 let query_req = QueryEventsRequest {
3563 entity_id: None,
3564 event_type: None,
3565 tenant_id: None,
3566 as_of: None,
3567 since: None,
3568 until: None,
3569 limit: None,
3570 event_type_prefix: req.event_type_prefix,
3571 exclude_event_type_prefix: None,
3572 payload_filter: req.payload_filter,
3573 };
3574 let events = store.query(&query_req).unwrap();
3575
3576 let mut entity_map: std::collections::HashMap<String, Vec<&Event>> =
3578 std::collections::HashMap::new();
3579 for event in &events {
3580 entity_map
3581 .entry(event.entity_id().to_string())
3582 .or_default()
3583 .push(event);
3584 }
3585
3586 assert_eq!(entity_map.len(), 3); assert_eq!(entity_map["idx-1"].len(), 2); assert_eq!(entity_map["idx-2"].len(), 1);
3589 assert_eq!(entity_map["idx-3"].len(), 1);
3590 }
3591
3592 #[tokio::test]
3595 async fn test_list_entities_order_and_pagination() {
3596 let store = create_test_store();
3597
3598 let base = chrono::Utc::now();
3600 for (i, eid) in ["org-a", "org-b", "org-c"].iter().enumerate() {
3601 let mut event = create_test_event(eid, "auth.org.created");
3602 event.timestamp = base + chrono::Duration::seconds(i as i64);
3603 store.ingest(&event).unwrap();
3604 }
3605 let prefix = || Some("auth.org.".to_string());
3606
3607 let desc = list_entities(
3609 State(store.clone()),
3610 Query(ListEntitiesRequest {
3611 event_type_prefix: prefix(),
3612 ..Default::default()
3613 }),
3614 )
3615 .await
3616 .unwrap();
3617 let desc_ids: Vec<&str> = desc
3618 .0
3619 .entities
3620 .iter()
3621 .map(|e| e.entity_id.as_str())
3622 .collect();
3623 assert_eq!(desc_ids, ["org-c", "org-b", "org-a"]);
3624
3625 let asc = list_entities(
3627 State(store.clone()),
3628 Query(ListEntitiesRequest {
3629 event_type_prefix: prefix(),
3630 order: Some("asc".to_string()),
3631 ..Default::default()
3632 }),
3633 )
3634 .await
3635 .unwrap();
3636 let asc_ids: Vec<&str> = asc
3637 .0
3638 .entities
3639 .iter()
3640 .map(|e| e.entity_id.as_str())
3641 .collect();
3642 assert_eq!(asc_ids, ["org-a", "org-b", "org-c"]);
3643
3644 let page2 = list_entities(
3646 State(store.clone()),
3647 Query(ListEntitiesRequest {
3648 event_type_prefix: prefix(),
3649 order: Some("asc".to_string()),
3650 limit: Some(1),
3651 offset: Some(1),
3652 ..Default::default()
3653 }),
3654 )
3655 .await
3656 .unwrap();
3657 assert_eq!(page2.0.entities.len(), 1);
3658 assert_eq!(page2.0.entities[0].entity_id, "org-b");
3659 assert_eq!(page2.0.total, 3);
3660 assert!(page2.0.has_more);
3661
3662 let err = list_entities(
3664 State(store.clone()),
3665 Query(ListEntitiesRequest {
3666 event_type_prefix: prefix(),
3667 order: Some("sideways".to_string()),
3668 ..Default::default()
3669 }),
3670 )
3671 .await;
3672 assert!(err.is_err(), "invalid order value must be rejected");
3673 }
3674
3675 fn create_test_event_with_payload(
3676 entity_id: &str,
3677 event_type: &str,
3678 payload: serde_json::Value,
3679 ) -> Event {
3680 Event::from_strings(
3681 event_type.to_string(),
3682 entity_id.to_string(),
3683 "test-stream".to_string(),
3684 payload,
3685 None,
3686 )
3687 .unwrap()
3688 }
3689
3690 #[tokio::test]
3691 async fn test_detect_duplicates_by_payload_fields() {
3692 let store = create_test_store();
3693
3694 store
3696 .ingest(&create_test_event_with_payload(
3697 "idx-1",
3698 "index.created",
3699 serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3700 ))
3701 .unwrap();
3702 store
3703 .ingest(&create_test_event_with_payload(
3704 "idx-2",
3705 "index.created",
3706 serde_json::json!({"name": "S&P 500", "user_id": "bob"}),
3707 ))
3708 .unwrap();
3709 store
3710 .ingest(&create_test_event_with_payload(
3711 "idx-3",
3712 "index.created",
3713 serde_json::json!({"name": "NASDAQ", "user_id": "alice"}),
3714 ))
3715 .unwrap();
3716 store
3717 .ingest(&create_test_event_with_payload(
3718 "idx-4",
3719 "index.created",
3720 serde_json::json!({"name": "NASDAQ", "user_id": "carol"}),
3721 ))
3722 .unwrap();
3723 store
3724 .ingest(&create_test_event_with_payload(
3725 "idx-5",
3726 "index.created",
3727 serde_json::json!({"name": "DAX", "user_id": "dave"}),
3728 ))
3729 .unwrap();
3730
3731 let query_req = QueryEventsRequest {
3733 entity_id: None,
3734 event_type: None,
3735 tenant_id: None,
3736 as_of: None,
3737 since: None,
3738 until: None,
3739 limit: None,
3740 event_type_prefix: Some("index.".to_string()),
3741 exclude_event_type_prefix: None,
3742 payload_filter: None,
3743 };
3744 let events = store.query(&query_req).unwrap();
3745
3746 let group_by_fields = vec!["name"];
3748 let mut entity_latest: std::collections::HashMap<String, &Event> =
3749 std::collections::HashMap::new();
3750 for event in &events {
3751 let eid = event.entity_id().to_string();
3752 entity_latest
3753 .entry(eid)
3754 .and_modify(|existing| {
3755 if event.timestamp() > existing.timestamp() {
3756 *existing = event;
3757 }
3758 })
3759 .or_insert(event);
3760 }
3761
3762 let mut groups: std::collections::HashMap<String, Vec<String>> =
3763 std::collections::HashMap::new();
3764 for (entity_id, event) in &entity_latest {
3765 let payload = event.payload();
3766 let mut key_parts = serde_json::Map::new();
3767 for field in &group_by_fields {
3768 let value = payload
3769 .get(*field)
3770 .cloned()
3771 .unwrap_or(serde_json::Value::Null);
3772 key_parts.insert((*field).to_string(), value);
3773 }
3774 let key_str = serde_json::to_string(&key_parts).unwrap_or_default();
3775 groups.entry(key_str).or_default().push(entity_id.clone());
3776 }
3777
3778 let duplicate_groups: Vec<_> = groups
3779 .into_iter()
3780 .filter(|(_, ids)| ids.len() > 1)
3781 .collect();
3782
3783 assert_eq!(duplicate_groups.len(), 2); for (_, ids) in &duplicate_groups {
3785 assert_eq!(ids.len(), 2);
3786 }
3787 }
3788
3789 #[tokio::test]
3790 async fn test_detect_duplicates_no_duplicates() {
3791 let store = create_test_store();
3792
3793 store
3795 .ingest(&create_test_event_with_payload(
3796 "idx-1",
3797 "index.created",
3798 serde_json::json!({"name": "A"}),
3799 ))
3800 .unwrap();
3801 store
3802 .ingest(&create_test_event_with_payload(
3803 "idx-2",
3804 "index.created",
3805 serde_json::json!({"name": "B"}),
3806 ))
3807 .unwrap();
3808
3809 let query_req = QueryEventsRequest {
3810 entity_id: None,
3811 event_type: None,
3812 tenant_id: None,
3813 as_of: None,
3814 since: None,
3815 until: None,
3816 limit: None,
3817 event_type_prefix: Some("index.".to_string()),
3818 exclude_event_type_prefix: None,
3819 payload_filter: None,
3820 };
3821 let events = store.query(&query_req).unwrap();
3822
3823 let mut entity_latest: std::collections::HashMap<String, &Event> =
3824 std::collections::HashMap::new();
3825 for event in &events {
3826 entity_latest
3827 .entry(event.entity_id().to_string())
3828 .or_insert(event);
3829 }
3830
3831 let mut groups: std::collections::HashMap<String, Vec<String>> =
3832 std::collections::HashMap::new();
3833 for (entity_id, event) in &entity_latest {
3834 let key_str =
3835 serde_json::to_string(&serde_json::json!({"name": event.payload().get("name")}))
3836 .unwrap();
3837 groups.entry(key_str).or_default().push(entity_id.clone());
3838 }
3839
3840 let duplicate_groups: Vec<_> = groups
3841 .into_iter()
3842 .filter(|(_, ids)| ids.len() > 1)
3843 .collect();
3844
3845 assert_eq!(duplicate_groups.len(), 0); }
3847
3848 #[tokio::test]
3849 async fn test_detect_duplicates_multi_field_group_by() {
3850 let store = create_test_store();
3851
3852 store
3854 .ingest(&create_test_event_with_payload(
3855 "idx-1",
3856 "index.created",
3857 serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3858 ))
3859 .unwrap();
3860 store
3861 .ingest(&create_test_event_with_payload(
3862 "idx-2",
3863 "index.created",
3864 serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3865 ))
3866 .unwrap();
3867 store
3869 .ingest(&create_test_event_with_payload(
3870 "idx-3",
3871 "index.created",
3872 serde_json::json!({"name": "S&P 500", "user_id": "bob"}),
3873 ))
3874 .unwrap();
3875
3876 let query_req = QueryEventsRequest {
3877 entity_id: None,
3878 event_type: None,
3879 tenant_id: None,
3880 as_of: None,
3881 since: None,
3882 until: None,
3883 limit: None,
3884 event_type_prefix: Some("index.".to_string()),
3885 exclude_event_type_prefix: None,
3886 payload_filter: None,
3887 };
3888 let events = store.query(&query_req).unwrap();
3889
3890 let group_by_fields = vec!["name", "user_id"];
3891 let mut entity_latest: std::collections::HashMap<String, &Event> =
3892 std::collections::HashMap::new();
3893 for event in &events {
3894 entity_latest
3895 .entry(event.entity_id().to_string())
3896 .and_modify(|existing| {
3897 if event.timestamp() > existing.timestamp() {
3898 *existing = event;
3899 }
3900 })
3901 .or_insert(event);
3902 }
3903
3904 let mut groups: std::collections::HashMap<String, Vec<String>> =
3905 std::collections::HashMap::new();
3906 for (entity_id, event) in &entity_latest {
3907 let payload = event.payload();
3908 let mut key_parts = serde_json::Map::new();
3909 for field in &group_by_fields {
3910 let value = payload
3911 .get(*field)
3912 .cloned()
3913 .unwrap_or(serde_json::Value::Null);
3914 key_parts.insert((*field).to_string(), value);
3915 }
3916 let key_str = serde_json::to_string(&key_parts).unwrap_or_default();
3917 groups.entry(key_str).or_default().push(entity_id.clone());
3918 }
3919
3920 let duplicate_groups: Vec<_> = groups
3921 .into_iter()
3922 .filter(|(_, ids)| ids.len() > 1)
3923 .collect();
3924
3925 assert_eq!(duplicate_groups.len(), 1);
3927 let (_, ref ids) = duplicate_groups[0];
3928 assert_eq!(ids.len(), 2);
3929 let mut sorted_ids = ids.clone();
3930 sorted_ids.sort();
3931 assert_eq!(sorted_ids, vec!["idx-1", "idx-2"]);
3932 }
3933
3934 #[tokio::test]
3935 async fn test_projection_state_cache() {
3936 let store = create_test_store();
3937
3938 let cache = store.projection_state_cache();
3940 cache.insert(
3941 "entity_snapshots:user-123".to_string(),
3942 serde_json::json!({"name": "Test User", "age": 30}),
3943 );
3944
3945 let state = cache.get("entity_snapshots:user-123");
3947 assert!(state.is_some());
3948 let state = state.unwrap();
3949 assert_eq!(state["name"], "Test User");
3950 assert_eq!(state["age"], 30);
3951 }
3952
3953 #[tokio::test]
3954 async fn test_projection_manager_list_projections() {
3955 let store = create_test_store();
3956
3957 let projection_manager = store.projection_manager();
3959 let projections = projection_manager.list_projections();
3960
3961 assert!(projections.len() >= 2);
3963
3964 let names: Vec<&str> = projections.iter().map(|(name, _)| name.as_str()).collect();
3965 assert!(names.contains(&"entity_snapshots"));
3966 assert!(names.contains(&"event_counters"));
3967 }
3968
3969 #[tokio::test]
3970 async fn test_projection_state_after_event_ingestion() {
3971 let store = create_test_store();
3972
3973 let event = create_test_event("user-456", "user.created");
3975 store.ingest(&event).unwrap();
3976
3977 let projection_manager = store.projection_manager();
3979 let snapshot_projection = projection_manager
3980 .get_projection("entity_snapshots")
3981 .unwrap();
3982
3983 let state = snapshot_projection.get_state("user-456");
3984 assert!(state.is_some());
3985 let state = state.unwrap();
3986 assert_eq!(state["name"], "Test");
3987 assert_eq!(state["value"], 42);
3988 }
3989
3990 #[tokio::test]
3991 async fn test_projection_state_cache_multiple_entities() {
3992 let store = create_test_store();
3993 let cache = store.projection_state_cache();
3994
3995 for i in 0..10 {
3997 cache.insert(
3998 format!("entity_snapshots:entity-{i}"),
3999 serde_json::json!({"id": i, "status": "active"}),
4000 );
4001 }
4002
4003 assert_eq!(cache.len(), 10);
4005
4006 for i in 0..10 {
4008 let key = format!("entity_snapshots:entity-{i}");
4009 let state = cache.get(&key);
4010 assert!(state.is_some());
4011 assert_eq!(state.unwrap()["id"], i);
4012 }
4013 }
4014
4015 #[tokio::test]
4016 async fn test_projection_state_update() {
4017 let store = create_test_store();
4018 let cache = store.projection_state_cache();
4019
4020 cache.insert(
4022 "entity_snapshots:user-789".to_string(),
4023 serde_json::json!({"balance": 100}),
4024 );
4025
4026 cache.insert(
4028 "entity_snapshots:user-789".to_string(),
4029 serde_json::json!({"balance": 150}),
4030 );
4031
4032 let state = cache.get("entity_snapshots:user-789").unwrap();
4034 assert_eq!(state["balance"], 150);
4035 }
4036
4037 #[tokio::test]
4038 async fn test_event_counter_projection() {
4039 let store = create_test_store();
4040
4041 store
4043 .ingest(&create_test_event("user-1", "user.created"))
4044 .unwrap();
4045 store
4046 .ingest(&create_test_event("user-2", "user.created"))
4047 .unwrap();
4048 store
4049 .ingest(&create_test_event("user-1", "user.updated"))
4050 .unwrap();
4051
4052 let projection_manager = store.projection_manager();
4054 let counter_projection = projection_manager.get_projection("event_counters").unwrap();
4055
4056 let created_state = counter_projection.get_state("user.created");
4058 assert!(created_state.is_some());
4059 assert_eq!(created_state.unwrap()["count"], 2);
4060
4061 let updated_state = counter_projection.get_state("user.updated");
4062 assert!(updated_state.is_some());
4063 assert_eq!(updated_state.unwrap()["count"], 1);
4064 }
4065
4066 #[tokio::test]
4067 async fn test_projection_state_cache_key_format() {
4068 let store = create_test_store();
4069 let cache = store.projection_state_cache();
4070
4071 let key = "orders:order-12345".to_string();
4073 cache.insert(key.clone(), serde_json::json!({"total": 99.99}));
4074
4075 let state = cache.get(&key).unwrap();
4076 assert_eq!(state["total"], 99.99);
4077 }
4078
4079 #[tokio::test]
4080 async fn test_projection_state_cache_removal() {
4081 let store = create_test_store();
4082 let cache = store.projection_state_cache();
4083
4084 cache.insert(
4086 "test:entity-1".to_string(),
4087 serde_json::json!({"data": "value"}),
4088 );
4089 assert_eq!(cache.len(), 1);
4090
4091 cache.remove("test:entity-1");
4092 assert_eq!(cache.len(), 0);
4093 assert!(cache.get("test:entity-1").is_none());
4094 }
4095
4096 #[tokio::test]
4097 async fn test_get_nonexistent_projection() {
4098 let store = create_test_store();
4099 let projection_manager = store.projection_manager();
4100
4101 let projection = projection_manager.get_projection("nonexistent_projection");
4103 assert!(projection.is_none());
4104 }
4105
4106 #[tokio::test]
4107 async fn test_get_nonexistent_entity_state() {
4108 let store = create_test_store();
4109 let projection_manager = store.projection_manager();
4110
4111 let snapshot_projection = projection_manager
4113 .get_projection("entity_snapshots")
4114 .unwrap();
4115 let state = snapshot_projection.get_state("nonexistent-entity-xyz");
4116 assert!(state.is_none());
4117 }
4118
4119 #[tokio::test]
4120 async fn test_projection_state_cache_concurrent_access() {
4121 let store = create_test_store();
4122 let cache = store.projection_state_cache();
4123
4124 let handles: Vec<_> = (0..10)
4126 .map(|i| {
4127 let cache_clone = cache.clone();
4128 tokio::spawn(async move {
4129 cache_clone.insert(
4130 format!("concurrent:entity-{i}"),
4131 serde_json::json!({"thread": i}),
4132 );
4133 })
4134 })
4135 .collect();
4136
4137 for handle in handles {
4138 handle.await.unwrap();
4139 }
4140
4141 assert_eq!(cache.len(), 10);
4143 }
4144
4145 #[tokio::test]
4146 async fn test_projection_state_large_payload() {
4147 let store = create_test_store();
4148 let cache = store.projection_state_cache();
4149
4150 let large_array: Vec<serde_json::Value> = (0..1000)
4152 .map(|i| serde_json::json!({"item": i, "description": "test item with some padding data to increase size"}))
4153 .collect();
4154
4155 cache.insert(
4156 "large:entity-1".to_string(),
4157 serde_json::json!({"items": large_array}),
4158 );
4159
4160 let state = cache.get("large:entity-1").unwrap();
4161 let items = state["items"].as_array().unwrap();
4162 assert_eq!(items.len(), 1000);
4163 }
4164
4165 #[tokio::test]
4166 async fn test_projection_state_complex_json() {
4167 let store = create_test_store();
4168 let cache = store.projection_state_cache();
4169
4170 let complex_state = serde_json::json!({
4172 "user": {
4173 "id": "user-123",
4174 "profile": {
4175 "name": "John Doe",
4176 "email": "john@example.com",
4177 "settings": {
4178 "theme": "dark",
4179 "notifications": true
4180 }
4181 },
4182 "roles": ["admin", "user"],
4183 "metadata": {
4184 "created_at": "2025-01-01T00:00:00Z",
4185 "last_login": null
4186 }
4187 }
4188 });
4189
4190 cache.insert("complex:user-123".to_string(), complex_state);
4191
4192 let state = cache.get("complex:user-123").unwrap();
4193 assert_eq!(state["user"]["profile"]["name"], "John Doe");
4194 assert_eq!(state["user"]["roles"][0], "admin");
4195 assert!(state["user"]["metadata"]["last_login"].is_null());
4196 }
4197
4198 #[tokio::test]
4199 async fn test_projection_state_cache_iteration() {
4200 let store = create_test_store();
4201 let cache = store.projection_state_cache();
4202
4203 for i in 0..5 {
4205 cache.insert(format!("iter:entity-{i}"), serde_json::json!({"index": i}));
4206 }
4207
4208 let entries: Vec<_> = cache.iter().map(|entry| entry.key().clone()).collect();
4210 assert_eq!(entries.len(), 5);
4211 }
4212
4213 #[tokio::test]
4214 async fn test_projection_manager_get_entity_snapshots() {
4215 let store = create_test_store();
4216 let projection_manager = store.projection_manager();
4217
4218 let projection = projection_manager.get_projection("entity_snapshots");
4220 assert!(projection.is_some());
4221 assert_eq!(projection.unwrap().name(), "entity_snapshots");
4222 }
4223
4224 #[tokio::test]
4225 async fn test_projection_manager_get_event_counters() {
4226 let store = create_test_store();
4227 let projection_manager = store.projection_manager();
4228
4229 let projection = projection_manager.get_projection("event_counters");
4231 assert!(projection.is_some());
4232 assert_eq!(projection.unwrap().name(), "event_counters");
4233 }
4234
4235 #[tokio::test]
4236 async fn test_projection_state_cache_overwrite() {
4237 let store = create_test_store();
4238 let cache = store.projection_state_cache();
4239
4240 cache.insert(
4242 "overwrite:entity-1".to_string(),
4243 serde_json::json!({"version": 1}),
4244 );
4245
4246 cache.insert(
4248 "overwrite:entity-1".to_string(),
4249 serde_json::json!({"version": 2}),
4250 );
4251
4252 cache.insert(
4254 "overwrite:entity-1".to_string(),
4255 serde_json::json!({"version": 3}),
4256 );
4257
4258 let state = cache.get("overwrite:entity-1").unwrap();
4259 assert_eq!(state["version"], 3);
4260
4261 assert_eq!(cache.len(), 1);
4263 }
4264
4265 #[tokio::test]
4266 async fn test_projection_state_multiple_projections() {
4267 let store = create_test_store();
4268 let cache = store.projection_state_cache();
4269
4270 cache.insert(
4272 "entity_snapshots:user-1".to_string(),
4273 serde_json::json!({"name": "Alice"}),
4274 );
4275 cache.insert(
4276 "event_counters:user.created".to_string(),
4277 serde_json::json!({"count": 5}),
4278 );
4279 cache.insert(
4280 "custom_projection:order-1".to_string(),
4281 serde_json::json!({"total": 150.0}),
4282 );
4283
4284 assert_eq!(
4286 cache.get("entity_snapshots:user-1").unwrap()["name"],
4287 "Alice"
4288 );
4289 assert_eq!(
4290 cache.get("event_counters:user.created").unwrap()["count"],
4291 5
4292 );
4293 assert_eq!(
4294 cache.get("custom_projection:order-1").unwrap()["total"],
4295 150.0
4296 );
4297 }
4298
4299 #[tokio::test]
4300 async fn test_bulk_projection_state_access() {
4301 let store = create_test_store();
4302
4303 for i in 0..5 {
4305 let event = create_test_event(&format!("bulk-user-{i}"), "user.created");
4306 store.ingest(&event).unwrap();
4307 }
4308
4309 let projection_manager = store.projection_manager();
4311 let snapshot_projection = projection_manager
4312 .get_projection("entity_snapshots")
4313 .unwrap();
4314
4315 for i in 0..5 {
4317 let state = snapshot_projection.get_state(&format!("bulk-user-{i}"));
4318 assert!(state.is_some(), "Entity bulk-user-{i} should have state");
4319 }
4320 }
4321
4322 #[tokio::test]
4323 async fn test_bulk_save_projection_states() {
4324 let store = create_test_store();
4325 let cache = store.projection_state_cache();
4326
4327 let states = vec![
4329 BulkSaveStateItem {
4330 entity_id: "bulk-entity-1".to_string(),
4331 state: serde_json::json!({"name": "Entity 1", "value": 100}),
4332 },
4333 BulkSaveStateItem {
4334 entity_id: "bulk-entity-2".to_string(),
4335 state: serde_json::json!({"name": "Entity 2", "value": 200}),
4336 },
4337 BulkSaveStateItem {
4338 entity_id: "bulk-entity-3".to_string(),
4339 state: serde_json::json!({"name": "Entity 3", "value": 300}),
4340 },
4341 ];
4342
4343 let projection_name = "test_projection";
4344
4345 for item in &states {
4347 cache.insert(
4348 format!("{projection_name}:{}", item.entity_id),
4349 item.state.clone(),
4350 );
4351 }
4352
4353 assert_eq!(cache.len(), 3);
4355
4356 let state1 = cache.get("test_projection:bulk-entity-1").unwrap();
4357 assert_eq!(state1["name"], "Entity 1");
4358 assert_eq!(state1["value"], 100);
4359
4360 let state2 = cache.get("test_projection:bulk-entity-2").unwrap();
4361 assert_eq!(state2["name"], "Entity 2");
4362 assert_eq!(state2["value"], 200);
4363
4364 let state3 = cache.get("test_projection:bulk-entity-3").unwrap();
4365 assert_eq!(state3["name"], "Entity 3");
4366 assert_eq!(state3["value"], 300);
4367 }
4368
4369 #[tokio::test]
4370 async fn test_bulk_save_empty_states() {
4371 let store = create_test_store();
4372 let cache = store.projection_state_cache();
4373
4374 cache.clear();
4376
4377 let states: Vec<BulkSaveStateItem> = vec![];
4379 assert_eq!(states.len(), 0);
4380
4381 assert_eq!(cache.len(), 0);
4383 }
4384
4385 #[tokio::test]
4386 async fn test_bulk_save_overwrites_existing() {
4387 let store = create_test_store();
4388 let cache = store.projection_state_cache();
4389
4390 cache.insert(
4392 "test:entity-1".to_string(),
4393 serde_json::json!({"version": 1, "data": "initial"}),
4394 );
4395
4396 let new_state = serde_json::json!({"version": 2, "data": "updated"});
4398 cache.insert("test:entity-1".to_string(), new_state);
4399
4400 let state = cache.get("test:entity-1").unwrap();
4402 assert_eq!(state["version"], 2);
4403 assert_eq!(state["data"], "updated");
4404 }
4405
4406 #[tokio::test]
4407 async fn test_bulk_save_high_volume() {
4408 let store = create_test_store();
4409 let cache = store.projection_state_cache();
4410
4411 for i in 0..1000 {
4413 cache.insert(
4414 format!("volume_test:entity-{i}"),
4415 serde_json::json!({"index": i, "status": "active"}),
4416 );
4417 }
4418
4419 assert_eq!(cache.len(), 1000);
4421
4422 assert_eq!(cache.get("volume_test:entity-0").unwrap()["index"], 0);
4424 assert_eq!(cache.get("volume_test:entity-500").unwrap()["index"], 500);
4425 assert_eq!(cache.get("volume_test:entity-999").unwrap()["index"], 999);
4426 }
4427
4428 #[tokio::test]
4429 async fn test_bulk_save_different_projections() {
4430 let store = create_test_store();
4431 let cache = store.projection_state_cache();
4432
4433 let projections = ["entity_snapshots", "event_counters", "custom_analytics"];
4435
4436 for proj in &projections {
4437 for i in 0..5 {
4438 cache.insert(
4439 format!("{proj}:entity-{i}"),
4440 serde_json::json!({"projection": proj, "id": i}),
4441 );
4442 }
4443 }
4444
4445 assert_eq!(cache.len(), 15);
4447
4448 for proj in &projections {
4450 let state = cache.get(&format!("{proj}:entity-0")).unwrap();
4451 assert_eq!(state["projection"], *proj);
4452 }
4453 }
4454
4455 #[tokio::test]
4466 async fn get_projection_state_falls_back_to_cache_when_unregistered() {
4467 let store = create_test_store();
4468 store.projection_state_cache().insert(
4469 "assets:BTC".to_string(),
4470 serde_json::json!({"symbol": "BTC", "altname": "Bitcoin"}),
4471 );
4472
4473 let resp = get_projection_state(
4474 State(Arc::clone(&store)),
4475 Path(("assets".to_string(), "BTC".to_string())),
4476 )
4477 .await
4478 .expect("should not error when projection is not registered");
4479
4480 assert_eq!(resp.0["found"], serde_json::Value::Bool(true));
4481 assert_eq!(resp.0["state"]["symbol"], "BTC");
4482 assert_eq!(resp.0["state"]["altname"], "Bitcoin");
4483 }
4484
4485 #[tokio::test]
4486 async fn get_projection_state_returns_not_found_when_absent_everywhere() {
4487 let store = create_test_store();
4488
4489 let resp = get_projection_state(
4490 State(Arc::clone(&store)),
4491 Path(("assets".to_string(), "UNKNOWN".to_string())),
4492 )
4493 .await
4494 .unwrap();
4495
4496 assert_eq!(resp.0["found"], serde_json::Value::Bool(false));
4497 assert_eq!(resp.0["state"], serde_json::Value::Null);
4498 }
4499
4500 #[tokio::test]
4501 async fn get_projection_state_registered_wins_over_cache() {
4502 let store = create_test_store();
4503
4504 let event = create_test_event("user-777", "user.created");
4506 store.ingest(&event).unwrap();
4507
4508 store.projection_state_cache().insert(
4510 "entity_snapshots:user-777".to_string(),
4511 serde_json::json!({"stolen": "value"}),
4512 );
4513
4514 let resp = get_projection_state(
4515 State(Arc::clone(&store)),
4516 Path(("entity_snapshots".to_string(), "user-777".to_string())),
4517 )
4518 .await
4519 .unwrap();
4520
4521 assert_eq!(resp.0["found"], serde_json::Value::Bool(true));
4524 assert!(
4525 resp.0["state"].get("stolen").is_none(),
4526 "cache entry must not shadow registered projection state: got {:?}",
4527 resp.0["state"]
4528 );
4529 }
4530
4531 #[tokio::test]
4532 async fn get_projection_state_summary_returns_cache_without_registration() {
4533 let store = create_test_store();
4534 let cache = store.projection_state_cache();
4535 cache.insert("assets:BTC".into(), serde_json::json!({"symbol": "BTC"}));
4536 cache.insert("assets:ETH".into(), serde_json::json!({"symbol": "ETH"}));
4537 cache.insert("trades:t-1".into(), serde_json::json!({"x": 1}));
4539
4540 let resp = get_projection_state_summary(
4541 State(Arc::clone(&store)),
4542 Path("assets".to_string()),
4543 Query(ProjectionStateSummaryParams::default()),
4544 )
4545 .await
4546 .unwrap();
4547
4548 assert_eq!(resp.0["total"], 2);
4549 let states = resp.0["states"].as_array().unwrap();
4550 let entity_ids: Vec<&str> = states
4551 .iter()
4552 .map(|s| s["entity_id"].as_str().unwrap())
4553 .collect();
4554 assert!(entity_ids.contains(&"BTC"));
4555 assert!(entity_ids.contains(&"ETH"));
4556 }
4557
4558 #[tokio::test]
4559 async fn bulk_get_projection_states_falls_back_to_cache() {
4560 let store = create_test_store();
4561 let cache = store.projection_state_cache();
4562 cache.insert("assets:BTC".into(), serde_json::json!({"symbol": "BTC"}));
4563 cache.insert("assets:ETH".into(), serde_json::json!({"symbol": "ETH"}));
4564
4565 let req = BulkGetStateRequest {
4566 entity_ids: vec!["BTC".into(), "ETH".into(), "MISSING".into()],
4567 };
4568
4569 let resp = bulk_get_projection_states(
4570 State(Arc::clone(&store)),
4571 Path("assets".to_string()),
4572 Json(req),
4573 )
4574 .await
4575 .unwrap();
4576
4577 assert_eq!(resp.0["total"], 3);
4578 let states = resp.0["states"].as_array().unwrap();
4579 let by_id: std::collections::HashMap<&str, &serde_json::Value> = states
4580 .iter()
4581 .map(|s| (s["entity_id"].as_str().unwrap(), s))
4582 .collect();
4583
4584 assert_eq!(by_id["BTC"]["found"], serde_json::Value::Bool(true));
4585 assert_eq!(by_id["BTC"]["state"]["symbol"], "BTC");
4586 assert_eq!(by_id["ETH"]["found"], serde_json::Value::Bool(true));
4587 assert_eq!(by_id["MISSING"]["found"], serde_json::Value::Bool(false));
4588 }
4589
4590 #[tokio::test]
4594 async fn poll_consumer_events_flattens_event_alongside_position() {
4595 let store = create_test_store();
4596 store
4597 .ingest(&create_test_event("user-1", "user.created"))
4598 .unwrap();
4599 store
4600 .ingest(&create_test_event("user-2", "user.updated"))
4601 .unwrap();
4602 store.consumer_registry().register("w1", &[]);
4603
4604 let resp = poll_consumer_events(
4605 State(Arc::clone(&store)),
4606 Path("w1".to_string()),
4607 Query(ConsumerPollQuery { limit: Some(10) }),
4608 )
4609 .await
4610 .unwrap();
4611
4612 let body = serde_json::to_value(&resp.0).unwrap();
4613 assert_eq!(body["count"], 2);
4614 let first = &body["events"][0];
4615 assert_eq!(first["position"], 1);
4616 assert!(
4617 first.get("event").is_none(),
4618 "event must be flattened, not nested: got {first:?}"
4619 );
4620 assert_eq!(first["event_type"], "user.created");
4621 assert_eq!(first["entity_id"], "user-1");
4622 }
4623}