Skip to main content

innate_core/kb/record/
mod.rs

1use super::*;
2
3mod evidence;
4
5/// Parameters for [`KnowledgeBase::record`].
6///
7/// Borrowed, `Default`-able: construct with `RecordParams { trace_id, source, ..Default::default() }`
8/// and set only the fields you need. Behavioural defaults are applied inside `record`:
9/// `used_attribution` empty → `"explicit"`, `feedback_kind` empty → `"user"`,
10/// `used_complete` `None` → `true` (a complete snapshot).
11#[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    /// 方案 H — 反事实审查:此 trace 若来自一次 appraise,且 actor **因警告回避了动作**
31    /// (没有真正采取该步,outcome 反映的是回避后的世界),则该结果不能当作直觉对错的
32    /// 证据 —— 标记为 `counterfactual_censored`,不计入校准映射 / ECE。默认 `false`
33    /// (actor 实际采取动作并观测到结果 → `observed`,计入校准)。
34    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            // usage_trace: used
182            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                // The parameter is validated above; this arm is only reachable if the
194                // stored used_attribution was tampered with outside the API.
195                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            // usage_trace: task_ok / task_fail
258            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                    // 方案 B/H:回填 verdict_log 的实际结果(若该 trace 是一次 appraise)。
272                    // observed_outcome = +1 坏结果发生(fail) / -1 好结果(ok)。
273                    // provenance='observed':actor 实际采取动作并观测到结果,计入校准。
274                    // verdict_heeded=true:因警告回避了动作,outcome 是反事实,不计校准
275                    // (counterfactual_censored,见原则 3 / 方案 H)。
276                    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            // Rebuild trace-derived evidence whenever either side of the pair arrives.
292            // This makes `outcome → used` and `used → outcome` equivalent and lets a
293            // later complete usage declaration replace an earlier one safely.
294            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            // Persist feedback facts before reducing them into confidence.
354            // INSERT OR IGNORE: skip all derived updates for duplicate (trace_id, chunk_id, signal).
355            // Track affected chunks so we only rebuild their context_stats (not the full table).
356            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            // Targeted rebuild — only update context_stats for chunks touched in this call.
454            self.rebuild_context_stats_for(&context_affected, &now)?;
455
456            // Fill in content fields (補写: output_summary, nomination, output, query) on existing log.
457            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            // Update episodic log
494            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            // Distillation eligibility depends on the session having FINISHED and
511            // left material behind — not on it having produced a confidence signal.
512            // A capture channel that cannot judge success (the Stop hook emits
513            // outcome="unknown" by design) still yields a usable summary, and that
514            // summary is worth distilling. Keeping the two coupled meant such a
515            // session was discarded outright. Confidence is untouched by this: it
516            // still moves only under `effective_outcome` ∈ {ok, fail} above, so an
517            // unknown outcome feeds the distiller without ever inflating a score.
518            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}