1use std::collections::BTreeMap;
2use time::OffsetDateTime;
3
4use super::BudgetLedgerEvent;
5use crate::value_objects::{
6 BudgetAccountId, BudgetBalance, BudgetDimension, BudgetLedgerVersion, BudgetLimits,
7 BudgetMeasurement, BudgetOperationId, BudgetQuantities, BudgetReconciliationId,
8 BudgetReservation, BudgetReservationEstimate, BudgetReservationId, BudgetTokenCount,
9 CostMicros, ExecutionDuration, MeasuredBudgetQuantities, ToolCallCount,
10};
11use crate::BudgetError;
12
13#[derive(Debug, Clone, PartialEq, Eq, Default)]
14pub struct BudgetLedger {
15 account_id: Option<BudgetAccountId>,
16 limits: Option<BudgetLimits>,
17 version: BudgetLedgerVersion,
18 reservations: BTreeMap<BudgetReservationId, BudgetReservation>,
19}
20
21impl BudgetLedger {
22 #[must_use]
23 pub fn empty() -> Self {
24 Self::default()
25 }
26 pub fn rehydrate(events: &[BudgetLedgerEvent]) -> Result<Self, BudgetError> {
27 let mut ledger = Self::empty();
28 for event in events {
29 ledger.apply(event.clone())?;
30 }
31 Ok(ledger)
32 }
33 #[must_use]
34 pub const fn version(&self) -> BudgetLedgerVersion {
35 self.version
36 }
37 #[must_use]
38 pub const fn account_id(&self) -> Option<&BudgetAccountId> {
39 self.account_id.as_ref()
40 }
41 #[must_use]
42 pub const fn limits(&self) -> Option<&BudgetLimits> {
43 self.limits.as_ref()
44 }
45 pub fn reservations(&self) -> impl Iterator<Item = &BudgetReservation> {
46 self.reservations.values()
47 }
48
49 pub fn decide_open(
50 &self,
51 account_id: BudgetAccountId,
52 limits: BudgetLimits,
53 opened_at: OffsetDateTime,
54 ) -> Result<Option<BudgetLedgerEvent>, BudgetError> {
55 match (&self.account_id, &self.limits) {
56 (None, None) => Ok(Some(BudgetLedgerEvent::Opened {
57 account_id,
58 limits,
59 opened_at,
60 })),
61 (Some(stored_id), Some(stored_limits))
62 if stored_id == &account_id && stored_limits == &limits =>
63 {
64 Ok(None)
65 }
66 _ => Err(BudgetError::Persistence(crate::DomainError::Conflict {
67 what: "budget_ledger",
68 })),
69 }
70 }
71
72 pub fn decide_reserve(
73 &self,
74 operation_id: BudgetOperationId,
75 estimate: BudgetReservationEstimate,
76 reserved_at: OffsetDateTime,
77 ) -> Result<Option<BudgetLedgerEvent>, BudgetError> {
78 let account_id = self.account_id.clone().ok_or(BudgetError::LedgerNotOpen)?;
79 let reservation_id = BudgetReservationId::for_operation(&account_id, &operation_id);
80 let quantities =
81 estimate.quantities(self.limits.as_ref().ok_or(BudgetError::LedgerNotOpen)?)?;
82 if let Some(existing) = self.reservations.get(&reservation_id) {
83 return if existing.operation_id() == &operation_id
84 && existing.quantities() == quantities
85 && existing.estimate() == estimate
86 {
87 Ok(None)
88 } else {
89 Err(BudgetError::ReservationConflict(reservation_id))
90 };
91 }
92 self.ensure_capacity(quantities)?;
93 Ok(Some(BudgetLedgerEvent::Reserved {
94 account_id,
95 reservation_id,
96 operation_id,
97 quantities,
98 estimate,
99 reserved_at,
100 }))
101 }
102
103 pub fn decide_reconcile(
104 &self,
105 reservation_id: BudgetReservationId,
106 reconciliation_id: BudgetReconciliationId,
107 measured: MeasuredBudgetQuantities,
108 reconciled_at: OffsetDateTime,
109 ) -> Result<Option<BudgetLedgerEvent>, BudgetError> {
110 let account_id = self.account_id.clone().ok_or(BudgetError::LedgerNotOpen)?;
111 let reservation = self
112 .reservations
113 .get(&reservation_id)
114 .ok_or_else(|| BudgetError::ReservationNotFound(reservation_id.clone()))?;
115 if let Some((stored_id, stored)) = reservation.reconciliation() {
116 if stored_id != &reconciliation_id {
117 return Err(BudgetError::ReconciliationConflict(reconciliation_id));
118 }
119 if stored == &measured {
120 return Ok(None);
121 }
122 if !measured.advances(*stored) {
123 return Err(BudgetError::ReconciliationConflict(reconciliation_id));
124 }
125 }
126 Ok(Some(BudgetLedgerEvent::Reconciled {
127 account_id,
128 reservation_id,
129 reconciliation_id,
130 measured,
131 reconciled_at,
132 }))
133 }
134
135 pub fn apply(&mut self, event: BudgetLedgerEvent) -> Result<(), BudgetError> {
136 if let Some(account_id) = &self.account_id {
137 if event.account_id() != account_id {
138 return Err(BudgetError::Persistence(
139 crate::DomainError::InvariantViolated {
140 reason: "budget event belongs to another account",
141 },
142 ));
143 }
144 }
145 match event {
146 BudgetLedgerEvent::Opened {
147 account_id, limits, ..
148 } => {
149 if self.account_id.is_some() {
150 return Err(BudgetError::Persistence(
151 crate::DomainError::AlreadyExists {
152 what: "budget_ledger",
153 },
154 ));
155 }
156 self.account_id = Some(account_id);
157 self.limits = Some(limits);
158 }
159 BudgetLedgerEvent::Reserved {
160 reservation_id,
161 operation_id,
162 quantities,
163 estimate,
164 reserved_at,
165 ..
166 } => {
167 if self.account_id.is_none() {
168 return Err(BudgetError::LedgerNotOpen);
169 }
170 let expected_id = BudgetReservationId::for_operation(
171 self.account_id.as_ref().expect("checked above"),
172 &operation_id,
173 );
174 if reservation_id != expected_id {
175 return Err(BudgetError::ReservationConflict(reservation_id));
176 }
177 if self.reservations.contains_key(&reservation_id) {
178 return Err(BudgetError::Persistence(
179 crate::DomainError::AlreadyExists {
180 what: "budget_reservation",
181 },
182 ));
183 }
184 self.ensure_capacity(quantities)?;
185 let expected_quantities =
186 estimate.quantities(self.limits.as_ref().ok_or(BudgetError::LedgerNotOpen)?)?;
187 if quantities != expected_quantities {
188 return Err(BudgetError::ReservationConflict(reservation_id));
189 }
190 self.reservations.insert(
191 reservation_id.clone(),
192 BudgetReservation::new(
193 reservation_id,
194 operation_id,
195 quantities,
196 estimate,
197 reserved_at,
198 ),
199 );
200 }
201 BudgetLedgerEvent::Reconciled {
202 reservation_id,
203 reconciliation_id,
204 measured,
205 ..
206 } => {
207 let reservation = self
208 .reservations
209 .get_mut(&reservation_id)
210 .ok_or_else(|| BudgetError::ReservationNotFound(reservation_id.clone()))?;
211 if let Some((stored_id, stored)) = reservation.reconciliation() {
212 if stored_id != &reconciliation_id || !measured.advances(*stored) {
213 return Err(BudgetError::ReconciliationConflict(reconciliation_id));
214 }
215 }
216 reservation.reconcile(reconciliation_id, measured);
217 }
218 }
219 self.version = self.version.next();
220 Ok(())
221 }
222
223 pub fn balance(&self) -> Result<BudgetBalance, BudgetError> {
224 let limits = self.limits.clone().ok_or(BudgetError::LedgerNotOpen)?;
225 let mut reserved = [0; 4];
226 let mut observed = [0; 4];
227 let mut estimated = [0; 4];
228 let mut unconfirmed = [0; 4];
229 for reservation in self.reservations.values() {
230 let requested = values(reservation.quantities());
231 match reservation.reconciliation() {
232 None => add(&mut reserved, requested)?,
233 Some((_, measured)) => classify(
234 &mut observed,
235 &mut estimated,
236 &mut unconfirmed,
237 requested,
238 *measured,
239 )?,
240 }
241 }
242 let maximum = limit_values(&limits);
243 let charged = sums([reserved, observed, estimated, unconfirmed])?;
244 let available = zip(maximum, charged, u64::saturating_sub);
245 let overrun = zip(charged, maximum, u64::saturating_sub);
246 Ok(BudgetBalance::new(
247 limits,
248 quantities(reserved),
249 quantities(observed),
250 quantities(estimated),
251 quantities(unconfirmed),
252 quantities(overrun),
253 quantities(available),
254 ))
255 }
256
257 fn ensure_capacity(&self, request: BudgetQuantities) -> Result<(), BudgetError> {
258 let account_id = self.account_id.clone().ok_or(BudgetError::LedgerNotOpen)?;
259 let balance = self.balance()?;
260 let overrun = values(balance.overrun());
261 let dimensions = [
262 BudgetDimension::Duration,
263 BudgetDimension::Tokens,
264 BudgetDimension::Cost,
265 BudgetDimension::ToolCalls,
266 ];
267 if let Some((index, _)) = overrun.iter().enumerate().find(|(_, value)| **value > 0) {
268 return Err(BudgetError::Exhausted {
269 account_id,
270 dimension: dimensions[index],
271 });
272 }
273 let available = values(balance.available());
274 let wanted = values(request);
275 for (index, dimension) in dimensions.into_iter().enumerate() {
276 if wanted[index] > available[index] {
277 return Err(BudgetError::Exhausted {
278 account_id,
279 dimension,
280 });
281 }
282 }
283 Ok(())
284 }
285}
286
287fn values(q: BudgetQuantities) -> [u64; 4] {
288 [
289 q.duration().as_micros(),
290 q.tokens().value(),
291 q.cost().value(),
292 q.tool_calls().value(),
293 ]
294}
295
296fn limit_values(limits: &BudgetLimits) -> [u64; 4] {
297 [
298 limits
299 .duration()
300 .map_or(u64::MAX, ExecutionDuration::as_micros),
301 limits.tokens().map_or(u64::MAX, BudgetTokenCount::value),
302 limits.cost().map_or(u64::MAX, CostMicros::value),
303 limits.tool_calls().map_or(u64::MAX, ToolCallCount::value),
304 ]
305}
306fn quantities(v: [u64; 4]) -> BudgetQuantities {
307 BudgetQuantities::new(
308 ExecutionDuration::from_micros(v[0]),
309 BudgetTokenCount::new(v[1]),
310 CostMicros::new(v[2]),
311 ToolCallCount::new(v[3]),
312 )
313}
314fn add(target: &mut [u64; 4], value: [u64; 4]) -> Result<(), BudgetError> {
315 for index in 0..4 {
316 target[index] = target[index]
317 .checked_add(value[index])
318 .ok_or(BudgetError::Persistence(
319 crate::DomainError::InvariantViolated {
320 reason: "budget total overflow",
321 },
322 ))?;
323 }
324 Ok(())
325}
326fn sums(values: [[u64; 4]; 4]) -> Result<[u64; 4], BudgetError> {
327 let mut result = [0; 4];
328 for value in values {
329 add(&mut result, value)?;
330 }
331 Ok(result)
332}
333fn zip(left: [u64; 4], right: [u64; 4], op: fn(u64, u64) -> u64) -> [u64; 4] {
334 [
335 op(left[0], right[0]),
336 op(left[1], right[1]),
337 op(left[2], right[2]),
338 op(left[3], right[3]),
339 ]
340}
341fn classify(
342 observed: &mut [u64; 4],
343 estimated: &mut [u64; 4],
344 unknown: &mut [u64; 4],
345 requested: [u64; 4],
346 measured: MeasuredBudgetQuantities,
347) -> Result<(), BudgetError> {
348 classify_one(
349 observed,
350 estimated,
351 unknown,
352 0,
353 requested[0],
354 measured.duration(),
355 ExecutionDuration::as_micros,
356 )?;
357 classify_one(
358 observed,
359 estimated,
360 unknown,
361 1,
362 requested[1],
363 measured.tokens(),
364 BudgetTokenCount::value,
365 )?;
366 classify_one(
367 observed,
368 estimated,
369 unknown,
370 2,
371 requested[2],
372 measured.cost(),
373 CostMicros::value,
374 )?;
375 classify_one(
376 observed,
377 estimated,
378 unknown,
379 3,
380 requested[3],
381 measured.tool_calls(),
382 ToolCallCount::value,
383 )
384}
385fn classify_one<T>(
386 observed: &mut [u64; 4],
387 estimated: &mut [u64; 4],
388 unknown: &mut [u64; 4],
389 index: usize,
390 requested: u64,
391 measurement: BudgetMeasurement<T>,
392 value_of: fn(T) -> u64,
393) -> Result<(), BudgetError> {
394 match measurement {
395 BudgetMeasurement::Observed(value) => add_at(observed, index, value_of(value)),
396 BudgetMeasurement::Estimated(value) => add_at(estimated, index, value_of(value)),
397 BudgetMeasurement::Unknown => add_at(unknown, index, requested),
398 }
399}
400fn add_at(target: &mut [u64; 4], index: usize, value: u64) -> Result<(), BudgetError> {
401 target[index] = target[index]
402 .checked_add(value)
403 .ok_or(BudgetError::Persistence(
404 crate::DomainError::InvariantViolated {
405 reason: "budget total overflow",
406 },
407 ))?;
408 Ok(())
409}