diff --git a/docs/EXECUTION.md b/docs/EXECUTION.md index 9166fa2..3d2d7e3 100644 --- a/docs/EXECUTION.md +++ b/docs/EXECUTION.md @@ -874,3 +874,24 @@ escopo da entrega atual nem a próxima tarefa aprovada. quebrada. - Próximo: #33 (Codex App Server autenticado), que também exige ensaio com conta real, ou #41 (recuperação do estado após falha). + +## Workspace — recuperação do estado após falha (#41) + +- Data: 2026-09-28. Decisão no [ADR 0019](decisions/0019-state-lock-and-recovery.md). +- `src/workspace/private.rs` concentra pasta privada, trava do sistema + (`File::try_lock`/`flock`, liberada até em `kill -9`, arquivo mantido com PID), + recuperação de `*.new` remanescente (preservado como `*.new.recovered-`, + nunca carregado nem apagado) e escrita atômica com fsync do arquivo e do + diretório. `ClaudeStore` e o `Store` da demo usam o módulo; ambos informam a + recuperação (aviso na sessão da demo e linha na conversa do Claude oculto). +- Evidências: testes do módulo (processo vivo recusado com PID, trava obsoleta + não bloqueia, symlink recusado, temporário remanescente bloqueia escrita até + ser recuperado, permissões 0600); teste da demo simulando writer morto entre + temporário e rename; scripts PTY verificam que a trava foi liberada via `flock` + não bloqueante; `check_claude_hidden_pty.py` faz `kill -9` da Bee, deixa um + `claude-native.new` parcial e reabre com `--resume` sem limpeza manual. + Latência medida está no ADR (máx. 130,6 ms para 16 MiB neste ambiente). + Checks canônicos, scripts PTY e `check_bundle.py` passaram localmente (Linux). + Autorrevisão. +- Limitações: sistemas de arquivos de rede não suportados; versões antigas não + respeitam a trava nova; persistência continua síncrona. diff --git a/docs/MANUAL_TESTS.md b/docs/MANUAL_TESTS.md index 7727426..b5d3da1 100644 --- a/docs/MANUAL_TESTS.md +++ b/docs/MANUAL_TESTS.md @@ -777,10 +777,10 @@ permitir consultar a origem intacta. Perfil não autentica nem modifica conta re Sair e executar o mesmo comando: histórico reaparece sem reenviar mensagens. Segundo processo com mesma pasta deve recusar. Para simular crash, usar apenas -estado descartável e encerrar seu processo abruptamente: conferir que não resta -processo, inspecionar/remover manualmente `workspace.lock`, reabrir e verificar -interrupção registrada. `workspace.new` remanescente também requer inspeção; -não remover arquivos de outra sessão. Corrupção de JSON e projeto diferente devem +estado descartável e encerrar seu processo com `kill -9`: reabrir sem remover +nada e verificar a interrupção registrada. Criar um `workspace.new` qualquer na +pasta antes de reabrir: ele deve virar `workspace.new.recovered-` e a sessão +mostrar o aviso, com o último estado completo carregado. Corrupção de JSON e projeto diferente devem recusar preservando os bytes originais. Rascunho/tema não são restaurados. ## 48. Exportar a simulação e revisar origem (pendente) diff --git a/docs/decisions/0018-unified-workspace.md b/docs/decisions/0018-unified-workspace.md index 0a87b68..3892e71 100644 --- a/docs/decisions/0018-unified-workspace.md +++ b/docs/decisions/0018-unified-workspace.md @@ -67,9 +67,9 @@ API com cobrança separada como alternativa, pois contradiz a escolha do mantene 3. Validar troca entre agente/perfil com contexto revisado e origem preservada. 4. Lançamento conjunto somente com essas evidências. A demo não satisfaz esse aceite. -Snapshot é substituído por rename após sync do arquivo, mas não promete durabilidade -contra perda de energia (sem fsync de diretório), proteção de ancestrais contra -escritores hostis ou recuperação automática de trava após encerramento abrupto. +Snapshot é substituído por rename após sync do arquivo; trava, recuperação após +encerramento abrupto e fsync de diretório foram revistos no [ADR 0019](0019-state-lock-and-recovery.md). +Não há proteção de ancestrais contra escritores hostis. I/O de persistência/exportação é síncrono; a demo não certifica responsividade sob disco lento. Uma sessão ativa por pasta de estado, não trava global de checkout para um futuro motor real. Colmeia estática, projeto explícito e tema não persistido. diff --git a/docs/decisions/0019-state-lock-and-recovery.md b/docs/decisions/0019-state-lock-and-recovery.md new file mode 100644 index 0000000..05114da --- /dev/null +++ b/docs/decisions/0019-state-lock-and-recovery.md @@ -0,0 +1,48 @@ +# 0019 — Trava do estado pelo sistema e recuperação de escrita + +Status: adotado em 2026-09-28 para a issue #41. Substitui, para as pastas de +estado do workspace, a trava por arquivo criado com `create_new` descrita no +[ADR 0018](0018-unified-workspace.md). + +## Contexto + +`workspace.lock` e `claude-native.lock` eram criados com `create_new` e +apagados na saída. Depois de um kill ou crash, a próxima abertura falhava até o +usuário confirmar que não havia processo ativo e remover o arquivo à mão. Um +`*.new` deixado por uma escrita interrompida também bloqueava todas as gravações +seguintes. O snapshot era sincronizado, mas o diretório não. + +## Decisão + +- Trava consultiva do sistema (`File::try_lock`, `flock` em Unix) sobre o arquivo + de trava, mantida aberta enquanto o estado está aberto. O arquivo continua no + disco e guarda só o PID, para a mensagem de "em uso". O kernel solta a trava + quando o processo termina, inclusive por `kill -9`; um processo vivo nunca é + desalojado. Arquivo de trava que não seja arquivo comum (symlink) é recusado. + Havendo disputa, a abertura tenta por até 500 ms: um `fork` de outra thread do + mesmo processo mantém uma cópia do descritor até o `exec` (visto como falha + intermitente de 7 em 60 execuções dos testes paralelos; 0 em 60 com a espera). +- Com a trava obtida, um `*.new` remanescente é de uma escrita que parou antes do + rename: o snapshot confirmado continua sendo o último estado completo. O + remanescente é renomeado para `*.new.recovered-` na mesma pasta, + nunca carregado nem apagado, e a interface informa onde ficou. +- Escrita: arquivo temporário novo 0600, `sync_all`, rename sobre o snapshot e + `sync_all` do diretório. Em macOS, `sync_all` usa `F_FULLFSYNC`. Erro antes do + rename remove o temporário e preserva o snapshot anterior. +- Código compartilhado em `src/workspace/private.rs` pelas duas pastas de estado. + +## Garantias e limites + +- Depois de `save` retornar sucesso, o snapshot sobrevive a queda de energia em + sistemas de arquivos locais que honram `fsync`. Sistemas de arquivos de rede + podem não implementar `flock` nem essa durabilidade; não são suportados. +- Proteção contra outro processo Memory Bee, não contra escritores hostis com o + mesmo usuário nem contra alteração dos diretórios ancestrais. +- Versões anteriores apagavam a trava; uma versão antiga aberta junto com a nova + na mesma pasta não respeita o `flock`. Não misturar versões numa pasta. +- Persistência continua síncrona. Medição local (Linux, disco virtual, `cargo + test --release -- --ignored snapshot_latency`): estado Claude de 64 KiB com + mediana 0,6 ms e máximo 23,6 ms; estado da demo no limite de 16 MiB com + mediana 28,5 ms e máximo 130,6 ms. Discos lentos podem congelar a interface + por mais tempo; mover a escrita para outra thread fica para quando um estado + real grande existir. diff --git a/docs/specs/workspace.md b/docs/specs/workspace.md index 4af1880..020c9fc 100644 --- a/docs/specs/workspace.md +++ b/docs/specs/workspace.md @@ -72,8 +72,8 @@ Em todos os casos, se o CLI não sair em cinco segundos, o processo é encerrado `claude-native.json` guarda projeto canônico, IDs e códigos de saída (até 128 sessões, 64 KiB), sem transcrição. A pasta é 0700 e arquivos são 0600 em Unix; -`claude-native.lock` impede segundo escritor. Após crash, verificar que não há -processo ativo antes de remover a trava antiga. `--resume` usa o último ID desta +`claude-native.lock` impede segundo escritor com trava do sistema, liberada até +em `kill -9`; não há limpeza manual ([ADR 0019](../decisions/0019-state-lock-and-recovery.md)). `--resume` usa o último ID desta pasta; não reenvia prompt. A sessão nativa pode não existir quando se sai antes de enviar uma mensagem. Isolamento de perfis e integração Codex estão pendentes. @@ -181,9 +181,12 @@ e eventos. Limites: 16 MiB, 128 sessões, 32 perfis, 100 mil eventos por sessão memória continua exportável. A sessão carregada como ativa é marcada interrompida, sem reexecutar/reencaminhar mensagens. Tema e rascunho não são persistidos. -`workspace.lock` recusa segundo escritor. Após kill/crash, confirmar que não existe -outro processo usando a pasta antes de remover manualmente a trava. `workspace.new` -remanescente também exige inspeção manual; não apagar automaticamente. JSON inválido, +`workspace.lock` recusa segundo escritor com trava do sistema, liberada quando o +processo termina, inclusive por kill; o arquivo fica no disco e não exige remoção. +Um `workspace.new` remanescente (escrita interrompida) é preservado como +`workspace.new.recovered-`, não é carregado, e a sessão registra o aviso; o +último snapshot completo é o carregado. Gravações sincronizam arquivo e diretório +([ADR 0019](../decisions/0019-state-lock-and-recovery.md)). JSON inválido, projeto diferente, versão desconhecida e symlink de estado são recusados sem substituir o arquivo original. Não é armazenamento cifrado nem trava de checkout. diff --git a/scripts/check_claude_hidden_pty.py b/scripts/check_claude_hidden_pty.py index 320c0b1..40987a2 100644 --- a/scripts/check_claude_hidden_pty.py +++ b/scripts/check_claude_hidden_pty.py @@ -12,6 +12,13 @@ import termios import time +def assert_unlocked(path): + """The lock file stays on disk; the process must have released it.""" + with open(path) as handle: + fcntl.flock(handle, fcntl.LOCK_EX | fcntl.LOCK_NB) + fcntl.flock(handle, fcntl.LOCK_UN) + + FAKE = r'''#!/usr/bin/env python3 import json, os, signal, subprocess, sys @@ -229,7 +236,7 @@ def untrusted(binary, project, state, env, quit_while_blocked): drain_until_exit(master, proc, output, 10) assert proc.returncode == 0, proc.returncode assert termios.tcgetattr(slave) == before - assert not (state / 'claude-native.lock').exists() + assert_unlocked(state / 'claude-native.lock') finally: if proc.poll() is None: proc.kill() @@ -278,7 +285,7 @@ def until(token): assert proc.returncode == 0, proc.returncode assert b'CLAUDE ORIGINAL SHOULD STAY HIDDEN' not in output assert termios.tcgetattr(slave) == before - assert not (state / 'claude-native.lock').exists() + assert_unlocked(state / 'claude-native.lock') finally: if proc.poll() is None: proc.kill() @@ -353,6 +360,21 @@ def main(): history_and_export(binary, project, state, env, True) # History is shown from the transcript, never typed back into Claude. assert stdin_log.read_text().splitlines() == ['/exit', '/exit'] + # A killed Bee leaves its lock file and a partial write behind; the + # next start needs no manual cleanup and keeps the partial file aside. + master, slave, _, proc, output, until = spawn(binary, project, state, env, True) + try: + until(b'decrescente') + proc.kill() + proc.wait() + finally: + os.close(master) + os.close(slave) + (state / 'claude-native.new').write_text('partial') + history_and_export(binary, project, state, env, True) + kept = [p for p in state.iterdir() if p.name.startswith('claude-native.new.recovered-')] + assert len(kept) == 1 and kept[0].read_text() == 'partial' + assert json.loads((state / 'claude-native.json').read_text())['sessions'] with tempfile.TemporaryDirectory(prefix='bee-editor-') as temp: root = Path(temp) fake_dir = root / 'bin' @@ -365,7 +387,7 @@ def main(): env = dict(os.environ, PATH=f'{fake_dir}:/usr/bin:/bin', BEE_BINARY=str(binary), BEE_FAKE_STDIN=str(root / 'stdin')) editor(binary, project, root / 'state-a', dict(env, BEE_BRACKETED='1'), (24, 80), True) editor(binary, project, root / 'state-b', env, (12, 40), False) # Documented minimum width. - print('PASS: Bee-only UI, hidden original output, hooks, deny/allow, Ctrl+Q denial, resume, blocked native setup via Ctrl+O, safe Ctrl+Q, history rehydration, verified export, multiline editor, scrolling, 40-column controls and terminal restoration') + print('PASS: Bee-only UI, hidden original output, hooks, deny/allow, Ctrl+Q denial, resume, blocked native setup via Ctrl+O, safe Ctrl+Q, history rehydration, verified export, multiline editor, scrolling, 40-column controls, recovery after kill and terminal restoration') if __name__ == '__main__': diff --git a/scripts/check_claude_native_pty.py b/scripts/check_claude_native_pty.py index 0c59608..a0f91d1 100644 --- a/scripts/check_claude_native_pty.py +++ b/scripts/check_claude_native_pty.py @@ -12,6 +12,13 @@ import termios import time +def assert_unlocked(path): + """The lock file stays on disk; the process must have released it.""" + with open(path) as handle: + fcntl.flock(handle, fcntl.LOCK_EX | fcntl.LOCK_NB) + fcntl.flock(handle, fcntl.LOCK_UN) + + def one_run(binary, project, state, env, resume=False): master, slave = pty.openpty() @@ -45,7 +52,7 @@ def drain_until(token, seconds=5): proc.wait(timeout=5) assert proc.returncode == 0 assert termios.tcgetattr(slave) == before, 'terminal state not restored' - assert not (state / 'claude-native.lock').exists() + assert_unlocked(state / 'claude-native.lock') return json.loads((state / 'claude-native.json').read_text()) finally: if proc.poll() is None: diff --git a/scripts/check_workspace_pty.py b/scripts/check_workspace_pty.py index 9ff92d8..b19584e 100644 --- a/scripts/check_workspace_pty.py +++ b/scripts/check_workspace_pty.py @@ -12,6 +12,13 @@ import termios import time +def assert_unlocked(path): + """The lock file stays on disk; the process must have released it.""" + with open(path) as handle: + fcntl.flock(handle, fcntl.LOCK_EX | fcntl.LOCK_NB) + fcntl.flock(handle, fcntl.LOCK_UN) + + def main(): repo = Path(__file__).resolve().parent.parent @@ -64,7 +71,8 @@ def wait_for(predicate): paste('synthetic task\nsecond line') assert not snapshot()['sessions'][0]['events'], 'Paste submitted a message' send('\r') - wait_for(lambda: snapshot()['sessions'][0]['events'][-1]['kind'] == 'approval') + # Saving is asynchronous to the keystroke; wait instead of indexing an empty list. + wait_for(lambda: snapshot()['sessions'][0]['events'][-1:] and snapshot()['sessions'][0]['events'][-1]['kind'] == 'approval') send('\t') send('\r') # Default deny. wait_for(lambda: not snapshot()['sessions'][0]['running']) @@ -96,7 +104,7 @@ def wait_for(predicate): proc.wait(timeout=5) assert proc.returncode == 0 assert termios.tcgetattr(slave) == before, 'Terminal was not restored' - assert not (state / 'workspace.lock').exists() + assert_unlocked(state / 'workspace.lock') assert list(project.iterdir()) == [], 'Demo modified project' print('PASS: paste, approval denial, confirmed export/schema/verify, resize, interruption, private state and terminal restoration') finally: diff --git a/src/tui/claude_native.rs b/src/tui/claude_native.rs index c7670f9..d5b113c 100644 --- a/src/tui/claude_native.rs +++ b/src/tui/claude_native.rs @@ -1707,6 +1707,12 @@ pub fn run_hidden( } let mut view = HiddenView::new(Instant::now()); view.resume = resume; + if let Some(kept) = &store.recovered { + view.push(format!( + "Gravação anterior interrompida; estado completo carregado e parcial preservado em {}.", + kept.display() + )); + } let mut status = None; let mut quit_at = None; let mut dirty = true; diff --git a/src/tui/draft.rs b/src/tui/draft.rs index cd577de..69237c3 100644 --- a/src/tui/draft.rs +++ b/src/tui/draft.rs @@ -68,18 +68,27 @@ impl Draft { pub fn end(&mut self) { self.cursor = self.line_end(); } - /// Byte offset of the `column`-th char of the line starting at `start`. + /// Byte offset in the line starting at `start` whose display column is + /// closest to `column` without passing it (wide chars take two cells). fn at_column(&self, start: usize, column: usize) -> usize { let line = &self.text[start..]; let line = &line[..line.find('\n').unwrap_or(line.len())]; - start - + line - .char_indices() - .nth(column) - .map_or(line.len(), |(i, _)| i) + let mut used = 0; + for (i, c) in line.char_indices() { + let w = c.width().unwrap_or(0); + if used + w > column { + return start + i; + } + used += w; + } + start + line.len() } + /// Display column of the cursor, in terminal cells. fn column(&self) -> usize { - self.text[self.line_start()..self.cursor].chars().count() + self.text[self.line_start()..self.cursor] + .chars() + .map(|c| c.width().unwrap_or(0)) + .sum() } /// Returns false on the first line, so the caller can use Up elsewhere. pub fn up(&mut self) -> bool { @@ -185,6 +194,21 @@ mod tests { assert_eq!(d.text(), ">primeir\nb\nerceira"); } + #[test] + fn vertical_moves_keep_the_visual_column_with_wide_chars() { + let mut d = draft("漢字漢\nabcdef"); + d.left(); + d.left(); + assert_eq!(d.layout(80).1, (1, 4)); + assert!(d.up()); + assert_eq!(d.layout(80).1, (0, 4), "lands after two wide chars"); + assert!(d.down()); + assert_eq!(d.layout(80).1, (1, 4)); + d.left(); + assert!(d.up()); + assert_eq!(d.layout(80).1, (0, 2), "never inside a wide char"); + } + #[test] fn paste_normalizes_breaks_and_drops_controls() { let mut d = draft("a\r\nb\rc\x1b[31m\td"); diff --git a/src/workspace/claude_native.rs b/src/workspace/claude_native.rs index 6b176a3..c739f31 100644 --- a/src/workspace/claude_native.rs +++ b/src/workspace/claude_native.rs @@ -3,8 +3,8 @@ use serde::{Deserialize, Serialize}; use std::{ collections::BTreeSet, - fs::{self, File, OpenOptions}, - io::{Read, Write}, + fs::{self, File}, + io::Read, path::{Path, PathBuf}, }; use uuid::Uuid; @@ -45,47 +45,22 @@ impl ClaudeState { pub struct ClaudeStore { root: PathBuf, -} - -fn private_new(path: &Path) -> Result { - let mut options = OpenOptions::new(); - options.write(true).create_new(true); - #[cfg(unix)] - { - use std::os::unix::fs::OpenOptionsExt; - options.mode(0o600); - } - options.open(path).map_err(|e| e.to_string()) + _lock: super::private::Lock, + /// Interrupted write found on open, kept aside and not loaded. + pub recovered: Option, } impl ClaudeStore { pub fn open(root: &Path, project: &Path) -> Result<(Self, ClaudeState), String> { - if !root.exists() { - let mut builder = fs::DirBuilder::new(); - #[cfg(unix)] - { - use std::os::unix::fs::DirBuilderExt; - builder.mode(0o700); - } - builder.create(root).map_err(|e| e.to_string())?; - } - let metadata = fs::symlink_metadata(root).map_err(|e| e.to_string())?; - if !metadata.is_dir() || metadata.file_type().is_symlink() { - return Err("Estado exige pasta real, sem symlink".into()); - } - #[cfg(unix)] - { - use std::os::unix::fs::PermissionsExt; - if metadata.permissions().mode() & 0o077 != 0 { - return Err("Pasta de estado deve ter permissão 0700".into()); - } - } - let root = root.canonicalize().map_err(|e| e.to_string())?; + let root = super::private::directory(root)?; let project = project.to_str().ok_or("Projeto deve ser UTF-8")?; - let store = Self { root }; - let mut lock = private_new(&store.root.join("claude-native.lock")) - .map_err(|_| "Estado Claude em uso ou trava antiga presente. Confirme que não há processo ativo antes de remover claude-native.lock.")?; - writeln!(lock, "{}", std::process::id()).map_err(|e| e.to_string())?; + let lock = super::private::Lock::acquire(&root, "claude-native.lock")?; + let recovered = super::private::recover(&root, "claude-native.new")?; + let store = Self { + root, + _lock: lock, + recovered, + }; let path = store.root.join("claude-native.json"); let state = match fs::symlink_metadata(&path) { Ok(meta) => { @@ -152,26 +127,14 @@ impl ClaudeStore { if bytes.len() as u64 > MAX_BYTES { return Err("Estado Claude excede 64 KiB".into()); } - let temp = self.root.join("claude-native.new"); - let mut file = - private_new(&temp).map_err(|_| "claude-native.new já existe; original preservado")?; - let result = file - .write_all(&bytes) - .and_then(|()| file.sync_all()) - .and_then(|()| fs::rename(&temp, self.root.join("claude-native.json"))); - if let Err(e) = result { - let _ = fs::remove_file(&temp); - return Err(e.to_string()); - } - Ok(()) + super::private::write_atomic( + &self.root, + "claude-native.new", + "claude-native.json", + &bytes, + ) } } -impl Drop for ClaudeStore { - fn drop(&mut self) { - let _ = fs::remove_file(self.root.join("claude-native.lock")); - } -} - /// Transcript reported by Claude's own `SessionStart` hook. Only a regular, /// non-symlink `.jsonl` is accepted; nothing is guessed. pub fn transcript_path(reported: &str, id: Uuid) -> Result { diff --git a/src/workspace/mod.rs b/src/workspace/mod.rs index 5f33702..860bcb8 100644 --- a/src/workspace/mod.rs +++ b/src/workspace/mod.rs @@ -2,6 +2,7 @@ pub mod adapter; pub mod claude_native; pub mod codex; +pub mod private; pub mod store; use serde::{Deserialize, Serialize}; diff --git a/src/workspace/private.rs b/src/workspace/private.rs new file mode 100644 index 0000000..c13b951 --- /dev/null +++ b/src/workspace/private.rs @@ -0,0 +1,224 @@ +//! Private state directory shared by the workspace stores: a single-writer +//! lock the operating system releases when the process dies, recovery of an +//! interrupted write, and atomic snapshots flushed before they count. +use std::{ + fs::{self, File, OpenOptions, TryLockError}, + io::{Read, Write}, + path::{Path, PathBuf}, +}; + +/// Creates the directory (0700) or checks an existing one; returns it canonical. +pub fn directory(root: &Path) -> Result { + if !root.exists() { + let mut builder = fs::DirBuilder::new(); + #[cfg(unix)] + { + use std::os::unix::fs::DirBuilderExt; + builder.mode(0o700); + } + builder.create(root).map_err(|e| e.to_string())?; + } + let meta = fs::symlink_metadata(root).map_err(|e| e.to_string())?; + if !meta.is_dir() || meta.file_type().is_symlink() { + return Err("Estado exige pasta real, sem symlink".into()); + } + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + if meta.permissions().mode() & 0o077 != 0 { + return Err("Pasta de estado deve ter permissão 0700".into()); + } + } + root.canonicalize().map_err(|e| e.to_string()) +} + +fn regular_or_missing(path: &Path) -> Result { + match fs::symlink_metadata(path) { + Ok(meta) if meta.is_file() && !meta.file_type().is_symlink() => Ok(true), + Ok(_) => Err(format!( + "{} deve ser arquivo comum, sem symlink", + path.file_name().unwrap_or_default().to_string_lossy() + )), + Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false), + Err(e) => Err(e.to_string()), + } +} + +/// Advisory lock held while the store is open. The file stays on disk; the +/// kernel drops the lock when the holder exits or is killed, so a crash never +/// requires manual cleanup and a live process is never displaced. +pub struct Lock { + _file: File, +} + +impl Lock { + pub fn acquire(root: &Path, name: &str) -> Result { + let path = root.join(name); + regular_or_missing(&path)?; + let mut options = OpenOptions::new(); + options.read(true).write(true).create(true).truncate(false); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.mode(0o600); + } + let mut file = options.open(&path).map_err(|e| e.to_string())?; + // A fork by another thread briefly holds a copy of a just-closed + // descriptor until its exec, so contention gets a short grace period. + // A process that really holds the lock is still refused. + let deadline = std::time::Instant::now() + std::time::Duration::from_millis(500); + let mut result = file.try_lock(); + while matches!(result, Err(TryLockError::WouldBlock)) + && std::time::Instant::now() < deadline + { + std::thread::sleep(std::time::Duration::from_millis(10)); + result = file.try_lock(); + } + match result { + Ok(()) => {} + Err(TryLockError::WouldBlock) => { + let mut holder = String::new(); + let _ = Read::by_ref(&mut file).take(32).read_to_string(&mut holder); + return Err(format!( + "Estado em uso por outro processo Memory Bee (pid {}); nada foi alterado.", + holder.trim() + )); + } + Err(TryLockError::Error(e)) => { + return Err(format!("Trava do estado indisponível: {e}")); + } + } + file.set_len(0) + .and_then(|()| writeln!(file, "{}", std::process::id())) + .map_err(|e| e.to_string())?; + Ok(Self { _file: file }) + } +} + +/// A leftover temporary file means a previous writer stopped before its +/// rename, so the committed snapshot is still the last complete state. The +/// leftover is kept beside it under a new name, never loaded or deleted. +/// Call only while holding the lock. +pub fn recover(root: &Path, temp: &str) -> Result, String> { + let path = root.join(temp); + if !regular_or_missing(&path)? { + return Ok(None); + } + let stamp = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map_or(0, |d| d.as_nanos()); + let kept = root.join(format!("{temp}.recovered-{stamp}")); + fs::rename(&path, &kept).map_err(|e| e.to_string())?; + Ok(Some(kept)) +} + +/// Writes `bytes` to a new private temporary file, flushes it, renames it +/// over `target` and flushes the directory. The old snapshot stays intact +/// until the rename; on error the temporary file is removed. +pub fn write_atomic(root: &Path, temp: &str, target: &str, bytes: &[u8]) -> Result<(), String> { + let temp_path = root.join(temp); + let mut options = OpenOptions::new(); + options.write(true).create_new(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.mode(0o600); + } + let mut file = options + .open(&temp_path) + .map_err(|_| format!("{temp} já existe ou está indisponível; original preservado"))?; + let result = file + .write_all(bytes) + .and_then(|()| file.sync_all()) + .and_then(|()| fs::rename(&temp_path, root.join(target))) + .map_err(|e| e.to_string()); + if result.is_err() { + let _ = fs::remove_file(&temp_path); + return result; + } + // The rename itself is durable only once the directory is flushed. + #[cfg(unix)] + File::open(root) + .and_then(|dir| dir.sync_all()) + .map_err(|e| format!("Estado gravado, mas a pasta não confirmou a gravação: {e}"))?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn temp_root() -> PathBuf { + let root = std::env::temp_dir().join(format!("bee-private-{}", uuid::Uuid::new_v4())); + directory(&root).unwrap() + } + + #[test] + fn lock_excludes_a_live_holder_and_survives_a_stale_file() { + let root = temp_root(); + let lock = Lock::acquire(&root, "x.lock").unwrap(); + let busy = Lock::acquire(&root, "x.lock").err().unwrap(); + assert!(busy.contains(&std::process::id().to_string()), "{busy}"); + drop(lock); + assert!( + root.join("x.lock").exists(), + "the lock file is never deleted" + ); + // A stale file (as after a kill) does not block the next open. + fs::write(root.join("x.lock"), "99999\n").unwrap(); + Lock::acquire(&root, "x.lock").unwrap(); + fs::remove_dir_all(root).unwrap(); + } + + /// Blocking time of one snapshot at each store's size limit; run with + /// `cargo test --locked -- --ignored --nocapture snapshot_latency`. + #[test] + #[ignore] + fn snapshot_latency() { + let root = temp_root(); + for (label, size) in [ + ("claude-native 64 KiB", 64 * 1024), + ("workspace 16 MiB", 16 << 20), + ] { + let bytes = vec![b'x'; size]; + let mut times = Vec::new(); + for _ in 0..5 { + let start = std::time::Instant::now(); + write_atomic(&root, "m.new", "m.json", &bytes).unwrap(); + times.push(start.elapsed().as_secs_f64() * 1000.0); + } + times.sort_by(f64::total_cmp); + println!( + "{label}: mediana {:.1} ms, máx {:.1} ms", + times[2], times[4] + ); + } + fs::remove_dir_all(root).unwrap(); + } + #[test] + fn interrupted_write_is_kept_aside_and_snapshots_replace_atomically() { + let root = temp_root(); + write_atomic(&root, "s.new", "s.json", b"old").unwrap(); + fs::write(root.join("s.new"), b"half").unwrap(); + assert!(write_atomic(&root, "s.new", "s.json", b"new").is_err()); + assert_eq!(fs::read(root.join("s.json")).unwrap(), b"old"); + let kept = recover(&root, "s.new").unwrap().unwrap(); + assert_eq!(fs::read(&kept).unwrap(), b"half"); + assert_eq!(recover(&root, "s.new").unwrap(), None); + write_atomic(&root, "s.new", "s.json", b"new").unwrap(); + assert_eq!(fs::read(root.join("s.json")).unwrap(), b"new"); + assert!(!root.join("s.new").exists()); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + let mode = fs::metadata(root.join("s.json")) + .unwrap() + .permissions() + .mode(); + assert_eq!(mode & 0o077, 0); + std::os::unix::fs::symlink(root.join("s.json"), root.join("l.lock")).unwrap(); + assert!(Lock::acquire(&root, "l.lock").is_err()); + } + fs::remove_dir_all(root).unwrap(); + } +} diff --git a/src/workspace/store.rs b/src/workspace/store.rs index 0743b57..f050aba 100644 --- a/src/workspace/store.rs +++ b/src/workspace/store.rs @@ -2,8 +2,8 @@ use super::{Agent, Event, Session}; use serde::{Deserialize, Serialize}; use std::{ - fs::{self, OpenOptions}, - io::{Read, Write}, + fs, + io::Read, path::{Path, PathBuf}, }; const MAX_BYTES: u64 = 16 * 1024 * 1024; @@ -96,43 +96,14 @@ impl State { pub struct Store { root: PathBuf, -} -fn private_file(path: &Path) -> Result { - let mut options = OpenOptions::new(); - options.write(true).create_new(true); - #[cfg(unix)] - { - use std::os::unix::fs::OpenOptionsExt; - options.mode(0o600); - } - options.open(path).map_err(|e| e.to_string()) + _lock: super::private::Lock, } impl Store { pub fn open(root: &Path, project: &str, agent: Agent) -> Result<(Self, State), String> { - if !root.exists() { - let mut builder = fs::DirBuilder::new(); - #[cfg(unix)] - { - use std::os::unix::fs::DirBuilderExt; - builder.mode(0o700); - } - builder.create(root).map_err(|e| e.to_string())?; - } - let meta = fs::symlink_metadata(root).map_err(|e| e.to_string())?; - if !meta.is_dir() || meta.file_type().is_symlink() { - return Err("Estado exige pasta real, sem symlink".into()); - } - #[cfg(unix)] - { - use std::os::unix::fs::PermissionsExt; - if meta.permissions().mode() & 0o077 != 0 { - return Err("Pasta de estado deve ter permissão 0700".into()); - } - } - let root = root.canonicalize().map_err(|e| e.to_string())?; - let mut lock = private_file(&root.join("workspace.lock")).map_err(|_| "Estado em uso ou trava anterior presente. Após confirmar que não existe outro processo, remova workspace.lock manualmente.")?; - let store = Self { root }; - writeln!(lock, "{}", std::process::id()).map_err(|e| e.to_string())?; + let root = super::private::directory(root)?; + let lock = super::private::Lock::acquire(&root, "workspace.lock")?; + let recovered = super::private::recover(&root, "workspace.new")?; + let store = Self { root, _lock: lock }; let path = store.root.join("workspace.json"); let metadata = match fs::symlink_metadata(&path) { Ok(meta) => Some(meta), @@ -158,6 +129,14 @@ impl Store { State::new(project.into(), agent) }; state.validate(project)?; + if let Some(kept) = recovered { + state.current_mut().record(Event::Notice { + message: format!( + "Gravação anterior interrompida; último estado completo carregado e o arquivo parcial preservado em {}.", + kept.display() + ), + }); + } for session in &mut state.sessions { if session.running { session.record(Event::Interrupted); @@ -175,25 +154,6 @@ impl Store { if bytes.len() as u64 > MAX_BYTES { return Err("Estado acima de 16 MiB; exporte o histórico antes de continuar".into()); } - let temp = self.root.join("workspace.new"); - let mut file = private_file(&temp).map_err( - |_| "Arquivo workspace.new já existe ou está indisponível; original preservado", - )?; - let result = (|| { - file.write_all(&bytes) - .and_then(|()| file.sync_all()) - .map_err(|e| e.to_string())?; - fs::rename(&temp, self.root.join("workspace.json")).map_err(|e| e.to_string())?; - Ok(()) - })(); - if result.is_err() { - let _ = fs::remove_file(&temp); - } - result - } -} -impl Drop for Store { - fn drop(&mut self) { - let _ = fs::remove_file(self.root.join("workspace.lock")); + super::private::write_atomic(&self.root, "workspace.new", "workspace.json", &bytes) } } diff --git a/tests/workspace.rs b/tests/workspace.rs index 7d96977..c09c438 100644 --- a/tests/workspace.rs +++ b/tests/workspace.rs @@ -97,7 +97,9 @@ fn store_refuses_foreign_project_and_preserves_malformed_data() { let (store, _) = Store::open(&path, "/one", Agent::Claude).unwrap(); drop(store); assert!(Store::open(&path, "/two", Agent::Codex).is_err()); - assert!(!path.join("workspace.lock").exists()); + // The lock file stays, but closing the store released it. + assert!(path.join("workspace.lock").exists()); + drop(Store::open(&path, "/one", Agent::Claude).unwrap()); fs::write(path.join("workspace.json"), "bad data").unwrap(); assert!(Store::open(&path, "/one", Agent::Codex).is_err()); assert_eq!( @@ -482,3 +484,33 @@ mod claude_transcript { assert!(excluded.history().lines().count() < prepared.history().lines().count()); } } + +#[test] +fn store_recovers_after_a_killed_writer_without_losing_data() { + let temp = Temp::new(); + let path = temp.0.join("state"); + let (store, state) = Store::open(&path, "/one", Agent::Claude).unwrap(); + let busy = Store::open(&path, "/one", Agent::Claude).err().unwrap(); + assert!(busy.contains("em uso"), "{busy}"); + drop(store); + // A kill between writing the temporary file and renaming it. + fs::write(path.join("workspace.new"), "partial").unwrap(); + let before = fs::read(path.join("workspace.json")).unwrap(); + let (_store, reopened) = Store::open(&path, "/one", Agent::Claude).unwrap(); + assert_eq!(reopened.sessions.len(), state.sessions.len()); + let kept: Vec<_> = fs::read_dir(&path) + .unwrap() + .filter_map(|e| e.ok()) + .filter(|e| { + e.file_name() + .to_string_lossy() + .starts_with("workspace.new.recovered-") + }) + .collect(); + assert_eq!(kept.len(), 1); + assert_eq!(fs::read_to_string(kept[0].path()).unwrap(), "partial"); + assert!(!path.join("workspace.new").exists()); + let notice = serde_json::to_string(&reopened.current().events).unwrap(); + assert!(notice.contains("arquivo parcial preservado"), "{notice}"); + assert_ne!(fs::read(path.join("workspace.json")).unwrap(), before); +}