use std::sync::atomic::{AtomicBool, AtomicI64, Ordering};
use std::sync::{Arc, Mutex, PoisonError};
use serde_json::{Map, Value};
use super::budget::{self, fold_balance, last_grant, Allocation, BudgetGrant};
use super::event::{
data_field, is_kernel_only, kernel_event, BUDGET_KINDS, FIELD_AMOUNT, FIELD_CHILD, FIELD_DESC,
FIELD_DETAIL, FIELD_KIND, FIELD_OPEN_CHILDREN, FIELD_OWNER, FIELD_PARENT, FIELD_REASON,
FIELD_REMAINING, FIELD_SCOPE_ID, FIELD_TAG, KIND_BUDGET_GRANTED, KIND_BUDGET_REFUSED,
KIND_BUDGET_RESERVED, KIND_BUDGET_SPENT, KIND_SESSION_CLOSED, KIND_SESSION_OPENED,
};
use super::event_store::{ChildScan, Current, CurrentStore, EventStore, Split};
use super::logs::Logs;
use super::projection::{tail_count, VIEW_TAIL};
use super::query::{self, QueryOpts, QueryParams, QueryRows};
use super::scope::{Scope, ScopeId};
use super::sqlite_store::SqliteEventStore;
use super::{projection, KnlError, KnlResult};
pub const DEFAULT_CLOSE_REASON: &str = "closed";
pub const CLOSE_REASON_SCOPE_EXIT: &str = "scope_exit";
pub const CLOSE_REASON_ERROR: &str = "error";
pub const CLOSE_REASON_DROPPED: &str = "dropped";
const CLOSED: &str = "session is closed";
async fn has_ended(store: &CurrentStore) -> KnlResult<bool> {
let ending = store.read_kinds(Some(&[KIND_SESSION_CLOSED]), 0, 1).await?;
Ok(!ending.is_empty())
}
pub const ANON: &str = "anon";
pub const SYSTEM: &str = "system";
fn granted_event(grant: &BudgetGrant, scope_id: &str) -> Map<String, Value> {
let mut data = Map::new();
data.insert(
FIELD_SCOPE_ID.to_string(),
Value::from(scope_id.to_string()),
);
data.insert(FIELD_AMOUNT.to_string(), Value::from(grant.amount));
if let Some(tag) = grant.tag.as_ref() {
data.insert(FIELD_TAG.to_string(), Value::from(tag.clone()));
}
if let Some(desc) = grant.desc.as_ref() {
data.insert(FIELD_DESC.to_string(), Value::from(desc.clone()));
}
kernel_event(KIND_BUDGET_GRANTED, data)
}
fn budget_move_event(
kind: &str,
amount: i64,
tag: Option<&str>,
scope_id: &str,
) -> Map<String, Value> {
kernel_event(kind, budget_move_data(amount, tag, scope_id))
}
fn budget_move_data(amount: i64, tag: Option<&str>, scope_id: &str) -> Map<String, Value> {
let mut data = Map::new();
data.insert(
FIELD_SCOPE_ID.to_string(),
Value::from(scope_id.to_string()),
);
data.insert(FIELD_AMOUNT.to_string(), Value::from(amount));
if let Some(tag) = tag {
data.insert(FIELD_TAG.to_string(), Value::from(tag.to_string()));
}
data
}
fn refused_event(
amount: i64,
remaining: i64,
tag: Option<&str>,
scope_id: &str,
) -> Map<String, Value> {
let mut data = budget_move_data(amount, tag, scope_id);
data.insert(FIELD_REMAINING.to_string(), Value::from(remaining));
kernel_event(KIND_BUDGET_REFUSED, data)
}
fn allocated_event(
amount: i64,
tag: Option<&str>,
scope_id: &str,
child: &str,
) -> Map<String, Value> {
let mut data = budget_move_data(amount, tag, scope_id);
data.insert(FIELD_CHILD.to_string(), Value::from(child.to_string()));
kernel_event(KIND_BUDGET_RESERVED, data)
}
fn allocation_refused_event(
amount: i64,
remaining: i64,
tag: Option<&str>,
scope_id: &str,
child: &str,
) -> Map<String, Value> {
let mut event = refused_event(amount, remaining, tag, scope_id);
if let Some(Value::Object(data)) = event.get_mut(super::event::FIELD_DATA) {
data.insert(FIELD_CHILD.to_string(), Value::from(child.to_string()));
}
event
}
fn with_meta(
mut event: Map<String, Value>,
meta: Option<Map<String, Value>>,
) -> Map<String, Value> {
if let Some(meta) = meta {
if !meta.is_empty() {
event.insert(super::event::FIELD_META.to_string(), Value::Object(meta));
}
}
event
}
fn child_opened_event(
owner: &str,
scope_id: &str,
parent: &str,
meta: Option<Map<String, Value>>,
) -> Map<String, Value> {
let mut data = Map::new();
data.insert(FIELD_OWNER.to_string(), Value::from(owner.to_string()));
data.insert(
FIELD_SCOPE_ID.to_string(),
Value::from(scope_id.to_string()),
);
data.insert(FIELD_PARENT.to_string(), Value::from(parent.to_string()));
with_meta(kernel_event(KIND_SESSION_OPENED, data), meta)
}
fn child_granted_event(
amount: i64,
tag: Option<&str>,
scope_id: &str,
parent: &str,
) -> Map<String, Value> {
let mut data = budget_move_data(amount, tag, scope_id);
data.insert(FIELD_PARENT.to_string(), Value::from(parent.to_string()));
kernel_event(KIND_BUDGET_GRANTED, data)
}
fn closing_event(
reason: Option<&str>,
detail: Option<&str>,
children: Vec<String>,
) -> Map<String, Value> {
let mut data = Map::new();
data.insert(
FIELD_REASON.to_string(),
Value::from(reason.unwrap_or(DEFAULT_CLOSE_REASON)),
);
if let Some(detail) = detail {
data.insert(FIELD_DETAIL.to_string(), Value::from(detail.to_string()));
}
if !children.is_empty() {
data.insert(
FIELD_OPEN_CHILDREN.to_string(),
Value::from(children.into_iter().map(Value::from).collect::<Vec<_>>()),
);
}
kernel_event(KIND_SESSION_CLOSED, data)
}
fn child_scan() -> ChildScan {
ChildScan {
opened: KIND_SESSION_OPENED.to_string(),
closed: KIND_SESSION_CLOSED.to_string(),
parent_field: FIELD_PARENT.to_string(),
}
}
const ALLOCATION_KINDS: &[&str] = &[
KIND_BUDGET_GRANTED,
KIND_BUDGET_RESERVED,
KIND_BUDGET_REFUSED,
KIND_BUDGET_SPENT,
KIND_SESSION_CLOSED,
];
fn kernel_only_hint(kind: &str) -> &'static str {
match kind {
KIND_SESSION_OPENED => "a session records its own opening",
KIND_SESSION_CLOSED => "use close",
_ => "use reserve / spend",
}
}
pub struct Session {
id: String,
scope: Scope,
store: CurrentStore,
balance: Mutex<(u64, Option<i64>)>,
closed: bool,
}
impl std::fmt::Debug for Session {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Session")
.field("id", &self.id)
.field("scope", &self.scope)
.field("closed", &self.closed)
.finish_non_exhaustive()
}
}
impl Session {
pub async fn new(
owner: String,
grant: Option<BudgetGrant>,
meta: Option<Map<String, Value>>,
logs: &Logs,
) -> KnlResult<Self> {
let stream = uuid::Uuid::new_v4().to_string();
let store = SqliteEventStore::open_memory(stream.clone(), logs).await?;
let mut session = Self::open_on(owner, grant, meta, Box::new(store)).await?;
session.adopt_id(stream);
Ok(session)
}
pub async fn open_on(
owner: String,
grant: Option<BudgetGrant>,
meta: Option<Map<String, Value>>,
store: Box<dyn EventStore>,
) -> KnlResult<Self> {
let store = CurrentStore::new(store, Vec::new());
let mut session = Self {
id: uuid::Uuid::new_v4().to_string(),
scope: Scope::new(owner, grant),
store,
balance: Mutex::new((0, None)),
closed: false,
};
let mut opened = Map::new();
opened.insert(
FIELD_OWNER.to_string(),
Value::from(session.scope.owner().to_string()),
);
opened.insert(
FIELD_SCOPE_ID.to_string(),
Value::from(session.scope.id().to_string()),
);
let started = with_meta(kernel_event(KIND_SESSION_OPENED, opened), meta);
let mut opening = vec![started];
if let Some(grant) = session.scope.grant().cloned() {
opening.push(granted_event(&grant, session.scope.id()));
}
session.store.append_many(opening).await?;
Ok(session)
}
pub async fn resume(grant: Option<BudgetGrant>, store: Box<dyn EventStore>) -> KnlResult<Self> {
Self::resume_on(grant, CurrentStore::new(store, Vec::new())).await
}
async fn resume_on(grant: Option<BudgetGrant>, store: CurrentStore) -> KnlResult<Self> {
let log = store.read(0, usize::MAX).await?;
let opened = log
.iter()
.find(|event| event.kind() == KIND_SESSION_OPENED)
.ok_or_else(|| {
KnlError::Validation(
"stream has no session to resume (no session_opened event)".to_string(),
)
})?;
if has_ended(&store).await? {
return Err(KnlError::Closed(format!(
"{CLOSED} (disposable; open a new session)"
)));
}
let owner = data_field(opened, FIELD_OWNER)
.and_then(Value::as_str)
.unwrap_or(ANON)
.to_string();
let scope_id: Option<ScopeId> = data_field(opened, FIELD_SCOPE_ID)
.and_then(Value::as_str)
.map(str::to_string);
let head = log.last().map(Current::seq).unwrap_or(0);
let mut session = Self {
id: uuid::Uuid::new_v4().to_string(),
scope: Scope::restore(scope_id, owner, last_grant(&log)),
store,
balance: Mutex::new((head, fold_balance(&log))),
closed: false,
};
if let Some(grant) = grant {
session.grant_on_resume(grant).await?;
}
Ok(session)
}
async fn grant_on_resume(&mut self, grant: BudgetGrant) -> KnlResult<()> {
budget::check_amount(grant.amount)?;
let scope_id = self.scope.id().to_string();
let recorded = grant.clone();
let committed = self
.store
.append_if(
Some(BUDGET_KINDS),
Box::new(move |events: Vec<Current>| {
last_grant(&events)?;
Some(granted_event(&recorded, &scope_id))
}),
)
.await?;
if committed.is_none() {
return Err(KnlError::Validation(
"this stream opened with no budget, and a resume does not give one: a session's \
quota is settled when it opens (open a new session with the grant, or resume \
without one)"
.to_string(),
));
}
self.scope.grant_more(grant)
}
pub async fn grant_more(&mut self, grant: BudgetGrant) -> KnlResult<()> {
let event = granted_event(&grant, self.scope.id());
self.append_kernel(event).await?;
self.scope.grant_more(grant)
}
pub async fn open_child(
&mut self,
child_stream: String,
owner: String,
allocation: Allocation,
meta: Option<Map<String, Value>>,
child_store: Box<dyn EventStore>,
) -> KnlResult<Self> {
if self.closed {
return Err(KnlError::Closed(format!(
"{CLOSED} (a child is opened from an open parent)"
)));
}
budget::check_amount(allocation.amount)?;
match (self.store.database(), child_store.database()) {
(Some(parent), Some(child)) if parent == child => {}
(Some(parent), Some(child)) => {
return Err(KnlError::Validation(format!(
"a child opens on its parent's database, and a tree is one log: the parent is \
on {parent:?} and the child was given {child:?}"
)));
}
_ => {
return Err(KnlError::Validation(
"a child opens on its parent's database, and one of the two stores keeps a \
single stream with no database to share"
.to_string(),
));
}
}
let amount = allocation.amount;
let parent_tag = self.scope.grant().and_then(|grant| grant.tag.clone());
let child_tag = allocation.tag.clone().or_else(|| parent_tag.clone());
let parent_scope = self.scope.id().to_string();
let parent_id = self.id.clone();
let child_id = child_stream.clone();
let child_scope_id = Scope::new(owner.clone(), None).id().to_string();
let ended = Arc::new(AtomicBool::new(false));
let occupied = Arc::new(AtomicBool::new(false));
let refused = Arc::new(AtomicBool::new(false));
let measured = Arc::new(AtomicI64::new(0));
let (found_ending, found_events, said_no, balance_seen) = (
Arc::clone(&ended),
Arc::clone(&occupied),
Arc::clone(&refused),
Arc::clone(&measured),
);
let committed = self
.store
.append_if_many(
&child_stream,
Some(ALLOCATION_KINDS),
Box::new(move |seen: Split<Current>| {
if !seen.other.is_empty() {
found_events.store(true, Ordering::Relaxed);
return None;
}
let events = seen.own;
if events.iter().any(|e| e.kind() == KIND_SESSION_CLOSED) {
found_ending.store(true, Ordering::Relaxed);
return None;
}
if let Some(balance) = fold_balance(&events) {
if balance < amount {
said_no.store(true, Ordering::Relaxed);
balance_seen.store(balance, Ordering::Relaxed);
return Some(Split::own(vec![allocation_refused_event(
amount,
balance,
parent_tag.as_deref(),
&parent_scope,
&child_id,
)]));
}
}
Some(Split {
own: vec![allocated_event(
amount,
parent_tag.as_deref(),
&parent_scope,
&child_id,
)],
other: vec![
child_opened_event(&owner, &child_scope_id, &parent_id, meta.clone()),
child_granted_event(
amount,
child_tag.as_deref(),
&child_scope_id,
&parent_id,
),
],
})
}),
)
.await?;
if occupied.load(Ordering::Relaxed) {
return Err(KnlError::Validation(format!(
"a child opens on a stream of its own, and {child_stream:?} already has events on \
it: pass a stream nothing has been written to (nothing was written here, on \
either side)"
)));
}
if committed.is_none() || ended.load(Ordering::Relaxed) {
return Err(KnlError::Closed(format!(
"{CLOSED} (the parent's log already carries its ending)"
)));
}
if refused.load(Ordering::Relaxed) {
let balance = measured.load(Ordering::Relaxed);
return Err(KnlError::Refused(format!(
"the parent's balance is {balance}, which does not cover an allocation of \
{amount}; the refusal is in the log and no child was opened"
)));
}
let mut child = Self::resume(None, child_store).await?;
child.adopt_id(child_stream);
Ok(child)
}
pub fn id(&self) -> &str {
&self.id
}
pub(crate) fn adopt_id(&mut self, id: impl Into<String>) {
self.id = id.into();
}
pub fn scope(&self) -> &Scope {
&self.scope
}
pub fn scope_id(&self) -> &str {
self.scope.id()
}
pub fn owner(&self) -> &str {
self.scope.owner()
}
pub fn database(&self) -> Option<&str> {
self.store.database()
}
pub async fn append(&mut self, event: Map<String, Value>) -> KnlResult<u64> {
let kind = event.get(FIELD_KIND).and_then(Value::as_str).unwrap_or("");
if is_kernel_only(kind) {
return Err(KnlError::Validation(format!(
"{kind:?} is written by the kernel only ({})",
kernel_only_hint(kind)
)));
}
self.append_kernel(event).await
}
async fn append_kernel(&mut self, event: Map<String, Value>) -> KnlResult<u64> {
if self.closed {
return Err(KnlError::Closed(CLOSED.to_string()));
}
let committed = self.store.append(event).await?;
Ok(committed.seq)
}
pub async fn events(&self, from: u64, limit: usize) -> KnlResult<Vec<Current>> {
self.store.read(from, limit).await
}
pub async fn len(&self) -> KnlResult<usize> {
self.store.len().await
}
pub async fn is_empty(&self) -> KnlResult<bool> {
self.store.is_empty().await
}
pub async fn reserve(&mut self, amount: i64) -> KnlResult<bool> {
if self.closed {
return Err(KnlError::Closed(CLOSED.to_string()));
}
budget::check_amount(amount)?;
let scope_id = self.scope.id().to_string();
let refused = Arc::new(AtomicBool::new(false));
let said_no = Arc::clone(&refused);
let decided_scope = scope_id.clone();
let committed = self
.store
.append_if(
Some(BUDGET_KINDS),
Box::new(move |events: Vec<Current>| {
let grant = last_grant(&events)?;
let balance = fold_balance(&events).unwrap_or(0);
if balance >= amount {
return Some(budget_move_event(
KIND_BUDGET_RESERVED,
amount,
grant.tag.as_deref(),
&decided_scope,
));
}
said_no.store(true, Ordering::Relaxed);
Some(refused_event(
amount,
balance,
grant.tag.as_deref(),
&decided_scope,
))
}),
)
.await?;
match (committed, refused.load(Ordering::Relaxed)) {
(Some(_), true) => Ok(false),
_ => Ok(true),
}
}
pub async fn spend(&mut self, amount: i64) -> KnlResult<()> {
if self.closed {
return Err(KnlError::Closed(CLOSED.to_string()));
}
budget::check_amount(amount)?;
let scope_id = self.scope.id().to_string();
self.store
.append_if(
Some(BUDGET_KINDS),
Box::new(move |events: Vec<Current>| {
let grant = last_grant(&events)?;
Some(budget_move_event(
KIND_BUDGET_SPENT,
amount,
grant.tag.as_deref(),
&scope_id,
))
}),
)
.await?;
Ok(())
}
pub fn grant(&self) -> Option<&BudgetGrant> {
self.scope.grant()
}
pub async fn remaining(&self) -> KnlResult<Option<i64>> {
let (folded_head_seq, cached) =
*self.balance.lock().unwrap_or_else(PoisonError::into_inner);
let head = self.store.head().await?.unwrap_or(0);
if head <= folded_head_seq {
return Ok(cached);
}
let ledger = self
.store
.read_kinds(Some(BUDGET_KINDS), 0, usize::MAX)
.await?;
let balance = fold_balance(&ledger);
*self.balance.lock().unwrap_or_else(PoisonError::into_inner) = (head, balance);
Ok(balance)
}
pub async fn exhausted(&self) -> KnlResult<bool> {
Ok(matches!(self.remaining().await?, Some(remaining) if remaining <= 0))
}
pub fn is_closed(&self) -> bool {
self.closed
}
pub async fn close(&mut self, reason: Option<&str>) -> KnlResult<()> {
self.close_with(reason, None).await
}
pub async fn close_with(
&mut self,
reason: Option<&str>,
detail: Option<&str>,
) -> KnlResult<()> {
if self.closed {
return Ok(());
}
let reason = reason.map(str::to_string);
let detail = detail.map(str::to_string);
self.store
.append_with_open_children(
&child_scan(),
Box::new(move |children| {
closing_event(reason.as_deref(), detail.as_deref(), children)
}),
)
.await?;
self.closed = true;
Ok(())
}
pub fn close_detached(&mut self, reason: &str) {
if self.closed {
return;
}
self.store
.detach_append(closing_event(Some(reason), None, Vec::new()));
self.closed = true;
}
pub async fn view(
&mut self,
name: &str,
opts: Option<&Map<String, Value>>,
) -> KnlResult<Value> {
match name {
VIEW_TAIL => {
let n = tail_count(opts)?;
let events = self.store.read_last(n).await?;
Ok(projection::tail_of(&events, n))
}
other => Err(KnlError::Validation(format!("unknown view {other:?}"))),
}
}
pub async fn query(
&self,
sql: &str,
params: QueryParams,
opts: &QueryOpts,
) -> KnlResult<QueryRows> {
let plan = query::plan(sql, params, opts, &self.id)?;
self.store.query(&plan).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::knl::event::{kind_of, FIELD_BEAT, FIELD_DATA, FIELD_META};
use crate::knl::event_store::MemEventStore;
use serde_json::json;
fn obj(value: Value) -> Map<String, Value> {
match value {
Value::Object(map) => map,
other => panic!("test fixture must be an object, got {other}"),
}
}
fn grant(amount: i64) -> BudgetGrant {
BudgetGrant {
amount,
tag: Some("tokens".to_string()),
desc: None,
}
}
fn test_logs() -> Logs {
static LOGS: std::sync::OnceLock<Logs> = std::sync::OnceLock::new();
LOGS.get_or_init(Logs::new).clone()
}
async fn new_session(budget: Option<i64>) -> Session {
Session::new(ANON.to_string(), budget.map(grant), None, &test_logs())
.await
.expect("open")
}
async fn folded(s: &Session) -> Option<i64> {
fold_balance(&s.events(0, usize::MAX).await.expect("events"))
}
fn as_current(events: Vec<Value>) -> Vec<Current> {
events.into_iter().map(Current::assume_current).collect()
}
fn kinds(events: &[Current]) -> Vec<&str> {
events.iter().map(Current::kind).collect()
}
async fn remaining(session: &Session) -> Option<i64> {
session.remaining().await.expect("the balance was readable")
}
async fn exhausted(session: &Session) -> bool {
session.exhausted().await.expect("the balance was readable")
}
async fn ledger(s: &Session) -> Vec<Current> {
s.events(0, usize::MAX)
.await
.expect("events")
.into_iter()
.filter(|e| e.kind().starts_with("budget_"))
.collect()
}
fn response(tokens: i64) -> Map<String, Value> {
obj(json!({
"kind": "llm_response",
"data": {
"content": [{ "type": "text", "text": "ok" }],
"usage": { "input_tokens": tokens },
"stop_reason": "end_turn"
}
}))
}
fn field<'a>(event: &'a Current, name: &str) -> &'a Value {
data_field(event, name).unwrap_or_else(|| panic!("data.{name} is missing: {event}"))
}
#[tokio::test]
async fn opening_carries_the_labels_it_was_opened_with() {
let meta = obj(json!({ "run": "r-1", "attempt": 2, "retried": true }));
let s = Session::new(ANON.to_string(), None, Some(meta.clone()), &Logs::new())
.await
.expect("open");
let events = s.events(0, usize::MAX).await.expect("events");
assert_eq!(events[0].kind(), KIND_SESSION_OPENED);
assert_eq!(
events[0].get(crate::knl::event::FIELD_META),
Some(&Value::Object(meta)),
"the labels ride the envelope, verbatim"
);
assert_eq!(
data_field(&events[0], "run"),
None,
"a label is not a field of the kind: data stays the kernel's"
);
}
#[tokio::test]
async fn an_unlabelled_opening_carries_no_labels() {
for session in [
new_session(None).await,
Session::new(ANON.to_string(), None, Some(Map::new()), &Logs::new())
.await
.expect("open"),
] {
let events = session.events(0, usize::MAX).await.expect("events");
let meta = events[0]
.get(crate::knl::event::FIELD_META)
.and_then(Value::as_object)
.expect("the store gives every event a meta object");
assert!(meta.is_empty(), "nothing was named, so nothing is recorded");
}
}
#[tokio::test]
async fn a_new_session_already_carries_session_opened() {
let s = new_session(None).await;
assert_eq!(s.len().await.expect("len"), 1);
let events = s.events(0, usize::MAX).await.expect("events");
assert_eq!(events[0].kind(), KIND_SESSION_OPENED);
assert_eq!(events[0].seq(), 1);
assert!(!s.is_closed());
assert!(!s.id().is_empty());
}
#[tokio::test]
async fn open_issues_a_scope_id_distinct_from_the_session_id() {
let a = new_session(None).await;
let b = new_session(None).await;
assert!(!a.scope_id().is_empty(), "a session opens under a scope");
assert!(!a.id().is_empty());
assert_ne!(
a.scope_id(),
a.id(),
"the scope id names the authority, the session id names the stream"
);
assert_eq!(a.scope().id(), a.scope_id(), "the delegate reads the scope");
assert_eq!(a.scope().owner(), a.owner(), "and so does the owner");
assert_ne!(a.scope_id(), b.scope_id(), "two sessions, two scopes");
assert_ne!(a.id(), b.id());
}
#[tokio::test]
async fn session_opened_and_every_budget_event_carry_the_scope_id() {
let mut s = new_session(Some(100)).await;
assert_eq!(s.reserve(30).await, Ok(true));
s.spend(10).await.expect("spend");
assert_eq!(s.reserve(10_000).await, Ok(false));
let scope_id = s.scope_id().to_string();
let events = s.events(0, usize::MAX).await.expect("events");
let opened = &events[0];
assert_eq!(opened.kind(), KIND_SESSION_OPENED);
assert_eq!(
field(opened, FIELD_SCOPE_ID).as_str(),
Some(scope_id.as_str()),
"the scope rides on session_opened: {opened}"
);
assert_eq!(
field(opened, FIELD_OWNER).as_str(),
Some(ANON),
"beside the owner: {opened}"
);
let moves = ledger(&s).await;
assert_eq!(
kinds(&moves),
vec![
KIND_BUDGET_GRANTED,
KIND_BUDGET_RESERVED,
KIND_BUDGET_SPENT,
KIND_BUDGET_REFUSED,
],
"every kind of move is exercised"
);
for event in &moves {
assert_eq!(
field(event, FIELD_SCOPE_ID).as_str(),
Some(scope_id.as_str()),
"a ledger entry must name the scope it was allowed under: {event}"
);
}
s.append(obj(json!({ "kind": "note" })))
.await
.expect("append");
let note = s
.events(0, usize::MAX)
.await
.expect("events")
.pop()
.expect("note");
assert_eq!(data_field(¬e, FIELD_SCOPE_ID), None, "{note}");
assert_eq!(note[FIELD_DATA], json!({}), "and no data of its own");
}
#[tokio::test]
async fn the_owner_is_total_and_read_back_verbatim() {
assert_eq!(new_session(None).await.owner(), ANON);
assert_eq!(
Session::new(SYSTEM.to_string(), None, None, &Logs::new())
.await
.expect("open")
.owner(),
SYSTEM
);
assert_eq!(
Session::new("user-42".to_string(), None, None, &Logs::new())
.await
.expect("open")
.owner(),
"user-42"
);
}
#[tokio::test]
async fn close_records_session_closed_once_with_the_given_reason() {
let mut s = new_session(None).await;
s.close(Some("budget_exhausted")).await.expect("close");
s.close(Some("ignored")).await.expect("close (idempotent)");
assert_eq!(s.len().await.expect("len"), 2, "close must be idempotent");
let last = s
.events(2, usize::MAX)
.await
.expect("events")
.pop()
.expect("session_closed");
assert_eq!(last.kind(), KIND_SESSION_CLOSED);
assert_eq!(*field(&last, FIELD_REASON), json!("budget_exhausted"));
assert!(s.is_closed());
}
#[tokio::test]
async fn close_without_a_reason_records_the_default() {
let mut s = new_session(None).await;
s.close(None).await.expect("close");
let last = s
.events(2, usize::MAX)
.await
.expect("events")
.pop()
.expect("session_closed");
assert_eq!(*field(&last, FIELD_REASON), json!(DEFAULT_CLOSE_REASON));
}
#[tokio::test]
async fn a_detached_close_records_the_same_boundary() {
let mut s = new_session(None).await;
s.close_detached(CLOSE_REASON_DROPPED);
assert!(s.is_closed(), "the handle is closed straight away");
s.close_detached(CLOSE_REASON_DROPPED);
s.close(Some("ignored")).await.expect("close is a no-op");
let last = s
.events(0, usize::MAX)
.await
.expect("events")
.pop()
.expect("session_closed");
assert_eq!(last.kind(), KIND_SESSION_CLOSED);
assert_eq!(
*field(&last, FIELD_REASON),
json!(CLOSE_REASON_DROPPED),
"exactly one boundary, carrying the backstop's reason"
);
assert_eq!(s.len().await.expect("len"), 2, "session_opened + closed");
}
#[tokio::test]
async fn a_closed_session_rejects_writes_but_keeps_serving_reads() {
let mut s = new_session(Some(10)).await;
s.append(obj(json!({ "kind": "note" })))
.await
.expect("append");
s.spend(4).await.expect("spend");
s.close(None).await.expect("close");
let err = s
.append(obj(json!({ "kind": "note" })))
.await
.expect_err("append after close");
assert_eq!(err.reason(), "session is closed");
let err = s.spend(1).await.expect_err("spend after close");
assert_eq!(err.reason(), "session is closed");
let err = s.reserve(1).await.expect_err("reserve after close");
assert_eq!(err.reason(), "session is closed");
assert_eq!(
s.len().await.expect("len"),
5,
"session_opened + budget_granted + note + budget_spent + session_closed"
);
assert_eq!(remaining(&s).await, Some(6));
assert_eq!(
folded(&s).await,
remaining(&s).await,
"the ledger is the balance"
);
assert!(!exhausted(&s).await);
assert_eq!(
s.events(0, usize::MAX).await.expect("events")[2].kind(),
"note"
);
}
#[tokio::test]
async fn two_sessions_share_nothing() {
let mut a = new_session(Some(100)).await;
let mut b = new_session(Some(100)).await;
assert_ne!(a.id(), b.id());
a.append(obj(json!({ "kind": "only_in_a" })))
.await
.expect("append");
a.spend(60).await.expect("spend");
assert_eq!(
a.len().await.expect("len"),
4,
"session_opened + budget_granted + only_in_a + budget_spent"
);
assert_eq!(
b.len().await.expect("len"),
2,
"session_opened + budget_granted"
);
assert_eq!(remaining(&a).await, Some(40));
assert_eq!(remaining(&b).await, Some(100));
assert_eq!(folded(&a).await, Some(40));
assert_eq!(folded(&b).await, Some(100));
a.close(None).await.expect("close");
assert!(b.append(obj(json!({ "kind": "still_open" }))).await.is_ok());
}
#[tokio::test]
async fn view_serves_the_one_named_projection_and_rejects_anything_else() {
let mut s = new_session(None).await;
s.append(obj(
json!({ "kind": "msg_user", "data": { "content": "hi" } }),
))
.await
.expect("append");
s.append(response(9)).await.expect("recorded");
let tail = s
.view(VIEW_TAIL, Some(&obj(json!({ "n": 1 }))))
.await
.expect("tail");
assert_eq!(tail.as_array().map(Vec::len), Some(1));
let err = s.view("nope", None).await.expect_err("unknown view");
assert_eq!(err.reason(), r#"unknown view "nope""#);
}
#[derive(Default)]
struct CountingStore {
inner: MemEventStore,
ranges: Arc<Mutex<Vec<(u64, usize)>>>,
tails: Arc<Mutex<Vec<(usize, usize)>>>,
}
#[async_trait::async_trait]
impl EventStore for CountingStore {
async fn append(&mut self, event: Map<String, Value>) -> KnlResult<crate::knl::Committed> {
self.inner.append(event).await
}
async fn append_if(
&mut self,
kinds: Option<&[&str]>,
decide: crate::knl::Decision,
) -> KnlResult<Option<crate::knl::Committed>> {
self.inner.append_if(kinds, decide).await
}
async fn read_kinds(
&self,
kinds: Option<&[&str]>,
from_seq: u64,
limit: usize,
) -> KnlResult<Vec<Value>> {
self.ranges
.lock()
.unwrap_or_else(PoisonError::into_inner)
.push((from_seq, limit));
self.inner.read_kinds(kinds, from_seq, limit).await
}
async fn read_last(&self, n: usize) -> KnlResult<Vec<Value>> {
let events = self.inner.read_last(n).await?;
self.tails
.lock()
.unwrap_or_else(PoisonError::into_inner)
.push((n, events.len()));
Ok(events)
}
async fn head(&self) -> KnlResult<Option<u64>> {
self.inner.head().await
}
async fn len(&self) -> KnlResult<usize> {
self.inner.len().await
}
}
#[tokio::test]
async fn tail_reads_the_end_of_the_stream_and_not_the_whole_of_it() {
let store = CountingStore::default();
let (ranges, tails) = (Arc::clone(&store.ranges), Arc::clone(&store.tails));
let mut s = Session::open_on(ANON.to_string(), None, None, Box::new(store))
.await
.expect("open");
for i in 0..1_000 {
s.append(obj(json!({ "kind": format!("e{i}") })))
.await
.expect("append");
}
assert_eq!(s.len().await.expect("len"), 1_001);
ranges
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clear();
tails.lock().unwrap_or_else(PoisonError::into_inner).clear();
let tail = s
.view(VIEW_TAIL, Some(&obj(json!({ "n": 5 }))))
.await
.expect("tail");
let tail = tail.as_array().expect("an array of events");
assert_eq!(tail.len(), 5);
assert_eq!(kind_of(&tail[4]), "e999", "the last event is the last one");
let taken = tails.lock().unwrap_or_else(PoisonError::into_inner).clone();
let scanned = ranges
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clone();
assert_eq!(
taken,
vec![(5, 5)],
"tail must ask the store for five rows and get five"
);
assert!(
scanned.is_empty(),
"no range read goes out behind it: {scanned:?}"
);
}
#[tokio::test]
async fn the_token_account_is_not_a_named_view() {
let mut s = new_session(None).await;
s.append(response(9)).await.expect("recorded");
let err = s.view("usage", None).await.expect_err("usage was served");
assert_eq!(err.reason(), r#"unknown view "usage""#);
assert_eq!(err.kind(), KnlError::VALIDATION);
let recorded = s
.events(2, usize::MAX)
.await
.expect("events")
.pop()
.expect("llm_response");
assert_eq!(*field(&recorded, "usage"), json!({ "input_tokens": 9 }));
}
#[tokio::test]
async fn the_conversation_is_not_a_named_view() {
let mut s = new_session(None).await;
s.append(obj(
json!({ "kind": "msg_user", "data": { "content": "hi" } }),
))
.await
.expect("append");
let err = s
.view("dialogue", None)
.await
.expect_err("dialogue was served");
assert_eq!(err.reason(), r#"unknown view "dialogue""#);
let events = s.events(0, usize::MAX).await.expect("events");
assert_eq!(events[1].kind(), "msg_user");
assert_eq!(*field(&events[1], "content"), json!("hi"));
}
#[tokio::test]
async fn appending_an_llm_response_records_it_verbatim_without_charging() {
let mut s = new_session(Some(100)).await;
let seq = s
.append(obj(json!({
"kind": "llm_response",
"meta": { "beat": "beat-7" },
"data": {
"content": [{ "type": "text", "text": "hi" }],
"usage": { "input_tokens": 20, "output_tokens": 10 }
}
})))
.await
.expect("append");
let recorded = s
.events(seq, usize::MAX)
.await
.expect("events")
.pop()
.expect("llm_response");
assert_eq!(recorded.kind(), "llm_response");
assert_eq!(
recorded[FIELD_META][FIELD_BEAT],
json!("beat-7"),
"the declared beat is recorded as given"
);
assert_eq!(remaining(&s).await, Some(100), "an append must not charge");
assert_eq!(
ledger(&s).await.len(),
1,
"only the opening grant is in the ledger"
);
assert_eq!(
*field(&recorded, "usage"),
json!({ "input_tokens": 20, "output_tokens": 10 }),
"the counts are stored as they came: {recorded}"
);
assert_eq!(
folded(&s).await,
Some(100),
"what was consumed and the balance are separate readings"
);
}
#[tokio::test]
async fn opening_with_a_grant_records_it() {
let s = Session::new(
ANON.to_string(),
Some(BudgetGrant {
amount: 500,
tag: Some("tokens".to_string()),
desc: Some("one nightly run".to_string()),
}),
None,
&Logs::new(),
)
.await
.expect("open");
let events = s.events(0, usize::MAX).await.expect("events");
assert_eq!(events.len(), 2, "session_opened + budget_granted");
assert_eq!(events[0].kind(), KIND_SESSION_OPENED);
assert_eq!(
data_field(&events[0], "budget"),
None,
"the grant is its own event, not a field on session_opened"
);
let granted = &events[1];
assert_eq!(granted.kind(), KIND_BUDGET_GRANTED);
assert_eq!(*field(granted, FIELD_AMOUNT), json!(500));
assert_eq!(*field(granted, FIELD_TAG), json!("tokens"));
assert_eq!(*field(granted, FIELD_DESC), json!("one nightly run"));
assert_eq!(remaining(&s).await, Some(500));
assert_eq!(folded(&s).await, remaining(&s).await);
let bare = new_session(None).await;
assert_eq!(bare.len().await.expect("len"), 1, "session_opened only");
assert!(ledger(&bare).await.is_empty());
assert_eq!(folded(&bare).await, None);
}
#[tokio::test]
async fn a_granted_reservation_is_one_event_and_one_deduction() {
let mut s = new_session(Some(100)).await;
assert_eq!(s.reserve(30).await, Ok(true));
let moves = ledger(&s).await;
assert_eq!(moves.len(), 2, "the grant and the reservation");
assert_eq!(moves[1].kind(), KIND_BUDGET_RESERVED);
assert_eq!(*field(&moves[1], FIELD_AMOUNT), json!(30));
assert_eq!(*field(&moves[1], FIELD_TAG), json!("tokens"));
assert_eq!(remaining(&s).await, Some(70));
assert_eq!(
folded(&s).await,
remaining(&s).await,
"the ledger is the balance"
);
}
#[tokio::test]
async fn a_refused_reservation_is_recorded_and_changes_no_balance() {
let mut s = new_session(Some(10)).await;
assert_eq!(s.reserve(11).await, Ok(false));
let moves = ledger(&s).await;
assert_eq!(moves.len(), 2, "the grant and the refusal");
assert_eq!(moves[1].kind(), KIND_BUDGET_REFUSED);
assert_eq!(*field(&moves[1], FIELD_AMOUNT), json!(11));
assert_eq!(
*field(&moves[1], FIELD_REMAINING),
json!(10),
"what there was"
);
assert_eq!(*field(&moves[1], FIELD_TAG), json!("tokens"));
assert_eq!(remaining(&s).await, Some(10), "a refusal must not deduct");
assert!(!exhausted(&s).await);
assert_eq!(folded(&s).await, remaining(&s).await);
assert_eq!(s.reserve(10).await, Ok(true));
assert_eq!(remaining(&s).await, Some(0));
assert_eq!(folded(&s).await, Some(0));
}
#[tokio::test]
async fn the_balance_is_the_fold_after_any_sequence_of_moves() {
let mut s = new_session(Some(1000)).await;
assert_eq!(s.reserve(200).await, Ok(true));
s.append(response(40)).await.expect("recorded");
s.spend(50).await.expect("spend");
assert_eq!(s.reserve(10_000).await, Ok(false));
assert_eq!(s.reserve(300).await, Ok(true));
s.spend(0).await.expect("spend");
let moves = ledger(&s).await;
assert_eq!(
kinds(&moves),
vec![
KIND_BUDGET_GRANTED,
KIND_BUDGET_RESERVED,
KIND_BUDGET_SPENT,
KIND_BUDGET_REFUSED,
KIND_BUDGET_RESERVED,
KIND_BUDGET_SPENT,
],
"every move left exactly one event"
);
assert_eq!(remaining(&s).await, Some(450), "1000 - 200 - 50 - 300");
assert_eq!(folded(&s).await, remaining(&s).await);
}
#[tokio::test]
async fn a_caller_cannot_append_the_budget_kinds() {
let mut s = new_session(Some(10)).await;
for event in [
json!({ "kind": "budget_granted", "data": { "amount": 1_000_000 } }),
json!({ "kind": "budget_reserved", "data": { "amount": 5 } }),
json!({ "kind": "budget_refused", "data": { "amount": 5, "remaining": 0 } }),
json!({ "kind": "budget_spent", "data": { "amount": 5 } }),
] {
let err = s
.append(obj(event.clone()))
.await
.expect_err("kernel-only kind");
assert!(
err.reason().contains("kernel only"),
"{event}: {}",
err.reason()
);
}
assert_eq!(
remaining(&s).await,
Some(10),
"no forged event moved the balance"
);
assert_eq!(ledger(&s).await.len(), 1, "nothing was recorded");
assert_eq!(folded(&s).await, remaining(&s).await);
}
#[tokio::test]
async fn a_caller_cannot_append_the_session_boundary_kinds() {
let mut s = new_session(Some(100)).await;
for event in [
json!({ "kind": "session_opened", "data": { "scope_id": "s", "owner": "me" } }),
json!({ "kind": "session_closed", "data": { "reason": "carried over" } }),
] {
let err = s
.append(obj(event.clone()))
.await
.expect_err("kernel-only kind");
assert!(
err.reason().contains("kernel only"),
"{event}: {}",
err.reason()
);
}
assert!(!s.is_closed(), "a refused append ended the session");
assert_eq!(s.len().await.expect("len"), 2, "nothing was recorded");
assert_eq!(s.append(obj(json!({ "kind": "note" }))).await, Ok(3));
assert_eq!(s.spend(10).await, Ok(()));
assert_eq!(remaining(&s).await, Some(90), "the settlement landed");
}
#[tokio::test]
async fn only_close_records_session_closed() {
let mut s = new_session(Some(100)).await;
s.append(obj(json!({ "kind": "note" })))
.await
.expect("append");
assert!(
!s.events(0, usize::MAX)
.await
.expect("events")
.iter()
.any(|e| e.kind() == KIND_SESSION_CLOSED),
"nothing but close writes the boundary"
);
s.close(Some("done")).await.expect("close");
assert!(s.is_closed());
let closed: Vec<Current> = s
.events(0, usize::MAX)
.await
.expect("events")
.into_iter()
.filter(|e| e.kind() == KIND_SESSION_CLOSED)
.collect();
assert_eq!(closed.len(), 1, "exactly one boundary: {closed:?}");
assert_eq!(*field(&closed[0], FIELD_REASON), json!("done"));
assert_eq!(
s.append(obj(json!({ "kind": "note" })))
.await
.expect_err("append after close")
.reason(),
"session is closed"
);
}
#[tokio::test]
async fn a_recorded_response_is_in_the_history_without_being_charged() {
let mut s = new_session(Some(100)).await;
s.append(response(30)).await.expect("recorded");
assert_eq!(remaining(&s).await, Some(100));
assert!(!exhausted(&s).await);
let recorded = s
.events(3, usize::MAX)
.await
.expect("events")
.pop()
.expect("llm_response");
assert_eq!(recorded.kind(), "llm_response");
assert_eq!(*field(&recorded, "stop_reason"), json!("end_turn"));
assert_eq!(field(&recorded, "usage")["input_tokens"], json!(30));
assert_eq!(
folded(&s).await,
Some(100),
"the ledger recorded no consumption"
);
}
#[tokio::test]
async fn beats_are_the_callers_word_and_the_kernel_adds_none() {
let mut s = new_session(None).await;
let seq = s.append(response(1)).await.expect("an undeclared beat");
let bare = s
.events(seq, usize::MAX)
.await
.expect("events")
.pop()
.expect("response");
assert_eq!(
bare[FIELD_META].get(FIELD_BEAT),
None,
"the kernel must not invent a beat: {bare}"
);
for event in [
json!({
"kind": "llm_response", "meta": { "beat": "b-1" },
"data": { "content": [], "usage": { "input_tokens": 1 } }
}),
json!({
"kind": "tool_call", "meta": { "beat": "b-1" },
"data": { "call_id": "c1", "name": "sh", "args": {} }
}),
json!({
"kind": "tool_result", "meta": { "beat": "b-1" },
"data": { "call_id": "c1", "ok": true, "result": "ok" }
}),
json!({
"kind": "llm_call_failed", "meta": { "beat": "b-1" },
"data": { "error": "boom" }
}),
] {
let seq = s.append(obj(event.clone())).await.expect("declared beat");
let recorded = s
.events(seq, usize::MAX)
.await
.expect("events")
.pop()
.expect("recorded");
assert_eq!(recorded[FIELD_META][FIELD_BEAT], json!("b-1"), "{event}");
}
let err = s
.append(obj(json!({ "kind": "note", "beat": "b-1" })))
.await
.expect_err("a beat at the top level");
assert!(err.reason().contains("meta.beat"), "{err}");
}
#[tokio::test]
async fn a_closed_session_records_nothing() {
let mut s = new_session(Some(100)).await;
s.close(None).await.expect("close");
let err = s.append(response(10)).await.expect_err("closed session");
assert_eq!(err.reason(), "session is closed");
assert_eq!(
s.len().await.expect("len"),
3,
"session_opened + budget_granted + session_closed only"
);
assert_eq!(remaining(&s).await, Some(100), "nothing was consumed");
}
#[tokio::test]
async fn the_budget_refuses_before_the_call_rather_than_flagging_after_it() {
let mut s = new_session(Some(10)).await;
async fn responses(s: &Session) -> usize {
s.events(0, usize::MAX)
.await
.expect("events")
.iter()
.filter(|e| e.kind() == "llm_response")
.count()
}
assert_eq!(s.reserve(10).await, Ok(true));
s.append(response(25)).await.expect("recorded");
assert_eq!(remaining(&s).await, Some(0), "the reservation took it all");
assert!(exhausted(&s).await);
assert_eq!(s.reserve(1).await, Ok(false));
assert_eq!(responses(&s).await, 1, "the refused beat made no call");
assert_eq!(remaining(&s).await, Some(0));
assert_eq!(folded(&s).await, remaining(&s).await);
s.append(response(5)).await.expect("recorded");
assert_eq!(responses(&s).await, 2, "stopping is the caller's decision");
}
#[tokio::test]
async fn without_a_budget_a_call_reports_no_remaining_and_is_never_exhausted() {
let mut s = new_session(None).await;
s.append(response(9_000)).await.expect("recorded");
assert_eq!(remaining(&s).await, None);
assert!(!exhausted(&s).await);
assert_eq!(s.reserve(1_000_000).await, Ok(true));
assert_eq!(s.spend(1_000_000).await, Ok(()));
assert!(
ledger(&s).await.is_empty(),
"a run with no quota keeps no ledger"
);
assert_eq!(remaining(&s).await, None);
assert!(!exhausted(&s).await);
}
#[tokio::test]
async fn views_stay_readable_and_correct_after_close() {
let mut s = new_session(None).await;
s.append(response(9)).await.expect("recorded");
let before = s
.view(VIEW_TAIL, Some(&obj(json!({ "n": 1 }))))
.await
.expect("tail");
s.close(None).await.expect("close");
let after = s
.view(VIEW_TAIL, Some(&obj(json!({ "n": 1 }))))
.await
.expect("tail after close");
assert_eq!(
kind_of(&before.as_array().expect("array")[0]),
"llm_response"
);
assert_eq!(
kind_of(&after.as_array().expect("array")[0]),
KIND_SESSION_CLOSED
);
assert_eq!(s.len().await.expect("len"), 3);
assert_eq!(
s.events(0, usize::MAX).await.expect("events")[2].kind(),
KIND_SESSION_CLOSED
);
}
#[tokio::test]
async fn open_on_records_the_owner_on_session_opened() {
use crate::knl::SqliteEventStore;
let store = SqliteEventStore::open_memory("owner-stream", &Logs::new())
.await
.expect("open");
let s = Session::open_on(
"user-7".to_string(),
Some(grant(100)),
None,
Box::new(store),
)
.await
.expect("open");
let events = s.events(0, usize::MAX).await.expect("events");
let opened = events.first().expect("session_opened");
assert_eq!(opened.kind(), KIND_SESSION_OPENED);
assert_eq!(
field(opened, FIELD_OWNER).as_str(),
Some("user-7"),
"owner rides on session_opened: {opened}"
);
assert_eq!(s.owner(), "user-7");
assert_eq!(events[1].kind(), KIND_BUDGET_GRANTED);
assert_eq!(*field(&events[1], FIELD_AMOUNT), json!(100));
}
#[tokio::test]
async fn resume_restores_the_owner_and_the_folded_balance() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "resume-stream";
let logs = Logs::new();
let before_close = {
let store = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open");
let mut s = Session::open_on(
"user-42".to_string(),
Some(grant(100)),
None,
Box::new(store),
)
.await
.expect("open");
assert_eq!(s.reserve(30).await, Ok(true));
s.append(response(30)).await.expect("first response");
s.append(obj(
json!({ "kind": "msg_user", "data": { "content": "more" } }),
))
.await
.expect("msg_user");
assert_eq!(s.reserve(15).await, Ok(true));
s.append(response(20)).await.expect("second response");
s.spend(5)
.await
.expect("the second call overran its estimate");
assert_eq!(remaining(&s).await, Some(50), "100 - 30 - 15 - 5");
assert_eq!(folded(&s).await, remaining(&s).await);
remaining(&s).await
};
let store = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen");
let mut resumed = Session::resume(None, Box::new(store))
.await
.expect("resume");
assert_eq!(
resumed.owner(),
"user-42",
"owner restored from session_opened"
);
assert_eq!(
remaining(&resumed).await,
before_close,
"the balance is what the ledger says it was"
);
assert_eq!(
resumed.grant().and_then(|g| g.tag.as_deref()),
Some("tokens"),
"the grant's words came back with it"
);
assert_eq!(
resumed.len().await.expect("len"),
8,
"session_opened + granted + reserved + response + msg_user \
+ reserved + response + spent — and nothing from resume itself"
);
let responses: Vec<Current> = resumed
.events(0, usize::MAX)
.await
.expect("events")
.into_iter()
.filter(|e| e.kind() == "llm_response")
.collect();
assert_eq!(responses.len(), 2, "{responses:?}");
assert_eq!(field(&responses[0], "usage")["input_tokens"], json!(30));
assert_eq!(field(&responses[1], "usage")["input_tokens"], json!(20));
assert_eq!(resumed.reserve(5).await, Ok(true));
let seq = resumed.append(response(5)).await.expect("third response");
let recorded = resumed
.events(seq, usize::MAX)
.await
.expect("events")
.pop()
.expect("llm_response");
assert_eq!(recorded.kind(), "llm_response");
assert_eq!(remaining(&resumed).await, Some(45), "5 reserved off the 50");
assert_eq!(folded(&resumed).await, remaining(&resumed).await);
}
#[tokio::test]
async fn resume_with_a_grant_records_it_and_raises_the_balance() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "regrant-stream";
let logs = Logs::new();
{
let store = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open");
let mut s = Session::open_on(
"user-9".to_string(),
Some(grant(100)),
None,
Box::new(store),
)
.await
.expect("open");
assert_eq!(s.reserve(80).await, Ok(true));
assert_eq!(remaining(&s).await, Some(20));
}
let store = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen");
let mut resumed = Session::resume(
Some(BudgetGrant {
amount: 50,
tag: Some("tokens".to_string()),
desc: Some("a little more".to_string()),
}),
Box::new(store),
)
.await
.expect("resume");
assert_eq!(remaining(&resumed).await, Some(70), "20 left + 50 granted");
assert_eq!(folded(&resumed).await, remaining(&resumed).await);
let moves = ledger(&resumed).await;
assert_eq!(
kinds(&moves),
vec![
KIND_BUDGET_GRANTED,
KIND_BUDGET_RESERVED,
KIND_BUDGET_GRANTED
],
"the second grant is a recorded fact"
);
assert_eq!(*field(&moves[2], FIELD_AMOUNT), json!(50));
assert_eq!(*field(&moves[2], FIELD_DESC), json!("a little more"));
assert_eq!(resumed.reserve(70).await, Ok(true));
assert_eq!(resumed.reserve(1).await, Ok(false));
assert_eq!(remaining(&resumed).await, Some(0));
assert_eq!(folded(&resumed).await, remaining(&resumed).await);
}
#[tokio::test]
async fn a_resume_does_not_give_a_stream_the_budget_it_opened_without() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "ungranted-stream";
let logs = Logs::new();
let store = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open");
let first = Session::open_on("user-1".to_string(), None, None, Box::new(store))
.await
.expect("open");
assert_eq!(
first.len().await.expect("len"),
1,
"session_opened, and no grant beside it"
);
let reopened = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen");
let err = Session::resume(Some(grant(100)), Box::new(reopened))
.await
.expect_err("a resume must not introduce a ledger");
assert_eq!(
err.kind(),
KnlError::VALIDATION,
"the caller's argument is what did not hold up: {err}"
);
assert!(
err.reason().contains("opened with no budget"),
"{}",
err.reason()
);
assert_eq!(
first.len().await.expect("len"),
1,
"the stream is untouched"
);
assert_eq!(remaining(&first).await, None, "and still has no ledger");
let reopened = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen");
let second = Session::resume(None, Box::new(reopened))
.await
.expect("a resume with no grant is the ordinary one");
assert_eq!(remaining(&second).await, None);
assert_eq!(second.len().await.expect("len"), 1);
}
#[tokio::test]
async fn a_handle_opened_without_a_grant_is_bound_by_the_grant_the_log_carries() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "late-grant-stream";
let logs = Logs::new();
let store = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open");
let mut opened_without =
Session::open_on("user-1".to_string(), None, None, Box::new(store))
.await
.expect("open");
let reopened = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen");
let mut owner_handle = Session::resume(None, Box::new(reopened))
.await
.expect("resume");
owner_handle
.grant_more(grant(100))
.await
.expect("the owner grants");
assert_eq!(
opened_without.grant(),
None,
"the cached grant is a hint, and this handle never got one"
);
assert_eq!(
opened_without.reserve(500).await,
Ok(false),
"500 does not fit in the 100 the log carries"
);
assert_eq!(opened_without.reserve(40).await, Ok(true));
assert_eq!(remaining(&opened_without).await, Some(60));
opened_without
.spend(60)
.await
.expect("a deduction on a ledger this handle did not open");
assert_eq!(remaining(&opened_without).await, Some(0));
assert_eq!(
folded(&opened_without).await,
remaining(&opened_without).await
);
let moves = ledger(&opened_without).await;
assert_eq!(
kinds(&moves),
vec![
KIND_BUDGET_GRANTED,
KIND_BUDGET_REFUSED,
KIND_BUDGET_RESERVED,
KIND_BUDGET_SPENT,
],
);
for event in &moves {
assert_eq!(
field(event, FIELD_TAG).as_str(),
Some("tokens"),
"the unit comes off the log's grant: {event}"
);
}
assert_eq!(*field(&moves[1], FIELD_REMAINING), json!(100));
}
async fn seed_an_opening_with_no_scope(path: &std::path::Path, stream: &str) {
use crate::knl::SqliteEventStore;
let logs = Logs::new();
drop(
SqliteEventStore::open(path, stream, &logs)
.await
.expect("open"),
);
assert!(logs.shutdown().await.is_empty(), "the writer joined");
let conn = rusqlite::Connection::open(path).expect("open the database directly");
conn.execute(
"INSERT INTO events \
(stream, seq, epoch_ms, kind, schema_version, meta, data) \
VALUES (?1, 1, 0, ?2, 2, '{}', '{}')",
rusqlite::params![stream, KIND_SESSION_OPENED],
)
.expect("seed the opening");
conn.execute(
"INSERT INTO stream_seq (stream, next_seq) VALUES (?1, 2) \
ON CONFLICT(stream) DO UPDATE SET next_seq = excluded.next_seq",
rusqlite::params![stream],
)
.expect("seed the counter");
}
#[tokio::test]
async fn resume_falls_back_to_anon_when_the_log_has_no_owner() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "legacy-stream";
seed_an_opening_with_no_scope(&path, stream).await;
let store = SqliteEventStore::open(&path, stream, &Logs::new())
.await
.expect("reopen");
let resumed = Session::resume(None, Box::new(store))
.await
.expect("resume");
assert_eq!(resumed.owner(), ANON);
assert_eq!(
remaining(&resumed).await,
None,
"resumed without a budget cap"
);
}
#[tokio::test]
async fn resume_restores_the_scope_id_and_owner_from_the_log() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "scope-resume-stream";
let logs = Logs::new();
let opened_scope = {
let store = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open");
let mut s = Session::open_on(
"user-11".to_string(),
Some(grant(100)),
None,
Box::new(store),
)
.await
.expect("open");
assert_eq!(s.reserve(40).await, Ok(true));
s.scope_id().to_string()
};
let store = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen");
let mut resumed = Session::resume(None, Box::new(store))
.await
.expect("resume");
assert_eq!(
resumed.scope_id(),
opened_scope,
"the scope id is restored from session_opened, not re-issued"
);
assert_eq!(resumed.owner(), "user-11");
assert_eq!(resumed.scope().owner(), "user-11");
assert_eq!(
remaining(&resumed).await,
Some(60),
"the balance is the fold's"
);
assert_eq!(resumed.reserve(10).await, Ok(true));
let last = ledger(&resumed).await.pop().expect("budget_reserved");
assert_eq!(last.kind(), KIND_BUDGET_RESERVED);
assert_eq!(
field(&last, FIELD_SCOPE_ID).as_str(),
Some(opened_scope.as_str()),
"{last}"
);
}
#[tokio::test]
async fn resume_issues_a_fresh_scope_id_when_the_log_records_none() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "legacy-scope-stream";
seed_an_opening_with_no_scope(&path, stream).await;
let store = SqliteEventStore::open(&path, stream, &Logs::new())
.await
.expect("reopen");
let mut resumed = Session::resume(None, Box::new(store))
.await
.expect("resume");
resumed
.grant_more(grant(50))
.await
.expect("the owner grants");
let opened = resumed
.events(0, usize::MAX)
.await
.expect("events")
.remove(0);
assert_eq!(opened.kind(), KIND_SESSION_OPENED);
assert_eq!(data_field(&opened, FIELD_SCOPE_ID), None, "{opened}");
assert_eq!(data_field(&opened, FIELD_OWNER), None, "{opened}");
assert!(
!resumed.scope_id().is_empty(),
"an older log must still resume under a scope"
);
assert_eq!(resumed.owner(), ANON, "the sibling fallback");
let granted = ledger(&resumed).await.pop().expect("budget_granted");
assert_eq!(granted.kind(), KIND_BUDGET_GRANTED);
assert_eq!(
field(&granted, FIELD_SCOPE_ID).as_str(),
Some(resumed.scope_id()),
"{granted}"
);
}
#[tokio::test]
async fn resume_of_an_empty_store_is_a_caller_error_not_an_anon_session() {
let err = Session::resume(Some(grant(100)), Box::new(MemEventStore::new()))
.await
.expect_err("an empty store has no session to resume");
assert!(
err.reason().contains("no session to resume"),
"{}",
err.reason()
);
}
#[tokio::test]
async fn resume_of_a_store_without_an_opening_is_a_caller_error() {
let mut store = MemEventStore::new();
store
.append(obj(json!({ "kind": "note" })))
.await
.expect("seed a non-opening event");
let err = Session::resume(None, Box::new(store))
.await
.expect_err("a log with no opening has no session to resume");
assert!(
err.reason().contains("no session to resume"),
"{}",
err.reason()
);
}
#[tokio::test]
async fn a_closed_stream_is_not_resumed() {
let mut store = MemEventStore::new();
store
.append(obj(json!({
"kind": "session_opened",
"data": { "scope_id": "scope-5", "owner": "user-5" }
})))
.await
.expect("seed the opening");
store
.append(obj(
json!({ "kind": "budget_granted", "data": { "amount": 100 } }),
))
.await
.expect("seed the grant");
store
.append(obj(
json!({ "kind": "session_closed", "data": { "reason": "done" } }),
))
.await
.expect("seed the ending");
let err = Session::resume(None, Box::new(store))
.await
.expect_err("a closed session must not be resumed");
assert!(
err.reason().contains("session is closed"),
"{}",
err.reason()
);
assert!(err.reason().contains("disposable"), "{}", err.reason());
}
struct BusyStore {
inner: MemEventStore,
injected: bool,
}
#[async_trait::async_trait]
impl EventStore for BusyStore {
async fn append(&mut self, event: Map<String, Value>) -> KnlResult<crate::knl::Committed> {
if !self.injected
&& event.get(FIELD_KIND).and_then(Value::as_str) == Some("llm_response")
{
self.injected = true;
self.inner
.append(obj(json!({ "kind": "sneaked_in" })))
.await
.expect("injected concurrent write");
}
self.inner.append(event).await
}
async fn append_if(
&mut self,
kinds: Option<&[&str]>,
decide: crate::knl::Decision,
) -> KnlResult<Option<crate::knl::Committed>> {
self.inner.append_if(kinds, decide).await
}
async fn read_kinds(
&self,
kinds: Option<&[&str]>,
from_seq: u64,
limit: usize,
) -> KnlResult<Vec<Value>> {
self.inner.read_kinds(kinds, from_seq, limit).await
}
async fn head(&self) -> KnlResult<Option<u64>> {
self.inner.head().await
}
async fn len(&self) -> KnlResult<usize> {
self.inner.len().await
}
}
#[tokio::test]
async fn an_append_lands_after_a_competing_write_rather_than_being_refused() {
let store = BusyStore {
inner: MemEventStore::new(),
injected: false,
};
let mut s = Session::open_on("user".to_string(), Some(grant(1000)), None, Box::new(store))
.await
.expect("open");
assert_eq!(
s.len().await.expect("len"),
2,
"session_opened + budget_granted so far"
);
let seq = s
.append(response(10))
.await
.expect("an append is not refused");
assert_eq!(seq, 4, "the seq is where the event really landed");
assert_eq!(s.len().await.expect("len"), 4, "both writes are in the log");
let log = s.events(0, usize::MAX).await.expect("events");
assert_eq!(
kinds(&log),
[
KIND_SESSION_OPENED,
KIND_BUDGET_GRANTED,
"sneaked_in",
"llm_response"
],
"the log interleaves in arrival order"
);
assert_eq!(
remaining(&s).await,
Some(1000),
"an append still charges nothing"
);
}
struct DecidedWritesOnlyStore {
inner: MemEventStore,
armed: Arc<AtomicBool>,
}
#[async_trait::async_trait]
impl EventStore for DecidedWritesOnlyStore {
async fn append(&mut self, event: Map<String, Value>) -> KnlResult<crate::knl::Committed> {
if self.armed.load(Ordering::Relaxed) {
return Err(KnlError::Storage(
"this store takes only what a decision wrote".to_string(),
));
}
self.inner.append(event).await
}
async fn append_if(
&mut self,
kinds: Option<&[&str]>,
decide: crate::knl::Decision,
) -> KnlResult<Option<crate::knl::Committed>> {
self.inner.append_if(kinds, decide).await
}
async fn read_kinds(
&self,
kinds: Option<&[&str]>,
from_seq: u64,
limit: usize,
) -> KnlResult<Vec<Value>> {
self.inner.read_kinds(kinds, from_seq, limit).await
}
async fn head(&self) -> KnlResult<Option<u64>> {
self.inner.head().await
}
async fn len(&self) -> KnlResult<usize> {
self.inner.len().await
}
}
#[tokio::test]
async fn a_refusal_lands_in_the_transaction_that_decided_it() {
let armed = Arc::new(AtomicBool::new(false));
let store = DecidedWritesOnlyStore {
inner: MemEventStore::new(),
armed: Arc::clone(&armed),
};
let mut s = Session::open_on("user".to_string(), Some(grant(10)), None, Box::new(store))
.await
.expect("open");
armed.store(true, Ordering::Relaxed);
assert_eq!(
s.reserve(50).await,
Ok(false),
"the refusal is the decision's own write, so it lands"
);
assert_eq!(
s.reserve(4).await,
Ok(true),
"and so is the reservation that fits"
);
let moves = ledger(&s).await;
assert_eq!(
kinds(&moves),
vec![
KIND_BUDGET_GRANTED,
KIND_BUDGET_REFUSED,
KIND_BUDGET_RESERVED
],
"exactly one entry per decision"
);
assert_eq!(*field(&moves[1], FIELD_AMOUNT), json!(50), "what was asked");
assert_eq!(
*field(&moves[1], FIELD_REMAINING),
json!(10),
"and the balance the decision measured it against"
);
assert_eq!(remaining(&s).await, Some(6), "a refusal moved nothing");
let err = s
.append(obj(json!({ "kind": "note" })))
.await
.expect_err("a plain append");
assert_eq!(err.kind(), KnlError::STORAGE);
}
struct HeadlessStore {
inner: MemEventStore,
}
#[async_trait::async_trait]
impl EventStore for HeadlessStore {
async fn append(&mut self, event: Map<String, Value>) -> KnlResult<crate::knl::Committed> {
self.inner.append(event).await
}
async fn append_if(
&mut self,
kinds: Option<&[&str]>,
decide: crate::knl::Decision,
) -> KnlResult<Option<crate::knl::Committed>> {
self.inner.append_if(kinds, decide).await
}
async fn read_kinds(
&self,
kinds: Option<&[&str]>,
from_seq: u64,
limit: usize,
) -> KnlResult<Vec<Value>> {
self.inner.read_kinds(kinds, from_seq, limit).await
}
async fn head(&self) -> KnlResult<Option<u64>> {
Err(KnlError::Busy("the head read is contended".to_string()))
}
async fn len(&self) -> KnlResult<usize> {
self.inner.len().await
}
}
#[tokio::test]
async fn a_balance_that_cannot_be_read_is_an_error_not_a_stale_fold() {
let store = HeadlessStore {
inner: MemEventStore::new(),
};
let s = Session::open_on("user".to_string(), Some(grant(100)), None, Box::new(store))
.await
.expect("the appends land; only the head read is down");
let err = s
.remaining()
.await
.expect_err("a failed read must not fold into a number");
assert_eq!(err.kind(), KnlError::BUSY, "the class travels out intact");
assert!(err.is_retryable(), "contention is the one retryable class");
let err = s.exhausted().await.expect_err("nor into a boolean");
assert_eq!(err.kind(), KnlError::BUSY);
assert_eq!(
s.len().await.expect("len"),
2,
"session_opened + budget_granted landed"
);
}
#[tokio::test]
async fn resume_of_a_nonexistent_sqlite_stream_is_a_caller_error() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let store = SqliteEventStore::open(&path, "ghost-stream", &Logs::new())
.await
.expect("open");
let err = Session::resume(Some(grant(100)), Box::new(store))
.await
.expect_err("an empty stream has no session to resume");
assert!(
err.reason().contains("no session to resume"),
"{}",
err.reason()
);
}
#[tokio::test]
async fn two_sessions_on_one_stream_both_append_and_the_log_interleaves() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "interleave-stream";
let logs = Logs::new();
let store_a = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open A");
let mut a = Session::open_on(
"user".to_string(),
Some(grant(1000)),
None,
Box::new(store_a),
)
.await
.expect("open A");
let store_b = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open B");
let mut b = Session::resume(None, Box::new(store_b))
.await
.expect("resume B");
assert_eq!(remaining(&b).await, Some(1000), "B resumed on A's ledger");
assert_eq!(
(a.len().await.expect("len"), b.len().await.expect("len")),
(2, 2),
"both see the same two events"
);
assert_eq!(a.append(response(10)).await.expect("A appends"), 3);
assert_eq!(
b.append(response(20)).await.expect("B appends too"),
4,
"B's write lands after A's, rather than being refused"
);
assert_eq!(a.append(response(30)).await.expect("A appends again"), 5);
let verify = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen to verify");
let log = as_current(verify.read(0, usize::MAX).await.expect("read log"));
let responses: Vec<u64> = log
.iter()
.filter(|e| e.kind() == "llm_response")
.map(Current::seq)
.collect();
assert_eq!(responses, [3, 4, 5], "every append landed, in order");
}
#[tokio::test]
async fn two_sessions_cannot_both_reserve_the_same_allowance() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "reserve-race-stream";
let logs = Logs::new();
let store_a = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open A");
let mut a = Session::open_on("user".to_string(), Some(grant(10)), None, Box::new(store_a))
.await
.expect("open A");
let store_b = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open B");
let mut b = Session::resume(None, Box::new(store_b))
.await
.expect("resume B");
assert_eq!(remaining(&b).await, Some(10), "both see the whole grant");
assert_eq!(remaining(&a).await, Some(10));
assert_eq!(a.reserve(6).await, Ok(true), "the first reservation fits");
assert_eq!(b.reserve(6).await, Ok(false), "the second does not");
assert_eq!(remaining(&b).await, Some(4), "B's balance is the ledger's");
assert_eq!(remaining(&a).await, Some(4), "and so is A's");
let verify = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen to verify");
let log = as_current(verify.read(0, usize::MAX).await.expect("read log"));
assert_eq!(fold_balance(&log), Some(4), "no allowance was taken twice");
let moves: Vec<&str> = kinds(&log)
.into_iter()
.filter(|k| k.starts_with("budget_"))
.collect();
assert_eq!(
moves,
[
KIND_BUDGET_GRANTED,
KIND_BUDGET_RESERVED,
KIND_BUDGET_REFUSED
],
"one grant, one reservation, one refusal"
);
let refused = log.last().expect("the refusal");
assert_eq!(*field(refused, FIELD_AMOUNT), json!(6));
assert_eq!(
*field(refused, FIELD_REMAINING),
json!(4),
"what there really was"
);
}
#[tokio::test]
async fn a_close_is_the_handles_and_the_log_records_what_arrives_after_it() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "close-race-stream";
let logs = Logs::new();
let store_a = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open A");
let mut a = Session::open_on(
"user".to_string(),
Some(grant(100)),
None,
Box::new(store_a),
)
.await
.expect("open A");
let store_b = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open B");
let mut b = Session::resume(None, Box::new(store_b))
.await
.expect("resume B");
let store_c = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open C");
let mut c = Session::resume(None, Box::new(store_c))
.await
.expect("resume C");
a.close(Some("done")).await.expect("A closes");
assert!(a.is_closed());
assert!(!b.is_closed(), "B's flag is its own and has not moved");
assert!(!c.is_closed());
assert_eq!(
b.append(obj(json!({ "kind": "note" }))).await,
Ok(4),
"an append after another handle's close is recorded"
);
assert!(!b.is_closed(), "landing a write closed nothing");
assert_eq!(c.reserve(5).await, Ok(true), "the ledger covers it");
assert_eq!(c.spend(10).await, Ok(()), "the settlement lands");
assert_eq!(
remaining(&c).await,
Some(85),
"100 − 5 − 10, folded in the tx"
);
assert!(!c.is_closed());
b.close(Some("late")).await.expect("B closes");
a.close(Some("again"))
.await
.expect("A is idempotent per handle");
let verify = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen to verify");
let log = as_current(verify.read(0, usize::MAX).await.expect("read log"));
assert_eq!(
kinds(&log),
[
KIND_SESSION_OPENED,
KIND_BUDGET_GRANTED,
KIND_SESSION_CLOSED,
"note",
KIND_BUDGET_RESERVED,
KIND_BUDGET_SPENT,
KIND_SESSION_CLOSED,
],
"everything that happened, in the order it arrived"
);
let endings: Vec<&Current> = log
.iter()
.filter(|event| event.kind() == KIND_SESSION_CLOSED)
.collect();
assert_eq!(endings.len(), 2, "two handles closed, two endings recorded");
assert_eq!(*field(endings[0], FIELD_REASON), json!("done"));
assert_eq!(*field(endings[1], FIELD_REASON), json!("late"));
assert_eq!(
fold_balance(&log),
Some(85),
"the ledger is what the moves that landed add up to"
);
}
#[tokio::test]
async fn a_settlement_records_the_move_and_the_balance_is_the_ledgers() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "spend-race-stream";
let logs = Logs::new();
let store_a = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open A");
let mut a = Session::open_on(
"user".to_string(),
Some(grant(100)),
None,
Box::new(store_a),
)
.await
.expect("open A");
let store_b = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open B");
let mut b = Session::resume(None, Box::new(store_b))
.await
.expect("resume B");
assert_eq!(
(remaining(&a).await, remaining(&b).await),
(Some(100), Some(100))
);
assert_eq!(a.spend(30).await, Ok(()), "A settles 30 of the 100");
assert_eq!(
remaining(&a).await,
Some(70),
"and reads the balance separately"
);
assert_eq!(
remaining(&b).await,
Some(70),
"B wrote nothing and still reads A's settlement off the ledger"
);
assert_eq!(b.spend(20).await, Ok(()), "B settles 20");
assert_eq!(remaining(&b).await, Some(50), "both settlements are in it");
let verify = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen to verify");
let log = as_current(verify.read(0, usize::MAX).await.expect("read log"));
assert_eq!(fold_balance(&log), Some(50), "the balance is the fold");
let moves: Vec<&str> = kinds(&log)
.into_iter()
.filter(|kind| kind.starts_with("budget_"))
.collect();
assert_eq!(
moves,
[KIND_BUDGET_GRANTED, KIND_BUDGET_SPENT, KIND_BUDGET_SPENT],
"one grant and two settlements"
);
assert_eq!(a.spend(1_000).await, Ok(()));
assert_eq!(b.spend(1).await, Ok(()));
assert_eq!(remaining(&a).await, Some(0));
assert_eq!(remaining(&b).await, Some(0));
let verify = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen to verify");
assert_eq!(
fold_balance(&as_current(
verify.read(0, usize::MAX).await.expect("read log")
)),
Some(0),
"the ledger floors at zero rather than going into debt"
);
}
#[tokio::test]
async fn a_handle_that_wrote_nothing_reports_what_the_other_spent() {
use crate::knl::SqliteEventStore;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "shared-balance-stream";
let logs = Logs::new();
let store_a = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open A");
let mut a = Session::open_on(
"user".to_string(),
Some(grant(100)),
None,
Box::new(store_a),
)
.await
.expect("open A");
let store_b = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open B");
let b = Session::resume(None, Box::new(store_b))
.await
.expect("resume B");
assert_eq!(
remaining(&b).await,
Some(100),
"both start on the same ledger"
);
assert_eq!(a.spend(30).await, Ok(()), "A settles 30");
assert_eq!(
remaining(&b).await,
Some(70),
"B sees the settlement it did not make"
);
assert_eq!(a.reserve(20).await, Ok(true), "A reserves 20");
assert_eq!(remaining(&b).await, Some(50), "and the reservation too");
assert!(!exhausted(&b).await);
assert_eq!(
remaining(&b).await,
Some(50),
"a second read is the same read"
);
assert_eq!(a.spend(1_000).await, Ok(()), "A overspends");
assert_eq!(remaining(&b).await, Some(0), "the floor is the ledger's");
assert!(exhausted(&b).await);
let verify = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen to verify");
assert_eq!(
fold_balance(&as_current(
verify.read(0, usize::MAX).await.expect("read log")
)),
remaining(&b).await
);
}
struct RenameLegacyKinds;
impl crate::knl::Upcaster for RenameLegacyKinds {
fn upcast(&self, mut event: Value) -> Value {
let Some(map) = event.as_object_mut() else {
return event;
};
let renamed = match map.get(FIELD_KIND).and_then(Value::as_str) {
Some("legacy_opened") => Some(KIND_SESSION_OPENED),
Some("legacy_response") => Some("llm_response"),
_ => None,
};
if let Some(kind) = renamed {
map.insert(FIELD_KIND.to_string(), Value::from(kind));
}
event
}
}
#[tokio::test]
async fn a_session_reads_every_path_through_the_upcaster_seam() {
use crate::knl::{
SqliteEventStore, Upcaster, CURRENT_SCHEMA_VERSION, SCHEMA_VERSION_FIELD,
};
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "seam-stream";
let logs = Logs::new();
{
let mut store = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open");
store
.append(obj(json!({
"kind": "legacy_opened",
"data": { "owner": "user-3", "scope_id": "scope-from-the-log" }
})))
.await
.expect("the opening");
store
.append(obj(json!({
"kind": "budget_granted",
"data": { "amount": 100, "tag": "tokens" }
})))
.await
.expect("the grant");
store
.append(obj(json!({
"kind": "legacy_response", "meta": { "beat": "b-1" },
"data": {
"content": [{ "type": "text", "text": "ok" }],
"usage": { "input_tokens": 7 }
}
})))
.await
.expect("the response");
}
let chain: Vec<Arc<dyn Upcaster>> = vec![Arc::new(RenameLegacyKinds)];
let seamed = CurrentStore::new(
Box::new(
SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen"),
),
chain,
);
let mut resumed = Session::resume_on(None, seamed)
.await
.expect("resume through the seam");
assert_eq!(resumed.owner(), "user-3", "the owner the step revealed");
assert_eq!(resumed.scope_id(), "scope-from-the-log");
assert_eq!(
resumed.grant().and_then(|g| g.tag.as_deref()),
Some("tokens"),
"and the grant with it"
);
let log = resumed.events(0, usize::MAX).await.expect("events");
assert_eq!(
kinds(&log),
[KIND_SESSION_OPENED, KIND_BUDGET_GRANTED, "llm_response"],
"every read is projected"
);
let tail = resumed
.view(VIEW_TAIL, Some(&obj(json!({ "n": 1 }))))
.await
.expect("tail");
let last = &tail.as_array().expect("array")[0];
assert_eq!(
kind_of(last),
"llm_response",
"the view read the projected kind: {last}"
);
assert_eq!(last[FIELD_DATA]["usage"]["input_tokens"], json!(7));
assert_eq!(
remaining(&resumed).await,
Some(100),
"the balance folds too"
);
let raw = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen raw");
let stored = raw.read(0, usize::MAX).await.expect("read raw");
assert_eq!(kind_of(&stored[0]), "legacy_opened", "{}", stored[0]);
assert_eq!(kind_of(&stored[2]), "legacy_response", "{}", stored[2]);
assert!(
raw.read_kinds(Some(&[KIND_SESSION_OPENED]), 0, usize::MAX)
.await
.expect("read raw")
.is_empty(),
"the kind filter selects on what is stored"
);
assert_eq!(
stored[0].get(SCHEMA_VERSION_FIELD).and_then(Value::as_u64),
Some(CURRENT_SCHEMA_VERSION),
"an untouched row keeps the version it was written under: {}",
stored[0]
);
}
#[tokio::test]
async fn a_closed_stream_seen_through_the_seam_is_still_refused() {
use crate::knl::{SqliteEventStore, Upcaster};
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let stream = "seam-closed-stream";
let logs = Logs::new();
{
let mut store = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open");
store
.append(obj(json!({
"kind": "legacy_opened",
"data": { "owner": "user-3", "scope_id": "scope-from-the-log" }
})))
.await
.expect("the opening");
store
.append(obj(json!({
"kind": "session_closed", "data": { "reason": "done" }
})))
.await
.expect("the ending");
}
let chain: Vec<Arc<dyn Upcaster>> = vec![Arc::new(RenameLegacyKinds)];
let seamed = CurrentStore::new(
Box::new(
SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen"),
),
chain,
);
let err = Session::resume_on(None, seamed)
.await
.expect_err("a stream that ended must not be resumed");
assert!(
err.reason().contains("session is closed"),
"{}",
err.reason()
);
}
#[tokio::test]
async fn an_in_memory_stream_is_resumable_while_it_is_open() {
let logs = Logs::new();
let mut s = Session::new(ANON.to_string(), Some(grant(100)), None, &logs)
.await
.expect("open");
assert_eq!(s.reserve(30).await, Ok(true));
s.append(obj(json!({ "kind": "note", "data": { "text": "hi" } })))
.await
.expect("append");
let store = SqliteEventStore::open_memory(s.id(), &logs)
.await
.expect("reopen the stream");
let resumed = Session::resume(None, Box::new(store))
.await
.expect("resume");
assert_eq!(resumed.owner(), ANON);
assert_eq!(remaining(&resumed).await, Some(70), "the ledger came back");
assert_eq!(
kinds(&resumed.events(0, usize::MAX).await.expect("events")),
vec![
KIND_SESSION_OPENED,
KIND_BUDGET_GRANTED,
KIND_BUDGET_RESERVED,
"note"
]
);
let other = Session::new(ANON.to_string(), Some(grant(100)), None, &logs)
.await
.expect("open");
assert_ne!(other.id(), s.id());
assert_eq!(other.database(), s.database(), "one database, two streams");
assert_eq!(
other.len().await.expect("len"),
2,
"opened + granted, and no note"
);
}
#[tokio::test]
async fn a_session_reads_its_own_log_with_sql() {
let mut s = new_session(None).await;
s.append(obj(
json!({ "kind": "msg_user", "data": { "content": "hi" } }),
))
.await
.expect("append");
s.append(response(9)).await.expect("recorded");
let found = s
.query(
"SELECT kind, seq FROM events WHERE stream = $stream ORDER BY seq",
QueryParams::None,
&QueryOpts::default(),
)
.await
.expect("query");
assert!(!found.truncated);
let kinds: Vec<&str> = found
.rows
.iter()
.map(|row| row["kind"].as_str().expect("a kind"))
.collect();
assert_eq!(kinds, [KIND_SESSION_OPENED, "msg_user", "llm_response"]);
let counted = s
.query(
"SELECT kind, COUNT(*) AS n FROM events WHERE stream = $stream \
GROUP BY kind ORDER BY kind",
QueryParams::None,
&QueryOpts::default(),
)
.await
.expect("query");
assert_eq!(counted.rows.len(), 3);
let mut other = new_session(None).await;
other
.append(obj(json!({ "kind": "only_theirs" })))
.await
.expect("append");
let mine = s
.query(
"SELECT kind FROM events WHERE stream IN $sessions",
QueryParams::None,
&QueryOpts::default(),
)
.await
.expect("query");
assert!(
!mine.rows.iter().any(|row| row["kind"] == "only_theirs"),
"{:?}",
mine.rows
);
s.close(None).await.expect("close");
assert!(s
.query("SELECT 1 AS one", QueryParams::None, &QueryOpts::default())
.await
.is_ok());
}
async fn store_beside(parent: &Session, stream: &str) -> Box<dyn EventStore> {
let db = parent.database().expect("the parent is on a database");
let log = test_logs()
.database(db)
.await
.expect("the parent's own log");
Box::new(SqliteEventStore::on(log, stream))
}
fn stream_id() -> String {
uuid::Uuid::new_v4().to_string()
}
async fn another_handle(of: &Session) -> Session {
let db = of.database().expect("a database");
let log = test_logs().database(db).await.expect("the session's log");
let store = SqliteEventStore::on(log, of.id());
let mut handle = Session::resume(None, Box::new(store))
.await
.expect("resume");
handle.adopt_id(of.id().to_string());
handle
}
async fn file_session(path: &std::path::Path, budget: i64, logs: &Logs) -> Session {
let stream = stream_id();
let store = SqliteEventStore::open(path, stream.clone(), logs)
.await
.expect("open the stream");
let mut session =
Session::open_on(ANON.to_string(), Some(grant(budget)), None, Box::new(store))
.await
.expect("open");
session.adopt_id(stream);
session
}
#[tokio::test]
async fn an_allocation_opens_the_child_and_moves_the_units() {
let mut parent = new_session(Some(100)).await;
let stream = stream_id();
let store = store_beside(&parent, &stream).await;
let mut child = parent
.open_child(
stream.clone(),
"user-42".to_string(),
Allocation::new(40),
None,
store,
)
.await
.expect("the parent's balance covers it");
assert_eq!(child.id(), stream, "the child is the stream it was given");
assert_eq!(child.owner(), "user-42");
assert_eq!(remaining(&parent).await, Some(60), "the parent paid");
assert_eq!(remaining(&child).await, Some(40), "and the child holds it");
let opened = child.events(0, usize::MAX).await.expect("events");
assert_eq!(
kinds(&opened),
vec![KIND_SESSION_OPENED, KIND_BUDGET_GRANTED]
);
assert_eq!(
field(&opened[0], FIELD_PARENT).as_str(),
Some(parent.id()),
"{}",
opened[0]
);
assert_eq!(
field(&opened[0], FIELD_SCOPE_ID).as_str(),
Some(child.scope_id())
);
assert_eq!(*field(&opened[1], FIELD_AMOUNT), json!(40));
assert_eq!(field(&opened[1], FIELD_PARENT).as_str(), Some(parent.id()));
assert_eq!(
field(&opened[1], FIELD_TAG).as_str(),
Some("tokens"),
"the child counts in the parent's unit unless it was renamed"
);
let moves = ledger(&parent).await;
assert_eq!(
kinds(&moves),
vec![KIND_BUDGET_GRANTED, KIND_BUDGET_RESERVED]
);
assert_eq!(*field(&moves[1], FIELD_AMOUNT), json!(40));
assert_eq!(
field(&moves[1], FIELD_CHILD).as_str(),
Some(stream.as_str())
);
child.close(Some("done")).await.expect("close the child");
assert_eq!(
remaining(&parent).await,
Some(60),
"an allocation is a spend"
);
}
#[tokio::test]
async fn an_allocation_may_rename_the_unit_for_the_child() {
let mut parent = new_session(Some(100)).await;
let stream = stream_id();
let store = store_beside(&parent, &stream).await;
let child = parent
.open_child(
stream,
"user-42".to_string(),
Allocation {
amount: 10,
tag: Some("turns".to_string()),
},
None,
store,
)
.await
.expect("the allocation");
let opened = child.events(0, usize::MAX).await.expect("events");
assert_eq!(field(&opened[1], FIELD_TAG).as_str(), Some("turns"));
let moves = ledger(&parent).await;
assert_eq!(
field(&moves[1], FIELD_TAG).as_str(),
Some("tokens"),
"the parent's entry counts what the parent counts"
);
}
#[tokio::test]
async fn an_allocation_the_balance_cannot_cover_is_refused_and_recorded() {
let mut parent = new_session(Some(10)).await;
let stream = stream_id();
let store = store_beside(&parent, &stream).await;
let err = parent
.open_child(
stream.clone(),
"user-42".to_string(),
Allocation::new(40),
None,
store,
)
.await
.expect_err("10 does not cover 40");
assert_eq!(err.kind(), KnlError::REFUSED, "{err}");
assert!(
!err.is_retryable(),
"the same balance answers the same: {err}"
);
assert!(err.reason().contains("40"), "{err}");
assert_eq!(remaining(&parent).await, Some(10));
let moves = ledger(&parent).await;
assert_eq!(
kinds(&moves),
vec![KIND_BUDGET_GRANTED, KIND_BUDGET_REFUSED]
);
assert_eq!(*field(&moves[1], FIELD_AMOUNT), json!(40));
assert_eq!(*field(&moves[1], FIELD_REMAINING), json!(10));
assert_eq!(
field(&moves[1], FIELD_CHILD).as_str(),
Some(stream.as_str())
);
let unused = store_beside(&parent, &stream).await;
assert_eq!(unused.len().await.expect("len"), 0);
}
#[tokio::test]
async fn a_parent_with_no_budget_allocates_without_a_balance_to_measure() {
let mut parent = new_session(None).await;
let stream = stream_id();
let store = store_beside(&parent, &stream).await;
let child = parent
.open_child(
stream,
"user-42".to_string(),
Allocation::new(7),
None,
store,
)
.await
.expect("there is no balance to refuse against");
assert_eq!(remaining(&parent).await, None, "still no budget here");
assert_eq!(remaining(&child).await, Some(7));
}
#[tokio::test]
async fn a_child_on_another_database_is_refused() {
let mut parent = new_session(Some(100)).await;
let stranger = stream_id();
let elsewhere = SqliteEventStore::open_memory(stranger.clone(), &Logs::new())
.await
.expect("another in-memory database");
let err = parent
.open_child(
stranger,
"user-42".to_string(),
Allocation::new(10),
None,
Box::new(elsewhere),
)
.await
.expect_err("that is a different log");
assert_eq!(err.kind(), KnlError::VALIDATION, "{err}");
assert!(err.reason().contains("one log"), "{err}");
assert_eq!(
parent.len().await.expect("len"),
2,
"opened + granted, and nothing else"
);
let mut single = Session::open_on(
ANON.to_string(),
Some(grant(50)),
None,
Box::new(MemEventStore::new()),
)
.await
.expect("open");
let err = single
.open_child(
stream_id(),
"user-42".to_string(),
Allocation::new(1),
None,
Box::new(MemEventStore::new()),
)
.await
.expect_err("no database to open a child on");
assert_eq!(err.kind(), KnlError::VALIDATION, "{err}");
}
#[tokio::test]
async fn a_child_does_not_open_on_a_stream_that_already_has_events() {
let mut parent = new_session(Some(100)).await;
let taken = stream_id();
let store = store_beside(&parent, &taken).await;
parent
.open_child(
taken.clone(),
"user-42".to_string(),
Allocation::new(10),
None,
store,
)
.await
.expect("the first allocation");
let before = parent.len().await.expect("len");
assert_eq!(remaining(&parent).await, Some(90));
let store = store_beside(&parent, &taken).await;
let err = parent
.open_child(
taken.clone(),
"user-42".to_string(),
Allocation::new(10),
None,
store,
)
.await
.expect_err("that stream is a session already");
assert_eq!(err.kind(), KnlError::VALIDATION, "{err}");
assert!(err.reason().contains("already has events"), "{err}");
assert_eq!(parent.len().await.expect("len"), before);
assert_eq!(remaining(&parent).await, Some(90));
let occupied = store_beside(&parent, &taken).await;
assert_eq!(occupied.len().await.expect("len"), 2);
let stranger = stream_id();
let mut seeded = store_beside(&parent, &stranger).await;
seeded.append(response(1)).await.expect("seed the stream");
let store = store_beside(&parent, &stranger).await;
let err = parent
.open_child(
stranger,
"user-42".to_string(),
Allocation::new(10),
None,
store,
)
.await
.expect_err("something is written there already");
assert_eq!(err.kind(), KnlError::VALIDATION, "{err}");
assert_eq!(parent.len().await.expect("len"), before);
assert_eq!(seeded.len().await.expect("len"), 1, "and nothing was added");
assert_eq!(
kinds(&ledger(&parent).await),
vec![KIND_BUDGET_GRANTED, KIND_BUDGET_RESERVED]
);
let fresh = stream_id();
let store = store_beside(&parent, &fresh).await;
let child = parent
.open_child(
fresh,
"user-42".to_string(),
Allocation::new(10),
None,
store,
)
.await
.expect("an empty stream is what a child opens on");
assert_eq!(remaining(&child).await, Some(10));
assert_eq!(remaining(&parent).await, Some(80));
}
#[tokio::test]
async fn a_closed_parent_opens_no_child() {
let mut parent = new_session(Some(100)).await;
let mut other = another_handle(&parent).await;
parent.close(Some("done")).await.expect("close");
let stream = stream_id();
let store = store_beside(&parent, &stream).await;
let err = parent
.open_child(
stream,
"user-42".to_string(),
Allocation::new(1),
None,
store,
)
.await
.expect_err("this handle closed");
assert_eq!(err.kind(), KnlError::CLOSED, "{err}");
let stream = stream_id();
let store = store_beside(&other, &stream).await;
let err = other
.open_child(
stream.clone(),
"user-42".to_string(),
Allocation::new(1),
None,
store,
)
.await
.expect_err("the log carries an ending");
assert_eq!(err.kind(), KnlError::CLOSED, "{err}");
let unused = store_beside(&other, &stream).await;
assert_eq!(unused.len().await.expect("len"), 0, "nothing was opened");
}
#[tokio::test]
async fn a_close_records_the_children_that_had_not_ended() {
let mut parent = new_session(Some(100)).await;
let still_open = stream_id();
let store = store_beside(&parent, &still_open).await;
let _running = parent
.open_child(
still_open.clone(),
"user-42".to_string(),
Allocation::new(10),
None,
store,
)
.await
.expect("the allocation");
let ended = stream_id();
let store = store_beside(&parent, &ended).await;
let mut done = parent
.open_child(
ended,
"user-42".to_string(),
Allocation::new(10),
None,
store,
)
.await
.expect("the allocation");
done.close(Some("done")).await.expect("close the child");
parent.close(Some("done")).await.expect("close");
let boundary = parent
.events(0, usize::MAX)
.await
.expect("events")
.pop()
.expect("the boundary");
assert_eq!(boundary.kind(), KIND_SESSION_CLOSED);
assert_eq!(
*field(&boundary, FIELD_OPEN_CHILDREN),
json!([still_open]),
"the child that had ended is not among them: {boundary}"
);
}
#[tokio::test]
async fn a_close_with_no_open_children_records_no_such_field() {
let mut s = new_session(None).await;
s.close(Some("done")).await.expect("close");
let boundary = s
.events(0, usize::MAX)
.await
.expect("events")
.pop()
.expect("the boundary");
assert_eq!(
data_field(&boundary, FIELD_OPEN_CHILDREN),
None,
"{boundary}"
);
}
#[tokio::test]
async fn two_children_allocating_at_once_never_over_allocate() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let logs = test_logs();
let mut one = file_session(&path, 100, &logs).await;
let mut two = another_handle(&one).await;
let first = stream_id();
let second = stream_id();
let first_store = store_beside(&one, &first).await;
let second_store = store_beside(&two, &second).await;
let (a, b) = tokio::join!(
one.open_child(
first,
"child-a".to_string(),
Allocation::new(60),
None,
first_store
),
two.open_child(
second,
"child-b".to_string(),
Allocation::new(60),
None,
second_store
),
);
let granted: i64 = [&a, &b].iter().filter(|outcome| outcome.is_ok()).count() as i64 * 60;
assert_eq!(granted, 60, "exactly one allocation may land");
assert!(
granted <= 100,
"the sum of the grants is within the balance"
);
let refused = match (&a, &b) {
(Err(e), Ok(_)) | (Ok(_), Err(e)) => e,
_ => panic!("one grant and one refusal, got {a:?} / {b:?}"),
};
assert_eq!(refused.kind(), KnlError::REFUSED, "{refused}");
assert_eq!(remaining(&one).await, Some(40));
assert_eq!(
kinds(&ledger(&one).await),
vec![
KIND_BUDGET_GRANTED,
KIND_BUDGET_RESERVED,
KIND_BUDGET_REFUSED
]
);
}
#[test]
fn an_allocation_folds_the_ledger_and_looks_for_the_ending() {
for kind in BUDGET_KINDS {
assert!(
ALLOCATION_KINDS.contains(kind),
"the ledger's {kind} must reach an allocation's decision"
);
}
assert!(ALLOCATION_KINDS.contains(&KIND_SESSION_CLOSED));
assert_eq!(ALLOCATION_KINDS.len(), BUDGET_KINDS.len() + 1);
}
}