trusty-memory 0.28.1

MCP server (stdio + Unix socket) for trusty-memory
Documentation
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
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
//! Shared helpers for the trusty-memory MCP tool surface.
//!
//! Why: Gate logic (content/blocklist/dedup), palace resolution, the shared
//! write-drawer pipeline, and attribution helpers are used across every
//! tool handler; collecting them in one submodule keeps the per-tool
//! handler files focused (issue #607).
//! What: free functions + small structs/consts moved verbatim out of the
//! former monolithic `tools.rs`. Visibility widened to `pub(crate)` where a
//! sibling submodule or the test module needs the item.
//! Test: `content_gate_*`, `blocklist_gate_*`, `dedup_*`, and the dispatch
//! tests in `tools::tests`.

use super::bm25::bm25_index_enqueue;
use crate::attribution::{session_tag_from_tags, CreatorInfo, CreatorSource, MCP_CLIENT_NAME};
use crate::kg_extract::{extract_triples, ExtractInput};
use crate::{ActivitySource, AppState, DaemonEvent};
use anyhow::{anyhow, Context, Result};
use serde_json::{json, Value};
use trusty_common::memory_core::filter::{FilterConfig, MCP_MIN_TOKENS};
use trusty_common::memory_core::palace::{PalaceId, RoomType};
use trusty_common::memory_core::retrieval::{admit_tier_c, RememberOptions, TierCAdmission};
use trusty_common::memory_core::timeouts::{self, OpBudget};
use uuid::Uuid;

/// Look up the friendly palace name (Palace.name) from the in-memory cache,
/// falling back to the id when the cache misses.
///
/// Why (issue #96): MCP-side emit calls need the same `palace_name` field
/// the HTTP path emits so the activity feed renders identical labels
/// regardless of origin.
/// Why (issue #228): the previous implementation called
/// `PalaceRegistry::list_palaces` — a synchronous filesystem walk — on every
/// `memory_remember` / `memory_note` write. With N palaces on disk that was
/// O(N) opendirs + `palace.json` reads per write, blocking the async runtime
/// thread (this helper had no `spawn_blocking` wrapper, unlike `palace_list`).
/// The replacement reads `state.palace_names` (a `DashMap` populated at
/// hydration / create time), which is a lock-free read and never touches
/// disk.
/// What: looks up `palace_id` in `state.palace_names`; on miss, returns the
/// id verbatim so emit calls never fail. Cache misses are non-fatal —
/// rename / create paths keep the cache in sync, but a fresh-after-restart
/// palace will hit the miss branch only until hydration completes.
/// Test: implicit — the MCP emit tests assert the `palace_id` matches; the
/// fallback is the same id-as-name behaviour the HTTP path uses. The cache
/// invariant is covered by `palace_name_cache_populated_after_hydration`
/// and `palace_name_cache_updates_on_create` in lib.rs.
pub(crate) fn lookup_palace_name(state: &AppState, palace_id: &str) -> String {
    state
        .palace_names
        .get(palace_id)
        .map(|entry| entry.value().clone())
        .unwrap_or_else(|| palace_id.to_string())
}

/// Minimum standalone-content word count enforced by [`content_gate`].
///
/// Why (issue #215): single-word user replies ("yes", "ok", "no thanks") have
/// no standalone memory value when the surrounding turn isn't captured
/// alongside them — they end up in the palace as orphan fragments that
/// pollute recall results. Requiring at least four whitespace-separated tokens
/// is a cheap heuristic that matches the natural boundary between "just a
/// reaction" and "an actual statement".
/// What: the threshold the gate compares against. Tokens are counted via
/// `split_whitespace().count()`, so punctuation does not inflate the count.
/// Test: `content_gate_blocks_short_no_context`, `content_gate_keeps_long`.
pub(crate) const CONTENT_GATE_MIN_WORDS: usize = 4;

