diff --git a/src/api/http_server.rs b/src/api/http_server.rs index 81ac9e2..f6060bb 100644 --- a/src/api/http_server.rs +++ b/src/api/http_server.rs @@ -696,10 +696,13 @@ 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. +/// Marks the (account, year, month) period as closed and captures a +/// frozen snapshot of its rollup totals + watermark at the moment of +/// closing. Subsequent `Usage` events inside the period are rejected +/// at ingest; `Correction`/`Retraction` events are still accepted and +/// surfaced via `get_period` as `pending_adjustments`. The frozen +/// numbers are what `get_period` will report from now on — that's the +/// whole point of closing a period for invoicing. async fn handle_close_period( State(state): State, Path((account_id, period)): Path<(String, String)>, @@ -710,7 +713,32 @@ async fn handle_close_period( let (year, month) = parse_period(&period) .map_err(|e| AppError(anyhow::anyhow!("invalid period: {}", e)))?; + // Idempotency: if already closed, just report it (and keep the + // original snapshot — re-closing shouldn't refresh it). + { + let manifest = state.manifest.read().await; + if let Some(existing) = crate::period::find_closed(&manifest, &account_id, year, month) { + return Ok(Json(serde_json::json!({ + "account_id": account_id, + "period": period, + "state": "Closed", + "already_closed": true, + "closed_at_ms": existing.closed_at_ms, + "frozen": frozen_json(existing), + }))); + } + } + + // Snapshot the rollup totals BEFORE we add the closed marker. The + // marker would block new Usage events but doesn't affect what's + // already in segments / rollups, so the snapshot reflects the + // state at close-time exactly. + let (from_ms, to_ms) = period_bounds(year, month)?; + let (frozen_quantity, frozen_event_count) = + snapshot_period_totals(&state, &account_id, from_ms, to_ms).await; + let mut manifest = state.manifest.write().await; + // Re-check inside the write lock — another caller might have raced us. if manifest.closed_periods.iter().any(|p| { p.account_id == account_id && p.year == year && p.month == month }) { @@ -722,12 +750,17 @@ async fn handle_close_period( }))); } let closed_at_ms = now_ms(); - manifest.closed_periods.push(ClosedPeriod { + let watermark_at_close_ms = manifest.watermarks.hourly_rollup_ms; + let entry = ClosedPeriod { account_id: account_id.clone(), year, month, closed_at_ms, - }); + frozen_quantity: Some(frozen_quantity), + frozen_event_count: Some(frozen_event_count), + watermark_at_close_ms: Some(watermark_at_close_ms), + }; + manifest.closed_periods.push(entry.clone()); manifest .save(&state.config.db_root) .map_err(|e| AppError(anyhow::anyhow!("manifest save: {}", e)))?; @@ -737,9 +770,82 @@ async fn handle_close_period( "period": period, "state": "Closed", "closed_at_ms": closed_at_ms, + "watermark_at_close_ms": watermark_at_close_ms, + "frozen": frozen_json(&entry), }))) } +/// Compute period bounds `[from_ms, to_ms)` from (year, month) using UTC. +fn period_bounds(year: u16, month: u8) -> Result<(i64, i64), AppError> { + use chrono::{NaiveDate, 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 year/month")))?; + let next = 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 year/month for end")))?; + Ok(( + Utc.from_utc_datetime(&from_dt).timestamp_millis(), + Utc.from_utc_datetime(&next).timestamp_millis(), + )) +} + +/// Run the rollup-source query for SUM(quantity) and COUNT over the +/// period and return (sum, count). Used by `close_period` to capture +/// the frozen snapshot. +async fn snapshot_period_totals( + state: &AppState, + account_id: &str, + from_ms: i64, + to_ms: i64, +) -> (i128, u64) { + use crate::query::executor::execute_plan; + use crate::query::plan::{AggregationFunction, QueryPlan, QuerySource}; + + let mut metrics = HashMap::new(); + metrics.insert("quantity".to_string(), AggregationFunction::Sum); + metrics.insert("count".to_string(), AggregationFunction::Count); + let plan = QueryPlan { + source: QuerySource::RollupHourly, + account_id: Some(account_id.to_string()), + from_ms, + to_ms, + filters: vec![], + group_by: vec![], + metrics, + limit: None, + }; + let result = execute_plan(state, &plan).await; + let sum = result + .iter() + .filter_map(|v| v.get("quantity")) + .filter_map(|v| v.as_str()) + .filter_map(|s| s.parse().ok()) + .next() + .unwrap_or(0); + let count = result + .iter() + .filter_map(|v| v.get("count")) + .filter_map(|v| v.as_u64()) + .next() + .unwrap_or(0); + (sum, count) +} + +fn frozen_json(entry: &crate::storage::manifest::ClosedPeriod) -> serde_json::Value { + match (entry.frozen_quantity, entry.frozen_event_count) { + (Some(q), Some(c)) => serde_json::json!({ + "quantity": q.to_string(), + "event_count": c, + }), + _ => serde_json::Value::Null, + } +} + /// `POST /v1/accounts/{account_id}/periods/{YYYY-MM}/reopen` /// /// Removes the closed marker. Future `Usage` events for the period are @@ -775,68 +881,149 @@ async fn handle_reopen_period( /// `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). +/// Returns the period's current state. +/// +/// **Open periods** return a live total — the current SUM(quantity) +/// over the period from rollups + memtable. +/// +/// **Closed periods** return the *frozen* snapshot captured at +/// close-time (`frozen.quantity`, `frozen.event_count`) plus a list +/// of `pending_adjustments` — Correction / Retraction events that +/// landed in the period after the close. `net_total = frozen + +/// pending_adjustments.quantity` is the value an invoice would +/// show today. Legacy entries closed before snapshot support fall +/// back to a live total with a `warning` field. 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}; + use crate::query::plan::{AggregationFunction, QueryFilter, QueryPlan, QuerySource}; let (year, month) = parse_period(&period) .map_err(|e| AppError(anyhow::anyhow!("invalid period: {}", e)))?; + let (from_ms, to_ms) = period_bounds(year, month)?; - // 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), + // Snapshot the closed-period entry (cloned) so the response build + // doesn't hold the lock across the adjustments query. + let closed_entry = state + .manifest + .read() + .await + .closed_periods + .iter() + .find(|p| { + p.account_id == account_id && p.year == year && p.month == month + }) + .cloned(); + + match closed_entry { + None => { + // Open: live total via rollup (with raw fallback for the + // open-period tail above the watermark). + 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); + Ok(Json(serde_json::json!({ + "account_id": account_id, + "period": period, + "state": "Open", + "closed_at_ms": serde_json::Value::Null, + "from_ms": from_ms, + "to_ms": to_ms, + "total_quantity": total.to_string(), + }))) } - }; - - 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(), - }))) + Some(entry) => { + // Closed: return the frozen snapshot + adjustments. + // Pending adjustments are Correction/Retraction events that + // landed in the period (raw scan, no aggregation). + let plan_adjustments = QueryPlan { + source: QuerySource::RawEvents, + account_id: Some(account_id.clone()), + from_ms, + to_ms, + filters: vec![QueryFilter { + field: "kind".into(), + values: vec!["Correction".into(), "Retraction".into()], + }], + group_by: vec![], + metrics: HashMap::new(), + limit: None, + }; + let adjustments = execute_plan(&state, &plan_adjustments).await; + let adjustments_quantity: i128 = adjustments + .iter() + .filter_map(|e| e.get("quantity")) + .filter_map(|v| v.as_i64()) + .map(|v| v as i128) + .sum(); + // The executor's QueryFilter is permissive about value type; + // event quantity is i128 serialized as a number by serde. + // For tests we read it as i64 which is fine for the values + // we use; production with larger quantities may need a + // tighter helper. + let _ = find_closed; // keep symbol importable + + // Build the response. + let frozen = frozen_json(&entry); + let net_total: Option = entry + .frozen_quantity + .map(|q| q.saturating_add(adjustments_quantity)); + + let mut body = serde_json::json!({ + "account_id": account_id, + "period": period, + "state": "Closed", + "closed_at_ms": entry.closed_at_ms, + "watermark_at_close_ms": entry.watermark_at_close_ms, + "from_ms": from_ms, + "to_ms": to_ms, + "frozen": frozen, + "pending_adjustments": adjustments, + "adjustments_quantity": adjustments_quantity.to_string(), + }); + if let Some(nt) = net_total { + body["net_total"] = serde_json::Value::String(nt.to_string()); + } + // Legacy entries closed before snapshot support: surface a + // warning and a live fallback so callers aren't confused + // by `frozen = null`. + if entry.frozen_quantity.is_none() { + 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 live = extract_quantity_sum(&execute_plan(&state, &plan).await); + body["live_total_fallback"] = serde_json::Value::String(live.to_string()); + body["warning"] = serde_json::Value::String( + "legacy ClosedPeriod (no snapshot at close time); live total returned for reference" + .into(), + ); + } + Ok(Json(body)) + } + } } /// Pull SUM(quantity) out of an executor result. Returns 0 when the diff --git a/src/storage/manifest.rs b/src/storage/manifest.rs index 73bb0f5..6a2b173 100644 --- a/src/storage/manifest.rs +++ b/src/storage/manifest.rs @@ -61,6 +61,18 @@ pub struct Watermarks { /// `Retraction` events are still accepted (they become post-close /// adjustments, per spec §13). An operator-driven `reopen_period` /// removes the entry from `closed_periods`. +/// +/// **Frozen snapshot.** At close time the rollup totals are captured +/// (`frozen_quantity`, `frozen_event_count`) along with the rollup +/// watermark (`watermark_at_close_ms`). Future queries on a closed +/// period return these values, NOT live data — that's the point of +/// "closing" a period for invoicing. Post-close adjustments are +/// queried separately as `pending_adjustments` (Correction + +/// Retraction events whose timestamp lands in the closed period). +/// +/// The frozen fields are `Option` for backward compat: a `ClosedPeriod` +/// written by an older build has no snapshot, and queries fall back to +/// the live total with a warning. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct ClosedPeriod { pub account_id: String, @@ -69,6 +81,16 @@ pub struct ClosedPeriod { /// 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, + /// Frozen SUM(quantity) for the period at close time. + #[serde(default)] + pub frozen_quantity: Option, + /// Frozen event count for the period at close time. + #[serde(default)] + pub frozen_event_count: Option, + /// Rollup watermark at close time. Lets operators reconcile with + /// the state the rollup system saw. + #[serde(default)] + pub watermark_at_close_ms: Option, } #[derive(Debug, Clone, Serialize, Deserialize, Default)] diff --git a/tests/period_lifecycle.rs b/tests/period_lifecycle.rs index 586e1e8..0e8d672 100644 --- a/tests/period_lifecycle.rs +++ b/tests/period_lifecycle.rs @@ -161,6 +161,9 @@ fn is_period_closed_finds_match() { year: 2026, month: 5, closed_at_ms: 1000, + frozen_quantity: None, + frozen_event_count: None, + watermark_at_close_ms: None, }); assert!(is_period_closed(&manifest, "acc", 2026, 5)); assert!(!is_period_closed(&manifest, "acc", 2026, 6)); @@ -341,6 +344,298 @@ async fn closed_period_only_rejects_in_period_timestamps() { assert_eq!(resp.rejected, 0); } +// ========================================================================= +// Frozen snapshot semantics +// ========================================================================= + +#[tokio::test] +async fn close_captures_frozen_snapshot() { + let root = tmp_root(); + let state = build_state(root.clone()); + + // Ingest two Usage events in 2026-05. + let payload = IngestBatchRequest { + events: vec![ + make_event("e1", "acc_s", ts_at(2026, 5, 10), 50), + make_event("e2", "acc_s", ts_at(2026, 5, 20), 75), + ], + }; + let _ = json_post(state.clone(), "/v1/usage/batch", serde_json::to_value(&payload).unwrap()).await; + + // Force a flush + rollup so the rollup query at close time has + // committed data to read. We drive the rollup worker directly. + use usagedb::rollup::worker::RollupWorker; + let memtable_events = { + let mut m = state.memtable.lock().await; + m.drain_all() + }; + // Write a segment for the events directly so the rollup worker has + // segments to scan. + use usagedb::ingest::flusher::build_segment_meta; + use usagedb::model::ids::bucket_for_account; + use usagedb::storage::segment_writer::RawSegmentWriter; + let bucket = bucket_for_account(&AccountId("acc_s".into()), 2); + let segment_id = format!("raw_{}", uuid::Uuid::new_v4().simple()); + let path = state.config.db_root.join(format!("{}.seg", segment_id)); + let mut w = RawSegmentWriter::new(path).unwrap(); + for e in &memtable_events { w.write_event(e).unwrap(); } + let (_rows, checksum) = w.finish().unwrap(); + let meta = build_segment_meta(&segment_id, &memtable_events, bucket, checksum); + { + let mut m = state.manifest.write().await; + m.raw_segments.push(meta); + m.save(&state.config.db_root).unwrap(); + } + let worker = RollupWorker::new( + state.clone(), + 0, + std::time::Duration::from_secs(30), + i64::MAX, + ); + worker + .tick(ts_at(2026, 6, 1) + 60_000) + .await + .unwrap(); + + // Close the period — should snapshot 50 + 75 = 125 with 2 events. + let (status, body) = + empty_post(state.clone(), "/v1/accounts/acc_s/periods/2026-05/close").await; + assert_eq!(status, StatusCode::OK); + let frozen = &body["frozen"]; + assert!(!frozen.is_null(), "close response must include frozen snapshot"); + assert_eq!(frozen["quantity"], serde_json::json!("125")); + assert_eq!(frozen["event_count"], serde_json::json!(2)); +} + +#[tokio::test] +async fn get_period_returns_frozen_snapshot_after_close() { + let root = tmp_root(); + let state = build_state(root.clone()); + + // Same setup as the snapshot test. + let payload = IngestBatchRequest { + events: vec![make_event("e1", "acc_g", ts_at(2026, 5, 10), 200)], + }; + let _ = json_post(state.clone(), "/v1/usage/batch", serde_json::to_value(&payload).unwrap()).await; + + use usagedb::ingest::flusher::build_segment_meta; + use usagedb::model::ids::bucket_for_account; + use usagedb::rollup::worker::RollupWorker; + use usagedb::storage::segment_writer::RawSegmentWriter; + let events = state.memtable.lock().await.drain_all(); + let bucket = bucket_for_account(&AccountId("acc_g".into()), 2); + let segment_id = format!("raw_{}", uuid::Uuid::new_v4().simple()); + let path = state.config.db_root.join(format!("{}.seg", segment_id)); + let mut w = RawSegmentWriter::new(path).unwrap(); + for e in &events { w.write_event(e).unwrap(); } + let (_rows, checksum) = w.finish().unwrap(); + let meta = build_segment_meta(&segment_id, &events, bucket, checksum); + { + let mut m = state.manifest.write().await; + m.raw_segments.push(meta); + m.save(&state.config.db_root).unwrap(); + } + let worker = RollupWorker::new(state.clone(), 0, std::time::Duration::from_secs(30), i64::MAX); + worker.tick(ts_at(2026, 6, 1) + 60_000).await.unwrap(); + + // Close. + let _ = empty_post(state.clone(), "/v1/accounts/acc_g/periods/2026-05/close").await; + + // GET should return the frozen snapshot, not a live total. + let (_status, body) = get_json(state.clone(), "/v1/accounts/acc_g/periods/2026-05").await; + assert_eq!(body["state"], serde_json::json!("Closed")); + assert_eq!(body["frozen"]["quantity"], serde_json::json!("200")); + assert_eq!(body["frozen"]["event_count"], serde_json::json!(1)); + assert_eq!( + body["adjustments_quantity"], + serde_json::json!("0"), + "no adjustments yet" + ); + assert_eq!(body["net_total"], serde_json::json!("200")); +} + +#[tokio::test] +async fn correction_after_close_surfaces_as_adjustment_not_frozen() { + let root = tmp_root(); + let state = build_state(root.clone()); + + // Setup: 1 event of 100, rolled up, then close. + let payload = IngestBatchRequest { + events: vec![make_event("orig", "acc_c", ts_at(2026, 5, 10), 100)], + }; + let _ = json_post(state.clone(), "/v1/usage/batch", serde_json::to_value(&payload).unwrap()).await; + + use usagedb::ingest::flusher::build_segment_meta; + use usagedb::model::ids::bucket_for_account; + use usagedb::rollup::worker::RollupWorker; + use usagedb::storage::segment_writer::RawSegmentWriter; + let events = state.memtable.lock().await.drain_all(); + let bucket = bucket_for_account(&AccountId("acc_c".into()), 2); + let segment_id = format!("raw_{}", uuid::Uuid::new_v4().simple()); + let path = state.config.db_root.join(format!("{}.seg", segment_id)); + let mut w = RawSegmentWriter::new(path).unwrap(); + for e in &events { w.write_event(e).unwrap(); } + let (_rows, checksum) = w.finish().unwrap(); + let meta = build_segment_meta(&segment_id, &events, bucket, checksum); + { + let mut m = state.manifest.write().await; + m.raw_segments.push(meta); + m.save(&state.config.db_root).unwrap(); + } + let worker = RollupWorker::new(state.clone(), 0, std::time::Duration::from_secs(30), i64::MAX); + worker.tick(ts_at(2026, 6, 1) + 60_000).await.unwrap(); + let _ = empty_post(state.clone(), "/v1/accounts/acc_c/periods/2026-05/close").await; + + // Now land a Correction event in the closed period. + let mut correction = make_event("corr", "acc_c", ts_at(2026, 5, 15), -40); + correction.kind = EventKind::Correction; + correction.correction_ref = Some(CorrectionRef { + original_event_id: EventId("orig".into()), + reason: "overcount".into(), + }); + let payload = IngestBatchRequest { events: vec![correction] }; + let (_status, body) = json_post( + state.clone(), + "/v1/usage/batch", + serde_json::to_value(&payload).unwrap(), + ) + .await; + let resp: IngestBatchResponse = serde_json::from_value(body).unwrap(); + assert_eq!(resp.accepted, 1, "correction should be accepted in closed period"); + + // GET should show: frozen still 100, pending_adjustments has the + // correction, net_total = 60. + let (_status, body) = get_json(state.clone(), "/v1/accounts/acc_c/periods/2026-05").await; + assert_eq!( + body["frozen"]["quantity"], + serde_json::json!("100"), + "frozen snapshot must NOT change after a post-close correction" + ); + let adjustments = body["pending_adjustments"].as_array().expect("adjustments array"); + assert_eq!(adjustments.len(), 1); + assert_eq!(adjustments[0]["event_id"], serde_json::json!("corr")); + assert_eq!(body["net_total"], serde_json::json!("60")); +} + +#[tokio::test] +async fn closing_twice_keeps_original_snapshot() { + let root = tmp_root(); + let state = build_state(root.clone()); + + // Setup with one event + rollup. + let payload = IngestBatchRequest { + events: vec![make_event("e1", "acc_t", ts_at(2026, 5, 10), 100)], + }; + let _ = json_post(state.clone(), "/v1/usage/batch", serde_json::to_value(&payload).unwrap()).await; + + use usagedb::ingest::flusher::build_segment_meta; + use usagedb::model::ids::bucket_for_account; + use usagedb::rollup::worker::RollupWorker; + use usagedb::storage::segment_writer::RawSegmentWriter; + let events = state.memtable.lock().await.drain_all(); + let bucket = bucket_for_account(&AccountId("acc_t".into()), 2); + let segment_id = format!("raw_{}", uuid::Uuid::new_v4().simple()); + let path = state.config.db_root.join(format!("{}.seg", segment_id)); + let mut w = RawSegmentWriter::new(path).unwrap(); + for e in &events { w.write_event(e).unwrap(); } + let (_rows, checksum) = w.finish().unwrap(); + let meta = build_segment_meta(&segment_id, &events, bucket, checksum); + { + let mut m = state.manifest.write().await; + m.raw_segments.push(meta); + m.save(&state.config.db_root).unwrap(); + } + let worker = RollupWorker::new(state.clone(), 0, std::time::Duration::from_secs(30), i64::MAX); + worker.tick(ts_at(2026, 6, 1) + 60_000).await.unwrap(); + + let (_, first) = empty_post(state.clone(), "/v1/accounts/acc_t/periods/2026-05/close").await; + let first_closed_at = first["closed_at_ms"].as_i64().unwrap(); + let (_, second) = empty_post(state.clone(), "/v1/accounts/acc_t/periods/2026-05/close").await; + assert_eq!(second["already_closed"], serde_json::json!(true)); + assert_eq!( + second["closed_at_ms"].as_i64().unwrap(), + first_closed_at, + "re-close must NOT overwrite the original closed_at_ms" + ); + assert_eq!( + second["frozen"]["quantity"], + serde_json::json!("100"), + "re-close must return the original snapshot value" + ); +} + +#[tokio::test] +async fn reopen_then_reclose_takes_a_fresh_snapshot() { + let root = tmp_root(); + let state = build_state(root.clone()); + + // 1 event of 100. + let payload = IngestBatchRequest { + events: vec![make_event("e1", "acc_rr", ts_at(2026, 5, 10), 100)], + }; + let _ = json_post(state.clone(), "/v1/usage/batch", serde_json::to_value(&payload).unwrap()).await; + + use usagedb::ingest::flusher::build_segment_meta; + use usagedb::model::ids::bucket_for_account; + use usagedb::rollup::worker::RollupWorker; + use usagedb::storage::segment_writer::RawSegmentWriter; + let events = state.memtable.lock().await.drain_all(); + let bucket = bucket_for_account(&AccountId("acc_rr".into()), 2); + let segment_id = format!("raw_{}", uuid::Uuid::new_v4().simple()); + let path = state.config.db_root.join(format!("{}.seg", segment_id)); + let mut w = RawSegmentWriter::new(path).unwrap(); + for e in &events { w.write_event(e).unwrap(); } + let (_rows, checksum) = w.finish().unwrap(); + let meta = build_segment_meta(&segment_id, &events, bucket, checksum); + { + let mut m = state.manifest.write().await; + m.raw_segments.push(meta); + m.save(&state.config.db_root).unwrap(); + } + let worker = RollupWorker::new(state.clone(), 0, std::time::Duration::from_secs(30), i64::MAX); + worker.tick(ts_at(2026, 6, 1) + 60_000).await.unwrap(); + + // First close: snapshot is 100. + let (_, first) = empty_post(state.clone(), "/v1/accounts/acc_rr/periods/2026-05/close").await; + assert_eq!(first["frozen"]["quantity"], serde_json::json!("100")); + + // Reopen. + let _ = empty_post(state.clone(), "/v1/accounts/acc_rr/periods/2026-05/reopen").await; + + // Ingest more. + let payload = IngestBatchRequest { + events: vec![make_event("e2", "acc_rr", ts_at(2026, 5, 20), 50)], + }; + let _ = json_post(state.clone(), "/v1/usage/batch", serde_json::to_value(&payload).unwrap()).await; + + // Flush + rollup again. + let events2 = state.memtable.lock().await.drain_all(); + let segment_id2 = format!("raw_{}", uuid::Uuid::new_v4().simple()); + let path2 = state.config.db_root.join(format!("{}.seg", segment_id2)); + let mut w2 = RawSegmentWriter::new(path2).unwrap(); + for e in &events2 { w2.write_event(e).unwrap(); } + let (_rows, checksum2) = w2.finish().unwrap(); + let meta2 = build_segment_meta(&segment_id2, &events2, bucket, checksum2); + { + let mut m = state.manifest.write().await; + m.raw_segments.push(meta2); + // Force the rollup worker to redo hour 10 by rebuilding rollups. + m.rollup_segments.retain(|s| s.bucket != bucket || s.min_timestamp_ms != ts_at(2026, 5, 10) / 3_600_000 * 3_600_000); + m.watermarks.hourly_rollup_ms = 0; + m.save(&state.config.db_root).unwrap(); + } + let worker2 = RollupWorker::new(state.clone(), 0, std::time::Duration::from_secs(30), i64::MAX); + worker2.tick(ts_at(2026, 6, 1) + 60_000).await.unwrap(); + + // Close again — snapshot should now be 150. + let (_, second) = empty_post(state.clone(), "/v1/accounts/acc_rr/periods/2026-05/close").await; + assert_eq!( + second["frozen"]["quantity"], + serde_json::json!("150"), + "reopen + reclose should capture a fresh snapshot" + ); +} + #[tokio::test] async fn closed_period_accepts_corrections() { let root = tmp_root();