use std::fs::{self, File, OpenOptions};
use std::io::{BufRead, BufReader, Lines, Read, Take, Write};
use std::path::{Path, PathBuf};
use serde::{Deserialize, Serialize};
use crate::{LedgerError, io_err};
const MANIFEST_FILE: &str = "manifest.json";
#[derive(Debug, Clone, Copy, Default)]
pub struct RolloverPolicy {
pub max_bytes: Option<u64>,
pub epoch_ms: Option<u64>,
}
impl RolloverPolicy {
pub fn max_bytes(max_bytes: u64) -> Self {
Self {
max_bytes: Some(max_bytes),
epoch_ms: None,
}
}
pub fn epoch_ms(epoch_ms: u64) -> Self {
Self {
max_bytes: None,
epoch_ms: Some(epoch_ms),
}
}
pub fn either(max_bytes: u64, epoch_ms: u64) -> Self {
Self {
max_bytes: Some(max_bytes),
epoch_ms: Some(epoch_ms),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct SegmentMeta {
pub file: String,
pub start_seq: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub end_seq: Option<u64>,
pub opened_ms: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct Manifest {
pub version: u32,
pub segments: Vec<SegmentMeta>,
}
impl Manifest {
pub(crate) fn active(&self) -> Option<&SegmentMeta> {
self.segments.iter().find(|s| s.end_seq.is_none())
}
pub(crate) fn active_mut(&mut self) -> Option<&mut SegmentMeta> {
self.segments.iter_mut().find(|s| s.end_seq.is_none())
}
}
pub(crate) fn segment_filename(index: u32) -> String {
format!("{index:08}.jsonl")
}
fn validate_segment_filename(file: &str) -> Result<u32, LedgerError> {
let bad = || LedgerError::Tamper {
seq: 0,
why: format!(
"manifest names a segment file with an invalid shape: {file:?} (expected exactly \
8 digits + \".jsonl\", e.g. 00000001.jsonl)"
),
};
let digits = file.strip_suffix(".jsonl").ok_or_else(bad)?;
if digits.len() != 8 || !digits.bytes().all(|b| b.is_ascii_digit()) {
return Err(bad());
}
digits.parse().map_err(|_| bad())
}
fn parse_index(filename: &str) -> Option<u32> {
validate_segment_filename(filename).ok()
}
fn validate_manifest_shape(manifest: &Manifest) -> Result<(), LedgerError> {
if manifest.segments.is_empty() {
return Err(LedgerError::Tamper {
seq: 0,
why: "manifest names zero segments — a valid segmented ledger always has at least \
one (segment::initialize never produces an empty list)"
.into(),
});
}
if manifest.segments[0].start_seq != 0 {
return Err(LedgerError::Tamper {
seq: 0,
why: format!(
"manifest's earliest segment {:?} starts at seq {}, not 0 — the true chain head \
is missing from this manifest",
manifest.segments[0].file, manifest.segments[0].start_seq
),
});
}
let active_count = manifest
.segments
.iter()
.filter(|s| s.end_seq.is_none())
.count();
if active_count > 1 {
return Err(LedgerError::Tamper {
seq: 0,
why: format!(
"manifest names {active_count} active (unsealed) segments — exactly one is valid"
),
});
}
if active_count == 1 && manifest.segments.last().is_none_or(|s| s.end_seq.is_some()) {
return Err(LedgerError::Tamper {
seq: 0,
why: "manifest's active segment must be the last entry".into(),
});
}
for pair in manifest.segments.windows(2) {
let (prev, next) = (&pair[0], &pair[1]);
if prev.end_seq != Some(next.start_seq) {
return Err(LedgerError::Tamper {
seq: 0,
why: format!(
"manifest segments {:?} (end_seq {:?}) and {:?} (start_seq {}) are not \
contiguous — segments must be listed in order with no gap, overlap, or \
reordering",
prev.file, prev.end_seq, next.file, next.start_seq
),
});
}
}
Ok(())
}
fn max_index(dir: &Path, manifest: &Manifest) -> Result<u32, LedgerError> {
let mut max = manifest
.segments
.iter()
.filter_map(|s| parse_index(&s.file))
.max()
.unwrap_or(0);
for entry in fs::read_dir(dir).map_err(|e| io_err(dir, e))? {
let entry = entry.map_err(|e| io_err(dir, e))?;
if let Some(name) = entry.file_name().to_str()
&& let Some(idx) = parse_index(name)
{
max = max.max(idx);
}
}
Ok(max)
}
pub(crate) fn load_manifest(dir: &Path) -> Result<Option<Manifest>, LedgerError> {
let path = dir.join(MANIFEST_FILE);
match fs::read(&path) {
Ok(bytes) => {
let manifest: Manifest =
serde_json::from_slice(&bytes).map_err(|e| LedgerError::Serde(e.to_string()))?;
for seg in &manifest.segments {
validate_segment_filename(&seg.file)?;
}
validate_manifest_shape(&manifest)?;
Ok(Some(manifest))
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
Err(e) => Err(io_err(&path, e)),
}
}
pub(crate) fn save_manifest(dir: &Path, manifest: &Manifest) -> Result<(), LedgerError> {
let path = dir.join(MANIFEST_FILE);
let tmp = dir.join("manifest.json.tmp");
let bytes =
serde_json::to_vec_pretty(manifest).map_err(|e| LedgerError::Serde(e.to_string()))?;
{
let mut f = File::create(&tmp).map_err(|e| io_err(&tmp, e))?;
f.write_all(&bytes).map_err(|e| io_err(&tmp, e))?;
f.sync_all().map_err(|e| io_err(&tmp, e))?;
}
fs::rename(&tmp, &path).map_err(|e| io_err(&path, e))?;
if let Ok(dirf) = File::open(dir) {
let _ = dirf.sync_all();
}
Ok(())
}
pub(crate) fn seal_file_permissions(path: &Path) {
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let _ = fs::set_permissions(path, fs::Permissions::from_mode(0o400));
}
#[cfg(not(unix))]
{
let _ = path;
}
}
pub(crate) fn unseal_file_permissions(path: &Path) {
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let _ = fs::set_permissions(path, fs::Permissions::from_mode(0o600));
}
#[cfg(not(unix))]
{
let _ = path;
}
}
pub(crate) fn segment_paths(dir: &Path) -> Result<Vec<PathBuf>, LedgerError> {
let manifest = load_manifest(dir)?.ok_or_else(|| LedgerError::Io {
path: dir.display().to_string(),
err: "segmented ledger directory has no manifest.json (not a valid segmented ledger — \
use Ledger::open_segmented to create one)"
.into(),
})?;
manifest
.segments
.iter()
.map(|s| {
let p = dir.join(&s.file);
if p.exists() {
Ok(p)
} else {
Err(LedgerError::Tamper {
seq: s.start_seq,
why: format!(
"segment {} is listed in the manifest but missing from disk (tail- or \
mid-truncation, or a corrupted deployment)",
s.file
),
})
}
})
.collect()
}
pub(crate) fn initialize(dir: &Path) -> Result<Manifest, LedgerError> {
let first = SegmentMeta {
file: segment_filename(1),
start_seq: 0,
end_seq: None,
opened_ms: 0,
};
let seg_path = dir.join(&first.file);
OpenOptions::new()
.create_new(true)
.write(true)
.open(&seg_path)
.map_err(|e| io_err(&seg_path, e))?;
let manifest = Manifest {
version: 1,
segments: vec![first],
};
save_manifest(dir, &manifest)?;
Ok(manifest)
}
pub(crate) fn roll_over(
dir: &Path,
manifest: &Manifest,
current_seq: u64,
next_ts_ms: u64,
) -> Result<(Manifest, PathBuf), LedgerError> {
let next_index = max_index(dir, manifest)?
.checked_add(1)
.ok_or_else(|| LedgerError::Io {
path: dir.display().to_string(),
err: "segment index exhausted (at u32::MAX)".into(),
})?;
if next_index > 99_999_999 {
return Err(LedgerError::Io {
path: dir.display().to_string(),
err: format!(
"segment index exhausted: the next segment index ({next_index}) no longer fits \
the 8-digit segment filename shape this ledger's segments use — this deployment \
has performed the maximum ~100,000,000 supported segment rollovers"
),
});
}
let new_filename = segment_filename(next_index);
let new_path = dir.join(&new_filename);
OpenOptions::new()
.create_new(true)
.write(true)
.open(&new_path)
.map_err(|e| io_err(&new_path, e))?;
let mut new_manifest = manifest.clone();
let old_active = new_manifest.active_mut().ok_or_else(|| LedgerError::Io {
path: dir.display().to_string(),
err: "segmented ledger manifest has no active segment to seal".into(),
})?;
old_active.end_seq = Some(current_seq);
let old_active_file = old_active.file.clone();
new_manifest.segments.push(SegmentMeta {
file: new_filename,
start_seq: current_seq,
end_seq: None,
opened_ms: next_ts_ms,
});
{
let old_path = dir.join(&old_active_file);
let f = OpenOptions::new()
.append(true)
.open(&old_path)
.map_err(|e| io_err(&old_path, e))?;
f.sync_all().map_err(|e| io_err(&old_path, e))?;
}
save_manifest(dir, &new_manifest)?; seal_file_permissions(&dir.join(&old_active_file));
Ok((new_manifest, new_path))
}
pub(crate) struct ChainedLines {
paths: std::vec::IntoIter<PathBuf>,
current: Option<(PathBuf, Lines<BufReader<Take<File>>>)>,
done: bool,
last_limit: u64,
}
impl Iterator for ChainedLines {
type Item = Result<String, LedgerError>;
fn next(&mut self) -> Option<Self::Item> {
if self.done {
return None;
}
loop {
if let Some((path, lines)) = self.current.as_mut() {
match lines.next() {
Some(Ok(line)) => return Some(Ok(line)),
Some(Err(e)) => {
self.done = true;
return Some(Err(io_err(path, e)));
}
None => {
self.current = None;
}
}
} else {
let next_path = self.paths.next()?;
let limit = if self.paths.len() == 0 {
self.last_limit
} else {
u64::MAX
};
match File::open(&next_path) {
Ok(f) => {
self.current = Some((next_path, BufReader::new(f.take(limit)).lines()));
}
Err(e) => {
self.done = true;
return Some(Err(io_err(&next_path, e)));
}
}
}
}
}
}
pub(crate) fn chained_lines(paths: Vec<PathBuf>) -> ChainedLines {
chained_lines_bounded(paths, u64::MAX)
}
pub(crate) fn chained_lines_bounded(paths: Vec<PathBuf>, last_limit: u64) -> ChainedLines {
ChainedLines {
paths: paths.into_iter(),
current: None,
done: false,
last_limit,
}
}