- 1
//! Durable attention layer (docs/design/29-personal-os.md P6): every - 2
//! `gateway::deliver` push also records an append-only entry at - 3
//! `<home>/inbox.jsonl`, so unattended signals survive with no chat - 4
//! channel configured. Read state is ack tombstone lines — nothing is - 5
//! ever deleted or rewritten (invariant-2 adjacent). Retention note: - 6
//! compaction of aged entries comes later; until then the ledger grows - 7
//! monotonically and reads stay bounded by [`MAX_SCAN`]. - 8
- 9
use std::collections::VecDeque; - 10
use std::io::{BufRead, BufReader, Write}; - 11
use std::path::{Path, PathBuf}; - 12
- 13
use chrono::{DateTime, Utc}; - 14
use serde::{Deserialize, Serialize}; - 15
- 16
/// Upper bound on raw lines examined per read. The ledger never shrinks - 17
/// today, so unbounded scans would eventually dominate surface latency; - 18
/// the newest window wins. - 19
pub const MAX_SCAN: usize = 10_000; - 20
- 21
#[derive(Debug, thiserror::Error)] - 22
pub enum InboxError { - 23
#[error("io error on {path}: {source}")] - 24
Io { - 25
path: PathBuf, - 26
source: std::io::Error, - 27
}, - 28
#[error("serialize inbox entry: {0}")] - 29
Serialize(String), - 30
} - 31
- 32
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] - 33
#[serde(rename_all = "snake_case")] - 34
pub enum Kind { - 35
TaskSummary, - 36
ApprovalPending, - 37
ApprovalDenied, - 38
BudgetAlert, - 39
Digest, - 40
Heartbeat, - 41
ProposalOpened, - 42
/// A scheduled task was due and could not start; the body says why and - 43
/// what to do about it. - 44
RoutineFailed, - 45
} - 46
- 47
impl Kind { - 48
/// The serde snake_case tag, also used in id derivation so ids agree - 49
/// with whatever the JSONL on disk shows. - 50
fn tag(self) -> String { - 51
// Unit-variant serialization cannot fail; the error path exists - 52
// only to satisfy the no-unwrap rule. - 53
serde_json::to_string(&self) - 54
.map(|s| s.trim_matches('"').to_string()) - 55
.unwrap_or_else(|_| format!("{self:?}")) - 56
} - 57
} - 58
- 59
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] - 60
pub struct Entry { - 61
pub id: String, - 62
pub ts: DateTime<Utc>, - 63
pub kind: Kind, - 64
pub title: String, - 65
pub body: String, - 66
#[serde(default)] - 67
pub session_id: Option<String>, - 68
#[serde(default)] - 69
pub task_id: Option<String>, - 70
#[serde(default)] - 71
pub result_id: Option<String>, - 72
#[serde(default)] - 73
pub dedupe_key: Option<String>, - 74
} - 75
- 76
#[derive(Serialize, Deserialize)] - 77
struct AckLine { - 78
ack_of: String, - 79
ts: DateTime<Utc>, - 80
} - 81
- 82
/// Entries plus a count of corrupt lines skipped while scanning — silent - 83
/// for the caller's eye, but never invisible. - 84
#[derive(Debug, Clone, Default, PartialEq)] - 85
pub struct Scan { - 86
pub entries: Vec<Entry>, - 87
pub corrupt: usize, - 88
} - 89
- 90
#[derive(Default)] - 91
struct Window { - 92
entries: Vec<Entry>, - 93
acked: Vec<String>, - 94
corrupt: usize, - 95
} - 96
- 97
/// FNV-1a 64-bit over `{ts}|{kind}|{title}`: fixed across toolchains like - 98
/// every other house id. Hash collisions are tolerated as house style — - 99
/// an id colliding merely merges two rows' read state. - 100
fn entry_id(ts: DateTime<Utc>, kind: Kind, title: &str) -> String { - 101
fnv1a(format!("{}|{}|{}", ts.to_rfc3339(), kind.tag(), title).as_bytes()) - 102
} - 103
- 104
fn fnv1a(data: &[u8]) -> String { - 105
let mut h: u64 = 0xcbf2_9ce4_8422_2325; - 106
for b in data { - 107
h ^= u64::from(*b); - 108
h = h.wrapping_mul(0x0000_0100_0000_01b3); - 109
} - 110
format!("{h:016x}") - 111
} - 112
- 113
pub fn inbox_path(home: &Path) -> PathBuf { - 114
home.join("inbox.jsonl") - 115
} - 116
- 117
/// Append one notification and return it. One formatted buffer + ONE - 118
/// write_all: O_APPEND makes a single write atomic, whereas `writeln!` - 119
/// emits several syscalls that concurrent deliverers could interleave - 120
/// mid-line (same discipline as gateway.rs deliver_log). - 121
pub fn record( - 122
home: &Path, - 123
kind: Kind, - 124
title: &str, - 125
body: &str, - 126
session_id: Option<&str>, - 127
task_id: Option<&str>, - 128
) -> Result<Entry, InboxError> { - 129
record_with_result_and_key(home, kind, title, body, session_id, task_id, None, None) - 130
} - 131
- 132
/// Append a notification linked to one immutable presentation result. - 133
pub fn record_with_result( - 134
home: &Path, - 135
kind: Kind, - 136
title: &str, - 137
body: &str, - 138
session_id: Option<&str>, - 139
task_id: Option<&str>, - 140
result_id: Option<&str>, - 141
) -> Result<Entry, InboxError> { - 142
record_with_result_and_key( - 143
home, kind, title, body, session_id, task_id, result_id, None, - 144
) - 145
} - 146
- 147
/// Append a notification only once for a stable source/destination identity. - 148
/// The check and append are serialized by the inbox file lock, so delivery - 149
/// retries cannot create duplicate unread entries. - 150
#[allow(clippy::too_many_arguments)] - 151
pub fn record_with_result_and_key( - 152
home: &Path, - 153
kind: Kind, - 154
title: &str, - 155
body: &str, - 156
session_id: Option<&str>, - 157
task_id: Option<&str>, - 158
result_id: Option<&str>, - 159
dedupe_key: Option<&str>, - 160
) -> Result<Entry, InboxError> { - 161
let _dedupe_lock = if dedupe_key.is_some() { - 162
Some(acquire_dedupe_lock(home)?) - 163
} else { - 164
None - 165
}; - 166
if let Some(key) = dedupe_key - 167
&& let Some(existing) = list(home, MAX_SCAN) - 168
.into_iter() - 169
.find(|entry| entry.dedupe_key.as_deref() == Some(key)) - 170
{ - 171
return Ok(existing); - 172
} - 173
let path = inbox_path(home); - 174
let ts = Utc::now(); - 175
let entry = Entry { - 176
id: entry_id(ts, kind, title), - 177
ts, - 178
kind, - 179
title: title.to_string(), - 180
body: body.to_string(), - 181
session_id: session_id.map(str::to_string), - 182
task_id: task_id.map(str::to_string), - 183
result_id: result_id.map(str::to_string), - 184
dedupe_key: dedupe_key.map(str::to_string), - 185
}; - 186
let line = serde_json::to_string(&entry).map_err(|e| InboxError::Serialize(e.to_string()))?; - 187
append_line(&path, &line)?; - 188
Ok(entry) - 189
} - 190
- 191
struct DedupeLock { - 192
path: PathBuf, - 193
} - 194
- 195
impl Drop for DedupeLock { - 196
fn drop(&mut self) { - 197
let _ = std::fs::remove_file(&self.path); - 198
} - 199
} - 200
- 201
fn acquire_dedupe_lock(home: &Path) -> Result<DedupeLock, InboxError> { - 202
let path = home.join("inbox.dedupe.lock"); - 203
std::fs::create_dir_all(home).map_err(|source| InboxError::Io { - 204
path: home.to_path_buf(), - 205
source, - 206
})?; - 207
for _ in 0..200 { - 208
match std::fs::OpenOptions::new() - 209
.write(true) - 210
.create_new(true) - 211
.open(&path) - 212
{ - 213
Ok(_) => return Ok(DedupeLock { path }), - 214
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => { - 215
std::thread::sleep(std::time::Duration::from_millis(5)); - 216
} - 217
Err(source) => return Err(InboxError::Io { path, source }), - 218
} - 219
} - 220
Err(InboxError::Io { - 221
path, - 222
source: std::io::Error::new(std::io::ErrorKind::TimedOut, "inbox dedupe lock timed out"), - 223
}) - 224
} - 225
- 226
/// Mark `id` read by appending a tombstone. Idempotent via - 227
/// read-before-write; returns false when already acked. - 228
pub fn ack(home: &Path, id: &str) -> Result<bool, InboxError> { - 229
let path = inbox_path(home); - 230
if scan_window(&path).acked.iter().any(|a| a == id) { - 231
return Ok(false); - 232
} - 233
let line = serde_json::to_string(&AckLine { - 234
ack_of: id.to_string(), - 235
ts: Utc::now(), - 236
}) - 237
.map_err(|e| InboxError::Serialize(e.to_string()))?; - 238
append_line(&path, &line)?; - 239
Ok(true) - 240
} - 241
- 242
pub fn unread(home: &Path) -> Vec<Entry> { - 243
unread_scanned(home).entries - 244
} - 245
- 246
/// Same as [`unread`] but reports how many corrupt lines were skipped. - 247
pub fn unread_scanned(home: &Path) -> Scan { - 248
let window = scan_window(&inbox_path(home)); - 249
let mut entries = window.entries; - 250
entries.retain(|e| !window.acked.contains(&e.id)); - 251
entries.reverse(); - 252
Scan { - 253
entries, - 254
corrupt: window.corrupt, - 255
} - 256
} - 257
- 258
/// All entries (acked included), newest first, capped at `limit`. - 259
pub fn list(home: &Path, limit: usize) -> Vec<Entry> { - 260
list_scanned(home, limit).entries - 261
} - 262
- 263
pub fn list_scanned(home: &Path, limit: usize) -> Scan { - 264
let mut window = scan_window(&inbox_path(home)); - 265
window.entries.reverse(); - 266
window.entries.truncate(limit); - 267
Scan { - 268
entries: window.entries, - 269
corrupt: window.corrupt, - 270
} - 271
} - 272
- 273
pub fn unread_count(home: &Path) -> usize { - 274
unread(home).len() - 275
} - 276
- 277
fn append_line(path: &Path, line: &str) -> Result<(), InboxError> { - 278
if let Some(parent) = path.parent() { - 279
std::fs::create_dir_all(parent).map_err(|e| InboxError::Io { - 280
path: parent.to_path_buf(), - 281
source: e, - 282
})?; - 283
} - 284
let mut buf = String::with_capacity(line.len() + 1); - 285
buf.push_str(line); - 286
buf.push('\n'); - 287
let mut f = std::fs::OpenOptions::new() - 288
.create(true) - 289
.append(true) - 290
.open(path) - 291
.map_err(|e| InboxError::Io { - 292
path: path.to_path_buf(), - 293
source: e, - 294
})?; - 295
f.write_all(buf.as_bytes()).map_err(|e| InboxError::Io { - 296
path: path.to_path_buf(), - 297
source: e, - 298
}) - 299
} - 300
- 301
/// Last MAX_SCAN raw lines (oldest→newest within the window); corrupt - 302
/// lines are counted, not fatal. Missing file reads as empty. Entries and - 303
/// ack tombstones come from the same window so read state stays - 304
/// consistent with what a surface can see. - 305
fn scan_window(path: &Path) -> Window { - 306
let mut window: VecDeque<String> = VecDeque::with_capacity(64); - 307
let mut torn_reads = 0usize; - 308
if let Ok(f) = std::fs::File::open(path) { - 309
for line in BufReader::new(f).lines() { - 310
match line { - 311
Ok(l) => { - 312
if window.len() == MAX_SCAN { - 313
window.pop_front(); - 314
} - 315
window.push_back(l); - 316
} - 317
Err(_) => { - 318
// A torn read mid-line cannot continue by contract of - 319
// BufRead::lines; count it and stop. - 320
torn_reads += 1; - 321
break; - 322
} - 323
} - 324
} - 325
} - 326
let mut out = Window::default(); - 327
out.corrupt += torn_reads; - 328
for line in &window { - 329
if line.trim().is_empty() { - 330
continue; - 331
} - 332
if let Ok(e) = serde_json::from_str::<Entry>(line) { - 333
out.entries.push(e); - 334
} else if let Ok(a) = serde_json::from_str::<AckLine>(line) { - 335
out.acked.push(a.ack_of); - 336
} else { - 337
out.corrupt += 1; - 338
} - 339
} - 340
out - 341
} - 342
- 343
#[cfg(test)] - 344
mod tests { - 345
#![allow(clippy::unwrap_used, clippy::expect_used)] - 346
- 347
use super::*; - 348
- 349
#[test] - 350
fn record_unread_roundtrip_newest_first() { - 351
let dir = tempfile::tempdir().unwrap(); - 352
let home = dir.path(); - 353
- 354
assert!(unread(home).is_empty()); - 355
assert_eq!(unread_count(home), 0); - 356
- 357
let a = record( - 358
home, - 359
Kind::TaskSummary, - 360
"first", - 361
"body a", - 362
Some("sess-1"), - 363
None, - 364
) - 365
.unwrap(); - 366
let b = record( - 367
home, - 368
Kind::BudgetAlert, - 369
"second", - 370
"body b", - 371
None, - 372
Some("task-9"), - 373
) - 374
.unwrap(); - 375
- 376
assert_eq!(a.session_id.as_deref(), Some("sess-1")); - 377
assert_eq!(b.task_id.as_deref(), Some("task-9")); - 378
assert_ne!(a.id, b.id); - 379
- 380
let items = unread(home); - 381
assert_eq!(items.len(), 2); - 382
assert_eq!(items[0], b); - 383
assert_eq!(items[1], a); - 384
assert_eq!(unread_count(home), 2); - 385
assert!(list(home, 10).contains(&a)); - 386
} - 387
- 388
#[test] - 389
fn ack_is_idempotent_and_shrinks_unread_once() { - 390
let dir = tempfile::tempdir().unwrap(); - 391
let home = dir.path(); - 392
let e = record(home, Kind::Digest, "daily", "", None, None).unwrap(); - 393
let other = record(home, Kind::Heartbeat, "beat", "", None, None).unwrap(); - 394
- 395
assert!(ack(home, &e.id).unwrap()); - 396
assert!(!ack(home, &e.id).unwrap()); - 397
- 398
let left = unread(home); - 399
assert_eq!(left, vec![other]); - 400
assert_eq!(unread_count(home), 1); - 401
- 402
assert!(!ack(home, &e.id).unwrap()); - 403
assert_eq!(unread_count(home), 1); - 404
} - 405
- 406
#[test] - 407
fn acked_excluded_from_unread_but_present_in_list() { - 408
let dir = tempfile::tempdir().unwrap(); - 409
let home = dir.path(); - 410
let e = record( - 411
home, - 412
Kind::ApprovalPending, - 413
"gate", - 414
"wants write", - 415
None, - 416
None, - 417
) - 418
.unwrap(); - 419
ack(home, &e.id).unwrap(); - 420
- 421
assert!(unread(home).is_empty()); - 422
let all = list(home, 10); - 423
assert_eq!(all.len(), 1); - 424
assert_eq!(all[0].id, e.id); - 425
assert_eq!(all[0].kind, Kind::ApprovalPending); - 426
} - 427
- 428
#[test] - 429
fn corrupt_middle_line_is_skipped_and_counted() { - 430
let dir = tempfile::tempdir().unwrap(); - 431
let home = dir.path(); - 432
let first = record(home, Kind::Digest, "one", "", None, None).unwrap(); - 433
let mut raw = std::fs::read_to_string(inbox_path(home)).unwrap(); - 434
raw.push_str("{\"id\": torn line no json\n"); - 435
std::fs::write(inbox_path(home), &raw).unwrap(); - 436
let last = record(home, Kind::BudgetAlert, "two", "", None, None).unwrap(); - 437
- 438
let scan = unread_scanned(home); - 439
assert_eq!(scan.corrupt, 1); - 440
assert_eq!(scan.entries, vec![last.clone(), first.clone()]); - 441
let listed = list_scanned(home, 10); - 442
assert_eq!(listed.corrupt, 1); - 443
// Ack still works with garbage in the middle. - 444
assert!(ack(home, &first.id).unwrap()); - 445
assert_eq!(unread_scanned(home).entries, vec![last]); - 446
} - 447
- 448
#[test] - 449
fn interleaved_writers_keep_line_integrity_and_order() { - 450
let dir = tempfile::tempdir().unwrap(); - 451
let home = dir.path(); - 452
let home_a = home.to_path_buf(); - 453
let home_b = home.to_path_buf(); - 454
- 455
let writer = |home: PathBuf, tag: &str| -> Vec<String> { - 456
(0..25) - 457
.map(|i| { - 458
record( - 459
home.as_path(), - 460
Kind::TaskSummary, - 461
&format!("{tag}-{i}"), - 462
"x", - 463
None, - 464
None, - 465
) - 466
.unwrap() - 467
.id - 468
}) - 469
.collect() - 470
}; - 471
let h1 = std::thread::spawn(move || writer(home_a, "a")); - 472
let h2 = std::thread::spawn(move || writer(home_b, "b")); - 473
let ids_a = h1.join().unwrap(); - 474
let ids_b = h2.join().unwrap(); - 475
- 476
let on_disk = unread_scanned(home); - 477
assert_eq!( - 478
on_disk.corrupt, 0, - 479
"no torn lines under O_APPEND single writes" - 480
); - 481
assert_eq!(on_disk.entries.len(), 50); - 482
- 483
// Each writer's own subsequence must appear in its recorded order. - 484
// `on_disk.entries` is newest-first, so later writes sit earlier. - 485
let position_of = |id: &str| { - 486
on_disk - 487
.entries - 488
.iter() - 489
.position(|e| e.id == id) - 490
.unwrap_or(usize::MAX) - 491
}; - 492
for pair in ids_a.windows(2) { - 493
assert!( - 494
position_of(&pair[1]) < position_of(&pair[0]), - 495
"writer A reordered" - 496
); - 497
} - 498
for pair in ids_b.windows(2) { - 499
assert!( - 500
position_of(&pair[1]) < position_of(&pair[0]), - 501
"writer B reordered" - 502
); - 503
} - 504
} - 505
- 506
#[test] - 507
fn ids_are_deterministic_fnv_over_ts_kind_title() { - 508
let fixed: DateTime<Utc> = "2026-08-25T00:00:00Z".parse().unwrap(); - 509
assert_eq!( - 510
entry_id(fixed, Kind::Digest, "weekly digest"), - 511
"76391567a848733e" - 512
); - 513
assert_eq!( - 514
entry_id(fixed, Kind::Digest, "weekly digest"), - 515
entry_id(fixed, Kind::Digest, "weekly digest") - 516
); - 517
assert_ne!( - 518
entry_id(fixed, Kind::Digest, "a"), - 519
entry_id(fixed, Kind::Digest, "b") - 520
); - 521
assert_ne!( - 522
entry_id(fixed, Kind::Heartbeat, "a"), - 523
entry_id(fixed, Kind::Digest, "a") - 524
); - 525
- 526
// The serde snake_case tags feed the id and the wire format. - 527
let line = serde_json::to_string(&Kind::ProposalOpened).unwrap(); - 528
assert_eq!(line, "\"proposal_opened\""); - 529
let back: Kind = serde_json::from_str(&line).unwrap(); - 530
assert_eq!(back, Kind::ProposalOpened); - 531
} - 532
- 533
#[test] - 534
fn result_delivery_deduplicates_by_source_key() { - 535
let dir = tempfile::tempdir().unwrap(); - 536
let first = record_with_result_and_key( - 537
dir.path(), - 538
Kind::TaskSummary, - 539
"finished", - 540
"answer", - 541
Some("session"), - 542
Some("task"), - 543
Some("result-1"), - 544
Some("surface:chat|result-1|0"), - 545
) - 546
.unwrap(); - 547
let second = record_with_result_and_key( - 548
dir.path(), - 549
Kind::TaskSummary, - 550
"finished", - 551
"answer changed", - 552
Some("session"), - 553
Some("task"), - 554
Some("result-1"), - 555
Some("surface:chat|result-1|0"), - 556
) - 557
.unwrap(); - 558
assert_eq!(first.id, second.id); - 559
assert_eq!(list(dir.path(), 10).len(), 1); - 560
} - 561
} - 562
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.