Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -59,14 +59,17 @@ 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
```

`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

Expand Down
182 changes: 182 additions & 0 deletions src/api/http_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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<crate::storage::manifest::ClosedPeriod> = {
let manifest = state.manifest.read().await;
manifest.closed_periods.clone()
};

let mut rejected = 0usize;
let mut classified: Vec<Classified> = Vec::with_capacity(payload.events.len());
for mut event in payload.events {
Expand All @@ -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 });
Expand Down Expand Up @@ -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<AppState>,
Path((account_id, period)): Path<(String, String)>,
) -> Result<Json<serde_json::Value>, 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<AppState>,
Path((account_id, period)): Path<(String, String)>,
) -> Result<Json<serde_json::Value>, 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<AppState>,
Path((account_id, period)): Path<(String, String)>,
) -> Result<Json<serde_json::Value>, 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 {
Expand Down
1 change: 1 addition & 0 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
61 changes: 61 additions & 0 deletions src/period.rs
Original file line number Diff line number Diff line change
@@ -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)
}
22 changes: 22 additions & 0 deletions src/storage/manifest.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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<ClosedPeriod>,

/// 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
Expand Down
Loading
Loading