1use super::*;
2
3mod evidence;
4
5#[derive(Debug, Clone, Default)]
12pub struct RecordParams<'a> {
13 pub trace_id: &'a str,
14 pub query: Option<&'a str>,
15 pub output: Option<&'a str>,
16 pub output_summary: Option<&'a str>,
17 pub outcome: Option<&'a str>,
18 pub used: Option<&'a [String]>,
19 pub used_attribution: &'a str,
20 pub used_complete: Option<bool>,
21 pub feedback_up: Option<&'a [String]>,
22 pub feedback_down: Option<&'a [String]>,
23 pub feedback_kind: &'a str,
24 pub feedback_actor: Option<&'a str>,
25 pub feedback_reason: Option<&'a str>,
26 pub nomination: Option<&'a str>,
27 pub priority: i64,
28 pub task_state: Option<&'a str>,
29 pub source: &'a str,
30 pub verdict_heeded: bool,
35}
36
37impl KnowledgeBase {
38 pub fn record(&self, params: RecordParams<'_>) -> Result<()> {
39 let tid = params.trace_id.to_string();
40 let src = params.source.to_string();
41 self.measure("record", Some(&src), Some(&tid), || self.record_inner(params))
42 }
43
44 fn record_inner(&self, params: RecordParams<'_>) -> Result<()> {
45 let RecordParams {
46 trace_id,
47 query,
48 output,
49 output_summary,
50 outcome,
51 used,
52 used_attribution,
53 used_complete,
54 feedback_up,
55 feedback_down,
56 feedback_kind,
57 feedback_actor,
58 feedback_reason,
59 nomination,
60 priority,
61 task_state,
62 source,
63 verdict_heeded,
64 } = params;
65 let used_attribution = if used_attribution.is_empty() {
66 "explicit"
67 } else {
68 used_attribution
69 };
70 let feedback_kind = if feedback_kind.is_empty() {
71 "user"
72 } else {
73 feedback_kind
74 };
75 let used_complete = used_complete.unwrap_or(true);
76 let dedupe_ids = |ids: &[String]| {
77 let mut seen = HashSet::new();
78 ids.iter()
79 .filter(|id| seen.insert((*id).clone()))
80 .cloned()
81 .collect::<Vec<_>>()
82 };
83 let normalized_used = used.map(dedupe_ids);
84 let normalized_feedback_up = feedback_up.map(dedupe_ids);
85 let normalized_feedback_down = feedback_down.map(dedupe_ids);
86 let used = normalized_used.as_deref();
87 let feedback_up = normalized_feedback_up.as_deref();
88 let feedback_down = normalized_feedback_down.as_deref();
89
90 if let Some(o) = outcome {
91 if !matches!(o, "ok" | "fail" | "unknown") {
92 return Err(InnateError::InvalidState(format!("invalid outcome: {o}")));
93 }
94 }
95 if !matches!(used_attribution, "explicit" | "cited" | "inferred") {
96 return Err(InnateError::InvalidState(format!(
97 "invalid used attribution: {used_attribution}"
98 )));
99 }
100 if !matches!(feedback_kind, "user" | "judge") {
101 return Err(InnateError::InvalidState(format!(
102 "invalid feedback kind: {feedback_kind}"
103 )));
104 }
105 if let Some(state) = task_state {
106 if !matches!(
107 state,
108 "recalled" | "running" | "completed" | "abandoned" | "timed_out"
109 ) {
110 return Err(InnateError::InvalidState(format!(
111 "invalid task state: {state}"
112 )));
113 }
114 }
115 validate_source(source)?;
116 if let (Some(ups), Some(downs)) = (feedback_up, feedback_down) {
117 let down_set: HashSet<&str> = downs.iter().map(String::as_str).collect();
118 if let Some(chunk_id) = ups.iter().find(|id| down_set.contains(id.as_str())) {
119 return Err(InnateError::InvalidState(format!(
120 "conflicting feedback for chunk {chunk_id}"
121 )));
122 }
123 }
124 let effective_priority = if nomination.is_some() && priority == 0 {
125 1
126 } else {
127 priority
128 };
129 let now = utc_now_iso();
130 let lib_id = self.storage.lib_id()?;
131
132 self.storage.begin_immediate()?;
133 let result = (|| -> Result<()> {
134 let log = self.storage.get_episodic_log(trace_id)?;
135 let mut is_fresh_insert = false;
136 let log = match log {
137 Some(l) => l,
138 None => {
139 let used_ids = used.map(serde_json::to_string).transpose()?;
140 let row = EpisodicLogRow {
141 id: gen_uuid(),
142 trace_id: trace_id.to_string(),
143 lib_id,
144 ts: now.clone(),
145 query: query.map(str::to_string).or_else(|| Some(String::new())),
146 output: output.map(str::to_string),
147 output_summary: output_summary.map(str::to_string),
148 outcome: outcome.map(str::to_string),
149 event_source: source.to_string(),
150 agent: agent_source(),
151 task_state: if matches!(outcome, Some("ok") | Some("fail")) {
152 "completed".to_string()
153 } else {
154 task_state.unwrap_or("running").to_string()
155 },
156 completed_at: matches!(outcome, Some("ok") | Some("fail"))
157 .then(|| now.clone()),
158 usage_state: usage_state(used).to_string(),
159 used_ids,
160 used_attribution: used.map(|_| used_attribution.to_string()),
161 used_complete,
162 context_key: query.map(|q| content_hash(&normalize_query(q))),
163 nomination: nomination.map(str::to_string),
164 priority: effective_priority,
165 distill_state: "open".to_string(),
166 ..Default::default()
167 };
168 self.storage.upsert_episodic_log(&row)?;
169 is_fresh_insert = true;
170 self.storage.get_episodic_log(trace_id)?.unwrap()
171 }
172 };
173 self.validate_trace_attribution(trace_id, used, "used")?;
174 self.validate_trace_attribution(trace_id, feedback_up, "feedback_up")?;
175 self.validate_trace_attribution(trace_id, feedback_down, "feedback_down")?;
176
177 let existing_outcome = log
178 .get("outcome")
179 .and_then(Value::as_str)
180 .map(str::to_string);
181 let effective_used_attribution = if used.is_some() {
183 used_attribution
184 } else {
185 log.get("used_attribution")
186 .and_then(Value::as_str)
187 .unwrap_or(used_attribution)
188 };
189 let used_strength = match effective_used_attribution {
190 "explicit" => 0.3,
191 "cited" => 0.25,
192 "inferred" => 0.15,
193 other => {
196 return Err(InnateError::InvalidState(format!(
197 "invalid stored used attribution: {other}"
198 )))
199 }
200 };
201 let existing_used_ids: Vec<String> = log
202 .get("used_ids")
203 .and_then(Value::as_str)
204 .and_then(|raw| serde_json::from_str(raw).ok())
205 .unwrap_or_default();
206 let existing_used_complete = log
207 .get("used_complete")
208 .and_then(Value::as_i64)
209 .unwrap_or(0)
210 != 0;
211 let effective_used_complete = used_complete || existing_used_complete;
212 let effective_used_ids = used.map(|reported| {
213 if used_complete {
214 reported.to_vec()
215 } else {
216 let mut merged = existing_used_ids.clone();
217 let mut seen: HashSet<String> = merged.iter().cloned().collect();
218 merged.extend(
219 reported
220 .iter()
221 .filter(|id| seen.insert((*id).clone()))
222 .cloned(),
223 );
224 merged
225 }
226 });
227 if let Some(used_ids) = effective_used_ids.as_deref() {
228 let previously_used: HashSet<String> = existing_used_ids.iter().cloned().collect();
229 if used_complete {
230 self.storage.replace_used_trace(
231 trace_id,
232 used_ids,
233 used_strength,
234 used_attribution,
235 source,
236 &now,
237 )?;
238 } else if let Some(reported) = used {
239 self.storage.merge_used_trace(
240 trace_id,
241 reported,
242 used_strength,
243 used_attribution,
244 source,
245 &now,
246 )?;
247 }
248 let affected: HashSet<String> = previously_used
249 .into_iter()
250 .chain(used_ids.iter().cloned())
251 .collect();
252 for cid in affected {
253 self.storage.refresh_chunk_last_used(&cid, &now)?;
254 }
255 }
256
257 if let Some(o) = outcome {
259 if matches!(o, "ok" | "fail") {
260 let event = if o == "ok" { "task_ok" } else { "task_fail" };
261 let strength = if event == "task_fail" { 0.15 } else { 1.0 };
262 self.storage.conn_execute(
263 "DELETE FROM usage_trace
264 WHERE trace_id=? AND event IN ('task_ok','task_fail')
265 AND chunk_id IS NULL",
266 rusqlite::params![trace_id],
267 )?;
268 self.storage.insert_usage_trace(
269 trace_id, None, event, strength, None, None, None, None, None, source, &now,
270 )?;
271 let observed_outcome = if event == "task_fail" { 1.0 } else { -1.0 };
277 let provenance = if verdict_heeded {
278 "counterfactual_censored"
279 } else {
280 "observed"
281 };
282 self.storage.backfill_verdict_outcome(
283 trace_id,
284 observed_outcome,
285 provenance,
286 &now,
287 )?;
288 }
289 }
290
291 let effective_outcome =
295 outcome
296 .filter(|value| *value != "unknown")
297 .or(existing_outcome
298 .as_deref()
299 .filter(|value| *value != "unknown"));
300 if let Some(o @ ("ok" | "fail")) = effective_outcome {
301 if used.is_some()
302 || (outcome.is_some_and(|value| value != "unknown")
303 && existing_outcome.as_deref() != outcome)
304 {
305 let fallback_ids: Vec<String>;
306 let effective_used: Option<&[String]> = if effective_used_ids.is_some() {
307 effective_used_ids.as_deref()
308 } else {
309 fallback_ids = log
310 .get("used_ids")
311 .and_then(Value::as_str)
312 .and_then(|s| serde_json::from_str(s).ok())
313 .unwrap_or_default();
314 if fallback_ids.is_empty() {
315 None
316 } else {
317 Some(&fallback_ids)
318 }
319 };
320 let effective_complete = if used.is_some() {
321 effective_used_complete
322 } else {
323 log.get("usage_state").and_then(Value::as_str) != Some("unknown")
324 && log
325 .get("used_complete")
326 .and_then(Value::as_i64)
327 .unwrap_or(1)
328 != 0
329 };
330 self.replace_outcome_evidence(
331 trace_id,
332 o,
333 effective_used,
334 effective_complete,
335 &now,
336 )?;
337 }
338 } else if used.is_some() && effective_used_complete {
339 self.replace_selected_unused_evidence(
340 trace_id,
341 effective_used_ids.as_deref().unwrap_or_default(),
342 &now,
343 )?;
344 }
345
346 let context_key = log
347 .get("context_key")
348 .and_then(Value::as_str)
349 .map(str::to_string)
350 .or_else(|| query.map(|q| content_hash(&normalize_query(q))));
351 let feedback_strength = if feedback_kind == "judge" { 0.6 } else { 1.0 };
352
353 let mut context_affected: HashSet<String> = HashSet::new();
357 if let Some(used_ids) = effective_used_ids.as_deref() {
358 for cid in used_ids {
359 context_affected.insert(cid.clone());
360 }
361 }
362 if let Some(ups) = feedback_up {
363 for cid in ups {
364 let corrected = self.storage.delete_feedback_event(trace_id, cid, "down")?;
365 self.storage.delete_chunk_trace_confidence_evidence(
366 trace_id,
367 cid,
368 "feedback_down",
369 )?;
370 let inserted = self.storage.insert_feedback_event(
371 &gen_uuid(),
372 trace_id,
373 cid,
374 "up",
375 feedback_strength,
376 source,
377 feedback_actor,
378 feedback_reason,
379 context_key.as_deref(),
380 &now,
381 )?;
382 if inserted > 0 {
383 self.upsert_trace_confidence_evidence(
384 trace_id,
385 cid,
386 "feedback_up",
387 1.0,
388 feedback_strength,
389 if feedback_kind == "judge" {
390 "judge_up"
391 } else {
392 "user_up"
393 },
394 context_key.as_deref(),
395 &now,
396 true,
397 )?;
398 self.storage.update_chunk_last_used(cid, &now)?;
399 self.refresh_governance_evidence(cid, &now)?;
400 context_affected.insert(cid.clone());
401 } else if corrected > 0 {
402 self.recompute_chunk_confidence(cid, &now)?;
403 self.refresh_governance_evidence(cid, &now)?;
404 context_affected.insert(cid.clone());
405 }
406 }
407 }
408 if let Some(downs) = feedback_down {
409 for cid in downs {
410 let corrected = self.storage.delete_feedback_event(trace_id, cid, "up")?;
411 self.storage.delete_chunk_trace_confidence_evidence(
412 trace_id,
413 cid,
414 "feedback_up",
415 )?;
416 let inserted = self.storage.insert_feedback_event(
417 &gen_uuid(),
418 trace_id,
419 cid,
420 "down",
421 feedback_strength,
422 source,
423 feedback_actor,
424 feedback_reason,
425 context_key.as_deref(),
426 &now,
427 )?;
428 if inserted > 0 {
429 self.upsert_trace_confidence_evidence(
430 trace_id,
431 cid,
432 "feedback_down",
433 0.0,
434 feedback_strength,
435 if feedback_kind == "judge" {
436 "judge_down"
437 } else {
438 "user_down"
439 },
440 context_key.as_deref(),
441 &now,
442 true,
443 )?;
444 self.refresh_governance_evidence(cid, &now)?;
445 context_affected.insert(cid.clone());
446 } else if corrected > 0 {
447 self.recompute_chunk_confidence(cid, &now)?;
448 self.refresh_governance_evidence(cid, &now)?;
449 context_affected.insert(cid.clone());
450 }
451 }
452 }
453 self.rebuild_context_stats_for(&context_affected, &now)?;
455
456 if !is_fresh_insert {
458 self.storage.patch_episodic_log_content(
459 trace_id,
460 query,
461 output,
462 output_summary,
463 nomination,
464 effective_priority,
465 )?;
466 }
467
468 let lifecycle_state = if effective_outcome.is_some() {
469 "completed"
470 } else {
471 task_state.unwrap_or_else(|| {
472 log.get("task_state")
473 .and_then(Value::as_str)
474 .unwrap_or("running")
475 })
476 };
477 let used_ids_json = effective_used_ids
478 .as_deref()
479 .map(serde_json::to_string)
480 .transpose()?;
481 self.storage.update_trace_lifecycle(
482 trace_id,
483 lifecycle_state,
484 (lifecycle_state == "completed").then_some(now.as_str()),
485 effective_used_ids
486 .as_deref()
487 .map(|ids| usage_state(Some(ids))),
488 used_ids_json.as_deref(),
489 used.map(|_| used_attribution),
490 used.map(|_| effective_used_complete),
491 )?;
492
493 let current_state = log
495 .get("distill_state")
496 .and_then(Value::as_str)
497 .unwrap_or("open");
498 let lifecycle_completed = lifecycle_state == "completed";
499 let has_material = output_summary.is_some()
500 || nomination.is_some()
501 || output.is_some()
502 || log.get("output_summary").and_then(Value::as_str).is_some()
503 || log.get("nomination").and_then(Value::as_str).is_some()
504 || log.get("output").and_then(Value::as_str).is_some();
505 let retryable_discard = current_state == "discarded"
506 && matches!(
507 log.get("distill_note").and_then(Value::as_str),
508 Some("insufficient_material" | "abandoned" | "timed_out")
509 );
510 let lifecycle_finished =
519 lifecycle_completed || matches!(lifecycle_state, "abandoned" | "timed_out");
520 let new_state = if lifecycle_finished
521 && (current_state == "open" || retryable_discard)
522 {
523 if has_material {
524 Some("new")
525 } else {
526 Some("discarded")
527 }
528 } else {
529 None
530 };
531 if let Some(state) = new_state {
532 let note = if state == "discarded" {
533 Some(if matches!(lifecycle_state, "abandoned" | "timed_out") {
534 lifecycle_state
535 } else {
536 "insufficient_material"
537 })
538 } else {
539 None
540 };
541 let outcome_str = outcome.map(str::to_string);
542 self.storage.update_episodic_log_state(
543 trace_id,
544 state,
545 note,
546 outcome_str.as_deref(),
547 )?;
548 if state == "new" && retryable_discard {
549 self.storage.conn_execute(
550 "UPDATE episodic_log SET distill_note=NULL WHERE trace_id=?",
551 rusqlite::params![trace_id],
552 )?;
553 }
554 } else if outcome.is_some() {
555 let outcome_str = outcome.map(str::to_string);
556 self.storage.update_episodic_log_state(
557 trace_id,
558 current_state,
559 None,
560 outcome_str.as_deref(),
561 )?;
562 }
563
564 self.storage.commit()
565 })();
566 if result.is_err() {
567 let _ = self.storage.rollback();
568 }
569 result?;
570 self.enqueue_evolve_if_needed(&now)?;
571 Ok(())
572 }
573}