trusty-common 0.54.3

Shared utilities and provider-agnostic streaming chat (ChatProvider, OllamaProvider, OpenRouter, tool-use) for trusty-* projects
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
//! Durable trail of the drawer deletions that maintenance paths make (#8732).
//!
//! Why: in #8729, 70 drawers left three live palaces within minutes and nothing
//! recorded why. The dream and purge summaries logged at `info`, below the
//! daemon's default `warn` filter, and none of them named a drawer. An operator
//! could not tell which drawer a removed one duplicated, or at what score.
//! What: every drawer a maintenance path deletes is appended as one JSON line to
//! `<palace data_dir>/maintenance_deletions.jsonl`: time, palace, drawer id,
//! reason, the surviving drawer id and cosine score where one exists, and the
//! pid of the deleting process. #8729: a deletion made through
//! [`PalaceHandle::forget_for_maintenance`] also copies the drawer's content,
//! room, tags, importance and creation time, so it can be recreated. Each deleting pass also logs one summary at
//! `warn`. #8729 (owner ruling): every removal is also logged on its own `warn`
//! line naming the palace, drawer id and reason, so the daemon log alone shows
//! which drawers went and why. `trusty-memory palace deletions` reads the file
//! back.
//! User-initiated deletions (`memory_forget`, the HTTP/UDS drawer delete, and so
//! `palace reclaim --apply`) call [`PalaceHandle::forget`]. #9283: each is
//! recorded as [`DeletionReason::UserForget`] carrying the drawer's content
//! hash and NO content copy, so the doctor drawer-count check can explain the
//! drop while the forgotten text stays gone. One exception (#9172): forgetting
//! a drawer this journal names as a dedup survivor is recorded as
//! [`DeletionReason::ForgetOfMergedSurvivor`] with its copy, because merged-in
//! text leaves with it.
//!
//! A failed append does not undo or block the deletion. The full record is
//! logged at `error` instead, which the default filter keeps. A palace with no
//! data dir (in-memory) logs the record at `warn`.
//! Test: `maintenance_log_tests::dream_dedup_records_the_removed_and_surviving_drawer`,
//! `maintenance_log_tests::every_maintenance_removal_logs_its_id_and_reason`,
//! `maintenance_log_tests::a_failed_record_write_logs_the_record_and_still_deletes`.

use crate::memory_core::content_hash::{ContentHash, memory_content_hash};
use crate::memory_core::palace::{Drawer, PalaceId};
use crate::memory_core::retrieval::{ForgetOutcome, PalaceHandle};
use anyhow::{Context, Result};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use std::io::{BufRead, Write};
use std::path::{Path, PathBuf};
use uuid::Uuid;

/// File name of the per-palace deletion journal, inside the palace data dir.
pub const MAINTENANCE_LOG_FILENAME: &str = "maintenance_deletions.jsonl";
/// File name the journal is rotated to once it passes [`ROTATE_AT_BYTES`].
pub const MAINTENANCE_LOG_ROTATED_FILENAME: &str = "maintenance_deletions.1.jsonl";
/// Journal size that triggers one rotation. #8729: a record carries the
/// removed drawer's content, so how many records fit depends on drawer size.
pub(crate) const ROTATE_AT_BYTES: u64 = 4 * 1024 * 1024;

/// Which maintenance path deleted a drawer.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum DeletionReason {
    /// Dream dedup merged the drawer into a near-duplicate and forgot it.
    DreamDedup,
    /// Dream content prune matched the blocklist or the word-count floor.
    DreamContentPrune,
    /// Dream prune: decayed importance at the floor and older than 30 days.
    DreamPrune,
    /// Room consolidation evicted an original a canonical drawer superseded.
    SemanticConsolidation,
    /// `PalaceHandle::purge_expired` reclaimed a drawer past its TTL.
    ExpiredPurge,
    /// The palace-open sweep reclaimed a drawer past its TTL.
    ExpiredPurgeAtOpen,
    /// #9172: a user forget removed a drawer a dedup merge had kept.
    ForgetOfMergedSurvivor,
    /// #9283: a user forget (MCP, HTTP/UDS delete, reclaim). Hash only, no copy.
    UserForget,
}

