use std::collections::{BTreeMap, BTreeSet};
use std::path::Path;
use anyhow::Result;
use serde_json::{json, Value};
use super::{ledger_path, read_jsonl, str_field, Vehicle};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum RedeliverSource {
Compact,
Resume,
}
impl RedeliverSource {
pub(crate) fn as_str(self) -> &'static str {
match self {
RedeliverSource::Compact => "compact",
RedeliverSource::Resume => "resume",
}
}
pub(crate) fn parse(s: &str) -> Option<Self> {
match s {
"compact" => Some(RedeliverSource::Compact),
"resume" => Some(RedeliverSource::Resume),
_ => None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum LedgerLine {
Emit {
id: String,
event: String,
slot: u32,
part: u32,
parts: u32,
vehicle: Vehicle,
ts_utc: String,
hook_session: Option<String>,
hook_agent_id: Option<String>,
block_count: Option<u32>,
},
Held {
id: String,
reason: String,
ts_utc: String,
},
Expired {
id: String,
ts_utc: String,
},
Ack {
id: String,
ts_utc: String,
},
Redelivered {
id: String,
source: RedeliverSource,
ts_utc: String,
},
Refused {
id: Option<String>,
reason: String,
ts_utc: String,
},
}
impl LedgerLine {
pub(crate) fn id(&self) -> Option<&str> {
match self {
LedgerLine::Emit { id, .. }
| LedgerLine::Held { id, .. }
| LedgerLine::Expired { id, .. }
| LedgerLine::Ack { id, .. }
| LedgerLine::Redelivered { id, .. } => Some(id),
LedgerLine::Refused { id, .. } => id.as_deref(),
}
}
pub(crate) fn ts_utc(&self) -> &str {
match self {
LedgerLine::Emit { ts_utc, .. }
| LedgerLine::Held { ts_utc, .. }
| LedgerLine::Expired { ts_utc, .. }
| LedgerLine::Ack { ts_utc, .. }
| LedgerLine::Redelivered { ts_utc, .. }
| LedgerLine::Refused { ts_utc, .. } => ts_utc,
}
}
pub(crate) fn to_json(&self) -> Value {
match self {
LedgerLine::Emit {
id,
event,
slot,
part,
parts,
vehicle,
ts_utc,
hook_session,
hook_agent_id,
block_count,
} => json!({
"id": id,
"kind": "emit",
"event": event,
"slot": slot,
"part": part,
"parts": parts,
"vehicle": vehicle.as_str(),
"ts_utc": ts_utc,
"hook_session": hook_session,
"hook_agent_id": hook_agent_id,
"block_count": block_count,
}),
LedgerLine::Held { id, reason, ts_utc } => json!({
"id": id, "kind": "held", "reason": reason, "ts_utc": ts_utc,
}),
LedgerLine::Expired { id, ts_utc } => json!({
"id": id, "kind": "expired", "ts_utc": ts_utc,
}),
LedgerLine::Ack { id, ts_utc } => json!({
"id": id, "kind": "ack", "ts_utc": ts_utc,
}),
LedgerLine::Redelivered { id, source, ts_utc } => json!({
"id": id, "kind": "redelivered", "source": source.as_str(), "ts_utc": ts_utc,
}),
LedgerLine::Refused { id, reason, ts_utc } => json!({
"id": id, "kind": "refused", "reason": reason, "ts_utc": ts_utc,
}),
}
}
pub(crate) fn from_json(v: &Value) -> Option<Self> {
let ts_utc = str_field(v, "ts_utc")?;
Some(match str_field(v, "kind")?.as_str() {
"emit" => LedgerLine::Emit {
id: str_field(v, "id")?,
event: str_field(v, "event")?,
slot: u32_field(v, "slot")?,
part: u32_field(v, "part")?,
parts: u32_field(v, "parts")?,
vehicle: Vehicle::parse(str_field(v, "vehicle")?.as_str())?,
ts_utc,
hook_session: str_field(v, "hook_session"),
hook_agent_id: str_field(v, "hook_agent_id"),
block_count: u32_field(v, "block_count"),
},
"held" => LedgerLine::Held {
id: str_field(v, "id")?,
reason: str_field(v, "reason").unwrap_or_default(),
ts_utc,
},
"expired" => LedgerLine::Expired {
id: str_field(v, "id")?,
ts_utc,
},
"ack" => LedgerLine::Ack {
id: str_field(v, "id")?,
ts_utc,
},
"redelivered" => LedgerLine::Redelivered {
id: str_field(v, "id")?,
source: RedeliverSource::parse(str_field(v, "source")?.as_str())?,
ts_utc,
},
"refused" => LedgerLine::Refused {
id: str_field(v, "id"),
reason: str_field(v, "reason").unwrap_or_default(),
ts_utc,
},
_ => return None,
})
}
}
fn u32_field(v: &Value, key: &str) -> Option<u32> {
v.get(key)
.and_then(Value::as_u64)
.and_then(|n| u32::try_from(n).ok())
}
#[derive(Debug, Clone, Default)]
pub(crate) struct MessageState {
pub(crate) emitted_parts: BTreeSet<u32>,
pub(crate) parts_expected: Option<u32>,
pub(crate) exit2_emits: usize,
pub(crate) held_reasons: Vec<String>,
pub(crate) redelivered: Vec<RedeliverSource>,
pub(crate) refused_reasons: Vec<String>,
pub(crate) expired: bool,
pub(crate) acked: bool,
pub(crate) first_ts_utc: Option<String>,
pub(crate) last_ts_utc: Option<String>,
}
impl MessageState {
pub(crate) fn first_part_emitted(&self) -> bool {
self.emitted_parts.contains(&1)
}
pub(crate) fn has_exit2(&self) -> bool {
self.exit2_emits > 0
}
fn absorb(&mut self, line: &LedgerLine) {
let ts = line.ts_utc().to_string();
if self.first_ts_utc.is_none() {
self.first_ts_utc = Some(ts.clone());
}
self.last_ts_utc = Some(ts);
match line {
LedgerLine::Emit {
part,
parts,
vehicle,
..
} => {
self.emitted_parts.insert(*part);
self.parts_expected = Some(*parts);
if *vehicle == Vehicle::Exit2 {
self.exit2_emits += 1;
}
}
LedgerLine::Held { reason, .. } => self.held_reasons.push(reason.clone()),
LedgerLine::Expired { .. } => self.expired = true,
LedgerLine::Ack { .. } => self.acked = true,
LedgerLine::Redelivered { source, .. } => self.redelivered.push(*source),
LedgerLine::Refused { reason, .. } => self.refused_reasons.push(reason.clone()),
}
}
}
pub(crate) fn states(lines: &[LedgerLine]) -> BTreeMap<String, MessageState> {
let mut out: BTreeMap<String, MessageState> = BTreeMap::new();
for line in lines {
let Some(id) = line.id() else { continue };
out.entry(id.to_string()).or_default().absorb(line);
}
out
}
pub(crate) fn consecutive_exit2_blocks(lines: &[LedgerLine]) -> usize {
let mut count = 0usize;
for line in lines.iter().rev() {
match line {
LedgerLine::Emit { vehicle, .. } => {
if *vehicle == Vehicle::Exit2 {
count += 1;
} else {
break;
}
}
LedgerLine::Ack { .. } => break,
_ => {}
}
}
count
}
pub(crate) fn append_ledger(root: &Path, lane: &str, line: &LedgerLine) -> Result<()> {
let path = ledger_path(root, lane)?;
super::append_line(&path, &serde_json::to_string(&line.to_json())?)
}
pub(crate) fn read_ledger(root: &Path, lane: &str) -> Result<(Vec<LedgerLine>, usize)> {
let path = ledger_path(root, lane)?;
let (values, mut skipped) = read_jsonl(&path)?;
let mut out = Vec::with_capacity(values.len());
for v in &values {
match LedgerLine::from_json(v) {
Some(line) => out.push(line),
None => skipped += 1,
}
}
Ok((out, skipped))
}