Skip to main content

trusty_memory/
kg_write.rs

1//! The one entry point for asserting a caller-supplied KG triple.
2//!
3//! Why: writing a triple that a caller chose is three obligations, not one —
4//! the Tier S admission gate (#4888), the write itself, and refreshing the
5//! always-injected prompt cache when the predicate is hot. Six surfaces owed
6//! all three and each carried its own copy, so the third drifted: the HTTP
7//! `POST /api/v1/palaces/{id}/kg` path (#5524) and the chat `kg_assert` tool
8//! (#4905) stored the triple and never refreshed, which made a standing rule
9//! written through either surface invisible to every later turn while both
10//! reported success. Which client a user happened to use silently decided
11//! whether their fact took effect. Centralising the sequence here removes the
12//! call-site obligation that the drift was made of — a new write surface gets
13//! all three by construction, and a behaviour fix lands once.
14//!
15//! What: [`assert_triple`] runs admission → assert → refresh in that order and
16//! returns [`KgWriteError`] on any step. The admission guard is held across the
17//! write, which is what makes the Tier S check-then-write sequence atomic.
18//! Refresh failures are reported, never swallowed: a triple that reached
19//! storage without reaching the prompt surface is precisely the defect this
20//! module exists to prevent, so it is a distinct error variant
21//! ([`KgWriteError::CacheRefresh`]) rather than a `warn!`.
22//!
23//! Test: `kg_write_refreshes_cache_for_hot_predicate`,
24//! `kg_write_skips_refresh_for_cold_predicate`,
25//! `kg_write_admission_refusal_leaves_storage_and_cache_untouched`,
26//! `kg_write_batch_policy_defers_refresh` in the inline `tests` module; the
27//! per-surface proofs are `http_kg_assert_endpoint_refreshes_prompt_cache`
28//! (`web::tests::prompt_tests`) and `chat_kg_assert_refreshes_prompt_cache`
29//! (`web::tests::chat_tests`).
30
31use std::sync::Arc;
32
33use trusty_common::memory_core::store::kg::Triple;
34use trusty_common::memory_core::PalaceHandle;
35
36use crate::AppState;
37
38/// Why a caller-supplied KG assert did not complete.
39///
40/// Why: the three failure points have genuinely different meanings to a
41/// caller. `Admission` means nothing was written and the caller should fix the
42/// request; `Assert` means nothing was written and the daemon is at fault;
43/// `CacheRefresh` means the triple IS in storage but the always-injected
44/// surface is stale. Collapsing them into one opaque error would leave HTTP
45/// unable to choose 400 vs 500, and would let the third be mistaken for a
46/// clean success — which is the bug in #5524 and #4905.
47/// What: a `thiserror` enum over `anyhow::Error`. `Admission` is transparent so
48/// the Tier S refusal text (which names the occupants and the tool that
49/// retires one) reaches the caller unchanged.
50/// Test: `kg_write_admission_refusal_leaves_storage_and_cache_untouched`.
51#[derive(Debug, thiserror::Error)]
52pub enum KgWriteError {
53    /// The Tier S gate refused the write. Nothing was stored.
54    #[error(transparent)]
55    Admission(anyhow::Error),
56    /// The KG write itself failed. Nothing was stored.
57    #[error("kg assert: {0:#}")]
58    Assert(anyhow::Error),
59    /// The triple was stored but the prompt cache could not be rebuilt, so the
60    /// fact is not yet on the always-injected surface.
61    #[error("triple written but prompt cache refresh failed: {0:#}")]
62    CacheRefresh(anyhow::Error),
63}
64
65/// When [`assert_triple`] refreshes the prompt cache after a hot write.
66///
67/// Why: a rebuild walks every registered palace, so a caller asserting N
68/// triples in one pass would pay N full walks for one net change. `Batch` lets
69/// such a caller defer to a single trailing refresh WITHOUT deciding for itself
70/// what counts as hot — that decision stays in [`assert_triple`] either way,
71/// which is the property this module exists to protect.
72/// What: `Inline` refreshes before returning; `Batch` skips the refresh and
73/// reports hotness in the receipt for [`refresh_after_batch`].
74/// Test: `kg_write_batch_policy_defers_refresh`.
75#[derive(Debug, Clone, Copy, PartialEq, Eq)]
76pub enum CachePolicy {
77    /// Refresh inline when the predicate is hot. Every single-write caller.
78    Inline,
79    /// Defer the refresh; the caller must pass the accumulated hotness to
80    /// [`refresh_after_batch`] when its loop ends.
81    Batch,
82}
83
84/// What [`assert_triple`] did, for callers that batch.
85#[derive(Debug, Clone, Copy)]
86pub struct AssertReceipt {
87    /// Whether the asserted predicate is on the Tier S prompt surface.
88    pub hot: bool,
89}
90
91/// Assert `triple` into `handle`'s KG, keeping the Tier S surface coherent.
92///
93/// Why: see the module doc — this is the single sequence every caller-supplied
94/// KG write must run, and the reason it is a function rather than a convention
95/// is that the convention already failed twice (#5524, #4905).
96///
97/// What: three steps in a fixed order.
98/// 1. [`crate::prompt_facts::check_tier_s_admission`], whose returned guard is
99///    held until after the write — dropping it earlier reopens the
100///    check-then-write race the lock closes.
101/// 2. `KnowledgeGraph::assert`.
102/// 3. When the predicate is hot AND `policy` is [`CachePolicy::Inline`],
103///    [`crate::prompt_facts::rebuild_prompt_cache`].
104///
105/// Hotness is read BEFORE the write because `triple` is moved into it. A
106/// failure at any step is returned; step 3 in particular is never downgraded to
107/// a log line, because "stored but invisible" is the exact user-visible symptom
108/// of the two issues this replaces.
109///
110/// Test: `kg_write_refreshes_cache_for_hot_predicate`,
111/// `kg_write_skips_refresh_for_cold_predicate`,
112/// `kg_write_admission_refusal_leaves_storage_and_cache_untouched`.
113pub async fn assert_triple(
114    state: &AppState,
115    handle: &Arc<PalaceHandle>,
116    triple: Triple,
117    policy: CachePolicy,
118) -> Result<AssertReceipt, KgWriteError> {
119    // #5524: `_admission` holds the Tier S lock across `kg.assert` below.
120    let _admission = crate::prompt_facts::check_tier_s_admission(
121        state,
122        handle,
123        &triple.subject,
124        &triple.predicate,
125        &triple.object,
126    )
127    .await
128    .map_err(KgWriteError::Admission)?;
129
130    // Read before the move — `assert` consumes `triple`.
131    let hot = crate::prompt_facts::is_hot_predicate(&triple.predicate);
132    handle
133        .kg
134        .assert(triple)
135        .await
136        .map_err(KgWriteError::Assert)?;
137
138    if hot && policy == CachePolicy::Inline {
139        refresh(state).await?;
140    }
141    Ok(AssertReceipt { hot })
142}
143
144/// Refresh the prompt cache after a [`CachePolicy::Batch`] loop.
145///
146/// Why: gives a batching caller the trailing half of [`assert_triple`] without
147/// letting it re-derive "was any of that hot?" from its own predicate list.
148/// What: no-op when `any_hot` is false; otherwise the same refresh
149/// [`assert_triple`] performs inline.
150/// Test: `kg_write_batch_policy_defers_refresh`.
151pub async fn refresh_after_batch(state: &AppState, any_hot: bool) -> Result<(), KgWriteError> {
152    if any_hot {
153        refresh(state).await
154    } else {
155        Ok(())
156    }
157}
158
159/// The single call to [`crate::prompt_facts::rebuild_prompt_cache`].
160///
161/// Note: this arm is not reachable today — `gather_hot_facts` logs and skips a
162/// palace it cannot read and then returns `Ok`, so the rebuild has no fallible
163/// step left. That skip silently truncates the cache (and the Tier S occupancy
164/// count that gates admission), which is a separate defect; the error is
165/// propagated here so that fixing it turns this into a reported failure rather
166/// than a silent one.
167async fn refresh(state: &AppState) -> Result<(), KgWriteError> {
168    crate::prompt_facts::rebuild_prompt_cache(state)
169        .await
170        .map_err(KgWriteError::CacheRefresh)
171}
172
173#[cfg(test)]
174mod tests {
175    use super::*;
176    use trusty_common::memory_core::palace::PalaceId;
177
178    /// Build an `AppState` over a fresh temp root with one created palace.
179    fn state_with_palace(name: &str) -> (AppState, Arc<PalaceHandle>) {
180        let tmp = tempfile::tempdir().expect("tempdir");
181        let root = tmp.path().to_path_buf();
182        std::mem::forget(tmp);
183        let state = AppState::new(root).with_default_palace(Some(name.to_string()));
184        let palace = trusty_common::memory_core::Palace {
185            id: PalaceId::new(name),
186            name: name.to_string(),
187            description: None,
188            created_at: chrono::Utc::now(),
189            data_dir: state.data_root.join(name),
190        };
191        let handle = state
192            .registry
193            .create_palace(&state.data_root, palace)
194            .expect("create palace");
195        (state, handle)
196    }
197
198    fn triple(subject: &str, predicate: &str, object: &str) -> Triple {
199        Triple {
200            subject: subject.to_string(),
201            predicate: predicate.to_string(),
202            object: object.to_string(),
203            valid_from: chrono::Utc::now(),
204            valid_to: None,
205            confidence: 1.0,
206            provenance: Some("test".to_string()),
207        }
208    }
209
210    /// The whole point of the module: a hot write reaches the prompt cache
211    /// without the caller doing anything beyond calling `assert_triple`.
212    #[tokio::test]
213    async fn kg_write_refreshes_cache_for_hot_predicate() {
214        let (state, handle) = state_with_palace("hotwrite");
215        let receipt = assert_triple(
216            &state,
217            &handle,
218            triple("rust", "has_convention", "no unwrap in library code"),
219            CachePolicy::Inline,
220        )
221        .await
222        .expect("assert");
223        assert!(receipt.hot);
224
225        let guard = state.prompt_context_cache.read().await;
226        assert!(
227            guard.formatted.contains("no unwrap in library code"),
228            "hot write missing from cache; got: {}",
229            guard.formatted
230        );
231    }
232
233    /// A cold predicate must not pay for a rebuild — the cache stays empty
234    /// rather than merely unchanged, proving the refresh was skipped and not
235    /// just uninformative.
236    #[tokio::test]
237    async fn kg_write_skips_refresh_for_cold_predicate() {
238        let (state, handle) = state_with_palace("coldwrite");
239        let receipt = assert_triple(
240            &state,
241            &handle,
242            triple("alice", "works_at", "Acme"),
243            CachePolicy::Inline,
244        )
245        .await
246        .expect("assert");
247        assert!(!receipt.hot);
248
249        let guard = state.prompt_context_cache.read().await;
250        assert!(
251            guard.triples.is_empty(),
252            "cold write should not populate the cache; got: {:?}",
253            guard.triples
254        );
255    }
256
257    /// Error arm: a refused write must leave BOTH storage and the cache
258    /// untouched. A gate that rejects while the row lands is the failure this
259    /// whole family keeps reproducing.
260    #[tokio::test]
261    async fn kg_write_admission_refusal_leaves_storage_and_cache_untouched() {
262        let (state, handle) = state_with_palace("refusal");
263        let over_long = "x".repeat(crate::prompt_facts::TIER_S_MAX_OBJECT_CHARS + 1);
264        let err = assert_triple(
265            &state,
266            &handle,
267            triple("subj", "has_convention", &over_long),
268            CachePolicy::Inline,
269        )
270        .await
271        .expect_err("over-long object must be refused");
272        assert!(
273            matches!(err, KgWriteError::Admission(_)),
274            "expected Admission, got: {err:?}"
275        );
276
277        let stored = handle.kg.query_active("subj").await.expect("query");
278        assert!(
279            stored.is_empty(),
280            "refused write reached storage: {stored:?}"
281        );
282        let guard = state.prompt_context_cache.read().await;
283        assert!(guard.triples.is_empty(), "refused write reached the cache");
284    }
285
286    /// `Batch` defers the refresh to `refresh_after_batch`, and the receipt is
287    /// what tells the caller a refresh is owed.
288    #[tokio::test]
289    async fn kg_write_batch_policy_defers_refresh() {
290        let (state, handle) = state_with_palace("batchwrite");
291        let receipt = assert_triple(
292            &state,
293            &handle,
294            triple("tm", "is_alias_for", "trusty-memory"),
295            CachePolicy::Batch,
296        )
297        .await
298        .expect("assert");
299        assert!(receipt.hot);
300        assert!(
301            state.prompt_context_cache.read().await.triples.is_empty(),
302            "Batch must not refresh inline"
303        );
304
305        refresh_after_batch(&state, receipt.hot)
306            .await
307            .expect("batch refresh");
308        let guard = state.prompt_context_cache.read().await;
309        assert!(
310            guard.formatted.contains("tm → trusty-memory"),
311            "batch refresh missed the write; got: {}",
312            guard.formatted
313        );
314    }
315
316    /// `refresh_after_batch(_, false)` must not touch the cache — otherwise a
317    /// cold-only batch would pay for a full registry walk.
318    #[tokio::test]
319    async fn kg_write_batch_refresh_is_a_noop_when_nothing_was_hot() {
320        let (state, handle) = state_with_palace("batchcold");
321        assert_triple(
322            &state,
323            &handle,
324            triple("bob", "lives_in", "Paris"),
325            CachePolicy::Batch,
326        )
327        .await
328        .expect("assert");
329        refresh_after_batch(&state, false).await.expect("noop");
330        assert!(state.prompt_context_cache.read().await.triples.is_empty());
331    }
332}