impl DeletionReason {
    /// The snake_case name written to the journal.
    pub fn as_str(self) -> &'static str {
        match self {
            Self::DreamDedup => "dream_dedup",
            Self::DreamContentPrune => "dream_content_prune",
            Self::DreamPrune => "dream_prune",
            Self::SemanticConsolidation => "semantic_consolidation",
            Self::ExpiredPurge => "expired_purge",
            Self::ExpiredPurgeAtOpen => "expired_purge_at_open",
            Self::ForgetOfMergedSurvivor => "forget_of_merged_survivor",
            Self::UserForget => "user_forget",
        }
    }
}

impl std::fmt::Display for DeletionReason {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        f.write_str(self.as_str())
    }
}

/// One journal line: a drawer a maintenance path deleted.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct MaintenanceDeletion {
    /// When the deletion was recorded (just after it committed).
    pub at: DateTime<Utc>,
    /// Palace id.
    pub palace: String,
    /// The deleted drawer.
    pub drawer_id: Uuid,
    /// The maintenance path that deleted it.
    pub reason: DeletionReason,
    /// The drawer that absorbed or superseded it (dedup, consolidation).
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub survivor_id: Option<Uuid>,
    /// Cosine similarity between the two drawers (dedup only).
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub score: Option<f32>,
    /// Pid of the process that deleted it; several processes may write one palace.
    pub pid: u32,
    /// #8729: a copy of the removed drawer, so it can be re-remembered.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub drawer: Option<RemovedDrawer>,
    /// #9283: the removed drawer's content hash. A user-forget record carries
    /// this in place of a content copy.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub content_hash: Option<ContentHash>,
}

/// What a journal record keeps of a removed drawer (#8729).
///
/// Why: an id alone names a lost drawer but cannot bring it back.
/// What: the fields `remember` needs to recreate the drawer.
/// Test: `maintenance_log_tests::every_dream_removal_journals_a_recoverable_copy`.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RemovedDrawer {
    /// The drawer body, verbatim.
    pub content: String,
    /// Room the drawer lived in.
    pub room_id: Uuid,
    /// Tags, in stored order.
    pub tags: Vec<String>,
    /// Stored importance.
    pub importance: f32,
    /// When the drawer was first written.
    pub created_at: DateTime<Utc>,
    /// Tier C slot the drawer held, if any.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub fact_key: Option<String>,
}

impl From<&Drawer> for RemovedDrawer {
    fn from(d: &Drawer) -> Self {
        Self {
            content: d.content().to_string(),
            room_id: d.room_id,
            tags: d.tags.clone(),
            importance: d.importance,
            created_at: d.created_at,
            fact_key: d.fact_key.clone(),
        }
    }
}

impl MaintenanceDeletion {
    /// A record stamped now with this process's pid and no survivor.
    pub fn new(palace: &PalaceId, drawer_id: Uuid, reason: DeletionReason) -> Self {
        Self {
            at: Utc::now(),
            palace: palace.as_str().to_string(),
            drawer_id,
            reason,
            survivor_id: None,
            score: None,
            pid: std::process::id(),
            drawer: None,
            content_hash: None,
        }
    }

    /// #9283: attach the drawer's content hash, never its content.
    pub fn with_content_hash_of(mut self, drawer: &Drawer) -> Self {
        let stored = drawer.content_hash();
        // A legacy drawer read back with no digest gets one computed here.
        self.content_hash = Some(if stored.is_unset() {
            memory_content_hash(drawer.content())
        } else {
            stored
        });
        self
    }

    /// #8729: attach a recoverable copy of the removed drawer.
    pub fn with_drawer(mut self, drawer: &Drawer) -> Self {
        self.drawer = Some(RemovedDrawer::from(drawer));
        self
    }

    /// Attach the surviving drawer and, for dedup, the similarity score.
    pub fn with_survivor(mut self, survivor_id: Uuid, score: Option<f32>) -> Self {
        self.survivor_id = Some(survivor_id);
        self.score = score;
        self
    }
}

