allsource-core 0.23.0

High-performance event store core built in Rust
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
use crate::{
    domain::{
        entities::{SchemaEnforcement, TenantQuotas, UsageMeter},
        value_objects::TenantId,
    },
    infrastructure::security::middleware::{Admin, Authenticated},
};
use axum::{Json, extract::State, http::StatusCode};
use serde::{Deserialize, Serialize};

// AppState is defined in api_v1.rs
use crate::infrastructure::web::api_v1::AppState;

// ============================================================================
// Request/Response Types
// ============================================================================

#[derive(Debug, Deserialize)]
pub struct CreateTenantRequest {
    pub id: String,
    pub name: String,
    pub description: Option<String>,
    pub quota_preset: Option<String>, // "trial", "free", "professional", "unlimited"
    pub quotas: Option<TenantQuotas>,
    #[serde(default)]
    pub is_demo: bool,
}

#[derive(Debug, Serialize)]
pub struct TenantResponse {
    pub id: String,
    pub name: String,
    pub description: Option<String>,
    pub quotas: TenantQuotas,
    pub created_at: chrono::DateTime<chrono::Utc>,
    pub updated_at: chrono::DateTime<chrono::Utc>,
    pub active: bool,
    pub is_demo: bool,
    pub schema_enforcement: SchemaEnforcement,
    /// Operational metadata (subscription tier, quotas, billing). Persisted via
    /// PUT /tenants/{id}; the Control Plane reads subscription/entitlement state
    /// from here. Must be serialized back out or the plan reads as "free".
    pub metadata: serde_json::Value,
}

impl TenantResponse {
    fn from_domain(tenant: &crate::domain::entities::Tenant) -> Self {
        Self {
            id: tenant.id().as_str().to_string(),
            name: tenant.name().to_string(),
            description: tenant.description().map(std::string::ToString::to_string),
            quotas: tenant.quotas().clone(),
            created_at: tenant.created_at(),
            updated_at: tenant.updated_at(),
            active: tenant.is_active(),
            is_demo: tenant.is_demo(),
            schema_enforcement: tenant.schema_enforcement(),
            metadata: tenant.metadata().clone(),
        }
    }
}

#[derive(Debug, Deserialize)]
pub struct UpdateTenantRequest {
    pub name: Option<String>,
    pub description: Option<String>,
    pub is_demo: Option<bool>,
    pub quotas: Option<TenantQuotas>,
    /// Operational metadata (subscription tier, quotas, billing). Replaces the
    /// tenant's metadata when present — callers send the full merged map. This
    /// is how the Control Plane persists subscription/entitlement state.
    pub metadata: Option<serde_json::Value>,
}

#[derive(Debug, Deserialize)]
pub struct UpdateQuotasRequest {
    pub quotas: TenantQuotas,
}

#[derive(Debug, Deserialize)]
pub struct UpdateSchemaEnforcementRequest {
    /// One of `permissive`, `warn`, `strict`.
    pub schema_enforcement: SchemaEnforcement,
}

/// Body for `POST /api/v1/tenants/{id}/usage/increment`.
///
/// `type` is the [`UsageMeter`] discriminator and DEFAULTS to `events` when
/// absent — that keeps the original untyped `{"count": N}` body the Query
/// Service sent before this field existed valid during a rolling deploy.
#[derive(Debug, Deserialize)]
pub struct IncrementUsageRequest {
    /// How much to add to the counter. Must be > 0. A `u64` field also makes
    /// serde reject negative numbers with a 400 before the handler runs.
    pub count: u64,
    /// Which meter to bump: `events` → `events_used`, `queries` →
    /// `queries_used`. Defaults to `events`.
    #[serde(default, rename = "type")]
    pub meter: UsageMeter,
}

/// Returned by the usage-increment endpoint: the meter that was bumped and its
/// new value, so the caller (and tests) can confirm the write landed.
#[derive(Debug, Serialize)]
pub struct IncrementUsageResponse {
    pub tenant_id: String,
    /// `events` or `queries`.
    #[serde(rename = "type")]
    pub meter: UsageMeter,
    /// The post-increment counter value (`metadata.quotas.<field>`).
    pub used: u64,
}

// ============================================================================
// Handlers
// ============================================================================

