use std::collections::{HashMap, HashSet, VecDeque};
use std::hash::{Hash, Hasher};
use std::path::Path;
use std::sync::Arc;
use futures_util::stream::BoxStream;
use futures_util::StreamExt;
use hotl_provider::{retry, CachePolicy, ProviderError, SamplingRequest, StreamEvent, ToolDef};
use hotl_tools::rules::Verdict;
use hotl_tools::{Permission, ToolOutcome};
use hotl_types::{
assistant_text, assistant_tool_uses, EntryPayload, Item, StopReason, SyntheticReason,
TokenUsage, ToolResultItem, ToolUse,
};
use serde_json::Value;
use tokio::sync::{mpsc, oneshot};
use tokio_util::sync::CancellationToken;
use crate::actor::SharedDeps;
use crate::ledger::Phase;
use crate::{AskReply, EngineEvent, Outcome, SessionCmd, TurnEnd};
const COMPACT_TRIGGER: f64 = 0.8;
const SPECULATE_TRIGGER: f64 = 0.6;
const SPECULATION_WAIT: std::time::Duration = std::time::Duration::from_secs(30);
const TURN_EXTENSION_MAX: u32 = 3;
const TODO_GATE_MAX: u32 = 2;
pub(crate) async fn run(
shared: Arc<SharedDeps>,
cmd_tx: mpsc::Sender<SessionCmd>,
events: mpsc::Sender<EngineEvent>,
cancel: CancellationToken,
cont: crate::TurnContinuation,
) {
let mut turn = Turn::new(shared, cmd_tx.clone(), events, cancel, cont);
let end = turn.drive().await;
let end = seal_end(end, turn.pipeline.drain().await);
let usage = turn.usage;
let report = turn.ledger.summary(crate::ledger::max_rss_bytes());
let _ = turn.events.send(EngineEvent::LedgerReport(report)).await;
let _ = cmd_tx.send(SessionCmd::TurnFinished { end, usage }).await;
}
fn seal_end(end: TurnEnd, commit: Commit) -> TurnEnd {
if commit.ok() {
return end;
}
match end {
TurnEnd::Outcome(Outcome::Cancelled) => TurnEnd::Outcome(Outcome::Cancelled),
TurnEnd::Outcome(Outcome::Error { message }) => {
TurnEnd::Outcome(Outcome::Error { message })
}
TurnEnd::Outcome(_) => TurnEnd::Outcome(commit.outcome()),
compact @ TurnEnd::Compact { .. } => compact,
}
}
enum Gate {
Ready {
input: Value,
summary: String,
},
Resolved {
outcome: ToolOutcome,
chargeable: bool,
},
}
struct Executed {
outcome: ToolOutcome,
chargeable: bool,
}
const ACK_WINDOW: usize = 16;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Commit {
Committed,
Sealed,
Aborted,
Unconfirmed,
StaleEpoch,
Gone,
}
impl Commit {
fn from_reply(reply: Option<crate::ProposeReply>) -> Self {
match reply {
Some(crate::ProposeReply::Committed) => Commit::Committed,
Some(crate::ProposeReply::Sealed) => Commit::Sealed,
Some(crate::ProposeReply::StaleEpoch) => Commit::StaleEpoch,
Some(crate::ProposeReply::Ticket(_)) => {
debug_assert!(
false,
"a ticket is not a commit: a Pipelined reply must go through TicketPipeline"
);
Commit::Unconfirmed
}
None => Commit::Gone,
}
}
fn ok(self) -> bool {
matches!(self, Commit::Committed)
}
fn and(self, next: Commit) -> Commit {
if self.ok() {
next
} else {
self
}
}
fn outcome(self) -> Outcome {
match self {
Commit::Aborted => Outcome::Cancelled,
other => Outcome::Error {
message: other.message().into(),
},
}
}
fn message(self) -> &'static str {
match self {
Commit::Committed => "committed",
Commit::Aborted => "the turn was superseded and will be restarted",
Commit::Unconfirmed => {
"an internal error left a commit unconfirmed — nothing can be trusted \
as recorded; start a new session to keep working"
}
Commit::Sealed => {
"the session log is sealed — nothing further can be recorded; \
start a new session to keep working"
}
Commit::StaleEpoch => {
"the masking rules changed twice while committing this entry; \
start a new session to keep working"
}
Commit::Gone => "the session is shutting down",
}
}
}
const PREPARE_OFFLOAD_BYTES: usize = 256 * 1024;
fn payload_size_estimate(payload: &EntryPayload) -> usize {
match payload {
EntryPayload::Item { item } => item_size_estimate(item),
_ => 0,
}
}
fn item_size_estimate(item: &Item) -> usize {
match item {
Item::User { text, .. } | Item::System { text } => text.len(),
Item::ToolResults { results } => results.iter().map(|r| r.content.len()).sum(),
Item::Assistant { blocks } => blocks.iter().map(value_size_estimate).sum(),
Item::Unknown => 0,
}
}
fn value_size_estimate(value: &Value) -> usize {
match value {
Value::String(s) => s.len(),
Value::Array(items) => items.iter().map(value_size_estimate).sum(),
Value::Object(map) => map.values().map(value_size_estimate).sum(),
Value::Null | Value::Bool(_) | Value::Number(_) => 0,
}
}
#[derive(Debug)]
enum PrepareError {
Serialize,
Offload,
}
impl From<serde_json::Error> for PrepareError {
fn from(_: serde_json::Error) -> Self {
PrepareError::Serialize
}
}
async fn prepare_entry(
payload: &EntryPayload,
masker: &Arc<hotl_store::Masker>,
rules_epoch: u32,
) -> Result<crate::PreparedEntry, PrepareError> {
let item = match payload {
EntryPayload::Item { item } => Some(item.clone()),
_ => None,
};
let kind = hotl_store::EntryKind::from(payload);
let bytes = if payload_size_estimate(payload) >= PREPARE_OFFLOAD_BYTES {
let masker = Arc::clone(masker);
let owned = payload.clone();
tokio::task::spawn_blocking(move || -> Result<hotl_store::MaskedBytes, PrepareError> {
let json = hotl_store::serialize_payload(&owned)?;
Ok(hotl_store::mask_json(&masker, json))
})
.await
.map_err(|_| PrepareError::Offload)??
} else {
let json = hotl_store::serialize_payload(payload)?;
hotl_store::mask_json(masker, json)
};
Ok(crate::PreparedEntry::new(
hotl_store::PreparedPayload::new(bytes, kind, rules_epoch),
item,
))
}
fn parallel_chunks<'a>(uses: &'a [ToolUse], registry: &hotl_tools::Registry) -> Vec<&'a [ToolUse]> {
let safe = |tu: &ToolUse| registry.get(&tu.name).is_some_and(|t| t.parallel_safe());
let mut chunks = Vec::new();
let mut start = 0;
while start < uses.len() {
let mut end = start + 1;
if safe(&uses[start]) {
while end < uses.len() && safe(&uses[end]) {
end += 1;
}
}
chunks.push(&uses[start..end]);
start = end;
}
chunks
}
fn spawn_speculation(
shared: &Arc<SharedDeps>,
snapshot: &Arc<Vec<Item>>,
cancel: &CancellationToken,
) -> Option<tokio::task::JoinHandle<Option<crate::SpecDigest>>> {
let tail_budget = (shared.config.context_window as f64 * crate::actor::TAIL_RATIO) as u64;
let plan = hotl_context::compaction::plan(snapshot, tail_budget)?;
let shared = Arc::clone(shared);
let snapshot = Arc::clone(snapshot);
let cancel = cancel.clone();
Some(tokio::spawn(async move {
let folded = &snapshot[plan.prefix_end..plan.kept_from];
let text = tokio::select! {
biased;
_ = cancel.cancelled() => None,
text = crate::actor::summarize(&shared, folded) => text,
}?;
Some(crate::SpecDigest {
prefix_end: plan.prefix_end,
kept_from: plan.kept_from,
text,
})
}))
}
async fn await_speculation(
handle: tokio::task::JoinHandle<Option<crate::SpecDigest>>,
cancel: &CancellationToken,
wait: std::time::Duration,
) -> Option<crate::SpecDigest> {
let abort = handle.abort_handle();
tokio::select! {
biased;
_ = cancel.cancelled() => { abort.abort(); None }
_ = tokio::time::sleep(wait) => { abort.abort(); None }
digest = handle => digest.ok().flatten(),
}
}
#[derive(Default)]
struct TicketPipeline {
tickets: VecDeque<crate::CommitTicket>,
last: Option<(String, u64)>,
}
impl TicketPipeline {
fn is_empty(&self) -> bool {
self.tickets.is_empty()
}
fn len(&self) -> usize {
self.tickets.len()
}
fn ack_seq(&self) -> Option<u64> {
self.last.as_ref().map(|(_, seq)| *seq)
}
fn leaf(&self) -> Option<&str> {
self.last.as_ref().map(|(id, _)| id.as_str())
}
async fn submit(&mut self, ticket: crate::CommitTicket) -> Commit {
let mut commit = Commit::Committed;
while self.len() >= ACK_WINDOW {
commit = commit.and(self.resolve_oldest().await);
}
self.last = Some((ticket.id.clone(), ticket.seq));
self.tickets.push_back(ticket);
commit
}
async fn drain(&mut self) -> Commit {
let mut commit = Commit::Committed;
while !self.tickets.is_empty() {
commit = commit.and(self.resolve_oldest().await);
}
commit
}
async fn resolve_oldest(&mut self) -> Commit {
let Some(ticket) = self.tickets.pop_front() else {
return Commit::Committed;
};
match ticket.ack.await {
Ok(Ok(_)) => Commit::Committed,
Ok(Err(crate::CommitFailed::LogSealed)) => Commit::Sealed,
Ok(Err(crate::CommitFailed::Aborted)) => Commit::Aborted,
Err(_) => Commit::Gone,
}
}
}
struct Speculation {
expected_leaf: String,
stream: SyncStream,
buffered: Vec<Result<StreamEvent, ProviderError>>,
terminal: bool,
}
struct SyncStream(std::sync::Mutex<BoxStream<'static, Result<StreamEvent, ProviderError>>>);
impl SyncStream {
fn new(stream: BoxStream<'static, Result<StreamEvent, ProviderError>>) -> Self {
Self(std::sync::Mutex::new(stream))
}
fn get_mut(&mut self) -> &mut BoxStream<'static, Result<StreamEvent, ProviderError>> {
self.0
.get_mut()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn into_inner(self) -> BoxStream<'static, Result<StreamEvent, ProviderError>> {
self.0
.into_inner()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
}
impl Speculation {
async fn fill(&mut self) {
while !self.terminal {
match self.stream.get_mut().next().await {
Some(event) => {
self.terminal = matches!(event, Ok(StreamEvent::Completed { .. }) | Err(_));
self.buffered.push(event);
}
None => self.terminal = true,
}
}
}
fn adopt(self) -> BoxStream<'static, Result<StreamEvent, ProviderError>> {
Box::pin(futures_util::stream::iter(self.buffered).chain(self.stream.into_inner()))
}
}
enum SampleEnd {
Completed {
stop: StopReason,
blocks: Vec<Value>,
},
Cancelled,
Unavailable(String),
ContextFull,
Fatal(String),
}
struct Turn {
shared: Arc<SharedDeps>,
cmd_tx: mpsc::Sender<SessionCmd>,
events: mpsc::Sender<EngineEvent>,
cancel: CancellationToken,
tool_defs: Arc<[ToolDef]>,
models: Vec<String>,
model_idx: usize,
call_sigs: VecDeque<CallSig>,
consecutive_failures: HashMap<String, u32>,
usage: TokenUsage,
anchor: Option<(u64, usize)>,
samples: u32,
injected_hints: HashSet<String>,
last_snapshot: Option<crate::actor::Snapshot>,
speculation: Option<tokio::task::JoinHandle<Option<crate::SpecDigest>>>,
turn_extensions: u32,
spent: i64,
samples_since_compact: u32,
ledger: crate::ledger::LoopLedger,
pipeline: TicketPipeline,
head: tokio::sync::watch::Receiver<Arc<crate::actor::ProjectionHead>>,
speculative: Option<Speculation>,
projected_tail: Vec<Item>,
}
impl Drop for Turn {
fn drop(&mut self) {
if let Some(handle) = self.speculation.take() {
handle.abort();
}
}
}
impl Turn {
fn new(
shared: Arc<SharedDeps>,
cmd_tx: mpsc::Sender<SessionCmd>,
events: mpsc::Sender<EngineEvent>,
cancel: CancellationToken,
cont: crate::TurnContinuation,
) -> Self {
let mut models = vec![shared.config.model.clone()];
models.extend(shared.config.fallback_models.iter().cloned());
let head = shared.head();
Self {
tool_defs: shared.registry.defs().into(),
shared,
cmd_tx,
events,
cancel,
models,
model_idx: cont.model_idx,
call_sigs: cont.call_sigs,
consecutive_failures: cont.consecutive_failures,
usage: TokenUsage::default(),
anchor: None,
samples: 0,
injected_hints: HashSet::new(),
last_snapshot: None,
speculation: None,
turn_extensions: cont.turn_extensions,
spent: cont.spent,
samples_since_compact: 0,
ledger: crate::ledger::LoopLedger::new(),
pipeline: TicketPipeline::default(),
head,
speculative: None,
projected_tail: Vec::new(),
}
}
fn continuation(&mut self) -> crate::TurnContinuation {
crate::TurnContinuation {
spent: self.spent,
model_idx: self.model_idx,
call_sigs: std::mem::take(&mut self.call_sigs),
consecutive_failures: std::mem::take(&mut self.consecutive_failures),
turn_extensions: self.turn_extensions,
samples_since_compact: self.samples_since_compact,
}
}
async fn drive(&mut self) -> TurnEnd {
let bound = self.shared.config.max_turns;
while bound < 0 || self.spent < bound {
self.spent += 1;
let (stop, blocks) = match self.sample().await {
SampleEnd::Completed { stop, blocks } => (stop, blocks),
SampleEnd::Cancelled => {
self.ledger.stamp(Phase::BoundaryEnd);
return TurnEnd::Outcome(Outcome::Cancelled);
}
SampleEnd::ContextFull => {
self.ledger.stamp(Phase::BoundaryEnd);
return TurnEnd::Compact {
spec: self.take_speculation().await,
cont: Box::new(self.continuation()),
};
}
SampleEnd::Unavailable(_) if self.model_idx + 1 < self.models.len() => {
self.model_idx += 1;
self.emit(EngineEvent::FallbackModel {
model: self.models[self.model_idx].clone(),
})
.await;
self.ledger.stamp(Phase::BoundaryEnd);
continue;
}
SampleEnd::Unavailable(m) | SampleEnd::Fatal(m) => {
self.ledger.stamp(Phase::BoundaryEnd);
return TurnEnd::Outcome(Outcome::Error { message: m });
}
};
match stop {
StopReason::ToolUse => {
let outcome = self.run_tool_phase(&blocks).await;
self.ledger.stamp(Phase::BoundaryEnd);
if let Some(outcome) = outcome {
return TurnEnd::Outcome(outcome);
}
}
StopReason::Refusal => {
self.ledger.stamp(Phase::BoundaryEnd);
return TurnEnd::Outcome(Outcome::Refused);
}
_ => {
let text = assistant_text(&blocks);
let todo_fires = self.todo_gate_should_fire();
let stop_reason = self.consult_stop(&text).await;
if todo_fires || stop_reason.is_some() {
self.turn_extensions += 1;
let commit = self.inject_gate_nudge(todo_fires, stop_reason).await;
if !commit.ok() {
self.ledger.stamp(Phase::BoundaryEnd);
return TurnEnd::Outcome(commit.outcome());
}
self.ledger.stamp(Phase::BoundaryEnd);
continue;
}
self.ledger.stamp(Phase::BoundaryEnd);
return TurnEnd::Outcome(Outcome::Done { text });
}
}
}
TurnEnd::Outcome(Outcome::TurnLimit)
}
fn todo_gate_should_fire(&self) -> bool {
self.turn_extensions < TODO_GATE_MAX.min(TURN_EXTENSION_MAX)
&& self
.last_snapshot
.as_ref()
.is_some_and(|s| unfinished_todos(&s.tail))
}
async fn consult_stop(&self, outcome_text: &str) -> Option<String> {
if self.turn_extensions >= TURN_EXTENSION_MAX {
return None;
}
crate::hooks::hook_gate!(
self.shared.hooks,
self.shared.hook_mask(),
crate::hooks::EventMask::STOP,
|hooks| match crate::hooks::call_stop(hooks, outcome_text).await {
crate::hooks::StopDecision::Block { reason } => Some(reason),
crate::hooks::StopDecision::Allow => None,
},
else None
)
}
async fn inject_gate_nudge(&mut self, todo_fires: bool, stop_reason: Option<String>) -> Commit {
let mut body = String::new();
if todo_fires {
body.push_str(
"You still have open items on your todo list. Continue working on them, \
or call todo_write to mark them completed or drop them, before ending \
your turn.",
);
}
if let Some(reason) = stop_reason {
if !body.is_empty() {
body.push('\n');
}
body.push_str(&reason);
}
self.propose_pipelined(
vec![EntryPayload::Item {
item: Item::User {
text: format!("<system-reminder>{body}</system-reminder>"),
synthetic: Some(SyntheticReason::SystemReminder),
},
}],
crate::SampleStage::AtBoundary,
)
.await
}
async fn sample(&mut self) -> SampleEnd {
self.ledger.start_sample();
self.ledger.stamp(Phase::BoundaryStart);
let (commit, head) = self.boundary().await;
if !commit.ok() {
self.speculative = None;
return match commit {
Commit::Aborted => SampleEnd::Cancelled,
other => SampleEnd::Fatal(other.message().into()),
};
}
let Some(head) = head else {
return SampleEnd::Fatal("session closed".into());
};
let snapshot = head.snapshot();
self.ledger.stamp(Phase::SnapshotReady);
self.samples += 1;
let adopted = self
.speculative
.take()
.filter(|spec| head.leaf() == Some(spec.expected_leaf.as_str()))
.map(Speculation::adopt);
let stream = match adopted {
Some(adopted) => {
let estimate = self.estimate_tokens(&snapshot);
self.maybe_speculate_digest(&snapshot, estimate);
adopted
}
None => {
let request = match self.build_request(&snapshot) {
Ok(request) => request,
Err(end) => return end,
};
self.shared.provider.stream(request)
}
};
self.ledger.stamp(Phase::RequestBuilt);
self.last_snapshot = Some(snapshot.clone());
self.projected_tail.clear();
let (stop, usage, blocks) = match self.collect_stream(stream).await {
Ok(completed) => completed,
Err(end) => return end,
};
self.usage += usage;
self.samples_since_compact += 1;
let reported = usage.input_tokens
+ usage.cache_read_input_tokens
+ usage.cache_creation_input_tokens
+ usage.output_tokens;
self.anchor = Some((reported, snapshot.durable.len() + 1));
let assistant = Item::Assistant {
blocks: blocks.clone(),
};
self.projected_tail.push(assistant.clone());
self.ledger.stamp(Phase::BatchProposed);
let commit = self
.propose_pipelined(
vec![
EntryPayload::Item { item: assistant },
EntryPayload::Usage { usage },
],
crate::SampleStage::AtBoundary,
)
.await;
self.ledger.stamp(Phase::WatermarkDurable);
if !commit.ok() {
return SampleEnd::Fatal(commit.message().into());
}
SampleEnd::Completed { stop, blocks }
}
async fn run_tool_phase(&mut self, blocks: &[Value]) -> Option<Outcome> {
let uses = assistant_tool_uses(blocks);
self.run_tool_batch(&uses).await
}
fn fold_doom_window(&mut self, uses: &[ToolUse]) -> Option<String> {
fold_call_sigs(&mut self.call_sigs, uses);
detect_doom_loop(self.call_sigs.make_contiguous())
}
async fn handle_doom_loop(&mut self, uses: &[ToolUse], pattern: String) -> Option<Outcome> {
let unattended = matches!(
self.shared.effective_mode(),
hotl_tools::rules::PermissionMode::Auto | hotl_tools::rules::PermissionMode::DontAsk
);
let stop = if unattended {
true
} else {
let cont = self
.ask(
format!("the agent keeps repeating: {pattern} — let it continue?"),
None,
)
.await;
!matches!(cont, AskReply::Allow | AskReply::AllowEdited { .. })
};
if stop {
let commit = self
.abort_batch(uses, "Stopped: a repeating tool-call loop was detected.")
.await;
if !commit.ok() {
return Some(commit.outcome());
}
return Some(Outcome::DoomLoop { pattern });
}
self.call_sigs.clear();
None
}
async fn run_tool_batch(&mut self, uses: &[ToolUse]) -> Option<Outcome> {
let commit = self.pipeline.drain().await;
if !commit.ok() {
return Some(commit.outcome());
}
if let Some(pattern) = self.fold_doom_window(uses) {
if let Some(outcome) = self.handle_doom_loop(uses, pattern).await {
return Some(outcome);
}
}
let mutating = uses.iter().any(|tu| {
!self
.shared
.registry
.get(&tu.name)
.is_some_and(|t| t.read_only())
});
if mutating {
self.snap(format!("pre batch {}", self.samples)).await;
}
let mut results = Vec::with_capacity(uses.len());
let mut budget_blown: Option<String> = None;
self.ledger.stamp(Phase::ToolsSpawned);
debug_assert!(
self.pipeline.is_empty(),
"barrier (a): no tool may run while a commit is unresolved"
);
if let [only] = uses {
if self.cancel.is_cancelled() {
results.push(pair(only, "Not executed (turn stopped).", true));
} else {
let gate = self.gate(only).await;
let executed = self.execute(only, gate).await;
results.push(self.finish_call(only, executed, &mut budget_blown).await);
}
} else {
for chunk in parallel_chunks(uses, &self.shared.registry) {
if self.cancel.is_cancelled() || budget_blown.is_some() {
for tu in chunk {
results.push(pair(tu, "Not executed (turn stopped).", true));
}
continue;
}
let mut gates = Vec::with_capacity(chunk.len());
for tu in chunk {
gates.push(self.gate(tu).await);
}
let outcomes = futures_util::future::join_all(
chunk
.iter()
.zip(gates)
.map(|(tu, gate)| self.execute(tu, gate)),
)
.await;
for (tu, executed) in chunk.iter().zip(outcomes) {
results.push(self.finish_call(tu, executed, &mut budget_blown).await);
}
}
}
self.ledger.stamp(Phase::ToolsJoined);
if mutating {
self.snap(format!("post batch {}", self.samples)).await;
}
let cancelled = self.cancel.is_cancelled();
let mut entries = vec![EntryPayload::Item {
item: Item::ToolResults { results },
}];
entries.extend(
self.subdir_hints(uses)
.into_iter()
.map(|item| EntryPayload::Item { item }),
);
self.projected_tail
.extend(entries.iter().filter_map(|e| match e {
EntryPayload::Item { item } => Some(item.clone()),
_ => None,
}));
self.ledger.restamp(Phase::BatchProposed);
let commit = self
.propose_pipelined(entries, crate::SampleStage::AtBoundary)
.await;
self.ledger.restamp(Phase::WatermarkDurable);
if !commit.ok() {
return Some(commit.outcome());
}
if cancelled {
return Some(Outcome::Cancelled);
}
let outcome = budget_blown.map(|tool| Outcome::ToolFailureBudget { tool });
if outcome.is_none() {
self.speculate();
}
outcome
}
async fn finish_call(
&mut self,
tu: &ToolUse,
mut executed: Executed,
budget_blown: &mut Option<String>,
) -> ToolResultItem {
self.maybe_evict(tu, &mut executed.outcome).await;
let (content, failed) = self.apply_failure_budget(tu, executed, budget_blown);
ToolResultItem {
tool_use_id: tu.id.clone(),
content,
is_error: failed,
}
}
fn apply_failure_budget(
&mut self,
tu: &ToolUse,
executed: Executed,
budget_blown: &mut Option<String>,
) -> (String, bool) {
let Executed {
outcome,
chargeable,
} = executed;
let mut content = outcome.content;
if outcome.is_error && !chargeable {
return (content, true);
}
if outcome.is_error {
let n = self
.consecutive_failures
.entry(tu.name.clone())
.or_insert(0);
*n += 1;
let left = self.shared.config.tool_failure_budget.saturating_sub(*n);
content.push_str(&format!("\n<retry attempts_left={left}>"));
if left == 0 && budget_blown.is_none() {
*budget_blown = Some(tu.name.clone());
}
} else {
self.consecutive_failures.remove(&tu.name);
}
(content, outcome.is_error)
}
async fn gate(&self, tu: &ToolUse) -> Gate {
let Some(tool) = self.shared.registry.get(&tu.name) else {
return Gate::Resolved {
outcome: unknown_tool(&self.tool_defs, &tu.name),
chargeable: true,
};
};
let mut input = tu.input.clone();
crate::hooks::hook_gate!(
self.shared.hooks,
self.shared.hook_mask(),
crate::hooks::EventMask::PRE_TOOL,
|hooks| {
let view = crate::hooks::cap_tool_input(&input);
match crate::hooks::call_pre_tool(hooks, &tu.name, &view, &self.cancel).await {
crate::hooks::PreToolDecision::Continue => {}
crate::hooks::PreToolDecision::Deny { message } => {
self.emit(EngineEvent::ToolDenied {
name: tu.name.clone(),
})
.await;
return Gate::Resolved {
outcome: ToolOutcome::err(format!(
"A hook blocked this tool call: {message}"
)),
chargeable: false,
};
}
crate::hooks::PreToolDecision::Rewrite { input: rewritten } => {
input = crate::hooks::restore_capped(&input, rewritten)
}
}
if self.cancel.is_cancelled() {
return Gate::Resolved {
outcome: ToolOutcome::err("Not executed (turn stopped)."),
chargeable: false,
};
}
},
else {}
);
let (summary, why) = match tool.permission(&input) {
Permission::None => (None, None),
Permission::Ask { summary } => (Some(summary), None),
Permission::AskProtected { summary, why } => (Some(summary), Some(why)),
};
let display = summary.clone().unwrap_or_else(|| tu.name.clone());
if let Some(summary) = summary {
match self.approve_input(tu, &input, summary, why).await {
AskReply::Allow => {}
AskReply::AllowEdited { input: edited } => input = edited, AskReply::Respond { content } => {
self.emit(EngineEvent::ToolDone {
name: tu.name.clone(),
ok: true,
})
.await;
return Gate::Resolved {
outcome: ToolOutcome::ok(content),
chargeable: true,
};
}
AskReply::Deny { message } => {
self.emit(EngineEvent::ToolDenied {
name: tu.name.clone(),
})
.await;
return Gate::Resolved {
outcome: match message {
Some(m) => ToolOutcome::err(format!("The user declined this tool call: {m}")),
None => ToolOutcome::err(
"The user declined this tool call. Ask what they'd like to do instead, or proceed another way.",
),
},
chargeable: false,
};
}
}
}
Gate::Ready {
input,
summary: display,
}
}
async fn execute(&self, tu: &ToolUse, gate: Gate) -> Executed {
let (input, summary) = match gate {
Gate::Ready { input, summary } => (input, summary),
Gate::Resolved {
outcome,
chargeable,
} => {
return Executed {
outcome,
chargeable,
}
}
};
self.emit(EngineEvent::ToolStart {
name: tu.name.clone(),
summary,
})
.await;
let Some(tool) = self.shared.registry.get(&tu.name) else {
return Executed {
outcome: unknown_tool(&self.tool_defs, &tu.name),
chargeable: true,
};
};
let mut outcome = tool.run(input, self.cancel.clone()).await;
if !outcome.is_error {
crate::hooks::hook_gate!(
self.shared.hooks,
self.shared.hook_mask(),
crate::hooks::EventMask::POST_TOOL,
|hooks| {
if let Some(replacement) = crate::hooks::call_post_tool(
hooks,
&tu.name,
&outcome.content,
&self.cancel,
)
.await
{
outcome.content = replacement;
}
},
else {}
);
}
self.emit(EngineEvent::ToolDone {
name: tu.name.clone(),
ok: !outcome.is_error,
})
.await;
Executed {
outcome,
chargeable: true,
}
}
async fn approve_input(
&self,
tu: &ToolUse,
input: &Value,
summary: String,
why: Option<String>,
) -> AskReply {
let protected = why.is_some();
let read_only = self
.shared
.registry
.get(&tu.name)
.is_some_and(|t| t.read_only());
match self.shared.rules.evaluate(
self.shared.effective_mode(),
&tu.name,
input,
self.shared.sandbox_enforced,
protected,
read_only,
) {
Verdict::Auto { rule } => {
self.emit(EngineEvent::ToolAutoAllowed {
name: tu.name.clone(),
rule,
})
.await;
AskReply::Allow
}
Verdict::Deny { rule } => AskReply::Deny {
message: Some(format!(
"a deny rule refused this call ({rule}); do not retry it"
)),
},
Verdict::Ask => self.ask(summary, why).await,
}
}
async fn abort_batch(&mut self, uses: &[ToolUse], message: &str) -> Commit {
let results = uses.iter().map(|tu| pair(tu, message, true)).collect();
self.ledger.restamp(Phase::BatchProposed);
let commit = self
.propose_pipelined(
vec![EntryPayload::Item {
item: Item::ToolResults { results },
}],
crate::SampleStage::AtBoundary,
)
.await;
self.ledger.restamp(Phase::WatermarkDurable);
commit
}
fn build_request(
&mut self,
snapshot: &crate::actor::Snapshot,
) -> Result<SamplingRequest, SampleEnd> {
let window = self.shared.config.context_window.max(1);
let estimate = self.estimate_tokens(snapshot);
if estimate > (window as f64 * COMPACT_TRIGGER) as u64 {
return Err(SampleEnd::ContextFull);
}
self.maybe_speculate_digest(snapshot, estimate);
Ok(self.compose_request(snapshot, estimate, self.samples))
}
fn maybe_speculate_digest(&mut self, snapshot: &crate::actor::Snapshot, estimate: u64) {
let window = self.shared.config.context_window.max(1);
if self.speculation.is_none()
&& !self.shared.config.compaction_reset
&& estimate > (window as f64 * SPECULATE_TRIGGER) as u64
{
self.speculation = spawn_speculation(&self.shared, &snapshot.durable, &self.cancel);
}
}
fn speculative_request(&self, snapshot: &crate::actor::Snapshot) -> Option<SamplingRequest> {
let window = self.shared.config.context_window.max(1);
let estimate = self.estimate_tokens(snapshot);
if estimate > (window as f64 * COMPACT_TRIGGER) as u64 {
return None;
}
Some(self.compose_request(snapshot, estimate, self.samples + 1))
}
fn compose_request(
&self,
snapshot: &crate::actor::Snapshot,
estimate: u64,
sample_no: u32,
) -> SamplingRequest {
let window = self.shared.config.context_window.max(1);
let used_pct = self
.shared
.config
.show_context_pct
.then(|| (estimate.saturating_mul(100) / window).min(100) as u8);
let turn_context = hotl_context::turn_context(
self.shared.clock.now_ms(),
&self.shared.cwd,
used_pct,
sample_no,
);
SamplingRequest {
model: self.models[self.model_idx].clone(),
max_tokens: self.shared.config.max_tokens,
system: Arc::clone(&self.shared.system),
items: Arc::clone(&snapshot.durable),
ephemeral_tail: Arc::clone(&snapshot.tail),
tools: Arc::clone(&self.tool_defs),
thinking: self.shared.config.thinking,
cache: if self.shared.config.cache_static {
CachePolicy::Static {
prefix_ttl: self.shared.config.cache_ttl,
}
} else {
CachePolicy::Off
},
turn_context: Some(turn_context),
}
}
async fn collect_stream(
&mut self,
stream: BoxStream<'static, Result<StreamEvent, ProviderError>>,
) -> Result<(StopReason, TokenUsage, Vec<Value>), SampleEnd> {
let mut stream = stream;
let mut completed = None;
loop {
tokio::select! {
biased;
_ = self.cancel.cancelled() => return Err(SampleEnd::Cancelled),
next = stream.next() => match next {
Some(Ok(event)) => {
self.ledger.stamp(Phase::FirstByte);
if let StreamEvent::Completed { stop, usage, blocks } = event {
self.ledger.stamp(Phase::LastBlockEnd);
completed = Some((stop, usage, blocks));
} else {
self.forward(event).await;
}
}
Some(Err(e)) if retry::is_context_overflow(&e) => {
drain_to_end(&mut stream).await;
return Err(SampleEnd::ContextFull);
}
Some(Err(e)) if retry::is_availability(&e) => {
drain_to_end(&mut stream).await;
return Err(SampleEnd::Unavailable(e.to_string()));
}
Some(Err(e)) => {
drain_to_end(&mut stream).await;
return Err(SampleEnd::Fatal(e.to_string()));
}
None => break,
}
}
}
completed.ok_or_else(|| SampleEnd::Fatal("stream ended without completion".into()))
}
fn estimate_tokens(&self, snapshot: &crate::actor::Snapshot) -> u64 {
anchored_estimate(self.anchor, &self.shared.system, &snapshot.durable)
+ hotl_context::tokens::estimate_items(&snapshot.tail)
}
fn subdir_hints(&mut self, uses: &[ToolUse]) -> Vec<Item> {
let mut out = Vec::new();
for tu in uses {
let Some(path) = tu.input.get("path").and_then(Value::as_str) else {
continue;
};
let Some((marker, item)) =
hotl_context::nested_instructions(&self.shared.cwd, Path::new(path))
else {
continue;
};
if self.injected_hints.contains(&marker) || self.in_projection(&marker) {
continue;
}
self.injected_hints.insert(marker);
out.push(item);
}
out
}
fn in_projection(&self, marker: &str) -> bool {
let Some(snapshot) = &self.last_snapshot else {
return false;
};
snapshot.durable.iter().any(|i| {
matches!(
i,
Item::User { text, synthetic: Some(hotl_types::SyntheticReason::SubdirInstructions) }
if text.contains(marker)
)
})
}
async fn snap(&self, label: String) {
if let Some(snapshots) = &self.shared.snapshots {
snapshots.snapshot(label).await;
}
}
async fn maybe_evict(&self, tu: &ToolUse, outcome: &mut ToolOutcome) {
let threshold = self.shared.config.evict_threshold_tokens;
if threshold == 0 || outcome.is_error {
return;
}
if hotl_context::tokens::estimate_text(&outcome.content) <= threshold {
return;
}
let content = std::mem::take(&mut outcome.content);
let total = content.len();
let head = clip(&content, 2048).to_string();
let (tx, rx) = oneshot::channel();
let cmd = SessionCmd::WriteBlob {
tool_use_id: tu.id.clone(),
content,
reply: tx,
};
if let Err(mpsc::error::SendError(cmd)) = self.cmd_tx.send(cmd).await {
if let SessionCmd::WriteBlob { content, .. } = cmd {
outcome.content = content;
}
return;
}
match rx.await {
Ok(Ok(path)) => {
outcome.content = format!(
"{head}\n<evicted total_bytes={total} file=\"{path}\">Full output saved. \
Read it with the read tool ({path}); use offset to page.</evicted>"
);
}
Ok(Err(content)) => outcome.content = content,
Err(_) => outcome.content = head,
}
}
async fn take_speculation(&mut self) -> Option<crate::SpecDigest> {
let handle = self.speculation.take()?;
await_speculation(handle, &self.cancel, SPECULATION_WAIT).await
}
async fn boundary(&mut self) -> (Commit, Option<Arc<crate::actor::ProjectionHead>>) {
let Turn {
pipeline,
speculative,
head,
..
} = self;
let refresh = async {
let commit = pipeline.drain().await;
if !commit.ok() {
return (commit, None);
}
let head = match pipeline.ack_seq() {
Some(my_ack_seq) => head
.wait_for(|head| head.epoch() >= my_ack_seq)
.await
.ok()
.map(|head| Arc::clone(&head)),
None => Some(Arc::clone(&head.borrow())),
};
(commit, head)
};
let mut refresh = std::pin::pin!(refresh);
match speculative.as_mut() {
Some(spec) => tokio::select! {
refreshed = &mut refresh => refreshed,
_ = spec.fill() => refresh.await,
},
None => refresh.await,
}
}
fn speculate(&mut self) {
let bound = self.shared.config.max_turns;
let will_sample_again = bound < 0 || self.spent < bound;
if self.cancel.is_cancelled() || self.speculative.is_some() || !will_sample_again {
return;
}
let Some(expected_leaf) = self.pipeline.leaf().map(str::to_string) else {
return;
};
let Some(predicted) = self.predicted_snapshot() else {
return;
};
let Some(request) = self.speculative_request(&predicted) else {
return;
};
let stream = self.shared.provider.stream(request);
self.speculative = Some(Speculation {
expected_leaf,
stream: SyncStream::new(stream),
buffered: Vec::new(),
terminal: false,
});
}
fn predicted_snapshot(&self) -> Option<crate::actor::Snapshot> {
let last = self.last_snapshot.as_ref()?;
let mut durable = Vec::with_capacity(last.durable.len() + self.projected_tail.len());
durable.extend(last.durable.iter().cloned());
durable.extend(self.projected_tail.iter().cloned());
Some(crate::actor::Snapshot {
durable: Arc::new(durable),
tail: Arc::clone(&last.tail),
})
}
async fn propose(&self, entries: Vec<EntryPayload>, stage: crate::SampleStage) -> Commit {
let epoch = self.shared.rules_epoch();
let commit = self.propose_at_epoch(&entries, epoch, stage).await;
if commit != Commit::StaleEpoch {
return commit;
}
let epoch = self.shared.rules_epoch();
self.propose_at_epoch(&entries, epoch, stage).await
}
async fn propose_at_epoch(
&self,
entries: &[EntryPayload],
epoch: u32,
stage: crate::SampleStage,
) -> Commit {
match self.send(entries, epoch, crate::AckMode::Sync, stage).await {
Ok(reply) => Commit::from_reply(Some(reply)),
Err(commit) => commit,
}
}
async fn propose_pipelined(
&mut self,
entries: Vec<EntryPayload>,
stage: crate::SampleStage,
) -> Commit {
if self.shared.config.ack_mode == crate::AckMode::Sync {
return self.propose(entries, stage).await;
}
let epoch = self.shared.rules_epoch();
let commit = self.submit_pipelined(&entries, epoch, stage).await;
if commit != Commit::StaleEpoch {
return commit;
}
let epoch = self.shared.rules_epoch();
self.submit_pipelined(&entries, epoch, stage).await
}
async fn submit_pipelined(
&mut self,
entries: &[EntryPayload],
epoch: u32,
stage: crate::SampleStage,
) -> Commit {
match self
.send(entries, epoch, crate::AckMode::Pipelined, stage)
.await
{
Ok(crate::ProposeReply::Ticket(ticket)) => self.pipeline.submit(ticket).await,
Ok(other) => Commit::from_reply(Some(other)),
Err(commit) => commit,
}
}
async fn send(
&self,
entries: &[EntryPayload],
epoch: u32,
mode: crate::AckMode,
stage: crate::SampleStage,
) -> Result<crate::ProposeReply, Commit> {
let masker = self.shared.masker();
let mut prepared = Vec::with_capacity(entries.len());
for payload in entries {
match prepare_entry(payload, masker, epoch).await {
Ok(entry) => prepared.push(entry),
Err(PrepareError::Serialize) => {
debug_assert!(
false,
"EntryPayload serialization must not fail for hotl's own payload shapes"
);
return Err(Commit::Gone);
}
Err(PrepareError::Offload) => return Err(Commit::Gone),
}
}
let (tx, rx) = oneshot::channel();
if self
.cmd_tx
.send(SessionCmd::ProposePrepared {
proposal: crate::EntryProposal::of(prepared),
stage,
mode,
reply: tx,
})
.await
.is_err()
{
return Err(Commit::Gone);
}
rx.await.map_err(|_| Commit::Gone)
}
async fn ask(&self, summary: String, why: Option<String>) -> AskReply {
let id = hotl_types::new_ulid();
let _ = self
.propose(
vec![EntryPayload::PendingAsk {
id: id.clone(),
summary: summary.clone(),
protected_why: why.clone(),
}],
crate::SampleStage::AtBoundary,
)
.await;
crate::hooks::hook_gate!(
self.shared.hooks,
self.shared.hook_mask(),
crate::hooks::EventMask::NOTIFICATION,
|hooks| {
crate::hooks::notify(
hooks,
&self.shared.notifications,
crate::hooks::NotificationKind::Blocked,
summary.clone(),
);
},
else {}
);
let (tx, rx) = oneshot::channel();
let event = EngineEvent::Ask {
summary,
protected_why: why,
reply: tx,
};
let reply = if self.events.send(event).await.is_err() {
AskReply::Deny { message: None }
} else {
tokio::select! {
biased;
_ = self.cancel.cancelled() => AskReply::Deny {
message: Some("the user interrupted the turn".into()),
},
reply = rx => reply.unwrap_or(AskReply::Deny { message: None }),
}
};
let allowed = matches!(
reply,
AskReply::Allow | AskReply::AllowEdited { .. } | AskReply::Respond { .. }
);
let _ = self
.propose(
vec![EntryPayload::AskResolved { id, allowed }],
crate::SampleStage::AtBoundary,
)
.await;
reply
}
async fn emit(&self, event: EngineEvent) {
let _ = self.events.send(event).await;
}
async fn forward(&self, event: StreamEvent) {
let mapped = match event {
StreamEvent::TextDelta { text, .. } => EngineEvent::TextDelta(text),
StreamEvent::ThinkingDelta { text, .. } => EngineEvent::ThinkingDelta(text),
StreamEvent::Retrying { attempt, reason } => EngineEvent::Retrying { attempt, reason },
_ => return,
};
self.emit(mapped).await;
}
}
pub(crate) async fn drain_to_end(
stream: &mut BoxStream<'static, Result<StreamEvent, ProviderError>>,
) {
for _ in 0..DRAIN_MAX {
if stream.next().await.is_none() {
return;
}
}
}
const DRAIN_MAX: usize = 64;
fn anchored_estimate(anchor: Option<(u64, usize)>, system: &str, durable: &[Item]) -> u64 {
use hotl_context::tokens;
match anchor {
Some((reported, len)) if durable.len() >= len => {
reported + tokens::estimate_items(&durable[len..])
}
_ => tokens::estimate_text(system) + tokens::estimate_items(durable),
}
}
fn clip(s: &str, max: usize) -> &str {
if s.len() <= max {
return s;
}
let mut end = max;
while !s.is_char_boundary(end) {
end -= 1;
}
&s[..end]
}
fn pair(tu: &ToolUse, message: &str, is_error: bool) -> ToolResultItem {
ToolResultItem {
tool_use_id: tu.id.clone(),
content: message.to_string(),
is_error,
}
}
fn unfinished_todos(tail: &[Item]) -> bool {
matches!(
tail.last(),
Some(Item::User {
text,
synthetic: Some(SyntheticReason::Todos),
}) if text.contains("[ ]") || text.contains("[~]")
)
}
fn unknown_tool(defs: &[ToolDef], name: &str) -> ToolOutcome {
let available: Vec<_> = defs.iter().map(|d| d.name.as_str()).collect();
ToolOutcome::err(format!(
"Unknown tool `{name}`. Available tools: {}.",
available.join(", ")
))
}
const DOOM_WINDOW: usize = 9;
#[derive(Debug)]
pub(crate) struct CallSig {
hash: u64,
display: String,
}
impl CallSig {
fn new(tu: &ToolUse) -> Self {
let display = format!("{}({})", tu.name, tu.input);
let mut hasher = std::collections::hash_map::DefaultHasher::new();
display.hash(&mut hasher);
Self {
hash: hasher.finish(),
display,
}
}
}
fn fold_call_sigs(window: &mut VecDeque<CallSig>, uses: &[ToolUse]) {
window.extend(uses.iter().map(CallSig::new));
while window.len() > DOOM_WINDOW {
window.pop_front();
}
}
fn detect_doom_loop(sigs: &[CallSig]) -> Option<String> {
const REPEATS: usize = 3;
for period in 1..=3usize {
let need = period * REPEATS;
if sigs.len() < need {
continue;
}
let tail = &sigs[sigs.len() - need..];
let block = &tail[..period];
let same = |a: &CallSig, b: &CallSig| a.hash == b.hash && a.display == b.display;
if tail
.chunks(period)
.all(|c| c.iter().zip(block).all(|(a, b)| same(a, b)))
{
return Some(
block
.iter()
.map(|s| s.display.as_str())
.collect::<Vec<_>>()
.join(" → "),
);
}
}
None
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn sig(name: &str, input: Value) -> CallSig {
CallSig::new(&ToolUse {
id: "t".into(),
name: name.into(),
input,
})
}
#[test]
fn doom_detector_finds_periods() {
let a = || sig("read", json!({"path":"x"}));
let b = || sig("bash", json!({"command":"ls"}));
assert!(detect_doom_loop(&[a(), a(), a()]).is_some());
let sigs = vec![a(), b(), a(), b(), a(), b()];
assert!(detect_doom_loop(&sigs).is_some());
assert!(detect_doom_loop(&[a(), a(), b()]).is_none());
assert!(detect_doom_loop(&[a(), a()]).is_none());
let pattern = detect_doom_loop(&[a(), a(), a()]).unwrap();
assert_eq!(pattern, "read({\"path\":\"x\"})");
}
fn uses(pairs: &[(&str, Value)]) -> Vec<ToolUse> {
pairs
.iter()
.enumerate()
.map(|(i, (name, input))| ToolUse {
id: i.to_string(),
name: (*name).into(),
input: input.clone(),
})
.collect()
}
#[test]
fn fold_call_sigs_preserves_source_order_and_evicts_from_the_front() {
let mut window: VecDeque<CallSig> = VecDeque::new();
let batch1 = uses(&[
("read", json!({"path": "a"})),
("bash", json!({"command": "ls"})),
]);
fold_call_sigs(&mut window, &batch1);
let displays: Vec<&str> = window.iter().map(|s| s.display.as_str()).collect();
assert_eq!(
displays,
vec![r#"read({"path":"a"})"#, r#"bash({"command":"ls"})"#]
);
let batch2 = uses(
&(0..DOOM_WINDOW)
.map(|i| ("t", json!(i)))
.collect::<Vec<_>>(),
);
fold_call_sigs(&mut window, &batch2);
assert_eq!(window.len(), DOOM_WINDOW);
let displays: Vec<String> = window.iter().map(|s| s.display.clone()).collect();
let expected: Vec<String> = (0..DOOM_WINDOW).map(|i| format!("t({i})")).collect();
assert_eq!(displays, expected);
}
struct DropFlag(Arc<std::sync::atomic::AtomicBool>);
impl Drop for DropFlag {
fn drop(&mut self) {
self.0.store(true, std::sync::atomic::Ordering::SeqCst);
}
}
#[tokio::test(start_paused = true)]
async fn the_speculation_bound_abandons_a_stalled_digest() {
let dropped = Arc::new(std::sync::atomic::AtomicBool::new(false));
let flag = DropFlag(Arc::clone(&dropped));
let handle = tokio::spawn(async move {
let _flag = flag;
std::future::pending::<()>().await;
None
});
let got = await_speculation(handle, &CancellationToken::new(), SPECULATION_WAIT).await;
assert!(got.is_none(), "a stalled digest must be abandoned");
tokio::task::yield_now().await;
assert!(
dropped.load(std::sync::atomic::Ordering::SeqCst),
"the abandoned speculation must be aborted, not left burning a provider call"
);
}
#[tokio::test(start_paused = true)]
async fn a_cancelled_turn_abandons_its_speculation_without_waiting_out_the_bound() {
let handle = tokio::spawn(async {
std::future::pending::<()>().await;
None
});
let cancel = CancellationToken::new();
cancel.cancel();
let start = tokio::time::Instant::now();
let got = await_speculation(handle, &cancel, SPECULATION_WAIT).await;
assert!(got.is_none());
assert_eq!(
start.elapsed(),
std::time::Duration::ZERO,
"an interrupt must not wait out the speculation bound"
);
}
#[tokio::test(start_paused = true)]
async fn a_finished_speculation_still_reaches_the_fold() {
let handle = tokio::spawn(async {
Some(crate::SpecDigest {
prefix_end: 1,
kept_from: 4,
text: "DIGEST".into(),
})
});
tokio::task::yield_now().await;
let got = await_speculation(handle, &CancellationToken::new(), SPECULATION_WAIT)
.await
.expect("a finished digest must reach the fold");
assert_eq!(got.text, "DIGEST");
assert_eq!((got.prefix_end, got.kept_from), (1, 4));
}
#[test]
fn a_hash_collision_alone_is_not_a_doom_loop() {
let collide = |display: &str| CallSig {
hash: 7,
display: display.into(),
};
assert!(detect_doom_loop(&[collide("a"), collide("b"), collide("c")]).is_none());
assert!(detect_doom_loop(&[collide("a"), collide("a"), collide("a")]).is_some());
}
fn todos_item(text: &str) -> Item {
Item::User {
text: text.into(),
synthetic: Some(SyntheticReason::Todos),
}
}
#[test]
fn a_dropped_actor_is_not_reported_as_a_sealed_log() {
assert_eq!(
Commit::from_reply(Some(crate::ProposeReply::Committed)),
Commit::Committed
);
assert_eq!(
Commit::from_reply(Some(crate::ProposeReply::Sealed)),
Commit::Sealed
);
assert_eq!(Commit::from_reply(None), Commit::Gone);
assert!(Commit::Sealed.message().contains("sealed"));
assert!(!Commit::Gone.message().contains("sealed"));
assert!(Commit::Gone.message().contains("shutting down"));
assert!(Commit::Committed.ok() && !Commit::Sealed.ok() && !Commit::Gone.ok());
}
fn tool_results(content: &str) -> Item {
Item::ToolResults {
results: vec![ToolResultItem {
tool_use_id: "t1".into(),
content: content.into(),
is_error: false,
}],
}
}
#[test]
fn the_anchor_counts_tool_results_even_with_todos_active() {
let prompt = Item::User {
text: "go".into(),
synthetic: None,
};
let sample1 = [prompt.clone()];
let anchor = Some((500, sample1.len() + 1));
let big = "x".repeat(3_000);
let sample2 = vec![
prompt,
Item::Assistant {
blocks: vec![json!({"type": "text", "text": "ok"})],
},
tool_results(&big),
];
let estimate = anchored_estimate(anchor, "sys", &sample2);
assert!(
estimate >= 500 + 1_000,
"tool results must be inside the estimate, got {estimate}"
);
}
#[test]
fn splitting_the_tail_out_leaves_the_estimate_unchanged() {
use hotl_context::tokens;
let durable = vec![
Item::User {
text: "go".into(),
synthetic: None,
},
tool_results(&"x".repeat(600)),
];
let reminder = todos_item("<todos>\n[~] wire the gate\n</todos>");
let flat: Vec<Item> = durable.iter().cloned().chain([reminder.clone()]).collect();
let tail = vec![reminder];
for anchor in [None, Some((500, 1)), Some((500, 2))] {
assert_eq!(
anchored_estimate(anchor, "sys", &durable) + tokens::estimate_items(&tail),
anchored_estimate(anchor, "sys", &flat),
"anchor {anchor:?}"
);
}
}
#[test]
fn unfinished_todos_reads_only_the_tagged_last_item() {
assert!(!unfinished_todos(&[]));
assert!(!unfinished_todos(&[Item::User {
text: "[ ] not a todo reminder".into(),
synthetic: None,
}]));
assert!(unfinished_todos(&[todos_item("<todos>\n[ ] a\n</todos>")]));
assert!(unfinished_todos(&[todos_item("<todos>\n[~] a\n</todos>")]));
assert!(!unfinished_todos(&[todos_item("<todos>\n[x] a\n</todos>")]));
}
#[tokio::test]
async fn prepare_entry_masks_and_stamps_the_epoch() {
std::env::set_var("HOTL_T8_TURN_TOKEN", "sk-super-secret-turntest-1");
let masker = Arc::new(hotl_store::Masker::from_env());
let payload = EntryPayload::Item {
item: Item::User {
text: "key: sk-super-secret-turntest-1".into(),
synthetic: None,
},
};
let prepared = prepare_entry(&payload, &masker, 7).await.expect("prepare");
assert_eq!(prepared.payload.rules_epoch(), 7);
assert_eq!(prepared.payload.kind(), hotl_store::EntryKind::Item);
let text = std::str::from_utf8(prepared.payload.bytes()).unwrap();
assert!(
!text.contains("sk-super-secret-turntest-1"),
"secret must be masked before it ever reaches the actor: {text}"
);
assert!(matches!(prepared.item, Some(Item::User { .. })));
std::env::remove_var("HOTL_T8_TURN_TOKEN");
}
#[tokio::test]
async fn prepare_entry_has_no_item_for_non_item_kinds() {
let masker = Arc::new(hotl_store::Masker::empty());
let payload = EntryPayload::Usage {
usage: TokenUsage::default(),
};
let prepared = prepare_entry(&payload, &masker, 0).await.expect("prepare");
assert!(prepared.item.is_none());
assert_eq!(prepared.payload.kind(), hotl_store::EntryKind::Usage);
}
#[tokio::test]
async fn prepare_entry_clean_input_matches_raw_serialization_byte_for_byte() {
let masker = Arc::new(
hotl_store::Masker::empty().with_value("HOTL_T8_UNUSED", "not-present-anywhere-12345"),
);
let payload = EntryPayload::Usage {
usage: TokenUsage::default(),
};
let prepared = prepare_entry(&payload, &masker, 0).await.expect("prepare");
let raw = hotl_store::serialize_payload(&payload).unwrap();
assert_eq!(prepared.payload.bytes().as_ref(), raw.as_bytes());
}
#[tokio::test]
async fn prepare_entry_dirty_input_matches_masker_apply() {
let masker = Arc::new(
hotl_store::Masker::empty().with_value("HOTL_T8_DIRTY", "present-secret-value-67890"),
);
let payload = EntryPayload::Item {
item: Item::User {
text: "token present-secret-value-67890 here".into(),
synthetic: None,
},
};
let prepared = prepare_entry(&payload, &masker, 0).await.expect("prepare");
let raw = hotl_store::serialize_payload(&payload).unwrap();
let expected = masker.apply(&raw);
assert_eq!(prepared.payload.bytes().as_ref(), expected.as_bytes());
}
#[tokio::test]
async fn prepare_entry_offloads_large_payloads_with_identical_bytes() {
let masker = Arc::new(hotl_store::Masker::empty());
let big_text = "x".repeat(PREPARE_OFFLOAD_BYTES + 4096);
let payload = EntryPayload::Item {
item: Item::User {
text: big_text,
synthetic: None,
},
};
let raw = hotl_store::serialize_payload(&payload).unwrap();
assert!(
raw.len() >= PREPARE_OFFLOAD_BYTES,
"fixture must actually cross the offload threshold"
);
let prepared = prepare_entry(&payload, &masker, 0).await.expect("prepare");
assert_eq!(prepared.payload.bytes().as_ref(), raw.as_bytes());
}
#[test]
fn payload_size_estimate_sums_tool_result_content_across_the_batch() {
let each = PREPARE_OFFLOAD_BYTES / 5;
let payload = EntryPayload::Item {
item: Item::ToolResults {
results: (0..10)
.map(|i| ToolResultItem {
tool_use_id: format!("t{i}"),
content: "x".repeat(each),
is_error: false,
})
.collect(),
},
};
assert!(
each < PREPARE_OFFLOAD_BYTES,
"fixture must keep each individual result under the threshold"
);
assert!(
payload_size_estimate(&payload) >= PREPARE_OFFLOAD_BYTES,
"ten results at PREPARE_OFFLOAD_BYTES/5 each must sum past the threshold"
);
}
#[test]
fn payload_size_estimate_is_zero_for_kinds_that_never_offload() {
assert_eq!(
payload_size_estimate(&EntryPayload::Usage {
usage: TokenUsage::default()
}),
0
);
}
#[tokio::test]
async fn prepare_entry_offload_decision_uses_the_pre_serialization_estimate() {
let masker = Arc::new(hotl_store::Masker::empty());
let payload = EntryPayload::Item {
item: Item::ToolResults {
results: vec![ToolResultItem {
tool_use_id: "t1".into(),
content: "y".repeat(PREPARE_OFFLOAD_BYTES + 4096),
is_error: false,
}],
},
};
assert!(payload_size_estimate(&payload) >= PREPARE_OFFLOAD_BYTES);
let raw = hotl_store::serialize_payload(&payload).unwrap();
let prepared = prepare_entry(&payload, &masker, 3).await.expect("prepare");
assert_eq!(prepared.payload.bytes().as_ref(), raw.as_bytes());
assert_eq!(prepared.payload.rules_epoch(), 3);
}
#[test]
fn commit_from_reply_distinguishes_stale_epoch_from_sealed() {
assert_eq!(
Commit::from_reply(Some(crate::ProposeReply::StaleEpoch)),
Commit::StaleEpoch
);
assert_eq!(
Commit::from_reply(Some(crate::ProposeReply::Sealed)),
Commit::Sealed
);
assert_eq!(
Commit::from_reply(Some(crate::ProposeReply::Committed)),
Commit::Committed
);
assert_eq!(Commit::from_reply(None), Commit::Gone);
}
fn ticket(
seq: u64,
) -> (
crate::CommitTicket,
oneshot::Sender<Result<crate::CommitAck, crate::CommitFailed>>,
) {
let (tx, rx) = oneshot::channel();
(
crate::CommitTicket {
id: format!("id-{seq}"),
seq,
ack: rx,
},
tx,
)
}
fn resolved(
seq: u64,
result: Result<crate::CommitAck, crate::CommitFailed>,
) -> crate::CommitTicket {
let (ticket, tx) = ticket(seq);
let _ = tx.send(result);
ticket
}
#[tokio::test]
async fn the_ticket_window_waits_on_the_oldest_before_it_overflows() {
let mut pipeline = TicketPipeline::default();
for seq in 0..ACK_WINDOW as u64 {
let commit = pipeline
.submit(resolved(seq, Ok(crate::CommitAck { offset: seq })))
.await;
assert_eq!(commit, Commit::Committed);
}
assert_eq!(pipeline.len(), ACK_WINDOW);
let commit = pipeline
.submit(resolved(99, Ok(crate::CommitAck { offset: 99 })))
.await;
assert_eq!(commit, Commit::Committed);
assert_eq!(
pipeline.len(),
ACK_WINDOW,
"the window is the bound, not a suggestion"
);
}
#[tokio::test]
async fn a_barrier_reports_the_first_failing_ticket_and_empties_the_window() {
let mut pipeline = TicketPipeline::default();
pipeline
.submit(resolved(1, Ok(crate::CommitAck { offset: 1 })))
.await;
pipeline
.submit(resolved(2, Err(crate::CommitFailed::LogSealed)))
.await;
pipeline
.submit(resolved(3, Ok(crate::CommitAck { offset: 3 })))
.await;
assert_eq!(pipeline.drain().await, Commit::Sealed);
assert!(
pipeline.is_empty(),
"a barrier resolves every ticket, failure included"
);
}
#[tokio::test]
async fn an_aborted_ticket_ends_the_turn_cancelled_not_failed() {
let mut pipeline = TicketPipeline::default();
pipeline
.submit(resolved(1, Err(crate::CommitFailed::Aborted)))
.await;
let commit = pipeline.drain().await;
assert_eq!(commit, Commit::Aborted);
assert_eq!(commit.outcome(), Outcome::Cancelled);
}
#[tokio::test]
async fn a_dropped_ticket_channel_is_a_shutdown_not_a_seal() {
let mut pipeline = TicketPipeline::default();
let (ticket, tx) = ticket(1);
drop(tx);
pipeline.submit(ticket).await;
assert_eq!(pipeline.drain().await, Commit::Gone);
}
#[test]
#[cfg_attr(
not(debug_assertions),
ignore = "drives a debug_assert, which release builds compile out"
)]
#[should_panic(expected = "a ticket is not a commit")]
fn from_reply_refuses_to_treat_a_ticket_as_a_durable_commit() {
let (_tx, rx) = oneshot::channel();
let _ = Commit::from_reply(Some(crate::ProposeReply::Ticket(crate::CommitTicket {
id: "x".into(),
seq: 1,
ack: rx,
})));
}
#[test]
fn an_unconfirmed_commit_is_never_ok() {
assert!(!Commit::Unconfirmed.ok());
assert!(matches!(
Commit::Unconfirmed.outcome(),
Outcome::Error { .. }
));
}
#[test]
fn a_failed_barrier_c_replaces_every_end_reason_it_invalidates() {
let sealed = |outcome| match seal_end(TurnEnd::Outcome(outcome), Commit::Sealed) {
TurnEnd::Outcome(o) => o,
TurnEnd::Compact { .. } => unreachable!(),
};
for outcome in [
Outcome::Done { text: "hi".into() },
Outcome::TurnLimit,
Outcome::Refused,
Outcome::DoomLoop {
pattern: "read".into(),
},
Outcome::ToolFailureBudget {
tool: "bash".into(),
},
] {
assert!(
matches!(sealed(outcome.clone()), Outcome::Error { .. }),
"{outcome:?} names a record the log never took"
);
}
assert_eq!(sealed(Outcome::Cancelled), Outcome::Cancelled);
let first = Outcome::Error {
message: "the original cause".into(),
};
assert_eq!(sealed(first.clone()), first);
}
#[test]
fn a_clean_barrier_c_leaves_the_end_alone() {
let end = seal_end(
TurnEnd::Outcome(Outcome::Done { text: "hi".into() }),
Commit::Committed,
);
assert!(matches!(end, TurnEnd::Outcome(Outcome::Done { .. })));
}
}