/// Where a record ended up.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RecordOutcome {
    /// Appended to the palace's journal file.
    Journaled,
    /// The journal was unavailable; the record went to the log only.
    LoggedOnly,
}

/// Path of the live journal for a palace data dir.
pub fn journal_path(data_dir: &Path) -> PathBuf {
    data_dir.join(MAINTENANCE_LOG_FILENAME)
}

/// How long an append waits for the journal lock before appending unlocked.
const JOURNAL_LOCK_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);

/// Append `rec` to the journal in `data_dir`, rotating once past `rotate_at`.
///
/// Why: up to ~16 processes append to one palace's journal (#8733). Two that
/// both saw the full size each renamed the live file over `.1`, and the second
/// rename discarded the whole previous generation.
/// What: the size check, the rename over the single `.1` generation, and the
/// append all run under the `file_lock` sidecar lock
/// (`maintenance_deletions.jsonl.lock`), a cross-process `flock`. The size is
/// read after the lock is held, so a writer that waited does not rotate a file
/// another writer just started. Each record is one `write_all` on an
/// `O_APPEND` handle. A lock that cannot be taken within
/// [`JOURNAL_LOCK_TIMEOUT`], or a failed rotation, falls back to a plain append
/// to the live file, so the record is kept and the journal only grows past its
/// bound.
/// Test: `maintenance_log_tests::concurrent_appends_across_the_rotation_boundary_lose_no_record`,
/// `maintenance_log_tests::a_lock_or_rotation_failure_still_appends_the_record`.
pub(crate) fn append(data_dir: &Path, rec: &MaintenanceDeletion, rotate_at: u64) -> Result<()> {
    let path = journal_path(data_dir);
    let mut line = serde_json::to_string(rec).context("serialize maintenance deletion")?;
    line.push('\n');
    let locked = crate::file_lock::with_exclusive_lock_timeout(&path, JOURNAL_LOCK_TIMEOUT, || {
        if let Err(e) = rotate_if_due(data_dir, &path, rotate_at) {
            tracing::warn!(palace = %rec.palace, "#8732: journal rotation failed; appending to the live file: {e:#}");
        }
        append_line(&path, &line)
    });
    match locked {
        Ok(appended) => appended,
        Err(e) => {
            tracing::warn!(palace = %rec.palace, "#8732: journal lock unavailable; appending without rotation: {e}");
            append_line(&path, &line)
        }
    }
}

/// Rename the live journal over `.1` when it has reached `rotate_at`.
fn rotate_if_due(data_dir: &Path, path: &Path, rotate_at: u64) -> Result<()> {
    if let Ok(meta) = std::fs::metadata(path)
        && meta.is_file()
        && meta.len() >= rotate_at
    {
        std::fs::rename(path, data_dir.join(MAINTENANCE_LOG_ROTATED_FILENAME))
            .with_context(|| format!("rotate {}", path.display()))?;
    }
    Ok(())
}

/// Write `line` to `path` in one `O_APPEND` write, creating the file if needed.
fn append_line(path: &Path, line: &str) -> Result<()> {
    let mut file = std::fs::OpenOptions::new()
        .create(true)
        .append(true)
        .open(path)
        .with_context(|| format!("open {}", path.display()))?;
    file.write_all(line.as_bytes())
        .with_context(|| format!("append {}", path.display()))
}