/// Resolve the quotas for a tenant being created.
///
/// Policy (prompt 048): NEW self-service tenants default to the TRIAL preset,
/// never free — a 14-day evaluation, not a permanent free tier (the Control
/// Plane owns the `trial_expires_at` clock + the expiry suspend sweep). The
/// `free_tier()` quota is RETAINED but only resolves when a caller EXPLICITLY
/// asks for `quota_preset: "free"` — that path exists for GRANDFATHERED tenants
/// (an operator re-provisioning an existing free tenant), not for new signups.
///
/// Precedence: an explicit `quotas` value wins; otherwise a recognized preset;
/// otherwise (unknown preset OR no preset at all) we fall back to the trial
/// tier, so an under-specified create can never silently mint a free-quota
/// tenant.
fn resolve_create_quotas(quotas: Option<TenantQuotas>, quota_preset: Option<&str>) -> TenantQuotas {
    if let Some(quotas) = quotas {
        return quotas;
    }
    match quota_preset {
        Some("trial") => TenantQuotas::trial_tier(),
        Some("free") => TenantQuotas::free_tier(), // grandfathered tenants only
        Some("professional") => TenantQuotas::professional(),
        Some("unlimited") => TenantQuotas::unlimited(),
        _ => TenantQuotas::trial_tier(),
    }
}

/// Create tenant (admin only)
/// POST /api/v1/tenants
pub async fn create_tenant_handler(
    State(state): State<AppState>,
    Admin(_): Admin,
    Json(req): Json<CreateTenantRequest>,
) -> Result<(StatusCode, Json<TenantResponse>), (StatusCode, String)> {
    // Determine quotas (see resolve_create_quotas for the policy + rationale).
    let quotas = resolve_create_quotas(req.quotas, req.quota_preset.as_deref());

    let tenant_id = TenantId::new(req.id).map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;

    let mut tenant = state
        .tenant_repo
        .create(tenant_id, req.name, quotas)
        .await
        .map_err(|e| {
            let status = match &e {
                crate::error::AllSourceError::TenantAlreadyExists(_) => StatusCode::CONFLICT,
                _ => StatusCode::BAD_REQUEST,
            };
            (status, e.to_string())
        })?;

    // Apply optional fields that aren't part of create()
    if let Some(desc) = req.description {
        tenant.update_description(Some(desc));
    }
    if req.is_demo {
        tenant.set_is_demo(true);
    }

    // Persist the updated fields
    state
        .tenant_repo
        .save(&tenant)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;

    Ok((
        StatusCode::CREATED,
        Json(TenantResponse::from_domain(&tenant)),
    ))
}

/// Get tenant
/// GET /api/v1/tenants/:id
pub async fn get_tenant_handler(
    State(state): State<AppState>,
    Authenticated(auth_ctx): Authenticated,
    axum::extract::Path(tenant_id): axum::extract::Path<String>,
) -> Result<Json<TenantResponse>, (StatusCode, String)> {
    // Users can only view their own tenant, admins can view any
    if tenant_id != auth_ctx.tenant_id() {
        auth_ctx
            .require_permission(crate::infrastructure::security::auth::Permission::Admin)
            .map_err(|_| {
                (
                    StatusCode::FORBIDDEN,
                    "Can only view own tenant".to_string(),
                )
            })?;
    }

    let tid =
        TenantId::new(tenant_id.clone()).map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;

    let tenant = state
        .tenant_repo
        .find_by_id(&tid)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?
        .ok_or_else(|| {
            (
                StatusCode::NOT_FOUND,
                format!("Tenant not found: {tenant_id}"),
            )
        })?;

    Ok(Json(TenantResponse::from_domain(&tenant)))
}

