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