Rust para sistemas distribuídos#
O delonix é um motor de um nó — mas foi desenhado com os padrões de um sistema distribuído, porque um nó já é um sistema distribuído: vários processos (CLI, supervisor, holder de rede, servidor CRI) a mexer no mesmo estado, a falhar a meio, a serem reiniciados por um systemd ou por um reboot. Tudo o que se segue escala de um nó para uma frota.
Os problemas são sempre os mesmos, e este capítulo dá uma resposta a cada:
| Problema | Padrão | Onde |
|---|---|---|
| «O que é que devo mudar?» | Reconciliação (estado desejado vs actual) | secção 1 |
| «E se aplicar duas vezes?» | Idempotência | secção 2 |
| «A rede falhou a meio» | Retentativas com backoff + retoma | secções 3-4 |
| «Dois processos escreveram ao mesmo tempo» | Escrita atómica + flock |
secção 5 |
| «Como sei o que correu mal sem ler a mensagem?» | Erros com classe | secção 6 |
Todos os exemplos vivem em examples/src/ch06_distribuido.rs, com testes.
1. Reconciliação: uma função pura#
Um sistema declarativo aceita «isto é o que eu quero» e descobre sozinho o que fazer. O coração é uma função que compara o desejado com o real e devolve as diferenças. A decisão de design que mais importa: essa função é pura.
//! **The declarative reconciler** — decides what has to change for the machine
//! to match the manifest, and decides it WITHOUT touching the machine.
//!
//! Until this module existed, `stack apply` only ever created: a resource that
//! already existed printed `already exists, nothing to do` and the command
//! returned 0. Changing the image in the manifest and re-applying did nothing
//! and reported success — the declarative twin of the dishonest reporting the
//! v0.37.0 audit removed from the imperative CLI, and worse here, because the
//! user changed the file on purpose and that intent was thrown away.
//!
//! # Why this is a pure function
//!
//! [`plan`] takes an already-read snapshot of both sides and returns a list of
//! [`Change`]. It never opens a store, never runs a command, never needs a
//! namespace or a privilege. That is what makes the interesting cases — a field
//! removed from the manifest, a resource owned by another stack, a change that
//! cannot be applied hot — testable as plain data, in milliseconds, with no
//! host state. Same discipline as `conditions::conditions_for`,
//! `vmbridge::bridge_plan` and `vm::resolve_vm_defaults`.
//!
//! # Field names are the ones the user typed
//!
//! [`Desired::fields`] and [`Actual::fields`] are keyed by the **manifest**
//! field name (`image`, `ports`, `restartPolicy`), never by the internal record
//! name (`memory_max`, `restart_policy`). The plan is read by whoever wrote the
//! YAML; naming a field they cannot find in their own file makes the diff
//! useless. Each Kind's module owns the translation.
//!
//! # Three-way, not two-way
//!
//! Comparing desired against actual cannot distinguish «the user deleted this
//! field from the manifest» from «this field was never ours». A two-way diff
//! has to pick one and is wrong in half the cases: either it keeps reverting
//! settings a human made with `container update`, or it never honours a removal.
//!
//! So the last spec we applied is kept on the resource itself (the
//! `delonix.io/last-applied` annotation — the mechanism `kubectl` uses, in the
//! place `kubectl` puts it) and the rule becomes:
//!
//! | in manifest | on machine | in last-applied | verdict |
//! |---|---|---|---|
//! | yes | differs | — | converge to the manifest |
//! | no | present | yes | **we** set it and it is gone from the file → revert |
//! | no | present | no | never ours → leave it alone |
//!
//! A resource with no `last-applied` at all (adopted, or created by hand) falls
//! into the last row for every field: the first apply only ever adds.Pura quer dizer: recebe os dois lados já lidos, devolve uma lista de mudanças, nunca abre um ficheiro nem corre um comando. É isso que torna os casos difíceis testáveis como dados, em microssegundos, sem estado do host. Uma versão minimal, com o mesmo raciocínio:
pub type Fields = BTreeMap<String, String>;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Action {
Create,
/// Campos a alterar, chave → novo valor.
Update(Fields),
/// Removido do manifesto e foi criado por nós: apagar.
Delete,
NoOp,
}
/// - `desired`: o manifesto de agora.
/// - `actual`: o que existe na máquina (`None` = não existe).
/// - `last_applied`: o que aplicámos da última vez (guardado no próprio recurso).
pub fn plan(desired: &Fields, actual: Option<&Fields>, last_applied: Option<&Fields>) -> Action {
let Some(actual) = actual else {
return Action::Create;
};
let mut changes = Fields::new();
for (k, want) in desired {
if actual.get(k) != Some(want) {
changes.insert(k.clone(), want.clone());
}
}
// Campo que já foi NOSSO (estava em last_applied) e saiu do manifesto → reverter.
// Campo que nunca foi nosso (alguém pôs à mão) → não tocar. É isto que o 3.º lado distingue.
if let Some(last) = last_applied {
for k in last.keys().filter(|k| !desired.contains_key(*k)) {
if actual.contains_key(k) {
changes.insert(k.clone(), String::new());
}
}
}
if changes.is_empty() { Action::NoOp } else { Action::Update(changes) }
}
/// Recurso removido do manifesto: só se apaga o que é NOSSO (tem dono = esta stack).
pub fn should_prune(owner: Option<&str>, stack: &str) -> bool {
owner == Some(stack)
}Porque 3 vias, e não 2#
Comparar só desejado contra actual não distingue duas situações opostas:
- «apaguei este campo do manifesto» → reverter;
- «alguém pôs este campo à mão» → não tocar.
Uma comparação de 2 vias tem de escolher uma — e erra na outra metade dos casos: ou reverte tudo o que um humano fez com container update, ou nunca honra uma remoção. O terceiro lado é o último estado que aplicámos, guardado no próprio recurso (a mesma anotação last-applied do kubectl). O teste mostra os três casos de uma vez:
let desired = f(&[("image", "nginx:2")]);
let actual = f(&[("image", "nginx:1"), ("memory", "64M"), ("debug", "on")]);
let last = f(&[("image", "nginx:1"), ("memory", "64M")]);
// image → mudou no manifesto → Update
// memory → era nosso e saiu → reverter (valor vazio)
// debug → nunca foi nosso → NÃO tocar
E o motor real — plan() do delonix-stack, com a mesma estrutura mais o que uma frota exige (posse por stack, campos «frios» que forçam recriação, adopção de recursos criados à mão):
/// Decides what has to change. Pure: both sides come in already read.
///
/// `stack` is the name that owns this apply; it is what separates «mine, gone
/// from the file, so remove it» from «someone else's, never touch it».
pub fn plan(desired: &[Desired], actual: &[Actual], stack: &str) -> Vec<Change> {
let mut out = Vec::new();
for d in desired {
let found = actual.iter().find(|a| a.kind == d.kind && a.name == d.name);
let Some(a) = found else {
out.push(Change::new(&d.kind, &d.name, Action::Create));
continue;
};
// Owned by another stack: refuse before computing anything. Two stacks
// converging the same resource would flap it between two shapes on
// every apply, and neither owner would understand why.
if let Some(owner) = &a.owner {
if owner != stack {
let mut c = Change::new(&d.kind, &d.name, Action::Conflict)
.with_reason(format!("owned by stack '{owner}'"));
c.owner = Some(owner.clone());
out.push(c);
continue;
}
}
if !d.converges {
out.push(
Change::new(&d.kind, &d.name, Action::NotConverged)
.with_reason("this Kind is ensure-present in this version"),
);
continue;
}
let diffs = diff_fields(&d.kind, d, a);
// A non-ownable resource is never «unowned and therefore to adopt»:
// having no owner is its normal state, not something to fix.
let unmanaged = d.ownable && a.owner.is_none();
let action = if let Some(cold) = diffs.iter().find(|x| !x.hot) {
let _ = cold;
Action::Replace
} else if !diffs.is_empty() {
Action::Update
} else if unmanaged {
Action::Adopt
} else {
Action::NoOp
};
// … (excerto)A posse é o que separa um apply seguro de um perigoso
Repara em if owner != stack { … Conflict }: dois stacks a convergir o mesmo recurso fariam-no oscilar entre duas formas a cada apply. E o --prune só apaga o que tem a nossa etiqueta — um recurso criado à mão é invisível, «porque um apply nunca deve considerar apagar o que não criou».
Fail-closed na recriação
Quando a diferença está num campo que só se muda recriando o recurso (-/+), o motor recusa sem --replace <Kind>/<nome>, e recusa antes da primeira criação: um apply que falhasse a meio deixaria a stack meio convergida e com erro.
2. Idempotência: a prova é «zero escritas»#
Aplicar N vezes tem de dar o mesmo estado que aplicar uma. Não basta o resultado ser igual — a segunda aplicação não deve escrever nada:
#[derive(Debug, Default)]
pub struct Machine {
pub resources: BTreeMap<String, Fields>,
pub writes: usize, // quantas vezes tocámos realmente na máquina
}
/// `apply` converge. Chamá-lo N vezes deixa o mesmo estado e — a prova — **zero escritas**
/// a partir da segunda.
pub fn apply(m: &mut Machine, name: &str, desired: &Fields) {
match plan(desired, m.resources.get(name), None) {
Action::Create => {
m.resources.insert(name.to_owned(), desired.clone());
m.writes += 1;
}
Action::Update(ch) => {
if let Some(r) = m.resources.get_mut(name) {
r.extend(ch);
}
m.writes += 1;
}
Action::Delete | Action::NoOp => {}
}
}O teste apply_is_idempotent afirma writes == 1 depois de três apply. Repara na pegadinha que o motor pagou: um apply que «não faz nada» também deixa o PID de um container intacto — por isso o cenário de caos do delonix não verifica só que o PID não mudou, verifica que o registo mudou (memory_max) e que o stack plan seguinte não tem nada a propor. Um teste que passa com o código apagado não prova nada.
3. Retentativas: backoff exponencial, com tecto#
/// 1s, 2s, 4s, 8s … até ao tecto. Sem tecto, uma falha longa dorme horas; sem *jitter*
/// (não incluído aqui por ser determinístico nos testes) uma frota inteira retenta em uníssono.
pub fn backoff(attempt: u32, base: Duration, cap: Duration) -> Duration {
let factor = 1u32.checked_shl(attempt).unwrap_or(u32::MAX);
base.saturating_mul(factor).min(cap)
}
/// Só se retenta o que faz sentido retentar: um 404 nunca melhora à 5.ª tentativa.
pub fn is_retryable(http_status: u16) -> bool {
matches!(http_status, 408 | 429 | 500 | 502 | 503 | 504)
}Três decisões escondidas em duas funções pequenas:
- Tecto (
cap): sem ele, uma falha longa dorme horas. checked_shl+saturating_mul:backoff(500, …)não dá pânico nem volta a zero — o teste prova-o.- Só se retenta o que faz sentido: um 404 não melhora à quinta tentativa; um 503 sim.
Em produção acrescenta-se jitter (aleatoriedade) para uma frota não retentar em uníssono — aqui fica de fora para o teste ser determinístico.
4. Retomar um download: o servidor responde a outra pergunta?#
O bug que deu origem a isto é real e ensina muito. delonix vm pull de 276 MiB morria aos 8 minutos numa ligação de 416 KB/s — e a tentativa seguinte recomeçava do byte zero. Abaixo de ~600 KB/s a imagem nunca acabava. A cura é Range: bytes=<n>- a partir do que já está em memória. O que torna isto subtil:
/// Valida um `Content-Range: bytes <inicio>-<fim>/<total>` contra o offset que pedimos.
/// Um 206 noutro offset «responde a outra pergunta»: colar duplicaria o prefixo e a corrupção
/// só apareceria no digest, depois de pagar o download inteiro.
pub fn parse_content_range(h: &str, expected_start: u64) -> Option<u64> {
let rest = h.strip_prefix("bytes ")?;
let (range, total) = rest.split_once('/')?;
let (start, _end) = range.split_once('-')?;
(start.parse::<u64>().ok()? == expected_start).then(|| total.parse().ok())?
}Um servidor pode responder ao pedido de retoma de três maneiras, e só uma é retomável: 206 no offset pedido (retoma), 206 noutro offset (responde a outra pergunta — colar duplicaria o prefixo e a corrupção só apareceria no digest, depois de pagar o download inteiro), ou 200 (ignorou o header: recomeça). E outra armadilha: o Content-Length de um 206 é o tamanho do fragmento, não do blob — o total vem do /<total> do Content-Range.
O que torna seguro costurar dois intervalos é o digest verificado no fim: bytes de duas respostas ou dão o hash publicado, ou o download é descartado. Comprova-se ao vivo no delonix contra um registo que corta a ligação a meio.
5. Estado em disco sem corridas#
Vários processos escrevem no mesmo estado. Duas ferramentas, e precisas das duas.
Escrita atómica#
/// Escrever num ficheiro temporário e fazer `rename`: um leitor vê o antigo inteiro
/// ou o novo inteiro, nunca metade. A ORDEM importa: `fsync` do conteúdo antes do `rename`.
pub fn write_atomic(path: &Path, bytes: &[u8]) -> std::io::Result<()> {
let dir = path.parent().unwrap_or_else(|| Path::new("."));
let tmp = dir.join(format!(
".{}.{}.tmp",
path.file_name().and_then(|n| n.to_str()).unwrap_or("f"),
std::process::id()
));
let mut f = fs::File::create(&tmp)?;
f.write_all(bytes)?;
f.sync_all()?;
fs::rename(&tmp, path)
}Escreve num temporário no mesmo directório e faz rename (atómico no POSIX): um leitor vê o ficheiro antigo inteiro ou o novo inteiro, nunca metade. Repara na ordem: fsync do conteúdo antes do rename. A versão de produção do delonix vai mais longe:
/// [`write_atomic`] with an explicit file mode, set **atomically at creation**.
///
/// For anything secret this is the only correct form. The alternative —
/// `fs::write` then `set_permissions` — creates the file under the ambient
/// umask and narrows it afterwards, leaving a window in which another local
/// user can open it. That is exactly the residual TOCTOU an earlier audit found
/// in the kubeconfig path and closed with `OpenOptions::mode`; the secret store
/// had the same shape and had not been converted.
pub fn write_atomic_mode(path: &Path, bytes: &[u8], mode: Option<u32>) -> Result<()> {
use std::io::Write;
use std::os::unix::fs::OpenOptionsExt;
let dir = path.parent().unwrap_or_else(|| Path::new("."));
let stem = path
.file_name()
.map(|n| n.to_string_lossy().into_owned())
.unwrap_or_else(|| "state".to_string());
// Unique per WRITER (pid + sequence): a fixed temp name lets two processes —
// or two threads of the CRI server — interleave their bytes in the same
// temp, and then `rename` faithfully publishes the corruption.
let seq = TMP_SEQ.fetch_add(1, Ordering::Relaxed);
let tmp = dir.join(format!(".{stem}.{}.{seq}.tmp", std::process::id()));
let write = || -> Result<()> {
let mut opts = fs::OpenOptions::new();
opts.write(true).create(true).truncate(true);
if let Some(m) = mode {
opts.mode(m); // atomic at creation — never widen-then-narrow
}
let mut f = opts.open(&tmp)?;
f.write_all(bytes)?;
// THE ORDER IS THE POINT: the content must be durable BEFORE the
// directory entry that publishes it exists.
f.sync_all()?;
drop(f);
fs::rename(&tmp, path)?;
// And the rename itself must be durable, or a crash can lose the entry
// even though the file's blocks are safely on disk.
if let Ok(d) = fs::File::open(dir) {
let _ = d.sync_all();
}
Ok(())
};
let r = write();
if r.is_err() {
let _ = fs::remove_file(&tmp); // never leave junk behind on failure
}
r
}Quatro detalhes que só se aprendem a perder dados: nome do temporário único por escritor (senão dois escritores intercalam bytes e o rename publica a corrupção), modo definido na criação (OpenOptions::mode) e não chmod depois (janela em que outro utilizador abre o ficheiro), fsync do directório depois do rename, e remover o temporário se falhar.
flock: read-modify-write#
A escrita atómica protege contra ficheiros rasgados, não contra o lost update: dois processos lêem v=1, ambos escrevem v=2, e um incremento perdeu-se. A cura é um lock e reler sob o lock:
/// **Safe read-modify-write** of an item: locks (`flock`), re-reads the
/// state ALREADY under the lock, applies `f` and writes — same pattern as
/// [`Store::update`], generalized to any `JsonStore<T>`.
///
/// Use this (and not `load` + mutate + `save`) whenever the change depends
/// on the CURRENT state and more than one process may touch the same key
/// concurrently — e.g. `delonix-vm`'s `status()` (background metrics
/// refresh) racing a CLI `vm start`/`stop`/`create` on the same VM.
///
/// `f` returns `false` to abort the write (nothing changes). The item
/// returned is the final state (or the one read, if it aborted).
pub fn update<F>(&self, key: &str, f: F) -> Result<T>
where
F: FnOnce(&mut T) -> bool,
{
let _lock = FileLock::acquire(&self.lock_path(key))?;
// Re-read UNDER the lock: between any earlier read and the `flock`
// another process may have written; using a stale value would
// reintroduce the lost update this exists to prevent.
let mut v = self.load(key)?;
if !f(&mut v) {
return Ok(v);
}
self.save(key, &v)?;
Ok(v)
}O comentário no código é o ponto: «Re-read UNDER the lock: between any earlier read and the flock another process may have written». O minicontainer tem o mesmo padrão em Store::update, com um teste que lança 16 threads a incrementar — sem o lock perdia escritas.
Um flock esquecido custou dados
Todos os caminhos de mutação da firewall do delonix contornavam o flock: load → mutar → aplicar no kernel → save. Dois comandos concorrentes aplicavam ambos no kernel, mas só o último save sobrevivia no disco — a regra «perdedora» ficava viva no nft e desaparecia em silêncio no próximo container start. Seis pontos de mutação, todos corrigidos para Store::update.
6. Erros com classe#
#[derive(Debug, thiserror::Error)]
pub enum EngineError {
#[error("not found: {0}")]
NotFound(String),
#[error("conflict: {0}")]
Conflict(String),
#[error("timeout after {0:?}")]
Timeout(Duration),
#[error("{0}")]
Other(String),
}
/// Uma tabela pequena, num só sítio. As mensagens são traduzidas (pt/en) e mudam;
/// o código de saída é contrato — um script que faça `grep 'no such'` parte-se com `--lang=pt`.
pub fn exit_code(e: &EngineError) -> i32 {
match e {
EngineError::NotFound(_) => 4,
EngineError::Conflict(_) => 5,
EngineError::Timeout(_) => 124,
EngineError::Other(_) => 1,
}
}O código de saída é contrato; a mensagem é para humanos e é traduzível. Uma tabela pequena, num só sítio — e um teste que a fixa. Para um lote (rm a b c onde vários falharam): uma classe única mantém-se; um lote misto cai no genérico, porque escolher a classe do primeiro faria o resultado depender da ordem em que escreveste os ids.
E async? tokio, gRPC e o resto#
Tudo acima é síncrono de propósito — o delonix é daemonless e a maioria do motor não precisa de um runtime assíncrono. Onde precisa (o servidor CRI, a API de gestão, o cliente de registos) usa tokio + tonic (gRPC) + hyper, e mantém-nos confinados aos crates de interface, longe da lógica pura. Regras que valem para qualquer sistema em Rust:
- Async só onde há I/O concorrente; lógica de domínio é síncrona e pura (testa-se sem runtime).
- Nunca bloquear dentro de uma task (
std::thread::sleep, I/O síncrono pesado): usaspawn_blocking. - Timeouts em tudo o que toca na rede — um pedido sem deadline é um hang à espera de acontecer. O delonix corrige
read_to_stringde sockets com tecto (o cliente de controlo do holder passou de 5 s para 30 s ao medir 30 attaches concorrentes: 15 falhavam). - Observabilidade por padrões abertos (
tracing→ OpenTelemetry, métricas Prometheus), nuncaprintln!numa biblioteca: quem imprime é a interface.
Verifica o que aprendeste#
cargo test -p examples ch06 # 9 testes
Exercício: torna o plan capaz de remover um campo que nunca esteve em last_applied mas que o manifesto marca explicitamente com null. Que teste escreves primeiro?