/// Merge a partial object into a tenant's metadata (tenant-scoped)
/// PATCH /api/v1/tenants/:id/metadata
///
/// Deep-merges the JSON object body into `metadata`, preserving every sibling
/// key, and persists atomically against concurrent quota bumps. The Query
/// Service uses this to store a tenant's opaque enabled-projection set
/// without clobbering `metadata.quotas` — Core does not interpret the merged
/// keys (see `docs/proposals/PER_TENANT_PROJECTIONS.md`). Unlike the admin-only
/// `PUT /tenants/:id` (which replaces the whole blob), this is scoped: a caller
/// may patch only its own tenant; admins may patch any.
pub async fn merge_tenant_metadata_handler(
    State(state): State<AppState>,
    Authenticated(auth_ctx): Authenticated,
    axum::extract::Path(tenant_id): axum::extract::Path<String>,
    Json(partial): Json<serde_json::Value>,
) -> Result<Json<TenantResponse>, (StatusCode, String)> {
    // Callers may patch only their own tenant; admins may patch any.
    if tenant_id != auth_ctx.tenant_id() {
        auth_ctx
            .require_permission(crate::infrastructure::security::auth::Permission::Admin)
            .map_err(|_| {
                (
                    StatusCode::FORBIDDEN,
                    "Can only modify own tenant".to_string(),
                )
            })?;
    }

    if !partial.is_object() {
        return Err((
            StatusCode::BAD_REQUEST,
            "metadata patch must be a JSON object".to_string(),
        ));
    }

    let tid =
        TenantId::new(tenant_id.clone()).map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;

    if state
        .tenant_repo
        .merge_metadata(&tid, partial)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?
        .is_none()
    {
        return Err((
            StatusCode::NOT_FOUND,
            format!("Tenant not found: {tenant_id}"),
        ));
    }

    // Re-fetch for the full updated representation (metadata was mutated and
    // persisted under the per-tenant lock inside merge_metadata).
    let tenant = state
        .tenant_repo
        .find_by_id(&tid)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?
        .ok_or_else(|| {
            (
                StatusCode::NOT_FOUND,
                format!("Tenant not found: {tenant_id}"),
            )
        })?;

    Ok(Json(TenantResponse::from_domain(&tenant)))
}

/// List all tenants (admin only)
/// GET /api/v1/tenants
pub async fn list_tenants_handler(
    State(state): State<AppState>,
    Admin(_): Admin,
) -> Result<Json<Vec<TenantResponse>>, (StatusCode, String)> {
    let tenants = state
        .tenant_repo
        .find_all(10_000, 0)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
    Ok(Json(
        tenants.iter().map(TenantResponse::from_domain).collect(),
    ))
}

/// Get tenant statistics
/// GET /api/v1/tenants/:id/stats
pub async fn get_tenant_stats_handler(
    State(state): State<AppState>,
    Authenticated(auth_ctx): Authenticated,
    axum::extract::Path(tenant_id): axum::extract::Path<String>,
) -> Result<Json<serde_json::Value>, (StatusCode, String)> {
    // Users can only view their own tenant stats
    if tenant_id != auth_ctx.tenant_id() {
        auth_ctx
            .require_permission(crate::infrastructure::security::auth::Permission::Admin)
            .map_err(|_| {
                (
                    StatusCode::FORBIDDEN,
                    "Can only view own tenant stats".to_string(),
                )
            })?;
    }

    let tid =
        TenantId::new(tenant_id.clone()).map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;

    let tenant = state
        .tenant_repo
        .find_by_id(&tid)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?
        .ok_or_else(|| {
            (
                StatusCode::NOT_FOUND,
                format!("Tenant not found: {tenant_id}"),
            )
        })?;

    let stats = build_tenant_stats(&tenant);
    Ok(Json(stats))
}

/// Update tenant quotas (admin only)
/// PUT /api/v1/tenants/:id/quotas
pub async fn update_quotas_handler(
    State(state): State<AppState>,
    Admin(_): Admin,
    axum::extract::Path(tenant_id): axum::extract::Path<String>,
    Json(req): Json<UpdateQuotasRequest>,
) -> Result<StatusCode, (StatusCode, String)> {
    let tid =
        TenantId::new(tenant_id.clone()).map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;

    let updated = state
        .tenant_repo
        .update_quotas(&tid, req.quotas)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;

    if !updated {
        return Err((
            StatusCode::NOT_FOUND,
            format!("Tenant not found: {tenant_id}"),
        ));
    }

    Ok(StatusCode::NO_CONTENT)
}