/// Record one maintenance deletion; never drops it silently.
///
/// Why/What: see the module doc. Appends to the journal when the palace has a
/// data dir, and logs the removal at `warn` (#8729). When the append fails, the
/// whole record goes to the log at `error`; with no data dir it goes at `warn`.
/// #8729: those two lines carry the record as its journal JSON, `drawer` copy
/// included, so the copy survives a failed append. Every arm reaches the
/// daemon's log at its default filter.
/// Test: `maintenance_log_tests::every_maintenance_removal_logs_its_id_and_reason`,
/// `maintenance_log_tests::a_failed_record_write_logs_the_record_and_still_deletes`.
pub fn record(data_dir: Option<&Path>, rec: &MaintenanceDeletion) -> RecordOutcome {
    // #8729: the stand-in record keeps the drawer copy, not only the ids.
    let as_json =
        || serde_json::to_string(rec).unwrap_or_else(|e| format!("<record not serializable: {e}>"));
    let Some(dir) = data_dir else {
        tracing::warn!(
            palace = %rec.palace, drawer_id = %rec.drawer_id, reason = %rec.reason,
            survivor_id = ?rec.survivor_id, score = ?rec.score, record = %as_json(),
            "#8732: maintenance deletion (palace has no data dir; this line is the record)"
        );
        return RecordOutcome::LoggedOnly;
    };
    match append(dir, rec, ROTATE_AT_BYTES) {
        Ok(()) => {
            // #8729: each removal is logged with its id and reason, not only
            // counted in the pass summary. #9283: a user forget is already
            // logged once by its caller (`log_user_forget`); no second line.
            if rec.reason != DeletionReason::UserForget {
                tracing::warn!(
                    palace = %rec.palace, drawer_id = %rec.drawer_id, reason = %rec.reason,
                    survivor_id = ?rec.survivor_id, score = ?rec.score,
                    "#8729: maintenance removed drawer {} ({})", rec.drawer_id, rec.reason
                );
            }
            RecordOutcome::Journaled
        }
        Err(e) => {
            tracing::error!(
                palace = %rec.palace, drawer_id = %rec.drawer_id, reason = %rec.reason,
                survivor_id = ?rec.survivor_id, score = ?rec.score, record = %as_json(),
                "#8732: maintenance deletion journal write failed; the drawer is \
                 deleted and this line is the record: {e:#}"
            );
            RecordOutcome::LoggedOnly
        }
    }
}

/// Log one `warn` summary for a pass that deleted `count` drawers.
pub fn warn_removed(palace: &PalaceId, pass: &str, count: usize) {
    if count == 0 {
        return;
    }
    tracing::warn!(
        palace = %palace, pass, count,
        "#8732: {pass} removed {count} drawer(s); per-drawer record: \
         `trusty-memory palace deletions {palace}`"
    );
}

/// The journal read back, oldest first.
#[derive(Debug, Default)]
pub struct JournalContents {
    /// Every parseable record, the rotated generation before the live one.
    pub records: Vec<MaintenanceDeletion>,
    /// Lines that did not parse (a torn write, or a newer schema).
    pub malformed: usize,
}

/// Read the rotated and live journals in `data_dir`. A missing file is empty.
///
/// Test: `maintenance_log_tests::the_journal_rotates_and_reads_back_oldest_first`.
pub fn read_journal(data_dir: &Path) -> Result<JournalContents> {
    let mut out = JournalContents::default();
    for name in [MAINTENANCE_LOG_ROTATED_FILENAME, MAINTENANCE_LOG_FILENAME] {
        let path = data_dir.join(name);
        let file = match std::fs::File::open(&path) {
            Ok(f) => f,
            Err(e) if e.kind() == std::io::ErrorKind::NotFound => continue,
            Err(e) => return Err(e).with_context(|| format!("open {}", path.display())),
        };
        for line in std::io::BufReader::new(file).lines() {
            let line = line.with_context(|| format!("read {}", path.display()))?;
            if line.trim().is_empty() {
                continue;
            }
            match serde_json::from_str::<MaintenanceDeletion>(&line) {
                Ok(rec) => out.records.push(rec),
                Err(_) => out.malformed += 1,
            }
        }
    }
    Ok(out)
}