/// Gate short standalone content unless a `context` wrapper is supplied or
/// the caller passes `force = true`.
///
/// Why: single-word or very-short standalone user responses ("yes", "ok")
/// have no standalone memory value (issue #215). Gate them unless a context
/// is provided. Issue #2442: `memory_remember`'s `force` operator override
/// must bypass this gate the same way it bypasses `dedup_gate` and the
/// deep `FilterConfig` noise/secret gate — previously `force` was parsed
/// AFTER this gate ran, so `force = true` writes of short standalone
/// content were still silently dropped, contradicting the tool's
/// documented "bypass all content-quality gates" contract.
/// What: returns `None` if `force` is `false`, `content` has fewer than
/// [`CONTENT_GATE_MIN_WORDS`] whitespace-separated tokens, AND `context` is
/// `None` (the write should be skipped). Returns `Some(combined)` where
/// `combined = "<context>\n\n---\n\n<content>"` when `context` is `Some` and
/// non-empty after trimming (context wrapping is formatting, not gating, so
/// it always applies regardless of `force`). Returns `Some(content)`
/// unchanged when `content` has at least [`CONTENT_GATE_MIN_WORDS`] tokens,
/// or when `force` is `true`. Tokens are counted on the trimmed `content` so
/// trailing whitespace doesn't inflate the count.
/// Test: `content_gate_blocks_short_no_context`,
/// `content_gate_wraps_short_with_context`,
/// `content_gate_keeps_long`, `content_gate_blank_context_treated_as_none`,
/// `content_gate_force_bypasses_short_content`.
pub(crate) fn content_gate(content: &str, context: Option<&str>, force: bool) -> Option<String> {
    let trimmed = content.trim();
    let word_count = trimmed.split_whitespace().count();
    // Treat a context that is empty or whitespace-only as "no context" — a
    // caller passing `""` should not unlock a write the gate would otherwise
    // drop, and the combined output would otherwise begin with a meaningless
    // separator.
    let context_clean = context.map(str::trim).filter(|s| !s.is_empty());
    if let Some(ctx) = context_clean {
        return Some(format!("{ctx}\n\n---\n\n{content}"));
    }
    if !force && word_count < CONTENT_GATE_MIN_WORDS {
        return None;
    }
    Some(content.to_string())
}

/// Rolling-window horizon for the dedup gate.
///
/// Why (issue #220): identical content is often emitted multiple times in
/// quick succession (auto-capture hook bursts, retries, copy-paste). A
/// 5-minute window catches the burst without rejecting deliberate user
/// re-statements hours later.
/// What: `chrono::Duration` value. Drawers created before
/// `now - DEDUP_WINDOW` are ignored by the dedup pass.
/// Test: indirect via `dedup_skips_near_duplicate` and
/// `dedup_allows_different_content` (use the helper directly).
pub(crate) const DEDUP_WINDOW_MINUTES: i64 = 5;

/// Maximum number of recent drawers the dedup pass scans.
///
/// Why: a palace can hold tens of thousands of drawers; we never need to
/// compare the new write against more than the most-recent handful to
/// catch the bursty-duplicate case. Capping the scan keeps the hot path
/// O(1) in the palace size.
/// What: ceiling on the candidate list pulled from
/// `PalaceHandle::list_drawers` before the time-window filter.
/// Test: `dedup_skips_near_duplicate` exercises the scan against a small
/// candidate set; the cap is enforced by `list_drawers`'s `limit` arg.
pub(crate) const DEDUP_SCAN_LIMIT: usize = 50;

/// Jaro-Winkler similarity threshold above which a candidate counts as a
/// near-duplicate of the new content.
///
/// Why: 0.92 is the empirically-chosen cutoff documented in the issue —
/// high enough to allow distinct facts to coexist, low enough to catch
/// trivial whitespace / punctuation / suffix variation. Jaro-Winkler is
/// preferred over plain Jaro because the auto-capture noise tends to share
/// the same prefix (`Tool use: …`, `Edit File: …`), which Jaro-Winkler
/// weights heavily.
/// What: `f64` threshold compared against `strsim::jaro_winkler`'s output.
/// Test: `dedup_skips_near_duplicate`, `dedup_allows_different_content`.
pub(crate) const DEDUP_SIMILARITY_THRESHOLD: f64 = 0.92;

/// Blocklist gate: returns the matched pattern when the content should be
/// silently skipped because it matches a known low-value auto-capture pattern.
///
/// Why (issue #220): Centralises the pattern-match logic so both
/// `memory_remember` and `memory_note` go through the same filter. Trims
/// leading whitespace before matching so indented variants still hit.
/// Why (issue #1481): returns *which* pattern matched (instead of a bare
/// `bool`) so the caller can name the trigger in the skip envelope rather than
/// emitting an opaque "blocked pattern" — callers need to know what to remove.
/// Why (issue #2442): delegates to
/// [`trusty_common::memory_core::filter::blocklist_match`], the single
/// prefix-anchored source of truth shared with the dream-cycle retroactive
/// prune pass. The old local copy matched via `str::contains`
/// (substring-anywhere), which fired on legitimate content that merely
/// QUOTED an auto-capture phrase mid-prose (e.g. a coding agent's turn
/// recapping `"Tool use: Bash"` from its own transcript) — thinning the
/// recall surface. `blocklist_match` anchors to the start of the trimmed
/// content instead.
/// What: thin wrapper around `blocklist_match`.
/// Test: `blocklist_gate_blocks_tool_use`,
/// `blocklist_gate_blocks_session_ended`,
/// `blocklist_gate_passes_normal_content`,
/// `blocklist_gate_names_matched_pattern`.
pub(crate) fn blocklist_gate(content: &str) -> Option<&'static str> {
    trusty_common::memory_core::filter::blocklist_match(content)
}

