1use super::*;
2
3impl KnowledgeBase {
4 pub fn evolve(&self, trigger: &str) -> Result<Value> {
5 self.measure("evolve", None, None, || self.evolve_inner(trigger))
6 }
7
8 fn evolve_inner(&self, trigger: &str) -> Result<Value> {
9 if !matches!(trigger, "manual" | "scheduled" | "threshold") {
10 return Err(InnateError::InvalidState(format!(
11 "invalid evolve trigger: {trigger}"
12 )));
13 }
14 let evolve_started_at = utc_now_iso();
15 let retry_cutoff = minutes_ago(&evolve_started_at, 5);
16 let recovered_failed = self.storage.conn_execute_count(
17 "UPDATE episodic_log
18 SET distill_state='new', distill_note='retry_failed',
19 distill_locked_at=NULL, distill_run_id=NULL
20 WHERE distill_state='failed'
21 AND distill_attempts < 3
22 AND COALESCE(distill_last_failed_at, distill_accounted_at, ts) < ?",
23 rusqlite::params![retry_cutoff],
24 )?;
25 if recovered_failed > 0 {
26 self.storage
27 .request_evolve(&gen_uuid(), "distill_retry", &evolve_started_at)?;
28 }
29 if trigger == "scheduled" {
30 let age_cutoff = hours_ago(&evolve_started_at, self.evolve_schedule_interval_hours);
31 let aged_new = count_query_params(
32 &self.storage,
33 "SELECT COUNT(*) FROM episodic_log
34 WHERE distill_state='new' AND ts <= ?",
35 rusqlite::params![age_cutoff],
36 )?;
37 if aged_new > 0 {
38 self.storage
39 .request_evolve(&gen_uuid(), "scheduled", &evolve_started_at)?;
40 }
41 }
42 let request = self.storage.claim_evolve_request_with_reason(
43 &evolve_started_at,
44 &minutes_ago(&evolve_started_at, self.screening_timeout_minutes),
45 )?;
46 let request_id = request.as_ref().map(|claim| claim.id.as_str());
47 let request_reason = request.as_ref().map(|claim| claim.reason.as_str());
48 if trigger == "scheduled" && request_id.is_none() {
50 let curator = Arc::clone(&self.curator);
51 let curate = curator.run(self, &CurateScope::default())?;
52 return Ok(json!({
53 "distilled": 0,
54 "curate": self.format_curate_report(&curate),
55 "skipped": "no_evolve_request"
56 }));
57 }
58
59 if trigger == "threshold" {
61 let rows = self.storage.query_chunks(
62 "SELECT COUNT(*) AS cnt FROM episodic_log WHERE distill_state='new'",
63 )?;
64 let cnt = rows
65 .first()
66 .and_then(|r| r.get("cnt"))
67 .and_then(Value::as_i64)
68 .unwrap_or(0);
69 if cnt < self.evolve_threshold {
70 let curator = Arc::clone(&self.curator);
71 let curate = curator.run(self, &CurateScope::default())?;
72 if let Some(id) = request_id {
73 if matches!(request_reason, Some("governance" | "governance_ready")) {
74 self.storage.finish_evolve_request(
75 id,
76 "completed",
77 Some("curate_only"),
78 &utc_now_iso(),
79 )?;
80 } else {
81 self.storage.defer_evolve_request(
82 id,
83 "below_threshold",
84 &hours_after(
85 &evolve_started_at,
86 self.evolve_schedule_interval_hours.max(1),
87 ),
88 )?;
89 }
90 }
91 return Ok(json!({
92 "distilled": 0,
93 "curate": self.format_curate_report(&curate),
94 "skipped": "below_threshold"
95 }));
96 }
97 }
98
99 if trigger != "manual" {
100 if let Some(limit) = self
101 .storage
102 .get_meta("max_distill_tokens_per_period")?
103 .and_then(|value| value.parse::<i64>().ok())
104 .filter(|value| *value > 0)
105 {
106 let period_start = self.distill_token_period_start(&evolve_started_at)?;
107 let rows = self.storage.query_chunks_params(
108 "SELECT COALESCE(SUM(prompt_tokens + completion_tokens),0) AS used
109 FROM distill_token_usage
110 WHERE accounted_at >= ?",
111 rusqlite::params![period_start],
112 )?;
113 let used_tokens = rows
114 .first()
115 .and_then(|row| row.get("used"))
116 .and_then(Value::as_i64)
117 .unwrap_or(0);
118 if used_tokens >= limit {
119 let curator = Arc::clone(&self.curator);
120 let curate = curator.run(self, &CurateScope::default())?;
121 if let Some(id) = request_id {
122 if matches!(request_reason, Some("governance" | "governance_ready")) {
123 self.storage.finish_evolve_request(
124 id,
125 "completed",
126 Some("curate_only"),
127 &utc_now_iso(),
128 )?;
129 } else {
130 self.storage.defer_evolve_request(
131 id,
132 "distill_token_limit",
133 &hours_after(&evolve_started_at, 1),
134 )?;
135 }
136 }
137 return Ok(json!({
138 "distilled": 0,
139 "curate": self.format_curate_report(&curate),
140 "skipped": "distill_token_limit",
141 "distill_tokens_used": used_tokens,
142 "distill_token_limit": limit,
143 "period_start": period_start,
144 }));
145 }
146 }
147 }
148
149 let result = (|| -> Result<Value> {
150 let distill = self.measure("distill", None, None, || self.distill_batch())?;
151 let curator = Arc::clone(&self.curator);
152 let curate = curator.run(self, &CurateScope::default())?;
153 Ok(json!({
154 "distilled": distill.distilled,
155 "distill_failed": distill.failed,
156 "curate": self.format_curate_report(&curate),
157 }))
158 })();
159 if let Some(id) = request_id {
160 let (state, note) = match &result {
161 Ok(_) => ("completed", None),
162 Err(error) => ("failed", Some(error.to_string())),
163 };
164 self.storage
165 .finish_evolve_request(id, state, note.as_deref(), &utc_now_iso())?;
166 }
167 if result.is_ok() {
168 self.storage
169 .finish_covered_evolve_requests(&evolve_started_at, &utc_now_iso())?;
170 }
171
172 let failed_remaining = count_query(
173 &self.storage,
174 "SELECT COUNT(*) FROM episodic_log
175 WHERE distill_state='failed' AND distill_attempts < 3",
176 )?;
177 if failed_remaining > 0 {
178 self.storage.request_evolve_at(
179 &gen_uuid(),
180 "distill_retry",
181 &utc_now_iso(),
182 Some(&minutes_after(&utc_now_iso(), 5)),
183 )?;
184 }
185
186 if result.is_ok() {
188 let remaining = count_query(
189 &self.storage,
190 "SELECT COUNT(*) FROM episodic_log WHERE distill_state='new'",
191 )?;
192 if remaining > 0 {
193 let _ = self
194 .storage
195 .request_evolve(&gen_uuid(), "batch_continue", &utc_now_iso());
196 }
197 }
198 result
199 }
200
201 fn format_curate_report(&self, curate: &CurateReport) -> Value {
202 json!({
203 "archived": curate.archived.len(),
204 "promoted": curate.promoted.len(),
205 "deduped": curate.deduped.len(),
206 "decayed": curate.decayed.len(),
207 "recovered": curate.recovered.len(),
208 "orphans": curate.orphans.len(),
209 "warnings": curate.warnings,
210 })
211 }
212
213 fn distill_batch(&self) -> Result<DistillBatchReport> {
214 let run_id = gen_uuid();
215 let now = utc_now_iso();
216
217 self.storage.begin_immediate()?;
219 let logs = match self
220 .storage
221 .claim_distill_batch(&run_id, self.distill_batch_size, &now)
222 {
223 Ok(l) => {
224 self.storage.commit()?;
225 l
226 }
227 Err(e) => {
228 let _ = self.storage.rollback();
229 return Err(e);
230 }
231 };
232
233 let mut chunks_by_log: HashMap<String, Vec<DistilledChunk>> = HashMap::new();
234 let mut failed_logs = HashSet::new();
235 let mut distill_errors = Vec::new();
236 for log in &logs {
237 let log_id = log.get("id").and_then(Value::as_str).unwrap_or("");
238 let context_key = log.get("context_key").and_then(Value::as_str);
239 let related_logs: Vec<Value> = logs
240 .iter()
241 .filter(|other| {
242 other.get("id").and_then(Value::as_str) == Some(log_id)
243 || (context_key.is_some()
244 && other.get("context_key").and_then(Value::as_str) == context_key)
245 })
246 .cloned()
247 .collect();
248 match self.distiller.distill_with_context(log, &related_logs) {
249 Ok(chunks) => {
250 if chunks.iter().any(|chunk| chunk.source_log_id != log_id) {
251 let error = "distiller returned a chunk for an unknown source log";
252 failed_logs.insert(log_id.to_string());
253 distill_errors.push(format!("{log_id}: {error}"));
254 self.finish_distill_log(
255 log_id,
256 "failed",
257 Some(&format!("distill_failed:{error}")),
258 estimate_distill_prompt_tokens(log, &related_logs),
259 0,
260 )?;
261 continue;
262 }
263 chunks_by_log.insert(log_id.to_string(), chunks);
264 }
265 Err(error) => {
266 let note = format!("distill_failed:{error}");
267 failed_logs.insert(log_id.to_string());
268 distill_errors.push(format!("{log_id}: {error}"));
269 self.finish_distill_log(
270 log_id,
271 "failed",
272 Some(¬e),
273 estimate_distill_prompt_tokens(log, &related_logs),
274 0,
275 )?;
276 }
277 }
278 }
279
280 let mut count = 0;
281 let provenance = self.distiller.provenance();
282 for log in &logs {
283 let log_id = log.get("id").and_then(Value::as_str).unwrap_or("");
284 if failed_logs.contains(log_id) {
285 continue;
286 }
287 let context_key = log.get("context_key").and_then(Value::as_str);
288 let related_logs: Vec<Value> = logs
289 .iter()
290 .filter(|other| {
291 other.get("id").and_then(Value::as_str) == Some(log_id)
292 || (context_key.is_some()
293 && other.get("context_key").and_then(Value::as_str) == context_key)
294 })
295 .cloned()
296 .collect();
297 let prompt_tokens = estimate_distill_prompt_tokens(log, &related_logs);
298 let chunks = chunks_by_log.remove(log_id).unwrap_or_default();
299 let completion_tokens = chunks
300 .iter()
301 .map(estimate_distilled_chunk_tokens)
302 .sum::<i64>();
303 if chunks.is_empty() {
304 self.finish_distill_log(
305 log_id,
306 "discarded",
307 Some("insufficient_material"),
308 prompt_tokens,
309 completion_tokens,
310 )?;
311 continue;
312 }
313
314 struct PreparedChunk {
319 row: ChunkRow,
320 cvec_bytes: Vec<u8>,
321 tvec_bytes: Vec<u8>,
322 }
323 let mut prepared: Vec<PreparedChunk> = Vec::with_capacity(chunks.len());
324 let mut embedding_failures = 0_usize;
325 for dc in chunks {
326 let (content, action) = self.sanitize_content(&dc.content);
327 if action == SanitizeAction::Discard {
328 continue; }
330 let h = content_hash(&content);
331 if self.storage.is_hash_invalidated(&h)? {
332 continue; }
334 let redacted = action == SanitizeAction::Redact;
335 let conf = if redacted { 0.4 } else { 0.55 };
336 let now2 = utc_now_iso();
337 let chunk_id = gen_uuid();
338 let tokens = estimate_tokens(&content) as i64;
339 let skill_name = dc
342 .skill_name
343 .clone()
344 .or_else(|| dc.trigger_desc.clone())
345 .filter(|s| !s.trim().is_empty());
346 let distill_provider = dc
349 .provider_override
350 .clone()
351 .or_else(|| provenance.provider.clone());
352 let row = ChunkRow {
353 id: chunk_id,
354 skill_name,
355 content: content.clone(),
356 trigger_desc: dc.trigger_desc.clone(),
357 anti_trigger_desc: dc.anti_trigger_desc,
358 content_hash: h,
359 token_count: Some(tokens),
360 origin: "distilled".to_string(),
361 agent: log
364 .get("agent")
365 .and_then(Value::as_str)
366 .map(str::to_string),
367 distilled_from: Some(dc.source_log_id),
368 distill_provider,
369 distill_model: provenance.model.clone(),
370 distill_prompt_version: provenance.prompt_version.clone(),
371 state: "pending".to_string(),
372 state_reason: Some("init:distilled".to_string()),
373 confidence: conf,
374 confidence_reason: Some("init:distilled".to_string()),
375 version: 1,
376 embed_version: 1,
377 created_at: now2.clone(),
378 updated_at: now2,
379 ..Default::default()
380 };
381 let (cvec_res, tvec_res) = self.embed_pair(
382 &content,
383 row.trigger_desc.as_deref().unwrap_or(&content),
384 "distill",
385 );
386 let cvec = match cvec_res {
387 Ok(v) if v.len() == self.embedding.content_dim() => v,
388 _ => {
392 embedding_failures += 1;
393 continue;
394 }
395 };
396 let tvec = match tvec_res {
397 Ok(v) if v.len() == self.embedding.trigger_dim() => v,
398 _ => {
399 embedding_failures += 1;
400 continue;
401 }
402 };
403 prepared.push(PreparedChunk {
404 row,
405 cvec_bytes: pack_embedding(&cvec),
406 tvec_bytes: pack_embedding(&tvec),
407 });
408 }
409
410 if prepared.is_empty() {
411 let note = if embedding_failures > 0 {
412 "embedding_failed"
413 } else {
414 "all_chunks_filtered"
415 };
416 self.finish_distill_log(
417 log_id,
418 if embedding_failures > 0 {
419 "failed"
420 } else {
421 "discarded"
422 },
423 Some(note),
424 prompt_tokens,
425 completion_tokens,
426 )?;
427 if embedding_failures > 0 {
428 failed_logs.insert(log_id.to_string());
429 }
430 continue;
431 }
432
433 let accounted_at = utc_now_iso();
435 self.storage.begin_immediate()?;
436 let write_result = (|| -> Result<()> {
437 for pc in &prepared {
438 self.storage.insert_chunk(&pc.row)?;
439 self.storage
440 .insert_vec_content(&pc.row.id, &pc.cvec_bytes)?;
441 self.storage
442 .insert_vec_trigger(&pc.row.id, &pc.tvec_bytes)?;
443 }
444 let note = (embedding_failures > 0)
445 .then(|| format!("partial_embedding_failures:{embedding_failures}"));
446 self.storage.finish_distill_log(
447 log_id,
448 "distilled",
449 note.as_deref(),
450 prompt_tokens,
451 completion_tokens,
452 &accounted_at,
453 )?;
454 self.storage.commit()
455 })();
456 if let Err(error) = write_result {
457 let _ = self.storage.rollback();
458 let note = format!("distill_write_failed:{error}");
459 self.finish_distill_log(
460 log_id,
461 "failed",
462 Some(¬e),
463 prompt_tokens,
464 completion_tokens,
465 )?;
466 failed_logs.insert(log_id.to_string());
467 continue;
468 }
469 count += 1;
470 }
471 if !distill_errors.is_empty() {
472 eprintln!(
476 "[innate] distillation partial failure ({} log(s)): {}",
477 distill_errors.len(),
478 distill_errors.join("; ")
479 );
480 }
481 Ok(DistillBatchReport {
482 distilled: count,
483 failed: failed_logs.len(),
484 })
485 }
486
487 fn finish_distill_log(
488 &self,
489 log_id: &str,
490 state: &str,
491 note: Option<&str>,
492 prompt_tokens: i64,
493 completion_tokens: i64,
494 ) -> Result<()> {
495 let accounted_at = utc_now_iso();
496 self.storage.begin_immediate()?;
497 let result = (|| -> Result<()> {
498 self.storage.finish_distill_log(
499 log_id,
500 state,
501 note,
502 prompt_tokens,
503 completion_tokens,
504 &accounted_at,
505 )?;
506 self.storage.commit()
507 })();
508 if result.is_err() {
509 let _ = self.storage.rollback();
510 }
511 result
512 }
513
514 pub(super) fn distill_token_period_start(&self, now: &str) -> Result<String> {
515 let window_hours = self
516 .storage
517 .get_meta("evolve.distill_token_window_hours")?
518 .and_then(|value| value.parse::<i64>().ok())
519 .unwrap_or(24)
520 .max(1);
521 Ok(hours_ago(now, window_hours))
522 }
523}