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