/// Dedup gate: returns true when the new content is a near-duplicate of a
/// drawer written to the same palace within the rolling window.
///
/// Why (issue #220): bursts of identical or near-identical content (auto-
/// capture retries, hook re-emissions, copy-paste artefacts) were
/// inflating the palace with no recall benefit. A short rolling window
/// catches the burst without rejecting deliberate re-statements hours
/// later.
/// What: pulls up to `DEDUP_SCAN_LIMIT` recent drawers from the live
/// in-memory table via `list_drawers` (a cheap snapshot, no I/O), filters
/// to those created within `DEDUP_WINDOW_MINUTES` of `now`, then computes
/// `strsim::jaro_winkler` against each. Returns `true` on the first match
/// above `DEDUP_SIMILARITY_THRESHOLD`. Returns `false` if `content` is
/// empty after trimming (the content gate handles that case separately)
/// or if the palace has no recent drawers.
/// Test: `dedup_skips_near_duplicate`, `dedup_allows_different_content`.
pub(crate) fn dedup_gate(handle: &trusty_common::memory_core::PalaceHandle, content: &str) -> bool {
    let trimmed = content.trim();
    if trimmed.is_empty() {
        return false;
    }
    let now = chrono::Utc::now();
    let window_start = now - chrono::Duration::minutes(DEDUP_WINDOW_MINUTES);
    let recent = handle.list_drawers(None, None, DEDUP_SCAN_LIMIT);
    recent
        .iter()
        .filter(|d| d.created_at >= window_start)
        .any(|d| strsim::jaro_winkler(trimmed, d.content().trim()) > DEDUP_SIMILARITY_THRESHOLD)
}

/// Build the strict MCP-level `RememberOptions`.
///
/// Why: Issue #61 — the MCP boundary is where auto-capture hooks deposit
/// raw tool/commit/prompt data; we want the 8-token threshold there even
/// though the library default is more permissive for direct callers.
/// Issue #1970: `defer_embedding` is threaded through separately so callers
/// can key it off the live `AppState::readiness()` at dispatch time rather
/// than baking a stale snapshot into this helper.
/// Issue #2520 (two-tier `force`): `allow_secret_like` is threaded through
/// separately from `force` — the MCP `force` arg only bypasses quality
/// gates now; a caller must set the distinct `allow_secret_like` arg to
/// also bypass secret detection.
/// What: Clones the default filter and bumps `min_tokens` to `MCP_MIN_TOKENS`.
/// Test: `dispatch_remember_rejects_short_content`;
/// `remember_succeeds_and_defers_embedding_while_state_is_warming` covers
/// `defer_embedding`; `dispatch_remember_force_still_blocks_secret` and
/// `dispatch_remember_allow_secret_like_bypasses_secret_gate` (this crate's
/// `tools::tests`) cover `allow_secret_like` end-to-end through MCP dispatch;
/// `remember_force_still_blocks_secret` and
/// `remember_force_and_allow_secret_like_stores_secret_shaped_content`
/// (`trusty_common::memory_core::retrieval::tests`) cover the same split at
/// the library layer.
pub(crate) fn mcp_remember_opts(
    force: bool,
    defer_embedding: bool,
    allow_secret_like: bool,
) -> RememberOptions {
    let filter = FilterConfig {
        min_tokens: MCP_MIN_TOKENS,
        ..FilterConfig::default()
    };
    RememberOptions {
        filter,
        force,
        defer_embedding,
        allow_secret_like,
        ..RememberOptions::default()
    }
}

