feat: add selectable routing groups and composite billing

Support per-model provider enablement and compact model editing. Capture request-time billing factors, charge customer costs separately, and preserve historical statistics without backfills.
This commit is contained in:
elky
2026-10-07 14:49:57 +08:00
parent 310098a853
commit 911c7f8875
110 changed files with 6524 additions and 559 deletions
@@ -8,9 +8,11 @@ use aether_data_contracts::repository::usage::{
sanitize_usage_request_metadata_object as project_usage_request_metadata_object,
sanitize_usage_request_metadata_ref as project_usage_request_metadata_ref,
usage_body_capture_is_authoritative, UsageBodyCaptureState,
PROVIDER_ACTUAL_SERVICE_TIER_METADATA_KEY, PROVIDER_CACHE_TTL_MINUTES_METADATA_KEY,
PROVIDER_REASONING_EFFORT_METADATA_KEY, PROVIDER_RESPONSE_MODEL_METADATA_KEY,
PROVIDER_SERVICE_TIER_METADATA_KEY, REQUESTED_REASONING_EFFORT_METADATA_KEY,
BILLING_MULTIPLIER_SNAPSHOT_METADATA_KEY, PROVIDER_ACTUAL_SERVICE_TIER_METADATA_KEY,
PROVIDER_CACHE_TTL_MINUTES_METADATA_KEY, PROVIDER_REASONING_EFFORT_METADATA_KEY,
PROVIDER_RESPONSE_MODEL_METADATA_KEY, PROVIDER_SERVICE_TIER_METADATA_KEY,
REQUESTED_REASONING_EFFORT_METADATA_KEY, ROUTING_GROUP_BILLING_MULTIPLIER_METADATA_KEY,
ROUTING_GROUP_ID_METADATA_KEY, ROUTING_GROUP_NAME_METADATA_KEY,
};
use serde_json::{Map, Value};
@@ -112,6 +114,10 @@ pub(crate) fn retain_first_byte_request_metadata(value: Option<Value>) -> Option
| "model_id"
| "global_model_id"
| "global_model_name"
| ROUTING_GROUP_BILLING_MULTIPLIER_METADATA_KEY
| BILLING_MULTIPLIER_SNAPSHOT_METADATA_KEY
| ROUTING_GROUP_ID_METADATA_KEY
| ROUTING_GROUP_NAME_METADATA_KEY
)
});
(!metadata.is_empty()).then_some(Value::Object(metadata))
@@ -542,6 +548,10 @@ mod tests {
"request_path": "/v1/chat/completions",
"upstream_is_stream": true,
"proxy": {"mode": "manual", "node_id": "proxy-1"},
"routing_group_billing_multiplier": 0.25,
"billing_multiplier_snapshot": {"version": 1, "factors": {"routing_group": 0.25, "user_group": 2.0}, "multiplier": 0.5},
"routing_group_id": "group-1",
"routing_group_name": "默认调度策略",
"billing_snapshot": {"dimensions": [1, 2, 3]},
"settlement_snapshot": {"status": "pending"},
"stage_timings_ms": {"planning": 12}
@@ -555,7 +565,11 @@ mod tests {
"client_ip": "203.0.113.8",
"request_path": "/v1/chat/completions",
"request_path_and_query": "/v1/chat/completions",
"upstream_is_stream": true
"upstream_is_stream": true,
"routing_group_billing_multiplier": 0.25,
"billing_multiplier_snapshot": {"version": 1, "factors": {"routing_group": 0.25, "user_group": 2.0}, "multiplier": 0.5},
"routing_group_id": "group-1",
"routing_group_name": "默认调度策略"
})
);
}
@@ -644,6 +658,52 @@ mod tests {
.is_none());
}
#[test]
fn routing_group_snapshot_survives_seed_and_sanitization() {
for multiplier in [0.0, 0.25, 1.0, 2.5] {
let context = json!({
"routing_group_billing_multiplier": multiplier,
"billing_multiplier_snapshot": {"version": 1, "factors": {"routing_group": multiplier, "user_group": 2.0}, "multiplier": multiplier * 2.0},
"routing_group_id": "group-1",
"routing_group_name": "请求时的分组",
"rate_multiplier": 0.75,
"routing_trace": {"untrusted": true}
});
let metadata = build_usage_request_metadata_seed(&sample_plan(), context.as_object())
.expect("group snapshot should survive projection");
assert_eq!(metadata["routing_group_billing_multiplier"], multiplier);
assert_eq!(
metadata["billing_multiplier_snapshot"],
context["billing_multiplier_snapshot"]
);
assert_eq!(metadata["routing_group_id"], "group-1");
assert_eq!(metadata["routing_group_name"], "请求时的分组");
assert_eq!(metadata["rate_multiplier"], 0.75);
assert!(metadata.get("routing_trace").is_none());
assert_eq!(
sanitize_usage_request_metadata(Some(metadata.clone())),
Some(metadata)
);
}
for multiplier in [json!(-1), json!("Infinity"), json!(f64::NAN)] {
let context = json!({"routing_group_billing_multiplier": multiplier});
let metadata = build_usage_request_metadata_seed(&sample_plan(), context.as_object())
.expect(
"invalid pricing must retain a marker instead of falling back to legacy billing",
);
assert_eq!(
metadata.get("billing_multiplier_snapshot"),
Some(&Value::Null)
);
assert!(
aether_data_contracts::repository::usage::billing_multiplier_snapshot(Some(
&metadata
))
.is_err()
);
}
}
#[test]
fn builds_seed_from_context_and_allowlisted_metadata_only() {
let metadata = build_usage_request_metadata_seed(
+2 -21
View File
@@ -26,9 +26,7 @@ use crate::request_metadata::{
request_body_derived_facts_action, retain_first_byte_request_metadata,
RequestBodyDerivedFactsAction,
};
use crate::settlement::{
reconcile_usage_policy_cost_for_event_with_result, settle_usage_with_reconciled_cost,
};
use crate::settlement::settle_usage_after_upsert;
use crate::shutdown::{UsageBackgroundTasks, UsageShutdownState};
use crate::worker::{
build_usage_queue_worker_with_record_gate, UsageWorkerControl, UsageWorkerObservation,
@@ -5156,21 +5154,6 @@ impl UsageRuntime {
where
T: UsageRuntimeAccess,
{
let reconciled = match reconcile_usage_policy_cost_for_event_with_result(data, event).await
{
Ok(reconciled) => reconciled,
Err(err) => {
warn!(
event_name = "usage_event_cost_reconciliation_failed",
log_type = "event",
usage_event_type = ?event.event_type,
request_id = %event.request_id,
error = %err,
"usage runtime failed to reconcile plan cost before direct usage upsert"
);
return false;
}
};
match build_upsert_usage_record_from_event(event) {
Ok(record) => match catch_usage_writer_panic(
"direct usage upsert",
@@ -5179,9 +5162,7 @@ impl UsageRuntime {
.await
{
Ok(Some(stored)) => {
if let Err(err) =
settle_usage_with_reconciled_cost(data, &stored, reconciled).await
{
if let Err(err) = settle_usage_after_upsert(data, &stored, event).await {
warn!(
event_name = "usage_terminal_settlement_failed",
log_type = "event",
+113 -14
View File
@@ -7,7 +7,7 @@ use aether_data_contracts::repository::settlement::{
};
use aether_data_contracts::repository::usage::PLAN_USAGE_RESERVATION_DEFERRED_METADATA_KEY;
use aether_data_contracts::repository::usage::{
cancelled_request_fee_is_billable, StoredRequestUsageAudit,
billing_multiplier_snapshot, cancelled_request_fee_is_billable, StoredRequestUsageAudit,
};
use aether_data_contracts::{DataLayerError, DataLayerError::InvalidInput};
use async_trait::async_trait;
@@ -72,16 +72,25 @@ pub(crate) async fn reconcile_usage_policy_cost_for_event_with_result(
return Ok(None);
};
let actual_cost_units = if terminal_state == UsagePolicyCostReservationState::Finalized {
let actual_cost_usd = event.data.actual_total_cost_usd.ok_or_else(|| {
let snapshot = billing_multiplier_snapshot(event.data.request_metadata.as_ref())?;
let cost = if snapshot.is_some() {
event.data.total_cost_usd
} else {
event.data.actual_total_cost_usd
};
let actual_cost_usd = cost.ok_or_else(|| {
InvalidInput(
"completed usage event with a plan reservation token is missing actual cost"
.to_string(),
)
})?;
nonnegative_usd_to_usage_policy_cost_units(finite_cost(actual_cost_usd)?.max(0.0))
.ok_or_else(|| {
InvalidInput("usage policy settlement cost exceeds the supported range".to_string())
})?
let actual_cost_usd = match snapshot {
Some(snapshot) => snapshot.cost(actual_cost_usd)?,
None => finite_cost(actual_cost_usd)?.max(0.0),
};
nonnegative_usd_to_usage_policy_cost_units(actual_cost_usd).ok_or_else(|| {
InvalidInput("usage policy settlement cost exceeds the supported range".to_string())
})?
} else {
0
};
@@ -115,6 +124,34 @@ pub async fn settle_usage_if_needed(
settle_usage_with_reconciled_cost(writer, usage, None).await
}
pub(crate) async fn settle_usage_after_upsert(
writer: &dyn UsageSettlementWriter,
usage: &StoredRequestUsageAudit,
event: &UsageEvent,
) -> Result<(), DataLayerError> {
// Different admissions can share a client request id. Finalize that event's
// own server-issued reservation without borrowing the other admission's rate.
if event_usage_policy_reservation_token(event).is_some()
&& event_usage_policy_reservation_token(event) != usage_policy_reservation_token(usage)
&& !plan_usage_reservation_reconciliation_is_deferred(event.data.request_metadata.as_ref())
{
let billable = event.event_type == UsageEventType::Completed
|| (event.event_type == UsageEventType::Cancelled
&& cancelled_request_fee_is_billable(event.data.request_metadata.as_ref()));
if billable
&& billing_multiplier_snapshot(event.data.request_metadata.as_ref())?.is_none()
&& billing_multiplier_snapshot(usage.request_metadata.as_ref())?.is_some()
{
return Err(InvalidInput(
"colliding usage admission is missing its own billing multiplier snapshot"
.to_string(),
));
}
reconcile_usage_policy_cost_for_event(writer, event).await?;
}
settle_usage_if_needed(writer, usage).await
}
pub(crate) async fn settle_usage_with_reconciled_cost(
writer: &dyn UsageSettlementWriter,
usage: &StoredRequestUsageAudit,
@@ -127,6 +164,11 @@ pub(crate) async fn settle_usage_with_reconciled_cost(
return Ok(());
}
let billing_cost_usd = match billing_multiplier_snapshot(usage.request_metadata.as_ref())? {
Some(snapshot) => snapshot.cost(usage.total_cost_usd)?,
None => finite_cost(usage.actual_total_cost_usd)?.max(0.0),
};
let finalized_at_unix_secs = usage
.finalized_at_unix_secs
.or(Some(usage.updated_at_unix_secs));
@@ -148,14 +190,14 @@ pub(crate) async fn settle_usage_with_reconciled_cost(
{
(
UsagePolicyCostReservationState::Finalized,
nonnegative_usd_to_usage_policy_cost_units(
finite_cost(usage.actual_total_cost_usd)?.max(0.0),
)
.ok_or_else(|| {
InvalidInput(
"usage policy settlement cost exceeds the supported range".to_string(),
)
})?,
nonnegative_usd_to_usage_policy_cost_units(billing_cost_usd).ok_or_else(
|| {
InvalidInput(
"usage policy settlement cost exceeds the supported range"
.to_string(),
)
},
)?,
)
} else {
(UsagePolicyCostReservationState::Released, 0)
@@ -194,6 +236,7 @@ pub(crate) async fn settle_usage_with_reconciled_cost(
billing_status: usage.billing_status.clone(),
total_cost_usd: finite_cost(usage.total_cost_usd)?,
actual_total_cost_usd: finite_cost(usage.actual_total_cost_usd)?,
billing_cost_usd: Some(billing_cost_usd),
finalized_at_unix_secs,
};
let _ = writer.settle_usage(input).await?;
@@ -449,6 +492,62 @@ mod tests {
);
}
#[tokio::test]
async fn composite_billing_rate_charges_customer_without_changing_provider_cost() {
for (group_rate, user_rate, expected_cost) in [(2.0, 0.75, 1.875), (0.0, 3.0, 0.0)] {
let writer = TestSettlementWriter {
has_writer: true,
..Default::default()
};
let mut usage = sample_usage();
let snapshot =
aether_data_contracts::repository::usage::BillingMultiplierSnapshot::from_factors(
std::collections::BTreeMap::from([
("routing_group".to_string(), group_rate),
("user_group".to_string(), user_rate),
]),
)
.unwrap();
usage.request_metadata.as_mut().unwrap()["billing_multiplier_snapshot"] =
json!(snapshot);
settle_usage_if_needed(&writer, &usage).await.unwrap();
let inputs = writer.inputs.lock().unwrap();
assert_eq!(inputs[0].billing_cost_usd, Some(expected_cost));
assert_eq!(inputs[0].total_cost_usd, 1.25);
assert_eq!(inputs[0].actual_total_cost_usd, 0.75);
assert_eq!(
writer.reconciliations.lock().unwrap()[0].actual_cost_units,
(expected_cost * 100_000_000.0).round() as u64
);
}
}
#[tokio::test]
async fn corrupt_or_overflowing_billing_rate_never_changes_wallet_or_reservation() {
for (base, snapshot) in [
(1.25, serde_json::Value::Null),
(
1.25,
json!({"version":1,"factors":{"routing_group":2},"multiplier":1}),
),
(
f64::MAX,
json!({"version":1,"factors":{"routing_group":2},"multiplier":2}),
),
] {
let writer = TestSettlementWriter {
has_writer: true,
..Default::default()
};
let mut usage = sample_usage();
usage.total_cost_usd = base;
usage.request_metadata.as_mut().unwrap()["billing_multiplier_snapshot"] = snapshot;
assert!(settle_usage_if_needed(&writer, &usage).await.is_err());
assert!(writer.inputs.lock().unwrap().is_empty());
assert!(writer.reconciliations.lock().unwrap().is_empty());
}
}
#[tokio::test]
async fn releases_pending_cancelled_usage_without_wallet_settlement() {
let writer = TestSettlementWriter {
@@ -66,6 +66,10 @@ impl UsageSettlementWriter for ReuseStore {
&self,
input: ReconcileUsagePolicyCostInput,
) -> Result<Option<StoredUsagePolicyCostReservation>, DataLayerError> {
assert!(
self.upserts.load(Ordering::Relaxed) > 0,
"a durable usage row must exist before a reservation is finalized"
);
input.validate()?;
self.reconciliations.lock().unwrap().push(input.clone());
tokio::task::yield_now().await;
@@ -185,7 +189,7 @@ async fn write(store: &ReuseStore, event: UsageEvent, direct: bool) {
}
#[tokio::test]
async fn worker_and_direct_writes_reuse_confirmed_reservation_and_still_settle_wallet() {
async fn worker_and_direct_writes_persist_before_reconciling_and_settling_wallet() {
for direct in [false, true] {
let store = ReuseStore::default();
write(&store, event(), direct).await;
@@ -198,11 +202,12 @@ async fn worker_and_direct_writes_reuse_confirmed_reservation_and_still_settle_w
assert_eq!(settlements.len(), 1);
assert_eq!(settlements[0].request_id, "req-1");
assert_eq!(settlements[0].actual_total_cost_usd, 0.75);
assert_eq!(settlements[0].billing_cost_usd, Some(0.75));
}
}
#[tokio::test]
async fn missing_or_different_reconciliation_results_keep_stored_usage_reconciliation() {
async fn missing_or_different_reconciliation_results_do_not_repeat_stored_usage_reconciliation() {
let changes: [fn(&mut StoredUsagePolicyCostReservation); 9] = [
|row| row.request_id = "other-request".to_string(),
|row| row.subject_id = "other-user".to_string(),
@@ -223,7 +228,7 @@ async fn missing_or_different_reconciliation_results_keep_stored_usage_reconcili
..Default::default()
};
write(&store, event(), direct).await;
assert_eq!(store.reconciliations.lock().unwrap().len(), 2);
assert_eq!(store.reconciliations.lock().unwrap().len(), 1);
assert_eq!(store.settlements.lock().unwrap().len(), 1);
}
}
@@ -252,19 +257,56 @@ async fn changed_stored_usage_is_reconciled_using_its_own_identity_cost_and_term
};
write(&store, event(), direct).await;
let reconciliations = store.reconciliations.lock().unwrap();
assert_eq!(reconciliations.len(), 2);
assert_eq!(reconciliations[1].request_id, stored.request_id);
let stored_token = stored.request_metadata.as_ref().unwrap()
["plan_usage_reservation_token"]
.as_str()
.unwrap();
assert_eq!(
reconciliations[1].subject_id,
reconciliations.len(),
if stored_token == RESERVATION_TOKEN {
1
} else {
2
}
);
if stored_token != RESERVATION_TOKEN {
assert_eq!(reconciliations[0].reservation_token, RESERVATION_TOKEN);
assert_eq!(reconciliations[0].actual_cost_units, 75_000_000);
}
let reconciliation = reconciliations.last().unwrap();
assert_eq!(reconciliation.request_id, stored.request_id);
assert_eq!(
reconciliation.subject_id,
stored.user_id.as_ref().unwrap().as_str()
);
assert_eq!(
reconciliations[1].reservation_token,
reconciliation.reservation_token,
stored.request_metadata.as_ref().unwrap()["plan_usage_reservation_token"]
.as_str()
.unwrap()
);
assert_ne!(reconciliations[0], reconciliations[1]);
assert_eq!(
reconciliation.actual_cost_units,
if stored.status == "failed" {
0
} else {
(stored.actual_total_cost_usd * 100_000_000.0) as u64
}
);
assert_eq!(
reconciliation.terminal_state,
if stored.status == "failed" {
UsagePolicyCostReservationState::Released
} else {
UsagePolicyCostReservationState::Finalized
}
);
assert_eq!(
reconciliation.finalized_at_unix_secs,
stored
.finalized_at_unix_secs
.unwrap_or(stored.updated_at_unix_secs)
);
let settlements = store.settlements.lock().unwrap();
assert_eq!(settlements.len(), 1);
assert_eq!(
@@ -276,6 +318,185 @@ async fn changed_stored_usage_is_reconciled_using_its_own_identity_cost_and_term
}
}
#[tokio::test]
async fn worker_and_direct_settle_customer_multiplier_snapshot_without_changing_provider_cost() {
for direct in [false, true] {
for (group_multiplier, promotion_multiplier, expected) in [
(0.0, 1.0, 0.0),
(1.0, 1.0, 1.25),
(2.0, 0.25, 0.625),
(3.0, 1.0, 3.75),
] {
let store = ReuseStore::default();
let mut event = event();
event.data.request_metadata.as_mut().unwrap()["billing_multiplier_snapshot"] = json!({
"version": 1,
"factors": {"routing_group": group_multiplier, "promotion": promotion_multiplier},
"multiplier": group_multiplier * promotion_multiplier
});
write(&store, event, direct).await;
let reconciliations = store.reconciliations.lock().unwrap();
assert_eq!(reconciliations.len(), 1, "direct={direct}");
assert_eq!(
reconciliations[0].actual_cost_units,
(expected * 100_000_000.0) as u64
);
let settlements = store.settlements.lock().unwrap();
assert_eq!(settlements.len(), 1);
assert_eq!(settlements[0].billing_cost_usd, Some(expected));
assert_eq!(settlements[0].total_cost_usd, 1.25);
assert_eq!(settlements[0].actual_total_cost_usd, 0.75);
}
}
}
#[tokio::test]
async fn sparse_terminal_event_uses_persisted_multiplier_and_cost_for_both_ledger_and_wallet() {
for direct in [false, true] {
let mut stored = sample_usage();
stored.total_cost_usd = 2.0;
stored.actual_total_cost_usd = 0.25;
stored.request_metadata.as_mut().unwrap()["billing_multiplier_snapshot"] = json!({
"version": 1,
"factors": {"routing_group": 3.0, "promotion": 0.5},
"multiplier": 1.5
});
let expected = stored.billing_cost().unwrap();
assert_eq!(expected, 3.0);
let store = ReuseStore {
stored_override: Some(stored),
..Default::default()
};
// Sparse asynchronous completion has neither the captured multiplier
// nor authoritative charges. Persistence restores the original snapshot.
let mut terminal = event();
terminal.data.total_cost_usd = None;
terminal.data.actual_total_cost_usd = None;
write(&store, terminal, direct).await;
let reconciliations = store.reconciliations.lock().unwrap();
assert_eq!(reconciliations.len(), 1, "direct={direct}");
assert_eq!(reconciliations[0].actual_cost_units, 300_000_000);
let settlements = store.settlements.lock().unwrap();
assert_eq!(settlements.len(), 1);
assert_eq!(settlements[0].billing_cost_usd, Some(expected));
assert_eq!(settlements[0].total_cost_usd, 2.0);
assert_eq!(settlements[0].actual_total_cost_usd, 0.25);
}
}
#[tokio::test]
async fn colliding_request_id_reconciles_new_token_with_its_own_snapshot_after_upsert() {
for direct in [false, true] {
let mut stored = sample_usage();
stored.total_cost_usd = 2.0;
stored.request_metadata = Some(json!({
"plan_usage_reservation_token": "previous-token",
"billing_multiplier_snapshot": {
"version": 1,
"factors": {"routing_group": 0.5},
"multiplier": 0.5
}
}));
let store = ReuseStore {
stored_override: Some(stored),
..Default::default()
};
let mut terminal = event();
terminal.data.request_metadata.as_mut().unwrap()["billing_multiplier_snapshot"] = json!({
"version": 1,
"factors": {"routing_group": 3.0},
"multiplier": 3.0
});
write(&store, terminal, direct).await;
assert_eq!(store.upserts.load(Ordering::Relaxed), 1);
let reconciliations = store.reconciliations.lock().unwrap();
assert_eq!(reconciliations.len(), 2, "direct={direct}");
assert_eq!(reconciliations[0].reservation_token, RESERVATION_TOKEN);
assert_eq!(reconciliations[0].actual_cost_units, 375_000_000);
assert_eq!(reconciliations[1].reservation_token, "previous-token");
assert_eq!(reconciliations[1].actual_cost_units, 100_000_000);
let settlements = store.settlements.lock().unwrap();
assert_eq!(settlements.len(), 1);
assert_eq!(settlements[0].billing_cost_usd, Some(1.0));
assert_eq!(settlements[0].actual_total_cost_usd, 0.75);
}
}
#[tokio::test]
async fn sparse_colliding_token_cannot_borrow_another_requests_multiplier() {
for direct in [false, true] {
for (event_type, billable_cancel) in [
(UsageEventType::Completed, false),
(UsageEventType::Cancelled, true),
] {
let mut stored = sample_usage();
stored.request_metadata = Some(json!({
"plan_usage_reservation_token": "previous-token",
"billing_multiplier_snapshot": {
"version": 1,
"factors": {"routing_group": 0.0},
"multiplier": 0.0
}
}));
let store = ReuseStore {
stored_override: Some(stored),
..Default::default()
};
let mut terminal = event();
terminal.event_type = event_type;
terminal.data.request_metadata.as_mut().unwrap()["cancelled_request_fee"] =
json!(billable_cancel);
if direct {
write(&store, terminal, true).await;
} else {
assert!(write_event_record(&store, &terminal).await.is_err());
}
assert_eq!(store.upserts.load(Ordering::Relaxed), 1);
assert!(
store.reconciliations.lock().unwrap().is_empty(),
"direct={direct}"
);
assert!(store.settlements.lock().unwrap().is_empty());
}
}
}
#[tokio::test]
async fn failed_colliding_token_is_released_without_borrowing_snapshot_or_charge() {
for direct in [false, true] {
let mut stored = sample_usage();
stored.billing_status = "settled".to_string();
stored.request_metadata = Some(json!({
"plan_usage_reservation_token": "previous-token",
"billing_multiplier_snapshot": {
"version": 1,
"factors": {"routing_group": 2.0},
"multiplier": 2.0
}
}));
let store = ReuseStore {
stored_override: Some(stored),
..Default::default()
};
let mut terminal = event();
terminal.event_type = UsageEventType::Failed;
terminal.data.total_cost_usd = None;
terminal.data.actual_total_cost_usd = None;
write(&store, terminal, direct).await;
let reconciliations = store.reconciliations.lock().unwrap();
assert_eq!(reconciliations.len(), 2, "direct={direct}");
assert_eq!(reconciliations[0].reservation_token, RESERVATION_TOKEN);
assert_eq!(
reconciliations[0].terminal_state,
UsagePolicyCostReservationState::Released
);
assert_eq!(reconciliations[0].actual_cost_units, 0);
assert_eq!(reconciliations[1].reservation_token, "previous-token");
assert_eq!(reconciliations[1].actual_cost_units, 250_000_000);
assert!(store.settlements.lock().unwrap().is_empty());
}
}
#[tokio::test]
async fn cancellation_release_billable_cancellation_and_zero_cost_preserve_settlement_rules() {
for direct in [false, true] {
@@ -329,7 +550,7 @@ async fn cancellation_release_billable_cancellation_and_zero_cost_preserve_settl
}
#[tokio::test]
async fn reconciliation_failure_stops_both_writes_before_upsert_and_wallet_settlement() {
async fn reconciliation_failure_keeps_durable_usage_but_stops_wallet_settlement() {
for direct in [false, true] {
let store = ReuseStore {
response: ReconcileResponse::Error,
@@ -341,13 +562,13 @@ async fn reconciliation_failure_stops_both_writes_before_upsert_and_wallet_settl
assert!(write_event_record(&store, &event()).await.is_err());
}
assert_eq!(store.reconciliations.lock().unwrap().len(), 1);
assert_eq!(store.upserts.load(Ordering::Relaxed), 0);
assert_eq!(store.upserts.load(Ordering::Relaxed), 1);
assert!(store.settlements.lock().unwrap().is_empty());
}
}
#[tokio::test]
async fn retry_after_upsert_failure_reconciles_again_before_settling() {
async fn retry_after_upsert_failure_reconciles_only_the_successfully_persisted_usage() {
for direct in [false, true] {
let store = ReuseStore {
fail_next_upsert: AtomicBool::new(true),
@@ -358,9 +579,10 @@ async fn retry_after_upsert_failure_reconciles_again_before_settling() {
} else {
assert!(write_event_record(&store, &event()).await.is_err());
}
assert!(store.reconciliations.lock().unwrap().is_empty());
assert!(store.settlements.lock().unwrap().is_empty());
write(&store, event(), direct).await;
assert_eq!(store.reconciliations.lock().unwrap().len(), 2);
assert_eq!(store.reconciliations.lock().unwrap().len(), 1);
assert_eq!(store.upserts.load(Ordering::Relaxed), 2);
assert_eq!(store.settlements.lock().unwrap().len(), 1);
}
@@ -481,11 +703,17 @@ async fn concurrent_duplicate_delivery_debits_real_memory_wallet_only_once() {
let store = store.clone();
let runtime = runtime.clone();
tasks.spawn(async move {
let mut terminal = event();
terminal.data.request_metadata.as_mut().unwrap()["billing_multiplier_snapshot"] = json!({
"version": 1,
"factors": {"routing_group": 3.0, "promotion": 0.5},
"multiplier": 1.5
});
if index % 2 == 0 {
write_event_record(store.as_ref(), &event()).await.unwrap();
write_event_record(store.as_ref(), &terminal).await.unwrap();
} else {
runtime
.record_terminal_event_direct(store.as_ref(), event())
.record_terminal_event_direct(store.as_ref(), terminal)
.await;
}
});
@@ -495,13 +723,25 @@ async fn concurrent_duplicate_delivery_debits_real_memory_wallet_only_once() {
}
assert_eq!(store.reconciliations.lock().unwrap().len(), 32);
assert_eq!(store.settlements.lock().unwrap().len(), 32);
assert!(store
.reconciliations
.lock()
.unwrap()
.iter()
.all(|input| input.actual_cost_units == 187_500_000));
assert!(store
.settlements
.lock()
.unwrap()
.iter()
.all(|input| input.billing_cost_usd == Some(1.875) && input.actual_total_cost_usd == 0.75));
let wallet = wallets
.find(WalletLookupKey::UserId("user-1"))
.await
.unwrap()
.unwrap();
assert_eq!(wallet.balance + wallet.gift_balance, 11.25);
assert_eq!(wallet.total_consumed, 0.75);
assert_eq!(wallet.balance + wallet.gift_balance, 10.125);
assert_eq!(wallet.total_consumed, 1.875);
assert!(matches!(
store
.repository
+66 -52
View File
@@ -16,9 +16,7 @@ use crate::queue::UsageDeadLetterOutcome;
use crate::runtime::{
UsageBillingEventEnricher, UsageRuntimeAccess, UsageWorkerRecordConcurrencyGate,
};
use crate::settlement::{
reconcile_usage_policy_cost_for_event_with_result, settle_usage_with_reconciled_cost,
};
use crate::settlement::settle_usage_after_upsert;
use crate::{
build_upsert_usage_record_from_event, UsageEvent, UsageEventType, UsageQueue,
UsageRuntimeConfig, UsageSettlementWriter,
@@ -796,10 +794,11 @@ pub async fn write_event_record<T>(data: &T, event: &UsageEvent) -> Result<(), D
where
T: UsageRecordWriter + UsageSettlementWriter + Send + Sync,
{
let reconciled = reconcile_usage_policy_cost_for_event_with_result(data, event).await?;
let record = build_upsert_usage_record_from_event(event)?;
if let Some(stored) = data.upsert_usage_record(record).await? {
settle_usage_with_reconciled_cost(data, &stored, reconciled).await?;
// Sparse terminal events may omit pricing factors. Only the stored request
// snapshot is authoritative before finalizing an immutable cost reservation.
settle_usage_after_upsert(data, &stored, event).await?;
}
// Manual proxy traffic is counted at the actual transport-attempt boundary. Usage events are
// replayable, so emitting that side effect here would count normal requests and reclaims twice.
@@ -1240,49 +1239,49 @@ mod tests {
.lock()
.expect("records lock")
.push(record.clone());
Ok(Some(
StoredRequestUsageAudit::new(
"usage-1".to_string(),
record.request_id,
record.user_id,
record.api_key_id,
record.username,
record.api_key_name,
record.provider_name,
record.model,
record.target_model,
record.provider_id,
record.provider_endpoint_id,
record.provider_api_key_id,
record.request_type,
record.api_format,
record.api_family,
record.endpoint_kind,
record.endpoint_api_format,
record.provider_api_family,
record.provider_endpoint_kind,
record.has_format_conversion.unwrap_or(false),
record.is_stream.unwrap_or(false),
record.input_tokens.unwrap_or_default() as i32,
record.output_tokens.unwrap_or_default() as i32,
record.total_tokens.unwrap_or_default() as i32,
record.total_cost_usd.unwrap_or_default(),
record.actual_total_cost_usd.unwrap_or_default(),
record.status_code.map(i32::from),
record.error_message,
record.error_category,
record.response_time_ms.map(|value| value as i32),
record.first_byte_time_ms.map(|value| value as i32),
record.status,
record.billing_status,
record
.created_at_unix_ms
.unwrap_or(record.updated_at_unix_secs) as i64,
record.updated_at_unix_secs as i64,
record.finalized_at_unix_secs.map(|value| value as i64),
)
.expect("stored usage should build"),
))
let mut stored = StoredRequestUsageAudit::new(
"usage-1".to_string(),
record.request_id,
record.user_id,
record.api_key_id,
record.username,
record.api_key_name,
record.provider_name,
record.model,
record.target_model,
record.provider_id,
record.provider_endpoint_id,
record.provider_api_key_id,
record.request_type,
record.api_format,
record.api_family,
record.endpoint_kind,
record.endpoint_api_format,
record.provider_api_family,
record.provider_endpoint_kind,
record.has_format_conversion.unwrap_or(false),
record.is_stream.unwrap_or(false),
record.input_tokens.unwrap_or_default() as i32,
record.output_tokens.unwrap_or_default() as i32,
record.total_tokens.unwrap_or_default() as i32,
record.total_cost_usd.unwrap_or_default(),
record.actual_total_cost_usd.unwrap_or_default(),
record.status_code.map(i32::from),
record.error_message,
record.error_category,
record.response_time_ms.map(|value| value as i32),
record.first_byte_time_ms.map(|value| value as i32),
record.status,
record.billing_status,
record
.created_at_unix_ms
.unwrap_or(record.updated_at_unix_secs) as i64,
record.updated_at_unix_secs as i64,
record.finalized_at_unix_secs.map(|value| value as i64),
)
.expect("stored usage should build");
stored.request_metadata = record.request_metadata;
Ok(Some(stored))
}
}
@@ -1510,7 +1509,7 @@ mod tests {
}
#[tokio::test]
async fn same_request_id_terminal_events_reconcile_each_reservation_token_before_upsert() {
async fn same_request_id_terminal_events_reconcile_each_persisted_reservation_token() {
let store = TestUsageStore::default();
let mut first = sample_event();
first.request_id = "shared-client-trace".to_string();
@@ -1619,7 +1618,12 @@ mod tests {
worker.queue.ensure_consumer_group().await.expect("group");
let mut event = sample_event();
event.data.request_metadata = Some(serde_json::json!({
"plan_usage_reservation_token": "pricing-retry-reservation"
"plan_usage_reservation_token": "pricing-retry-reservation",
"billing_multiplier_snapshot": {
"version": 1,
"factors": {"routing_group": 3.0, "promotion": 0.5},
"multiplier": 1.5
}
}));
worker.queue.enqueue(&event).await.expect("enqueue");
let batch = worker
@@ -1676,17 +1680,27 @@ mod tests {
assert_eq!(records[0].total_cost_usd, Some(0.456));
assert_eq!(records[0].actual_total_cost_usd, Some(0.123));
assert_eq!(records[0].total_tokens, Some(10));
assert_eq!(
records[0].request_metadata.as_ref().unwrap()["billing_multiplier_snapshot"]
["multiplier"],
1.5
);
}
{
let reconciliations = store.reconciliations.lock().expect("reconciliations lock");
assert_eq!(reconciliations.len(), 1);
assert_eq!(reconciliations[0].actual_cost_units, 12_300_000);
assert_eq!(reconciliations[0].actual_cost_units, 68_400_000);
assert_eq!(
reconciliations[0].reservation_token,
"pricing-retry-reservation"
);
}
assert_eq!(store.settlements.lock().expect("settlements lock").len(), 1);
{
let settlements = store.settlements.lock().expect("settlements lock");
assert_eq!(settlements.len(), 1);
assert_eq!(settlements[0].billing_cost_usd, Some(0.456 * 1.5));
assert_eq!(settlements[0].actual_total_cost_usd, 0.123);
}
assert_eq!(
store.enrich_calls.lock().expect("enrich calls lock").len(),
2