/// Increment a tenant's forward usage counter (admin only).
/// POST /api/v1/tenants/:id/usage/increment
///
/// Write half of forward usage-metering. The Query Service POSTs one increment
/// per metered activity — events ingested (`type: "events"`) or queries run
/// (`type: "queries"`) — and Core atomically bumps the matching
/// `metadata.quotas.{events_used,queries_used}` counter the dashboard reads.
///
/// Body: `{ "count": <u64>, "type": "events" | "queries" }`. `type` defaults
/// to `"events"` so the pre-typing `{count}` body keeps working during a
/// rolling deploy. Returns the post-increment value. `404` if the tenant is
/// unknown; `400` (via serde) for a missing/negative `count`; `400` for an
/// explicit `count: 0` (a no-op increment is a client bug, not a write).
///
/// Atomicity is the repository's job — see
/// [`TenantRepository::increment_usage`]; concurrent batched increments for
/// the same tenant are serialized so no counts are lost.
pub async fn increment_usage_handler(
    State(state): State<AppState>,
    Admin(_): Admin,
    axum::extract::Path(tenant_id): axum::extract::Path<String>,
    Json(req): Json<IncrementUsageRequest>,
) -> Result<Json<IncrementUsageResponse>, (StatusCode, String)> {
    if req.count == 0 {
        return Err((
            StatusCode::BAD_REQUEST,
            "count must be greater than zero".to_string(),
        ));
    }

    let tid =
        TenantId::new(tenant_id.clone()).map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;

    let new_value = state
        .tenant_repo
        .increment_usage(&tid, req.meter, req.count)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?
        .ok_or_else(|| {
            (
                StatusCode::NOT_FOUND,
                format!("Tenant not found: {tenant_id}"),
            )
        })?;

    Ok(Json(IncrementUsageResponse {
        tenant_id,
        meter: req.meter,
        used: new_value,
    }))
}

/// Set a tenant's schema-enforcement mode (admin only).
/// PUT /api/v1/tenants/:id/schema-enforcement
///
/// Gap 3 toggle: flips registered-schema validation on ingest between
/// `permissive` (default), `warn`, and `strict`. The gateway exposes this to a
/// tenant for their own tenant; Core keeps it admin-gated and internal.
pub async fn update_schema_enforcement_handler(
    State(state): State<AppState>,
    Admin(_): Admin,
    axum::extract::Path(tenant_id): axum::extract::Path<String>,
    Json(req): Json<UpdateSchemaEnforcementRequest>,
) -> Result<StatusCode, (StatusCode, String)> {
    let tid =
        TenantId::new(tenant_id.clone()).map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;

    let updated = state
        .tenant_repo
        .update_schema_enforcement(&tid, req.schema_enforcement)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;

    if !updated {
        return Err((
            StatusCode::NOT_FOUND,
            format!("Tenant not found: {tenant_id}"),
        ));
    }

    Ok(StatusCode::NO_CONTENT)
}

/// Deactivate tenant (admin only)
/// POST /api/v1/tenants/:id/deactivate
pub async fn deactivate_tenant_handler(
    State(state): State<AppState>,
    Admin(_): Admin,
    axum::extract::Path(tenant_id): axum::extract::Path<String>,
) -> Result<StatusCode, (StatusCode, String)> {
    let tid =
        TenantId::new(tenant_id.clone()).map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;

    let deactivated = state
        .tenant_repo
        .deactivate(&tid)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;

    if !deactivated {
        return Err((
            StatusCode::BAD_REQUEST,
            format!("Tenant not found: {tenant_id}"),
        ));
    }

    Ok(StatusCode::NO_CONTENT)
}

/// Activate tenant (admin only)
/// POST /api/v1/tenants/:id/activate
pub async fn activate_tenant_handler(
    State(state): State<AppState>,
    Admin(_): Admin,
    axum::extract::Path(tenant_id): axum::extract::Path<String>,
) -> Result<StatusCode, (StatusCode, String)> {
    let tid =
        TenantId::new(tenant_id.clone()).map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;

    let activated = state
        .tenant_repo
        .activate(&tid)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;

    if !activated {
        return Err((
            StatusCode::NOT_FOUND,
            format!("Tenant not found: {tenant_id}"),
        ));
    }

    Ok(StatusCode::NO_CONTENT)
}

/// Update tenant (admin only)
/// PUT /api/v1/tenants/:id
pub async fn update_tenant_handler(
    State(state): State<AppState>,
    Admin(_): Admin,
    axum::extract::Path(tenant_id): axum::extract::Path<String>,
    Json(req): Json<UpdateTenantRequest>,
) -> Result<Json<TenantResponse>, (StatusCode, String)> {
    let tid =
        TenantId::new(tenant_id.clone()).map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;

    let mut tenant = state
        .tenant_repo
        .find_by_id(&tid)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?
        .ok_or_else(|| {
            (
                StatusCode::NOT_FOUND,
                format!("Tenant not found: {tenant_id}"),
            )
        })?;

    if let Some(name) = req.name {
        tenant
            .update_name(name)
            .map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;
    }
    if let Some(desc) = req.description {
        tenant.update_description(Some(desc));
    }
    if let Some(is_demo) = req.is_demo {
        tenant.set_is_demo(is_demo);
    }
    if let Some(quotas) = req.quotas {
        tenant.update_quotas(quotas);
    }
    if let Some(metadata) = req.metadata {
        tenant.update_metadata(metadata);
    }

    state
        .tenant_repo
        .save(&tenant)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;

    Ok(Json(TenantResponse::from_domain(&tenant)))
}