/// Resolve the ADR-0028 Tier C arguments and run the admission gate (#4886).
///
/// Why: the fail-closed degradation D4 requires has to be OBSERVABLE at the
/// tool boundary. A write that asked for a slot and silently became an ordinary
/// drawer teaches the author nothing, and they will keep writing keys that
/// never take effect. Running the (pure) gate here lets the handler report the
/// tier it actually got, and passing the RESOLVED values down means the library
/// re-runs the same decision on already-valid input rather than re-deriving a
/// default against a slightly later clock.
/// What: reads `fact_key` and the RFC 3339 `expires_at` from `args` and returns
/// the [`TierCAdmission`]. A malformed `expires_at` STRING is a hard error, not
/// a degradation: an unparseable timestamp is a caller typo, and silently
/// substituting the 24-hour default would hide it. Semantic refusals (a key
/// that breaks the D5 grammar, a TTL that has already elapsed) degrade, because
/// those are the cases D4 says must fall back to today's behaviour.
/// Test: `dispatch_remember_admits_a_tier_c_slot`,
/// `dispatch_remember_reports_a_refused_slot_as_tier_e`,
/// `dispatch_remember_rejects_an_unparseable_expires_at`.
pub(crate) fn resolve_tier_c(args: &Value, tool: &str) -> Result<TierCAdmission> {
    let expires_at = match args.get("expires_at").and_then(|v| v.as_str()) {
        Some(raw) => Some(
            chrono::DateTime::parse_from_rfc3339(raw)
                .map(|t| t.with_timezone(&chrono::Utc))
                .with_context(|| {
                    format!("{tool}: 'expires_at' must be an RFC 3339 timestamp, got {raw:?}")
                })?,
        ),
        None => None,
    };
    Ok(admit_tier_c(
        args.get("fact_key").and_then(|v| v.as_str()),
        expires_at,
        chrono::Utc::now(),
    ))
}

/// Fold an admission decision into the write options and the response envelope.
///
/// Why: both `memory_remember` and `memory_note` need the identical
/// "carry the resolved slot down, report the tier back up" step; duplicating it
/// is how the two tools would drift on what a refusal looks like.
/// What: on admission, sets `opts.fact_key` / `opts.expires_at` to the resolved
/// values. Otherwise leaves them unset so the library writes an ordinary
/// drawer. Returns `(tier_label, refusal_reason)` for the JSON envelope.
/// Test: same tests as [`resolve_tier_c`].
pub(crate) fn apply_tier_c(
    admission: &TierCAdmission,
    opts: &mut RememberOptions,
) -> (&'static str, Option<String>) {
    match admission {
        TierCAdmission::Admitted {
            fact_key,
            expires_at,
        } => {
            opts.fact_key = Some(fact_key.clone());
            opts.expires_at = Some(*expires_at);
            ("C", None)
        }
        TierCAdmission::Refused(refusal) => ("E", Some(refusal.to_string())),
        TierCAdmission::NotRequested => ("E", None),
    }
}

/// Reverse of `RoomType::parse`: produce a stable label for KG `in-room`
/// extraction.
///
/// Why: The auto-extractor wants the same friendly label the caller passed
/// (`"Backend"`, `"General"`, …) so the graph stays consistent across
/// remember calls regardless of how the MCP client spelled the argument.
/// ADR-0027 T3: the projection itself now lives in `trusty-common`
/// (`room_identity::room_label`) so the room registry and the KG extractor
/// cannot drift apart on what a room is called; this stays as the
/// `Option`-shaped adapter the extractor call sites expect.
/// What: Delegates to the shared projection — the canonical enum-name string
/// for the built-in variants, the inner string for `Custom`.
/// Test: Indirect — `auto_kg_extraction_hooks_into_memory_remember`
/// round-trips a known room label.
pub(crate) fn room_label(room: &RoomType) -> Option<String> {
    Some(trusty_common::memory_core::room_identity::room_label(room))
}

/// Refuse a maintenance tool call unless this process holds the data root's
/// maintenance lease (#8733).
///
/// Why: the MCP dream and compact tools evict, delete and rewrite rows, which
/// only the elected maintainer may do; `dream_run` already refuses this way.
/// What: `Ok` when `may_run_maintenance` allows it; otherwise an error carrying
/// the same text as `dream_run`'s `Conflict`. An unavailable lease refuses.
/// Test: `palace_dream_is_refused_without_the_maintenance_lease`,
/// `dream_consolidate_room_is_refused_without_the_maintenance_lease`,
/// `palace_compact_is_refused_without_the_maintenance_lease`.
pub(crate) fn require_maintenance_lease(state: &AppState) -> Result<()> {
    if state.registry.may_run_maintenance() {
        Ok(())
    } else {
        Err(anyhow!(crate::service::core_kg::MAINTENANCE_LEASE_NOT_HELD))
    }
}

