use std::collections::VecDeque;
use std::path::PathBuf;
use std::sync::atomic::{AtomicU8, Ordering};
use std::sync::{Arc, Mutex};
use futures_util::StreamExt;
use hotl_context::compaction;
use hotl_platform::Clock;
use hotl_provider::{Provider, SamplingRequest, StreamEvent};
use hotl_store::SessionLog;
use hotl_tools::{
rules::{PermissionMode, Rules},
Registry,
};
use hotl_types::{assistant_text, EntryPayload, Item, SyntheticReason, Todo, TokenUsage};
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use crate::{turn, EngineConfig, EngineEvent, Outcome, SessionCmd, SessionDeps, TurnEnd};
pub(crate) const TAIL_RATIO: f64 = 0.3;
const SUMMARIZE_ATTEMPTS: u32 = 2;
const SUMMARIZE_MAX_TOKENS: u32 = 2_000;
const MAX_COMPACT_STREAK: u32 = 2;
const COMPACT_SUMMARIZE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(120);
const QUEUE_MAX: usize = 64;
const HELD_BYTES_MAX: usize = 64 * 1024;
const FOLD_MARK: &str = "\n[… later text truncated]";
struct PendingAck {
ack: tokio::sync::oneshot::Receiver<std::io::Result<hotl_store::Ack>>,
items: Vec<Item>,
id: String,
seq: u64,
stage: crate::SampleStage,
ticket: Option<tokio::sync::oneshot::Sender<Result<crate::CommitAck, crate::CommitFailed>>>,
}
pub struct ProjectionHead {
items: Arc<Vec<Item>>,
todos: Arc<Vec<Todo>>,
leaf: Option<String>,
epoch: u64,
}
impl ProjectionHead {
pub fn leaf(&self) -> Option<&str> {
self.leaf.as_deref()
}
pub fn epoch(&self) -> u64 {
self.epoch
}
pub fn snapshot(&self) -> Snapshot {
Snapshot {
durable: Arc::clone(&self.items),
tail: match hotl_tools::todo::render_reminder(&self.todos) {
Some(reminder) => Arc::new(vec![reminder]),
None => empty_tail(),
},
}
}
}
#[derive(Clone, Debug)]
pub struct Snapshot {
pub durable: Arc<Vec<Item>>,
pub tail: Arc<Vec<Item>>,
}
pub(crate) fn empty_tail() -> Arc<Vec<Item>> {
static EMPTY: std::sync::OnceLock<Arc<Vec<Item>>> = std::sync::OnceLock::new();
Arc::clone(EMPTY.get_or_init(|| Arc::new(Vec::new())))
}
pub(crate) fn head_channel() -> (
tokio::sync::watch::Sender<Arc<ProjectionHead>>,
tokio::sync::watch::Receiver<Arc<ProjectionHead>>,
) {
tokio::sync::watch::channel(Arc::new(ProjectionHead {
items: Arc::new(Vec::new()),
todos: Arc::new(Vec::new()),
leaf: None,
epoch: 0,
}))
}
struct Head {
tx: tokio::sync::watch::Sender<Arc<ProjectionHead>>,
items: Arc<Vec<Item>>,
todos: Arc<Vec<Todo>>,
leaf: Option<String>,
epoch: u64,
}
impl Head {
fn new(
tx: tokio::sync::watch::Sender<Arc<ProjectionHead>>,
items: Vec<Item>,
todos: Vec<Todo>,
) -> Self {
let mut head = Self {
tx,
items: Arc::new(items),
todos: Arc::new(todos),
leaf: None,
epoch: 0,
};
head.publish();
head
}
fn items(&self) -> &Arc<Vec<Item>> {
&self.items
}
fn todos(&self) -> &Arc<Vec<Todo>> {
&self.todos
}
fn apply(&mut self, item: Item) {
Arc::make_mut(&mut self.items).push(item);
}
fn advance(&mut self, leaf: String, seq: u64) {
self.leaf = Some(leaf);
self.epoch = seq;
}
fn set_todos(&mut self, todos: Vec<Todo>) {
self.todos = Arc::new(todos);
}
fn repoint(&mut self, items: Vec<Item>) {
self.items = Arc::new(items);
self.publish();
}
fn publish(&mut self) {
let _ = self.tx.send(Arc::new(ProjectionHead {
items: Arc::clone(&self.items),
todos: Arc::clone(&self.todos),
leaf: self.leaf.clone(),
epoch: self.epoch,
}));
}
}
#[derive(Clone, Copy, Debug)]
enum Boundary {
CommitSettled { stage: crate::SampleStage },
TurnEnded,
}
impl Boundary {
fn is_between_samples(self) -> bool {
match self {
Self::CommitSettled { stage } => stage == crate::SampleStage::AtBoundary,
Self::TurnEnded => true,
}
}
}
#[derive(Clone, Copy)]
enum Resolution {
Ack,
Abort,
}
#[derive(Default)]
struct Pipeline {
fifo: VecDeque<PendingAck>,
seq: u64,
}
impl Pipeline {
fn is_empty(&self) -> bool {
self.fifo.is_empty()
}
fn next_seq(&mut self) -> u64 {
self.seq += 1;
self.seq
}
async fn drain(&mut self, head: &mut Head, resolution: Resolution) {
while !self.is_empty() {
let acked = {
let entry = self.fifo.front_mut().expect("just checked non-empty");
await_ack(&mut entry.ack).await
};
let entry = self.fifo.pop_front().expect("just checked non-empty");
let settled = apply_ack(entry, acked, head, resolution);
head.publish();
settled.resolve();
}
}
}
enum Woke {
Ack(std::io::Result<hotl_store::Ack>),
Cmd(Option<SessionCmd>),
}
async fn next_ack(front: &mut Option<PendingAck>) -> std::io::Result<hotl_store::Ack> {
match front {
Some(entry) => await_ack(&mut entry.ack).await,
None => std::future::pending().await,
}
}
async fn await_ack(
ack: &mut tokio::sync::oneshot::Receiver<std::io::Result<hotl_store::Ack>>,
) -> std::io::Result<hotl_store::Ack> {
match ack.await {
Ok(result) => result,
Err(_) => Err(std::io::Error::other(
"the log writer stopped before the entry was committed",
)),
}
}
#[must_use = "a decided ticket must be resolved, or its proposer waits forever"]
struct Settled {
ticket: Option<tokio::sync::oneshot::Sender<Result<crate::CommitAck, crate::CommitFailed>>>,
resolved: Result<crate::CommitAck, crate::CommitFailed>,
}
impl Settled {
fn resolve(self) {
if let Some(ticket) = self.ticket {
let _ = ticket.send(self.resolved);
}
}
}
fn apply_ack(
entry: PendingAck,
acked: std::io::Result<hotl_store::Ack>,
head: &mut Head,
resolution: Resolution,
) -> Settled {
let resolved = match acked {
Ok(ack) => {
for item in entry.items {
head.apply(item);
}
head.advance(entry.id, entry.seq);
match resolution {
Resolution::Ack => Ok(crate::CommitAck { offset: ack.offset }),
Resolution::Abort => Err(crate::CommitFailed::Aborted),
}
}
Err(_) => Err(crate::CommitFailed::LogSealed),
};
Settled {
ticket: entry.ticket,
resolved,
}
}
pub(crate) struct SharedDeps {
pub provider: Arc<dyn Provider>,
pub registry: Arc<Registry>,
pub rules: Arc<Rules>,
mode: AtomicU8,
pub sandbox_enforced: bool,
pub clock: Arc<dyn Clock>,
pub system: Arc<str>,
pub cwd: PathBuf,
pub config: EngineConfig,
pub snapshots: Option<Arc<dyn crate::Snapshotter>>,
pub hooks: Option<Arc<dyn crate::hooks::Hooks>>,
hook_mask: Arc<AtomicU8>,
pub notifications: crate::hooks::NotificationDrain,
masker: Arc<hotl_store::Masker>,
rules_epoch: std::sync::atomic::AtomicU32,
head_rx: tokio::sync::watch::Receiver<Arc<ProjectionHead>>,
}
fn mode_to_u8(mode: PermissionMode) -> u8 {
match mode {
PermissionMode::Ask => 0,
PermissionMode::Auto => 1,
PermissionMode::Plan => 2,
PermissionMode::DontAsk => 3,
}
}
fn u8_to_mode(v: u8) -> PermissionMode {
match v {
1 => PermissionMode::Auto,
2 => PermissionMode::Plan,
3 => PermissionMode::DontAsk,
_ => PermissionMode::Ask,
}
}
impl SharedDeps {
fn new(
deps: SessionDeps,
notifications: crate::hooks::NotificationDrain,
head_rx: tokio::sync::watch::Receiver<Arc<ProjectionHead>>,
) -> (Self, SessionLog) {
let mode = AtomicU8::new(mode_to_u8(deps.rules.mode()));
let hook_mask = deps
.hooks
.as_ref()
.and_then(|h| h.mask_handle())
.unwrap_or_else(|| {
Arc::new(AtomicU8::new(
deps.hooks
.as_ref()
.map_or(crate::hooks::EventMask::NONE, |h| h.event_mask())
.bits(),
))
});
let masker = deps.log.masker_handle();
let shared = Self {
provider: deps.provider,
registry: deps.registry,
rules: deps.rules,
mode,
sandbox_enforced: deps.sandbox_enforced,
clock: deps.clock,
system: deps.system.into(),
cwd: deps.cwd,
config: deps.config,
snapshots: deps.snapshots,
hooks: deps.hooks,
hook_mask,
notifications,
masker,
rules_epoch: std::sync::atomic::AtomicU32::new(0),
head_rx,
};
(shared, deps.log)
}
pub(crate) fn head(&self) -> tokio::sync::watch::Receiver<Arc<ProjectionHead>> {
self.head_rx.clone()
}
pub(crate) fn hook_mask(&self) -> crate::hooks::EventMask {
crate::hooks::mask_of(&self.hook_mask)
}
pub(crate) fn effective_mode(&self) -> PermissionMode {
u8_to_mode(self.mode.load(Ordering::Relaxed))
}
fn set_mode(&self, mode: PermissionMode) -> PermissionMode {
let mode = hotl_tools::rules::enforced_mode(mode);
self.mode.store(mode_to_u8(mode), Ordering::Relaxed);
mode
}
async fn append(
&self,
log: &mut SessionLog,
pipeline: &mut Pipeline,
head: &mut Head,
payload: EntryPayload,
) -> bool {
pipeline.drain(head, Resolution::Ack).await;
let seq = pipeline.next_seq();
if log
.append_acked(&payload, self.clock.now_ms())
.await
.is_err()
{
return false;
}
if let EntryPayload::Item { item } = payload {
head.apply(item);
}
if let Some(id) = log.last_id() {
head.advance(id.to_string(), seq);
}
head.publish();
true
}
pub(crate) fn masker(&self) -> &Arc<hotl_store::Masker> {
&self.masker
}
pub(crate) fn rules_epoch(&self) -> u32 {
self.rules_epoch.load(Ordering::Relaxed)
}
async fn append_prepared(
&self,
log: &mut SessionLog,
prepared: hotl_store::PreparedPayload,
) -> bool {
log.append_prepared(prepared, self.clock.now_ms())
.await
.is_ok()
}
fn forward_prepared(
&self,
log: &mut SessionLog,
prepared: hotl_store::PreparedPayload,
) -> std::io::Result<hotl_store::Forwarded> {
log.forward_prepared(prepared, self.clock.now_ms())
}
fn forward_group(
&self,
log: &mut SessionLog,
group: Vec<hotl_store::PreparedPayload>,
) -> std::io::Result<hotl_store::Forwarded> {
log.forward_group(group, self.clock.now_ms())
}
}
pub(crate) async fn run(
mut deps: SessionDeps,
mut cmd_rx: mpsc::Receiver<SessionCmd>,
cmd_tx: mpsc::WeakSender<SessionCmd>,
events: mpsc::Sender<EngineEvent>,
current_turn: Arc<Mutex<CancellationToken>>,
notifications: crate::hooks::NotificationDrain,
head_tx: tokio::sync::watch::Sender<Arc<ProjectionHead>>,
) {
let head_rx = head_tx.subscribe();
let mut head = Head::new(
head_tx,
pair_tool_results(std::mem::take(&mut deps.initial_items)),
std::mem::take(&mut deps.initial_todos),
);
let mut running = false;
let mut queue: VecDeque<(String, Option<SyntheticReason>)> = VecDeque::new();
let mut held_steers: Vec<String> = Vec::new();
let (shared, mut log) = SharedDeps::new(deps, notifications, head_rx);
let shared = Arc::new(shared);
let mut carry_usage = TokenUsage::default();
let mut compact_streak: u32 = 0;
let mut pipeline = Pipeline::default();
loop {
let mut front = pipeline.fifo.pop_front();
let woke = tokio::select! {
biased;
acked = next_ack(&mut front) => Woke::Ack(acked),
cmd = cmd_rx.recv() => Woke::Cmd(cmd),
};
let cmd = match woke {
Woke::Ack(acked) => {
let entry = front.expect("the ack arm only runs with a front entry");
let stage = entry.stage;
let settled = apply_ack(entry, acked, &mut head, Resolution::Ack);
release_steers(
&shared,
&mut log,
&mut head,
&mut pipeline,
&mut held_steers,
Boundary::CommitSettled { stage },
)
.await;
head.publish();
settled.resolve();
continue;
}
Woke::Cmd(cmd) => {
if let Some(entry) = front {
pipeline.fifo.push_front(entry);
}
match cmd {
Some(cmd) => cmd,
None => break,
}
}
};
match cmd {
SessionCmd::Prompt(text) => {
running = admit_prompt(
&shared,
&mut log,
&mut head,
&mut pipeline,
&mut queue,
running,
text,
None,
&cmd_tx,
&events,
¤t_turn,
)
.await;
}
SessionCmd::PromptTagged { text, synthetic } => {
running = admit_prompt(
&shared,
&mut log,
&mut head,
&mut pipeline,
&mut queue,
running,
text,
Some(synthetic),
&cmd_tx,
&events,
¤t_turn,
)
.await;
}
SessionCmd::Continue => {
if !running && crate::needs_continuation(head.items()) {
spawn_turn(&shared, &cmd_tx, &events, ¤t_turn);
running = true;
}
}
SessionCmd::Steer(text) => {
admit_steer(
&shared,
&mut log,
&mut head,
&mut pipeline,
&mut held_steers,
running,
text,
)
.await
}
SessionCmd::Rename(name) => {
let _ = shared
.append(
&mut log,
&mut pipeline,
&mut head,
EntryPayload::Rename { name },
)
.await;
}
SessionCmd::SetMode(mode) => {
let mode = shared.set_mode(mode);
let _ = shared
.append(
&mut log,
&mut pipeline,
&mut head,
EntryPayload::ModeSet {
mode: mode.as_str().into(),
},
)
.await;
}
SessionCmd::SetTodos(new_todos) => {
head.set_todos(new_todos);
let items = (**head.todos()).clone();
let _ = shared
.append(
&mut log,
&mut pipeline,
&mut head,
EntryPayload::Todos {
items: items.clone(),
},
)
.await;
let _ = events.send(EngineEvent::TodosChanged { items }).await;
}
SessionCmd::Propose { entries, reply } => {
let committed = commit(&shared, &mut log, &mut head, &mut pipeline, entries).await;
let _ = reply.send(committed);
}
SessionCmd::ProposePrepared {
proposal,
stage,
mode,
reply,
} => {
let result = commit_prepared(
&shared,
&mut log,
&mut head,
&mut pipeline,
proposal,
mode,
stage,
)
.await;
if mode == crate::AckMode::Sync
&& !matches!(result, crate::ProposeReply::StaleEpoch)
{
release_steers(
&shared,
&mut log,
&mut head,
&mut pipeline,
&mut held_steers,
Boundary::CommitSettled { stage },
)
.await;
}
let _ = reply.send(result);
}
SessionCmd::WriteBlob {
tool_use_id,
content,
reply,
} => {
let result = match log.write_blob_acked(&tool_use_id, &content).await {
Ok(path) => Ok(path.display().to_string()),
Err(_) => Err(content), };
let _ = reply.send(result);
}
SessionCmd::TurnFinished { end, usage } => {
close_open_batch(&shared, &mut log, &mut head, &mut pipeline).await;
release_steers(
&shared,
&mut log,
&mut head,
&mut pipeline,
&mut held_steers,
Boundary::TurnEnded,
)
.await;
on_turn_finished(
TurnFinishedCtx {
shared: &shared,
log: &mut log,
head: &mut head,
pipeline: &mut pipeline,
queue: &mut queue,
running: &mut running,
carry_usage: &mut carry_usage,
compact_streak: &mut compact_streak,
cmd_tx: &cmd_tx,
events: &events,
current_turn: ¤t_turn,
},
end,
usage,
)
.await;
}
SessionCmd::BumpRulesEpoch => {
shared.rules_epoch.fetch_add(1, Ordering::Relaxed);
}
}
}
crate::hooks::hook_gate!(
shared.hooks,
shared.hook_mask(),
crate::hooks::EventMask::SESSION_END,
|hooks| {
crate::hooks::call_session_end(hooks).await;
},
else {}
);
}
struct TurnFinishedCtx<'a> {
shared: &'a Arc<SharedDeps>,
log: &'a mut SessionLog,
head: &'a mut Head,
pipeline: &'a mut Pipeline,
queue: &'a mut VecDeque<(String, Option<SyntheticReason>)>,
running: &'a mut bool,
carry_usage: &'a mut TokenUsage,
compact_streak: &'a mut u32,
cmd_tx: &'a mpsc::WeakSender<SessionCmd>,
events: &'a mpsc::Sender<EngineEvent>,
current_turn: &'a Arc<Mutex<CancellationToken>>,
}
async fn on_turn_finished(ctx: TurnFinishedCtx<'_>, end: TurnEnd, mut usage: TokenUsage) {
let outcome = match end {
TurnEnd::Outcome(outcome) => Some(outcome),
TurnEnd::Compact { spec, cont } => {
*ctx.carry_usage += usage;
usage = TokenUsage::default();
try_compact(
ctx.shared,
ctx.log,
ctx.head,
ctx.pipeline,
ctx.compact_streak,
spec,
cont,
ctx.cmd_tx,
ctx.events,
ctx.current_turn,
)
.await
}
};
if let Some(outcome) = outcome {
*ctx.compact_streak = 0;
let mut total = usage;
total += std::mem::take(ctx.carry_usage);
*ctx.running = end_turn(
ctx.shared,
ctx.log,
ctx.head,
ctx.pipeline,
ctx.queue,
outcome,
total,
ctx.cmd_tx,
ctx.events,
ctx.current_turn,
)
.await;
}
}
#[allow(clippy::too_many_arguments)]
async fn try_compact(
shared: &Arc<SharedDeps>,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
compact_streak: &mut u32,
spec: Option<crate::SpecDigest>,
cont: Box<crate::TurnContinuation>,
cmd_tx: &mpsc::WeakSender<SessionCmd>,
events: &mpsc::Sender<EngineEvent>,
current_turn: &Arc<Mutex<CancellationToken>>,
) -> Option<Outcome> {
if cont.samples_since_compact > 0 {
*compact_streak = 0;
}
*compact_streak += 1;
let cancel = current_turn
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
let compacted = if *compact_streak > MAX_COMPACT_STREAK {
Err("context window exhausted — compaction can no longer make room".into())
} else {
tokio::select! {
biased;
_ = cancel.cancelled() => return Some(Outcome::Cancelled),
compacted = compact(shared, log, head, pipeline, spec) => compacted,
}
};
match compacted {
Ok(degraded) => {
let _ = events.send(EngineEvent::Compacted { degraded }).await;
if cancel.is_cancelled() {
return Some(Outcome::Cancelled);
}
respawn_turn(shared, cmd_tx, events, cancel, *cont);
None }
Err(message) => Some(Outcome::Error { message }),
}
}
#[allow(clippy::too_many_arguments)]
async fn end_turn(
shared: &Arc<SharedDeps>,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
queue: &mut VecDeque<(String, Option<SyntheticReason>)>,
outcome: Outcome,
usage: TokenUsage,
cmd_tx: &mpsc::WeakSender<SessionCmd>,
events: &mpsc::Sender<EngineEvent>,
current_turn: &Arc<Mutex<CancellationToken>>,
) -> bool {
annotate(shared, log, head, pipeline, &outcome).await;
crate::hooks::hook_gate!(
shared.hooks,
shared.hook_mask(),
crate::hooks::EventMask::NOTIFICATION,
|hooks| {
crate::hooks::notify(
hooks,
&shared.notifications,
crate::hooks::NotificationKind::Done,
outcome_detail(&outcome),
);
},
else {}
);
let _ = events.send(EngineEvent::TurnDone { outcome, usage }).await;
match queue.pop_front() {
Some((next, synthetic)) => {
start_turn(
shared,
log,
head,
pipeline,
next,
synthetic,
cmd_tx,
events,
current_turn,
)
.await
}
None => {
crate::hooks::hook_gate!(
shared.hooks,
shared.hook_mask(),
crate::hooks::EventMask::NOTIFICATION,
|hooks| {
crate::hooks::notify(
hooks,
&shared.notifications,
crate::hooks::NotificationKind::Idle,
"awaiting a prompt",
);
},
else {}
);
false
}
}
}
fn outcome_detail(outcome: &Outcome) -> String {
match outcome {
Outcome::Done { text } => text.clone(),
other => format!("{other:?}"),
}
}
fn fold_into(dst: &mut String, text: &str, max_bytes: usize) {
if dst.ends_with(FOLD_MARK) {
dst.truncate(dst.len() - FOLD_MARK.len());
}
if !dst.is_empty() {
dst.push_str("\n\n");
}
dst.push_str(text);
if dst.len() > max_bytes {
let mut end = max_bytes;
while !dst.is_char_boundary(end) {
end -= 1;
}
dst.truncate(end);
dst.push_str(FOLD_MARK);
}
}
fn awaiting_tool_results(items: &[Item]) -> bool {
matches!(
items.last(),
Some(Item::Assistant { blocks }) if !hotl_types::assistant_tool_uses(blocks).is_empty()
)
}
async fn admit_steer(
shared: &SharedDeps,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
held: &mut Vec<String>,
running: bool,
text: String,
) {
if running || awaiting_tool_results(head.items()) {
let total: usize = held.iter().map(String::len).sum();
match held.last_mut() {
Some(last) if total >= HELD_BYTES_MAX => fold_into(last, &text, HELD_BYTES_MAX),
_ => held.push(text),
}
return;
}
append_steer(shared, log, head, pipeline, text).await;
}
async fn append_steer(
shared: &SharedDeps,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
text: String,
) {
debug_assert!(
!awaiting_tool_results(head.items()),
"a steer must never land while a tool batch is open"
);
shared
.append(
log,
pipeline,
head,
EntryPayload::Item {
item: Item::User {
text,
synthetic: Some(SyntheticReason::Steer),
},
},
)
.await;
}
async fn release_steers(
shared: &SharedDeps,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
held: &mut Vec<String>,
at: Boundary,
) {
debug_assert!(
at.is_between_samples(),
"a held steer may only land between samples: releasing behind an in-sample \
commit would put it ahead of the assistant item the model is still \
producing — the inversion 72a6f1b fixed ({at:?})"
);
if held.is_empty() || awaiting_tool_results(head.items()) {
return;
}
for text in std::mem::take(held) {
append_steer(shared, log, head, pipeline, text).await;
}
}
async fn close_open_batch(
shared: &SharedDeps,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
) {
let Some(Item::Assistant { blocks }) = head.items().last() else {
return;
};
let uses = hotl_types::assistant_tool_uses(blocks);
if uses.is_empty() {
return;
}
let payload = EntryPayload::Item {
item: Item::ToolResults {
results: uses
.iter()
.map(|tu| hotl_types::ToolResultItem {
tool_use_id: tu.id.clone(),
content: "Not executed (the turn ended first).".into(),
is_error: true,
})
.collect(),
},
};
shared.append(log, pipeline, head, payload).await;
}
pub(crate) fn pair_tool_results(items: Vec<Item>) -> Vec<Item> {
let mut out: Vec<Item> = Vec::with_capacity(items.len());
let mut stranded: Vec<Item> = Vec::new();
for item in items {
if !awaiting_tool_results(&out) && stranded.is_empty() {
out.push(item);
continue;
}
match item {
Item::ToolResults { .. } => {
out.push(item);
out.append(&mut stranded);
}
Item::Assistant { .. } => {
out.append(&mut stranded);
out.push(item);
}
_ => stranded.push(item),
}
}
out.append(&mut stranded);
out
}
async fn commit(
shared: &SharedDeps,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
entries: Vec<EntryPayload>,
) -> bool {
for payload in entries {
if !shared.append(log, pipeline, head, payload).await {
return false;
}
}
true
}
async fn commit_prepared(
shared: &SharedDeps,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
proposal: crate::EntryProposal,
mode: crate::AckMode,
stage: crate::SampleStage,
) -> crate::ProposeReply {
let current_epoch = shared.rules_epoch();
if proposal
.entries()
.iter()
.any(|e| e.payload().rules_epoch() < current_epoch)
{
return crate::ProposeReply::StaleEpoch;
}
if proposal.is_empty() {
return crate::ProposeReply::Committed;
}
match mode {
crate::AckMode::Sync => {
pipeline.drain(head, Resolution::Ack).await;
for entry in into_entries(proposal) {
let (payload, item) = entry.into_parts();
let seq = pipeline.next_seq();
if !shared.append_prepared(log, payload).await {
return crate::ProposeReply::Sealed;
}
if let Some(item) = item {
head.apply(item);
}
if let Some(id) = log.last_id() {
head.advance(id.to_string(), seq);
}
head.publish();
}
crate::ProposeReply::Committed
}
crate::AckMode::Pipelined => match proposal {
crate::EntryProposal::Single(entry) => {
let (payload, item) = entry.into_parts();
let seq = pipeline.next_seq();
match shared.forward_prepared(log, payload) {
Ok(forwarded) => crate::ProposeReply::Ticket(push_pending(
pipeline,
forwarded,
item.into_iter().collect(),
seq,
stage,
)),
Err(_) => crate::ProposeReply::Sealed,
}
}
crate::EntryProposal::Group(entries) => {
let mut payloads = Vec::with_capacity(entries.len());
let mut items = Vec::with_capacity(entries.len());
let mut seq = 0;
for entry in entries {
let (payload, item) = entry.into_parts();
seq = pipeline.next_seq();
payloads.push(payload);
items.extend(item);
}
match shared.forward_group(log, payloads) {
Ok(forwarded) => crate::ProposeReply::Ticket(push_pending(
pipeline, forwarded, items, seq, stage,
)),
Err(_) => crate::ProposeReply::Sealed,
}
}
},
}
}
fn into_entries(proposal: crate::EntryProposal) -> Vec<crate::PreparedEntry> {
match proposal {
crate::EntryProposal::Single(entry) => vec![entry],
crate::EntryProposal::Group(entries) => entries,
}
}
fn push_pending(
pipeline: &mut Pipeline,
forwarded: hotl_store::Forwarded,
items: Vec<Item>,
seq: u64,
stage: crate::SampleStage,
) -> crate::CommitTicket {
let (tx, rx) = tokio::sync::oneshot::channel();
let ticket = crate::CommitTicket {
id: forwarded.id.clone(),
seq,
ack: rx,
};
pipeline.fifo.push_back(PendingAck {
ack: forwarded.ack,
items,
id: forwarded.id,
seq,
stage,
ticket: Some(tx),
});
ticket
}
async fn annotate(
shared: &SharedDeps,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
outcome: &Outcome,
) {
let reason = match outcome {
Outcome::Cancelled => Some("user interrupt".to_string()),
Outcome::TurnLimit => Some(format!("max_turns ({}) reached", shared.config.max_turns)),
Outcome::DoomLoop { pattern } => Some(format!("doom loop: {pattern}")),
Outcome::ToolFailureBudget { tool } => Some(format!("tool failure budget: {tool}")),
Outcome::Error { message } => Some(format!("error: {message}")),
Outcome::Done { .. } | Outcome::Refused => None,
};
if let Some(reason) = reason {
shared
.append(log, pipeline, head, EntryPayload::Cancelled { reason })
.await;
}
}
async fn compact(
shared: &SharedDeps,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
spec: Option<crate::SpecDigest>,
) -> Result<bool, String> {
pipeline.drain(head, Resolution::Abort).await;
if !shared.config.compaction_reset {
if let Some(spec) = spec {
if spec.prefix_end < spec.kept_from && spec.kept_from <= head.items().len() {
let digest = vec![compaction::digest_item(&spec.text)];
let payload = EntryPayload::Compaction {
digest: digest.clone(),
prefix_end: spec.prefix_end,
kept_from: spec.kept_from,
degraded: false,
};
if !shared.append(log, pipeline, head, payload).await {
return Err("session log is sealed".into());
}
let plan = compaction::Plan {
prefix_end: spec.prefix_end,
kept_from: spec.kept_from,
};
head.repoint(compaction::apply(head.items(), &plan, &digest));
return Ok(false);
}
}
}
let tail_budget = (shared.config.context_window as f64 * TAIL_RATIO) as u64;
let Some(plan) = compaction::plan(head.items(), tail_budget) else {
return Err("context window exhausted — nothing left to compact".into());
};
let plan = if shared.config.compaction_reset {
compaction::Plan {
prefix_end: plan.prefix_end,
kept_from: head.items().len(),
}
} else {
plan
};
let snapshot = Arc::clone(head.items());
let folded = &snapshot[plan.prefix_end..plan.kept_from];
let (digest, degraded) =
match summarize_bounded(summarize(shared, folded), COMPACT_SUMMARIZE_TIMEOUT).await {
Some(text) => (vec![compaction::digest_item(&text)], false),
None => (vec![compaction::floor_digest()], true),
};
let payload = EntryPayload::Compaction {
digest: digest.clone(),
prefix_end: plan.prefix_end,
kept_from: plan.kept_from,
degraded,
};
if !shared.append(log, pipeline, head, payload).await {
return Err("session log is sealed".into());
}
head.repoint(compaction::apply(head.items(), &plan, &digest));
Ok(degraded)
}
async fn summarize_bounded(
fut: impl std::future::Future<Output = Option<String>>,
bound: std::time::Duration,
) -> Option<String> {
tokio::time::timeout(bound, fut).await.ok().flatten()
}
pub(crate) async fn summarize(shared: &SharedDeps, folded: &[Item]) -> Option<String> {
let model = shared
.config
.fast_model
.clone()
.unwrap_or_else(|| shared.config.model.clone());
let request = SamplingRequest {
model,
max_tokens: SUMMARIZE_MAX_TOKENS,
system: compaction::SUMMARIZE_SYSTEM.into(),
items: Arc::new(vec![Item::User {
text: compaction::summarize_prompt(folded),
synthetic: None,
}]),
ephemeral_tail: empty_tail(),
tools: Vec::new().into(),
thinking: false,
cache: hotl_provider::CachePolicy::Off,
turn_context: None,
};
for _ in 0..SUMMARIZE_ATTEMPTS {
let mut stream = shared.provider.stream(request.clone());
let mut text: Option<String> = None;
while let Some(event) = stream.next().await {
match event {
Ok(StreamEvent::Completed { blocks, .. }) => text = Some(assistant_text(&blocks)),
Ok(_) => {}
Err(_) => {
text = None;
crate::turn::drain_to_end(&mut stream).await;
break;
}
}
}
if let Some(t) = text.filter(|t| !t.trim().is_empty()) {
return Some(t);
}
}
None
}
#[allow(clippy::too_many_arguments)]
async fn admit_prompt(
shared: &Arc<SharedDeps>,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
queue: &mut VecDeque<(String, Option<SyntheticReason>)>,
running: bool,
text: String,
synthetic: Option<SyntheticReason>,
cmd_tx: &mpsc::WeakSender<SessionCmd>,
events: &mpsc::Sender<EngineEvent>,
current_turn: &Arc<Mutex<CancellationToken>>,
) -> bool {
if running {
let full = queue.len() >= QUEUE_MAX;
match queue.back_mut() {
Some(last) if full => fold_into(&mut last.0, &text, HELD_BYTES_MAX),
_ => queue.push_back((text, synthetic)),
}
let _ = events.send(EngineEvent::PromptQueued).await;
return true;
}
start_turn(
shared,
log,
head,
pipeline,
text,
synthetic,
cmd_tx,
events,
current_turn,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn start_turn(
shared: &Arc<SharedDeps>,
log: &mut SessionLog,
head: &mut Head,
pipeline: &mut Pipeline,
text: String,
synthetic: Option<SyntheticReason>,
cmd_tx: &mpsc::WeakSender<SessionCmd>,
events: &mpsc::Sender<EngineEvent>,
current_turn: &Arc<Mutex<CancellationToken>>,
) -> bool {
let prompt_for_hooks = text.clone();
let payload = EntryPayload::Item {
item: Item::User { text, synthetic },
};
if !shared.append(log, pipeline, head, payload).await {
let _ = events
.send(EngineEvent::TurnDone {
outcome: Outcome::Error {
message: "session log is sealed".into(),
},
usage: TokenUsage::default(),
})
.await;
return false;
}
crate::hooks::hook_gate!(
shared.hooks,
shared.hook_mask(),
crate::hooks::EventMask::USER_PROMPT,
|hooks| {
if let Some(context) = crate::hooks::call_user_prompt(hooks, &prompt_for_hooks).await {
let reminder = EntryPayload::Item {
item: Item::User {
text: format!("<system-reminder>{context}</system-reminder>"),
synthetic: Some(SyntheticReason::SystemReminder),
},
};
shared.append(log, pipeline, head, reminder).await;
}
},
else {}
);
spawn_turn(shared, cmd_tx, events, current_turn);
true
}
fn spawn_turn(
shared: &Arc<SharedDeps>,
cmd_tx: &mpsc::WeakSender<SessionCmd>,
events: &mpsc::Sender<EngineEvent>,
current_turn: &Arc<Mutex<CancellationToken>>,
) {
let token = CancellationToken::new();
*current_turn
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = token.clone();
respawn_turn(
shared,
cmd_tx,
events,
token,
crate::TurnContinuation::default(),
);
}
fn respawn_turn(
shared: &Arc<SharedDeps>,
cmd_tx: &mpsc::WeakSender<SessionCmd>,
events: &mpsc::Sender<EngineEvent>,
token: CancellationToken,
cont: crate::TurnContinuation,
) {
let Some(cmd_tx) = cmd_tx.upgrade() else {
return;
};
let supervisor_tx = cmd_tx.clone();
let handle = tokio::spawn(turn::run(
shared.clone(),
cmd_tx,
events.clone(),
token,
cont,
));
tokio::spawn(async move {
if handle.await.is_err() {
let _ = supervisor_tx
.send(SessionCmd::TurnFinished {
end: TurnEnd::Outcome(Outcome::Error {
message: "the turn ended unexpectedly (internal error). \
The session is intact — retry, or rephrase the request."
.into(),
}),
usage: TokenUsage::default(),
})
.await;
}
});
}
#[cfg(test)]
mod tests {
use super::COMPACT_SUMMARIZE_TIMEOUT;
use super::{
awaiting_tool_results, commit_prepared, compact, fold_into, pair_tool_results,
release_steers, summarize_bounded, Pipeline, Resolution, SharedDeps,
};
use hotl_store::SessionLog;
use hotl_types::{EntryPayload, Item, SyntheticReason, ToolResultItem};
use serde_json::json;
use std::sync::atomic::Ordering;
use std::sync::Arc;
#[tokio::test(start_paused = true)]
async fn a_hung_inline_summarize_degrades_instead_of_wedging() {
let hung = summarize_bounded(std::future::pending(), COMPACT_SUMMARIZE_TIMEOUT).await;
assert!(
hung.is_none(),
"a hung summarize must degrade to the floor digest, not stall the command loop"
);
let answered = summarize_bounded(
std::future::ready(Some("DIGEST".to_string())),
COMPACT_SUMMARIZE_TIMEOUT,
)
.await;
assert_eq!(
answered.as_deref(),
Some("DIGEST"),
"a summarize that answers inside the bound must still be used"
);
}
fn user(text: &str) -> Item {
Item::User {
text: text.into(),
synthetic: Some(SyntheticReason::Steer),
}
}
fn calls(id: &str) -> Item {
Item::Assistant {
blocks: vec![json!({"type": "tool_use", "id": id, "name": "read", "input": {}})],
}
}
fn says(text: &str) -> Item {
Item::Assistant {
blocks: vec![json!({"type": "text", "text": text})],
}
}
fn answers(id: &str) -> Item {
Item::ToolResults {
results: vec![ToolResultItem {
tool_use_id: id.into(),
content: "ok".into(),
is_error: false,
}],
}
}
#[test]
fn folding_bounds_the_buffer_and_discloses_the_truncation() {
let mut dst = String::from("first");
fold_into(&mut dst, &"x".repeat(10_000), 128);
assert!(
dst.len() <= 128 + 64,
"fold must bound the entry, got {}",
dst.len()
);
assert!(
dst.starts_with("first"),
"the oldest text is kept, not clobbered"
);
assert!(
dst.contains("truncated"),
"truncation must be disclosed in-band"
);
for _ in 0..1_000 {
fold_into(&mut dst, "more", 128);
}
assert!(dst.len() <= 128 + 64, "got {}", dst.len());
assert!(dst.starts_with("first"));
}
#[test]
fn folding_under_the_cap_keeps_every_word() {
let mut dst = String::from("first");
fold_into(&mut dst, "second", 1_024);
assert!(dst.contains("first") && dst.contains("second"));
assert!(!dst.contains("truncated"), "nothing was dropped: {dst}");
}
#[test]
fn only_unanswered_tool_calls_hold_the_batch_open() {
assert!(awaiting_tool_results(&[calls("t1")]));
assert!(!awaiting_tool_results(&[says("hello")]));
assert!(!awaiting_tool_results(&[calls("t1"), answers("t1")]));
assert!(!awaiting_tool_results(&[]));
}
#[test]
fn a_stranded_steer_moves_behind_the_results_it_interrupted() {
let repaired = pair_tool_results(vec![calls("t1"), user("wait"), answers("t1")]);
assert_eq!(repaired, vec![calls("t1"), answers("t1"), user("wait")]);
}
#[test]
fn several_stranded_items_keep_their_order() {
let repaired = pair_tool_results(vec![
calls("t1"),
user("one"),
user("two"),
answers("t1"),
says("done"),
]);
assert_eq!(
repaired,
vec![
calls("t1"),
answers("t1"),
user("one"),
user("two"),
says("done"),
]
);
}
#[test]
fn already_paired_history_is_left_alone() {
let good = vec![
user("start"),
calls("t1"),
answers("t1"),
user("next"),
says("done"),
];
assert_eq!(pair_tool_results(good.clone()), good);
}
#[test]
fn a_gap_with_no_results_coming_is_not_reordered() {
let orphaned = vec![calls("t1"), user("never answered"), says("moved on")];
assert_eq!(pair_tool_results(orphaned.clone()), orphaned);
}
#[test]
fn a_trailing_gap_survives_repair() {
let trailing = vec![calls("t1"), user("last word")];
assert_eq!(pair_tool_results(trailing.clone()), trailing);
}
fn test_deps(dir: &std::path::Path, log: hotl_store::SessionLog) -> crate::SessionDeps {
crate::SessionDeps {
provider: Arc::new(hotl_provider::ScriptedProvider::new(vec![])),
registry: Arc::new(hotl_tools::Registry::builtin()),
rules: Arc::new(hotl_tools::rules::Rules::default()),
sandbox_enforced: false,
clock: Arc::new(hotl_platform::SystemClock),
log,
system: "sys".into(),
cwd: dir.to_path_buf(),
snapshots: None,
hooks: None,
initial_items: Vec::new(),
initial_todos: Vec::new(),
config: crate::EngineConfig::default(),
}
}
fn test_shared(dir: &std::path::Path) -> (Arc<SharedDeps>, SessionLog, super::Head) {
let log = SessionLog::create(dir, "m", None, hotl_store::Masker::empty(), 0).expect("log");
let (head_tx, head_rx) = super::head_channel();
let head = super::Head::new(head_tx, Vec::new(), Vec::new());
let (shared, log) = SharedDeps::new(
test_deps(dir, log),
crate::hooks::NotificationDrain::new(),
head_rx,
);
(Arc::new(shared), log, head)
}
#[tokio::test]
async fn commit_prepared_rejects_an_entry_whose_epoch_predates_current_and_commits_nothing() {
let dir = tempfile::tempdir().unwrap();
let (shared, mut log, mut head) = test_shared(dir.path());
let before = std::fs::read_to_string(log.path()).unwrap();
let old_epoch = shared.rules_epoch();
shared.rules_epoch.fetch_add(1, Ordering::Relaxed);
assert!(shared.rules_epoch() > old_epoch);
let payload = hotl_types::EntryPayload::Usage {
usage: hotl_types::TokenUsage::default(),
};
let prepared =
hotl_store::prepare_payload(&payload, &hotl_store::Masker::empty(), old_epoch)
.expect("prepare");
let entries = vec![crate::PreparedEntry::new(prepared, None)];
let result = commit_prepared(
&shared,
&mut log,
&mut head,
&mut Pipeline::default(),
crate::EntryProposal::of(entries),
crate::AckMode::Sync,
crate::SampleStage::AtBoundary,
)
.await;
assert!(matches!(result, crate::ProposeReply::StaleEpoch));
assert!(
head.items().is_empty(),
"a stale proposal must not touch the projection"
);
let after = std::fs::read_to_string(log.path()).unwrap();
assert_eq!(before, after, "a stale proposal must not reach the log");
}
#[tokio::test]
async fn commit_prepared_accepts_an_entry_whose_epoch_is_not_older_than_current() {
let dir = tempfile::tempdir().unwrap();
let (shared, mut log, mut head) = test_shared(dir.path());
let newer_epoch = shared.rules_epoch() + 1;
let payload = hotl_types::EntryPayload::Usage {
usage: hotl_types::TokenUsage::default(),
};
let prepared =
hotl_store::prepare_payload(&payload, &hotl_store::Masker::empty(), newer_epoch)
.expect("prepare");
let entries = vec![crate::PreparedEntry::new(prepared, None)];
let result = commit_prepared(
&shared,
&mut log,
&mut head,
&mut Pipeline::default(),
crate::EntryProposal::of(entries),
crate::AckMode::Sync,
crate::SampleStage::AtBoundary,
)
.await;
assert!(matches!(result, crate::ProposeReply::Committed));
}
fn prepared(shared: &SharedDeps, item: Item) -> crate::PreparedEntry {
let payload = EntryPayload::Item { item: item.clone() };
let prepared = hotl_store::prepare_payload(
&payload,
&hotl_store::Masker::empty(),
shared.rules_epoch(),
)
.expect("prepare");
crate::PreparedEntry::new(prepared, Some(item))
}
#[tokio::test]
async fn a_pipelined_proposal_answers_with_a_ticket_before_the_projection_moves() {
let dir = tempfile::tempdir().unwrap();
let (shared, mut log, mut head) = test_shared(dir.path());
let mut pipeline = Pipeline::default();
let reply = commit_prepared(
&shared,
&mut log,
&mut head,
&mut pipeline,
crate::EntryProposal::of(vec![prepared(&shared, user("hi"))]),
crate::AckMode::Pipelined,
crate::SampleStage::AtBoundary,
)
.await;
let crate::ProposeReply::Ticket(ticket) = reply else {
panic!("Pipelined must answer with a ticket, got {reply:?}")
};
assert_eq!(ticket.seq, 1, "seq is assigned at validation, eagerly");
assert!(!ticket.id.is_empty(), "so is the ulid");
assert!(
head.items().is_empty(),
"the projection advances only on ack, never on forward"
);
pipeline.drain(&mut head, Resolution::Ack).await;
assert_eq!(head.items().len(), 1, "…and it advances when the ack lands");
let ack = ticket
.ack
.await
.expect("the actor resolves the ticket")
.expect("committed");
assert!(ack.offset > 0, "the ticket carries the byte offset");
}
#[tokio::test]
async fn the_pipeline_advances_the_projection_in_fifo_order() {
let dir = tempfile::tempdir().unwrap();
let (shared, mut log, mut head) = test_shared(dir.path());
let mut pipeline = Pipeline::default();
let mut tickets = Vec::new();
for text in ["one", "two", "three"] {
let reply = commit_prepared(
&shared,
&mut log,
&mut head,
&mut pipeline,
crate::EntryProposal::of(vec![prepared(&shared, user(text))]),
crate::AckMode::Pipelined,
crate::SampleStage::AtBoundary,
)
.await;
let crate::ProposeReply::Ticket(ticket) = reply else {
panic!("expected a ticket")
};
tickets.push(ticket);
}
assert!(head.items().is_empty());
pipeline.drain(&mut head, Resolution::Ack).await;
assert_eq!(
head.items().as_slice(),
[user("one"), user("two"), user("three")].as_slice()
);
let mut last = 0;
for (i, ticket) in tickets.into_iter().enumerate() {
assert_eq!(ticket.seq, i as u64 + 1);
let ack = ticket.ack.await.expect("resolved").expect("committed");
assert!(ack.offset > last, "offsets follow disk order");
last = ack.offset;
}
}
#[tokio::test]
async fn an_unacked_pipelined_entry_never_advances_the_projection() {
let dir = tempfile::tempdir().unwrap();
let (shared, mut log, mut head) = test_shared(dir.path());
let mut pipeline = Pipeline::default();
log.inject_fault(hotl_store::WriteFault::DropAckBeforeFsync);
let reply = commit_prepared(
&shared,
&mut log,
&mut head,
&mut pipeline,
crate::EntryProposal::of(vec![prepared(&shared, user("doomed"))]),
crate::AckMode::Pipelined,
crate::SampleStage::AtBoundary,
)
.await;
let crate::ProposeReply::Ticket(ticket) = reply else {
panic!("expected a ticket")
};
pipeline.drain(&mut head, Resolution::Ack).await;
assert!(
head.items().is_empty(),
"a crash may leave the log ahead of the projection, never the reverse"
);
assert_eq!(
ticket.ack.await.expect("resolved"),
Err(crate::CommitFailed::LogSealed)
);
}
#[tokio::test]
async fn seq_is_the_session_wide_commit_order_not_the_pipelined_subset() {
let dir = tempfile::tempdir().unwrap();
let (shared, mut log, mut head) = test_shared(dir.path());
let mut pipeline = Pipeline::default();
commit_prepared(
&shared,
&mut log,
&mut head,
&mut pipeline,
crate::EntryProposal::of(vec![prepared(&shared, user("sync"))]),
crate::AckMode::Sync,
crate::SampleStage::AtBoundary,
)
.await;
assert!(
shared
.append(
&mut log,
&mut pipeline,
&mut head,
EntryPayload::Rename {
name: "inline".into()
},
)
.await
);
let reply = commit_prepared(
&shared,
&mut log,
&mut head,
&mut pipeline,
crate::EntryProposal::of(vec![prepared(&shared, user("pipelined"))]),
crate::AckMode::Pipelined,
crate::SampleStage::AtBoundary,
)
.await;
let crate::ProposeReply::Ticket(ticket) = reply else {
panic!("expected a ticket")
};
assert_eq!(
ticket.seq, 3,
"the third commit of the session carries seq 3, not seq 1"
);
}
#[tokio::test]
async fn a_compaction_drains_the_pipeline_before_it_builds_and_mints() {
let dir = tempfile::tempdir().unwrap();
let (shared, mut log, mut head) = test_shared(dir.path());
let mut pipeline = Pipeline::default();
let mut tickets = Vec::new();
for text in ["first", "second"] {
let reply = commit_prepared(
&shared,
&mut log,
&mut head,
&mut pipeline,
crate::EntryProposal::of(vec![prepared(&shared, user(text))]),
crate::AckMode::Pipelined,
crate::SampleStage::AtBoundary,
)
.await;
let crate::ProposeReply::Ticket(ticket) = reply else {
panic!("expected a ticket")
};
tickets.push(ticket);
}
assert!(
head.items().is_empty(),
"two entries forwarded, none projected yet"
);
let spec = crate::SpecDigest {
prefix_end: 0,
kept_from: 2,
text: "folded".into(),
};
let degraded = compact(&shared, &mut log, &mut head, &mut pipeline, Some(spec))
.await
.expect("the fold must see the drained projection");
assert!(!degraded);
for ticket in tickets {
assert_eq!(
ticket.ack.await.expect("resolved"),
Err(crate::CommitFailed::Aborted),
"an aborted turn loses its claim on the log, never the bytes"
);
}
let entries: Vec<hotl_types::Entry> = std::fs::read_to_string(log.path())
.unwrap()
.lines()
.map(|l| serde_json::from_str(l).expect("entry"))
.collect();
let last = entries.len() - 1;
assert!(
matches!(entries[last].payload, EntryPayload::Compaction { .. }),
"the fold is minted last: {:?}",
entries[last].payload
);
assert_eq!(
entries[last].parent_id.as_deref(),
Some(entries[last - 1].id.as_str()),
"the aborting entry chains onto the drained leaf"
);
}
#[tokio::test]
#[cfg_attr(
not(debug_assertions),
ignore = "drives a debug_assert, which release builds compile out"
)]
#[should_panic(expected = "a held steer may only land between samples")]
async fn an_in_sample_commit_is_not_a_boundary_a_held_steer_may_land_at() {
let dir = tempfile::tempdir().unwrap();
let (shared, mut log, mut head) = test_shared(dir.path());
let mut held = vec!["hold me".to_string()];
release_steers(
&shared,
&mut log,
&mut head,
&mut Pipeline::default(),
&mut held,
super::Boundary::CommitSettled {
stage: crate::SampleStage::InSample,
},
)
.await;
}
#[tokio::test]
async fn a_commit_that_closed_its_sample_releases_the_steer_it_held() {
let dir = tempfile::tempdir().unwrap();
let (shared, mut log, mut head) = test_shared(dir.path());
let mut held = vec!["hold me".to_string()];
release_steers(
&shared,
&mut log,
&mut head,
&mut Pipeline::default(),
&mut held,
super::Boundary::CommitSettled {
stage: crate::SampleStage::AtBoundary,
},
)
.await;
assert!(held.is_empty(), "the steer must have landed");
assert_eq!(head.items().len(), 1, "…and reached the projection");
}
#[tokio::test]
async fn commit_prepared_commits_a_fresh_entry_and_updates_the_projection() {
let dir = tempfile::tempdir().unwrap();
let (shared, mut log, mut head) = test_shared(dir.path());
let epoch = shared.rules_epoch();
let payload = EntryPayload::Item { item: user("hi") };
let prepared = hotl_store::prepare_payload(&payload, &hotl_store::Masker::empty(), epoch)
.expect("prepare");
let entries = vec![crate::PreparedEntry::new(prepared, Some(user("hi")))];
let result = commit_prepared(
&shared,
&mut log,
&mut head,
&mut Pipeline::default(),
crate::EntryProposal::of(entries),
crate::AckMode::Sync,
crate::SampleStage::AtBoundary,
)
.await;
assert!(matches!(result, crate::ProposeReply::Committed));
assert_eq!(
head.items().len(),
1,
"a fresh proposal must reach the projection"
);
let replayed = hotl_store::replay(log.path()).expect("replay");
assert_eq!(
replayed.items.len(),
1,
"and the disk, chained after the header"
);
}
}