🦀 bions-rust vague 1 : 6 briques build-your-own-x — on ne les rebuild plus jamais
Principe RS-7 : « dès qu'on build un truc, plus personne n'a à le rebuild —
la seule chose à faire est l'optimisation. » (nexus/RepoVerse)
- bion-vc : horloges vectorielles + MvReg fork-visible (LA spec
xion-relativiste-v0 enfin codée — CRDT testé par permutations)
- bion-triplet : l'Adressage Génératif (gen_hash BLAKE3, coords, résidu ;
résidu vide quand déjà-su ; align décidable au bit)
- bion-tsoinlog: journal append-only rejouable (CRC32 maison, crash-recovery)
- bion-kv : magasin clé-valeur bitcask (compaction atomique, tombstones)
- bion-regex : moteur Thompson NFA linéaire (jamais exponentiel — Russ Cox)
- bion-git : mini-git content-addressed (SHA-1 maison + vecteurs officiels,
branches divergentes = le fork visible)
129 tests verts, clippy 0 warning, doc française = chaque bion est un cours.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
8
bion-tsoinlog/Cargo.toml
Normal file
8
bion-tsoinlog/Cargo.toml
Normal file
@@ -0,0 +1,8 @@
|
||||
[package]
|
||||
name = "bion-tsoinlog"
|
||||
version = "0.1.0"
|
||||
edition = "2021"
|
||||
description = "Journal d'événements append-only rejouable — la primitive de la machine à tsoins."
|
||||
license = "MIT"
|
||||
|
||||
[dependencies]
|
||||
91
bion-tsoinlog/README.md
Normal file
91
bion-tsoinlog/README.md
Normal file
@@ -0,0 +1,91 @@
|
||||
# bion-tsoinlog — le journal append-only rejouable
|
||||
|
||||
## Quoi
|
||||
|
||||
La **primitive de la machine à tsoins** : un journal d'événements sur disque où
|
||||
l'on ne fait qu'**appender** (jamais modifier, jamais effacer) et que l'on peut
|
||||
**rejouer** — en entier ou par tranche `[from, to)`. Chaque événement = un
|
||||
`topic` (UTF-8) + un `payload` (octets opaques), et reçoit un numéro de séquence
|
||||
`Seq` strictement croissant.
|
||||
|
||||
- **std-only**, zéro dépendance, zéro `unsafe`.
|
||||
- Format binaire v1 minuscule et documenté : `magic "TSOINLG1"` puis
|
||||
`len(u32 LE) · topic_len(u16 LE) · topic · payload · crc32(u32 LE)`.
|
||||
- **CRC-32 (IEEE) implémenté maison** (table générée à la compilation, vecteur
|
||||
canonique `"123456789" → 0xCBF43926` testé) — pas de lib, build-your-own.
|
||||
- **Récupération après crash** : à l'ouverture, la queue tronquée ou au CRC faux
|
||||
est détectée et retaillée ; l'historique valide survit toujours (testé en
|
||||
tronquant et en corrompant le fichier à la main).
|
||||
- `fsync` configurable (`Sync::Always` / `Sync::Never` + `sync()` manuel).
|
||||
- Index en mémoire (offset de chaque record) reconstruit au scan d'ouverture →
|
||||
`iter_from(seq)` et `replay` démarrent en seek direct, pas de re-scan.
|
||||
|
||||
## Pourquoi
|
||||
|
||||
C'est la structure au cœur de Kafka, des WAL de bases de données, de l'event
|
||||
sourcing, de git — et de la machine à tsoins : **enregistrer le réel dans
|
||||
l'ordre, pouvoir le revivre**. Le complément exact de `tsoin-codec`
|
||||
(`~/xerboxion-rt/crates/tsoin-codec`) : le codec transforme un contenu en
|
||||
coordonnée de Babel, le log mémorise durablement la *séquence* des événements.
|
||||
On appende volontiers des tsoins-de-fil comme payloads ; le log, lui, ne
|
||||
présuppose rien sur le contenu.
|
||||
|
||||
## Exemple
|
||||
|
||||
```rust
|
||||
use bion_tsoinlog::{TsoinLog, Seq, Sync};
|
||||
|
||||
let mut log = TsoinLog::open("journal.tsoinlog")?; // Sync::Never par défaut
|
||||
// ou : TsoinLog::open_with("journal.tsoinlog", Sync::Always)? // durable à chaque append
|
||||
|
||||
let s0 = log.append("capteur/temp", b"21.5")?; // → Seq(0)
|
||||
let s1 = log.append("bus/emit", b"{\"topic\":\"leds\"}")?; // → Seq(1)
|
||||
|
||||
// Tout relire :
|
||||
for rec in log.iter()? {
|
||||
let rec = rec?;
|
||||
println!("#{} [{}] {} octets", rec.seq.0, rec.topic, rec.payload.len());
|
||||
}
|
||||
|
||||
// Rejouer une tranche [1, 2) (from inclus, to exclu, comme un Range) :
|
||||
log.replay(Seq(1), Seq(2), |rec| { /* ré-appliquer l'événement */ })?;
|
||||
# std::io::Result::Ok(())
|
||||
```
|
||||
|
||||
## API publique (stable — rétrocompatibilité éternelle)
|
||||
|
||||
| Élément | Rôle |
|
||||
|---|---|
|
||||
| `TsoinLog::open(path)` / `open_with(path, Sync)` | ouvre/crée + scan + réparation |
|
||||
| `append(topic, payload) -> io::Result<Seq>` | appende un événement |
|
||||
| `iter()` / `iter_from(Seq)` | itère (instantané cohérent, seek O(1)) |
|
||||
| `replay(from, to, f) -> io::Result<u64>` | rejoue `[from, to)`, retourne le compte |
|
||||
| `sync()` / `set_sync(Sync)` | contrôle du `fsync` |
|
||||
| `len()` / `is_empty()` / `next_seq()` | état du journal |
|
||||
| `Seq(u64)` / `Record { seq, topic, payload }` / `Sync` | types de données |
|
||||
| `crc32(&[u8]) -> u32` / `MAGIC` / `MAX_TOPIC_LEN` | briques exposées |
|
||||
|
||||
## Build-your-own-x correspondants
|
||||
|
||||
Dans [build-your-own-x](https://github.com/codecrafters-io/build-your-own-x) :
|
||||
|
||||
- **Build your own Database** — le log est un *write-ahead log* (WAL) minimal ;
|
||||
- **Build your own Git** — un historique append-only adressé par position ;
|
||||
- le CRC-32 maison est la brique commune à zip/gzip/PNG/Ethernet (« build your
|
||||
own checksum » depuis la division polynomiale dans GF(2)).
|
||||
|
||||
## Comment l'optimiser (l'invitation au fork)
|
||||
|
||||
Le format v1 est volontairement le plus simple qui soit correct. Pistes, dans
|
||||
l'ordre de rentabilité, **sans jamais casser la lecture des fichiers v1** :
|
||||
|
||||
1. **`Sync::EveryN(n)` / group commit** — amortir le `fsync` sur n appends ;
|
||||
2. **index persistant** (fichier `.idx` side-car, régénérable) — ouverture O(1)
|
||||
sur les très gros journaux au lieu du scan complet ;
|
||||
3. **segments + compaction** — découper en fichiers de taille bornée, archiver
|
||||
ou fusionner les vieux segments (le chemin vers Kafka) ;
|
||||
4. **mmap en lecture** — itération zéro-copie ;
|
||||
5. **CRC vectorisé** (slicing-by-8, ou `crc32` matériel SSE4.2) — même résultat,
|
||||
~10× plus vite ;
|
||||
6. **compression des payloads** — brancher `tsoin-codec` : appender la
|
||||
coordonnée de Babel au lieu des octets bruts.
|
||||
721
bion-tsoinlog/src/lib.rs
Normal file
721
bion-tsoinlog/src/lib.rs
Normal file
@@ -0,0 +1,721 @@
|
||||
//! `bion-tsoinlog` — le **journal d'événements append-only rejouable**.
|
||||
//!
|
||||
//! C'est la primitive de la *machine à tsoins* : tout ce qui arrive est **appendé**
|
||||
//! (jamais modifié, jamais effacé), et le passé peut être **rejoué** à l'identique,
|
||||
//! en entier ou par tranche. Là où `tsoin-codec` (dans `xerboxion-rt`) sait *encoder*
|
||||
//! un contenu en coordonnée de Babel, ce bion-ci est le **LOG** : le fichier durable
|
||||
//! qui mémorise la séquence des événements bruts. Les deux se composent : on peut
|
||||
//! très bien appender des tsoins-de-fil comme payloads — mais le log, lui, ne
|
||||
//! présuppose RIEN sur le contenu (des octets opaques + un topic).
|
||||
//!
|
||||
//! # Le cours : pourquoi un log append-only ?
|
||||
//!
|
||||
//! Un log append-only est la structure de données la plus simple qui donne à la fois :
|
||||
//!
|
||||
//! 1. **La durabilité** — on n'écrit qu'à la fin du fichier ; un crash ne peut abîmer
|
||||
//! que le *dernier* enregistrement, jamais l'historique.
|
||||
//! 2. **L'ordre total** — chaque enregistrement reçoit un numéro de séquence
|
||||
//! ([`Seq`]) strictement croissant : le temps du journal.
|
||||
//! 3. **Le rejeu** — l'état de n'importe quel système peut être *reconstruit* en
|
||||
//! rejouant le log du début (ou d'un point connu) : c'est l'*event sourcing*,
|
||||
//! le principe de Kafka, des WAL de bases de données, de git… et de la machine
|
||||
//! à tsoins : enregistrer le réel, pouvoir le revivre.
|
||||
//!
|
||||
//! # Format binaire (v1) — simple et documenté
|
||||
//!
|
||||
//! Le fichier commence par un en-tête de 8 octets, puis une suite d'enregistrements :
|
||||
//!
|
||||
//! ```text
|
||||
//! fichier := magic(8 = "TSOINLG1") record*
|
||||
//! record := len(u32 LE) body crc32(u32 LE)
|
||||
//! body := topic_len(u16 LE) topic(UTF-8) payload(octets bruts)
|
||||
//! ```
|
||||
//!
|
||||
//! - `len` = taille du `body` en octets (donc `record` = 4 + len + 4 octets).
|
||||
//! - `crc32` = CRC-32 (IEEE 802.3, implémenté ici même — voir [`crc32`]) du `body`.
|
||||
//! - le numéro de séquence n'est **pas** stocké : il est *positionnel* (le i-ème
|
||||
//! enregistrement du fichier a `Seq(i)`), donc impossible à désynchroniser.
|
||||
//!
|
||||
//! # Récupération après crash
|
||||
//!
|
||||
//! À l'ouverture, le fichier est scanné : le premier enregistrement tronqué ou dont
|
||||
//! le CRC ne colle pas marque la **fin valide** du journal. Tout ce qui suit est
|
||||
//! ignoré et le fichier est retaillé à cette frontière — le dernier écrit partiel
|
||||
//! d'un crash disparaît proprement, l'historique intact reste. (Testé en tronquant
|
||||
//! et en corrompant à la main.)
|
||||
//!
|
||||
//! # Exemple
|
||||
//!
|
||||
//! ```
|
||||
//! use bion_tsoinlog::{TsoinLog, Seq};
|
||||
//! # let dir = std::env::temp_dir().join(format!("tsoinlog-doc-{}", std::process::id()));
|
||||
//! # std::fs::create_dir_all(&dir).unwrap();
|
||||
//! # let path = dir.join("journal.tsoinlog");
|
||||
//! # let _ = std::fs::remove_file(&path);
|
||||
//! let mut log = TsoinLog::open(&path).unwrap();
|
||||
//! let s0 = log.append("capteur/temp", b"21.5").unwrap();
|
||||
//! let s1 = log.append("capteur/temp", b"21.7").unwrap();
|
||||
//! assert_eq!((s0, s1), (Seq(0), Seq(1)));
|
||||
//!
|
||||
//! // Rejouer la tranche [0, 2) :
|
||||
//! let mut vus = Vec::new();
|
||||
//! log.replay(Seq(0), Seq(2), |rec| vus.push(rec.payload.clone())).unwrap();
|
||||
//! assert_eq!(vus, vec![b"21.5".to_vec(), b"21.7".to_vec()]);
|
||||
//! # std::fs::remove_file(&path).unwrap();
|
||||
//! ```
|
||||
|
||||
use std::fs::{File, OpenOptions};
|
||||
use std::io::{self, BufReader, Read, Seek, SeekFrom, Write};
|
||||
use std::path::{Path, PathBuf};
|
||||
|
||||
/// Les 8 octets magiques en tête de fichier : identifient le format + sa version.
|
||||
/// Si le format devait évoluer un jour, ce serait `TSOINLG2` — jamais une rupture
|
||||
/// silencieuse (rétrocompatibilité éternelle : un lecteur v2 lira toujours le v1).
|
||||
pub const MAGIC: [u8; 8] = *b"TSOINLG1";
|
||||
|
||||
/// Taille maximale du topic (il est préfixé par un `u16`).
|
||||
pub const MAX_TOPIC_LEN: usize = u16::MAX as usize;
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// CRC-32 (IEEE 802.3) — implémenté depuis les principes (build-your-own).
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Table des 256 restes précalculés du CRC-32, générée **à la compilation**.
|
||||
///
|
||||
/// Le CRC-32 est une division polynomiale dans GF(2) : le message est vu comme un
|
||||
/// grand polynôme à coefficients binaires, divisé par le polynôme générateur
|
||||
/// IEEE `0x04C11DB7`. Ici on utilise sa forme *réfléchie* `0xEDB88320` (bits
|
||||
/// inversés), ce qui permet de traiter les octets LSB-d'abord — la convention
|
||||
/// standard (zip, gzip, PNG, Ethernet). La table mémorise le reste de la division
|
||||
/// pour chacune des 256 valeurs d'octet : on avance alors octet par octet au lieu
|
||||
/// de bit par bit.
|
||||
const CRC_TABLE: [u32; 256] = {
|
||||
let mut table = [0u32; 256];
|
||||
let mut i = 0;
|
||||
while i < 256 {
|
||||
let mut c = i as u32;
|
||||
let mut k = 0;
|
||||
while k < 8 {
|
||||
// Si le bit sortant est 1, on « soustrait » (XOR) le générateur.
|
||||
c = if c & 1 != 0 {
|
||||
0xEDB8_8320 ^ (c >> 1)
|
||||
} else {
|
||||
c >> 1
|
||||
};
|
||||
k += 1;
|
||||
}
|
||||
table[i] = c;
|
||||
i += 1;
|
||||
}
|
||||
table
|
||||
};
|
||||
|
||||
/// CRC-32 (IEEE) d'un buffer. Pré/post-conditionnement standard : registre initial
|
||||
/// tout à 1 (`!0`), résultat inversé — c'est ce qui rend le CRC sensible aux zéros
|
||||
/// de tête et de queue. Vecteur de test canonique : `crc32(b"123456789") == 0xCBF43926`.
|
||||
pub fn crc32(data: &[u8]) -> u32 {
|
||||
let mut c = !0u32;
|
||||
for &b in data {
|
||||
c = CRC_TABLE[((c ^ b as u32) & 0xFF) as usize] ^ (c >> 8);
|
||||
}
|
||||
!c
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Types publics
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Numéro de séquence d'un enregistrement : le « temps » du journal.
|
||||
///
|
||||
/// Le premier enregistrement a `Seq(0)`, le suivant `Seq(1)`, etc. C'est un
|
||||
/// identifiant *positionnel* : il n'est pas stocké dans le fichier, il ne peut
|
||||
/// donc jamais être incohérent avec lui.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
|
||||
pub struct Seq(pub u64);
|
||||
|
||||
/// Un enregistrement relu depuis le journal : son numéro, son topic, ses octets.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct Record {
|
||||
/// Position dans le journal (0-indexée, strictement croissante).
|
||||
pub seq: Seq,
|
||||
/// Canal logique de l'événement (UTF-8, ≤ 65535 octets).
|
||||
pub topic: String,
|
||||
/// Contenu opaque : le log ne l'interprète jamais.
|
||||
pub payload: Vec<u8>,
|
||||
}
|
||||
|
||||
/// Politique de synchronisation disque après chaque `append`.
|
||||
///
|
||||
/// - [`Sync::Always`] : `fsync` après chaque écriture — durabilité maximale
|
||||
/// (l'enregistrement survit à une coupure de courant dès que `append` retourne),
|
||||
/// débit minimal.
|
||||
/// - [`Sync::Never`] : on laisse l'OS vider ses caches — débit maximal ; en cas de
|
||||
/// crash machine, les derniers enregistrements peuvent manquer, mais grâce à la
|
||||
/// récupération le journal reste *cohérent* (jamais corrompu, juste plus court).
|
||||
///
|
||||
/// C'est LE compromis classique des journaux (cf. `innodb_flush_log_at_trx_commit`,
|
||||
/// `fsync` de Redis AOF…). Par défaut : `Never` (on peut toujours appeler
|
||||
/// [`TsoinLog::sync`] aux moments importants).
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
|
||||
pub enum Sync {
|
||||
/// `fsync` à chaque `append`.
|
||||
Always,
|
||||
/// Jamais de `fsync` automatique (défaut).
|
||||
#[default]
|
||||
Never,
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Le journal
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Le journal append-only. Une instance = un fichier ouvert en écriture (append),
|
||||
/// plus un **index en mémoire** (offset de chaque enregistrement) reconstruit au
|
||||
/// scan d'ouverture — c'est lui qui rend `iter_from`/`replay` en accès direct
|
||||
/// (seek O(1)) au lieu d'un re-scan.
|
||||
#[derive(Debug)]
|
||||
pub struct TsoinLog {
|
||||
path: PathBuf,
|
||||
file: File,
|
||||
/// `index[i]` = offset du début de l'enregistrement `Seq(i)` dans le fichier.
|
||||
index: Vec<u64>,
|
||||
/// Fin valide du journal = offset où écrire le prochain enregistrement.
|
||||
end: u64,
|
||||
sync: Sync,
|
||||
}
|
||||
|
||||
impl TsoinLog {
|
||||
/// Ouvre (ou crée) le journal à `path`, avec la politique par défaut
|
||||
/// ([`Sync::Never`]). Scanne le fichier, reconstruit l'index, et **répare**
|
||||
/// une éventuelle fin tronquée/corrompue (voir la doc du module).
|
||||
pub fn open<P: AsRef<Path>>(path: P) -> io::Result<Self> {
|
||||
Self::open_with(path, Sync::default())
|
||||
}
|
||||
|
||||
/// Comme [`TsoinLog::open`] mais en choisissant la politique de `fsync`.
|
||||
pub fn open_with<P: AsRef<Path>>(path: P, sync: Sync) -> io::Result<Self> {
|
||||
let path = path.as_ref().to_path_buf();
|
||||
let mut file = OpenOptions::new()
|
||||
.read(true)
|
||||
.write(true)
|
||||
.create(true)
|
||||
.truncate(false) // append-only : on ne détruit JAMAIS l'existant
|
||||
.open(&path)?;
|
||||
|
||||
let file_len = file.metadata()?.len();
|
||||
if file_len == 0 {
|
||||
// Fichier neuf : on pose l'en-tête.
|
||||
file.write_all(&MAGIC)?;
|
||||
file.sync_all()?;
|
||||
return Ok(Self {
|
||||
path,
|
||||
file,
|
||||
index: Vec::new(),
|
||||
end: MAGIC.len() as u64,
|
||||
sync,
|
||||
});
|
||||
}
|
||||
if file_len < MAGIC.len() as u64 {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::InvalidData,
|
||||
"fichier trop court pour être un tsoinlog (en-tête absent)",
|
||||
));
|
||||
}
|
||||
let mut magic = [0u8; 8];
|
||||
file.seek(SeekFrom::Start(0))?;
|
||||
file.read_exact(&mut magic)?;
|
||||
if magic != MAGIC {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::InvalidData,
|
||||
"mauvais magic : ce fichier n'est pas un tsoinlog v1",
|
||||
));
|
||||
}
|
||||
|
||||
// Scan : on avance enregistrement par enregistrement tant que tout est
|
||||
// valide ; le premier accroc marque la fin réelle du journal.
|
||||
let mut reader = BufReader::new(&mut file);
|
||||
reader.seek(SeekFrom::Start(MAGIC.len() as u64))?;
|
||||
let mut index = Vec::new();
|
||||
let mut off = MAGIC.len() as u64;
|
||||
// Fin de boucle = fin propre OU queue tronquée/corrompue : on s'arrête là.
|
||||
while let Some((next_off, _body)) = read_record_at(&mut reader, off, file_len) {
|
||||
index.push(off);
|
||||
off = next_off;
|
||||
}
|
||||
drop(reader);
|
||||
|
||||
// Réparation : si des octets invalides traînent après la fin valide
|
||||
// (crash en pleine écriture), on retaille — le journal redevient sain.
|
||||
if off < file_len {
|
||||
file.set_len(off)?;
|
||||
file.sync_all()?;
|
||||
}
|
||||
Ok(Self {
|
||||
path,
|
||||
file,
|
||||
index,
|
||||
end: off,
|
||||
sync,
|
||||
})
|
||||
}
|
||||
|
||||
/// Change la politique de `fsync` (prend effet dès le prochain `append`).
|
||||
pub fn set_sync(&mut self, sync: Sync) {
|
||||
self.sync = sync;
|
||||
}
|
||||
|
||||
/// Appende un événement et retourne son numéro de séquence.
|
||||
///
|
||||
/// Erreurs : `InvalidInput` si `topic` dépasse [`MAX_TOPIC_LEN`] octets ou si
|
||||
/// `topic + payload` dépasse `u32::MAX - 2` octets ; sinon les erreurs d'E/S.
|
||||
pub fn append(&mut self, topic: &str, payload: &[u8]) -> io::Result<Seq> {
|
||||
let t = topic.as_bytes();
|
||||
if t.len() > MAX_TOPIC_LEN {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::InvalidInput,
|
||||
"topic > 65535 octets",
|
||||
));
|
||||
}
|
||||
let body_len = 2usize
|
||||
.checked_add(t.len())
|
||||
.and_then(|n| n.checked_add(payload.len()))
|
||||
.filter(|&n| n <= u32::MAX as usize)
|
||||
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "enregistrement > 4 Gio"))?;
|
||||
|
||||
// On assemble le record complet en mémoire puis UNE écriture : si le
|
||||
// processus meurt au milieu, on obtient au pire un suffixe partiel —
|
||||
// exactement le cas que la récupération d'ouverture sait effacer.
|
||||
let mut buf = Vec::with_capacity(4 + body_len + 4);
|
||||
buf.extend_from_slice(&(body_len as u32).to_le_bytes());
|
||||
buf.extend_from_slice(&(t.len() as u16).to_le_bytes());
|
||||
buf.extend_from_slice(t);
|
||||
buf.extend_from_slice(payload);
|
||||
let crc = crc32(&buf[4..]);
|
||||
buf.extend_from_slice(&crc.to_le_bytes());
|
||||
|
||||
self.file.seek(SeekFrom::Start(self.end))?;
|
||||
self.file.write_all(&buf)?;
|
||||
if self.sync == Sync::Always {
|
||||
self.file.sync_data()?;
|
||||
}
|
||||
let seq = Seq(self.index.len() as u64);
|
||||
self.index.push(self.end);
|
||||
self.end += buf.len() as u64;
|
||||
Ok(seq)
|
||||
}
|
||||
|
||||
/// Force un `fsync` maintenant (utile avec [`Sync::Never`] aux points clés).
|
||||
pub fn sync(&mut self) -> io::Result<()> {
|
||||
self.file.sync_data()
|
||||
}
|
||||
|
||||
/// Nombre d'enregistrements valides dans le journal.
|
||||
pub fn len(&self) -> u64 {
|
||||
self.index.len() as u64
|
||||
}
|
||||
|
||||
/// `true` si le journal ne contient encore aucun enregistrement.
|
||||
pub fn is_empty(&self) -> bool {
|
||||
self.index.is_empty()
|
||||
}
|
||||
|
||||
/// Le numéro que recevra le **prochain** `append`.
|
||||
pub fn next_seq(&self) -> Seq {
|
||||
Seq(self.index.len() as u64)
|
||||
}
|
||||
|
||||
/// Itère sur tous les enregistrements, du premier au dernier.
|
||||
pub fn iter(&self) -> io::Result<Iter> {
|
||||
self.iter_from(Seq(0))
|
||||
}
|
||||
|
||||
/// Itère à partir de `from` (inclus). Grâce à l'index, le départ est un
|
||||
/// `seek` direct — pas de re-scan du fichier. Si `from` est au-delà de la
|
||||
/// fin, l'itérateur est simplement vide.
|
||||
pub fn iter_from(&self, from: Seq) -> io::Result<Iter> {
|
||||
// Handle de lecture indépendant : on peut itérer sans gêner l'écriture.
|
||||
let file = File::open(&self.path)?;
|
||||
let mut reader = BufReader::new(file);
|
||||
let start = self.index.get(from.0 as usize).copied().unwrap_or(self.end);
|
||||
reader.seek(SeekFrom::Start(start))?;
|
||||
Ok(Iter {
|
||||
reader,
|
||||
offset: start,
|
||||
end: self.end,
|
||||
next_seq: from,
|
||||
})
|
||||
}
|
||||
|
||||
/// Rejoue la tranche `[from, to)` (from inclus, to exclu — comme un `Range`)
|
||||
/// en appelant `f` sur chaque enregistrement, dans l'ordre. Retourne le
|
||||
/// nombre d'enregistrements rejoués. `to` peut dépasser la fin : on s'arrête
|
||||
/// au dernier enregistrement existant (rejouer « jusqu'au bout » = passer
|
||||
/// `Seq(u64::MAX)` ou `log.next_seq()`).
|
||||
pub fn replay<F>(&self, from: Seq, to: Seq, mut f: F) -> io::Result<u64>
|
||||
where
|
||||
F: FnMut(&Record),
|
||||
{
|
||||
let mut n = 0u64;
|
||||
for rec in self.iter_from(from)? {
|
||||
let rec = rec?;
|
||||
if rec.seq >= to {
|
||||
break;
|
||||
}
|
||||
f(&rec);
|
||||
n += 1;
|
||||
}
|
||||
Ok(n)
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Lecture bas niveau + itérateur
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Tente de lire un enregistrement complet et valide commençant à `off`
|
||||
/// (le reader doit déjà y être positionné). Retourne `Some((offset_suivant, body))`
|
||||
/// si tout est bon, `None` si la fin du fichier arrive avant, si `len` déborde
|
||||
/// de la zone `limit`, ou si le CRC ne correspond pas — les trois visages d'une
|
||||
/// queue de fichier morte.
|
||||
fn read_record_at<R: Read>(reader: &mut R, off: u64, limit: u64) -> Option<(u64, Vec<u8>)> {
|
||||
// 4 (len) + 4 (crc) au minimum.
|
||||
if limit.saturating_sub(off) < 8 {
|
||||
return None;
|
||||
}
|
||||
let mut len4 = [0u8; 4];
|
||||
reader.read_exact(&mut len4).ok()?;
|
||||
let body_len = u32::from_le_bytes(len4) as u64;
|
||||
if body_len < 2 || off + 4 + body_len + 4 > limit {
|
||||
return None; // longueur invalide ou enregistrement tronqué
|
||||
}
|
||||
let mut body = vec![0u8; body_len as usize];
|
||||
reader.read_exact(&mut body).ok()?;
|
||||
let mut crc4 = [0u8; 4];
|
||||
reader.read_exact(&mut crc4).ok()?;
|
||||
if u32::from_le_bytes(crc4) != crc32(&body) {
|
||||
return None; // octets abîmés : ce record (et tout ce qui suit) est mort
|
||||
}
|
||||
// Le topic_len doit être cohérent avec la taille du body.
|
||||
let topic_len = u16::from_le_bytes([body[0], body[1]]) as u64;
|
||||
if 2 + topic_len > body_len {
|
||||
return None;
|
||||
}
|
||||
Some((off + 4 + body_len + 4, body))
|
||||
}
|
||||
|
||||
/// Décode un `body` validé en [`Record`].
|
||||
fn decode_body(seq: Seq, body: Vec<u8>) -> io::Result<Record> {
|
||||
let topic_len = u16::from_le_bytes([body[0], body[1]]) as usize;
|
||||
let topic = std::str::from_utf8(&body[2..2 + topic_len])
|
||||
.map_err(|_| io::Error::new(io::ErrorKind::InvalidData, "topic non UTF-8"))?
|
||||
.to_owned();
|
||||
let payload = body[2 + topic_len..].to_vec();
|
||||
Ok(Record {
|
||||
seq,
|
||||
topic,
|
||||
payload,
|
||||
})
|
||||
}
|
||||
|
||||
/// Itérateur sur les enregistrements du journal. Borné à la fin valide connue au
|
||||
/// moment de sa création : un `append` postérieur n'est pas visible par un
|
||||
/// itérateur déjà ouvert (instantané cohérent).
|
||||
pub struct Iter {
|
||||
reader: BufReader<File>,
|
||||
offset: u64,
|
||||
end: u64,
|
||||
next_seq: Seq,
|
||||
}
|
||||
|
||||
impl Iterator for Iter {
|
||||
type Item = io::Result<Record>;
|
||||
|
||||
fn next(&mut self) -> Option<Self::Item> {
|
||||
let (next_off, body) = read_record_at(&mut self.reader, self.offset, self.end)?;
|
||||
self.offset = next_off;
|
||||
let seq = self.next_seq;
|
||||
self.next_seq = Seq(seq.0 + 1);
|
||||
Some(decode_body(seq, body))
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Tests
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
|
||||
/// Chemin de fichier de test unique (répertoire temp de l'OS).
|
||||
fn tmp_path(tag: &str) -> PathBuf {
|
||||
static N: AtomicU64 = AtomicU64::new(0);
|
||||
let n = N.fetch_add(1, Ordering::Relaxed);
|
||||
let dir = std::env::temp_dir().join(format!("bion-tsoinlog-tests-{}", std::process::id()));
|
||||
std::fs::create_dir_all(&dir).unwrap();
|
||||
dir.join(format!("{tag}-{n}.tsoinlog"))
|
||||
}
|
||||
|
||||
// 1. Le vecteur canonique du CRC-32 IEEE : la preuve que notre implémentation
|
||||
// maison est LA bonne (interopérable avec zip/gzip/PNG).
|
||||
#[test]
|
||||
fn crc32_vecteur_canonique() {
|
||||
assert_eq!(crc32(b"123456789"), 0xCBF4_3926);
|
||||
assert_eq!(crc32(b""), 0);
|
||||
assert_ne!(crc32(b"a"), crc32(b"b"));
|
||||
}
|
||||
|
||||
// 2. Aller-retour élémentaire : ce qu'on appende est ce qu'on relit.
|
||||
#[test]
|
||||
fn append_read_roundtrip() {
|
||||
let p = tmp_path("roundtrip");
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
assert!(log.is_empty());
|
||||
let s = log
|
||||
.append("hello/monde", b"payload \x00\xff binaire")
|
||||
.unwrap();
|
||||
assert_eq!(s, Seq(0));
|
||||
let recs: Vec<Record> = log.iter().unwrap().map(|r| r.unwrap()).collect();
|
||||
assert_eq!(recs.len(), 1);
|
||||
assert_eq!(recs[0].seq, Seq(0));
|
||||
assert_eq!(recs[0].topic, "hello/monde");
|
||||
assert_eq!(recs[0].payload, b"payload \x00\xff binaire");
|
||||
}
|
||||
|
||||
// 3. Plusieurs topics, ordre et numéros de séquence préservés.
|
||||
#[test]
|
||||
fn topics_multiples_et_ordre() {
|
||||
let p = tmp_path("topics");
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
for i in 0..5u8 {
|
||||
let s = log.append(&format!("t/{i}"), &[i, i, i]).unwrap();
|
||||
assert_eq!(s, Seq(i as u64));
|
||||
}
|
||||
let recs: Vec<Record> = log.iter().unwrap().map(|r| r.unwrap()).collect();
|
||||
assert_eq!(recs.len(), 5);
|
||||
for (i, r) in recs.iter().enumerate() {
|
||||
assert_eq!(r.seq, Seq(i as u64));
|
||||
assert_eq!(r.topic, format!("t/{i}"));
|
||||
assert_eq!(r.payload, vec![i as u8; 3]);
|
||||
}
|
||||
}
|
||||
|
||||
// 4. Cas limites : topic vide, payload vide, gros payload — tout est légal.
|
||||
#[test]
|
||||
fn cas_limites_vides_et_gros() {
|
||||
let p = tmp_path("limites");
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
log.append("", b"").unwrap();
|
||||
log.append("juste-topic", b"").unwrap();
|
||||
let gros = vec![0xABu8; 1 << 20]; // 1 Mio
|
||||
log.append("", &gros).unwrap();
|
||||
let recs: Vec<Record> = log.iter().unwrap().map(|r| r.unwrap()).collect();
|
||||
assert_eq!(
|
||||
recs[0],
|
||||
Record {
|
||||
seq: Seq(0),
|
||||
topic: String::new(),
|
||||
payload: vec![]
|
||||
}
|
||||
);
|
||||
assert_eq!(recs[1].topic, "juste-topic");
|
||||
assert_eq!(recs[2].payload, gros);
|
||||
// Et un topic trop long est refusé proprement.
|
||||
let trop = "x".repeat(MAX_TOPIC_LEN + 1);
|
||||
assert_eq!(
|
||||
log.append(&trop, b"").unwrap_err().kind(),
|
||||
io::ErrorKind::InvalidInput
|
||||
);
|
||||
}
|
||||
|
||||
// 5. Persistance : on ferme, on rouvre, la séquence continue où elle était.
|
||||
#[test]
|
||||
fn reouverture_continue_la_sequence() {
|
||||
let p = tmp_path("reopen");
|
||||
{
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
log.append("a", b"1").unwrap();
|
||||
log.append("a", b"2").unwrap();
|
||||
}
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
assert_eq!(log.len(), 2);
|
||||
assert_eq!(log.next_seq(), Seq(2));
|
||||
assert_eq!(log.append("a", b"3").unwrap(), Seq(2));
|
||||
let recs: Vec<Record> = log.iter().unwrap().map(|r| r.unwrap()).collect();
|
||||
assert_eq!(
|
||||
recs.iter().map(|r| r.payload[0]).collect::<Vec<_>>(),
|
||||
b"123".to_vec()
|
||||
);
|
||||
}
|
||||
|
||||
// 6. CRASH-RECOVERY (troncature) : on coupe le fichier au milieu du dernier
|
||||
// enregistrement — à la réouverture il est ignoré, le reste est intact,
|
||||
// et on peut ré-appender par-dessus.
|
||||
#[test]
|
||||
fn recovery_fichier_tronque() {
|
||||
let p = tmp_path("tronque");
|
||||
{
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
log.append("ok", b"garde-moi").unwrap();
|
||||
log.append("ok", b"garde-moi aussi").unwrap();
|
||||
log.append("boom", b"je serai coupe en plein vol").unwrap();
|
||||
}
|
||||
// Simule le crash : on tronque 5 octets dans le dernier record.
|
||||
let len = std::fs::metadata(&p).unwrap().len();
|
||||
let f = OpenOptions::new().write(true).open(&p).unwrap();
|
||||
f.set_len(len - 5).unwrap();
|
||||
drop(f);
|
||||
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
assert_eq!(log.len(), 2, "le record tronqué doit être ignoré");
|
||||
// Le fichier a été retaillé à la frontière valide.
|
||||
let recs: Vec<Record> = log.iter().unwrap().map(|r| r.unwrap()).collect();
|
||||
assert_eq!(recs[1].payload, b"garde-moi aussi");
|
||||
// Et la vie continue : le prochain append prend Seq(2).
|
||||
assert_eq!(log.append("neuf", b"apres le crash").unwrap(), Seq(2));
|
||||
let log2 = TsoinLog::open(&p).unwrap();
|
||||
assert_eq!(log2.len(), 3);
|
||||
}
|
||||
|
||||
// 7. CRASH-RECOVERY (corruption) : un octet du dernier record est abîmé →
|
||||
// le CRC le détecte, le record est écarté, l'historique d'avant survit.
|
||||
#[test]
|
||||
fn recovery_crc_corrompu() {
|
||||
let p = tmp_path("corrompu");
|
||||
{
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
log.append("sain", b"aaaa").unwrap();
|
||||
log.append("abime", b"bbbb").unwrap();
|
||||
}
|
||||
// Flip d'un octet dans le payload du 2e record (l'avant-dernier octet
|
||||
// avant le CRC final).
|
||||
let mut bytes = std::fs::read(&p).unwrap();
|
||||
let n = bytes.len();
|
||||
bytes[n - 6] ^= 0xFF;
|
||||
std::fs::write(&p, &bytes).unwrap();
|
||||
|
||||
let log = TsoinLog::open(&p).unwrap();
|
||||
assert_eq!(log.len(), 1, "le record au CRC faux doit être écarté");
|
||||
let recs: Vec<Record> = log.iter().unwrap().map(|r| r.unwrap()).collect();
|
||||
assert_eq!(recs[0].topic, "sain");
|
||||
}
|
||||
|
||||
// 8. Un fichier qui n'est pas un tsoinlog est refusé (pas de lecture hasardeuse).
|
||||
#[test]
|
||||
fn mauvais_magic_refuse() {
|
||||
let p = tmp_path("magic");
|
||||
std::fs::write(&p, b"PASUNLOGDUTOUT!!").unwrap();
|
||||
let err = TsoinLog::open(&p).unwrap_err();
|
||||
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
|
||||
}
|
||||
|
||||
// 9. iter_from : départ au milieu (seek direct via l'index), et au-delà de
|
||||
// la fin → itérateur vide, pas d'erreur.
|
||||
#[test]
|
||||
fn iter_from_milieu_et_apres_fin() {
|
||||
let p = tmp_path("iterfrom");
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
for i in 0..10u64 {
|
||||
log.append("n", &i.to_le_bytes()).unwrap();
|
||||
}
|
||||
let recs: Vec<Record> = log.iter_from(Seq(7)).unwrap().map(|r| r.unwrap()).collect();
|
||||
assert_eq!(recs.len(), 3);
|
||||
assert_eq!(recs[0].seq, Seq(7));
|
||||
assert_eq!(recs[0].payload, 7u64.to_le_bytes());
|
||||
assert_eq!(log.iter_from(Seq(10)).unwrap().count(), 0);
|
||||
assert_eq!(log.iter_from(Seq(9999)).unwrap().count(), 0);
|
||||
}
|
||||
|
||||
// 10. Replay partiel : la tranche [from, to) exactement, dans l'ordre.
|
||||
#[test]
|
||||
fn replay_partiel() {
|
||||
let p = tmp_path("replay");
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
for i in 0..10u8 {
|
||||
log.append("ev", &[i]).unwrap();
|
||||
}
|
||||
let mut vus = Vec::new();
|
||||
let n = log
|
||||
.replay(Seq(3), Seq(7), |r| vus.push(r.payload[0]))
|
||||
.unwrap();
|
||||
assert_eq!(n, 4);
|
||||
assert_eq!(vus, vec![3, 4, 5, 6]);
|
||||
// to au-delà de la fin : on rejoue jusqu'au bout sans erreur.
|
||||
let n = log.replay(Seq(8), Seq(u64::MAX), |_| {}).unwrap();
|
||||
assert_eq!(n, 2);
|
||||
// tranche vide.
|
||||
let n = log
|
||||
.replay(Seq(5), Seq(5), |_| panic!("ne doit pas être appelé"))
|
||||
.unwrap();
|
||||
assert_eq!(n, 0);
|
||||
}
|
||||
|
||||
// 11. GROS VOLUME : 10 000 records, roundtrip intégral + index après réouverture.
|
||||
#[test]
|
||||
fn volume_10k_records() {
|
||||
let p = tmp_path("volume");
|
||||
{
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
for i in 0..10_000u64 {
|
||||
let payload = [i.to_le_bytes().as_slice(), &[(i % 251) as u8; 17]].concat();
|
||||
let s = log.append(&format!("vol/{}", i % 7), &payload).unwrap();
|
||||
assert_eq!(s, Seq(i));
|
||||
}
|
||||
}
|
||||
// Réouverture : le scan reconstruit l'index sur les 10k records.
|
||||
let log = TsoinLog::open(&p).unwrap();
|
||||
assert_eq!(log.len(), 10_000);
|
||||
let mut count = 0u64;
|
||||
for rec in log.iter().unwrap() {
|
||||
let rec = rec.unwrap();
|
||||
assert_eq!(rec.seq, Seq(count));
|
||||
let i = u64::from_le_bytes(rec.payload[..8].try_into().unwrap());
|
||||
assert_eq!(i, count);
|
||||
assert_eq!(rec.topic, format!("vol/{}", count % 7));
|
||||
count += 1;
|
||||
}
|
||||
assert_eq!(count, 10_000);
|
||||
// Accès direct profond via l'index.
|
||||
let r = log.iter_from(Seq(9_999)).unwrap().next().unwrap().unwrap();
|
||||
assert_eq!(
|
||||
u64::from_le_bytes(r.payload[..8].try_into().unwrap()),
|
||||
9_999
|
||||
);
|
||||
}
|
||||
|
||||
// 12. fsync configurable : les deux politiques écrivent des journaux identiques
|
||||
// (la durabilité change, pas le format) ; sync() manuel disponible.
|
||||
#[test]
|
||||
fn politiques_de_sync() {
|
||||
let p1 = tmp_path("sync-always");
|
||||
let p2 = tmp_path("sync-never");
|
||||
let mut a = TsoinLog::open_with(&p1, Sync::Always).unwrap();
|
||||
let mut b = TsoinLog::open_with(&p2, Sync::Never).unwrap();
|
||||
for i in 0..20u8 {
|
||||
a.append("s", &[i]).unwrap();
|
||||
b.append("s", &[i]).unwrap();
|
||||
}
|
||||
b.sync().unwrap();
|
||||
b.set_sync(Sync::Always);
|
||||
b.append("s", &[99]).unwrap(); // record de 4 + (2+1+1) + 4 = 12 octets
|
||||
let a_bytes = std::fs::read(&p1).unwrap();
|
||||
let b_bytes = std::fs::read(&p2).unwrap();
|
||||
assert_eq!(a_bytes.len() + 12, b_bytes.len());
|
||||
// Les 20 premiers records sont octet-pour-octet identiques.
|
||||
assert_eq!(a_bytes[..], b_bytes[..a_bytes.len()]);
|
||||
}
|
||||
|
||||
// 13. Un itérateur ouvert est un instantané : les appends postérieurs ne
|
||||
// s'y invitent pas (cohérence de lecture).
|
||||
#[test]
|
||||
fn iterateur_est_un_instantane() {
|
||||
let p = tmp_path("snapshot");
|
||||
let mut log = TsoinLog::open(&p).unwrap();
|
||||
log.append("x", b"1").unwrap();
|
||||
let it = log.iter().unwrap();
|
||||
log.append("x", b"2").unwrap();
|
||||
assert_eq!(it.count(), 1);
|
||||
assert_eq!(log.iter().unwrap().count(), 2);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user