/// Resolve (or lazily open) the palace handle for a tool call.
///
/// Why the liveness guard (issue #4001): this is the single choke point every
/// `memory_recall` / `memory_remember` passes through, and it is exactly where
/// the #3992 wedge occurred — `open_palace` serialises on a per-palace mutex
/// and, on a miss, drops into `concurrent_open`'s redb retry/backoff loop.
/// Registering the operation here is what lets `/health` (and therefore both
/// doctors) observe that work has stopped moving. The guard is RAII, so the
/// `?` on the line below releases it just as reliably as the success path.
/// What: unchanged behaviour, plus a [`crate::worker_liveness`] registration
/// held for the duration of the open. Writes register once for their whole
/// duration in [`begin_budgeted_write`] instead.
/// Test: `worker_liveness::tests`, `web::tests::health_tests`.
pub(crate) fn open_palace_handle(
    state: &AppState,
    palace_id: &str,
) -> Result<std::sync::Arc<trusty_common::memory_core::PalaceHandle>> {
    let _tracked = state.worker_liveness.track();
    open_palace_handle_within(state, palace_id, OpBudget::start(state.write_op_budget))
}

/// Open a palace handle inside an operation that has already spent part of its
/// budget (issue #4002).
///
/// Why: [`open_palace_handle`] starts a fresh
/// [`trusty_common::memory_core::timeouts::open_queue_timeout`] window every
/// time it is called. On the write path the caller has already waited up to
/// `write_lock_timeout` for the per-palace write mutex, so the two waits
/// composed into their sum — 60 s + ~63 s before any error surfaced. Write
/// callers pass their in-flight budget here so this leg spends only the
/// remainder.
/// What: identical to [`open_palace_handle`] except the open-queue wait is
/// clamped by `budget` and it registers no liveness guard of its own. Every
/// caller runs inside [`begin_budgeted_write`], whose [`WriteGuard`] already
/// counts the write once. Read paths keep calling [`open_palace_handle`], which
/// stamps a full budget and therefore behaves exactly as before.
/// Test: `tools::tests::write_budget_tests`, `tools::tests::write_liveness_tests`.
pub(crate) fn open_palace_handle_within(
    state: &AppState,
    palace_id: &str,
    budget: OpBudget,
) -> Result<std::sync::Arc<trusty_common::memory_core::PalaceHandle>> {
    // #4001: untracked here so one write is not counted twice in `in_flight`.
    let pid = PalaceId::new(palace_id);
    state
        .registry
        .open_palace_within(&state.data_root, &pid, budget)
        .with_context(|| format!("open palace {palace_id}"))
}

/// Acquire the per-palace write mutex as the FIRST leg of a budgeted write
/// (issue #4002).
///
/// Why: `memory_remember`, `memory_note` and `task_add` all open with the same
/// three lines, and all three previously handed `write_lock_timeout()` straight
/// to the lock. That is the leg the open-queue wait then stacked on top of.
/// Routing them through one helper is what guarantees the budget is stamped
/// before the first wait — a per-handler copy is exactly how one of them would
/// later drift back to the additive shape.
/// What: stamps an [`OpBudget`] of [`AppState::write_op_budget`], registers the
/// write with [`crate::worker_liveness`], waits at most
/// `min(write_lock_timeout(), budget)` for the palace's write mutex, and
/// returns the budget so the caller can pass its remainder to
/// [`open_palace_handle_within`]. The registration lives inside the returned
/// [`WriteGuard`], so the write counts from the start of its wait until the
/// lock is released. `tool` prefixes the error so the caller sees which
/// handler gave up.
/// Test: `tools::tests::write_budget_tests`, `tools::tests::write_liveness_tests`.
pub(crate) async fn begin_budgeted_write<'a>(
    state: &'a AppState,
    write_lock: &'a std::sync::Arc<tokio::sync::Mutex<()>>,
    palace_id: &str,
    tool: &str,
) -> Result<(WriteGuard<'a>, OpBudget)> {
    let budget = OpBudget::start(state.write_op_budget);
    // #4001: register before the wait so a queued writer is visible, and hold
    // the registration until the lock is released so a stalled holder is too.
    let tracked = state.worker_liveness.track();
    let lock = timeouts::lock_with_timeout(
        write_lock,
        budget.leg(timeouts::write_lock_timeout()),
        palace_id,
    )
    .await
    .map_err(|e| anyhow::anyhow!("{tool}: {e:#}"))?;
    Ok((
        WriteGuard {
            _lock: lock,
            _tracked: tracked,
        },
        budget,
    ))
}

