Skip to main content

allsource_core/infrastructure/web/
api.rs

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/// Wait for follower ACK(s) in semi-sync/sync replication modes.
64///
65/// In async mode (default), returns immediately. In semi-sync mode, waits for
66/// at least 1 follower to ACK the current WAL offset. In sync mode, waits for
67/// all followers. If the timeout expires, logs a warning and continues (degraded mode).
68#[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 the read guard before the async wait to avoid holding it across await
84        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    // No-op in community edition
107}
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)) // v0.6: Prometheus metrics endpoint
113        .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)) // v0.2: WebSocket streaming
118        // v0.10: Stream and event type discovery endpoints
119        .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        // v0.2: Advanced analytics endpoints
129        .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        // v0.2: Snapshot management endpoints
133        .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        // v0.2: Compaction endpoints
140        .route("/api/v1/compaction/trigger", post(trigger_compaction))
141        .route("/api/v1/compaction/stats", get(compaction_stats))
142        // v0.5: Schema registry endpoints
143        .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        // v0.5: Replay and projection rebuild endpoints
156        .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        // v0.5: Stream processing pipeline endpoints
165        .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        // v0.7: Projection State API for Query Service integration
179        .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        // v0.11: Webhook management endpoints
211        .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        // v2.0: Advanced query features
224        .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
263// v0.6: Prometheus metrics endpoint
264pub async fn prometheus_metrics(State(store): State<SharedStore>) -> impl IntoResponse {
265    let metrics = store.metrics();
266
267    match metrics.encode() {
268        Ok(encoded) => Response::builder()
269            .status(200)
270            .header("Content-Type", "text/plain; version=0.0.4")
271            .body(encoded)
272            .unwrap()
273            .into_response(),
274        Err(e) => Response::builder()
275            .status(500)
276            .body(format!("Error encoding metrics: {e}"))
277            .unwrap()
278            .into_response(),
279    }
280}
281
282pub async fn ingest_event(
283    State(store): State<SharedStore>,
284    Json(req): Json<IngestEventRequest>,
285) -> Result<Json<IngestEventResponse>> {
286    let expected_version = req.expected_version;
287
288    let tenant_id = req.tenant_id.unwrap_or_else(|| "default".to_string());
289    let event = Event::from_strings(
290        req.event_type,
291        req.entity_id,
292        tenant_id,
293        req.payload,
294        req.metadata,
295    )?;
296
297    let event_id = event.id;
298    let timestamp = event.timestamp;
299
300    let new_version =
301        super::archive_work::append(Arc::clone(&store), event, expected_version).await?;
302
303    tracing::info!("Event ingested: {}", event_id);
304
305    Ok(Json(IngestEventResponse {
306        event_id,
307        timestamp,
308        version: Some(new_version),
309    }))
310}
311
312/// Ingest a single event with semi-sync/sync replication ACK waiting.
313///
314/// Used by the v1 API. Tenant comes from `req.tenant_id` — the Control Plane
315/// delegation layer sets it from the authenticated caller before forwarding.
316/// Core is internal-only and does not authenticate public traffic.
317/// Look up the tenant's `SchemaEnforcement` mode and, if non-permissive,
318/// validate the event payload against any registered schema for the
319/// event_type.
320///
321/// Fast path: `Permissive` (default for unconfigured tenants AND for tenants
322/// not present in the repo, like dev/test setups) returns `Ok(())` without
323/// touching the schema registry — preserves the pre-v0.21.5 ingest cost.
324///
325/// `Warn`: validation runs, violations log at WARN, the write proceeds.
326/// `Strict`: violations return `AllSourceError::SchemaViolation` → 422 with
327/// the structured body the HTTP layer builds in `error.rs`.
328///
329/// If no schema is registered for the event_type, validation is a no-op
330/// regardless of mode — this matches the principle that schemas are
331/// opt-in per event_type, not per tenant.
332async fn enforce_schema_if_configured(
333    state: &AppState,
334    tenant_id: &str,
335    event: &Event,
336) -> Result<()> {
337    // Cheapest possible lookup: parse the tenant_id; if it fails, treat as
338    // permissive (defensive — Event::from_strings already validated, but
339    // this keeps the contract clear).
340    let Ok(parsed) = TenantId::new(tenant_id.to_string()) else {
341        return Ok(());
342    };
343    let mode = match state.tenant_repo.find_by_id(&parsed).await {
344        Ok(Some(t)) => t.schema_enforcement(),
345        // Unknown tenant → permissive. Avoids breaking the default tenant
346        // and any dev setups that ingest without pre-registering tenants.
347        _ => SchemaEnforcement::Permissive,
348    };
349    if matches!(mode, SchemaEnforcement::Permissive) {
350        return Ok(());
351    }
352
353    // Schema lookup keyed by event_type. Latest version only (None) — schema
354    // evolution is a separate concern; tenants pin a version via their own
355    // registration flow if they need to.
356    let registry = state.store.schema_registry();
357    let Ok(schema) = registry.get_schema(event.event_type.as_str(), None) else {
358        // No schema registered for this event_type → fast path, regardless
359        // of enforcement mode. Schemas are opt-in.
360        return Ok(());
361    };
362
363    let result = registry
364        .validate(
365            event.event_type.as_str(),
366            Some(schema.version),
367            &event.payload,
368        )
369        .map_err(|e| crate::error::AllSourceError::InternalError(e.to_string()))?;
370
371    if result.valid {
372        return Ok(());
373    }
374
375    match mode {
376        SchemaEnforcement::Strict => Err(crate::error::AllSourceError::SchemaViolation {
377            event_type: event.event_type.as_str().to_string(),
378            schema_version: result.schema_version,
379            errors: result.errors,
380        }),
381        SchemaEnforcement::Warn => {
382            tracing::warn!(
383                tenant = %tenant_id,
384                event_type = %event.event_type.as_str(),
385                schema_version = result.schema_version,
386                errors = ?result.errors,
387                "schema violation (warn mode — write accepted)"
388            );
389            Ok(())
390        }
391        // Already handled above
392        SchemaEnforcement::Permissive => Ok(()),
393    }
394}
395
396pub async fn ingest_event_v1(
397    State(state): State<AppState>,
398    Json(req): Json<IngestEventRequest>,
399) -> Result<Json<IngestEventResponse>> {
400    let expected_version = req.expected_version;
401
402    let tenant_id = req.tenant_id.unwrap_or_else(|| "default".to_string());
403
404    let event = Event::from_strings(
405        req.event_type,
406        req.entity_id,
407        tenant_id.clone(),
408        req.payload,
409        req.metadata,
410    )?;
411
412    // Per-tenant schema enforcement (Permissive is the fast path — no
413    // tenant lookup, no schema query). See neotoma-gaps bead t-0795.
414    enforce_schema_if_configured(&state, &tenant_id, &event).await?;
415
416    let event_id = event.id;
417    let timestamp = event.timestamp;
418
419    let new_version =
420        super::archive_work::append(Arc::clone(&state.store), event, expected_version).await?;
421
422    // Semi-sync/sync: wait for follower ACK(s) before returning
423    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
434/// Batch ingest multiple events in a single request
435///
436/// This endpoint allows ingesting multiple events atomically, which is more
437/// efficient than making individual requests for each event.
438pub 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        // Note: the non-v1 batch path uses State<SharedStore>, not AppState,
458        // so it has no tenant_repo to check enforcement against. Schema
459        // enforcement is wired through the v1 batch handler below — this
460        // path remains permissive (it predates the tenant_repo wiring).
461
462        let event_id = event.id;
463        let timestamp = event.timestamp;
464
465        let new_version =
466            super::archive_work::append(Arc::clone(&store), event, expected_version).await?;
467
468        ingested_events.push(IngestEventResponse {
469            event_id,
470            timestamp,
471            version: Some(new_version),
472        });
473    }
474
475    let ingested = ingested_events.len();
476    tracing::info!("Batch ingested {} events", ingested);
477
478    Ok(Json(IngestEventsBatchResponse {
479        total,
480        ingested,
481        events: ingested_events,
482    }))
483}
484
485/// Batch ingest with semi-sync/sync replication ACK waiting.
486///
487/// Used by the v1 API. Per-event tenant comes from `event_req.tenant_id` — the
488/// Control Plane delegation layer sets it from the authenticated caller before
489/// forwarding. Core is internal-only and does not authenticate public traffic.
490pub async fn ingest_events_batch_v1(
491    State(state): State<AppState>,
492    Json(req): Json<IngestEventsBatchRequest>,
493) -> Result<Json<IngestEventsBatchResponse>> {
494    let total = req.events.len();
495    let mut ingested_events = Vec::with_capacity(total);
496
497    for event_req in req.events {
498        let tenant_id = event_req.tenant_id.unwrap_or_else(|| "default".to_string());
499        let expected_version = event_req.expected_version;
500
501        let event = Event::from_strings(
502            event_req.event_type,
503            event_req.entity_id,
504            tenant_id.clone(),
505            event_req.payload,
506            event_req.metadata,
507        )?;
508
509        enforce_schema_if_configured(&state, &tenant_id, &event).await?;
510
511        let event_id = event.id;
512        let timestamp = event.timestamp;
513
514        let new_version =
515            super::archive_work::append(Arc::clone(&state.store), event, expected_version).await?;
516
517        ingested_events.push(IngestEventResponse {
518            event_id,
519            timestamp,
520            version: Some(new_version),
521        });
522    }
523
524    // Semi-sync/sync: wait for follower ACK(s) after all events are ingested
525    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/// Sort-order query parameter for `GET /api/v1/events/query`. Kept separate
538/// from `QueryEventsRequest` so ordering is purely an HTTP-layer concern and
539/// the internal query DTO (used by ~50 call sites) stays untouched.
540#[derive(Debug, Deserialize)]
541pub struct EventOrderParam {
542    /// `asc` (oldest first, the default) or `desc` (newest first).
543    pub order: Option<String>,
544}
545
546/// Pagination offset for `GET /api/v1/events/query`. Separate from
547/// `QueryEventsRequest` for the same reason as `EventOrderParam`: offset is an
548/// HTTP-layer windowing concern, not part of the internal query predicate that
549/// `store.query()` and its ~50 call sites evaluate.
550#[derive(Debug, Deserialize)]
551pub struct EventOffsetParam {
552    /// Number of matching events to skip before applying `limit`. Default 0.
553    pub offset: Option<usize>,
554}
555
556/// Opt-in integrity protocol. Omission preserves the legacy tolerant query.
557#[derive(Debug, Default, Deserialize)]
558pub struct EventIntegrityParam {
559    pub integrity: Option<String>,
560}
561
562pub async fn query_events(
563    OptionalAuth(auth): OptionalAuth,
564    Query(req): Query<QueryEventsRequest>,
565    Query(order_param): Query<EventOrderParam>,
566    Query(offset_param): Query<EventOffsetParam>,
567    Query(integrity_param): Query<EventIntegrityParam>,
568    State(store): State<SharedStore>,
569) -> Result<Json<QueryEventsResponse>> {
570    let offset = offset_param.offset.unwrap_or(0);
571    let queried_entity_id = req.entity_id.clone();
572
573    // Sort order. Default is ascending (oldest first) to preserve replay
574    // semantics for existing consumers. `order=desc` returns newest first,
575    // so `?entity_id=<id>&limit=1&order=desc` reliably yields the latest
576    // event for an entity (issue #177).
577    let descending = match order_param.order.as_deref() {
578        None => false,
579        Some(o) if o.eq_ignore_ascii_case("asc") => false,
580        Some(o) if o.eq_ignore_ascii_case("desc") => true,
581        Some(other) => {
582            return Err(crate::error::AllSourceError::InvalidInput(format!(
583                "invalid 'order' value '{other}': expected 'asc' or 'desc'"
584            )));
585        }
586    };
587
588    // Tenant resolution. Since Core is internal-only (bead t-0ff8), the only
589    // callers are Control Plane's delegation layer and other internal Fly
590    // services. Request param wins — that's what the gateway sets authoritatively
591    // from the authenticated caller's identity. Auth-context fallback is kept
592    // as a defense-in-depth for any internal caller that forgets to pass
593    // tenant_id but is authenticated (legacy; audit and remove once all
594    // internal callers are confirmed to set tenant_id explicitly).
595    let enforced_tenant = req
596        .tenant_id
597        .clone()
598        .or_else(|| auth.as_ref().map(|a| a.tenant_id().to_string()));
599
600    if let Some(protocol) = integrity_param.integrity {
601        return super::retained_query::query(
602            store,
603            QueryEventsRequest {
604                tenant_id: enforced_tenant,
605                ..req
606            },
607            offset,
608            descending,
609            &protocol,
610        )
611        .await
612        .map(Json);
613    }
614
615    // FAIL CLOSED (tenant isolation): the public events query must NEVER return
616    // cross-tenant results. The gateway always injects an auth-derived tenant_id
617    // (and overwrites any client-supplied one); if neither a request tenant nor
618    // an auth tenant is present we return an empty result rather than scanning
619    // across tenants. A genuine cross-tenant/admin scan is a separate, explicit
620    // internal path — it does not ride this endpoint.
621    if enforced_tenant.as_deref().unwrap_or("").is_empty() {
622        return Ok(Json(QueryEventsResponse {
623            events: Vec::new(),
624            count: 0,
625            total_count: 0,
626            has_more: false,
627            entity_version: None,
628            archive_integrity: None,
629        }));
630    }
631
632    // One windowed pass: the store sorts borrowed matches, applies `offset` and
633    // `limit`, and clones only the page — while still reporting the pre-window
634    // match count for `total_count`/`has_more`. Asking for the total used to
635    // mean a second, unlimited query that cloned the whole history, so
636    // `?entity_id=X&limit=1&order=desc` cost as much as fetching everything
637    // (issue #251). Offset is applied before limit so `offset=N&limit=N` walks
638    // pages instead of returning page one forever (issue #250).
639    let scoped_req = QueryEventsRequest {
640        tenant_id: enforced_tenant,
641        ..req
642    };
643    let (limited_events, total_count) = store.query_window(&scoped_req, offset, descending)?;
644
645    let count = limited_events.len();
646    // `has_more` is relative to the window actually served, not to the page
647    // size — a paginator that trusts a bare `count < total_count` never
648    // terminates once an offset is in play.
649    let has_more = offset + count < total_count;
650    let events: Vec<EventDto> = limited_events.iter().map(EventDto::from).collect();
651
652    // Include entity_version only when filtering by a single entity_id
653    let entity_version = queried_entity_id
654        .as_deref()
655        .map(|eid| store.get_entity_version(eid));
656
657    tracing::debug!("Query returned {} events (total: {})", count, total_count);
658
659    Ok(Json(QueryEventsResponse {
660        events,
661        count,
662        total_count,
663        has_more,
664        entity_version,
665        archive_integrity: None,
666    }))
667}
668
669pub async fn list_entities(
670    State(store): State<SharedStore>,
671    Query(req): Query<ListEntitiesRequest>,
672) -> Result<Json<ListEntitiesResponse>> {
673    use std::collections::HashMap;
674
675    // Get all events matching the filters
676    let query_req = QueryEventsRequest {
677        entity_id: None,
678        event_type: None,
679        tenant_id: None,
680        as_of: None,
681        since: None,
682        until: None,
683        limit: None,
684        event_type_prefix: req.event_type_prefix,
685        exclude_event_type_prefix: None,
686        payload_filter: req.payload_filter,
687    };
688    let events = store.query(&query_req)?;
689
690    // Group by entity_id
691    let mut entity_map: HashMap<String, Vec<&Event>> = HashMap::new();
692    for event in &events {
693        entity_map
694            .entry(event.entity_id().to_string())
695            .or_default()
696            .push(event);
697    }
698
699    // Sort direction by last-event time. Default is `desc` (newest activity
700    // first) to preserve the endpoint's long-standing behavior; `order=asc`
701    // returns oldest activity first (issue #178).
702    let ascending = match req.order.as_deref() {
703        None => false,
704        Some(o) if o.eq_ignore_ascii_case("desc") => false,
705        Some(o) if o.eq_ignore_ascii_case("asc") => true,
706        Some(other) => {
707            return Err(crate::error::AllSourceError::InvalidInput(format!(
708                "invalid 'order' value '{other}': expected 'asc' or 'desc'"
709            )));
710        }
711    };
712
713    // Build entity summaries, then sort by last-event time in the requested
714    // direction. Ties are broken by `entity_id` (always ascending) so the
715    // total order is deterministic — required for stable offset pagination.
716    let mut summaries: Vec<EntitySummary> = entity_map
717        .into_iter()
718        .map(|(entity_id, events)| {
719            let last = events.iter().max_by_key(|e| e.timestamp()).unwrap();
720            EntitySummary {
721                entity_id,
722                event_count: events.len(),
723                last_event_type: last.event_type_str().to_string(),
724                last_event_at: last.timestamp(),
725            }
726        })
727        .collect();
728    summaries.sort_by(|a, b| {
729        let by_time = a.last_event_at.cmp(&b.last_event_at);
730        let by_time = if ascending {
731            by_time
732        } else {
733            by_time.reverse()
734        };
735        by_time.then_with(|| a.entity_id.cmp(&b.entity_id))
736    });
737
738    let total = summaries.len();
739
740    // Apply offset and limit
741    let offset = req.offset.unwrap_or(0);
742    let summaries: Vec<EntitySummary> = summaries.into_iter().skip(offset).collect::<Vec<_>>();
743    let summaries = if let Some(limit) = req.limit {
744        let has_more = summaries.len() > limit;
745        let truncated: Vec<EntitySummary> = summaries.into_iter().take(limit).collect();
746        return Ok(Json(ListEntitiesResponse {
747            entities: truncated,
748            total,
749            has_more,
750        }));
751    } else {
752        summaries
753    };
754
755    Ok(Json(ListEntitiesResponse {
756        entities: summaries,
757        total,
758        has_more: false,
759    }))
760}
761
762pub async fn detect_duplicates(
763    State(store): State<SharedStore>,
764    Query(req): Query<DetectDuplicatesRequest>,
765) -> Result<Json<DetectDuplicatesResponse>> {
766    use std::collections::HashMap;
767
768    let group_by_fields: Vec<&str> = req.group_by.split(',').map(str::trim).collect();
769
770    // Query events scoped by the required prefix
771    let query_req = QueryEventsRequest {
772        entity_id: None,
773        event_type: None,
774        tenant_id: None,
775        as_of: None,
776        since: None,
777        until: None,
778        limit: None,
779        event_type_prefix: Some(req.event_type_prefix),
780        exclude_event_type_prefix: None,
781        payload_filter: None,
782    };
783    let events = store.query(&query_req)?;
784
785    // For each entity, extract the latest event's payload fields specified by group_by
786    // Then group entities by those field values
787    let mut entity_latest: HashMap<String, &Event> = HashMap::new();
788    for event in &events {
789        let eid = event.entity_id().to_string();
790        entity_latest
791            .entry(eid)
792            .and_modify(|existing| {
793                if event.timestamp() > existing.timestamp() {
794                    *existing = event;
795                }
796            })
797            .or_insert(event);
798    }
799
800    // Group entities by their payload field values
801    let mut groups: HashMap<String, Vec<String>> = HashMap::new();
802    for (entity_id, event) in &entity_latest {
803        let payload = event.payload();
804        let mut key_parts = serde_json::Map::new();
805        for field in &group_by_fields {
806            let value = payload
807                .get(*field)
808                .cloned()
809                .unwrap_or(serde_json::Value::Null);
810            key_parts.insert((*field).to_string(), value);
811        }
812        let key_str = serde_json::to_string(&key_parts).unwrap_or_default();
813        groups.entry(key_str).or_default().push(entity_id.clone());
814    }
815
816    // Filter to groups with count > 1 (actual duplicates)
817    let mut duplicate_groups: Vec<DuplicateGroup> = groups
818        .into_iter()
819        .filter(|(_, ids)| ids.len() > 1)
820        .map(|(key_str, mut ids)| {
821            ids.sort();
822            let key: serde_json::Value =
823                serde_json::from_str(&key_str).unwrap_or(serde_json::Value::Null);
824            let count = ids.len();
825            DuplicateGroup {
826                key,
827                entity_ids: ids,
828                count,
829            }
830        })
831        .collect();
832
833    // Sort by count descending for consistent output
834    duplicate_groups.sort_by_key(|a| std::cmp::Reverse(a.count));
835
836    let total = duplicate_groups.len();
837
838    // Apply offset and limit
839    let offset = req.offset.unwrap_or(0);
840    let duplicate_groups: Vec<DuplicateGroup> = duplicate_groups.into_iter().skip(offset).collect();
841
842    if let Some(limit) = req.limit {
843        let has_more = duplicate_groups.len() > limit;
844        let truncated: Vec<DuplicateGroup> = duplicate_groups.into_iter().take(limit).collect();
845        return Ok(Json(DetectDuplicatesResponse {
846            duplicates: truncated,
847            total,
848            has_more,
849        }));
850    }
851
852    Ok(Json(DetectDuplicatesResponse {
853        duplicates: duplicate_groups,
854        total,
855        has_more: false,
856    }))
857}
858
859#[derive(Deserialize)]
860pub struct EntityStateParams {
861    as_of: Option<chrono::DateTime<chrono::Utc>>,
862    /// Tenant to scope to (the gateway injects the authenticated tenant).
863    ///
864    /// When present, reconstruction folds only that tenant's events and skips
865    /// the snapshot fast path, which carries no tenant dimension (#230).
866    tenant_id: Option<String>,
867}
868
869pub async fn get_entity_state(
870    State(store): State<SharedStore>,
871    Path(entity_id): Path<String>,
872    Query(params): Query<EntityStateParams>,
873) -> Result<Json<serde_json::Value>> {
874    let state = match params.tenant_id.as_deref() {
875        Some(tenant_id) => {
876            store.reconstruct_state_for_tenant(&entity_id, params.as_of, tenant_id)?
877        }
878        None => store.reconstruct_state(&entity_id, params.as_of)?,
879    };
880
881    tracing::info!("State reconstructed for entity: {}", entity_id);
882
883    Ok(Json(state))
884}
885
886pub async fn get_entity_snapshot(
887    State(store): State<SharedStore>,
888    Path(entity_id): Path<String>,
889    Query(params): Query<EntityStateParams>,
890) -> Result<Json<serde_json::Value>> {
891    // Snapshots are keyed by entity_id with no tenant dimension, so a scoped
892    // caller cannot be served one safely — two tenants sharing an entity_id
893    // would see each other's state. Fold that tenant's events instead (#230).
894    let snapshot = match params.tenant_id.as_deref() {
895        Some(tenant_id) => store.reconstruct_state_for_tenant(&entity_id, None, tenant_id)?,
896        None => store.get_snapshot(&entity_id)?,
897    };
898
899    tracing::debug!("Snapshot retrieved for entity: {}", entity_id);
900
901    Ok(Json(snapshot))
902}
903
904/// Query parameters for the stats endpoint.
905#[derive(Debug, Deserialize)]
906pub struct StatsParams {
907    /// Tenant to scope to (the gateway injects the authenticated tenant).
908    ///
909    /// Absent = global, whole-store totals. That form is internal/admin only and
910    /// must never be reachable by a tenant — the gateway routes its public
911    /// `GET /api/v1/stats` through here with this parameter forced (#230).
912    pub tenant_id: Option<String>,
913}
914
915pub async fn get_stats(
916    State(store): State<SharedStore>,
917    Query(params): Query<StatsParams>,
918) -> impl IntoResponse {
919    match params.tenant_id.as_deref() {
920        Some(tenant_id) => {
921            Json(serde_json::to_value(store.stats_for_tenant(tenant_id)).unwrap_or_default())
922        }
923        None => Json(serde_json::to_value(store.stats()).unwrap_or_default()),
924    }
925}
926
927// v0.10: List all streams (entity_ids) in the event store
928/// Query parameters for listing streams
929#[derive(Debug, Deserialize)]
930pub struct ListStreamsParams {
931    /// Tenant to scope to (the gateway injects the authenticated tenant).
932    pub tenant_id: Option<String>,
933    /// Optional limit on number of streams to return
934    pub limit: Option<usize>,
935    /// Optional offset for pagination
936    pub offset: Option<usize>,
937}
938
939/// Response for listing streams
940#[derive(Debug, serde::Serialize)]
941pub struct ListStreamsResponse {
942    pub streams: Vec<StreamInfo>,
943    pub total: usize,
944}
945
946pub async fn list_streams(
947    OptionalAuth(auth): OptionalAuth,
948    State(store): State<SharedStore>,
949    Query(params): Query<ListStreamsParams>,
950) -> Json<ListStreamsResponse> {
951    // Tenant-scoped + fail closed: no tenant → empty, never a cross-tenant list.
952    let tenant = params
953        .tenant_id
954        .clone()
955        .or_else(|| auth.as_ref().map(|a| a.tenant_id().to_string()))
956        .filter(|t| !t.is_empty());
957    let Some(tenant) = tenant else {
958        return Json(ListStreamsResponse {
959            streams: vec![],
960            total: 0,
961        });
962    };
963    let mut streams = store.list_streams_for_tenant(&tenant);
964    let total = streams.len();
965
966    // Sort by last_event_at descending (most recent first)
967    streams.sort_by_key(|a| std::cmp::Reverse(a.last_event_at));
968
969    // Apply pagination
970    if let Some(offset) = params.offset {
971        if offset < streams.len() {
972            streams = streams[offset..].to_vec();
973        } else {
974            streams = vec![];
975        }
976    }
977
978    if let Some(limit) = params.limit {
979        streams.truncate(limit);
980    }
981
982    tracing::debug!("Listed {} streams (total: {})", streams.len(), total);
983
984    Json(ListStreamsResponse { streams, total })
985}
986
987// v0.10: List all event types in the event store
988/// Query parameters for listing event types
989#[derive(Debug, Deserialize)]
990pub struct ListEventTypesParams {
991    /// Tenant to scope to (the gateway injects the authenticated tenant).
992    pub tenant_id: Option<String>,
993    /// Optional limit on number of event types to return
994    pub limit: Option<usize>,
995    /// Optional offset for pagination
996    pub offset: Option<usize>,
997}
998
999/// Response for listing event types
1000#[derive(Debug, serde::Serialize)]
1001pub struct ListEventTypesResponse {
1002    pub event_types: Vec<EventTypeInfo>,
1003    pub total: usize,
1004}
1005
1006pub async fn list_event_types(
1007    OptionalAuth(auth): OptionalAuth,
1008    State(store): State<SharedStore>,
1009    Query(params): Query<ListEventTypesParams>,
1010) -> Json<ListEventTypesResponse> {
1011    // Tenant-scoped + fail closed: no tenant → empty, never cross-tenant types.
1012    let tenant = params
1013        .tenant_id
1014        .clone()
1015        .or_else(|| auth.as_ref().map(|a| a.tenant_id().to_string()))
1016        .filter(|t| !t.is_empty());
1017    let Some(tenant) = tenant else {
1018        return Json(ListEventTypesResponse {
1019            event_types: vec![],
1020            total: 0,
1021        });
1022    };
1023    let mut event_types = store.list_event_types_for_tenant(&tenant);
1024    let total = event_types.len();
1025
1026    // Sort by event_count descending (most used first)
1027    event_types.sort_by_key(|a| std::cmp::Reverse(a.event_count));
1028
1029    // Apply pagination
1030    if let Some(offset) = params.offset {
1031        if offset < event_types.len() {
1032            event_types = event_types[offset..].to_vec();
1033        } else {
1034            event_types = vec![];
1035        }
1036    }
1037
1038    if let Some(limit) = params.limit {
1039        event_types.truncate(limit);
1040    }
1041
1042    tracing::debug!(
1043        "Listed {} event types (total: {})",
1044        event_types.len(),
1045        total
1046    );
1047
1048    Json(ListEventTypesResponse { event_types, total })
1049}
1050
1051// v0.2: WebSocket endpoint for real-time event streaming
1052#[derive(Debug, Deserialize)]
1053pub struct WebSocketParams {
1054    pub consumer_id: Option<String>,
1055}
1056
1057pub async fn events_websocket(
1058    ws: WebSocketUpgrade,
1059    State(store): State<SharedStore>,
1060    Query(params): Query<WebSocketParams>,
1061) -> Response {
1062    let websocket_manager = store.websocket_manager();
1063
1064    ws.on_upgrade(move |socket| async move {
1065        if let Some(consumer_id) = params.consumer_id {
1066            websocket_manager
1067                .handle_socket_with_consumer(socket, consumer_id, store)
1068                .await;
1069        } else {
1070            websocket_manager.handle_socket(socket).await;
1071        }
1072    })
1073}
1074
1075// v0.2: Event frequency analytics endpoint
1076pub async fn analytics_frequency(
1077    State(store): State<SharedStore>,
1078    Query(req): Query<EventFrequencyRequest>,
1079) -> Result<Json<EventFrequencyResponse>> {
1080    let response = AnalyticsEngine::event_frequency(&store, &req)?;
1081
1082    tracing::debug!(
1083        "Frequency analysis returned {} buckets",
1084        response.buckets.len()
1085    );
1086
1087    Ok(Json(response))
1088}
1089
1090// v0.2: Statistical summary endpoint
1091pub async fn analytics_summary(
1092    State(store): State<SharedStore>,
1093    Query(req): Query<StatsSummaryRequest>,
1094) -> Result<Json<StatsSummaryResponse>> {
1095    let response = AnalyticsEngine::stats_summary(&store, &req)?;
1096
1097    tracing::debug!(
1098        "Stats summary: {} events across {} entities",
1099        response.total_events,
1100        response.unique_entities
1101    );
1102
1103    Ok(Json(response))
1104}
1105
1106// v0.2: Event correlation analysis endpoint
1107pub async fn analytics_correlation(
1108    State(store): State<SharedStore>,
1109    Query(req): Query<CorrelationRequest>,
1110) -> Result<Json<CorrelationResponse>> {
1111    let response = AnalyticsEngine::analyze_correlation(&store, req)?;
1112
1113    tracing::debug!(
1114        "Correlation analysis: {}/{} correlated pairs ({:.2}%)",
1115        response.correlated_pairs,
1116        response.total_a,
1117        response.correlation_percentage
1118    );
1119
1120    Ok(Json(response))
1121}
1122
1123// v0.2: Create a snapshot for an entity
1124pub async fn create_snapshot(
1125    State(store): State<SharedStore>,
1126    Json(req): Json<CreateSnapshotRequest>,
1127) -> Result<Json<CreateSnapshotResponse>> {
1128    store.create_snapshot(&req.entity_id)?;
1129
1130    let snapshot_manager = store.snapshot_manager();
1131    let snapshot = snapshot_manager
1132        .get_latest_snapshot(&req.entity_id)
1133        .ok_or_else(|| crate::error::AllSourceError::EntityNotFound(req.entity_id.clone()))?;
1134
1135    tracing::info!("📸 Created snapshot for entity: {}", req.entity_id);
1136
1137    Ok(Json(CreateSnapshotResponse {
1138        snapshot_id: snapshot.id,
1139        entity_id: snapshot.entity_id,
1140        created_at: snapshot.created_at,
1141        event_count: snapshot.event_count,
1142        size_bytes: snapshot.metadata.size_bytes,
1143    }))
1144}
1145
1146// v0.2: List snapshots
1147pub async fn list_snapshots(
1148    State(store): State<SharedStore>,
1149    Query(req): Query<ListSnapshotsRequest>,
1150) -> Result<Json<ListSnapshotsResponse>> {
1151    let snapshot_manager = store.snapshot_manager();
1152
1153    let snapshots: Vec<SnapshotInfo> = if let Some(entity_id) = req.entity_id {
1154        snapshot_manager
1155            .get_all_snapshots(&entity_id)
1156            .into_iter()
1157            .map(SnapshotInfo::from)
1158            .collect()
1159    } else {
1160        // List all entities with snapshots
1161        let entities = snapshot_manager.list_entities();
1162        entities
1163            .iter()
1164            .flat_map(|entity_id| {
1165                snapshot_manager
1166                    .get_all_snapshots(entity_id)
1167                    .into_iter()
1168                    .map(SnapshotInfo::from)
1169            })
1170            .collect()
1171    };
1172
1173    let total = snapshots.len();
1174
1175    tracing::debug!("Listed {} snapshots", total);
1176
1177    Ok(Json(ListSnapshotsResponse { snapshots, total }))
1178}
1179
1180// v0.2: Get latest snapshot for an entity
1181pub async fn get_latest_snapshot(
1182    State(store): State<SharedStore>,
1183    Path(entity_id): Path<String>,
1184) -> Result<Json<serde_json::Value>> {
1185    let snapshot_manager = store.snapshot_manager();
1186
1187    let snapshot = snapshot_manager
1188        .get_latest_snapshot(&entity_id)
1189        .ok_or_else(|| crate::error::AllSourceError::EntityNotFound(entity_id.clone()))?;
1190
1191    tracing::debug!("Retrieved latest snapshot for entity: {}", entity_id);
1192
1193    Ok(Json(serde_json::json!({
1194        "snapshot_id": snapshot.id,
1195        "entity_id": snapshot.entity_id,
1196        "created_at": snapshot.created_at,
1197        "as_of": snapshot.as_of,
1198        "event_count": snapshot.event_count,
1199        "size_bytes": snapshot.metadata.size_bytes,
1200        "snapshot_type": snapshot.metadata.snapshot_type,
1201        "state": snapshot.state
1202    })))
1203}
1204
1205// v0.2: Trigger manual compaction
1206pub async fn trigger_compaction(
1207    State(store): State<SharedStore>,
1208) -> Result<Json<CompactionResult>> {
1209    let compaction_manager = store.compaction_manager().ok_or_else(|| {
1210        crate::error::AllSourceError::InternalError(
1211            "Compaction not enabled (no Parquet storage)".to_string(),
1212        )
1213    })?;
1214
1215    tracing::info!("📦 Manual compaction triggered via API");
1216
1217    let result = compaction_manager.compact_now()?;
1218
1219    Ok(Json(result))
1220}
1221
1222// v0.2: Get compaction statistics
1223pub async fn compaction_stats(State(store): State<SharedStore>) -> Result<Json<serde_json::Value>> {
1224    let compaction_manager = store.compaction_manager().ok_or_else(|| {
1225        crate::error::AllSourceError::InternalError(
1226            "Compaction not enabled (no Parquet storage)".to_string(),
1227        )
1228    })?;
1229
1230    let stats = compaction_manager.stats();
1231    let config = compaction_manager.config();
1232
1233    Ok(Json(serde_json::json!({
1234        "stats": stats,
1235        "config": {
1236            "min_files_to_compact": config.min_files_to_compact,
1237            "target_file_size": config.target_file_size,
1238            "max_file_size": config.max_file_size,
1239            "small_file_threshold": config.small_file_threshold,
1240            "compaction_interval_seconds": config.compaction_interval_seconds,
1241            "auto_compact": config.auto_compact,
1242            "strategy": config.strategy
1243        }
1244    })))
1245}
1246
1247// v0.5: Register a new schema
1248pub async fn register_schema(
1249    State(store): State<SharedStore>,
1250    Json(req): Json<RegisterSchemaRequest>,
1251) -> Result<Json<RegisterSchemaResponse>> {
1252    let schema_registry = store.schema_registry();
1253
1254    let response =
1255        schema_registry.register_schema(req.subject, req.schema, req.description, req.tags)?;
1256
1257    tracing::info!(
1258        "📋 Schema registered: v{} for '{}'",
1259        response.version,
1260        response.subject
1261    );
1262
1263    Ok(Json(response))
1264}
1265
1266// v0.5: Get a schema by subject and optional version
1267#[derive(Deserialize)]
1268pub struct GetSchemaParams {
1269    version: Option<u32>,
1270}
1271
1272pub async fn get_schema(
1273    State(store): State<SharedStore>,
1274    Path(subject): Path<String>,
1275    Query(params): Query<GetSchemaParams>,
1276) -> Result<Json<serde_json::Value>> {
1277    let schema_registry = store.schema_registry();
1278
1279    let schema = schema_registry.get_schema(&subject, params.version)?;
1280
1281    tracing::debug!("Retrieved schema v{} for '{}'", schema.version, subject);
1282
1283    Ok(Json(serde_json::json!({
1284        "id": schema.id,
1285        "subject": schema.subject,
1286        "version": schema.version,
1287        "schema": schema.schema,
1288        "created_at": schema.created_at,
1289        "description": schema.description,
1290        "tags": schema.tags
1291    })))
1292}
1293
1294// v0.5: List all versions of a schema subject
1295pub async fn list_schema_versions(
1296    State(store): State<SharedStore>,
1297    Path(subject): Path<String>,
1298) -> Result<Json<serde_json::Value>> {
1299    let schema_registry = store.schema_registry();
1300
1301    let versions = schema_registry.list_versions(&subject)?;
1302
1303    Ok(Json(serde_json::json!({
1304        "subject": subject,
1305        "versions": versions
1306    })))
1307}
1308
1309// v0.5: List all schema subjects
1310pub async fn list_subjects(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1311    let schema_registry = store.schema_registry();
1312
1313    let subjects = schema_registry.list_subjects();
1314
1315    Json(serde_json::json!({
1316        "subjects": subjects,
1317        "total": subjects.len()
1318    }))
1319}
1320
1321// v0.5: Validate an event against a schema
1322pub async fn validate_event_schema(
1323    State(store): State<SharedStore>,
1324    Json(req): Json<ValidateEventRequest>,
1325) -> Result<Json<ValidateEventResponse>> {
1326    let schema_registry = store.schema_registry();
1327
1328    let response = schema_registry.validate(&req.subject, req.version, &req.payload)?;
1329
1330    if response.valid {
1331        tracing::debug!(
1332            "✅ Event validated against schema '{}' v{}",
1333            req.subject,
1334            response.schema_version
1335        );
1336    } else {
1337        tracing::warn!(
1338            "❌ Event validation failed for '{}': {:?}",
1339            req.subject,
1340            response.errors
1341        );
1342    }
1343
1344    Ok(Json(response))
1345}
1346
1347// v0.5: Set compatibility mode for a subject
1348#[derive(Deserialize)]
1349pub struct SetCompatibilityRequest {
1350    compatibility: CompatibilityMode,
1351}
1352
1353pub async fn set_compatibility_mode(
1354    State(store): State<SharedStore>,
1355    Path(subject): Path<String>,
1356    Json(req): Json<SetCompatibilityRequest>,
1357) -> Json<serde_json::Value> {
1358    let schema_registry = store.schema_registry();
1359
1360    schema_registry.set_compatibility_mode(subject.clone(), req.compatibility);
1361
1362    tracing::info!(
1363        "🔧 Set compatibility mode for '{}' to {:?}",
1364        subject,
1365        req.compatibility
1366    );
1367
1368    Json(serde_json::json!({
1369        "subject": subject,
1370        "compatibility": req.compatibility
1371    }))
1372}
1373
1374// v0.5: Start a replay operation
1375pub async fn start_replay(
1376    State(store): State<SharedStore>,
1377    Json(req): Json<StartReplayRequest>,
1378) -> Result<Json<StartReplayResponse>> {
1379    let replay_manager = store.replay_manager();
1380
1381    let response = replay_manager.start_replay(store, req)?;
1382
1383    tracing::info!(
1384        "🔄 Started replay {} with {} events",
1385        response.replay_id,
1386        response.total_events
1387    );
1388
1389    Ok(Json(response))
1390}
1391
1392// v0.5: Get replay progress
1393pub async fn get_replay_progress(
1394    State(store): State<SharedStore>,
1395    Path(replay_id): Path<uuid::Uuid>,
1396) -> Result<Json<ReplayProgress>> {
1397    let replay_manager = store.replay_manager();
1398
1399    let progress = replay_manager.get_progress(replay_id)?;
1400
1401    Ok(Json(progress))
1402}
1403
1404// v0.5: List all replay operations
1405pub async fn list_replays(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1406    let replay_manager = store.replay_manager();
1407
1408    let replays = replay_manager.list_replays();
1409
1410    Json(serde_json::json!({
1411        "replays": replays,
1412        "total": replays.len()
1413    }))
1414}
1415
1416// v0.5: Cancel a running replay
1417pub async fn cancel_replay(
1418    State(store): State<SharedStore>,
1419    Path(replay_id): Path<uuid::Uuid>,
1420) -> Result<Json<serde_json::Value>> {
1421    let replay_manager = store.replay_manager();
1422
1423    replay_manager.cancel_replay(replay_id)?;
1424
1425    tracing::info!("🛑 Cancelled replay {}", replay_id);
1426
1427    Ok(Json(serde_json::json!({
1428        "replay_id": replay_id,
1429        "status": "cancelled"
1430    })))
1431}
1432
1433// v0.5: Delete a completed replay
1434pub async fn delete_replay(
1435    State(store): State<SharedStore>,
1436    Path(replay_id): Path<uuid::Uuid>,
1437) -> Result<Json<serde_json::Value>> {
1438    let replay_manager = store.replay_manager();
1439
1440    let deleted = replay_manager.delete_replay(replay_id)?;
1441
1442    if deleted {
1443        tracing::info!("🗑️  Deleted replay {}", replay_id);
1444    }
1445
1446    Ok(Json(serde_json::json!({
1447        "replay_id": replay_id,
1448        "deleted": deleted
1449    })))
1450}
1451
1452// v0.5: Register a new pipeline
1453pub async fn register_pipeline(
1454    State(store): State<SharedStore>,
1455    Json(config): Json<PipelineConfig>,
1456) -> Result<Json<serde_json::Value>> {
1457    let pipeline_manager = store.pipeline_manager();
1458
1459    let pipeline_id = pipeline_manager.register(config.clone());
1460
1461    tracing::info!(
1462        "🔀 Pipeline registered: {} (name: {})",
1463        pipeline_id,
1464        config.name
1465    );
1466
1467    Ok(Json(serde_json::json!({
1468        "pipeline_id": pipeline_id,
1469        "name": config.name,
1470        "enabled": config.enabled
1471    })))
1472}
1473
1474// v0.5: List all pipelines
1475pub async fn list_pipelines(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1476    let pipeline_manager = store.pipeline_manager();
1477
1478    let pipelines = pipeline_manager.list();
1479
1480    tracing::debug!("Listed {} pipelines", pipelines.len());
1481
1482    Json(serde_json::json!({
1483        "pipelines": pipelines,
1484        "total": pipelines.len()
1485    }))
1486}
1487
1488// v0.5: Get a specific pipeline
1489pub async fn get_pipeline(
1490    State(store): State<SharedStore>,
1491    Path(pipeline_id): Path<uuid::Uuid>,
1492) -> Result<Json<PipelineConfig>> {
1493    let pipeline_manager = store.pipeline_manager();
1494
1495    let pipeline = pipeline_manager.get(pipeline_id).ok_or_else(|| {
1496        crate::error::AllSourceError::ValidationError(format!("Pipeline not found: {pipeline_id}"))
1497    })?;
1498
1499    Ok(Json(pipeline.config().clone()))
1500}
1501
1502// v0.5: Remove a pipeline
1503pub async fn remove_pipeline(
1504    State(store): State<SharedStore>,
1505    Path(pipeline_id): Path<uuid::Uuid>,
1506) -> Result<Json<serde_json::Value>> {
1507    let pipeline_manager = store.pipeline_manager();
1508
1509    let removed = pipeline_manager.remove(pipeline_id);
1510
1511    if removed {
1512        tracing::info!("🗑️  Removed pipeline {}", pipeline_id);
1513    }
1514
1515    Ok(Json(serde_json::json!({
1516        "pipeline_id": pipeline_id,
1517        "removed": removed
1518    })))
1519}
1520
1521// v0.5: Get statistics for all pipelines
1522pub async fn all_pipeline_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1523    let pipeline_manager = store.pipeline_manager();
1524
1525    let stats = pipeline_manager.all_stats();
1526
1527    Json(serde_json::json!({
1528        "stats": stats,
1529        "total": stats.len()
1530    }))
1531}
1532
1533// v0.5: Get statistics for a specific pipeline
1534pub async fn get_pipeline_stats(
1535    State(store): State<SharedStore>,
1536    Path(pipeline_id): Path<uuid::Uuid>,
1537) -> Result<Json<PipelineStats>> {
1538    let pipeline_manager = store.pipeline_manager();
1539
1540    let pipeline = pipeline_manager.get(pipeline_id).ok_or_else(|| {
1541        crate::error::AllSourceError::ValidationError(format!("Pipeline not found: {pipeline_id}"))
1542    })?;
1543
1544    Ok(Json(pipeline.stats()))
1545}
1546
1547// v0.5: Reset a pipeline's state
1548pub async fn reset_pipeline(
1549    State(store): State<SharedStore>,
1550    Path(pipeline_id): Path<uuid::Uuid>,
1551) -> Result<Json<serde_json::Value>> {
1552    let pipeline_manager = store.pipeline_manager();
1553
1554    let pipeline = pipeline_manager.get(pipeline_id).ok_or_else(|| {
1555        crate::error::AllSourceError::ValidationError(format!("Pipeline not found: {pipeline_id}"))
1556    })?;
1557
1558    pipeline.reset();
1559
1560    tracing::info!("🔄 Reset pipeline {}", pipeline_id);
1561
1562    Ok(Json(serde_json::json!({
1563        "pipeline_id": pipeline_id,
1564        "reset": true
1565    })))
1566}
1567
1568// =============================================================================
1569// v0.11: Single Event Lookup by ID
1570// =============================================================================
1571
1572/// Get a single event by UUID
1573pub async fn get_event_by_id(
1574    State(store): State<SharedStore>,
1575    Path(event_id): Path<uuid::Uuid>,
1576) -> Result<Json<serde_json::Value>> {
1577    let event = store.get_event_by_id(&event_id)?.ok_or_else(|| {
1578        crate::error::AllSourceError::EntityNotFound(format!("Event '{event_id}' not found"))
1579    })?;
1580
1581    let dto = EventDto::from(&event);
1582
1583    tracing::debug!("Event retrieved by ID: {}", event_id);
1584
1585    Ok(Json(serde_json::json!({
1586        "event": dto,
1587        "found": true
1588    })))
1589}
1590
1591// =============================================================================
1592// v0.7: Projection State API for Query Service Integration
1593// =============================================================================
1594
1595/// List all registered projections
1596pub async fn list_projections(State(store): State<SharedStore>) -> Json<serde_json::Value> {
1597    let projection_manager = store.projection_manager();
1598    let status_map = store.projection_status();
1599
1600    let projections: Vec<serde_json::Value> = projection_manager
1601        .list_projections()
1602        .iter()
1603        .map(|(name, projection)| {
1604            let status = status_map
1605                .get(name)
1606                .map_or_else(|| "running".to_string(), |s| s.value().clone());
1607            serde_json::json!({
1608                "name": name,
1609                "type": format!("{:?}", projection.name()),
1610                "status": status,
1611            })
1612        })
1613        .collect();
1614
1615    tracing::debug!("Listed {} projections", projections.len());
1616
1617    Json(serde_json::json!({
1618        "projections": projections,
1619        "total": projections.len()
1620    }))
1621}
1622
1623/// Get projection metadata by name
1624pub async fn get_projection(
1625    State(store): State<SharedStore>,
1626    Path(name): Path<String>,
1627) -> Result<Json<serde_json::Value>> {
1628    let projection_manager = store.projection_manager();
1629
1630    let projection = projection_manager.get_projection(&name).ok_or_else(|| {
1631        crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1632    })?;
1633
1634    Ok(Json(serde_json::json!({
1635        "name": projection.name(),
1636        "found": true
1637    })))
1638}
1639
1640/// Get projection state for a specific entity.
1641///
1642/// Resolution order:
1643/// 1. **Registered projection** — if `name` is registered with the projection
1644///    manager, return the projection's own `get_state(entity_id)` output.
1645/// 2. **Projection state cache** — otherwise fall back to whatever was written
1646///    via `save_projection_state` / `bulk_save_projection_states`. This supports
1647///    SDK-managed projections (e.g. the Rust SDK's `ProjectionWorker`) that
1648///    compute state client-side and push it back without registering a
1649///    projection in Core's manager.
1650///
1651/// Returns `found: false` with `state: null` when neither source has state.
1652pub async fn get_projection_state(
1653    State(store): State<SharedStore>,
1654    Path((name, entity_id)): Path<(String, String)>,
1655) -> Result<Json<serde_json::Value>> {
1656    let state = store
1657        .projection_manager()
1658        .get_projection(&name)
1659        .and_then(|p| p.get_state(&entity_id))
1660        .or_else(|| {
1661            store
1662                .projection_state_cache()
1663                .get(&format!("{name}:{entity_id}"))
1664                .map(|entry| entry.value().clone())
1665        });
1666
1667    tracing::debug!("Projection state retrieved: {} / {}", name, entity_id);
1668
1669    Ok(Json(serde_json::json!({
1670        "projection": name,
1671        "entity_id": entity_id,
1672        "state": state,
1673        "found": state.is_some()
1674    })))
1675}
1676
1677/// Delete (clear) a projection by name
1678///
1679/// Removes all state from the projection. The projection definition remains
1680/// registered but its accumulated state is cleared.
1681pub async fn delete_projection(
1682    State(store): State<SharedStore>,
1683    Path(name): Path<String>,
1684) -> Result<Json<serde_json::Value>> {
1685    let projection_manager = store.projection_manager();
1686
1687    let projection = projection_manager.get_projection(&name).ok_or_else(|| {
1688        crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1689    })?;
1690
1691    projection.clear();
1692
1693    // Also clear any cached state for this projection
1694    let cache = store.projection_state_cache();
1695    let prefix = format!("{name}:");
1696    let keys_to_remove: Vec<String> = cache
1697        .iter()
1698        .filter(|entry| entry.key().starts_with(&prefix))
1699        .map(|entry| entry.key().clone())
1700        .collect();
1701    for key in keys_to_remove {
1702        cache.remove(&key);
1703    }
1704
1705    tracing::info!("Projection deleted (cleared): {}", name);
1706
1707    Ok(Json(serde_json::json!({
1708        "projection": name,
1709        "deleted": true
1710    })))
1711}
1712
1713/// Query parameters for `GET /api/v1/projections/{name}/state`.
1714///
1715/// All optional: with none of them set the endpoint keeps its historical
1716/// behaviour of returning every cached entity for the projection, so existing
1717/// callers (the Query Service's `ProjectionServer` hydration, the SDKs) are
1718/// unaffected. Names match `ListEntitiesRequest` for consistency.
1719#[derive(Debug, Default, Deserialize)]
1720pub struct ProjectionStateSummaryParams {
1721    /// Maximum number of entity states to return. Unbounded when absent.
1722    pub limit: Option<usize>,
1723    /// Number of matching states to skip before applying `limit`. Default 0.
1724    pub offset: Option<usize>,
1725    /// Return only entities whose id starts with this prefix — lets a caller
1726    /// walk one shard of the keyspace without enumerating the whole projection.
1727    pub entity_id_prefix: Option<String>,
1728}
1729
1730/// Get aggregate projection state (all entities).
1731///
1732/// Returns the cached states written via `save_projection_state` /
1733/// `bulk_save_projection_states`. The projection does NOT need to be
1734/// registered with the projection manager — this supports SDK-managed
1735/// projections that push state without server-side registration.
1736///
1737/// Supports `limit`, `offset` and `entity_id_prefix` (issue #249): this is the
1738/// only endpoint that can *enumerate* a projection — `bulk_get_projection_states`
1739/// needs the ids up front — so a projection with one entry per tenant needs a
1740/// way to bound and resume a request. Entities are ordered by `entity_id` so
1741/// offset paging is stable; `total` is the full match set and `has_more` tells
1742/// a paginator when to stop.
1743///
1744/// Returns an empty list when no state has been written.
1745pub async fn get_projection_state_summary(
1746    State(store): State<SharedStore>,
1747    Path(name): Path<String>,
1748    Query(params): Query<ProjectionStateSummaryParams>,
1749) -> Result<Json<serde_json::Value>> {
1750    let cache = store.projection_state_cache();
1751    let prefix = format!("{name}:");
1752    let offset = params.offset.unwrap_or(0);
1753
1754    // Collect the matching keys first and sort them. DashMap iteration order is
1755    // arbitrary, so offset paging is only coherent against a total order — and
1756    // windowing ids instead of values means only the returned page's states are
1757    // cloned, not the whole projection.
1758    let mut entity_ids: Vec<String> = cache
1759        .iter()
1760        .filter_map(|entry| entry.key().strip_prefix(&prefix).map(ToString::to_string))
1761        .filter(|entity_id| {
1762            params
1763                .entity_id_prefix
1764                .as_ref()
1765                .is_none_or(|p| entity_id.starts_with(p))
1766        })
1767        .collect();
1768    entity_ids.sort_unstable();
1769
1770    let total = entity_ids.len();
1771
1772    let page = entity_ids.into_iter().skip(offset);
1773    let page: Vec<String> = match params.limit {
1774        Some(limit) => page.take(limit).collect(),
1775        None => page.collect(),
1776    };
1777
1778    let states: Vec<serde_json::Value> = page
1779        .into_iter()
1780        .filter_map(|entity_id| {
1781            // Skip entries deleted between the key scan and the value read.
1782            cache.get(&format!("{prefix}{entity_id}")).map(|entry| {
1783                serde_json::json!({
1784                    "entity_id": entity_id,
1785                    "state": entry.value().clone()
1786                })
1787            })
1788        })
1789        .collect();
1790
1791    let count = states.len();
1792    // Relative to the window actually served — a paginator that trusts a bare
1793    // `count < total` never terminates once an offset is in play (cf. #250).
1794    let has_more = offset + count < total;
1795
1796    tracing::debug!(
1797        "Projection state summary: {} ({} of {} entities, offset {})",
1798        name,
1799        count,
1800        total,
1801        offset
1802    );
1803
1804    Ok(Json(serde_json::json!({
1805        "projection": name,
1806        "states": states,
1807        "count": count,
1808        "total": total,
1809        "has_more": has_more
1810    })))
1811}
1812
1813/// Reset a projection to its initial state
1814///
1815/// Clears all accumulated state and reprocesses events from the beginning.
1816pub async fn reset_projection(
1817    State(store): State<SharedStore>,
1818    Path(name): Path<String>,
1819) -> Result<Json<serde_json::Value>> {
1820    let reprocessed = store.reset_projection(&name)?;
1821
1822    tracing::info!(
1823        "Projection reset: {} ({} events reprocessed)",
1824        name,
1825        reprocessed
1826    );
1827
1828    Ok(Json(serde_json::json!({
1829        "projection": name,
1830        "reset": true,
1831        "events_reprocessed": reprocessed
1832    })))
1833}
1834
1835/// Pause a projection
1836///
1837/// Sets the projection status to "paused" so it stops processing new events.
1838pub async fn pause_projection(
1839    State(store): State<SharedStore>,
1840    Path(name): Path<String>,
1841) -> Result<Json<serde_json::Value>> {
1842    let projection_manager = store.projection_manager();
1843
1844    // Verify projection exists
1845    let _projection = projection_manager.get_projection(&name).ok_or_else(|| {
1846        crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1847    })?;
1848
1849    store
1850        .projection_status()
1851        .insert(name.clone(), "paused".to_string());
1852
1853    tracing::info!("Projection paused: {}", name);
1854
1855    Ok(Json(serde_json::json!({
1856        "projection": name,
1857        "status": "paused"
1858    })))
1859}
1860
1861/// Start (resume) a projection
1862///
1863/// Sets the projection status to "running" so it resumes processing events.
1864pub async fn start_projection(
1865    State(store): State<SharedStore>,
1866    Path(name): Path<String>,
1867) -> Result<Json<serde_json::Value>> {
1868    let projection_manager = store.projection_manager();
1869
1870    // Verify projection exists
1871    let _projection = projection_manager.get_projection(&name).ok_or_else(|| {
1872        crate::error::AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1873    })?;
1874
1875    store
1876        .projection_status()
1877        .insert(name.clone(), "running".to_string());
1878
1879    tracing::info!("Projection started: {}", name);
1880
1881    Ok(Json(serde_json::json!({
1882        "projection": name,
1883        "status": "running"
1884    })))
1885}
1886
1887/// Request body for saving projection state
1888#[derive(Debug, Deserialize)]
1889pub struct SaveProjectionStateRequest {
1890    pub state: serde_json::Value,
1891}
1892
1893/// Save/update projection state for an entity
1894///
1895/// This endpoint allows external services (like Elixir Query Service) to
1896/// store computed projection state back to the Core for persistence.
1897pub async fn save_projection_state(
1898    State(store): State<SharedStore>,
1899    Path((name, entity_id)): Path<(String, String)>,
1900    Json(req): Json<SaveProjectionStateRequest>,
1901) -> Result<Json<serde_json::Value>> {
1902    let projection_cache = store.projection_state_cache();
1903
1904    // Store in the projection state cache
1905    projection_cache.insert(format!("{name}:{entity_id}"), req.state.clone());
1906
1907    tracing::info!("Projection state saved: {} / {}", name, entity_id);
1908
1909    Ok(Json(serde_json::json!({
1910        "projection": name,
1911        "entity_id": entity_id,
1912        "saved": true
1913    })))
1914}
1915
1916/// Bulk get projection states for multiple entities
1917///
1918/// Efficient endpoint for fetching multiple entity states in a single request.
1919#[derive(Debug, Deserialize)]
1920pub struct BulkGetStateRequest {
1921    pub entity_ids: Vec<String>,
1922}
1923
1924/// Bulk save projection states for multiple entities
1925///
1926/// Efficient endpoint for saving multiple entity states in a single request.
1927#[derive(Debug, Deserialize)]
1928pub struct BulkSaveStateRequest {
1929    pub states: Vec<BulkSaveStateItem>,
1930}
1931
1932#[derive(Debug, Deserialize)]
1933pub struct BulkSaveStateItem {
1934    pub entity_id: String,
1935    pub state: serde_json::Value,
1936}
1937
1938pub async fn bulk_get_projection_states(
1939    State(store): State<SharedStore>,
1940    Path(name): Path<String>,
1941    Json(req): Json<BulkGetStateRequest>,
1942) -> Result<Json<serde_json::Value>> {
1943    // Same fallback rule as `get_projection_state`: registered projection
1944    // wins, cache is the fallback. This lets SDK-managed projections read
1945    // their pushed-back state without registering in Core's projection manager.
1946    let projection = store.projection_manager().get_projection(&name);
1947    let cache = store.projection_state_cache();
1948
1949    let states: Vec<serde_json::Value> = req
1950        .entity_ids
1951        .iter()
1952        .map(|entity_id| {
1953            let state = projection
1954                .as_ref()
1955                .and_then(|p| p.get_state(entity_id))
1956                .or_else(|| {
1957                    cache
1958                        .get(&format!("{name}:{entity_id}"))
1959                        .map(|entry| entry.value().clone())
1960                });
1961            serde_json::json!({
1962                "entity_id": entity_id,
1963                "state": state,
1964                "found": state.is_some()
1965            })
1966        })
1967        .collect();
1968
1969    tracing::debug!(
1970        "Bulk projection state retrieved: {} entities from {}",
1971        states.len(),
1972        name
1973    );
1974
1975    Ok(Json(serde_json::json!({
1976        "projection": name,
1977        "states": states,
1978        "total": states.len()
1979    })))
1980}
1981
1982/// Bulk save projection states for multiple entities
1983///
1984/// This endpoint allows efficient batch saving of projection states,
1985/// critical for high-throughput event processing pipelines.
1986pub async fn bulk_save_projection_states(
1987    State(store): State<SharedStore>,
1988    Path(name): Path<String>,
1989    Json(req): Json<BulkSaveStateRequest>,
1990) -> Result<Json<serde_json::Value>> {
1991    let projection_cache = store.projection_state_cache();
1992
1993    let mut saved_count = 0;
1994    for item in &req.states {
1995        projection_cache.insert(format!("{name}:{}", item.entity_id), item.state.clone());
1996        saved_count += 1;
1997    }
1998
1999    tracing::info!(
2000        "Bulk projection state saved: {} entities for {}",
2001        saved_count,
2002        name
2003    );
2004
2005    Ok(Json(serde_json::json!({
2006        "projection": name,
2007        "saved": saved_count,
2008        "total": req.states.len()
2009    })))
2010}
2011
2012// =============================================================================
2013// v0.11: Webhook Management API
2014// =============================================================================
2015
2016/// Query parameters for listing webhooks
2017#[derive(Debug, Deserialize)]
2018pub struct ListWebhooksParams {
2019    pub tenant_id: Option<String>,
2020}
2021
2022/// Register a new webhook subscription
2023pub async fn register_webhook(
2024    State(store): State<SharedStore>,
2025    Json(req): Json<RegisterWebhookRequest>,
2026) -> Json<serde_json::Value> {
2027    let registry = store.webhook_registry();
2028    let webhook = registry.register(req);
2029
2030    tracing::info!("Webhook registered: {} -> {}", webhook.id, webhook.url);
2031
2032    Json(serde_json::json!({
2033        "webhook": webhook,
2034        "created": true
2035    }))
2036}
2037
2038/// List webhooks, optionally filtered by tenant_id
2039pub async fn list_webhooks(
2040    State(store): State<SharedStore>,
2041    Query(params): Query<ListWebhooksParams>,
2042) -> Json<serde_json::Value> {
2043    let registry = store.webhook_registry();
2044
2045    let webhooks = if let Some(tenant_id) = params.tenant_id {
2046        registry.list_by_tenant(&tenant_id)
2047    } else {
2048        // Without tenant filter, return empty (tenants should always filter)
2049        vec![]
2050    };
2051
2052    let total = webhooks.len();
2053
2054    Json(serde_json::json!({
2055        "webhooks": webhooks,
2056        "total": total
2057    }))
2058}
2059
2060/// Get a specific webhook by ID
2061pub async fn get_webhook(
2062    State(store): State<SharedStore>,
2063    Path(webhook_id): Path<uuid::Uuid>,
2064) -> Result<Json<serde_json::Value>> {
2065    let registry = store.webhook_registry();
2066
2067    let webhook = registry.get(webhook_id).ok_or_else(|| {
2068        crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2069    })?;
2070
2071    Ok(Json(serde_json::json!({
2072        "webhook": webhook,
2073        "found": true
2074    })))
2075}
2076
2077/// Update a webhook subscription
2078pub async fn update_webhook(
2079    State(store): State<SharedStore>,
2080    Path(webhook_id): Path<uuid::Uuid>,
2081    Json(req): Json<UpdateWebhookRequest>,
2082) -> Result<Json<serde_json::Value>> {
2083    let registry = store.webhook_registry();
2084
2085    let webhook = registry.update(webhook_id, req).ok_or_else(|| {
2086        crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2087    })?;
2088
2089    tracing::info!("Webhook updated: {}", webhook_id);
2090
2091    Ok(Json(serde_json::json!({
2092        "webhook": webhook,
2093        "updated": true
2094    })))
2095}
2096
2097/// Delete a webhook subscription
2098pub async fn delete_webhook(
2099    State(store): State<SharedStore>,
2100    Path(webhook_id): Path<uuid::Uuid>,
2101) -> Result<Json<serde_json::Value>> {
2102    let registry = store.webhook_registry();
2103
2104    let webhook = registry.delete(webhook_id).ok_or_else(|| {
2105        crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2106    })?;
2107
2108    tracing::info!("Webhook deleted: {} ({})", webhook_id, webhook.url);
2109
2110    Ok(Json(serde_json::json!({
2111        "webhook_id": webhook_id,
2112        "deleted": true
2113    })))
2114}
2115
2116/// Query parameters for listing webhook deliveries
2117#[derive(Debug, Deserialize)]
2118pub struct ListDeliveriesParams {
2119    pub limit: Option<usize>,
2120}
2121
2122/// List delivery history for a webhook
2123pub async fn list_webhook_deliveries(
2124    State(store): State<SharedStore>,
2125    Path(webhook_id): Path<uuid::Uuid>,
2126    Query(params): Query<ListDeliveriesParams>,
2127) -> Result<Json<serde_json::Value>> {
2128    let registry = store.webhook_registry();
2129
2130    // Verify webhook exists
2131    registry.get(webhook_id).ok_or_else(|| {
2132        crate::error::AllSourceError::EntityNotFound(format!("Webhook '{webhook_id}' not found"))
2133    })?;
2134
2135    let limit = params.limit.unwrap_or(50);
2136    let deliveries = registry.get_deliveries(webhook_id, limit);
2137    let total = deliveries.len();
2138
2139    Ok(Json(serde_json::json!({
2140        "webhook_id": webhook_id,
2141        "deliveries": deliveries,
2142        "total": total
2143    })))
2144}
2145
2146// =============================================================================
2147// v2.0: Advanced Query Features
2148// =============================================================================
2149
2150/// EventQL: Execute SQL queries over events using DataFusion
2151#[cfg(feature = "analytics")]
2152pub async fn eventql_query(
2153    State(store): State<SharedStore>,
2154    Json(req): Json<crate::infrastructure::query::eventql::EventQLRequest>,
2155) -> Result<Json<serde_json::Value>> {
2156    let events = store.snapshot_events();
2157    match crate::infrastructure::query::eventql::execute_eventql(&events, &req).await {
2158        Ok(response) => Ok(Json(serde_json::json!({
2159            "columns": response.columns,
2160            "rows": response.rows,
2161            "row_count": response.row_count,
2162        }))),
2163        Err(e) => Err(crate::error::AllSourceError::InvalidQuery(e)),
2164    }
2165}
2166
2167/// GraphQL: Execute GraphQL queries
2168pub async fn graphql_query(
2169    State(store): State<SharedStore>,
2170    Json(req): Json<GraphQLRequest>,
2171) -> Json<serde_json::Value> {
2172    let fields = match crate::infrastructure::query::graphql::parse_query(&req.query) {
2173        Ok(f) => f,
2174        Err(e) => {
2175            return Json(
2176                serde_json::to_value(GraphQLResponse {
2177                    data: None,
2178                    errors: vec![GraphQLError { message: e }],
2179                })
2180                .unwrap(),
2181            );
2182        }
2183    };
2184
2185    let mut data = serde_json::Map::new();
2186    let mut errors = Vec::new();
2187
2188    for field in &fields {
2189        match field.name.as_str() {
2190            "events" => {
2191                let request = crate::application::dto::QueryEventsRequest {
2192                    entity_id: field.arguments.get("entity_id").cloned(),
2193                    event_type: field.arguments.get("event_type").cloned(),
2194                    tenant_id: field.arguments.get("tenant_id").cloned(),
2195                    limit: field.arguments.get("limit").and_then(|l| l.parse().ok()),
2196                    as_of: None,
2197                    since: None,
2198                    until: None,
2199                    event_type_prefix: None,
2200                    exclude_event_type_prefix: None,
2201                    payload_filter: None,
2202                };
2203                match store.query(&request) {
2204                    Ok(events) => {
2205                        let json_events: Vec<serde_json::Value> = events
2206                            .iter()
2207                            .map(|e| {
2208                                crate::infrastructure::query::graphql::event_to_json(
2209                                    e,
2210                                    &field.fields,
2211                                )
2212                            })
2213                            .collect();
2214                        data.insert("events".to_string(), serde_json::Value::Array(json_events));
2215                    }
2216                    Err(e) => errors.push(GraphQLError {
2217                        message: format!("events query failed: {e}"),
2218                    }),
2219                }
2220            }
2221            "event" => {
2222                if let Some(id_str) = field.arguments.get("id") {
2223                    if let Ok(id) = uuid::Uuid::parse_str(id_str) {
2224                        match store.get_event_by_id(&id) {
2225                            Ok(Some(event)) => {
2226                                data.insert(
2227                                    "event".to_string(),
2228                                    crate::infrastructure::query::graphql::event_to_json(
2229                                        &event,
2230                                        &field.fields,
2231                                    ),
2232                                );
2233                            }
2234                            Ok(None) => {
2235                                data.insert("event".to_string(), serde_json::Value::Null);
2236                            }
2237                            Err(e) => errors.push(GraphQLError {
2238                                message: format!("event lookup failed: {e}"),
2239                            }),
2240                        }
2241                    } else {
2242                        errors.push(GraphQLError {
2243                            message: format!("Invalid UUID: {id_str}"),
2244                        });
2245                    }
2246                } else {
2247                    errors.push(GraphQLError {
2248                        message: "event query requires 'id' argument".to_string(),
2249                    });
2250                }
2251            }
2252            "projections" => {
2253                let pm = store.projection_manager();
2254                let names: Vec<serde_json::Value> = pm
2255                    .list_projections()
2256                    .iter()
2257                    .map(|(name, _)| serde_json::Value::String(name.clone()))
2258                    .collect();
2259                data.insert("projections".to_string(), serde_json::Value::Array(names));
2260            }
2261            "stats" => {
2262                let stats = store.stats();
2263                data.insert(
2264                    "stats".to_string(),
2265                    serde_json::json!({
2266                        "total_events": stats.total_events,
2267                        "total_entities": stats.total_entities,
2268                        "total_event_types": stats.total_event_types,
2269                    }),
2270                );
2271            }
2272            "__schema" => {
2273                data.insert(
2274                    "__schema".to_string(),
2275                    crate::infrastructure::query::graphql::introspection_schema(),
2276                );
2277            }
2278            other => {
2279                errors.push(GraphQLError {
2280                    message: format!("Unknown field: {other}"),
2281                });
2282            }
2283        }
2284    }
2285
2286    Json(
2287        serde_json::to_value(GraphQLResponse {
2288            data: Some(serde_json::Value::Object(data)),
2289            errors,
2290        })
2291        .unwrap(),
2292    )
2293}
2294
2295/// Geospatial: Query events by location
2296pub async fn geo_query(
2297    State(store): State<SharedStore>,
2298    Json(req): Json<GeoQueryRequest>,
2299) -> Json<serde_json::Value> {
2300    let events = store.snapshot_events();
2301    let geo_index = store.geo_index();
2302    let results =
2303        crate::infrastructure::query::geospatial::execute_geo_query(&events, &geo_index, &req);
2304    let total = results.len();
2305    Json(serde_json::json!({
2306        "results": results,
2307        "total": total,
2308    }))
2309}
2310
2311/// Geospatial index stats
2312pub async fn geo_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
2313    let stats = store.geo_index().stats();
2314    Json(serde_json::json!(stats))
2315}
2316
2317/// Exactly-once processing stats
2318pub async fn exactly_once_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
2319    let stats = store.exactly_once().stats();
2320    Json(serde_json::json!(stats))
2321}
2322
2323/// Schema evolution history for an event type
2324pub async fn schema_evolution_history(
2325    State(store): State<SharedStore>,
2326    Path(event_type): Path<String>,
2327) -> Json<serde_json::Value> {
2328    let mgr = store.schema_evolution();
2329    let history = mgr.get_history(&event_type);
2330    let version = mgr.get_version(&event_type);
2331    Json(serde_json::json!({
2332        "event_type": event_type,
2333        "current_version": version,
2334        "history": history,
2335    }))
2336}
2337
2338/// Current inferred schema for an event type
2339pub async fn schema_evolution_schema(
2340    State(store): State<SharedStore>,
2341    Path(event_type): Path<String>,
2342) -> Json<serde_json::Value> {
2343    let mgr = store.schema_evolution();
2344    if let Some(schema) = mgr.get_schema(&event_type) {
2345        let json_schema = crate::application::services::schema_evolution::to_json_schema(&schema);
2346        Json(serde_json::json!({
2347            "event_type": event_type,
2348            "version": mgr.get_version(&event_type),
2349            "inferred_schema": schema,
2350            "json_schema": json_schema,
2351        }))
2352    } else {
2353        Json(serde_json::json!({
2354            "event_type": event_type,
2355            "error": "No schema inferred for this event type"
2356        }))
2357    }
2358}
2359
2360/// Schema evolution stats
2361pub async fn schema_evolution_stats(State(store): State<SharedStore>) -> Json<serde_json::Value> {
2362    let stats = store.schema_evolution().stats();
2363    let event_types = store.schema_evolution().list_event_types();
2364    Json(serde_json::json!({
2365        "stats": stats,
2366        "tracked_event_types": event_types,
2367    }))
2368}
2369
2370// =============================================================================
2371// Sync Protocol Endpoints (v0.11: embedded↔server bidirectional sync)
2372// =============================================================================
2373
2374/// POST /api/v1/sync/pull — Client sends version vector, server returns delta events.
2375#[cfg(feature = "embedded-sync")]
2376pub async fn sync_pull_handler(
2377    State(state): State<AppState>,
2378    Json(request): Json<crate::embedded::sync_types::SyncPullRequest>,
2379) -> Result<Json<crate::embedded::sync_types::SyncPullResponse>> {
2380    use crate::infrastructure::cluster::{crdt::ReplicatedEvent, hlc::HlcTimestamp};
2381
2382    let store = &state.store;
2383
2384    // Compute "since" threshold from the client's version vector
2385    // We return all events the client hasn't seen yet
2386    let since = request
2387        .version_vector
2388        .values()
2389        .map(|ts| ts.physical_ms)
2390        .min()
2391        .and_then(|ms| chrono::DateTime::from_timestamp_millis(ms as i64));
2392
2393    let events = store.query(&crate::application::dto::QueryEventsRequest {
2394        entity_id: None,
2395        event_type: None,
2396        tenant_id: None,
2397        as_of: None,
2398        since,
2399        until: None,
2400        limit: None,
2401        event_type_prefix: None,
2402        exclude_event_type_prefix: None,
2403        payload_filter: None,
2404    })?;
2405
2406    // Convert domain events to ReplicatedEvent wire format
2407    let mut replicated = Vec::with_capacity(events.len());
2408    let mut last_ms = 0u64;
2409    let mut logical = 0u32;
2410
2411    for event in &events {
2412        let event_ms = event.timestamp().timestamp_millis() as u64;
2413        if event_ms == last_ms {
2414            logical += 1;
2415        } else {
2416            last_ms = event_ms;
2417            logical = 0;
2418        }
2419
2420        replicated.push(ReplicatedEvent {
2421            event_id: event.id().to_string(),
2422            hlc_timestamp: HlcTimestamp::new(event_ms, logical, 0),
2423            origin_region: "server".to_string(),
2424            event_data: serde_json::json!({
2425                "event_type": event.event_type_str(),
2426                "entity_id": event.entity_id_str(),
2427                "tenant_id": event.tenant_id_str(),
2428                "payload": event.payload,
2429                "metadata": event.metadata,
2430            }),
2431        });
2432    }
2433
2434    Ok(Json(crate::embedded::sync_types::SyncPullResponse {
2435        events: replicated,
2436        version_vector: std::collections::BTreeMap::new(),
2437    }))
2438}
2439
2440/// POST /api/v1/sync/push — Client pushes events, server applies CRDT resolution.
2441#[cfg(feature = "embedded-sync")]
2442pub async fn sync_push_handler(
2443    State(state): State<AppState>,
2444    Json(request): Json<crate::embedded::sync_types::SyncPushRequest>,
2445) -> Result<Json<crate::embedded::sync_types::SyncPushResponse>> {
2446    let store = &state.store;
2447
2448    let mut accepted = 0usize;
2449    let mut skipped = 0usize;
2450
2451    for rep_event in &request.events {
2452        let event_data = &rep_event.event_data;
2453        let event_type = event_data
2454            .get("event_type")
2455            .and_then(|v| v.as_str())
2456            .unwrap_or("unknown")
2457            .to_string();
2458        let entity_id = event_data
2459            .get("entity_id")
2460            .and_then(|v| v.as_str())
2461            .unwrap_or("unknown")
2462            .to_string();
2463        let tenant_id = event_data
2464            .get("tenant_id")
2465            .and_then(|v| v.as_str())
2466            .unwrap_or("default")
2467            .to_string();
2468        let payload = event_data
2469            .get("payload")
2470            .cloned()
2471            .unwrap_or(serde_json::json!({}));
2472        let metadata = event_data.get("metadata").cloned();
2473
2474        match Event::from_strings(event_type, entity_id, tenant_id, payload, metadata) {
2475            Ok(domain_event) => {
2476                store.ingest(&domain_event)?;
2477                accepted += 1;
2478            }
2479            Err(_) => {
2480                skipped += 1;
2481            }
2482        }
2483    }
2484
2485    Ok(Json(crate::embedded::sync_types::SyncPushResponse {
2486        accepted,
2487        skipped,
2488        version_vector: std::collections::BTreeMap::new(),
2489    }))
2490}
2491
2492// =============================================================================
2493// Consumer endpoints for durable subscriptions (v0.14)
2494// =============================================================================
2495
2496/// POST /api/v1/consumers — Register a durable consumer
2497pub async fn register_consumer(
2498    State(store): State<SharedStore>,
2499    Json(req): Json<RegisterConsumerRequest>,
2500) -> Result<Json<ConsumerResponse>> {
2501    let consumer = store
2502        .consumer_registry()
2503        .register(&req.consumer_id, &req.event_type_filters);
2504
2505    Ok(Json(ConsumerResponse {
2506        consumer_id: consumer.consumer_id,
2507        event_type_filters: consumer.event_type_filters,
2508        cursor_position: consumer.cursor_position,
2509    }))
2510}
2511
2512/// GET /api/v1/consumers/{consumer_id} — Get consumer metadata and cursor position
2513pub async fn get_consumer(
2514    State(store): State<SharedStore>,
2515    Path(consumer_id): Path<String>,
2516) -> Result<Json<ConsumerResponse>> {
2517    let consumer = store.consumer_registry().get_or_create(&consumer_id);
2518
2519    Ok(Json(ConsumerResponse {
2520        consumer_id: consumer.consumer_id,
2521        event_type_filters: consumer.event_type_filters,
2522        cursor_position: consumer.cursor_position,
2523    }))
2524}
2525
2526/// GET /api/v1/consumers/{consumer_id}/events — Poll for events since last ack
2527#[derive(Debug, Deserialize)]
2528pub struct ConsumerPollQuery {
2529    pub limit: Option<usize>,
2530}
2531
2532pub async fn poll_consumer_events(
2533    State(store): State<SharedStore>,
2534    Path(consumer_id): Path<String>,
2535    Query(query): Query<ConsumerPollQuery>,
2536) -> Result<Json<ConsumerEventsResponse>> {
2537    let consumer = store.consumer_registry().get_or_create(&consumer_id);
2538    let offset = consumer.cursor_position.unwrap_or(0);
2539    let limit = query.limit.unwrap_or(100);
2540
2541    let events = store.events_after_offset(offset, &consumer.event_type_filters, limit);
2542    let count = events.len();
2543
2544    let consumer_events: Vec<ConsumerEventDto> = events
2545        .into_iter()
2546        .map(|(position, event)| ConsumerEventDto {
2547            position,
2548            event: EventDto::from(&event),
2549        })
2550        .collect();
2551
2552    Ok(Json(ConsumerEventsResponse {
2553        events: consumer_events,
2554        count,
2555    }))
2556}
2557
2558/// POST /api/v1/consumers/{consumer_id}/ack — Acknowledge processed events
2559pub async fn ack_consumer(
2560    State(store): State<SharedStore>,
2561    Path(consumer_id): Path<String>,
2562    Json(req): Json<AckRequest>,
2563) -> Result<Json<serde_json::Value>> {
2564    let max_offset = store.total_events() as u64;
2565
2566    store
2567        .consumer_registry()
2568        .ack(&consumer_id, req.position, max_offset)
2569        .map_err(crate::error::AllSourceError::InvalidInput)?;
2570
2571    Ok(Json(serde_json::json!({
2572        "status": "ok",
2573        "consumer_id": consumer_id,
2574        "position": req.position,
2575    })))
2576}
2577
2578#[cfg(test)]
2579mod tests {
2580    use super::*;
2581    use crate::{domain::entities::Event, store::EventStore};
2582
2583    fn create_test_store() -> Arc<EventStore> {
2584        Arc::new(EventStore::new())
2585    }
2586
2587    /// Call the REAL `GET /api/v1/events/query` handler with `query` as the
2588    /// query string, parsed through the same extractors the router uses.
2589    ///
2590    /// Tests that hand the handler DTOs they built themselves cannot see
2591    /// parameters the DTOs never declare (issue #250) and cannot exercise the
2592    /// ordering/windowing composition the handler delegates to the store
2593    /// (issue #251), so ordering and pagination guards go through here.
2594    /// `tenant_id` is defaulted to the one `create_test_event` stamps.
2595    async fn query_page(store: &SharedStore, query: &str) -> QueryEventsResponse {
2596        use axum::extract::{Query, State};
2597
2598        let uri: axum::http::Uri = format!("/api/v1/events/query?tenant_id=test-stream&{query}")
2599            .parse()
2600            .unwrap();
2601        query_events(
2602            OptionalAuth(None),
2603            Query::try_from_uri(&uri).unwrap(),
2604            Query::try_from_uri(&uri).unwrap(),
2605            Query::try_from_uri(&uri).unwrap(),
2606            Query::try_from_uri(&uri).unwrap(),
2607            State(store.clone()),
2608        )
2609        .await
2610        .unwrap()
2611        .0
2612    }
2613
2614    fn create_test_event(entity_id: &str, event_type: &str) -> Event {
2615        Event::from_strings(
2616            event_type.to_string(),
2617            entity_id.to_string(),
2618            "test-stream".to_string(),
2619            serde_json::json!({
2620                "name": "Test",
2621                "value": 42
2622            }),
2623            None,
2624        )
2625        .unwrap()
2626    }
2627
2628    #[tokio::test]
2629    async fn test_query_events_has_more_and_total_count() {
2630        let store = create_test_store();
2631
2632        // Ingest 50 events
2633        for i in 0..50 {
2634            store
2635                .ingest(&create_test_event(&format!("entity-{i}"), "user.created"))
2636                .unwrap();
2637        }
2638
2639        // Query with limit=10 — should get has_more=true, total_count=50
2640        let req = QueryEventsRequest {
2641            entity_id: None,
2642            event_type: None,
2643            tenant_id: None,
2644            as_of: None,
2645            since: None,
2646            until: None,
2647            limit: Some(10),
2648            event_type_prefix: None,
2649            exclude_event_type_prefix: None,
2650            payload_filter: None,
2651        };
2652
2653        let requested_limit = req.limit;
2654        let unlimited_req = QueryEventsRequest {
2655            limit: None,
2656            ..QueryEventsRequest {
2657                entity_id: req.entity_id,
2658                event_type: req.event_type,
2659                tenant_id: req.tenant_id,
2660                as_of: req.as_of,
2661                since: req.since,
2662                until: req.until,
2663                limit: None,
2664                event_type_prefix: req.event_type_prefix,
2665                exclude_event_type_prefix: None,
2666                payload_filter: req.payload_filter,
2667            }
2668        };
2669        let all_events = store.query(&unlimited_req).unwrap();
2670        let total_count = all_events.len();
2671        let limited_events: Vec<Event> = if let Some(limit) = requested_limit {
2672            all_events.into_iter().take(limit).collect()
2673        } else {
2674            all_events
2675        };
2676        let count = limited_events.len();
2677        let has_more = count < total_count;
2678
2679        assert_eq!(count, 10);
2680        assert_eq!(total_count, 50);
2681        assert!(has_more);
2682    }
2683
2684    // Regression guard for issue #250: `GET /api/v1/events/query` must honour
2685    // `offset`. It used to be dropped silently (the DTO did not declare it), so
2686    // every page returned the same first `limit` events and `has_more` stayed
2687    // true — the Rust SDK's `EventPaginator` (and the Go SDK's `QueryOptions`)
2688    // both send `offset`, so `collect_all()` looped forever accumulating
2689    // duplicates. Drives the REAL handler through query-string deserialization,
2690    // because the bug lived in the wire layer: a hand-rolled test that builds
2691    // the DTO in code cannot see a field the DTO never declares.
2692    #[tokio::test]
2693    async fn query_events_honours_offset_pagination() {
2694        use axum::extract::{Query, State};
2695
2696        let store = create_test_store();
2697        for i in 0..25 {
2698            store
2699                .ingest(&create_test_event(&format!("e-{i:02}"), "user.created"))
2700                .unwrap();
2701        }
2702
2703        // Parse a real query string through the same extractors the router uses.
2704        async fn page(store: &SharedStore, limit: usize, offset: usize) -> QueryEventsResponse {
2705            let uri: axum::http::Uri =
2706                format!("/api/v1/events/query?tenant_id=test-stream&limit={limit}&offset={offset}")
2707                    .parse()
2708                    .unwrap();
2709            let req: Query<QueryEventsRequest> = Query::try_from_uri(&uri).unwrap();
2710            let order: Query<EventOrderParam> = Query::try_from_uri(&uri).unwrap();
2711            let off: Query<EventOffsetParam> = Query::try_from_uri(&uri).unwrap();
2712            assert_eq!(off.0.offset, Some(offset), "offset must deserialize");
2713            query_events(
2714                OptionalAuth(None),
2715                req,
2716                order,
2717                off,
2718                Query::try_from_uri(&uri).unwrap(),
2719                State(store.clone()),
2720            )
2721            .await
2722            .unwrap()
2723            .0
2724        }
2725
2726        let p1 = page(&store, 10, 0).await;
2727        let p2 = page(&store, 10, 10).await;
2728        let p3 = page(&store, 10, 20).await;
2729
2730        assert_eq!(p1.count, 10);
2731        assert_eq!(p2.count, 10);
2732        assert_eq!(p3.count, 5, "last page returns the remainder");
2733        assert_eq!(p1.total_count, 25);
2734
2735        // The core failure: page 2 must not be page 1 again.
2736        let ids = |r: &QueryEventsResponse| -> Vec<String> {
2737            r.events.iter().map(|e| e.entity_id.clone()).collect()
2738        };
2739        assert_ne!(ids(&p1), ids(&p2), "offset=10 must skip the first page");
2740
2741        let mut all = ids(&p1);
2742        all.extend(ids(&p2));
2743        all.extend(ids(&p3));
2744        let unique: std::collections::HashSet<_> = all.iter().cloned().collect();
2745        assert_eq!(
2746            unique.len(),
2747            25,
2748            "paging the whole set must yield 25 distinct entities, not duplicates"
2749        );
2750
2751        // `has_more` must account for the offset, otherwise a paginator that
2752        // trusts it never terminates.
2753        assert!(p1.has_more, "25 events, page 1 of 10 → more remain");
2754        assert!(p2.has_more, "25 events, page 2 of 10 → more remain");
2755        assert!(!p3.has_more, "offset=20 + count=5 == total → exhausted");
2756
2757        // Past the end: empty page, and exhausted rather than "more".
2758        let past = page(&store, 10, 100).await;
2759        assert_eq!(past.count, 0);
2760        assert!(!past.has_more, "offset beyond the match set is exhausted");
2761    }
2762
2763    // Regression guard for issue #249: `GET /api/v1/projections/{name}/state`
2764    // must honour `limit`, `offset` and `entity_id_prefix`. The handler used to
2765    // take only `Path(name)`, so query parameters were dropped on the floor and
2766    // the response grew linearly with the number of cached entities — a caller
2767    // with one entry per tenant had no way to bound a request or resume one.
2768    // Driven through a real router so query-string deserialization is exercised:
2769    // a test that hands the handler a DTO it built itself cannot fail on
2770    // parameters the handler never declares.
2771    #[tokio::test]
2772    async fn projection_state_summary_honours_limit_offset_and_prefix() {
2773        use axum::{
2774            body::{Body, to_bytes},
2775            http::Request,
2776        };
2777        use tower::ServiceExt; // for `oneshot`
2778
2779        let store = create_test_store();
2780        let cache = store.projection_state_cache();
2781        for i in 0..25 {
2782            cache.insert(
2783                format!("demo:tenant-{i:02}"),
2784                serde_json::json!({ "n": i as u64 }),
2785            );
2786        }
2787        // A neighbouring projection and a same-prefix-looking key must not leak.
2788        cache.insert(
2789            "other:tenant-99".to_string(),
2790            serde_json::json!({ "n": 99 }),
2791        );
2792
2793        let app = Router::new()
2794            .route(
2795                "/api/v1/projections/{name}/state",
2796                get(get_projection_state_summary),
2797            )
2798            .with_state(store.clone());
2799
2800        async fn fetch(app: &Router, uri: &str) -> serde_json::Value {
2801            let resp = app
2802                .clone()
2803                .oneshot(Request::builder().uri(uri).body(Body::empty()).unwrap())
2804                .await
2805                .unwrap();
2806            assert_eq!(resp.status(), axum::http::StatusCode::OK, "GET {uri}");
2807            let bytes = to_bytes(resp.into_body(), usize::MAX).await.unwrap();
2808            serde_json::from_slice(&bytes).unwrap()
2809        }
2810
2811        let ids = |body: &serde_json::Value| -> Vec<String> {
2812            body["states"]
2813                .as_array()
2814                .unwrap()
2815                .iter()
2816                .map(|s| s["entity_id"].as_str().unwrap().to_string())
2817                .collect()
2818        };
2819
2820        // No params: unchanged behaviour — the whole projection, scoped to it.
2821        let all = fetch(&app, "/api/v1/projections/demo/state").await;
2822        assert_eq!(all["total"], 25);
2823        assert_eq!(ids(&all).len(), 25);
2824
2825        // limit bounds the body; total still reports the full match set.
2826        let p1 = fetch(&app, "/api/v1/projections/demo/state?limit=10").await;
2827        assert_eq!(ids(&p1).len(), 10, "limit must bound the response");
2828        assert_eq!(p1["total"], 25);
2829        assert_eq!(p1["count"], 10);
2830        assert_eq!(p1["has_more"], true);
2831
2832        // offset resumes where the previous page stopped.
2833        let p2 = fetch(&app, "/api/v1/projections/demo/state?limit=10&offset=10").await;
2834        let p3 = fetch(&app, "/api/v1/projections/demo/state?limit=10&offset=20").await;
2835        assert_eq!(ids(&p3).len(), 5, "last page returns the remainder");
2836        assert_eq!(p3["has_more"], false, "offset + count == total → exhausted");
2837        assert_ne!(ids(&p1), ids(&p2), "offset=10 must skip the first page");
2838
2839        let mut walked = ids(&p1);
2840        walked.extend(ids(&p2));
2841        walked.extend(ids(&p3));
2842        let unique: std::collections::HashSet<_> = walked.iter().cloned().collect();
2843        assert_eq!(
2844            unique.len(),
2845            25,
2846            "paging the whole projection must yield 25 distinct entities"
2847        );
2848
2849        // Paging is only meaningful over a stable order — DashMap iteration is not.
2850        let mut sorted = walked.clone();
2851        sorted.sort();
2852        assert_eq!(walked, sorted, "pages must be ordered by entity_id");
2853
2854        // Past the end: empty page, exhausted rather than "more".
2855        let past = fetch(&app, "/api/v1/projections/demo/state?limit=10&offset=100").await;
2856        assert_eq!(past["count"], 0);
2857        assert_eq!(past["has_more"], false);
2858        assert_eq!(past["total"], 25);
2859
2860        // entity_id_prefix narrows to one shard of the keyspace.
2861        let shard = fetch(
2862            &app,
2863            "/api/v1/projections/demo/state?entity_id_prefix=tenant-1",
2864        )
2865        .await;
2866        assert_eq!(shard["total"], 10, "tenant-10..tenant-19");
2867        assert!(
2868            ids(&shard).iter().all(|id| id.starts_with("tenant-1")),
2869            "entity_id_prefix must filter"
2870        );
2871
2872        // The projection scope itself still holds under paging.
2873        let other = fetch(&app, "/api/v1/projections/other/state?limit=10").await;
2874        assert_eq!(ids(&other), vec!["tenant-99".to_string()]);
2875    }
2876
2877    // Regression guard for issue #251: a bounded page must cost the page, not the
2878    // whole match set. The handler used to re-run the query with `limit: None`
2879    // just to compute `total_count`, so `?entity_id=E&limit=1&order=desc` — the
2880    // documented "latest event for an entity" read — cloned and sorted the
2881    // entity's entire history on every request.
2882    //
2883    // Cost is measured by counting `Event` CLONES (`crate::clone_probe`), NOT
2884    // with Core's `query_results_total` metric: that counter is incremented with
2885    // `results.len()`, i.e. rows RETURNED, so it reads 1 whether the store
2886    // clones one event or clones 200 and discards 199 — it cannot fail on a
2887    // revert of the store-side windowing. Drives the REAL handler through
2888    // query-string deserialization so the count is what an HTTP caller pays.
2889    #[test]
2890    fn query_events_limit_does_not_materialize_whole_history() {
2891        const HISTORY: usize = 200;
2892
2893        let store = create_test_store();
2894        for _ in 0..HISTORY {
2895            store
2896                .ingest(&create_test_event("entity-hot", "user.updated"))
2897                .unwrap();
2898        }
2899
2900        // `clone_probe` is thread-local, so the handler future is driven to
2901        // completion on THIS thread (current-thread runtime, inside the measured
2902        // closure) — every clone the request makes is therefore counted.
2903        let runtime = tokio::runtime::Builder::new_current_thread()
2904            .enable_all()
2905            .build()
2906            .unwrap();
2907        let (resp, materialized) = crate::clone_probe::measure(|| {
2908            runtime.block_on(query_page(
2909                &store,
2910                "entity_id=entity-hot&limit=1&order=desc",
2911            ))
2912        });
2913
2914        // The page itself stays correct: one event, and the total/has_more pair
2915        // still describes the full match set.
2916        assert_eq!(resp.count, 1);
2917        assert_eq!(resp.events.len(), 1);
2918        assert_eq!(resp.total_count, HISTORY);
2919        assert!(resp.has_more);
2920
2921        assert_eq!(
2922            materialized, 1,
2923            "limit=1 materialized {materialized} events out of {HISTORY}: \
2924             `limit` must bound what a request materializes, not just what it \
2925             returns"
2926        );
2927    }
2928
2929    /// Like [`query_page`] but surfaces the handler's error instead of
2930    /// unwrapping — for the parameter values that must be REJECTED.
2931    async fn query_page_result(
2932        store: &SharedStore,
2933        query: &str,
2934    ) -> Result<Json<QueryEventsResponse>> {
2935        use axum::extract::{Query, State};
2936
2937        let uri: axum::http::Uri = format!("/api/v1/events/query?tenant_id=test-stream&{query}")
2938            .parse()
2939            .unwrap();
2940        query_events(
2941            OptionalAuth(None),
2942            Query::try_from_uri(&uri).unwrap(),
2943            Query::try_from_uri(&uri).unwrap(),
2944            Query::try_from_uri(&uri).unwrap(),
2945            Query::try_from_uri(&uri).unwrap(),
2946            State(store.clone()),
2947        )
2948        .await
2949    }
2950
2951    // A filter the server cannot apply must be an ERROR, not silence. The
2952    // payload filter is parsed inside the per-event predicate with `if let
2953    // Ok(..)`, so an unparseable one simply never matched anything and the
2954    // query answered as if no filter had been sent — the caller asked for
2955    // "events where user_id = alice" and got the tenant's whole stream, with a
2956    // `total_count` that agreed. Fails OPEN, which is the dangerous direction
2957    // for a filter, and the endpoint already rejects an unusable `order`, so
2958    // silence here was also inconsistent.
2959    #[tokio::test]
2960    async fn query_events_rejects_a_payload_filter_it_cannot_apply() {
2961        let store = create_test_store();
2962        for name in ["alice", "bob", "carol"] {
2963            let mut event = create_test_event(name, "user.created");
2964            event.payload = serde_json::json!({ "user_id": name });
2965            store.ingest(&event).unwrap();
2966        }
2967
2968        // A well-formed filter still filters — the guard must not reject the
2969        // shape callers actually send (`{"user_id":"alice"}`, percent-encoded).
2970        let ok = query_page(&store, "payload_filter=%7B%22user_id%22%3A%22alice%22%7D").await;
2971        assert_eq!(ok.count, 1, "a valid payload_filter must still work");
2972        assert_eq!(ok.total_count, 1);
2973
2974        for bad in [
2975            "not-json",                    // not JSON at all
2976            "%7B%22user_id%22%3A%22alice", // truncated object
2977            "%5B%22alice%22%5D",           // valid JSON, but an array, not an object
2978            "42",                          // valid JSON, but a scalar
2979        ] {
2980            let result = query_page_result(&store, &format!("payload_filter={bad}")).await;
2981            let Err(err) = result else {
2982                let resp = result.unwrap().0;
2983                panic!(
2984                    "payload_filter={bad} was silently ignored: returned {} of {} \
2985                     events unfiltered instead of rejecting a filter the server \
2986                     cannot apply",
2987                    resp.count, resp.total_count
2988                );
2989            };
2990            assert!(
2991                matches!(err, crate::error::AllSourceError::InvalidInput(_)),
2992                "payload_filter={bad} must be a 400, got {err:?}"
2993            );
2994        }
2995    }
2996
2997    // `order` is the other parameter whose only legal values are a closed set.
2998    // It is validated, but nothing pinned that: collapsing the match to a
2999    // `_ => false` default arm would make `?order=descending` silently return
3000    // OLDEST-first — a paginator would read the wrong end of the stream and
3001    // nothing in the response would say so.
3002    #[tokio::test]
3003    async fn query_events_rejects_an_unusable_order_value() {
3004        let store = create_test_store();
3005        store
3006            .ingest(&create_test_event("e-1", "user.created"))
3007            .unwrap();
3008
3009        for bad in ["descending", "DESCENDING", "newest", "1", "asc%20"] {
3010            let result = query_page_result(&store, &format!("order={bad}")).await;
3011            let Err(err) = result else {
3012                panic!("order={bad} must be rejected, not silently defaulted");
3013            };
3014            assert!(
3015                matches!(err, crate::error::AllSourceError::InvalidInput(_)),
3016                "order={bad} must be a 400, got {err:?}"
3017            );
3018        }
3019
3020        // The accepted spellings stay accepted, in any case.
3021        for good in ["asc", "ASC", "desc", "DeSc"] {
3022            let accepted = query_page_result(&store, &format!("order={good}"))
3023                .await
3024                .unwrap_or_else(|e| panic!("order={good} must be accepted: {e:?}"));
3025            assert_eq!(accepted.0.count, 1, "order={good}");
3026        }
3027    }
3028
3029    // `exclude_event_type_prefix` had no test anywhere in Core — not at this
3030    // handler, not in the store — despite being an HTTP-exposed filter the
3031    // Query Service forwards verbatim, and despite #251 having just rewritten
3032    // the code that decides WHEN it runs relative to the window. Its contract
3033    // is more than "drop these events": the DTO promises exclusion happens
3034    // BEFORE the limit, so excluded events never consume the result window, and
3035    // the prefixes are comma-separated. Both claims are only observable in
3036    // composition with `limit`/`offset`, so this drives the real handler.
3037    #[tokio::test]
3038    async fn query_events_exclude_prefix_applies_before_the_window() {
3039        const TYPES: [&str; 4] = [
3040            "audit.write",
3041            "user.created",
3042            "service.ping",
3043            "user.updated",
3044        ];
3045        let store = create_test_store();
3046        let base = chrono::Utc::now() - chrono::Duration::hours(24);
3047        let mut kept = Vec::new();
3048        for i in 0..12i64 {
3049            let mut event = create_test_event("org-1", TYPES[i as usize % 4]);
3050            event.timestamp = base + chrono::Duration::minutes(i);
3051            event.version = i + 1;
3052            if i % 2 == 1 {
3053                kept.push(event.id);
3054            }
3055            store.ingest(&event).unwrap();
3056        }
3057        assert_eq!(kept.len(), 6, "half the stream is user.*");
3058
3059        // A page of 3 must be 3 SURVIVING events. Excluding after the window
3060        // instead would serve 3 minus however many the page happened to hit —
3061        // here 1 — while still reporting a plausible-looking count.
3062        let page = query_page(
3063            &store,
3064            "exclude_event_type_prefix=audit.,%20service.&limit=3",
3065        )
3066        .await;
3067        assert_eq!(
3068            page.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3069            kept[..3].to_vec(),
3070            "limit=3 must yield 3 non-excluded events: exclusion runs before \
3071             the window, and the prefix list is comma-separated (whitespace \
3072             trimmed)"
3073        );
3074        assert!(
3075            page.events
3076                .iter()
3077                .all(|e| !e.event_type.starts_with("audit.")
3078                    && !e.event_type.starts_with("service.")),
3079            "excluded namespaces must not appear: {:?}",
3080            page.events
3081                .iter()
3082                .map(|e| &e.event_type)
3083                .collect::<Vec<_>>()
3084        );
3085        assert_eq!(
3086            page.total_count, 6,
3087            "total_count is the post-exclusion match set, not the 12 ingested"
3088        );
3089        assert!(page.has_more);
3090
3091        // Paging the excluded view walks the 6 survivors exactly once and stops.
3092        let mut walked = Vec::new();
3093        for offset in [0, 3, 6] {
3094            let p = query_page(
3095                &store,
3096                &format!("exclude_event_type_prefix=audit.,service.&limit=3&offset={offset}"),
3097            )
3098            .await;
3099            assert_eq!(
3100                p.has_more,
3101                offset + p.count < 6,
3102                "has_more must terminate on the excluded view (offset={offset})"
3103            );
3104            walked.extend(p.events.iter().map(|e| e.id));
3105        }
3106        assert_eq!(
3107            walked, kept,
3108            "exclusion + paging must cover the survivors once"
3109        );
3110
3111        // Composes with order=desc: newest survivor first, not newest event.
3112        let desc = query_page(
3113            &store,
3114            "exclude_event_type_prefix=audit.,service.&limit=1&order=desc",
3115        )
3116        .await;
3117        assert_eq!(
3118            desc.events[0].id,
3119            *kept.last().unwrap(),
3120            "order=desc&limit=1 over an excluded view is the newest SURVIVOR"
3121        );
3122
3123        // One prefix excludes only its namespace; a prefix matching nothing
3124        // excludes nothing.
3125        let audit_only = query_page(&store, "exclude_event_type_prefix=audit.").await;
3126        assert_eq!(audit_only.total_count, 9, "12 minus the 3 audit.* events");
3127        let nothing = query_page(&store, "exclude_event_type_prefix=nosuch.").await;
3128        assert_eq!(nothing.total_count, 12);
3129    }
3130
3131    // `since`/`until`/`as_of` must narrow the result set on EVERY query shape,
3132    // including the one no index narrows: scoped by tenant only, which is what
3133    // the gateway forwards for "this tenant's activity since T" (the Query
3134    // Service passes all three straight through — `@core_compat_filters`).
3135    // Those three filters are evaluated against index entries, and the
3136    // full-scan branch of the store never consults the index, so a tenant-only
3137    // query silently ignored the window and answered with the whole history —
3138    // with `total_count`/`has_more` describing that history too, so a paginator
3139    // walked events the caller had explicitly excluded.
3140    //
3141    // Driven through the real handler and real query-string deserialization:
3142    // the parameters exist on the DTO and parse fine, so nothing at the wire
3143    // layer reveals that no code downstream reads them on this path.
3144    #[tokio::test]
3145    async fn query_events_honours_time_window_without_an_entity_or_type_filter() {
3146        use chrono::{SecondsFormat, SubsecRound};
3147
3148        let store = create_test_store();
3149        // FIXED base, not `Utc::now()`, and truncated to whole seconds.
3150        //
3151        // This test previously seeded from `Utc::now()`, which made it depend on
3152        // the host clock's resolution: Linux hands back nanoseconds, macOS does
3153        // not. `SecondsFormat::Micros` truncates DOWNWARDS, so with a
3154        // nanosecond-bearing base the string `until=<t>` names an instant
3155        // fractionally BEFORE the event stamped at `t`, and an inclusive window
3156        // silently drops its boundary event. It passed on every developer
3157        // machine and failed only in CI — the worst kind of flake.
3158        //
3159        // The nanoseconds below are deliberate: they make the round-trip
3160        // assertion in `at` fail loudly if the `trunc_subsecs(0)` is ever
3161        // removed, on every platform rather than only on one.
3162        let base = chrono::DateTime::parse_from_rfc3339("2026-01-01T00:00:00.123456789Z")
3163            .unwrap()
3164            .to_utc()
3165            .trunc_subsecs(0);
3166        let mut ids = Vec::new();
3167        for i in 0..5i64 {
3168            let mut event = create_test_event(&format!("e-{i}"), "user.created");
3169            event.timestamp = base + chrono::Duration::hours(i);
3170            event.version = i + 1;
3171            ids.push(event.id);
3172            store.ingest(&event).unwrap();
3173        }
3174        // `Z`-suffixed so the timestamp survives a query string — an offset of
3175        // `+00:00` would be decoded as a space.
3176        let at = |h: i64| {
3177            let t = base + chrono::Duration::hours(h);
3178            let s = t.to_rfc3339_opts(SecondsFormat::Micros, true);
3179            // Pin the invariant the whole test rests on: the formatted string
3180            // must name the SAME instant the event carries. If a future edit
3181            // drops the `trunc_subsecs(0)` above, this fails on every platform
3182            // rather than only on the ones with a nanosecond clock.
3183            assert_eq!(
3184                chrono::DateTime::parse_from_rfc3339(&s).unwrap().to_utc(),
3185                t,
3186                "query-string timestamp must round-trip exactly, else the window \
3187                 boundary silently excludes the event stamped at it"
3188            );
3189            s
3190        };
3191
3192        for (qs, expected) in [
3193            (format!("since={}", at(2)), vec![ids[2], ids[3], ids[4]]),
3194            (format!("until={}", at(1)), vec![ids[0], ids[1]]),
3195            (format!("as_of={}", at(1)), vec![ids[0], ids[1]]),
3196            (
3197                format!("since={}&until={}", at(1), at(3)),
3198                vec![ids[1], ids[2], ids[3]],
3199            ),
3200        ] {
3201            let resp = query_page(&store, &qs).await;
3202            let got: Vec<_> = resp.events.iter().map(|e| e.id).collect();
3203            assert_eq!(got, expected, "?{qs} must return only the window");
3204            assert_eq!(resp.count, expected.len(), "?{qs}");
3205            assert_eq!(
3206                resp.total_count,
3207                expected.len(),
3208                "?{qs}: total_count must count the window, not the history"
3209            );
3210            assert!(!resp.has_more, "?{qs}: the whole window was served");
3211        }
3212
3213        // The window composes with paging: page 2 of a `since` window is the
3214        // second page OF THAT WINDOW, and `has_more` terminates on it.
3215        let page1 = query_page(&store, &format!("since={}&limit=2", at(2))).await;
3216        let page2 = query_page(&store, &format!("since={}&limit=2&offset=2", at(2))).await;
3217        assert_eq!(
3218            page1.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3219            vec![ids[2], ids[3]]
3220        );
3221        assert!(page1.has_more, "3 in the window, 2 served");
3222        assert_eq!(
3223            page2.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3224            vec![ids[4]]
3225        );
3226        assert!(!page2.has_more, "offset 2 + count 1 == the window's 3");
3227        assert_eq!(page2.total_count, 3);
3228
3229        // …and with `order=desc`.
3230        let desc = query_page(&store, &format!("since={}&order=desc", at(2))).await;
3231        assert_eq!(
3232            desc.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3233            vec![ids[4], ids[3], ids[2]]
3234        );
3235
3236        // An empty window is empty, not "everything".
3237        let empty = query_page(&store, &format!("since={}", at(99))).await;
3238        assert_eq!(empty.count, 0);
3239        assert_eq!(empty.total_count, 0);
3240        assert!(!empty.has_more);
3241    }
3242
3243    // Tenant-isolation gate: the public events query must fail CLOSED — a request
3244    // with no auth context and no tenant_id returns nothing, never a cross-tenant
3245    // scan. Calls the real handler so the boundary check is exercised.
3246    #[tokio::test]
3247    async fn query_events_fails_closed_without_tenant() {
3248        use axum::extract::{Query, State};
3249
3250        let store = create_test_store();
3251        for i in 0..5 {
3252            store
3253                .ingest(&create_test_event(&format!("e-{i}"), "user.created"))
3254                .unwrap();
3255        }
3256
3257        let resp = query_events(
3258            OptionalAuth(None),
3259            Query(QueryEventsRequest::default()),
3260            Query(EventOrderParam { order: None }),
3261            Query(EventOffsetParam { offset: None }),
3262            Query(EventIntegrityParam::default()),
3263            State(store.clone()),
3264        )
3265        .await
3266        .unwrap();
3267        assert_eq!(
3268            resp.0.total_count, 0,
3269            "a no-tenant query must NOT return cross-tenant events"
3270        );
3271        assert_eq!(resp.0.count, 0);
3272
3273        // The same query scoped to the events' tenant returns them.
3274        // (create_test_event stamps tenant "test-stream" — the 3rd from_strings arg.)
3275        let scoped = query_events(
3276            OptionalAuth(None),
3277            Query(QueryEventsRequest {
3278                tenant_id: Some("test-stream".to_string()),
3279                ..QueryEventsRequest::default()
3280            }),
3281            Query(EventOrderParam { order: None }),
3282            Query(EventOffsetParam { offset: None }),
3283            Query(EventIntegrityParam::default()),
3284            State(store),
3285        )
3286        .await
3287        .unwrap();
3288        assert_eq!(
3289            scoped.0.total_count, 5,
3290            "tenant-scoped query returns its events"
3291        );
3292    }
3293
3294    // The dashboard's streams + event-types counts must be per-tenant. These
3295    // endpoints used to scan ALL tenants (platform totals shown as "yours", and a
3296    // cross-tenant spill). Assert each tenant sees only its own, and no-tenant
3297    // fails closed.
3298    #[tokio::test]
3299    async fn list_streams_and_types_are_tenant_scoped() {
3300        use crate::domain::entities::Event;
3301        use axum::extract::{Query, State};
3302
3303        let store = create_test_store();
3304        let ev = |entity: &str, etype: &str, tenant: &str| {
3305            Event::from_strings(
3306                etype.to_string(),
3307                entity.to_string(),
3308                tenant.to_string(),
3309                serde_json::json!({}),
3310                None,
3311            )
3312            .unwrap()
3313        };
3314        // tenant A: 2 entities, 2 types. tenant B: 1 entity, 1 type.
3315        store.ingest(&ev("e1", "order.placed", "tenant-a")).unwrap();
3316        store.ingest(&ev("e2", "user.created", "tenant-a")).unwrap();
3317        store
3318            .ingest(&ev("e9", "thing.happened", "tenant-b"))
3319            .unwrap();
3320
3321        let streams = |tid: Option<&str>| {
3322            list_streams(
3323                OptionalAuth(None),
3324                State(store.clone()),
3325                Query(ListStreamsParams {
3326                    tenant_id: tid.map(String::from),
3327                    limit: None,
3328                    offset: None,
3329                }),
3330            )
3331        };
3332        assert_eq!(
3333            streams(Some("tenant-a")).await.0.total,
3334            2,
3335            "tenant-a streams"
3336        );
3337        assert_eq!(
3338            streams(Some("tenant-b")).await.0.total,
3339            1,
3340            "tenant-b streams"
3341        );
3342        assert_eq!(
3343            streams(None).await.0.total,
3344            0,
3345            "no tenant -> no streams (fail closed)"
3346        );
3347
3348        let types = |tid: Option<&str>| {
3349            list_event_types(
3350                OptionalAuth(None),
3351                State(store.clone()),
3352                Query(ListEventTypesParams {
3353                    tenant_id: tid.map(String::from),
3354                    limit: None,
3355                    offset: None,
3356                }),
3357            )
3358        };
3359        assert_eq!(
3360            types(Some("tenant-a")).await.0.total,
3361            2,
3362            "tenant-a event types"
3363        );
3364        assert_eq!(
3365            types(Some("tenant-b")).await.0.total,
3366            1,
3367            "tenant-b event types"
3368        );
3369        assert_eq!(
3370            types(None).await.0.total,
3371            0,
3372            "no tenant -> no types (fail closed)"
3373        );
3374    }
3375
3376    #[tokio::test]
3377    async fn test_query_events_no_more_results() {
3378        let store = create_test_store();
3379
3380        // Ingest 5 events
3381        for i in 0..5 {
3382            store
3383                .ingest(&create_test_event(&format!("entity-{i}"), "user.created"))
3384                .unwrap();
3385        }
3386
3387        // Query with limit=100 — should get has_more=false, total_count=5
3388        let all_events = store
3389            .query(&QueryEventsRequest {
3390                entity_id: None,
3391                event_type: None,
3392                tenant_id: None,
3393                as_of: None,
3394                since: None,
3395                until: None,
3396                limit: None,
3397                event_type_prefix: None,
3398                exclude_event_type_prefix: None,
3399                payload_filter: None,
3400            })
3401            .unwrap();
3402        let total_count = all_events.len();
3403        let limited_events: Vec<Event> = all_events.into_iter().take(100).collect();
3404        let count = limited_events.len();
3405        let has_more = count < total_count;
3406
3407        assert_eq!(count, 5);
3408        assert_eq!(total_count, 5);
3409        assert!(!has_more);
3410    }
3411
3412    // Regression test for issue #177: `order=desc` + `limit=1` must return
3413    // the NEWEST event for an entity, not the oldest. Mirrors the ordering
3414    // logic in `query_events` (store.query → reverse-if-desc → take(limit)).
3415    #[tokio::test]
3416    async fn test_query_events_order_desc_returns_latest() {
3417        // Drives the REAL handler. An earlier version of this test rebuilt the
3418        // ordering inline (`ascending.clone(); reverse(); take(1)`) and never
3419        // called `query_events`, so it could not fail on a mis-composition in
3420        // the code that actually serves `order=desc` — which since issue #251
3421        // lives in `EventStore::query_window`, not in the handler.
3422        let store = create_test_store();
3423
3424        // Five events for the same entity with strictly increasing timestamps —
3425        // mimics a backfill appending corrected events.
3426        let base = chrono::Utc::now();
3427        let mut ascending_ids = Vec::new();
3428        for i in 0..5i64 {
3429            let mut event = create_test_event("org-1", "auth.org.updated");
3430            event.timestamp = base + chrono::Duration::seconds(i);
3431            event.version = i + 1;
3432            ascending_ids.push(event.id);
3433            store.ingest(&event).unwrap();
3434        }
3435        let newest_ts = base + chrono::Duration::seconds(4);
3436
3437        // The documented "latest event for an entity" read.
3438        let latest = query_page(&store, "entity_id=org-1&limit=1&order=desc").await;
3439        assert_eq!(latest.count, 1);
3440        assert_eq!(
3441            latest.events[0].id,
3442            ascending_ids[4],
3443            "order=desc&limit=1 must yield the NEWEST event, got the one at \
3444             ascending position {:?}",
3445            ascending_ids
3446                .iter()
3447                .position(|id| *id == latest.events[0].id)
3448        );
3449        assert_eq!(latest.events[0].timestamp, newest_ts);
3450        assert_eq!(latest.total_count, 5, "total is the full match set");
3451        assert!(latest.has_more);
3452
3453        // Default order (and an explicit `asc`) still yields the OLDEST.
3454        for qs in [
3455            "entity_id=org-1&limit=1",
3456            "entity_id=org-1&limit=1&order=asc",
3457        ] {
3458            let oldest = query_page(&store, qs).await;
3459            assert_eq!(oldest.events[0].id, ascending_ids[0], "{qs}");
3460            assert_eq!(oldest.events[0].timestamp, base);
3461        }
3462
3463        // An unbounded desc page is the exact reverse of the ascending one —
3464        // reversal must apply to the whole match set, not just to the page.
3465        let all_desc = query_page(&store, "entity_id=org-1&order=desc").await;
3466        let got: Vec<_> = all_desc.events.iter().map(|e| e.id).collect();
3467        let expected: Vec<_> = ascending_ids.iter().rev().copied().collect();
3468        assert_eq!(got, expected, "order=desc must return newest-first");
3469    }
3470
3471    // Regression guard for issue #251's relocation of the ordering: `order=desc`
3472    // moved out of the handler and into `EventStore::query_window`, where it now
3473    // composes with `offset` and `limit`. The contract is reverse-THEN-skip-THEN-
3474    // take: `order=desc&offset=1&limit=2` is "the 2nd and 3rd newest". Skipping
3475    // before reversing (or reversing only the page) returns a different, quietly
3476    // wrong page — with the same count, total_count and has_more.
3477    #[tokio::test]
3478    async fn query_events_desc_composes_with_offset_and_limit() {
3479        let store = create_test_store();
3480        let base = chrono::Utc::now();
3481        let mut ascending_ids = Vec::new();
3482        for i in 0..5i64 {
3483            let mut event = create_test_event("org-1", "auth.org.updated");
3484            event.timestamp = base + chrono::Duration::seconds(i);
3485            event.version = i + 1;
3486            ascending_ids.push(event.id);
3487            store.ingest(&event).unwrap();
3488        }
3489        let newest_first: Vec<_> = ascending_ids.iter().rev().copied().collect();
3490
3491        for (offset, limit) in [(0, 2), (1, 2), (2, 2), (3, 2), (4, 2), (5, 2), (1, 4)] {
3492            let page = query_page(
3493                &store,
3494                &format!("entity_id=org-1&order=desc&offset={offset}&limit={limit}"),
3495            )
3496            .await;
3497            let got: Vec<_> = page.events.iter().map(|e| e.id).collect();
3498            let expected: Vec<_> = newest_first
3499                .iter()
3500                .skip(offset)
3501                .take(limit)
3502                .copied()
3503                .collect();
3504            assert_eq!(
3505                got, expected,
3506                "order=desc&offset={offset}&limit={limit} must reverse, then \
3507                 skip, then take"
3508            );
3509            assert_eq!(page.count, expected.len());
3510            assert_eq!(page.total_count, 5);
3511            assert_eq!(
3512                page.has_more,
3513                offset + expected.len() < 5,
3514                "has_more must account for the offset (offset={offset})"
3515            );
3516        }
3517
3518        // Walking the whole entity newest-first must visit every event exactly
3519        // once — the property a `order=desc` paginator depends on.
3520        let mut walked = Vec::new();
3521        for offset in (0..5).step_by(2) {
3522            let page = query_page(
3523                &store,
3524                &format!("entity_id=org-1&order=desc&offset={offset}&limit=2"),
3525            )
3526            .await;
3527            walked.extend(page.events.iter().map(|e| e.id));
3528        }
3529        assert_eq!(walked, newest_first, "desc paging must cover the set once");
3530    }
3531
3532    #[tokio::test]
3533    async fn test_list_entities_by_type_prefix() {
3534        let store = create_test_store();
3535
3536        // 3 index entities
3537        store
3538            .ingest(&create_test_event("idx-1", "index.created"))
3539            .unwrap();
3540        store
3541            .ingest(&create_test_event("idx-1", "index.updated"))
3542            .unwrap();
3543        store
3544            .ingest(&create_test_event("idx-2", "index.created"))
3545            .unwrap();
3546        store
3547            .ingest(&create_test_event("idx-3", "index.created"))
3548            .unwrap();
3549        // 2 trade entities
3550        store
3551            .ingest(&create_test_event("trade-1", "trade.created"))
3552            .unwrap();
3553        store
3554            .ingest(&create_test_event("trade-2", "trade.created"))
3555            .unwrap();
3556
3557        // List entities for index.*
3558        let req = ListEntitiesRequest {
3559            event_type_prefix: Some("index.".to_string()),
3560            ..Default::default()
3561        };
3562        let query_req = QueryEventsRequest {
3563            entity_id: None,
3564            event_type: None,
3565            tenant_id: None,
3566            as_of: None,
3567            since: None,
3568            until: None,
3569            limit: None,
3570            event_type_prefix: req.event_type_prefix,
3571            exclude_event_type_prefix: None,
3572            payload_filter: req.payload_filter,
3573        };
3574        let events = store.query(&query_req).unwrap();
3575
3576        // Group and verify
3577        let mut entity_map: std::collections::HashMap<String, Vec<&Event>> =
3578            std::collections::HashMap::new();
3579        for event in &events {
3580            entity_map
3581                .entry(event.entity_id().to_string())
3582                .or_default()
3583                .push(event);
3584        }
3585
3586        assert_eq!(entity_map.len(), 3); // idx-1, idx-2, idx-3
3587        assert_eq!(entity_map["idx-1"].len(), 2); // 2 events for idx-1
3588        assert_eq!(entity_map["idx-2"].len(), 1);
3589        assert_eq!(entity_map["idx-3"].len(), 1);
3590    }
3591
3592    // Issue #178: `list_entities` accepts an `order` param and pages
3593    // deterministically over the resulting sort.
3594    #[tokio::test]
3595    async fn test_list_entities_order_and_pagination() {
3596        let store = create_test_store();
3597
3598        // Three entities with strictly increasing last-event times.
3599        let base = chrono::Utc::now();
3600        for (i, eid) in ["org-a", "org-b", "org-c"].iter().enumerate() {
3601            let mut event = create_test_event(eid, "auth.org.created");
3602            event.timestamp = base + chrono::Duration::seconds(i as i64);
3603            store.ingest(&event).unwrap();
3604        }
3605        let prefix = || Some("auth.org.".to_string());
3606
3607        // Default: newest activity first (desc) — preserves prior behavior.
3608        let desc = list_entities(
3609            State(store.clone()),
3610            Query(ListEntitiesRequest {
3611                event_type_prefix: prefix(),
3612                ..Default::default()
3613            }),
3614        )
3615        .await
3616        .unwrap();
3617        let desc_ids: Vec<&str> = desc
3618            .0
3619            .entities
3620            .iter()
3621            .map(|e| e.entity_id.as_str())
3622            .collect();
3623        assert_eq!(desc_ids, ["org-c", "org-b", "org-a"]);
3624
3625        // order=asc: oldest activity first.
3626        let asc = list_entities(
3627            State(store.clone()),
3628            Query(ListEntitiesRequest {
3629                event_type_prefix: prefix(),
3630                order: Some("asc".to_string()),
3631                ..Default::default()
3632            }),
3633        )
3634        .await
3635        .unwrap();
3636        let asc_ids: Vec<&str> = asc
3637            .0
3638            .entities
3639            .iter()
3640            .map(|e| e.entity_id.as_str())
3641            .collect();
3642        assert_eq!(asc_ids, ["org-a", "org-b", "org-c"]);
3643
3644        // Offset pagination over the deterministic asc order: page 2, size 1.
3645        let page2 = list_entities(
3646            State(store.clone()),
3647            Query(ListEntitiesRequest {
3648                event_type_prefix: prefix(),
3649                order: Some("asc".to_string()),
3650                limit: Some(1),
3651                offset: Some(1),
3652                ..Default::default()
3653            }),
3654        )
3655        .await
3656        .unwrap();
3657        assert_eq!(page2.0.entities.len(), 1);
3658        assert_eq!(page2.0.entities[0].entity_id, "org-b");
3659        assert_eq!(page2.0.total, 3);
3660        assert!(page2.0.has_more);
3661
3662        // Invalid order value is rejected.
3663        let err = list_entities(
3664            State(store.clone()),
3665            Query(ListEntitiesRequest {
3666                event_type_prefix: prefix(),
3667                order: Some("sideways".to_string()),
3668                ..Default::default()
3669            }),
3670        )
3671        .await;
3672        assert!(err.is_err(), "invalid order value must be rejected");
3673    }
3674
3675    fn create_test_event_with_payload(
3676        entity_id: &str,
3677        event_type: &str,
3678        payload: serde_json::Value,
3679    ) -> Event {
3680        Event::from_strings(
3681            event_type.to_string(),
3682            entity_id.to_string(),
3683            "test-stream".to_string(),
3684            payload,
3685            None,
3686        )
3687        .unwrap()
3688    }
3689
3690    #[tokio::test]
3691    async fn test_detect_duplicates_by_payload_fields() {
3692        let store = create_test_store();
3693
3694        // Create entities with duplicate "name" field values
3695        store
3696            .ingest(&create_test_event_with_payload(
3697                "idx-1",
3698                "index.created",
3699                serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3700            ))
3701            .unwrap();
3702        store
3703            .ingest(&create_test_event_with_payload(
3704                "idx-2",
3705                "index.created",
3706                serde_json::json!({"name": "S&P 500", "user_id": "bob"}),
3707            ))
3708            .unwrap();
3709        store
3710            .ingest(&create_test_event_with_payload(
3711                "idx-3",
3712                "index.created",
3713                serde_json::json!({"name": "NASDAQ", "user_id": "alice"}),
3714            ))
3715            .unwrap();
3716        store
3717            .ingest(&create_test_event_with_payload(
3718                "idx-4",
3719                "index.created",
3720                serde_json::json!({"name": "NASDAQ", "user_id": "carol"}),
3721            ))
3722            .unwrap();
3723        store
3724            .ingest(&create_test_event_with_payload(
3725                "idx-5",
3726                "index.created",
3727                serde_json::json!({"name": "DAX", "user_id": "dave"}),
3728            ))
3729            .unwrap();
3730
3731        // Group by name — should find 2 groups: "S&P 500" (idx-1, idx-2) and "NASDAQ" (idx-3, idx-4)
3732        let query_req = QueryEventsRequest {
3733            entity_id: None,
3734            event_type: None,
3735            tenant_id: None,
3736            as_of: None,
3737            since: None,
3738            until: None,
3739            limit: None,
3740            event_type_prefix: Some("index.".to_string()),
3741            exclude_event_type_prefix: None,
3742            payload_filter: None,
3743        };
3744        let events = store.query(&query_req).unwrap();
3745
3746        // Manually replicate the handler logic for testing
3747        let group_by_fields = vec!["name"];
3748        let mut entity_latest: std::collections::HashMap<String, &Event> =
3749            std::collections::HashMap::new();
3750        for event in &events {
3751            let eid = event.entity_id().to_string();
3752            entity_latest
3753                .entry(eid)
3754                .and_modify(|existing| {
3755                    if event.timestamp() > existing.timestamp() {
3756                        *existing = event;
3757                    }
3758                })
3759                .or_insert(event);
3760        }
3761
3762        let mut groups: std::collections::HashMap<String, Vec<String>> =
3763            std::collections::HashMap::new();
3764        for (entity_id, event) in &entity_latest {
3765            let payload = event.payload();
3766            let mut key_parts = serde_json::Map::new();
3767            for field in &group_by_fields {
3768                let value = payload
3769                    .get(*field)
3770                    .cloned()
3771                    .unwrap_or(serde_json::Value::Null);
3772                key_parts.insert((*field).to_string(), value);
3773            }
3774            let key_str = serde_json::to_string(&key_parts).unwrap_or_default();
3775            groups.entry(key_str).or_default().push(entity_id.clone());
3776        }
3777
3778        let duplicate_groups: Vec<_> = groups
3779            .into_iter()
3780            .filter(|(_, ids)| ids.len() > 1)
3781            .collect();
3782
3783        assert_eq!(duplicate_groups.len(), 2); // S&P 500 and NASDAQ groups
3784        for (_, ids) in &duplicate_groups {
3785            assert_eq!(ids.len(), 2);
3786        }
3787    }
3788
3789    #[tokio::test]
3790    async fn test_detect_duplicates_no_duplicates() {
3791        let store = create_test_store();
3792
3793        // All unique names
3794        store
3795            .ingest(&create_test_event_with_payload(
3796                "idx-1",
3797                "index.created",
3798                serde_json::json!({"name": "A"}),
3799            ))
3800            .unwrap();
3801        store
3802            .ingest(&create_test_event_with_payload(
3803                "idx-2",
3804                "index.created",
3805                serde_json::json!({"name": "B"}),
3806            ))
3807            .unwrap();
3808
3809        let query_req = QueryEventsRequest {
3810            entity_id: None,
3811            event_type: None,
3812            tenant_id: None,
3813            as_of: None,
3814            since: None,
3815            until: None,
3816            limit: None,
3817            event_type_prefix: Some("index.".to_string()),
3818            exclude_event_type_prefix: None,
3819            payload_filter: None,
3820        };
3821        let events = store.query(&query_req).unwrap();
3822
3823        let mut entity_latest: std::collections::HashMap<String, &Event> =
3824            std::collections::HashMap::new();
3825        for event in &events {
3826            entity_latest
3827                .entry(event.entity_id().to_string())
3828                .or_insert(event);
3829        }
3830
3831        let mut groups: std::collections::HashMap<String, Vec<String>> =
3832            std::collections::HashMap::new();
3833        for (entity_id, event) in &entity_latest {
3834            let key_str =
3835                serde_json::to_string(&serde_json::json!({"name": event.payload().get("name")}))
3836                    .unwrap();
3837            groups.entry(key_str).or_default().push(entity_id.clone());
3838        }
3839
3840        let duplicate_groups: Vec<_> = groups
3841            .into_iter()
3842            .filter(|(_, ids)| ids.len() > 1)
3843            .collect();
3844
3845        assert_eq!(duplicate_groups.len(), 0); // No duplicates
3846    }
3847
3848    #[tokio::test]
3849    async fn test_detect_duplicates_multi_field_group_by() {
3850        let store = create_test_store();
3851
3852        // Two entities with same name AND user_id = true duplicate
3853        store
3854            .ingest(&create_test_event_with_payload(
3855                "idx-1",
3856                "index.created",
3857                serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3858            ))
3859            .unwrap();
3860        store
3861            .ingest(&create_test_event_with_payload(
3862                "idx-2",
3863                "index.created",
3864                serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3865            ))
3866            .unwrap();
3867        // Same name but different user_id = NOT a duplicate in multi-field group
3868        store
3869            .ingest(&create_test_event_with_payload(
3870                "idx-3",
3871                "index.created",
3872                serde_json::json!({"name": "S&P 500", "user_id": "bob"}),
3873            ))
3874            .unwrap();
3875
3876        let query_req = QueryEventsRequest {
3877            entity_id: None,
3878            event_type: None,
3879            tenant_id: None,
3880            as_of: None,
3881            since: None,
3882            until: None,
3883            limit: None,
3884            event_type_prefix: Some("index.".to_string()),
3885            exclude_event_type_prefix: None,
3886            payload_filter: None,
3887        };
3888        let events = store.query(&query_req).unwrap();
3889
3890        let group_by_fields = vec!["name", "user_id"];
3891        let mut entity_latest: std::collections::HashMap<String, &Event> =
3892            std::collections::HashMap::new();
3893        for event in &events {
3894            entity_latest
3895                .entry(event.entity_id().to_string())
3896                .and_modify(|existing| {
3897                    if event.timestamp() > existing.timestamp() {
3898                        *existing = event;
3899                    }
3900                })
3901                .or_insert(event);
3902        }
3903
3904        let mut groups: std::collections::HashMap<String, Vec<String>> =
3905            std::collections::HashMap::new();
3906        for (entity_id, event) in &entity_latest {
3907            let payload = event.payload();
3908            let mut key_parts = serde_json::Map::new();
3909            for field in &group_by_fields {
3910                let value = payload
3911                    .get(*field)
3912                    .cloned()
3913                    .unwrap_or(serde_json::Value::Null);
3914                key_parts.insert((*field).to_string(), value);
3915            }
3916            let key_str = serde_json::to_string(&key_parts).unwrap_or_default();
3917            groups.entry(key_str).or_default().push(entity_id.clone());
3918        }
3919
3920        let duplicate_groups: Vec<_> = groups
3921            .into_iter()
3922            .filter(|(_, ids)| ids.len() > 1)
3923            .collect();
3924
3925        // Only 1 duplicate group: name=S&P 500, user_id=alice (idx-1, idx-2)
3926        assert_eq!(duplicate_groups.len(), 1);
3927        let (_, ref ids) = duplicate_groups[0];
3928        assert_eq!(ids.len(), 2);
3929        let mut sorted_ids = ids.clone();
3930        sorted_ids.sort();
3931        assert_eq!(sorted_ids, vec!["idx-1", "idx-2"]);
3932    }
3933
3934    #[tokio::test]
3935    async fn test_projection_state_cache() {
3936        let store = create_test_store();
3937
3938        // Test cache insertion
3939        let cache = store.projection_state_cache();
3940        cache.insert(
3941            "entity_snapshots:user-123".to_string(),
3942            serde_json::json!({"name": "Test User", "age": 30}),
3943        );
3944
3945        // Test cache retrieval
3946        let state = cache.get("entity_snapshots:user-123");
3947        assert!(state.is_some());
3948        let state = state.unwrap();
3949        assert_eq!(state["name"], "Test User");
3950        assert_eq!(state["age"], 30);
3951    }
3952
3953    #[tokio::test]
3954    async fn test_projection_manager_list_projections() {
3955        let store = create_test_store();
3956
3957        // List projections (built-in projections should be available)
3958        let projection_manager = store.projection_manager();
3959        let projections = projection_manager.list_projections();
3960
3961        // Should have entity_snapshots and event_counters
3962        assert!(projections.len() >= 2);
3963
3964        let names: Vec<&str> = projections.iter().map(|(name, _)| name.as_str()).collect();
3965        assert!(names.contains(&"entity_snapshots"));
3966        assert!(names.contains(&"event_counters"));
3967    }
3968
3969    #[tokio::test]
3970    async fn test_projection_state_after_event_ingestion() {
3971        let store = create_test_store();
3972
3973        // Ingest an event
3974        let event = create_test_event("user-456", "user.created");
3975        store.ingest(&event).unwrap();
3976
3977        // Get projection state
3978        let projection_manager = store.projection_manager();
3979        let snapshot_projection = projection_manager
3980            .get_projection("entity_snapshots")
3981            .unwrap();
3982
3983        let state = snapshot_projection.get_state("user-456");
3984        assert!(state.is_some());
3985        let state = state.unwrap();
3986        assert_eq!(state["name"], "Test");
3987        assert_eq!(state["value"], 42);
3988    }
3989
3990    #[tokio::test]
3991    async fn test_projection_state_cache_multiple_entities() {
3992        let store = create_test_store();
3993        let cache = store.projection_state_cache();
3994
3995        // Insert multiple entities
3996        for i in 0..10 {
3997            cache.insert(
3998                format!("entity_snapshots:entity-{i}"),
3999                serde_json::json!({"id": i, "status": "active"}),
4000            );
4001        }
4002
4003        // Verify all insertions
4004        assert_eq!(cache.len(), 10);
4005
4006        // Verify each entity
4007        for i in 0..10 {
4008            let key = format!("entity_snapshots:entity-{i}");
4009            let state = cache.get(&key);
4010            assert!(state.is_some());
4011            assert_eq!(state.unwrap()["id"], i);
4012        }
4013    }
4014
4015    #[tokio::test]
4016    async fn test_projection_state_update() {
4017        let store = create_test_store();
4018        let cache = store.projection_state_cache();
4019
4020        // Initial state
4021        cache.insert(
4022            "entity_snapshots:user-789".to_string(),
4023            serde_json::json!({"balance": 100}),
4024        );
4025
4026        // Update state
4027        cache.insert(
4028            "entity_snapshots:user-789".to_string(),
4029            serde_json::json!({"balance": 150}),
4030        );
4031
4032        // Verify update
4033        let state = cache.get("entity_snapshots:user-789").unwrap();
4034        assert_eq!(state["balance"], 150);
4035    }
4036
4037    #[tokio::test]
4038    async fn test_event_counter_projection() {
4039        let store = create_test_store();
4040
4041        // Ingest events of different types
4042        store
4043            .ingest(&create_test_event("user-1", "user.created"))
4044            .unwrap();
4045        store
4046            .ingest(&create_test_event("user-2", "user.created"))
4047            .unwrap();
4048        store
4049            .ingest(&create_test_event("user-1", "user.updated"))
4050            .unwrap();
4051
4052        // Get event counter projection
4053        let projection_manager = store.projection_manager();
4054        let counter_projection = projection_manager.get_projection("event_counters").unwrap();
4055
4056        // Check counts
4057        let created_state = counter_projection.get_state("user.created");
4058        assert!(created_state.is_some());
4059        assert_eq!(created_state.unwrap()["count"], 2);
4060
4061        let updated_state = counter_projection.get_state("user.updated");
4062        assert!(updated_state.is_some());
4063        assert_eq!(updated_state.unwrap()["count"], 1);
4064    }
4065
4066    #[tokio::test]
4067    async fn test_projection_state_cache_key_format() {
4068        let store = create_test_store();
4069        let cache = store.projection_state_cache();
4070
4071        // Test standard key format: {projection_name}:{entity_id}
4072        let key = "orders:order-12345".to_string();
4073        cache.insert(key.clone(), serde_json::json!({"total": 99.99}));
4074
4075        let state = cache.get(&key).unwrap();
4076        assert_eq!(state["total"], 99.99);
4077    }
4078
4079    #[tokio::test]
4080    async fn test_projection_state_cache_removal() {
4081        let store = create_test_store();
4082        let cache = store.projection_state_cache();
4083
4084        // Insert and then remove
4085        cache.insert(
4086            "test:entity-1".to_string(),
4087            serde_json::json!({"data": "value"}),
4088        );
4089        assert_eq!(cache.len(), 1);
4090
4091        cache.remove("test:entity-1");
4092        assert_eq!(cache.len(), 0);
4093        assert!(cache.get("test:entity-1").is_none());
4094    }
4095
4096    #[tokio::test]
4097    async fn test_get_nonexistent_projection() {
4098        let store = create_test_store();
4099        let projection_manager = store.projection_manager();
4100
4101        // Requesting a non-existent projection should return None
4102        let projection = projection_manager.get_projection("nonexistent_projection");
4103        assert!(projection.is_none());
4104    }
4105
4106    #[tokio::test]
4107    async fn test_get_nonexistent_entity_state() {
4108        let store = create_test_store();
4109        let projection_manager = store.projection_manager();
4110
4111        // Get state for non-existent entity
4112        let snapshot_projection = projection_manager
4113            .get_projection("entity_snapshots")
4114            .unwrap();
4115        let state = snapshot_projection.get_state("nonexistent-entity-xyz");
4116        assert!(state.is_none());
4117    }
4118
4119    #[tokio::test]
4120    async fn test_projection_state_cache_concurrent_access() {
4121        let store = create_test_store();
4122        let cache = store.projection_state_cache();
4123
4124        // Simulate concurrent writes
4125        let handles: Vec<_> = (0..10)
4126            .map(|i| {
4127                let cache_clone = cache.clone();
4128                tokio::spawn(async move {
4129                    cache_clone.insert(
4130                        format!("concurrent:entity-{i}"),
4131                        serde_json::json!({"thread": i}),
4132                    );
4133                })
4134            })
4135            .collect();
4136
4137        for handle in handles {
4138            handle.await.unwrap();
4139        }
4140
4141        // All 10 entries should be present
4142        assert_eq!(cache.len(), 10);
4143    }
4144
4145    #[tokio::test]
4146    async fn test_projection_state_large_payload() {
4147        let store = create_test_store();
4148        let cache = store.projection_state_cache();
4149
4150        // Create a large JSON payload (~10KB)
4151        let large_array: Vec<serde_json::Value> = (0..1000)
4152            .map(|i| serde_json::json!({"item": i, "description": "test item with some padding data to increase size"}))
4153            .collect();
4154
4155        cache.insert(
4156            "large:entity-1".to_string(),
4157            serde_json::json!({"items": large_array}),
4158        );
4159
4160        let state = cache.get("large:entity-1").unwrap();
4161        let items = state["items"].as_array().unwrap();
4162        assert_eq!(items.len(), 1000);
4163    }
4164
4165    #[tokio::test]
4166    async fn test_projection_state_complex_json() {
4167        let store = create_test_store();
4168        let cache = store.projection_state_cache();
4169
4170        // Complex nested JSON structure
4171        let complex_state = serde_json::json!({
4172            "user": {
4173                "id": "user-123",
4174                "profile": {
4175                    "name": "John Doe",
4176                    "email": "john@example.com",
4177                    "settings": {
4178                        "theme": "dark",
4179                        "notifications": true
4180                    }
4181                },
4182                "roles": ["admin", "user"],
4183                "metadata": {
4184                    "created_at": "2025-01-01T00:00:00Z",
4185                    "last_login": null
4186                }
4187            }
4188        });
4189
4190        cache.insert("complex:user-123".to_string(), complex_state);
4191
4192        let state = cache.get("complex:user-123").unwrap();
4193        assert_eq!(state["user"]["profile"]["name"], "John Doe");
4194        assert_eq!(state["user"]["roles"][0], "admin");
4195        assert!(state["user"]["metadata"]["last_login"].is_null());
4196    }
4197
4198    #[tokio::test]
4199    async fn test_projection_state_cache_iteration() {
4200        let store = create_test_store();
4201        let cache = store.projection_state_cache();
4202
4203        // Insert entries
4204        for i in 0..5 {
4205            cache.insert(format!("iter:entity-{i}"), serde_json::json!({"index": i}));
4206        }
4207
4208        // Iterate over all entries
4209        let entries: Vec<_> = cache.iter().map(|entry| entry.key().clone()).collect();
4210        assert_eq!(entries.len(), 5);
4211    }
4212
4213    #[tokio::test]
4214    async fn test_projection_manager_get_entity_snapshots() {
4215        let store = create_test_store();
4216        let projection_manager = store.projection_manager();
4217
4218        // Get entity_snapshots projection specifically
4219        let projection = projection_manager.get_projection("entity_snapshots");
4220        assert!(projection.is_some());
4221        assert_eq!(projection.unwrap().name(), "entity_snapshots");
4222    }
4223
4224    #[tokio::test]
4225    async fn test_projection_manager_get_event_counters() {
4226        let store = create_test_store();
4227        let projection_manager = store.projection_manager();
4228
4229        // Get event_counters projection specifically
4230        let projection = projection_manager.get_projection("event_counters");
4231        assert!(projection.is_some());
4232        assert_eq!(projection.unwrap().name(), "event_counters");
4233    }
4234
4235    #[tokio::test]
4236    async fn test_projection_state_cache_overwrite() {
4237        let store = create_test_store();
4238        let cache = store.projection_state_cache();
4239
4240        // Initial value
4241        cache.insert(
4242            "overwrite:entity-1".to_string(),
4243            serde_json::json!({"version": 1}),
4244        );
4245
4246        // Overwrite with new value
4247        cache.insert(
4248            "overwrite:entity-1".to_string(),
4249            serde_json::json!({"version": 2}),
4250        );
4251
4252        // Overwrite again
4253        cache.insert(
4254            "overwrite:entity-1".to_string(),
4255            serde_json::json!({"version": 3}),
4256        );
4257
4258        let state = cache.get("overwrite:entity-1").unwrap();
4259        assert_eq!(state["version"], 3);
4260
4261        // Should still be only 1 entry
4262        assert_eq!(cache.len(), 1);
4263    }
4264
4265    #[tokio::test]
4266    async fn test_projection_state_multiple_projections() {
4267        let store = create_test_store();
4268        let cache = store.projection_state_cache();
4269
4270        // Store states for different projections
4271        cache.insert(
4272            "entity_snapshots:user-1".to_string(),
4273            serde_json::json!({"name": "Alice"}),
4274        );
4275        cache.insert(
4276            "event_counters:user.created".to_string(),
4277            serde_json::json!({"count": 5}),
4278        );
4279        cache.insert(
4280            "custom_projection:order-1".to_string(),
4281            serde_json::json!({"total": 150.0}),
4282        );
4283
4284        // Verify each projection's state
4285        assert_eq!(
4286            cache.get("entity_snapshots:user-1").unwrap()["name"],
4287            "Alice"
4288        );
4289        assert_eq!(
4290            cache.get("event_counters:user.created").unwrap()["count"],
4291            5
4292        );
4293        assert_eq!(
4294            cache.get("custom_projection:order-1").unwrap()["total"],
4295            150.0
4296        );
4297    }
4298
4299    #[tokio::test]
4300    async fn test_bulk_projection_state_access() {
4301        let store = create_test_store();
4302
4303        // Ingest multiple events for different entities
4304        for i in 0..5 {
4305            let event = create_test_event(&format!("bulk-user-{i}"), "user.created");
4306            store.ingest(&event).unwrap();
4307        }
4308
4309        // Get projection and verify bulk access
4310        let projection_manager = store.projection_manager();
4311        let snapshot_projection = projection_manager
4312            .get_projection("entity_snapshots")
4313            .unwrap();
4314
4315        // Verify we can access all entities
4316        for i in 0..5 {
4317            let state = snapshot_projection.get_state(&format!("bulk-user-{i}"));
4318            assert!(state.is_some(), "Entity bulk-user-{i} should have state");
4319        }
4320    }
4321
4322    #[tokio::test]
4323    async fn test_bulk_save_projection_states() {
4324        let store = create_test_store();
4325        let cache = store.projection_state_cache();
4326
4327        // Simulate bulk save request
4328        let states = vec![
4329            BulkSaveStateItem {
4330                entity_id: "bulk-entity-1".to_string(),
4331                state: serde_json::json!({"name": "Entity 1", "value": 100}),
4332            },
4333            BulkSaveStateItem {
4334                entity_id: "bulk-entity-2".to_string(),
4335                state: serde_json::json!({"name": "Entity 2", "value": 200}),
4336            },
4337            BulkSaveStateItem {
4338                entity_id: "bulk-entity-3".to_string(),
4339                state: serde_json::json!({"name": "Entity 3", "value": 300}),
4340            },
4341        ];
4342
4343        let projection_name = "test_projection";
4344
4345        // Save states to cache (simulating bulk_save_projection_states handler)
4346        for item in &states {
4347            cache.insert(
4348                format!("{projection_name}:{}", item.entity_id),
4349                item.state.clone(),
4350            );
4351        }
4352
4353        // Verify all states were saved
4354        assert_eq!(cache.len(), 3);
4355
4356        let state1 = cache.get("test_projection:bulk-entity-1").unwrap();
4357        assert_eq!(state1["name"], "Entity 1");
4358        assert_eq!(state1["value"], 100);
4359
4360        let state2 = cache.get("test_projection:bulk-entity-2").unwrap();
4361        assert_eq!(state2["name"], "Entity 2");
4362        assert_eq!(state2["value"], 200);
4363
4364        let state3 = cache.get("test_projection:bulk-entity-3").unwrap();
4365        assert_eq!(state3["name"], "Entity 3");
4366        assert_eq!(state3["value"], 300);
4367    }
4368
4369    #[tokio::test]
4370    async fn test_bulk_save_empty_states() {
4371        let store = create_test_store();
4372        let cache = store.projection_state_cache();
4373
4374        // Clear cache
4375        cache.clear();
4376
4377        // Empty states should work fine
4378        let states: Vec<BulkSaveStateItem> = vec![];
4379        assert_eq!(states.len(), 0);
4380
4381        // Cache should remain empty
4382        assert_eq!(cache.len(), 0);
4383    }
4384
4385    #[tokio::test]
4386    async fn test_bulk_save_overwrites_existing() {
4387        let store = create_test_store();
4388        let cache = store.projection_state_cache();
4389
4390        // Insert initial state
4391        cache.insert(
4392            "test:entity-1".to_string(),
4393            serde_json::json!({"version": 1, "data": "initial"}),
4394        );
4395
4396        // Bulk save with updated state
4397        let new_state = serde_json::json!({"version": 2, "data": "updated"});
4398        cache.insert("test:entity-1".to_string(), new_state);
4399
4400        // Verify overwrite
4401        let state = cache.get("test:entity-1").unwrap();
4402        assert_eq!(state["version"], 2);
4403        assert_eq!(state["data"], "updated");
4404    }
4405
4406    #[tokio::test]
4407    async fn test_bulk_save_high_volume() {
4408        let store = create_test_store();
4409        let cache = store.projection_state_cache();
4410
4411        // Simulate high volume save (1000 entities)
4412        for i in 0..1000 {
4413            cache.insert(
4414                format!("volume_test:entity-{i}"),
4415                serde_json::json!({"index": i, "status": "active"}),
4416            );
4417        }
4418
4419        // Verify count
4420        assert_eq!(cache.len(), 1000);
4421
4422        // Spot check some entries
4423        assert_eq!(cache.get("volume_test:entity-0").unwrap()["index"], 0);
4424        assert_eq!(cache.get("volume_test:entity-500").unwrap()["index"], 500);
4425        assert_eq!(cache.get("volume_test:entity-999").unwrap()["index"], 999);
4426    }
4427
4428    #[tokio::test]
4429    async fn test_bulk_save_different_projections() {
4430        let store = create_test_store();
4431        let cache = store.projection_state_cache();
4432
4433        // Save to multiple projections in bulk
4434        let projections = ["entity_snapshots", "event_counters", "custom_analytics"];
4435
4436        for proj in &projections {
4437            for i in 0..5 {
4438                cache.insert(
4439                    format!("{proj}:entity-{i}"),
4440                    serde_json::json!({"projection": proj, "id": i}),
4441                );
4442            }
4443        }
4444
4445        // Verify total count (3 projections * 5 entities)
4446        assert_eq!(cache.len(), 15);
4447
4448        // Verify each projection
4449        for proj in &projections {
4450            let state = cache.get(&format!("{proj}:entity-0")).unwrap();
4451            assert_eq!(state["projection"], *proj);
4452        }
4453    }
4454
4455    // -------------------------------------------------------------------------
4456    // Cache-fallback tests for the projection-state read handlers (v0.19.1).
4457    //
4458    // SDK-managed projections write state via save_projection_state /
4459    // bulk_save_projection_states without registering in projection_manager.
4460    // These tests verify the read handlers fall back to the cache instead of
4461    // 404-ing. Registered projection wins where both exist; cache fills the
4462    // gap when not.
4463    // -------------------------------------------------------------------------
4464
4465    #[tokio::test]
4466    async fn get_projection_state_falls_back_to_cache_when_unregistered() {
4467        let store = create_test_store();
4468        store.projection_state_cache().insert(
4469            "assets:BTC".to_string(),
4470            serde_json::json!({"symbol": "BTC", "altname": "Bitcoin"}),
4471        );
4472
4473        let resp = get_projection_state(
4474            State(Arc::clone(&store)),
4475            Path(("assets".to_string(), "BTC".to_string())),
4476        )
4477        .await
4478        .expect("should not error when projection is not registered");
4479
4480        assert_eq!(resp.0["found"], serde_json::Value::Bool(true));
4481        assert_eq!(resp.0["state"]["symbol"], "BTC");
4482        assert_eq!(resp.0["state"]["altname"], "Bitcoin");
4483    }
4484
4485    #[tokio::test]
4486    async fn get_projection_state_returns_not_found_when_absent_everywhere() {
4487        let store = create_test_store();
4488
4489        let resp = get_projection_state(
4490            State(Arc::clone(&store)),
4491            Path(("assets".to_string(), "UNKNOWN".to_string())),
4492        )
4493        .await
4494        .unwrap();
4495
4496        assert_eq!(resp.0["found"], serde_json::Value::Bool(false));
4497        assert_eq!(resp.0["state"], serde_json::Value::Null);
4498    }
4499
4500    #[tokio::test]
4501    async fn get_projection_state_registered_wins_over_cache() {
4502        let store = create_test_store();
4503
4504        // Ingest an event so entity_snapshots (a registered projection) has state.
4505        let event = create_test_event("user-777", "user.created");
4506        store.ingest(&event).unwrap();
4507
4508        // Plant a conflicting cache entry for the same (projection, entity).
4509        store.projection_state_cache().insert(
4510            "entity_snapshots:user-777".to_string(),
4511            serde_json::json!({"stolen": "value"}),
4512        );
4513
4514        let resp = get_projection_state(
4515            State(Arc::clone(&store)),
4516            Path(("entity_snapshots".to_string(), "user-777".to_string())),
4517        )
4518        .await
4519        .unwrap();
4520
4521        // Registered projection wins — cache fallback is only consulted when
4522        // the registered projection has no state for this entity.
4523        assert_eq!(resp.0["found"], serde_json::Value::Bool(true));
4524        assert!(
4525            resp.0["state"].get("stolen").is_none(),
4526            "cache entry must not shadow registered projection state: got {:?}",
4527            resp.0["state"]
4528        );
4529    }
4530
4531    #[tokio::test]
4532    async fn get_projection_state_summary_returns_cache_without_registration() {
4533        let store = create_test_store();
4534        let cache = store.projection_state_cache();
4535        cache.insert("assets:BTC".into(), serde_json::json!({"symbol": "BTC"}));
4536        cache.insert("assets:ETH".into(), serde_json::json!({"symbol": "ETH"}));
4537        // Different projection name — must not appear in the summary.
4538        cache.insert("trades:t-1".into(), serde_json::json!({"x": 1}));
4539
4540        let resp = get_projection_state_summary(
4541            State(Arc::clone(&store)),
4542            Path("assets".to_string()),
4543            Query(ProjectionStateSummaryParams::default()),
4544        )
4545        .await
4546        .unwrap();
4547
4548        assert_eq!(resp.0["total"], 2);
4549        let states = resp.0["states"].as_array().unwrap();
4550        let entity_ids: Vec<&str> = states
4551            .iter()
4552            .map(|s| s["entity_id"].as_str().unwrap())
4553            .collect();
4554        assert!(entity_ids.contains(&"BTC"));
4555        assert!(entity_ids.contains(&"ETH"));
4556    }
4557
4558    #[tokio::test]
4559    async fn bulk_get_projection_states_falls_back_to_cache() {
4560        let store = create_test_store();
4561        let cache = store.projection_state_cache();
4562        cache.insert("assets:BTC".into(), serde_json::json!({"symbol": "BTC"}));
4563        cache.insert("assets:ETH".into(), serde_json::json!({"symbol": "ETH"}));
4564
4565        let req = BulkGetStateRequest {
4566            entity_ids: vec!["BTC".into(), "ETH".into(), "MISSING".into()],
4567        };
4568
4569        let resp = bulk_get_projection_states(
4570            State(Arc::clone(&store)),
4571            Path("assets".to_string()),
4572            Json(req),
4573        )
4574        .await
4575        .unwrap();
4576
4577        assert_eq!(resp.0["total"], 3);
4578        let states = resp.0["states"].as_array().unwrap();
4579        let by_id: std::collections::HashMap<&str, &serde_json::Value> = states
4580            .iter()
4581            .map(|s| (s["entity_id"].as_str().unwrap(), s))
4582            .collect();
4583
4584        assert_eq!(by_id["BTC"]["found"], serde_json::Value::Bool(true));
4585        assert_eq!(by_id["BTC"]["state"]["symbol"], "BTC");
4586        assert_eq!(by_id["ETH"]["found"], serde_json::Value::Bool(true));
4587        assert_eq!(by_id["MISSING"]["found"], serde_json::Value::Bool(false));
4588    }
4589
4590    /// Pins the poll wire shape the Rust SDK's `poll_consumer_events` decodes:
4591    /// `ConsumerEventDto` flattens the event, so its fields sit next to
4592    /// `position` rather than nested under an `event` key.
4593    #[tokio::test]
4594    async fn poll_consumer_events_flattens_event_alongside_position() {
4595        let store = create_test_store();
4596        store
4597            .ingest(&create_test_event("user-1", "user.created"))
4598            .unwrap();
4599        store
4600            .ingest(&create_test_event("user-2", "user.updated"))
4601            .unwrap();
4602        store.consumer_registry().register("w1", &[]);
4603
4604        let resp = poll_consumer_events(
4605            State(Arc::clone(&store)),
4606            Path("w1".to_string()),
4607            Query(ConsumerPollQuery { limit: Some(10) }),
4608        )
4609        .await
4610        .unwrap();
4611
4612        let body = serde_json::to_value(&resp.0).unwrap();
4613        assert_eq!(body["count"], 2);
4614        let first = &body["events"][0];
4615        assert_eq!(first["position"], 1);
4616        assert!(
4617            first.get("event").is_none(),
4618            "event must be flattened, not nested: got {first:?}"
4619        );
4620        assert_eq!(first["event_type"], "user.created");
4621        assert_eq!(first["entity_id"], "user-1");
4622    }
4623}