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(|a, b| b.count.cmp(&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(|a, b| b.last_event_at.cmp(&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(|a, b| b.event_count.cmp(&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;
3114
3115        let store = create_test_store();
3116        let base = chrono::Utc::now() - chrono::Duration::hours(24);
3117        let mut ids = Vec::new();
3118        for i in 0..5i64 {
3119            let mut event = create_test_event(&format!("e-{i}"), "user.created");
3120            event.timestamp = base + chrono::Duration::hours(i);
3121            event.version = i + 1;
3122            ids.push(event.id);
3123            store.ingest(&event).unwrap();
3124        }
3125        // `Z`-suffixed so the timestamp survives a query string — an offset of
3126        // `+00:00` would be decoded as a space.
3127        let at = |h: i64| {
3128            (base + chrono::Duration::hours(h)).to_rfc3339_opts(SecondsFormat::Micros, true)
3129        };
3130
3131        for (qs, expected) in [
3132            (format!("since={}", at(2)), vec![ids[2], ids[3], ids[4]]),
3133            (format!("until={}", at(1)), vec![ids[0], ids[1]]),
3134            (format!("as_of={}", at(1)), vec![ids[0], ids[1]]),
3135            (
3136                format!("since={}&until={}", at(1), at(3)),
3137                vec![ids[1], ids[2], ids[3]],
3138            ),
3139        ] {
3140            let resp = query_page(&store, &qs).await;
3141            let got: Vec<_> = resp.events.iter().map(|e| e.id).collect();
3142            assert_eq!(got, expected, "?{qs} must return only the window");
3143            assert_eq!(resp.count, expected.len(), "?{qs}");
3144            assert_eq!(
3145                resp.total_count,
3146                expected.len(),
3147                "?{qs}: total_count must count the window, not the history"
3148            );
3149            assert!(!resp.has_more, "?{qs}: the whole window was served");
3150        }
3151
3152        // The window composes with paging: page 2 of a `since` window is the
3153        // second page OF THAT WINDOW, and `has_more` terminates on it.
3154        let page1 = query_page(&store, &format!("since={}&limit=2", at(2))).await;
3155        let page2 = query_page(&store, &format!("since={}&limit=2&offset=2", at(2))).await;
3156        assert_eq!(
3157            page1.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3158            vec![ids[2], ids[3]]
3159        );
3160        assert!(page1.has_more, "3 in the window, 2 served");
3161        assert_eq!(
3162            page2.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3163            vec![ids[4]]
3164        );
3165        assert!(!page2.has_more, "offset 2 + count 1 == the window's 3");
3166        assert_eq!(page2.total_count, 3);
3167
3168        // …and with `order=desc`.
3169        let desc = query_page(&store, &format!("since={}&order=desc", at(2))).await;
3170        assert_eq!(
3171            desc.events.iter().map(|e| e.id).collect::<Vec<_>>(),
3172            vec![ids[4], ids[3], ids[2]]
3173        );
3174
3175        // An empty window is empty, not "everything".
3176        let empty = query_page(&store, &format!("since={}", at(99))).await;
3177        assert_eq!(empty.count, 0);
3178        assert_eq!(empty.total_count, 0);
3179        assert!(!empty.has_more);
3180    }
3181
3182    // Tenant-isolation gate: the public events query must fail CLOSED — a request
3183    // with no auth context and no tenant_id returns nothing, never a cross-tenant
3184    // scan. Calls the real handler so the boundary check is exercised.
3185    #[tokio::test]
3186    async fn query_events_fails_closed_without_tenant() {
3187        use axum::extract::{Query, State};
3188
3189        let store = create_test_store();
3190        for i in 0..5 {
3191            store
3192                .ingest(&create_test_event(&format!("e-{i}"), "user.created"))
3193                .unwrap();
3194        }
3195
3196        let resp = query_events(
3197            OptionalAuth(None),
3198            Query(QueryEventsRequest::default()),
3199            Query(EventOrderParam { order: None }),
3200            Query(EventOffsetParam { offset: None }),
3201            State(store.clone()),
3202        )
3203        .await
3204        .unwrap();
3205        assert_eq!(
3206            resp.0.total_count, 0,
3207            "a no-tenant query must NOT return cross-tenant events"
3208        );
3209        assert_eq!(resp.0.count, 0);
3210
3211        // The same query scoped to the events' tenant returns them.
3212        // (create_test_event stamps tenant "test-stream" — the 3rd from_strings arg.)
3213        let scoped = query_events(
3214            OptionalAuth(None),
3215            Query(QueryEventsRequest {
3216                tenant_id: Some("test-stream".to_string()),
3217                ..QueryEventsRequest::default()
3218            }),
3219            Query(EventOrderParam { order: None }),
3220            Query(EventOffsetParam { offset: None }),
3221            State(store),
3222        )
3223        .await
3224        .unwrap();
3225        assert_eq!(
3226            scoped.0.total_count, 5,
3227            "tenant-scoped query returns its events"
3228        );
3229    }
3230
3231    // The dashboard's streams + event-types counts must be per-tenant. These
3232    // endpoints used to scan ALL tenants (platform totals shown as "yours", and a
3233    // cross-tenant spill). Assert each tenant sees only its own, and no-tenant
3234    // fails closed.
3235    #[tokio::test]
3236    async fn list_streams_and_types_are_tenant_scoped() {
3237        use crate::domain::entities::Event;
3238        use axum::extract::{Query, State};
3239
3240        let store = create_test_store();
3241        let ev = |entity: &str, etype: &str, tenant: &str| {
3242            Event::from_strings(
3243                etype.to_string(),
3244                entity.to_string(),
3245                tenant.to_string(),
3246                serde_json::json!({}),
3247                None,
3248            )
3249            .unwrap()
3250        };
3251        // tenant A: 2 entities, 2 types. tenant B: 1 entity, 1 type.
3252        store.ingest(&ev("e1", "order.placed", "tenant-a")).unwrap();
3253        store.ingest(&ev("e2", "user.created", "tenant-a")).unwrap();
3254        store
3255            .ingest(&ev("e9", "thing.happened", "tenant-b"))
3256            .unwrap();
3257
3258        let streams = |tid: Option<&str>| {
3259            list_streams(
3260                OptionalAuth(None),
3261                State(store.clone()),
3262                Query(ListStreamsParams {
3263                    tenant_id: tid.map(String::from),
3264                    limit: None,
3265                    offset: None,
3266                }),
3267            )
3268        };
3269        assert_eq!(
3270            streams(Some("tenant-a")).await.0.total,
3271            2,
3272            "tenant-a streams"
3273        );
3274        assert_eq!(
3275            streams(Some("tenant-b")).await.0.total,
3276            1,
3277            "tenant-b streams"
3278        );
3279        assert_eq!(
3280            streams(None).await.0.total,
3281            0,
3282            "no tenant -> no streams (fail closed)"
3283        );
3284
3285        let types = |tid: Option<&str>| {
3286            list_event_types(
3287                OptionalAuth(None),
3288                State(store.clone()),
3289                Query(ListEventTypesParams {
3290                    tenant_id: tid.map(String::from),
3291                    limit: None,
3292                    offset: None,
3293                }),
3294            )
3295        };
3296        assert_eq!(
3297            types(Some("tenant-a")).await.0.total,
3298            2,
3299            "tenant-a event types"
3300        );
3301        assert_eq!(
3302            types(Some("tenant-b")).await.0.total,
3303            1,
3304            "tenant-b event types"
3305        );
3306        assert_eq!(
3307            types(None).await.0.total,
3308            0,
3309            "no tenant -> no types (fail closed)"
3310        );
3311    }
3312
3313    #[tokio::test]
3314    async fn test_query_events_no_more_results() {
3315        let store = create_test_store();
3316
3317        // Ingest 5 events
3318        for i in 0..5 {
3319            store
3320                .ingest(&create_test_event(&format!("entity-{i}"), "user.created"))
3321                .unwrap();
3322        }
3323
3324        // Query with limit=100 — should get has_more=false, total_count=5
3325        let all_events = store
3326            .query(&QueryEventsRequest {
3327                entity_id: None,
3328                event_type: None,
3329                tenant_id: None,
3330                as_of: None,
3331                since: None,
3332                until: None,
3333                limit: None,
3334                event_type_prefix: None,
3335                exclude_event_type_prefix: None,
3336                payload_filter: None,
3337            })
3338            .unwrap();
3339        let total_count = all_events.len();
3340        let limited_events: Vec<Event> = all_events.into_iter().take(100).collect();
3341        let count = limited_events.len();
3342        let has_more = count < total_count;
3343
3344        assert_eq!(count, 5);
3345        assert_eq!(total_count, 5);
3346        assert!(!has_more);
3347    }
3348
3349    // Regression test for issue #177: `order=desc` + `limit=1` must return
3350    // the NEWEST event for an entity, not the oldest. Mirrors the ordering
3351    // logic in `query_events` (store.query → reverse-if-desc → take(limit)).
3352    #[tokio::test]
3353    async fn test_query_events_order_desc_returns_latest() {
3354        // Drives the REAL handler. An earlier version of this test rebuilt the
3355        // ordering inline (`ascending.clone(); reverse(); take(1)`) and never
3356        // called `query_events`, so it could not fail on a mis-composition in
3357        // the code that actually serves `order=desc` — which since issue #251
3358        // lives in `EventStore::query_window`, not in the handler.
3359        let store = create_test_store();
3360
3361        // Five events for the same entity with strictly increasing timestamps —
3362        // mimics a backfill appending corrected events.
3363        let base = chrono::Utc::now();
3364        let mut ascending_ids = Vec::new();
3365        for i in 0..5i64 {
3366            let mut event = create_test_event("org-1", "auth.org.updated");
3367            event.timestamp = base + chrono::Duration::seconds(i);
3368            event.version = i + 1;
3369            ascending_ids.push(event.id);
3370            store.ingest(&event).unwrap();
3371        }
3372        let newest_ts = base + chrono::Duration::seconds(4);
3373
3374        // The documented "latest event for an entity" read.
3375        let latest = query_page(&store, "entity_id=org-1&limit=1&order=desc").await;
3376        assert_eq!(latest.count, 1);
3377        assert_eq!(
3378            latest.events[0].id,
3379            ascending_ids[4],
3380            "order=desc&limit=1 must yield the NEWEST event, got the one at \
3381             ascending position {:?}",
3382            ascending_ids
3383                .iter()
3384                .position(|id| *id == latest.events[0].id)
3385        );
3386        assert_eq!(latest.events[0].timestamp, newest_ts);
3387        assert_eq!(latest.total_count, 5, "total is the full match set");
3388        assert!(latest.has_more);
3389
3390        // Default order (and an explicit `asc`) still yields the OLDEST.
3391        for qs in [
3392            "entity_id=org-1&limit=1",
3393            "entity_id=org-1&limit=1&order=asc",
3394        ] {
3395            let oldest = query_page(&store, qs).await;
3396            assert_eq!(oldest.events[0].id, ascending_ids[0], "{qs}");
3397            assert_eq!(oldest.events[0].timestamp, base);
3398        }
3399
3400        // An unbounded desc page is the exact reverse of the ascending one —
3401        // reversal must apply to the whole match set, not just to the page.
3402        let all_desc = query_page(&store, "entity_id=org-1&order=desc").await;
3403        let got: Vec<_> = all_desc.events.iter().map(|e| e.id).collect();
3404        let expected: Vec<_> = ascending_ids.iter().rev().copied().collect();
3405        assert_eq!(got, expected, "order=desc must return newest-first");
3406    }
3407
3408    // Regression guard for issue #251's relocation of the ordering: `order=desc`
3409    // moved out of the handler and into `EventStore::query_window`, where it now
3410    // composes with `offset` and `limit`. The contract is reverse-THEN-skip-THEN-
3411    // take: `order=desc&offset=1&limit=2` is "the 2nd and 3rd newest". Skipping
3412    // before reversing (or reversing only the page) returns a different, quietly
3413    // wrong page — with the same count, total_count and has_more.
3414    #[tokio::test]
3415    async fn query_events_desc_composes_with_offset_and_limit() {
3416        let store = create_test_store();
3417        let base = chrono::Utc::now();
3418        let mut ascending_ids = Vec::new();
3419        for i in 0..5i64 {
3420            let mut event = create_test_event("org-1", "auth.org.updated");
3421            event.timestamp = base + chrono::Duration::seconds(i);
3422            event.version = i + 1;
3423            ascending_ids.push(event.id);
3424            store.ingest(&event).unwrap();
3425        }
3426        let newest_first: Vec<_> = ascending_ids.iter().rev().copied().collect();
3427
3428        for (offset, limit) in [(0, 2), (1, 2), (2, 2), (3, 2), (4, 2), (5, 2), (1, 4)] {
3429            let page = query_page(
3430                &store,
3431                &format!("entity_id=org-1&order=desc&offset={offset}&limit={limit}"),
3432            )
3433            .await;
3434            let got: Vec<_> = page.events.iter().map(|e| e.id).collect();
3435            let expected: Vec<_> = newest_first
3436                .iter()
3437                .skip(offset)
3438                .take(limit)
3439                .copied()
3440                .collect();
3441            assert_eq!(
3442                got, expected,
3443                "order=desc&offset={offset}&limit={limit} must reverse, then \
3444                 skip, then take"
3445            );
3446            assert_eq!(page.count, expected.len());
3447            assert_eq!(page.total_count, 5);
3448            assert_eq!(
3449                page.has_more,
3450                offset + expected.len() < 5,
3451                "has_more must account for the offset (offset={offset})"
3452            );
3453        }
3454
3455        // Walking the whole entity newest-first must visit every event exactly
3456        // once — the property a `order=desc` paginator depends on.
3457        let mut walked = Vec::new();
3458        for offset in (0..5).step_by(2) {
3459            let page = query_page(
3460                &store,
3461                &format!("entity_id=org-1&order=desc&offset={offset}&limit=2"),
3462            )
3463            .await;
3464            walked.extend(page.events.iter().map(|e| e.id));
3465        }
3466        assert_eq!(walked, newest_first, "desc paging must cover the set once");
3467    }
3468
3469    #[tokio::test]
3470    async fn test_list_entities_by_type_prefix() {
3471        let store = create_test_store();
3472
3473        // 3 index entities
3474        store
3475            .ingest(&create_test_event("idx-1", "index.created"))
3476            .unwrap();
3477        store
3478            .ingest(&create_test_event("idx-1", "index.updated"))
3479            .unwrap();
3480        store
3481            .ingest(&create_test_event("idx-2", "index.created"))
3482            .unwrap();
3483        store
3484            .ingest(&create_test_event("idx-3", "index.created"))
3485            .unwrap();
3486        // 2 trade entities
3487        store
3488            .ingest(&create_test_event("trade-1", "trade.created"))
3489            .unwrap();
3490        store
3491            .ingest(&create_test_event("trade-2", "trade.created"))
3492            .unwrap();
3493
3494        // List entities for index.*
3495        let req = ListEntitiesRequest {
3496            event_type_prefix: Some("index.".to_string()),
3497            ..Default::default()
3498        };
3499        let query_req = QueryEventsRequest {
3500            entity_id: None,
3501            event_type: None,
3502            tenant_id: None,
3503            as_of: None,
3504            since: None,
3505            until: None,
3506            limit: None,
3507            event_type_prefix: req.event_type_prefix,
3508            exclude_event_type_prefix: None,
3509            payload_filter: req.payload_filter,
3510        };
3511        let events = store.query(&query_req).unwrap();
3512
3513        // Group and verify
3514        let mut entity_map: std::collections::HashMap<String, Vec<&Event>> =
3515            std::collections::HashMap::new();
3516        for event in &events {
3517            entity_map
3518                .entry(event.entity_id().to_string())
3519                .or_default()
3520                .push(event);
3521        }
3522
3523        assert_eq!(entity_map.len(), 3); // idx-1, idx-2, idx-3
3524        assert_eq!(entity_map["idx-1"].len(), 2); // 2 events for idx-1
3525        assert_eq!(entity_map["idx-2"].len(), 1);
3526        assert_eq!(entity_map["idx-3"].len(), 1);
3527    }
3528
3529    // Issue #178: `list_entities` accepts an `order` param and pages
3530    // deterministically over the resulting sort.
3531    #[tokio::test]
3532    async fn test_list_entities_order_and_pagination() {
3533        let store = create_test_store();
3534
3535        // Three entities with strictly increasing last-event times.
3536        let base = chrono::Utc::now();
3537        for (i, eid) in ["org-a", "org-b", "org-c"].iter().enumerate() {
3538            let mut event = create_test_event(eid, "auth.org.created");
3539            event.timestamp = base + chrono::Duration::seconds(i as i64);
3540            store.ingest(&event).unwrap();
3541        }
3542        let prefix = || Some("auth.org.".to_string());
3543
3544        // Default: newest activity first (desc) — preserves prior behavior.
3545        let desc = list_entities(
3546            State(store.clone()),
3547            Query(ListEntitiesRequest {
3548                event_type_prefix: prefix(),
3549                ..Default::default()
3550            }),
3551        )
3552        .await
3553        .unwrap();
3554        let desc_ids: Vec<&str> = desc
3555            .0
3556            .entities
3557            .iter()
3558            .map(|e| e.entity_id.as_str())
3559            .collect();
3560        assert_eq!(desc_ids, ["org-c", "org-b", "org-a"]);
3561
3562        // order=asc: oldest activity first.
3563        let asc = list_entities(
3564            State(store.clone()),
3565            Query(ListEntitiesRequest {
3566                event_type_prefix: prefix(),
3567                order: Some("asc".to_string()),
3568                ..Default::default()
3569            }),
3570        )
3571        .await
3572        .unwrap();
3573        let asc_ids: Vec<&str> = asc
3574            .0
3575            .entities
3576            .iter()
3577            .map(|e| e.entity_id.as_str())
3578            .collect();
3579        assert_eq!(asc_ids, ["org-a", "org-b", "org-c"]);
3580
3581        // Offset pagination over the deterministic asc order: page 2, size 1.
3582        let page2 = list_entities(
3583            State(store.clone()),
3584            Query(ListEntitiesRequest {
3585                event_type_prefix: prefix(),
3586                order: Some("asc".to_string()),
3587                limit: Some(1),
3588                offset: Some(1),
3589                ..Default::default()
3590            }),
3591        )
3592        .await
3593        .unwrap();
3594        assert_eq!(page2.0.entities.len(), 1);
3595        assert_eq!(page2.0.entities[0].entity_id, "org-b");
3596        assert_eq!(page2.0.total, 3);
3597        assert!(page2.0.has_more);
3598
3599        // Invalid order value is rejected.
3600        let err = list_entities(
3601            State(store.clone()),
3602            Query(ListEntitiesRequest {
3603                event_type_prefix: prefix(),
3604                order: Some("sideways".to_string()),
3605                ..Default::default()
3606            }),
3607        )
3608        .await;
3609        assert!(err.is_err(), "invalid order value must be rejected");
3610    }
3611
3612    fn create_test_event_with_payload(
3613        entity_id: &str,
3614        event_type: &str,
3615        payload: serde_json::Value,
3616    ) -> Event {
3617        Event::from_strings(
3618            event_type.to_string(),
3619            entity_id.to_string(),
3620            "test-stream".to_string(),
3621            payload,
3622            None,
3623        )
3624        .unwrap()
3625    }
3626
3627    #[tokio::test]
3628    async fn test_detect_duplicates_by_payload_fields() {
3629        let store = create_test_store();
3630
3631        // Create entities with duplicate "name" field values
3632        store
3633            .ingest(&create_test_event_with_payload(
3634                "idx-1",
3635                "index.created",
3636                serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3637            ))
3638            .unwrap();
3639        store
3640            .ingest(&create_test_event_with_payload(
3641                "idx-2",
3642                "index.created",
3643                serde_json::json!({"name": "S&P 500", "user_id": "bob"}),
3644            ))
3645            .unwrap();
3646        store
3647            .ingest(&create_test_event_with_payload(
3648                "idx-3",
3649                "index.created",
3650                serde_json::json!({"name": "NASDAQ", "user_id": "alice"}),
3651            ))
3652            .unwrap();
3653        store
3654            .ingest(&create_test_event_with_payload(
3655                "idx-4",
3656                "index.created",
3657                serde_json::json!({"name": "NASDAQ", "user_id": "carol"}),
3658            ))
3659            .unwrap();
3660        store
3661            .ingest(&create_test_event_with_payload(
3662                "idx-5",
3663                "index.created",
3664                serde_json::json!({"name": "DAX", "user_id": "dave"}),
3665            ))
3666            .unwrap();
3667
3668        // Group by name — should find 2 groups: "S&P 500" (idx-1, idx-2) and "NASDAQ" (idx-3, idx-4)
3669        let query_req = QueryEventsRequest {
3670            entity_id: None,
3671            event_type: None,
3672            tenant_id: None,
3673            as_of: None,
3674            since: None,
3675            until: None,
3676            limit: None,
3677            event_type_prefix: Some("index.".to_string()),
3678            exclude_event_type_prefix: None,
3679            payload_filter: None,
3680        };
3681        let events = store.query(&query_req).unwrap();
3682
3683        // Manually replicate the handler logic for testing
3684        let group_by_fields = vec!["name"];
3685        let mut entity_latest: std::collections::HashMap<String, &Event> =
3686            std::collections::HashMap::new();
3687        for event in &events {
3688            let eid = event.entity_id().to_string();
3689            entity_latest
3690                .entry(eid)
3691                .and_modify(|existing| {
3692                    if event.timestamp() > existing.timestamp() {
3693                        *existing = event;
3694                    }
3695                })
3696                .or_insert(event);
3697        }
3698
3699        let mut groups: std::collections::HashMap<String, Vec<String>> =
3700            std::collections::HashMap::new();
3701        for (entity_id, event) in &entity_latest {
3702            let payload = event.payload();
3703            let mut key_parts = serde_json::Map::new();
3704            for field in &group_by_fields {
3705                let value = payload
3706                    .get(*field)
3707                    .cloned()
3708                    .unwrap_or(serde_json::Value::Null);
3709                key_parts.insert((*field).to_string(), value);
3710            }
3711            let key_str = serde_json::to_string(&key_parts).unwrap_or_default();
3712            groups.entry(key_str).or_default().push(entity_id.clone());
3713        }
3714
3715        let duplicate_groups: Vec<_> = groups
3716            .into_iter()
3717            .filter(|(_, ids)| ids.len() > 1)
3718            .collect();
3719
3720        assert_eq!(duplicate_groups.len(), 2); // S&P 500 and NASDAQ groups
3721        for (_, ids) in &duplicate_groups {
3722            assert_eq!(ids.len(), 2);
3723        }
3724    }
3725
3726    #[tokio::test]
3727    async fn test_detect_duplicates_no_duplicates() {
3728        let store = create_test_store();
3729
3730        // All unique names
3731        store
3732            .ingest(&create_test_event_with_payload(
3733                "idx-1",
3734                "index.created",
3735                serde_json::json!({"name": "A"}),
3736            ))
3737            .unwrap();
3738        store
3739            .ingest(&create_test_event_with_payload(
3740                "idx-2",
3741                "index.created",
3742                serde_json::json!({"name": "B"}),
3743            ))
3744            .unwrap();
3745
3746        let query_req = QueryEventsRequest {
3747            entity_id: None,
3748            event_type: None,
3749            tenant_id: None,
3750            as_of: None,
3751            since: None,
3752            until: None,
3753            limit: None,
3754            event_type_prefix: Some("index.".to_string()),
3755            exclude_event_type_prefix: None,
3756            payload_filter: None,
3757        };
3758        let events = store.query(&query_req).unwrap();
3759
3760        let mut entity_latest: std::collections::HashMap<String, &Event> =
3761            std::collections::HashMap::new();
3762        for event in &events {
3763            entity_latest
3764                .entry(event.entity_id().to_string())
3765                .or_insert(event);
3766        }
3767
3768        let mut groups: std::collections::HashMap<String, Vec<String>> =
3769            std::collections::HashMap::new();
3770        for (entity_id, event) in &entity_latest {
3771            let key_str =
3772                serde_json::to_string(&serde_json::json!({"name": event.payload().get("name")}))
3773                    .unwrap();
3774            groups.entry(key_str).or_default().push(entity_id.clone());
3775        }
3776
3777        let duplicate_groups: Vec<_> = groups
3778            .into_iter()
3779            .filter(|(_, ids)| ids.len() > 1)
3780            .collect();
3781
3782        assert_eq!(duplicate_groups.len(), 0); // No duplicates
3783    }
3784
3785    #[tokio::test]
3786    async fn test_detect_duplicates_multi_field_group_by() {
3787        let store = create_test_store();
3788
3789        // Two entities with same name AND user_id = true duplicate
3790        store
3791            .ingest(&create_test_event_with_payload(
3792                "idx-1",
3793                "index.created",
3794                serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3795            ))
3796            .unwrap();
3797        store
3798            .ingest(&create_test_event_with_payload(
3799                "idx-2",
3800                "index.created",
3801                serde_json::json!({"name": "S&P 500", "user_id": "alice"}),
3802            ))
3803            .unwrap();
3804        // Same name but different user_id = NOT a duplicate in multi-field group
3805        store
3806            .ingest(&create_test_event_with_payload(
3807                "idx-3",
3808                "index.created",
3809                serde_json::json!({"name": "S&P 500", "user_id": "bob"}),
3810            ))
3811            .unwrap();
3812
3813        let query_req = QueryEventsRequest {
3814            entity_id: None,
3815            event_type: None,
3816            tenant_id: None,
3817            as_of: None,
3818            since: None,
3819            until: None,
3820            limit: None,
3821            event_type_prefix: Some("index.".to_string()),
3822            exclude_event_type_prefix: None,
3823            payload_filter: None,
3824        };
3825        let events = store.query(&query_req).unwrap();
3826
3827        let group_by_fields = vec!["name", "user_id"];
3828        let mut entity_latest: std::collections::HashMap<String, &Event> =
3829            std::collections::HashMap::new();
3830        for event in &events {
3831            entity_latest
3832                .entry(event.entity_id().to_string())
3833                .and_modify(|existing| {
3834                    if event.timestamp() > existing.timestamp() {
3835                        *existing = event;
3836                    }
3837                })
3838                .or_insert(event);
3839        }
3840
3841        let mut groups: std::collections::HashMap<String, Vec<String>> =
3842            std::collections::HashMap::new();
3843        for (entity_id, event) in &entity_latest {
3844            let payload = event.payload();
3845            let mut key_parts = serde_json::Map::new();
3846            for field in &group_by_fields {
3847                let value = payload
3848                    .get(*field)
3849                    .cloned()
3850                    .unwrap_or(serde_json::Value::Null);
3851                key_parts.insert((*field).to_string(), value);
3852            }
3853            let key_str = serde_json::to_string(&key_parts).unwrap_or_default();
3854            groups.entry(key_str).or_default().push(entity_id.clone());
3855        }
3856
3857        let duplicate_groups: Vec<_> = groups
3858            .into_iter()
3859            .filter(|(_, ids)| ids.len() > 1)
3860            .collect();
3861
3862        // Only 1 duplicate group: name=S&P 500, user_id=alice (idx-1, idx-2)
3863        assert_eq!(duplicate_groups.len(), 1);
3864        let (_, ref ids) = duplicate_groups[0];
3865        assert_eq!(ids.len(), 2);
3866        let mut sorted_ids = ids.clone();
3867        sorted_ids.sort();
3868        assert_eq!(sorted_ids, vec!["idx-1", "idx-2"]);
3869    }
3870
3871    #[tokio::test]
3872    async fn test_projection_state_cache() {
3873        let store = create_test_store();
3874
3875        // Test cache insertion
3876        let cache = store.projection_state_cache();
3877        cache.insert(
3878            "entity_snapshots:user-123".to_string(),
3879            serde_json::json!({"name": "Test User", "age": 30}),
3880        );
3881
3882        // Test cache retrieval
3883        let state = cache.get("entity_snapshots:user-123");
3884        assert!(state.is_some());
3885        let state = state.unwrap();
3886        assert_eq!(state["name"], "Test User");
3887        assert_eq!(state["age"], 30);
3888    }
3889
3890    #[tokio::test]
3891    async fn test_projection_manager_list_projections() {
3892        let store = create_test_store();
3893
3894        // List projections (built-in projections should be available)
3895        let projection_manager = store.projection_manager();
3896        let projections = projection_manager.list_projections();
3897
3898        // Should have entity_snapshots and event_counters
3899        assert!(projections.len() >= 2);
3900
3901        let names: Vec<&str> = projections.iter().map(|(name, _)| name.as_str()).collect();
3902        assert!(names.contains(&"entity_snapshots"));
3903        assert!(names.contains(&"event_counters"));
3904    }
3905
3906    #[tokio::test]
3907    async fn test_projection_state_after_event_ingestion() {
3908        let store = create_test_store();
3909
3910        // Ingest an event
3911        let event = create_test_event("user-456", "user.created");
3912        store.ingest(&event).unwrap();
3913
3914        // Get projection state
3915        let projection_manager = store.projection_manager();
3916        let snapshot_projection = projection_manager
3917            .get_projection("entity_snapshots")
3918            .unwrap();
3919
3920        let state = snapshot_projection.get_state("user-456");
3921        assert!(state.is_some());
3922        let state = state.unwrap();
3923        assert_eq!(state["name"], "Test");
3924        assert_eq!(state["value"], 42);
3925    }
3926
3927    #[tokio::test]
3928    async fn test_projection_state_cache_multiple_entities() {
3929        let store = create_test_store();
3930        let cache = store.projection_state_cache();
3931
3932        // Insert multiple entities
3933        for i in 0..10 {
3934            cache.insert(
3935                format!("entity_snapshots:entity-{i}"),
3936                serde_json::json!({"id": i, "status": "active"}),
3937            );
3938        }
3939
3940        // Verify all insertions
3941        assert_eq!(cache.len(), 10);
3942
3943        // Verify each entity
3944        for i in 0..10 {
3945            let key = format!("entity_snapshots:entity-{i}");
3946            let state = cache.get(&key);
3947            assert!(state.is_some());
3948            assert_eq!(state.unwrap()["id"], i);
3949        }
3950    }
3951
3952    #[tokio::test]
3953    async fn test_projection_state_update() {
3954        let store = create_test_store();
3955        let cache = store.projection_state_cache();
3956
3957        // Initial state
3958        cache.insert(
3959            "entity_snapshots:user-789".to_string(),
3960            serde_json::json!({"balance": 100}),
3961        );
3962
3963        // Update state
3964        cache.insert(
3965            "entity_snapshots:user-789".to_string(),
3966            serde_json::json!({"balance": 150}),
3967        );
3968
3969        // Verify update
3970        let state = cache.get("entity_snapshots:user-789").unwrap();
3971        assert_eq!(state["balance"], 150);
3972    }
3973
3974    #[tokio::test]
3975    async fn test_event_counter_projection() {
3976        let store = create_test_store();
3977
3978        // Ingest events of different types
3979        store
3980            .ingest(&create_test_event("user-1", "user.created"))
3981            .unwrap();
3982        store
3983            .ingest(&create_test_event("user-2", "user.created"))
3984            .unwrap();
3985        store
3986            .ingest(&create_test_event("user-1", "user.updated"))
3987            .unwrap();
3988
3989        // Get event counter projection
3990        let projection_manager = store.projection_manager();
3991        let counter_projection = projection_manager.get_projection("event_counters").unwrap();
3992
3993        // Check counts
3994        let created_state = counter_projection.get_state("user.created");
3995        assert!(created_state.is_some());
3996        assert_eq!(created_state.unwrap()["count"], 2);
3997
3998        let updated_state = counter_projection.get_state("user.updated");
3999        assert!(updated_state.is_some());
4000        assert_eq!(updated_state.unwrap()["count"], 1);
4001    }
4002
4003    #[tokio::test]
4004    async fn test_projection_state_cache_key_format() {
4005        let store = create_test_store();
4006        let cache = store.projection_state_cache();
4007
4008        // Test standard key format: {projection_name}:{entity_id}
4009        let key = "orders:order-12345".to_string();
4010        cache.insert(key.clone(), serde_json::json!({"total": 99.99}));
4011
4012        let state = cache.get(&key).unwrap();
4013        assert_eq!(state["total"], 99.99);
4014    }
4015
4016    #[tokio::test]
4017    async fn test_projection_state_cache_removal() {
4018        let store = create_test_store();
4019        let cache = store.projection_state_cache();
4020
4021        // Insert and then remove
4022        cache.insert(
4023            "test:entity-1".to_string(),
4024            serde_json::json!({"data": "value"}),
4025        );
4026        assert_eq!(cache.len(), 1);
4027
4028        cache.remove("test:entity-1");
4029        assert_eq!(cache.len(), 0);
4030        assert!(cache.get("test:entity-1").is_none());
4031    }
4032
4033    #[tokio::test]
4034    async fn test_get_nonexistent_projection() {
4035        let store = create_test_store();
4036        let projection_manager = store.projection_manager();
4037
4038        // Requesting a non-existent projection should return None
4039        let projection = projection_manager.get_projection("nonexistent_projection");
4040        assert!(projection.is_none());
4041    }
4042
4043    #[tokio::test]
4044    async fn test_get_nonexistent_entity_state() {
4045        let store = create_test_store();
4046        let projection_manager = store.projection_manager();
4047
4048        // Get state for non-existent entity
4049        let snapshot_projection = projection_manager
4050            .get_projection("entity_snapshots")
4051            .unwrap();
4052        let state = snapshot_projection.get_state("nonexistent-entity-xyz");
4053        assert!(state.is_none());
4054    }
4055
4056    #[tokio::test]
4057    async fn test_projection_state_cache_concurrent_access() {
4058        let store = create_test_store();
4059        let cache = store.projection_state_cache();
4060
4061        // Simulate concurrent writes
4062        let handles: Vec<_> = (0..10)
4063            .map(|i| {
4064                let cache_clone = cache.clone();
4065                tokio::spawn(async move {
4066                    cache_clone.insert(
4067                        format!("concurrent:entity-{i}"),
4068                        serde_json::json!({"thread": i}),
4069                    );
4070                })
4071            })
4072            .collect();
4073
4074        for handle in handles {
4075            handle.await.unwrap();
4076        }
4077
4078        // All 10 entries should be present
4079        assert_eq!(cache.len(), 10);
4080    }
4081
4082    #[tokio::test]
4083    async fn test_projection_state_large_payload() {
4084        let store = create_test_store();
4085        let cache = store.projection_state_cache();
4086
4087        // Create a large JSON payload (~10KB)
4088        let large_array: Vec<serde_json::Value> = (0..1000)
4089            .map(|i| serde_json::json!({"item": i, "description": "test item with some padding data to increase size"}))
4090            .collect();
4091
4092        cache.insert(
4093            "large:entity-1".to_string(),
4094            serde_json::json!({"items": large_array}),
4095        );
4096
4097        let state = cache.get("large:entity-1").unwrap();
4098        let items = state["items"].as_array().unwrap();
4099        assert_eq!(items.len(), 1000);
4100    }
4101
4102    #[tokio::test]
4103    async fn test_projection_state_complex_json() {
4104        let store = create_test_store();
4105        let cache = store.projection_state_cache();
4106
4107        // Complex nested JSON structure
4108        let complex_state = serde_json::json!({
4109            "user": {
4110                "id": "user-123",
4111                "profile": {
4112                    "name": "John Doe",
4113                    "email": "john@example.com",
4114                    "settings": {
4115                        "theme": "dark",
4116                        "notifications": true
4117                    }
4118                },
4119                "roles": ["admin", "user"],
4120                "metadata": {
4121                    "created_at": "2025-01-01T00:00:00Z",
4122                    "last_login": null
4123                }
4124            }
4125        });
4126
4127        cache.insert("complex:user-123".to_string(), complex_state);
4128
4129        let state = cache.get("complex:user-123").unwrap();
4130        assert_eq!(state["user"]["profile"]["name"], "John Doe");
4131        assert_eq!(state["user"]["roles"][0], "admin");
4132        assert!(state["user"]["metadata"]["last_login"].is_null());
4133    }
4134
4135    #[tokio::test]
4136    async fn test_projection_state_cache_iteration() {
4137        let store = create_test_store();
4138        let cache = store.projection_state_cache();
4139
4140        // Insert entries
4141        for i in 0..5 {
4142            cache.insert(format!("iter:entity-{i}"), serde_json::json!({"index": i}));
4143        }
4144
4145        // Iterate over all entries
4146        let entries: Vec<_> = cache.iter().map(|entry| entry.key().clone()).collect();
4147        assert_eq!(entries.len(), 5);
4148    }
4149
4150    #[tokio::test]
4151    async fn test_projection_manager_get_entity_snapshots() {
4152        let store = create_test_store();
4153        let projection_manager = store.projection_manager();
4154
4155        // Get entity_snapshots projection specifically
4156        let projection = projection_manager.get_projection("entity_snapshots");
4157        assert!(projection.is_some());
4158        assert_eq!(projection.unwrap().name(), "entity_snapshots");
4159    }
4160
4161    #[tokio::test]
4162    async fn test_projection_manager_get_event_counters() {
4163        let store = create_test_store();
4164        let projection_manager = store.projection_manager();
4165
4166        // Get event_counters projection specifically
4167        let projection = projection_manager.get_projection("event_counters");
4168        assert!(projection.is_some());
4169        assert_eq!(projection.unwrap().name(), "event_counters");
4170    }
4171
4172    #[tokio::test]
4173    async fn test_projection_state_cache_overwrite() {
4174        let store = create_test_store();
4175        let cache = store.projection_state_cache();
4176
4177        // Initial value
4178        cache.insert(
4179            "overwrite:entity-1".to_string(),
4180            serde_json::json!({"version": 1}),
4181        );
4182
4183        // Overwrite with new value
4184        cache.insert(
4185            "overwrite:entity-1".to_string(),
4186            serde_json::json!({"version": 2}),
4187        );
4188
4189        // Overwrite again
4190        cache.insert(
4191            "overwrite:entity-1".to_string(),
4192            serde_json::json!({"version": 3}),
4193        );
4194
4195        let state = cache.get("overwrite:entity-1").unwrap();
4196        assert_eq!(state["version"], 3);
4197
4198        // Should still be only 1 entry
4199        assert_eq!(cache.len(), 1);
4200    }
4201
4202    #[tokio::test]
4203    async fn test_projection_state_multiple_projections() {
4204        let store = create_test_store();
4205        let cache = store.projection_state_cache();
4206
4207        // Store states for different projections
4208        cache.insert(
4209            "entity_snapshots:user-1".to_string(),
4210            serde_json::json!({"name": "Alice"}),
4211        );
4212        cache.insert(
4213            "event_counters:user.created".to_string(),
4214            serde_json::json!({"count": 5}),
4215        );
4216        cache.insert(
4217            "custom_projection:order-1".to_string(),
4218            serde_json::json!({"total": 150.0}),
4219        );
4220
4221        // Verify each projection's state
4222        assert_eq!(
4223            cache.get("entity_snapshots:user-1").unwrap()["name"],
4224            "Alice"
4225        );
4226        assert_eq!(
4227            cache.get("event_counters:user.created").unwrap()["count"],
4228            5
4229        );
4230        assert_eq!(
4231            cache.get("custom_projection:order-1").unwrap()["total"],
4232            150.0
4233        );
4234    }
4235
4236    #[tokio::test]
4237    async fn test_bulk_projection_state_access() {
4238        let store = create_test_store();
4239
4240        // Ingest multiple events for different entities
4241        for i in 0..5 {
4242            let event = create_test_event(&format!("bulk-user-{i}"), "user.created");
4243            store.ingest(&event).unwrap();
4244        }
4245
4246        // Get projection and verify bulk access
4247        let projection_manager = store.projection_manager();
4248        let snapshot_projection = projection_manager
4249            .get_projection("entity_snapshots")
4250            .unwrap();
4251
4252        // Verify we can access all entities
4253        for i in 0..5 {
4254            let state = snapshot_projection.get_state(&format!("bulk-user-{i}"));
4255            assert!(state.is_some(), "Entity bulk-user-{i} should have state");
4256        }
4257    }
4258
4259    #[tokio::test]
4260    async fn test_bulk_save_projection_states() {
4261        let store = create_test_store();
4262        let cache = store.projection_state_cache();
4263
4264        // Simulate bulk save request
4265        let states = vec![
4266            BulkSaveStateItem {
4267                entity_id: "bulk-entity-1".to_string(),
4268                state: serde_json::json!({"name": "Entity 1", "value": 100}),
4269            },
4270            BulkSaveStateItem {
4271                entity_id: "bulk-entity-2".to_string(),
4272                state: serde_json::json!({"name": "Entity 2", "value": 200}),
4273            },
4274            BulkSaveStateItem {
4275                entity_id: "bulk-entity-3".to_string(),
4276                state: serde_json::json!({"name": "Entity 3", "value": 300}),
4277            },
4278        ];
4279
4280        let projection_name = "test_projection";
4281
4282        // Save states to cache (simulating bulk_save_projection_states handler)
4283        for item in &states {
4284            cache.insert(
4285                format!("{projection_name}:{}", item.entity_id),
4286                item.state.clone(),
4287            );
4288        }
4289
4290        // Verify all states were saved
4291        assert_eq!(cache.len(), 3);
4292
4293        let state1 = cache.get("test_projection:bulk-entity-1").unwrap();
4294        assert_eq!(state1["name"], "Entity 1");
4295        assert_eq!(state1["value"], 100);
4296
4297        let state2 = cache.get("test_projection:bulk-entity-2").unwrap();
4298        assert_eq!(state2["name"], "Entity 2");
4299        assert_eq!(state2["value"], 200);
4300
4301        let state3 = cache.get("test_projection:bulk-entity-3").unwrap();
4302        assert_eq!(state3["name"], "Entity 3");
4303        assert_eq!(state3["value"], 300);
4304    }
4305
4306    #[tokio::test]
4307    async fn test_bulk_save_empty_states() {
4308        let store = create_test_store();
4309        let cache = store.projection_state_cache();
4310
4311        // Clear cache
4312        cache.clear();
4313
4314        // Empty states should work fine
4315        let states: Vec<BulkSaveStateItem> = vec![];
4316        assert_eq!(states.len(), 0);
4317
4318        // Cache should remain empty
4319        assert_eq!(cache.len(), 0);
4320    }
4321
4322    #[tokio::test]
4323    async fn test_bulk_save_overwrites_existing() {
4324        let store = create_test_store();
4325        let cache = store.projection_state_cache();
4326
4327        // Insert initial state
4328        cache.insert(
4329            "test:entity-1".to_string(),
4330            serde_json::json!({"version": 1, "data": "initial"}),
4331        );
4332
4333        // Bulk save with updated state
4334        let new_state = serde_json::json!({"version": 2, "data": "updated"});
4335        cache.insert("test:entity-1".to_string(), new_state);
4336
4337        // Verify overwrite
4338        let state = cache.get("test:entity-1").unwrap();
4339        assert_eq!(state["version"], 2);
4340        assert_eq!(state["data"], "updated");
4341    }
4342
4343    #[tokio::test]
4344    async fn test_bulk_save_high_volume() {
4345        let store = create_test_store();
4346        let cache = store.projection_state_cache();
4347
4348        // Simulate high volume save (1000 entities)
4349        for i in 0..1000 {
4350            cache.insert(
4351                format!("volume_test:entity-{i}"),
4352                serde_json::json!({"index": i, "status": "active"}),
4353            );
4354        }
4355
4356        // Verify count
4357        assert_eq!(cache.len(), 1000);
4358
4359        // Spot check some entries
4360        assert_eq!(cache.get("volume_test:entity-0").unwrap()["index"], 0);
4361        assert_eq!(cache.get("volume_test:entity-500").unwrap()["index"], 500);
4362        assert_eq!(cache.get("volume_test:entity-999").unwrap()["index"], 999);
4363    }
4364
4365    #[tokio::test]
4366    async fn test_bulk_save_different_projections() {
4367        let store = create_test_store();
4368        let cache = store.projection_state_cache();
4369
4370        // Save to multiple projections in bulk
4371        let projections = ["entity_snapshots", "event_counters", "custom_analytics"];
4372
4373        for proj in &projections {
4374            for i in 0..5 {
4375                cache.insert(
4376                    format!("{proj}:entity-{i}"),
4377                    serde_json::json!({"projection": proj, "id": i}),
4378                );
4379            }
4380        }
4381
4382        // Verify total count (3 projections * 5 entities)
4383        assert_eq!(cache.len(), 15);
4384
4385        // Verify each projection
4386        for proj in &projections {
4387            let state = cache.get(&format!("{proj}:entity-0")).unwrap();
4388            assert_eq!(state["projection"], *proj);
4389        }
4390    }
4391
4392    // -------------------------------------------------------------------------
4393    // Cache-fallback tests for the projection-state read handlers (v0.19.1).
4394    //
4395    // SDK-managed projections write state via save_projection_state /
4396    // bulk_save_projection_states without registering in projection_manager.
4397    // These tests verify the read handlers fall back to the cache instead of
4398    // 404-ing. Registered projection wins where both exist; cache fills the
4399    // gap when not.
4400    // -------------------------------------------------------------------------
4401
4402    #[tokio::test]
4403    async fn get_projection_state_falls_back_to_cache_when_unregistered() {
4404        let store = create_test_store();
4405        store.projection_state_cache().insert(
4406            "assets:BTC".to_string(),
4407            serde_json::json!({"symbol": "BTC", "altname": "Bitcoin"}),
4408        );
4409
4410        let resp = get_projection_state(
4411            State(Arc::clone(&store)),
4412            Path(("assets".to_string(), "BTC".to_string())),
4413        )
4414        .await
4415        .expect("should not error when projection is not registered");
4416
4417        assert_eq!(resp.0["found"], serde_json::Value::Bool(true));
4418        assert_eq!(resp.0["state"]["symbol"], "BTC");
4419        assert_eq!(resp.0["state"]["altname"], "Bitcoin");
4420    }
4421
4422    #[tokio::test]
4423    async fn get_projection_state_returns_not_found_when_absent_everywhere() {
4424        let store = create_test_store();
4425
4426        let resp = get_projection_state(
4427            State(Arc::clone(&store)),
4428            Path(("assets".to_string(), "UNKNOWN".to_string())),
4429        )
4430        .await
4431        .unwrap();
4432
4433        assert_eq!(resp.0["found"], serde_json::Value::Bool(false));
4434        assert_eq!(resp.0["state"], serde_json::Value::Null);
4435    }
4436
4437    #[tokio::test]
4438    async fn get_projection_state_registered_wins_over_cache() {
4439        let store = create_test_store();
4440
4441        // Ingest an event so entity_snapshots (a registered projection) has state.
4442        let event = create_test_event("user-777", "user.created");
4443        store.ingest(&event).unwrap();
4444
4445        // Plant a conflicting cache entry for the same (projection, entity).
4446        store.projection_state_cache().insert(
4447            "entity_snapshots:user-777".to_string(),
4448            serde_json::json!({"stolen": "value"}),
4449        );
4450
4451        let resp = get_projection_state(
4452            State(Arc::clone(&store)),
4453            Path(("entity_snapshots".to_string(), "user-777".to_string())),
4454        )
4455        .await
4456        .unwrap();
4457
4458        // Registered projection wins — cache fallback is only consulted when
4459        // the registered projection has no state for this entity.
4460        assert_eq!(resp.0["found"], serde_json::Value::Bool(true));
4461        assert!(
4462            resp.0["state"].get("stolen").is_none(),
4463            "cache entry must not shadow registered projection state: got {:?}",
4464            resp.0["state"]
4465        );
4466    }
4467
4468    #[tokio::test]
4469    async fn get_projection_state_summary_returns_cache_without_registration() {
4470        let store = create_test_store();
4471        let cache = store.projection_state_cache();
4472        cache.insert("assets:BTC".into(), serde_json::json!({"symbol": "BTC"}));
4473        cache.insert("assets:ETH".into(), serde_json::json!({"symbol": "ETH"}));
4474        // Different projection name — must not appear in the summary.
4475        cache.insert("trades:t-1".into(), serde_json::json!({"x": 1}));
4476
4477        let resp = get_projection_state_summary(
4478            State(Arc::clone(&store)),
4479            Path("assets".to_string()),
4480            Query(ProjectionStateSummaryParams::default()),
4481        )
4482        .await
4483        .unwrap();
4484
4485        assert_eq!(resp.0["total"], 2);
4486        let states = resp.0["states"].as_array().unwrap();
4487        let entity_ids: Vec<&str> = states
4488            .iter()
4489            .map(|s| s["entity_id"].as_str().unwrap())
4490            .collect();
4491        assert!(entity_ids.contains(&"BTC"));
4492        assert!(entity_ids.contains(&"ETH"));
4493    }
4494
4495    #[tokio::test]
4496    async fn bulk_get_projection_states_falls_back_to_cache() {
4497        let store = create_test_store();
4498        let cache = store.projection_state_cache();
4499        cache.insert("assets:BTC".into(), serde_json::json!({"symbol": "BTC"}));
4500        cache.insert("assets:ETH".into(), serde_json::json!({"symbol": "ETH"}));
4501
4502        let req = BulkGetStateRequest {
4503            entity_ids: vec!["BTC".into(), "ETH".into(), "MISSING".into()],
4504        };
4505
4506        let resp = bulk_get_projection_states(
4507            State(Arc::clone(&store)),
4508            Path("assets".to_string()),
4509            Json(req),
4510        )
4511        .await
4512        .unwrap();
4513
4514        assert_eq!(resp.0["total"], 3);
4515        let states = resp.0["states"].as_array().unwrap();
4516        let by_id: std::collections::HashMap<&str, &serde_json::Value> = states
4517            .iter()
4518            .map(|s| (s["entity_id"].as_str().unwrap(), s))
4519            .collect();
4520
4521        assert_eq!(by_id["BTC"]["found"], serde_json::Value::Bool(true));
4522        assert_eq!(by_id["BTC"]["state"]["symbol"], "BTC");
4523        assert_eq!(by_id["ETH"]["found"], serde_json::Value::Bool(true));
4524        assert_eq!(by_id["MISSING"]["found"], serde_json::Value::Bool(false));
4525    }
4526
4527    /// Pins the poll wire shape the Rust SDK's `poll_consumer_events` decodes:
4528    /// `ConsumerEventDto` flattens the event, so its fields sit next to
4529    /// `position` rather than nested under an `event` key.
4530    #[tokio::test]
4531    async fn poll_consumer_events_flattens_event_alongside_position() {
4532        let store = create_test_store();
4533        store
4534            .ingest(&create_test_event("user-1", "user.created"))
4535            .unwrap();
4536        store
4537            .ingest(&create_test_event("user-2", "user.updated"))
4538            .unwrap();
4539        store.consumer_registry().register("w1", &[]);
4540
4541        let resp = poll_consumer_events(
4542            State(Arc::clone(&store)),
4543            Path("w1".to_string()),
4544            Query(ConsumerPollQuery { limit: Some(10) }),
4545        )
4546        .await
4547        .unwrap();
4548
4549        let body = serde_json::to_value(&resp.0).unwrap();
4550        assert_eq!(body["count"], 2);
4551        let first = &body["events"][0];
4552        assert_eq!(first["position"], 1);
4553        assert!(
4554            first.get("event").is_none(),
4555            "event must be flattened, not nested: got {first:?}"
4556        );
4557        assert_eq!(first["event_type"], "user.created");
4558        assert_eq!(first["entity_id"], "user-1");
4559    }
4560}