/// A held palace write lock plus the write's liveness registration (#4001).
///
/// Why: on 2026-09-13 a writer held this lock and never finished, while
/// `memory.health` reported `ok` with 0 in flight. Queued writers give up at
/// their bound, which sits below the wedge threshold, so only the holder's age
/// can trip it. Binding the registration to the lock guard is what keeps the
/// holder visible, and it means neither can be dropped without the other.
/// A write that holds the lock past the threshold reads as wedged even if it
/// later finishes: every writer queued behind it has already failed.
/// What: the `tokio` lock guard and a [`crate::worker_liveness::WorkGuard`];
/// both release on drop.
/// Test: `tools::tests::write_liveness_tests`.
pub(crate) struct WriteGuard<'a> {
    _lock: tokio::sync::MutexGuard<'a, ()>,
    _tracked: crate::worker_liveness::WorkGuard<'a>,
}

/// Run deterministic KG extraction over a freshly-written drawer and assert
/// any resulting triples through the palace's `KnowledgeGraph`.
///
/// Why: Issue #97 — `memory_remember` and `memory_note` should auto-populate
/// the KG so palaces with drawers always have a graph. The extractor is pure
/// and offline so the write hot path stays fast; failures *must never* fail
/// the parent write (the drawer is already on disk), so this function logs
/// and swallows every error.
/// What: Builds an `ExtractInput`, runs `extract_triples`, then calls
/// `handle.kg.assert` for each triple. Any failure during assertion is
/// captured as a `tracing::warn!` and the rest of the triples are still
/// attempted; the function returns nothing.
/// Test: `auto_kg_extraction_hooks_into_memory_remember`,
/// `auto_kg_extraction_no_op_does_not_fail_remember`,
/// `web::tests::http_create_drawer_runs_auto_kg_extraction`.
pub(crate) async fn auto_extract_and_assert(
    handle: &std::sync::Arc<trusty_common::memory_core::PalaceHandle>,
    drawer_id: Uuid,
    content: &str,
    tags: &[String],
    room: Option<&str>,
) {
    let input = ExtractInput {
        drawer_id,
        content,
        tags,
        room,
    };
    let triples = extract_triples(&input);
    if triples.is_empty() {
        return;
    }
    for triple in triples {
        let s = triple.subject.clone();
        let p = triple.predicate.clone();
        if let Err(e) = handle.kg.assert(triple).await {
            tracing::warn!(
                drawer_id = %drawer_id,
                subject = %s,
                predicate = %p,
                "auto kg extraction: assert failed (non-fatal): {e:#}",
            );
        }
    }
}

/// Resolve a palace argument, falling back to `state.default_palace` when
/// the caller omitted `palace`.
///
/// Why: `serve --palace <name>` lets the operator bind a process to a single
/// project namespace; tool calls then no longer need to repeat the palace
/// every time. This helper centralises the precedence rule (explicit arg
/// wins over default) and produces a uniform error when neither is set.
/// What: Returns the explicit `args["palace"]` string if present, otherwise
/// `state.default_palace`. Errors with a helpful message if both are absent.
/// Test: `default_palace_used_when_arg_omitted` and
/// `dispatch_unknown_tool_errors`.
pub(crate) fn resolve_palace<'a>(
    state: &'a AppState,
    args: &'a Value,
    tool: &str,
) -> Result<String> {
    if let Some(p) = args.get("palace").and_then(|v| v.as_str()) {
        return Ok(p.to_string());
    }
    state
        .default_palace
        .clone()
        .ok_or_else(|| anyhow!("{tool}: missing 'palace' (no --palace default configured)"))
}

/// Inputs to the shared write-drawer pipeline.
///
/// Why (issue #227): `memory_remember` and `memory_note` share the same
/// "open palace → write drawer → fan-out side effects" tail. Capturing those
/// inputs in one struct keeps the handler call sites flat and makes the
/// shared pipeline a single function — every behavioural divergence between
/// the two tools is now visible in their handlers, not buried in a
/// 60-line block of duplicated post-write fan-out.
/// What: bundles every value the post-gate pipeline needs. `room_label_for_kg`
/// is pre-computed by the handler (memory_note pins it to `"General"`;
/// memory_remember derives it from `RoomType` via [`room_label`]).
/// Test: exercised end-to-end by `dispatch_remember_then_recall`,
/// `dispatch_remember_with_context_writes_combined`, and the note tests.
pub(crate) struct WriteDrawerParams<'a> {
    pub(crate) palace_id: &'a str,
    pub(crate) content: String,
    pub(crate) tags: Vec<String>,
    pub(crate) room: RoomType,
    pub(crate) importance: f32,
    pub(crate) opts: RememberOptions,
    pub(crate) room_label_for_kg: Option<String>,
    /// Remaining budget of the write that is already in flight (issue #4002).
    ///
    /// Why: `write_drawer` opens the palace handle, which is the second leg of
    /// the wait the caller started when it took the write mutex. Without the
    /// caller's budget it would start a fresh open-queue window and the two
    /// legs would sum again.
    /// What: the [`OpBudget`] returned by [`begin_budgeted_write`].
    /// Test: `tools::tests::write_budget_tests`.
    pub(crate) budget: OpBudget,
}