/// Delete tenant (admin only)
/// DELETE /api/v1/tenants/:id
pub async fn delete_tenant_handler(
    State(state): State<AppState>,
    Admin(_): Admin,
    axum::extract::Path(tenant_id): axum::extract::Path<String>,
) -> Result<StatusCode, (StatusCode, String)> {
    let tid =
        TenantId::new(tenant_id.clone()).map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?;

    let deleted = state
        .tenant_repo
        .delete(&tid)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;

    if !deleted {
        return Err((
            StatusCode::BAD_REQUEST,
            format!("Tenant not found: {tenant_id}"),
        ));
    }

    Ok(StatusCode::NO_CONTENT)
}

// ============================================================================
// Helpers
// ============================================================================

/// Read a non-negative meter (`events_used` / `queries_used`) from the tenant's
/// `metadata.quotas` blob.
///
/// WHY this exists: the durable per-tenant event/query totals live in
/// `metadata.quotas.{events_used,queries_used}` — the counter the forward-metering
/// path bumps (`POST /api/v1/tenants/{id}/usage/increment`,
/// [`crate::domain::entities::UsageMeter`]) and the backfill reconciles. The
/// in-memory [`crate::domain::entities::TenantUsage`] struct's `total_events`
/// field is only ever moved by `record_event()`, which the real ingest path never
/// calls — so `usage.total_events` is structurally always 0. `/stats` must report
/// the metered counter, the same number the Query Service dashboard reads
/// (`tenant_controller.ex` `events_used`), or every downstream consumer (the admin
/// tenant list/detail event_count, fleet-health's has-data gate, cluster status)
/// reads 0 for tenants that demonstrably have events. JSON numbers arrive as
/// u64/f64 depending on the writer; accept either and treat anything else as 0.
fn metered_quota_counter(tenant: &crate::domain::entities::Tenant, field: &str) -> u64 {
    tenant
        .metadata()
        .get("quotas")
        .and_then(|q| q.get(field))
        .and_then(|v| v.as_u64().or_else(|| v.as_f64().map(|f| f.max(0.0) as u64)))
        .unwrap_or(0)
}

