Skip to main content

made_core/entities/
budget_ledger.rs

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}