/// Run the shared write pipeline after content has been gated and attribution
/// applied.
///
/// Why (issue #227): centralises the open-palace → remember → BM25 → emit →
/// auto-KG-extract tail that `memory_remember` and `memory_note` both run.
/// Hosting it in one place keeps the side-effect ordering identical across
/// the two tools and makes future write-side hooks land in one location.
/// What: opens the palace handle, calls `remember_with_options`, fires the
/// BM25 index task, emits `DrawerAdded` + the aggregate status event, and
/// runs the auto-KG-extraction pass (best-effort). Returns the new drawer
/// id on success; any underlying error propagates via `anyhow::Result`.
///
/// **A write-pipeline timeout error does not mean the write did not land**
/// (#6366): the durable commit, once dispatched, completes regardless of the
/// caller's budget, so retrying blind duplicates content that carries no
/// `fact_key`. Slotted (Tier C) writes are idempotent per slot and safe to
/// retry. See `PalaceHandle::remember_with_options_within`.
/// Test: covered through `dispatch_remember_then_recall`,
/// `dispatch_remember_with_context_writes_combined`,
/// `dispatch_note_skips_short_no_context` (negative path before this runs),
/// and `auto_kg_extraction_hooks_into_memory_remember`.
pub(crate) async fn write_drawer(state: &AppState, params: WriteDrawerParams<'_>) -> Result<Uuid> {
    let WriteDrawerParams {
        palace_id,
        content,
        tags,
        room,
        importance,
        opts,
        room_label_for_kg,
        budget,
    } = params;

    // #4002: second leg of the caller's budget, not a fresh open-queue window.
    let handle = open_palace_handle_within(state, palace_id, budget)?;
    // Snapshot the preview before `content` is moved into the write so the
    // activity feed shows what landed on disk (matches the HTTP path).
    let preview = crate::service::drawer_content_preview(&content);
    // Issue #97: keep originals so the auto-KG extractor sees the same
    // content / tags that landed in the drawer. `remember_with_options`
    // consumes them, so clone before the call.
    let content_for_kg = content.clone();
    let tags_for_kg = tags.clone();
    // #6366: the write mutex is held for this whole call. Pass the daemon's
    // configured ceiling so a slow commit fails with a named reason and
    // releases the mutex, instead of stalling every other writer on this palace.
    let drawer_id = handle
        .remember_with_options_within(
            content,
            room,
            tags,
            importance,
            opts,
            state.write_pipeline_budget,
        )
        .await
        .context("PalaceHandle::remember_with_options")?;
    // Issue #156 + #231: opt-in BM25 lexical lane. Enqueue onto the
    // bounded indexer channel so the redb write returns immediately;
    // a full queue is dropped + logged rather than allowed to grow
    // unbounded behind a slow daemon (#231). Daemon errors observed
    // by the worker are logged but never block the MCP response.
    // #5036: index under the RESOLVED palace id. `open_palace` follows
    // aliases, so writing under the requested slug files the drawer into a
    // palace the reader never searches.
    bm25_index_enqueue(state, handle.id.as_str(), drawer_id, &content_for_kg);
    // Issue #96: emit a DrawerAdded so the activity feed shows
    // MCP-origin writes with `source = Mcp`.
    let palace_name = lookup_palace_name(state, palace_id);
    let drawer_count = handle.drawers.read().len();
    state.emit(DaemonEvent::DrawerAdded {
        palace_id: palace_id.to_string(),
        palace_name,
        drawer_count,
        timestamp: chrono::Utc::now(),
        content_preview: preview,
        source: ActivitySource::Mcp,
    });
    // Issue #228: do NOT emit `StatusChanged` on every write — the
    // aggregate-recompute was O(N palaces) of disk I/O on the hot path.
    // The periodic ticker spawned by `run_http_on` refreshes dashboard
    // totals on a fixed cadence; mutations themselves still surface via
    // the `DrawerAdded` SSE frame above.
    // Issue #97: best-effort auto-extraction. Failures never fail the
    // write — the drawer is already on disk.
    auto_extract_and_assert(
        &handle,
        drawer_id,
        &content_for_kg,
        &tags_for_kg,
        room_label_for_kg.as_deref(),
    )
    .await;
    Ok(drawer_id)
}