/// Build tenant statistics JSON (presentation concern).
fn build_tenant_stats(tenant: &crate::domain::entities::Tenant) -> serde_json::Value {
    let quotas = tenant.quotas();
    let usage = tenant.usage();

    // The authoritative lifetime totals come from the metered metadata counter,
    // NOT the in-memory TenantUsage struct (see `metered_quota_counter`). These are
    // the numbers the dashboard + admin console show.
    let events_used = metered_quota_counter(tenant, "events_used");
    let queries_used = metered_quota_counter(tenant, "queries_used");

    let events_pct = if quotas.max_events_per_day() > 0 {
        (usage.events_today() as f64 / quotas.max_events_per_day() as f64) * 100.0
    } else {
        0.0
    };

    let storage_pct = if quotas.max_storage_bytes() > 0 {
        (usage.storage_bytes() as f64 / quotas.max_storage_bytes() as f64) * 100.0
    } else {
        0.0
    };

    let queries_pct = if quotas.max_queries_per_hour() > 0 {
        (usage.queries_this_hour() as f64 / quotas.max_queries_per_hour() as f64) * 100.0
    } else {
        0.0
    };

    // Serialize the live usage struct, then overlay the real metered lifetime
    // totals so `usage.total_events` reflects the durable counter instead of the
    // always-0 in-memory field. The Control Plane's tolerant stats decoder
    // (tenant_stats.go) reads `usage.total_events` (then `event_count`), so this
    // overlay alone fixes every CP consumer with no CP change required.
    let mut usage_json = serde_json::to_value(usage).unwrap_or_else(|_| serde_json::json!({}));
    if let Some(obj) = usage_json.as_object_mut() {
        obj.insert("total_events".to_string(), serde_json::json!(events_used));
        obj.insert("queries_used".to_string(), serde_json::json!(queries_used));
    }

    serde_json::json!({
        "tenant_id": tenant.id().as_str(),
        "name": tenant.name(),
        "active": tenant.is_active(),
        "is_demo": tenant.is_demo(),
        // Flat lifetime totals — the primary fields the admin Events/Queries
        // columns want. Sourced from the durable metered counter.
        "event_count": events_used,
        "query_count": queries_used,
        "usage": usage_json,
        "quotas": tenant.quotas(),
        "utilization": {
            "events_today": {
                "used": usage.events_today(),
                "limit": quotas.max_events_per_day(),
                "percentage": events_pct.min(100.0)
            },
            "storage": {
                "used_bytes": usage.storage_bytes(),
                "limit_bytes": quotas.max_storage_bytes(),
                "percentage": storage_pct.min(100.0)
            },
            "queries_this_hour": {
                "used": usage.queries_this_hour(),
                "limit": quotas.max_queries_per_hour(),
                "percentage": queries_pct.min(100.0)
            }
        },
        "created_at": tenant.created_at(),
        "updated_at": tenant.updated_at()
    })
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn increment_request_defaults_meter_to_events() {
        // Backward compat: the pre-typing Query Service sent just `{count}`.
        // That body must still parse and target the events meter, so an old
        // QS hitting a new Core during a rolling deploy keeps metering events.
        let req: IncrementUsageRequest = serde_json::from_str(r#"{"count": 7}"#).unwrap();
        assert_eq!(req.count, 7);
        assert_eq!(req.meter, UsageMeter::Events);
        assert_eq!(req.meter.quota_field(), "events_used");
    }

    #[test]
    fn increment_request_parses_typed_events_and_queries() {
        let ev: IncrementUsageRequest =
            serde_json::from_str(r#"{"count": 3, "type": "events"}"#).unwrap();
        assert_eq!(ev.meter, UsageMeter::Events);

        let qs: IncrementUsageRequest =
            serde_json::from_str(r#"{"count": 9, "type": "queries"}"#).unwrap();
        assert_eq!(qs.meter, UsageMeter::Queries);
        assert_eq!(qs.meter.quota_field(), "queries_used");
    }

    #[test]
    fn increment_request_rejects_negative_count() {
        // `count: u64` makes serde reject negatives before the handler runs,
        // which axum surfaces as a 400 — no separate validation needed.
        let err = serde_json::from_str::<IncrementUsageRequest>(r#"{"count": -1}"#);
        assert!(err.is_err(), "negative count must fail to deserialize");
    }

    #[test]
    fn new_tenant_defaults_to_trial_quota_not_free() {
        // Prompt 048: a create with no quotas and no preset must resolve to the
        // TRIAL tier (the new-signup default), never free and never the old
        // standard fallback — so an under-specified create can't mint free.
        assert_eq!(
            resolve_create_quotas(None, None),
            TenantQuotas::trial_tier()
        );
        // An unknown preset also falls back to trial, never free.
        assert_eq!(
            resolve_create_quotas(None, Some("bogus")),
            TenantQuotas::trial_tier()
        );
        // The explicit "trial" preset resolves to the trial tier.
        assert_eq!(
            resolve_create_quotas(None, Some("trial")),
            TenantQuotas::trial_tier()
        );
    }

    #[test]
    fn free_preset_still_resolves_to_free_for_grandfathered_tenants() {
        // The free_tier() quota is RETAINED: an operator explicitly asking for
        // `quota_preset:"free"` (re-provisioning a grandfathered tenant) still
        // gets the free quota. This change stops NEW free, it does not delete it.
        assert_eq!(
            resolve_create_quotas(None, Some("free")),
            TenantQuotas::free_tier()
        );
    }

    #[test]
    fn explicit_quotas_win_over_preset_default() {
        // A caller-supplied quotas value always wins (e.g. the demo flow sending
        // unlimited), regardless of the trial default.
        let custom = TenantQuotas::unlimited();
        assert_eq!(resolve_create_quotas(Some(custom.clone()), None), custom);
    }

    #[test]
    fn increment_response_serializes_type_and_used() {
        let resp = IncrementUsageResponse {
            tenant_id: "acme".to_string(),
            meter: UsageMeter::Queries,
            used: 42,
        };
        let v = serde_json::to_value(&resp).unwrap();
        assert_eq!(v["tenant_id"], "acme");
        assert_eq!(v["type"], "queries");
        assert_eq!(v["used"], 42);
    }

    use crate::domain::{
        entities::{Tenant, TenantQuotas},
        value_objects::TenantId,
    };

    fn tenant_with_quota_usage(events_used: u64, queries_used: u64) -> Tenant {
        let mut t = Tenant::new(
            TenantId::new("acme".to_string()).unwrap(),
            "Acme".to_string(),
            TenantQuotas::free_tier(),
        )
        .unwrap();
        // Mirror exactly what the metering increment path writes:
        // metadata.quotas.{events_used,queries_used}.
        t.update_metadata(serde_json::json!({
            "quotas": { "events_used": events_used, "queries_used": queries_used }
        }));
        t
    }

    #[test]
    fn stats_event_count_reflects_metered_counter_not_inmemory_usage() {
        // The in-memory TenantUsage.total_events is never bumped by the real
        // ingest path, so /stats must report the durable metered counter from
        // metadata.quotas.events_used — the same number the QS dashboard shows.
        let tenant = tenant_with_quota_usage(257, 12);
        let stats = build_tenant_stats(&tenant);

        // Flat fields the admin Events/Queries columns prefer.
        assert_eq!(stats["event_count"], 257);
        assert_eq!(stats["query_count"], 12);
        // Overlaid into the usage block so the CP's tolerant decoder
        // (tenant_stats.go reads usage.total_events) also picks it up.
        assert_eq!(stats["usage"]["total_events"], 257);
        assert_eq!(stats["usage"]["queries_used"], 12);
    }

    #[test]
    fn stats_event_count_is_zero_when_unmetered() {
        // A tenant whose metadata has no quota counters reports 0 (guarded),
        // never a missing field or a crash.
        let tenant = Tenant::new(
            TenantId::new("empty".to_string()).unwrap(),
            "Empty".to_string(),
            TenantQuotas::free_tier(),
        )
        .unwrap();
        let stats = build_tenant_stats(&tenant);
        assert_eq!(stats["event_count"], 0);
        assert_eq!(stats["usage"]["total_events"], 0);
    }

    #[tokio::test]
    async fn stats_reflects_real_metering_path_end_to_end() {
        // The full chain the admin console depends on, in one test:
        //   increment_usage (the real QS metering path, writes
        //   metadata.quotas.events_used) -> find_by_id -> build_tenant_stats.
        // Proves the number the metering path records is the number /stats reports
        // as event_count — closing the counts=0 gap at the source.
        use crate::{
            domain::repositories::TenantRepository,
            infrastructure::repositories::InMemoryTenantRepository,
        };

        let repo = InMemoryTenantRepository::new();
        let id = TenantId::new("acme".to_string()).unwrap();
        repo.create(id.clone(), "ACME".to_string(), TenantQuotas::free_tier())
            .await
            .unwrap();

        // Meter 257 events + 12 queries exactly as the Query Service does.
        repo.increment_usage(&id, UsageMeter::Events, 257)
            .await
            .unwrap();
        repo.increment_usage(&id, UsageMeter::Queries, 12)
            .await
            .unwrap();

        let tenant = repo.find_by_id(&id).await.unwrap().unwrap();
        let stats = build_tenant_stats(&tenant);

        assert_eq!(
            stats["event_count"], 257,
            "stats must report the metered events"
        );
        assert_eq!(
            stats["query_count"], 12,
            "stats must report the metered queries"
        );
        assert_eq!(stats["usage"]["total_events"], 257);
    }

    #[test]
    fn stats_metered_counter_accepts_float_json() {
        // JSON numbers can arrive as f64 depending on the writer; the reader
        // must coerce, not drop the value to 0.
        let mut tenant = Tenant::new(
            TenantId::new("floaty".to_string()).unwrap(),
            "Floaty".to_string(),
            TenantQuotas::free_tier(),
        )
        .unwrap();
        tenant.update_metadata(serde_json::json!({
            "quotas": { "events_used": 5.0, "queries_used": 3.0 }
        }));
        let stats = build_tenant_stats(&tenant);
        assert_eq!(stats["event_count"], 5);
        assert_eq!(stats["query_count"], 3);
    }
}