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}