/// Build a JSON "skipped" envelope used by both write handlers when a gate
/// rejects the input.
///
/// Why (issue #227): keeps the three skip reasons (`blocked pattern`,
/// `short prompt, no context`, `duplicate within window`) emitted as a
/// uniform shape so callers can parse the envelope without per-tool
/// branching.
/// What: returns `{"palace": <id>, "status": "skipped", "reason": <reason>}`.
/// Test: exercised by `dispatch_remember_skips_short_no_context`,
/// `dispatch_note_skips_short_no_context`,
/// `dispatch_remember_blocks_blocklist_pattern`.
pub(crate) fn skipped_envelope(palace_id: &str, reason: &str) -> Value {
    json!({
        "palace": palace_id,
        "status": "skipped",
        "reason": reason,
    })
}

/// Extract a `tags` argument (JSON array of strings) into a `Vec<String>`.
///
/// Why: every write-side handler accepts an optional `tags` argument with
/// identical shape; centralising the parse keeps the handlers focused on
/// their tool-specific logic.
/// What: returns the strings in order; non-string entries are silently
/// dropped (matches pre-refactor behaviour).
/// Test: covered indirectly by `dispatch_remember_then_recall` and
/// `auto_kg_extraction_hooks_into_memory_remember`.
pub(crate) fn parse_tags(args: &Value) -> Vec<String> {
    args.get("tags")
        .and_then(|v| v.as_array())
        .map(|arr| {
            arr.iter()
                .filter_map(|t| t.as_str().map(|s| s.to_string()))
                .collect()
        })
        .unwrap_or_default()
}

/// Attach the MCP attribution tags (`creator:*` and the bare-UUID session
/// projection) to the caller-supplied tag list.
///
/// Why (Submission-logging Part B + issue #202): every MCP-origin drawer must
/// carry the writer identity so the activity panel and audit logs can attribute
/// the write. Issue #202 also projects a bare-UUID session tag into the
/// reserved `creator:session=<first-8>` slot when present.
///
/// Why (CRITICAL fix, DOC-53 §4.3): this handler runs inside the SHARED
/// `trusty-memory` daemon serving every concurrently-attached session from
/// one process — it previously used `CreatorInfo::new_self`, which resolves
/// `std::env::current_dir()`/`TM_WORKSTREAM_NAME` from the DAEMON's own
/// process, stamping the identical `ws:`/`creator:workstream=` tag onto
/// every caller's writes regardless of which session actually wrote them.
/// `args["cwd"]`/`args["workstream"]` are populated per-request by the MCP
/// stdio bridge (`commands::serve_stdio_bridge::inject_caller_context`,
/// mirroring the existing `args["cwd"]` precedent in
/// `tools::palace_ops::handle_palace_create`) — [`CreatorInfo::new_for_caller`]
/// trusts only that per-request value, never the daemon's own identity.
/// What: appends the session-tag projection (when one is found in the input
/// tags) then merges `CreatorInfo::new_for_caller(MCP, Mcp, args["cwd"],
/// args["workstream"])` into the vec via [`CreatorInfo::merge_into_deduped`]
/// — deduped (not plain `merge_into`) because a hand-written claim drawer
/// (DOC-53 §3.1) already carries `ws:<name>` in its own caller-supplied tags
/// by convention, and a plain merge would duplicate it (MEDIUM 1).
/// Test: `dispatch_remember_then_recall` (attribution present);
/// `mcp_writes_carry_distinct_ws_tags_per_caller_over_rpc` (the caller-vs-
/// daemon isolation this exists to fix — `web::tests::attribution_tests`);
/// `attach_mcp_attribution_dedupes_hand_written_ws_claim_tag`.
pub(crate) fn attach_mcp_attribution(tags: &mut Vec<String>, args: &Value) {
    if let Some(session_tag) = session_tag_from_tags(tags) {
        tags.push(session_tag);
    }
    let caller_cwd = args
        .get("cwd")
        .and_then(|v| v.as_str())
        .filter(|s| !s.is_empty());
    let caller_workstream = args
        .get("workstream")
        .and_then(|v| v.as_str())
        .filter(|s| !s.is_empty());
    CreatorInfo::new_for_caller(
        MCP_CLIENT_NAME,
        CreatorSource::Mcp,
        caller_cwd,
        caller_workstream,
    )
    .merge_into_deduped(tags);
}