From 5a07a5239e5c84bd26b159e02a4476eae3ca7be0 Mon Sep 17 00:00:00 2001 From: przemek Date: Sat, 16 May 2026 19:48:36 +0200 Subject: [PATCH] =?UTF-8?q?Period=20lifecycle=20(minimal:=20Open=20?= =?UTF-8?q?=E2=86=94=20Closed)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds the spec §13 / §19.10 closed-period semantics that operators need before invoicing: once a billing period is closed, new Usage events for that (account, year, month) tuple are rejected at ingest; Correction and Retraction events are still accepted and become post-close adjustments. Future intermediate states (Closing / Invoiced / Adjusted) and frozen-snapshot query semantics are explicitly out of scope for this PR — they need design + a separate effort. This lands the foundation the rest will build on. API: POST /v1/accounts/{id}/periods/{YYYY-MM}/close Marks the period closed. Idempotent: re-closing returns already_closed=true without changing state. POST /v1/accounts/{id}/periods/{YYYY-MM}/reopen Removes the marker. Returns removed=true if anything was undone. GET /v1/accounts/{id}/periods/{YYYY-MM} Returns: state: "Open" | "Closed" closed_at_ms: i64 or null from_ms, to_ms: period bounds in unix ms total_quantity: current SUM(quantity) for the period (live, not a frozen snapshot — see TODO) Plumbing: - Manifest gains `closed_periods: Vec` with #[serde(default)] so old manifests deserialize unchanged. - New src/period.rs module: period_for_ts, parse_period (YYYY-MM), is_period_closed, find_closed. - handle_ingest snapshots the closed_periods list under the manifest read lock at the start of the batch, then rejects each Usage event whose (account, year, month) matches an entry. Per-event check is O(closed_periods.len()) — tiny in practice. Correction/Retraction events bypass the check. Tests (tests/period_lifecycle.rs, 12 tests): - period_for_ts_returns_utc_year_month - parse_period_round_trips (incl. validation errors) - is_period_closed_finds_match - close_period_persists_in_manifest - close_period_is_idempotent (no duplicate entries) - reopen_period_removes_marker - get_period_returns_state_and_total (open then closed) - close_period_rejects_invalid_format - closed_period_rejects_usage_events - closed_period_only_rejects_target_account (other accounts OK) - closed_period_only_rejects_in_period_timestamps (other months OK) - closed_period_accepts_corrections (adjustments path) Total tests: 134 (was 122; +12). Clean under RUSTFLAGS=-D warnings. Co-Authored-By: Claude Opus 4.7 (1M context) --- README.md | 5 +- src/api/http_server.rs | 182 ++++++++++++++++++ src/lib.rs | 1 + src/period.rs | 61 ++++++ src/storage/manifest.rs | 22 +++ tests/period_lifecycle.rs | 378 ++++++++++++++++++++++++++++++++++++++ 6 files changed, 648 insertions(+), 1 deletion(-) create mode 100644 src/period.rs create mode 100644 tests/period_lifecycle.rs diff --git a/README.md b/README.md index 28a6024..67531b6 100644 --- a/README.md +++ b/README.md @@ -59,6 +59,9 @@ GET /v1/accounts/{account_id}/usage ?from&to&group_by&product_id&met GET /v1/accounts/{account_id}/usage/events ?from&to&meter_id&product_id GET /v1/accounts/{account_id}/explain ?from&to — breakdown + segment provenance + corrections GET /v1/accounts/{account_id}/verify ?from&to — raw-vs-rollup drift check +GET /v1/accounts/{account_id}/periods/{YYYY-MM} — state + total +POST /v1/accounts/{account_id}/periods/{YYYY-MM}/close — mark closed +POST /v1/accounts/{account_id}/periods/{YYYY-MM}/reopen — mark open POST /v1/query/json { "source", "account_id", "from", "to", "group_by", "filters", "metrics" } POST /v1/query/sql { "query": "SELECT meter_id, SUM(quantity) FROM usage_events WHERE account_id = '...' GROUP BY meter_id" } GET /health @@ -66,7 +69,7 @@ GET /health `from` / `to` are RFC 3339. Supported `metrics`: `sum`, `count`. Supported group keys: column names (`account_id`, `product_id`, `meter_id`, `model_id`, `source`, `unit`), `hour_start_ms`, `day`, or any dimension key. The account-usage GET defaults `source=rollup` for fast monthly totals; pass `source=raw` to force a raw scan. -Ingest response counts `accepted`, `duplicates` (same id + same payload), `conflicts` (same id + different payload — surfaces silent collector bugs), and `rejected` (validation failures: missing required IDs, non-positive timestamp, >16 dimensions, Correction/Retraction without `correction_ref`). +Ingest response counts `accepted`, `duplicates` (same id + same payload), `conflicts` (same id + different payload — surfaces silent collector bugs), and `rejected` (validation failures: missing required IDs, non-positive timestamp, >16 dimensions, Correction/Retraction without `correction_ref`, or `Usage` event landing in a closed period). ## Building and running diff --git a/src/api/http_server.rs b/src/api/http_server.rs index 2f97305..81ac9e2 100644 --- a/src/api/http_server.rs +++ b/src/api/http_server.rs @@ -47,6 +47,10 @@ pub fn build_router(state: AppState) -> Router { // Phase D operability — explain a total and verify rollup-vs-raw drift. .route("/v1/accounts/{account_id}/explain", get(handle_explain)) .route("/v1/accounts/{account_id}/verify", get(handle_verify)) + // Period lifecycle (minimal: Open ↔ Closed). + .route("/v1/accounts/{account_id}/periods/{period}", get(handle_get_period)) + .route("/v1/accounts/{account_id}/periods/{period}/close", post(handle_close_period)) + .route("/v1/accounts/{account_id}/periods/{period}/reopen", post(handle_reopen_period)) // Flexible POST query for arbitrary filter shapes. .route("/v1/query/json", post(handle_query_json)) // SQL subset endpoint. @@ -140,6 +144,17 @@ async fn handle_ingest( // wonky clock can't poison TTL eviction). Rejected events never // reach the WAL or dedupe. let ingest_now = now_ms(); + + // Snapshot the closed-periods list under the manifest read lock so + // the per-event check below is cheap (no async lock in the loop). + // The window between snapshot and validation is small enough that a + // racing close_period call doesn't matter — the operator should + // wait for expected ingest to drain before closing anyway. + let closed_snapshot: Vec = { + let manifest = state.manifest.read().await; + manifest.closed_periods.clone() + }; + let mut rejected = 0usize; let mut classified: Vec = Vec::with_capacity(payload.events.len()); for mut event in payload.events { @@ -148,6 +163,28 @@ async fn handle_ingest( tracing::warn!(?reason, event_id = %event.event_id.0, "rejected event"); continue; } + // Period-closed check: reject `Usage` events landing in a + // closed period. Corrections / retractions are intentionally + // allowed through — they become post-close adjustments. + if matches!(event.kind, EventKind::Usage) { + if let Some((year, month)) = crate::period::period_for_ts(event.timestamp_ms) { + if closed_snapshot + .iter() + .any(|p| p.account_id == event.account_id.0 + && p.year == year + && p.month == month) + { + rejected += 1; + tracing::warn!( + event_id = %event.event_id.0, + account = %event.account_id.0, + year, month, + "rejected: Usage event in closed period" + ); + continue; + } + } + } event.ingested_at_ms = ingest_now; let (event_id_hash, payload_hash) = compute_event_hashes(&event); classified.push(Classified { event, event_id_hash, payload_hash }); @@ -657,6 +694,151 @@ async fn handle_verify( }))) } +/// `POST /v1/accounts/{account_id}/periods/{YYYY-MM}/close` +/// +/// Marks the (account, year, month) period as closed. Subsequent `Usage` +/// events with a timestamp inside this period will be rejected at +/// ingest. `Correction` and `Retraction` events for closed periods are +/// still accepted — they become adjustments. +async fn handle_close_period( + State(state): State, + Path((account_id, period)): Path<(String, String)>, +) -> Result, AppError> { + use crate::period::parse_period; + use crate::storage::manifest::ClosedPeriod; + + let (year, month) = parse_period(&period) + .map_err(|e| AppError(anyhow::anyhow!("invalid period: {}", e)))?; + + let mut manifest = state.manifest.write().await; + if manifest.closed_periods.iter().any(|p| { + p.account_id == account_id && p.year == year && p.month == month + }) { + return Ok(Json(serde_json::json!({ + "account_id": account_id, + "period": period, + "state": "Closed", + "already_closed": true, + }))); + } + let closed_at_ms = now_ms(); + manifest.closed_periods.push(ClosedPeriod { + account_id: account_id.clone(), + year, + month, + closed_at_ms, + }); + manifest + .save(&state.config.db_root) + .map_err(|e| AppError(anyhow::anyhow!("manifest save: {}", e)))?; + + Ok(Json(serde_json::json!({ + "account_id": account_id, + "period": period, + "state": "Closed", + "closed_at_ms": closed_at_ms, + }))) +} + +/// `POST /v1/accounts/{account_id}/periods/{YYYY-MM}/reopen` +/// +/// Removes the closed marker. Future `Usage` events for the period are +/// accepted again. +async fn handle_reopen_period( + State(state): State, + Path((account_id, period)): Path<(String, String)>, +) -> Result, AppError> { + use crate::period::parse_period; + + let (year, month) = parse_period(&period) + .map_err(|e| AppError(anyhow::anyhow!("invalid period: {}", e)))?; + + let mut manifest = state.manifest.write().await; + let before = manifest.closed_periods.len(); + manifest.closed_periods.retain(|p| { + !(p.account_id == account_id && p.year == year && p.month == month) + }); + let removed = before - manifest.closed_periods.len(); + if removed > 0 { + manifest + .save(&state.config.db_root) + .map_err(|e| AppError(anyhow::anyhow!("manifest save: {}", e)))?; + } + + Ok(Json(serde_json::json!({ + "account_id": account_id, + "period": period, + "state": "Open", + "removed": removed > 0, + }))) +} + +/// `GET /v1/accounts/{account_id}/periods/{YYYY-MM}` +/// +/// Returns the period's current state (`Open`/`Closed`), `closed_at_ms` +/// if applicable, and the current SUM(quantity) over the period. The +/// SUM is live (not a frozen snapshot) — future snapshot semantics +/// will land alongside intermediate states (Closing/Invoiced/Adjusted). +async fn handle_get_period( + State(state): State, + Path((account_id, period)): Path<(String, String)>, +) -> Result, AppError> { + use crate::period::{find_closed, parse_period}; + use crate::query::executor::execute_plan; + use crate::query::plan::{AggregationFunction, QueryPlan, QuerySource}; + + let (year, month) = parse_period(&period) + .map_err(|e| AppError(anyhow::anyhow!("invalid period: {}", e)))?; + + // Period bounds: [first day of month UTC, first day of next month UTC). + use chrono::{NaiveDate, NaiveDateTime, TimeZone, Utc}; + let from_dt = NaiveDate::from_ymd_opt(year as i32, month as u32, 1) + .and_then(|d| d.and_hms_opt(0, 0, 0)) + .ok_or_else(|| AppError(anyhow::anyhow!("invalid date for period {}", period)))?; + let next_month_dt: NaiveDateTime = if month == 12 { + NaiveDate::from_ymd_opt(year as i32 + 1, 1, 1) + } else { + NaiveDate::from_ymd_opt(year as i32, month as u32 + 1, 1) + } + .and_then(|d| d.and_hms_opt(0, 0, 0)) + .ok_or_else(|| AppError(anyhow::anyhow!("invalid date for period {}", period)))?; + let from_ms = Utc.from_utc_datetime(&from_dt).timestamp_millis(); + let to_ms = Utc.from_utc_datetime(&next_month_dt).timestamp_millis(); + + let mut metrics = HashMap::new(); + metrics.insert("quantity".to_string(), AggregationFunction::Sum); + let plan = QueryPlan { + source: QuerySource::RollupHourly, + account_id: Some(account_id.clone()), + from_ms, + to_ms, + filters: vec![], + group_by: vec![], + metrics, + limit: None, + }; + let result = execute_plan(&state, &plan).await; + let total = extract_quantity_sum(&result); + + let (manifest_state, closed_at_ms) = { + let manifest = state.manifest.read().await; + match find_closed(&manifest, &account_id, year, month) { + Some(p) => ("Closed", Some(p.closed_at_ms)), + None => ("Open", None), + } + }; + + Ok(Json(serde_json::json!({ + "account_id": account_id, + "period": period, + "state": manifest_state, + "closed_at_ms": closed_at_ms, + "from_ms": from_ms, + "to_ms": to_ms, + "total_quantity": total.to_string(), + }))) +} + /// Pull SUM(quantity) out of an executor result. Returns 0 when the /// result is empty (e.g., no events in range). fn extract_quantity_sum(result: &[serde_json::Value]) -> i128 { diff --git a/src/lib.rs b/src/lib.rs index d27e3eb..0a632f9 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -4,6 +4,7 @@ pub mod compact; pub mod export; pub mod ingest; pub mod model; +pub mod period; pub mod query; pub mod rollup; pub mod runtime; diff --git a/src/period.rs b/src/period.rs new file mode 100644 index 0000000..b8345e8 --- /dev/null +++ b/src/period.rs @@ -0,0 +1,61 @@ +//! Billing period helpers (Phase D). +//! +//! A period is a (year, month) tuple. The lifecycle is minimal — +//! `Open` (default) or `Closed`. Closed periods reject new `Usage` +//! events at ingest; `Correction` and `Retraction` events for closed +//! periods are still accepted and become post-close adjustments. +//! +//! Future extension points (left for a later PR): +//! - intermediate states (Closing / Invoiced / Adjusted) +//! - frozen snapshot semantics (closed-period queries return the +//! totals as they stood at close-time, not live) +//! - non-month periods (weekly / quarterly) + +use chrono::{DateTime, Datelike}; + +use crate::storage::manifest::{ClosedPeriod, Manifest}; + +/// Compute the (year, month) period for an event timestamp. Uses UTC, +/// matching the spec §21 simplification. +pub fn period_for_ts(ts_ms: i64) -> Option<(u16, u8)> { + let dt = DateTime::from_timestamp_millis(ts_ms)?; + Some((dt.year() as u16, dt.month() as u8)) +} + +/// True if the (account, year, month) tuple has been closed. +pub fn is_period_closed(manifest: &Manifest, account: &str, year: u16, month: u8) -> bool { + manifest + .closed_periods + .iter() + .any(|p| p.account_id == account && p.year == year && p.month == month) +} + +/// Parse the `YYYY-MM` URL/CLI param form into (year, month). +pub fn parse_period(s: &str) -> Result<(u16, u8), String> { + let (year_s, month_s) = s + .split_once('-') + .ok_or_else(|| format!("period must be YYYY-MM, got `{}`", s))?; + let year: u16 = year_s + .parse() + .map_err(|e| format!("invalid year `{}`: {}", year_s, e))?; + let month: u8 = month_s + .parse() + .map_err(|e| format!("invalid month `{}`: {}", month_s, e))?; + if !(1..=12).contains(&month) { + return Err(format!("month must be 1..=12, got {}", month)); + } + Ok((year, month)) +} + +/// Look up the matching ClosedPeriod entry, if any. +pub fn find_closed<'a>( + manifest: &'a Manifest, + account: &str, + year: u16, + month: u8, +) -> Option<&'a ClosedPeriod> { + manifest + .closed_periods + .iter() + .find(|p| p.account_id == account && p.year == year && p.month == month) +} diff --git a/src/storage/manifest.rs b/src/storage/manifest.rs index 5ad6bc3..73bb0f5 100644 --- a/src/storage/manifest.rs +++ b/src/storage/manifest.rs @@ -56,6 +56,21 @@ pub struct Watermarks { pub hourly_rollup_ms: i64, } +/// A finalized billing period for one account. New `Usage` events with a +/// timestamp inside this period are rejected at ingest. `Correction` and +/// `Retraction` events are still accepted (they become post-close +/// adjustments, per spec §13). An operator-driven `reopen_period` +/// removes the entry from `closed_periods`. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct ClosedPeriod { + pub account_id: String, + pub year: u16, + pub month: u8, + /// Wall-clock when the close happened (unix ms). Surfaced in the + /// GET endpoint so operators can correlate with their billing run. + pub closed_at_ms: i64, +} + #[derive(Debug, Clone, Serialize, Deserialize, Default)] pub struct Manifest { pub db_version: u32, @@ -70,6 +85,13 @@ pub struct Manifest { #[serde(default)] pub last_sealed_wal_id: u64, + /// Per-account billing periods that have been closed. `Usage` events + /// timestamped inside any of these are rejected at ingest. + /// `#[serde(default)]` so manifests written before period lifecycle + /// deserialize with an empty list. + #[serde(default)] + pub closed_periods: Vec, + /// Generation number — bumps on every save. Used by the manifest- /// generation recovery to walk backwards through historical versions /// when the current one is corrupt. `#[serde(skip)]` because it's diff --git a/tests/period_lifecycle.rs b/tests/period_lifecycle.rs new file mode 100644 index 0000000..586e1e8 --- /dev/null +++ b/tests/period_lifecycle.rs @@ -0,0 +1,378 @@ +//! Integration tests for the minimal period lifecycle (Phase D): +//! - POST /v1/accounts/{id}/periods/{YYYY-MM}/close → Closed state +//! - POST .../reopen → Open state +//! - GET .../{period} → state + closed_at_ms + total_quantity +//! - Ingest: `Usage` events in a closed period are rejected +//! - Ingest: `Correction` / `Retraction` events are still accepted (adjustments) +//! - Period helpers: `period_for_ts`, `parse_period`, `is_period_closed` + +use std::path::PathBuf; +use std::sync::Arc; + +use axum::body::Body; +use axum::http::{Method, Request, StatusCode}; +use http_body_util::BodyExt; +use tokio::sync::{Mutex, RwLock}; +use tower::ServiceExt; + +use usagedb::api::http::{IngestBatchRequest, IngestBatchResponse}; +use usagedb::api::http_server::build_router; +use usagedb::ingest::dedupe::HotDedupe; +use usagedb::ingest::memtable::Memtable; +use usagedb::ingest::wal::Wal; +use usagedb::model::dimensions::SmallDimensions; +use usagedb::model::event::{CorrectionRef, EventKind, UsageEvent}; +use usagedb::model::ids::{ + AccountId, EventId, MeterId, ModelId, ProductId, SourceId, SubscriptionId, Unit, +}; +use usagedb::period::{is_period_closed, parse_period, period_for_ts}; +use usagedb::runtime::config::Config; +use usagedb::runtime::state::{AppState, AppStateInner}; +use usagedb::storage::manifest::{ClosedPeriod, Manifest}; + +fn tmp_root() -> PathBuf { + let dir = tempfile::tempdir().expect("tempdir"); + let p = dir.path().to_path_buf(); + std::mem::forget(dir); + p +} + +fn build_state(db_root: PathBuf) -> AppState { + let config = Config { + db_root: db_root.clone(), + default_bucket_count: 2, + memtable_max_age_ms: i64::MAX, + ..Config::default() + }; + std::fs::create_dir_all(&config.db_root).unwrap(); + let wal = Wal::open(db_root.join("wal"), 0).unwrap(); + let manifest = Manifest { bucket_count: 2, ..Manifest::default() }; + let (flush_sender, _r) = tokio::sync::mpsc::channel(4); + Arc::new(AppStateInner { + config, + dedupe: Mutex::new(HotDedupe::new(1000)), + wal: Mutex::new(wal), + memtable: Mutex::new(Memtable::new()), + manifest: RwLock::new(manifest), + flush_sender, + }) +} + +fn make_event(id: &str, account: &str, ts: i64, qty: i128) -> UsageEvent { + UsageEvent { + event_id: EventId(id.to_string()), + kind: EventKind::Usage, + correction_ref: None, + account_id: AccountId(account.to_string()), + subscription_id: Some(SubscriptionId("sub".into())), + product_id: ProductId("prod".into()), + meter_id: MeterId("meter".into()), + timestamp_ms: ts, + quantity: qty, + unit: Unit("token".into()), + source: SourceId("test".into()), + model_id: Some(ModelId("m1".into())), + dimensions: SmallDimensions::default(), + ingested_at_ms: ts, + } +} + +/// Timestamp for the first second of a UTC (year, month, day). +fn ts_at(year: i32, month: u32, day: u32) -> i64 { + use chrono::{NaiveDate, TimeZone, Utc}; + let dt = NaiveDate::from_ymd_opt(year, month, day) + .unwrap() + .and_hms_opt(0, 0, 0) + .unwrap(); + Utc.from_utc_datetime(&dt).timestamp_millis() +} + +async fn json_post(state: AppState, uri: &str, body: serde_json::Value) -> (StatusCode, serde_json::Value) { + let app = build_router(state); + let req = Request::builder() + .method(Method::POST) + .uri(uri) + .header("content-type", "application/json") + .body(Body::from(serde_json::to_vec(&body).unwrap())) + .unwrap(); + let response = app.oneshot(req).await.unwrap(); + let status = response.status(); + let body = response.into_body().collect().await.unwrap().to_bytes(); + let v = serde_json::from_slice(&body).unwrap_or(serde_json::Value::Null); + (status, v) +} + +async fn empty_post(state: AppState, uri: &str) -> (StatusCode, serde_json::Value) { + let app = build_router(state); + let req = Request::builder() + .method(Method::POST) + .uri(uri) + .body(Body::empty()) + .unwrap(); + let response = app.oneshot(req).await.unwrap(); + let status = response.status(); + let body = response.into_body().collect().await.unwrap().to_bytes(); + let v = serde_json::from_slice(&body).unwrap_or(serde_json::Value::Null); + (status, v) +} + +async fn get_json(state: AppState, uri: &str) -> (StatusCode, serde_json::Value) { + let app = build_router(state); + let req = Request::builder().uri(uri).body(Body::empty()).unwrap(); + let response = app.oneshot(req).await.unwrap(); + let status = response.status(); + let body = response.into_body().collect().await.unwrap().to_bytes(); + let v = serde_json::from_slice(&body).unwrap_or(serde_json::Value::Null); + (status, v) +} + +// ========================================================================= +// Period helpers +// ========================================================================= + +#[test] +fn period_for_ts_returns_utc_year_month() { + // 2026-05-16T12:00:00Z + let ts = ts_at(2026, 5, 16); + assert_eq!(period_for_ts(ts), Some((2026, 5))); + // Boundary: 2026-01-01 → (2026, 1) + assert_eq!(period_for_ts(ts_at(2026, 1, 1)), Some((2026, 1))); + // Negative timestamps are pre-1970; valid per chrono but unusual. + // We don't need to assert specific behavior, just that nothing panics. + let _ = period_for_ts(-1); +} + +#[test] +fn parse_period_round_trips() { + assert_eq!(parse_period("2026-05"), Ok((2026, 5))); + assert_eq!(parse_period("2030-12"), Ok((2030, 12))); + assert!(parse_period("2026").is_err()); + assert!(parse_period("2026-13").is_err()); + assert!(parse_period("2026-0").is_err()); + assert!(parse_period("abc-de").is_err()); +} + +#[test] +fn is_period_closed_finds_match() { + let mut manifest = Manifest::default(); + assert!(!is_period_closed(&manifest, "acc", 2026, 5)); + manifest.closed_periods.push(ClosedPeriod { + account_id: "acc".into(), + year: 2026, + month: 5, + closed_at_ms: 1000, + }); + assert!(is_period_closed(&manifest, "acc", 2026, 5)); + assert!(!is_period_closed(&manifest, "acc", 2026, 6)); + assert!(!is_period_closed(&manifest, "different", 2026, 5)); +} + +// ========================================================================= +// HTTP: close / reopen / get +// ========================================================================= + +#[tokio::test] +async fn close_period_persists_in_manifest() { + let root = tmp_root(); + let state = build_state(root.clone()); + + let (status, body) = empty_post( + state.clone(), + "/v1/accounts/acc_x/periods/2026-05/close", + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body["state"], serde_json::json!("Closed")); + assert!(body["closed_at_ms"].as_i64().unwrap() > 0); + + // Verify persisted in manifest. + let manifest = state.manifest.read().await; + assert_eq!(manifest.closed_periods.len(), 1); + assert_eq!(manifest.closed_periods[0].account_id, "acc_x"); + assert_eq!(manifest.closed_periods[0].year, 2026); + assert_eq!(manifest.closed_periods[0].month, 5); +} + +#[tokio::test] +async fn close_period_is_idempotent() { + let root = tmp_root(); + let state = build_state(root.clone()); + + let _ = empty_post(state.clone(), "/v1/accounts/acc/periods/2026-05/close").await; + let (status, body) = empty_post( + state.clone(), + "/v1/accounts/acc/periods/2026-05/close", + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body["already_closed"], serde_json::json!(true)); + + // Manifest still has exactly one entry. + let m = state.manifest.read().await; + assert_eq!(m.closed_periods.len(), 1); +} + +#[tokio::test] +async fn reopen_period_removes_marker() { + let root = tmp_root(); + let state = build_state(root.clone()); + + let _ = empty_post(state.clone(), "/v1/accounts/acc/periods/2026-05/close").await; + let (status, body) = empty_post( + state.clone(), + "/v1/accounts/acc/periods/2026-05/reopen", + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body["state"], serde_json::json!("Open")); + assert_eq!(body["removed"], serde_json::json!(true)); + + let m = state.manifest.read().await; + assert!(m.closed_periods.is_empty()); +} + +#[tokio::test] +async fn get_period_returns_state_and_total() { + let root = tmp_root(); + let state = build_state(root.clone()); + + let (status, body) = get_json( + state.clone(), + "/v1/accounts/acc/periods/2026-05", + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body["state"], serde_json::json!("Open")); + assert!(body["closed_at_ms"].is_null()); + assert_eq!(body["total_quantity"], serde_json::json!("0")); + + // After closing, the GET reports Closed + the close timestamp. + let _ = empty_post(state.clone(), "/v1/accounts/acc/periods/2026-05/close").await; + let (_status, body) = get_json(state.clone(), "/v1/accounts/acc/periods/2026-05").await; + assert_eq!(body["state"], serde_json::json!("Closed")); + assert!(body["closed_at_ms"].as_i64().unwrap() > 0); +} + +#[tokio::test] +async fn close_period_rejects_invalid_format() { + let root = tmp_root(); + let state = build_state(root.clone()); + let (status, _) = + empty_post(state, "/v1/accounts/acc/periods/not-a-period/close").await; + assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); +} + +// ========================================================================= +// Ingest semantics +// ========================================================================= + +#[tokio::test] +async fn closed_period_rejects_usage_events() { + let root = tmp_root(); + let state = build_state(root.clone()); + + // Close 2026-05 for acc_x. + let _ = empty_post(state.clone(), "/v1/accounts/acc_x/periods/2026-05/close").await; + + // Ingest a Usage event in that period — should be rejected. + let in_period_ts = ts_at(2026, 5, 10); + let event = make_event("evt_blocked", "acc_x", in_period_ts, 100); + let payload = IngestBatchRequest { events: vec![event] }; + + let (status, body) = json_post( + state.clone(), + "/v1/usage/batch", + serde_json::to_value(&payload).unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK); + let resp: IngestBatchResponse = serde_json::from_value(body).unwrap(); + assert_eq!(resp.accepted, 0); + assert_eq!(resp.rejected, 1); +} + +#[tokio::test] +async fn closed_period_only_rejects_target_account() { + let root = tmp_root(); + let state = build_state(root.clone()); + + let _ = empty_post(state.clone(), "/v1/accounts/acc_x/periods/2026-05/close").await; + + // Same period but a DIFFERENT account — should be accepted. + let in_period_ts = ts_at(2026, 5, 10); + let event = make_event("evt_other", "acc_y", in_period_ts, 100); + let payload = IngestBatchRequest { events: vec![event] }; + + let (_status, body) = json_post( + state, + "/v1/usage/batch", + serde_json::to_value(&payload).unwrap(), + ) + .await; + let resp: IngestBatchResponse = serde_json::from_value(body).unwrap(); + assert_eq!(resp.accepted, 1); + assert_eq!(resp.rejected, 0); +} + +#[tokio::test] +async fn closed_period_only_rejects_in_period_timestamps() { + let root = tmp_root(); + let state = build_state(root.clone()); + + let _ = empty_post(state.clone(), "/v1/accounts/acc_x/periods/2026-05/close").await; + + // Same account, but a DIFFERENT period — should be accepted. + let ts_in_april = ts_at(2026, 4, 28); + let ts_in_june = ts_at(2026, 6, 1); + let payload = IngestBatchRequest { + events: vec![ + make_event("evt_apr", "acc_x", ts_in_april, 1), + make_event("evt_jun", "acc_x", ts_in_june, 1), + ], + }; + let (_status, body) = json_post( + state, + "/v1/usage/batch", + serde_json::to_value(&payload).unwrap(), + ) + .await; + let resp: IngestBatchResponse = serde_json::from_value(body).unwrap(); + assert_eq!(resp.accepted, 2, "events outside closed period must be accepted"); + assert_eq!(resp.rejected, 0); +} + +#[tokio::test] +async fn closed_period_accepts_corrections() { + let root = tmp_root(); + let state = build_state(root.clone()); + + let _ = empty_post(state.clone(), "/v1/accounts/acc_x/periods/2026-05/close").await; + + let in_period_ts = ts_at(2026, 5, 10); + let mut correction = make_event("evt_correction", "acc_x", in_period_ts, -50); + correction.kind = EventKind::Correction; + correction.correction_ref = Some(CorrectionRef { + original_event_id: EventId("evt_orig".into()), + reason: "overcount".into(), + }); + let mut retraction = make_event("evt_retraction", "acc_x", in_period_ts, 0); + retraction.kind = EventKind::Retraction; + retraction.correction_ref = Some(CorrectionRef { + original_event_id: EventId("evt_orig".into()), + reason: "retract".into(), + }); + + let payload = IngestBatchRequest { events: vec![correction, retraction] }; + let (_status, body) = json_post( + state, + "/v1/usage/batch", + serde_json::to_value(&payload).unwrap(), + ) + .await; + let resp: IngestBatchResponse = serde_json::from_value(body).unwrap(); + assert_eq!( + resp.accepted, 2, + "Correction + Retraction in closed period must be accepted as adjustments" + ); + assert_eq!(resp.rejected, 0); +}