impl PalaceHandle {
    /// Forget `id` on behalf of a maintenance pass, then record the deletion.
    ///
    /// Why: a maintenance deletion must leave a trail (#8732); a user's
    /// `forget` must not be reclassified as one, so the two stay separate
    /// entry points.
    /// What: runs `PalaceHandle::forget_removing`. Only a real delete
    /// ([`ForgetOutcome::Deleted`]) is recorded; `survivor` carries the
    /// surviving drawer id and, for dedup, the score. #8729: a failed L1
    /// snapshot save after the delete is returned only after the record.
    /// Test: `maintenance_log_tests::dream_dedup_records_the_removed_and_surviving_drawer`,
    /// `maintenance_log_tests::a_user_forget_writes_a_user_forget_journal_record`,
    /// `maintenance_log_tests::a_failed_snapshot_save_after_the_delete_still_journals_the_copy`.
    pub async fn forget_for_maintenance(
        &self,
        id: Uuid,
        reason: DeletionReason,
        survivor: Option<(Uuid, Option<f32>)>,
    ) -> Result<ForgetOutcome> {
        // #9172: `forget_removing`, not `forget`, so a maintenance deletion of
        // a survivor writes this one record rather than two.
        let (removed, l1_saved) = self.forget_removing(id).await?;
        let Some(removed) = removed else {
            l1_saved?;
            return Ok(ForgetOutcome::NotFound);
        };
        // #8729: the record carries the removed drawer, not only its id.
        let mut rec = MaintenanceDeletion::new(&self.id, id, reason).with_drawer(&removed);
        if let Some((survivor_id, score)) = survivor {
            rec = rec.with_survivor(survivor_id, score);
        }
        record(self.data_dir.as_deref(), &rec);
        // #8729: the drawer is gone from redb; its copy is recorded first.
        l1_saved?;
        Ok(ForgetOutcome::Deleted)
    }
}

/// Record one user forget of `drawer`.
///
/// Why: #9283 — the doctor drawer-count check explains a count drop from this
/// journal, and a user forget left no line, so every one read as silent loss.
/// The record must not defeat the forget, so it carries no content (owner
/// ruling 2026-10-06). #9172: a forgotten dedup survivor keeps its copy,
/// because text merged into it from other drawers leaves with it.
/// What: a no-op without a data dir. When [`names_survivor`] says yes, records
/// [`DeletionReason::ForgetOfMergedSurvivor`] with the drawer copy; otherwise
/// records [`DeletionReason::UserForget`] with the content hash only.
/// Test: `maintenance_log_tests::a_user_forget_writes_a_user_forget_journal_record`,
/// `maintenance_log_tests::forgotten_content_cannot_be_recovered_from_the_journal`,
/// `dedup_survivor_tests::forgetting_a_dedup_survivor_writes_a_journal_record`.
pub(crate) fn record_user_forget(handle: &PalaceHandle, drawer: &Drawer) {
    let Some(dir) = handle.data_dir.as_deref() else {
        return;
    };
    let rec = if names_survivor(dir, drawer) {
        MaintenanceDeletion::new(
            &handle.id,
            drawer.id,
            DeletionReason::ForgetOfMergedSurvivor,
        )
        .with_drawer(drawer)
    } else {
        // #9283: id, reason, time and hash — never the content.
        MaintenanceDeletion::new(&handle.id, drawer.id, DeletionReason::UserForget)
            .with_content_hash_of(drawer)
    };
    record(Some(dir), &rec);
}

/// Whether a journal record in `dir` names `drawer` as its survivor.
///
/// What: answers `false` without reading when neither journal file was written
/// after `drawer` was created — no record can name a drawer younger than the
/// file. An unreadable journal answers `true`: a spare record costs less than
/// a missing one.
fn names_survivor(dir: &Path, drawer: &Drawer) -> bool {
    let created = std::time::SystemTime::from(drawer.created_at);
    let written_since = [MAINTENANCE_LOG_ROTATED_FILENAME, MAINTENANCE_LOG_FILENAME]
        .iter()
        .filter_map(|name| std::fs::metadata(dir.join(name)).ok()?.modified().ok())
        .any(|modified| modified >= created);
    if !written_since {
        return false;
    }
    match read_journal(dir) {
        Ok(journal) => journal
            .records
            .iter()
            .any(|r| r.survivor_id == Some(drawer.id)),
        Err(e) => {
            tracing::warn!(drawer_id = %drawer.id, "#9172: journal unreadable; recording the forget: {e:#}");
            true
        }
    }
}