| 1 | //! The on-disk ring buffer, the tombstone, and the wipe. |
| 2 | //! |
| 3 | //! Everything lives under `$CODEWHALE_HOME/telemetry/`, created `0700`, with |
| 4 | //! every file `0600`: |
| 5 | //! |
| 6 | //! | file | role | |
| 7 | //! |---|---| |
| 8 | //! | `buffer.jsonl` | one JSON event per line, awaiting a flush | |
| 9 | //! | `buffer.jsonl.lock` | a **sibling** lock file; never the data file | |
| 10 | //! | `dryrun.jsonl` | the sink when the endpoint resolves empty, same ring policy | |
| 11 | //! | `state.json` | last version seen and last flush attempt | |
| 12 | //! | `install_id.json` | the random install id | |
| 13 | //! | `disabled` | the tombstone: present ⇒ nothing is appended, drained, or sent | |
| 14 | //! |
| 15 | //! **Appends never take a lock.** One `O_APPEND` `write(2)` under `PIPE_BUF` is |
| 16 | //! atomic on every filesystem this ships to, and taking `fd_lock` here would be |
| 17 | //! a *blocking* acquisition on the panic hook and the SIGINT path. `flock` is |
| 18 | //! per-fd within a process, so an actor panic while holding the compaction lock |
| 19 | //! would self-deadlock the hook — `catch_unwind` runs *after* the hook, so it |
| 20 | //! cannot save this — and a second Codewhale process sharing `CODEWHALE_HOME` |
| 21 | //! would hang Ctrl-C, breaking the second-signal contract in `main.rs`. |
| 22 | //! |
| 23 | //! Compaction is the only lock holder and uses `try_write()`: on contention it |
| 24 | //! skips this cycle. Appenders re-open per append, so a compaction rewrite |
| 25 | //! cannot leave anyone writing to a stale inode. |
| 26 | |
| 27 | use std::fs::{self, DirBuilder, File, OpenOptions}; |
| 28 | use std::io::Write as _; |
| 29 | use std::path::{Path, PathBuf}; |
| 30 | |
| 31 | use anyhow::{Context, Result}; |
| 32 | |
| 33 | /// Newest events retained in either sink. |
| 34 | pub const MAX_EVENTS: usize = 512; |
| 35 | /// Byte ceiling for either sink. |
| 36 | pub const MAX_BYTES: u64 = 256 * 1024; |
| 37 | /// A single append must fit in one atomic `write(2)`. |
| 38 | pub const MAX_LINE_BYTES: usize = 4096; |
| 39 | |
| 40 | /// Below this size a sink cannot possibly hold [`MAX_EVENTS`] lines, so an |
| 41 | /// append skips the count probe entirely. The shortest serializable event line |
| 42 | /// is well over 8 bytes, and `512 * 9 > 4096`, so this bound is safe by |
| 43 | /// construction — `probe_threshold_cannot_hide_an_over_cap_buffer` pins it. |
| 44 | const PROBE_BYTES: u64 = 4096; |
| 45 | |
| 46 | /// `buffer.jsonl` — the pending-event sink. |
| 47 | #[must_use] |
| 48 | pub fn buffer_path(root: &Path) -> PathBuf { |
| 49 | root.join("buffer.jsonl") |
| 50 | } |
| 51 | |
| 52 | /// `dryrun.jsonl` — where batches go when the endpoint resolves to `None`. |
| 53 | /// |
| 54 | /// Reached by configuring `telemetry_endpoint` empty; an unconfigured endpoint |
| 55 | /// resolves to `codewhale_config::DEFAULT_TELEMETRY_ENDPOINT` instead. |
| 56 | #[must_use] |
| 57 | pub fn dryrun_path(root: &Path) -> PathBuf { |
| 58 | root.join("dryrun.jsonl") |
| 59 | } |
| 60 | |
| 61 | /// `buffer.jsonl.lock` — the sibling lock file. Never the data file, and never |
| 62 | /// unlinked: replacing it would leave appenders and compactors holding |
| 63 | /// different inodes and serialising against nothing. |
| 64 | #[must_use] |
| 65 | pub fn lock_path(root: &Path) -> PathBuf { |
| 66 | root.join("buffer.jsonl.lock") |
| 67 | } |
| 68 | |
| 69 | /// `disabled` — the tombstone. |
| 70 | #[must_use] |
| 71 | pub fn tombstone_path(root: &Path) -> PathBuf { |
| 72 | root.join("disabled") |
| 73 | } |
| 74 | |
| 75 | /// `install_id.json`. |
| 76 | #[must_use] |
| 77 | pub fn install_id_path(root: &Path) -> PathBuf { |
| 78 | root.join("install_id.json") |
| 79 | } |
| 80 | |
| 81 | /// `state.json`. |
| 82 | #[must_use] |
| 83 | pub fn state_path(root: &Path) -> PathBuf { |
| 84 | root.join("state.json") |
| 85 | } |
| 86 | |
| 87 | /// Whether the tombstone is present. |
| 88 | /// |
| 89 | /// Re-checked on **every** append and immediately before **every** send. This is |
| 90 | /// what makes `codewhale config set telemetry false` — an external write by |
| 91 | /// another process — observable to a session that is already running. |
| 92 | #[must_use] |
| 93 | pub fn tombstone_present(root: &Path) -> bool { |
| 94 | tombstone_path(root).exists() |
| 95 | } |
| 96 | |
| 97 | /// Create the telemetry directory `0700`, if it is missing. |
| 98 | pub fn ensure_dir(root: &Path) -> Result<()> { |
| 99 | if root.is_dir() { |
| 100 | return Ok(()); |
| 101 | } |
| 102 | let mut builder = DirBuilder::new(); |
| 103 | builder.recursive(true); |
| 104 | #[cfg(unix)] |
| 105 | { |
| 106 | use std::os::unix::fs::DirBuilderExt as _; |
| 107 | builder.mode(0o700); |
| 108 | } |
| 109 | builder |
| 110 | .create(root) |
| 111 | .with_context(|| format!("failed to create {}", root.display())) |
| 112 | } |
| 113 | |
| 114 | #[cfg(unix)] |
| 115 | fn secure(file: &File) -> Result<()> { |
| 116 | use std::os::unix::fs::PermissionsExt as _; |
| 117 | file.set_permissions(fs::Permissions::from_mode(0o600)) |
| 118 | .context("failed to restrict telemetry file permissions") |
| 119 | } |
| 120 | |
| 121 | #[cfg(not(unix))] |
| 122 | fn secure(_file: &File) -> Result<()> { |
| 123 | Ok(()) |
| 124 | } |
| 125 | |
| 126 | /// Append one serialized event or batch to `path`. |
| 127 | /// |
| 128 | /// Returns `None` — never an error — when the tombstone is present, when the |
| 129 | /// line would not fit in one atomic write, or when any filesystem step fails. |
| 130 | /// Telemetry is fail-open by construction: it never returns an error to a |
| 131 | /// caller and never blocks a turn, a tool, or process exit. |
| 132 | pub fn append(root: &Path, path: &Path, line: &str) -> Option<()> { |
| 133 | if tombstone_present(root) { |
| 134 | return None; |
| 135 | } |
| 136 | let bytes = line.as_bytes(); |
| 137 | if bytes.is_empty() || bytes.len() + 1 > MAX_LINE_BYTES { |
| 138 | return None; |
| 139 | } |
| 140 | ensure_dir(root).ok()?; |
| 141 | |
| 142 | let mut buf = Vec::with_capacity(bytes.len() + 1); |
| 143 | buf.extend_from_slice(bytes); |
| 144 | buf.push(b'\n'); |
| 145 | |
| 146 | let file = OpenOptions::new() |
| 147 | .create(true) |
| 148 | .append(true) |
| 149 | .open(path) |
| 150 | .ok()?; |
| 151 | secure(&file).ok()?; |
| 152 | // One `write(2)`, not `write_fmt` and not two calls: a split write is what |
| 153 | // a concurrent appender would interleave with. |
| 154 | (&file).write_all(&buf).ok()?; |
| 155 | file.sync_data().ok()?; |
| 156 | drop(file); |
| 157 | |
| 158 | enforce_ring(root, path); |
| 159 | Some(()) |
| 160 | } |
| 161 | |
| 162 | /// Append a line that is too large for one atomic `write(2)`, serialising |
| 163 | /// against other writers with the compaction lock instead. |
| 164 | /// |
| 165 | /// Only the dry-run sink uses this: a whole batch does not fit under |
| 166 | /// `PIPE_BUF`, and the flush path is neither the panic hook nor the signal |
| 167 | /// handler, so a **non-blocking** `try_write` is safe there. On contention the |
| 168 | /// batch is dropped, which is the same fail-open behavior as a failed POST. |
| 169 | pub fn append_locked(root: &Path, path: &Path, line: &str) -> Option<()> { |
| 170 | if tombstone_present(root) { |
| 171 | return None; |
| 172 | } |
| 173 | let bytes = line.as_bytes(); |
| 174 | if bytes.is_empty() || bytes.len() as u64 + 1 > MAX_BYTES { |
| 175 | return None; |
| 176 | } |
| 177 | ensure_dir(root).ok()?; |
| 178 | |
| 179 | let mut buf = Vec::with_capacity(bytes.len() + 1); |
| 180 | buf.extend_from_slice(bytes); |
| 181 | buf.push(b'\n'); |
| 182 | |
| 183 | let wrote = try_with_lock(root, || { |
| 184 | if tombstone_present(root) { |
| 185 | return Ok(false); |
| 186 | } |
| 187 | let file = OpenOptions::new() |
| 188 | .create(true) |
| 189 | .append(true) |
| 190 | .open(path) |
| 191 | .with_context(|| format!("failed to open {}", path.display()))?; |
| 192 | secure(&file)?; |
| 193 | (&file) |
| 194 | .write_all(&buf) |
| 195 | .with_context(|| format!("failed to append to {}", path.display()))?; |
| 196 | file.sync_data() |
| 197 | .with_context(|| format!("failed to sync {}", path.display()))?; |
| 198 | Ok(true) |
| 199 | }) |
| 200 | .ok() |
| 201 | .flatten() |
| 202 | .unwrap_or(false); |
| 203 | |
| 204 | if !wrote { |
| 205 | return None; |
| 206 | } |
| 207 | enforce_ring(root, path); |
| 208 | Some(()) |
| 209 | } |
| 210 | |
| 211 | /// Keep the newest [`MAX_EVENTS`] lines and at most [`MAX_BYTES`], under the |
| 212 | /// compaction lock. On lock contention this cycle is skipped: the next append |
| 213 | /// tries again, and the cap is a ceiling on disk footprint, not an invariant |
| 214 | /// that must hold at every instant. |
| 215 | fn enforce_ring(root: &Path, path: &Path) { |
| 216 | let Ok(meta) = fs::metadata(path) else { |
| 217 | return; |
| 218 | }; |
| 219 | let len = meta.len(); |
| 220 | if len < PROBE_BYTES { |
| 221 | return; |
| 222 | } |
| 223 | let Ok(contents) = fs::read_to_string(path) else { |
| 224 | return; |
| 225 | }; |
| 226 | let lines: Vec<&str> = contents.lines().filter(|l| !l.trim().is_empty()).collect(); |
| 227 | if lines.len() <= MAX_EVENTS && len <= MAX_BYTES { |
| 228 | return; |
| 229 | } |
| 230 | |
| 231 | let _ = try_with_lock(root, || { |
| 232 | let mut kept: Vec<&str> = lines |
| 233 | .iter() |
| 234 | .rev() |
| 235 | .take(MAX_EVENTS) |
| 236 | .rev() |
| 237 | .copied() |
| 238 | .collect::<Vec<_>>(); |
| 239 | // Byte ceiling second: drop from the oldest end until the survivors fit. |
| 240 | while kept.len() > 1 && byte_len(&kept) > MAX_BYTES { |
| 241 | kept.remove(0); |
| 242 | } |
| 243 | let mut body = kept.join("\n"); |
| 244 | if !body.is_empty() { |
| 245 | body.push('\n'); |
| 246 | } |
| 247 | rewrite(path, body.as_bytes()) |
| 248 | }); |
| 249 | } |
| 250 | |
| 251 | fn byte_len(lines: &[&str]) -> u64 { |
| 252 | lines.iter().map(|l| l.len() as u64 + 1).sum() |
| 253 | } |
| 254 | |
| 255 | /// Replace `path` atomically through a sibling temp file in the same directory. |
| 256 | fn rewrite(path: &Path, bytes: &[u8]) -> Result<()> { |
| 257 | let dir = path.parent().unwrap_or_else(|| Path::new(".")); |
| 258 | let mut tmp = tempfile::NamedTempFile::new_in(dir) |
| 259 | .with_context(|| format!("failed to stage a rewrite of {}", path.display()))?; |
| 260 | tmp.write_all(bytes) |
| 261 | .with_context(|| format!("failed to write a rewrite of {}", path.display()))?; |
| 262 | tmp.flush() |
| 263 | .with_context(|| format!("failed to flush a rewrite of {}", path.display()))?; |
| 264 | secure(tmp.as_file())?; |
| 265 | tmp.persist(path) |
| 266 | .map_err(|error| error.error) |
| 267 | .with_context(|| format!("failed to persist {}", path.display()))?; |
| 268 | Ok(()) |
| 269 | } |
| 270 | |
| 271 | /// Open (creating if needed) the sibling lock file. |
| 272 | fn open_lock(root: &Path) -> Result<File> { |
| 273 | ensure_dir(root)?; |
| 274 | let path = lock_path(root); |
| 275 | let file = OpenOptions::new() |
| 276 | .create(true) |
| 277 | .read(true) |
| 278 | .write(true) |
| 279 | // The file is only a lock handle; its contents are never read and |
| 280 | // truncating it would race other holders for no benefit. |
| 281 | .truncate(false) |
| 282 | .open(&path) |
| 283 | .with_context(|| format!("failed to open {}", path.display()))?; |
| 284 | secure(&file)?; |
| 285 | Ok(file) |
| 286 | } |
| 287 | |
| 288 | /// Run `operation` holding the exclusive compaction lock, **blocking**. |
| 289 | /// |
| 290 | /// Only the opt-out wipe uses this. It is not an exit path, so blocking is |
| 291 | /// fine there and nowhere else. |
| 292 | pub fn with_lock<T>(root: &Path, operation: impl FnOnce() -> Result<T>) -> Result<T> { |
| 293 | let file = open_lock(root)?; |
| 294 | let mut lock = fd_lock::RwLock::new(file); |
| 295 | let _guard = lock.write().context("failed to take the telemetry lock")?; |
| 296 | operation() |
| 297 | } |
| 298 | |
| 299 | /// Run `operation` holding the exclusive compaction lock if it is free. |
| 300 | /// |
| 301 | /// Returns `Ok(None)` when the lock is held elsewhere. Never blocks. |
| 302 | pub fn try_with_lock<T>(root: &Path, operation: impl FnOnce() -> Result<T>) -> Result<Option<T>> { |
| 303 | let file = open_lock(root)?; |
| 304 | let mut lock = fd_lock::RwLock::new(file); |
| 305 | match lock.try_write() { |
| 306 | Ok(_guard) => operation().map(Some), |
| 307 | Err(_) => Ok(None), |
| 308 | } |
| 309 | } |
| 310 | |
| 311 | /// Read every intact line from `path`, dropping a torn trailing line. |
| 312 | /// |
| 313 | /// `std::process::exit` on the signal path can truncate a concurrent write, so |
| 314 | /// the last line may be a partial JSON document. Skipping unparseable lines is |
| 315 | /// the whole tolerance: a drain must never fail because one record was cut. |
| 316 | #[must_use] |
| 317 | pub fn read_lines(path: &Path) -> Vec<String> { |
| 318 | let Ok(contents) = fs::read_to_string(path) else { |
| 319 | return Vec::new(); |
| 320 | }; |
| 321 | contents |
| 322 | .lines() |
| 323 | .filter(|line| !line.trim().is_empty()) |
| 324 | .map(str::to_string) |
| 325 | .collect() |
| 326 | } |
| 327 | |
| 328 | /// Take every buffered line and truncate the buffer, under the compaction lock. |
| 329 | /// |
| 330 | /// Returns an empty vector when the tombstone is present or the lock is held |
| 331 | /// elsewhere. Truncates rather than unlinks — `crates/tui/src/fleet/ledger.rs` |
| 332 | /// documents the rule: replacing the file leaves appenders holding the old |
| 333 | /// inode. |
| 334 | #[must_use] |
| 335 | pub fn drain(root: &Path) -> Vec<String> { |
| 336 | if tombstone_present(root) { |
| 337 | return Vec::new(); |
| 338 | } |
| 339 | let path = buffer_path(root); |
| 340 | let drained = try_with_lock(root, || { |
| 341 | // Re-check under the lock: a wipe may have landed between the check |
| 342 | // above and the acquisition. |
| 343 | if tombstone_present(root) { |
| 344 | return Ok(Vec::new()); |
| 345 | } |
| 346 | let lines = read_lines(&path); |
| 347 | if !lines.is_empty() { |
| 348 | truncate(&path)?; |
| 349 | } |
| 350 | Ok(lines) |
| 351 | }); |
| 352 | drained.ok().flatten().unwrap_or_default() |
| 353 | } |
| 354 | |
| 355 | /// Truncate a file to zero length, leaving the inode in place. A missing file |
| 356 | /// is not an error. |
| 357 | pub fn truncate(path: &Path) -> Result<()> { |
| 358 | if !path.exists() { |
| 359 | return Ok(()); |
| 360 | } |
| 361 | let file = OpenOptions::new() |
| 362 | .write(true) |
| 363 | .truncate(true) |
| 364 | .open(path) |
| 365 | .with_context(|| format!("failed to truncate {}", path.display()))?; |
| 366 | secure(&file)?; |
| 367 | Ok(()) |
| 368 | } |
| 369 | |
| 370 | /// Wipe every trace of collection, leaving a permanent tombstone. |
| 371 | /// |
| 372 | /// Order matters and is the whole of the guarantee: |
| 373 | /// |
| 374 | /// 1. take the blocking lock — this is not an exit path; |
| 375 | /// 2. write the tombstone **first**, and never remove it here; |
| 376 | /// 3. truncate `buffer.jsonl` and `dryrun.jsonl` — do **not** unlink them, and |
| 377 | /// never unlink the lock file; |
| 378 | /// 4. remove `install_id.json` and `state.json`. |
| 379 | /// |
| 380 | /// If any step after the tombstone fails, the error is returned and the caller |
| 381 | /// logs it — but the tombstone alone already makes the buffer permanently |
| 382 | /// undrainable, so a failed wipe fails **closed**. |
| 383 | pub fn wipe(root: &Path) -> Result<()> { |
| 384 | with_lock(root, || { |
| 385 | let tombstone = tombstone_path(root); |
| 386 | let file = OpenOptions::new() |
| 387 | .create(true) |
| 388 | .write(true) |
| 389 | .truncate(true) |
| 390 | .open(&tombstone) |
| 391 | .with_context(|| format!("failed to write {}", tombstone.display()))?; |
| 392 | secure(&file)?; |
| 393 | drop(file); |
| 394 | |
| 395 | let mut failure: Option<anyhow::Error> = None; |
| 396 | for path in [buffer_path(root), dryrun_path(root)] { |
| 397 | if let Err(error) = truncate(&path) { |
| 398 | failure.get_or_insert(error); |
| 399 | } |
| 400 | } |
| 401 | for path in [install_id_path(root), state_path(root)] { |
| 402 | if path.exists() |
| 403 | && let Err(error) = fs::remove_file(&path) |
| 404 | { |
| 405 | failure.get_or_insert( |
| 406 | anyhow::Error::new(error) |
| 407 | .context(format!("failed to remove {}", path.display())), |
| 408 | ); |
| 409 | } |
| 410 | } |
| 411 | match failure { |
| 412 | Some(error) => Err(error), |
| 413 | None => Ok(()), |
| 414 | } |
| 415 | }) |
| 416 | } |
| 417 | |
| 418 | /// Clear the tombstone and drop anything buffered before consent. |
| 419 | /// |
| 420 | /// Called by `init` on every arming. No event recorded before the user said yes |
| 421 | /// can be in the batch that follows it — a stale buffer left by an earlier |
| 422 | /// consenting run, or by a bug, is not evidence of this user's consent. |
| 423 | pub fn arm(root: &Path) -> Result<()> { |
| 424 | ensure_dir(root)?; |
| 425 | with_lock(root, || { |
| 426 | let tombstone = tombstone_path(root); |
| 427 | if tombstone.exists() { |
| 428 | fs::remove_file(&tombstone) |
| 429 | .with_context(|| format!("failed to remove {}", tombstone.display()))?; |
| 430 | } |
| 431 | truncate(&buffer_path(root)) |
| 432 | }) |
| 433 | } |
| 434 |