use std::path::Path;
use anyhow::Result;
use serde_json::Value;
use super::{
is_expired, is_queue_event, ledger_path, read_inbox, read_message, render, states, HookInput,
InboxLine, LedgerLine, MessageState, Mode, RedeliverSource, SlotChain, CHUNK_BUDGET,
};
const BASELINE_FILE: &str = "ledger.base";
pub(crate) const HELD_SOURCE_MISSING: &str = "source-missing";
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct Chunk {
pub(crate) id: String,
pub(crate) part: u32,
pub(crate) parts: u32,
pub(crate) mode: Mode,
pub(crate) text: String,
pub(crate) prior_exit2: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum Bookkeeping {
Expired(String),
Redelivered(String, RedeliverSource),
SourceMissing(String),
}
#[derive(Debug, Clone, Default)]
pub(crate) struct Plan {
pub(crate) chunks: Vec<Chunk>,
pub(crate) bookkeeping: Vec<Bookkeeping>,
}
pub(crate) fn fold_baseline(chain: &SlotChain, ledger_len: u64) -> u64 {
let path = chain.dir.join(BASELINE_FILE);
if chain.slot <= 1 {
let _ = std::fs::write(&path, ledger_len.to_string());
return ledger_len;
}
std::fs::read_to_string(&path)
.ok()
.and_then(|s| s.trim().parse::<u64>().ok())
.map_or(ledger_len, |n| n.min(ledger_len))
}
pub(crate) fn ledger_len(root: &Path, lane: &str) -> Result<u64> {
let path = ledger_path(root, lane)?;
Ok(std::fs::metadata(&path).map_or(0, |m| m.len()))
}
pub(crate) fn read_ledger_prefix(
root: &Path,
lane: &str,
baseline: u64,
) -> Result<(Vec<LedgerLine>, usize)> {
let path = ledger_path(root, lane)?;
let Ok(raw) = std::fs::read(&path) else {
return Ok((Vec::new(), 0));
};
let end = usize::try_from(baseline)
.unwrap_or(usize::MAX)
.min(raw.len());
let text = String::from_utf8_lossy(&raw[..end]).into_owned();
let complete = text.ends_with('\n');
let lines: Vec<&str> = text.lines().collect();
let mut out = Vec::with_capacity(lines.len());
let mut skipped = 0usize;
for (i, line) in lines.iter().enumerate() {
if line.trim().is_empty() || (!complete && i + 1 == lines.len()) {
continue;
}
match serde_json::from_str::<Value>(line)
.ok()
.as_ref()
.and_then(LedgerLine::from_json)
{
Some(parsed) => out.push(parsed),
None => skipped += 1,
}
}
Ok((out, skipped))
}
pub(crate) fn plan(
root: &Path,
lane: &str,
hook: &HookInput,
ledger: &[LedgerLine],
now_utc: &str,
) -> Result<Plan> {
let (inbox, _) = read_inbox(root, lane)?;
let folded = states(ledger);
let mut out = Plan::default();
let mut seen: Vec<&str> = Vec::new();
for line in &inbox {
if seen.contains(&line.id.as_str()) {
continue;
}
seen.push(line.id.as_str());
let state = folded.get(&line.id);
match verdict(line, state, hook, ledger, now_utc) {
Verdict::Skip => {}
Verdict::Expire => out.bookkeeping.push(Bookkeeping::Expired(line.id.clone())),
Verdict::Send { redeliver } => {
if let Some(source) = redeliver {
out.bookkeeping
.push(Bookkeeping::Redelivered(line.id.clone(), source));
}
append_chunks(root, line, state, &mut out)?;
}
}
}
Ok(out)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Verdict {
Skip,
Expire,
Send {
redeliver: Option<RedeliverSource>,
},
}
fn verdict(
line: &InboxLine,
state: Option<&MessageState>,
hook: &HookInput,
ledger: &[LedgerLine],
now_utc: &str,
) -> Verdict {
if state.is_some_and(|s| s.acked) {
return Verdict::Skip;
}
if line
.expires_utc
.as_deref()
.is_some_and(|e| is_expired(e, now_utc))
{
return if state.is_some_and(|s| s.expired) {
Verdict::Skip
} else {
Verdict::Expire
};
}
if !mode_eligible(line.mode, hook) {
return Verdict::Skip;
}
match state {
None => Verdict::Send { redeliver: None },
Some(s) if !s.first_part_emitted() => Verdict::Send { redeliver: None },
Some(_)
if hook.session_start_source("compact")
&& !redelivered_since_emit(ledger, &line.id) =>
{
Verdict::Send {
redeliver: Some(RedeliverSource::Compact),
}
}
Some(_) => Verdict::Skip,
}
}
fn redelivered_since_emit(ledger: &[LedgerLine], id: &str) -> bool {
for line in ledger.iter().rev() {
if line.id() != Some(id) {
continue;
}
match line {
LedgerLine::Redelivered { .. } => return true,
LedgerLine::Emit { .. } => return false,
_ => {}
}
}
false
}
fn mode_eligible(mode: Mode, hook: &HookInput) -> bool {
match mode {
Mode::Steer => true,
Mode::Queue => is_queue_event(&hook.hook_event_name, hook.source.as_deref()),
}
}
fn append_chunks(
root: &Path,
line: &InboxLine,
state: Option<&MessageState>,
out: &mut Plan,
) -> Result<()> {
let Some(msg) = read_message(root, &line.id)? else {
let already =
state.is_some_and(|s| s.held_reasons.iter().any(|r| r == HELD_SOURCE_MISSING));
if !already {
out.bookkeeping
.push(Bookkeeping::SourceMissing(line.id.clone()));
}
return Ok(());
};
let rendered = render(&msg, CHUNK_BUDGET)?;
let parts = u32::try_from(rendered.len()).unwrap_or(u32::MAX);
let prior_exit2 = state.is_some_and(MessageState::has_exit2);
for (i, text) in rendered.into_iter().enumerate() {
out.chunks.push(Chunk {
id: msg.id.clone(),
part: u32::try_from(i + 1).unwrap_or(u32::MAX),
parts,
mode: line.mode,
text,
prior_exit2,
});
}
Ok(())
}