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 "deduped": curate.deduped.len(),
205 "decayed": curate.decayed.len(),
206 "recovered": curate.recovered.len(),
207 "orphans": curate.orphans.len(),
208 "warnings": curate.warnings,
209 })
210 }
211
212 fn distill_batch(&self) -> Result<DistillBatchReport> {
213 let run_id = gen_uuid();
214 let now = utc_now_iso();
215
216 self.storage.begin_immediate()?;
218 let logs = match self
219 .storage
220 .claim_distill_batch(&run_id, self.distill_batch_size, &now)
221 {
222 Ok(l) => {
223 self.storage.commit()?;
224 l
225 }
226 Err(e) => {
227 let _ = self.storage.rollback();
228 return Err(e);
229 }
230 };
231
232 let mut chunks_by_log: HashMap<String, Vec<DistilledChunk>> = HashMap::new();
233 let mut failed_logs = HashSet::new();
234 let mut distill_errors = Vec::new();
235 for log in &logs {
236 let log_id = log.get("id").and_then(Value::as_str).unwrap_or("");
237 let context_key = log.get("context_key").and_then(Value::as_str);
238 let related_logs: Vec<Value> = logs
239 .iter()
240 .filter(|other| {
241 other.get("id").and_then(Value::as_str) == Some(log_id)
242 || (context_key.is_some()
243 && other.get("context_key").and_then(Value::as_str) == context_key)
244 })
245 .cloned()
246 .collect();
247 match self.distiller.distill_with_context(log, &related_logs) {
248 Ok(chunks) => {
249 if chunks.iter().any(|chunk| chunk.source_log_id != log_id) {
250 let error = "distiller returned a chunk for an unknown source log";
251 failed_logs.insert(log_id.to_string());
252 distill_errors.push(format!("{log_id}: {error}"));
253 self.finish_distill_log(
254 log_id,
255 "failed",
256 Some(&format!("distill_failed:{error}")),
257 estimate_distill_prompt_tokens(log, &related_logs),
258 0,
259 )?;
260 continue;
261 }
262 chunks_by_log.insert(log_id.to_string(), chunks);
263 }
264 Err(error) => {
265 let note = format!("distill_failed:{error}");
266 failed_logs.insert(log_id.to_string());
267 distill_errors.push(format!("{log_id}: {error}"));
268 self.finish_distill_log(
269 log_id,
270 "failed",
271 Some(¬e),
272 estimate_distill_prompt_tokens(log, &related_logs),
273 0,
274 )?;
275 }
276 }
277 }
278
279 let mut count = 0;
280 let provenance = self.distiller.provenance();
281 for log in &logs {
282 let log_id = log.get("id").and_then(Value::as_str).unwrap_or("");
283 if failed_logs.contains(log_id) {
284 continue;
285 }
286 let context_key = log.get("context_key").and_then(Value::as_str);
287 let related_logs: Vec<Value> = logs
288 .iter()
289 .filter(|other| {
290 other.get("id").and_then(Value::as_str) == Some(log_id)
291 || (context_key.is_some()
292 && other.get("context_key").and_then(Value::as_str) == context_key)
293 })
294 .cloned()
295 .collect();
296 let prompt_tokens = estimate_distill_prompt_tokens(log, &related_logs);
297 let chunks = chunks_by_log.remove(log_id).unwrap_or_default();
298 let completion_tokens = chunks
299 .iter()
300 .map(estimate_distilled_chunk_tokens)
301 .sum::<i64>();
302 if chunks.is_empty() {
303 self.finish_distill_log(
304 log_id,
305 "discarded",
306 Some("insufficient_material"),
307 prompt_tokens,
308 completion_tokens,
309 )?;
310 continue;
311 }
312
313 struct PreparedChunk {
318 row: ChunkRow,
319 cvec_bytes: Vec<u8>,
320 tvec_bytes: Vec<u8>,
321 }
322 let mut prepared: Vec<PreparedChunk> = Vec::with_capacity(chunks.len());
323 let mut embedding_failures = 0_usize;
324 for dc in chunks {
325 let (content, action) = self.sanitize_content(&dc.content);
326 if action == SanitizeAction::Discard {
327 continue; }
329 let h = content_hash(&content);
330 if self.storage.is_hash_invalidated(&h)? {
331 continue; }
333 let redacted = action == SanitizeAction::Redact;
334 let conf = if redacted { 0.4 } else { 0.55 };
335 let now2 = utc_now_iso();
336 let chunk_id = gen_uuid();
337 let tokens = estimate_tokens(&content) as i64;
338 let skill_name = dc
341 .skill_name
342 .clone()
343 .or_else(|| dc.trigger_desc.clone())
344 .filter(|s| !s.trim().is_empty());
345 let distill_provider = dc
348 .provider_override
349 .clone()
350 .or_else(|| provenance.provider.clone());
351 let row = ChunkRow {
352 id: chunk_id,
353 skill_name,
354 content: content.clone(),
355 trigger_desc: dc.trigger_desc.clone(),
356 anti_trigger_desc: dc.anti_trigger_desc,
357 content_hash: h,
358 token_count: Some(tokens),
359 origin: "distilled".to_string(),
360 agent: log
363 .get("agent")
364 .and_then(Value::as_str)
365 .map(str::to_string),
366 distilled_from: Some(dc.source_log_id),
367 distill_provider,
368 distill_model: provenance.model.clone(),
369 distill_prompt_version: provenance.prompt_version.clone(),
370 state: "pending".to_string(),
371 state_reason: Some("init:distilled".to_string()),
372 confidence: conf,
373 confidence_reason: Some("init:distilled".to_string()),
374 version: 1,
375 embed_version: 1,
376 created_at: now2.clone(),
377 updated_at: now2,
378 ..Default::default()
379 };
380 let (cvec_res, tvec_res) = self.embed_pair(
381 &content,
382 row.trigger_desc.as_deref().unwrap_or(&content),
383 "distill",
384 );
385 let cvec = match cvec_res {
386 Ok(v) if v.len() == self.embedding.content_dim() => v,
387 _ => {
391 embedding_failures += 1;
392 continue;
393 }
394 };
395 let tvec = match tvec_res {
396 Ok(v) if v.len() == self.embedding.trigger_dim() => v,
397 _ => {
398 embedding_failures += 1;
399 continue;
400 }
401 };
402 prepared.push(PreparedChunk {
403 row,
404 cvec_bytes: pack_embedding(&cvec),
405 tvec_bytes: pack_embedding(&tvec),
406 });
407 }
408
409 if prepared.is_empty() {
410 let note = if embedding_failures > 0 {
411 "embedding_failed"
412 } else {
413 "all_chunks_filtered"
414 };
415 self.finish_distill_log(
416 log_id,
417 if embedding_failures > 0 {
418 "failed"
419 } else {
420 "discarded"
421 },
422 Some(note),
423 prompt_tokens,
424 completion_tokens,
425 )?;
426 if embedding_failures > 0 {
427 failed_logs.insert(log_id.to_string());
428 }
429 continue;
430 }
431
432 let accounted_at = utc_now_iso();
434 self.storage.begin_immediate()?;
435 let write_result = (|| -> Result<()> {
436 for pc in &prepared {
437 self.storage.insert_chunk(&pc.row)?;
438 self.storage
439 .insert_vec_content(&pc.row.id, &pc.cvec_bytes)?;
440 self.storage
441 .insert_vec_trigger(&pc.row.id, &pc.tvec_bytes)?;
442 }
443 let note = (embedding_failures > 0)
444 .then(|| format!("partial_embedding_failures:{embedding_failures}"));
445 self.storage.finish_distill_log(
446 log_id,
447 "distilled",
448 note.as_deref(),
449 prompt_tokens,
450 completion_tokens,
451 &accounted_at,
452 )?;
453 self.storage.commit()
454 })();
455 if let Err(error) = write_result {
456 let _ = self.storage.rollback();
457 let note = format!("distill_write_failed:{error}");
458 self.finish_distill_log(
459 log_id,
460 "failed",
461 Some(¬e),
462 prompt_tokens,
463 completion_tokens,
464 )?;
465 failed_logs.insert(log_id.to_string());
466 continue;
467 }
468 count += 1;
469 }
470 if !distill_errors.is_empty() {
471 eprintln!(
475 "[innate] distillation partial failure ({} log(s)): {}",
476 distill_errors.len(),
477 distill_errors.join("; ")
478 );
479 }
480 Ok(DistillBatchReport {
481 distilled: count,
482 failed: failed_logs.len(),
483 })
484 }
485
486 fn finish_distill_log(
487 &self,
488 log_id: &str,
489 state: &str,
490 note: Option<&str>,
491 prompt_tokens: i64,
492 completion_tokens: i64,
493 ) -> Result<()> {
494 let accounted_at = utc_now_iso();
495 self.storage.begin_immediate()?;
496 let result = (|| -> Result<()> {
497 self.storage.finish_distill_log(
498 log_id,
499 state,
500 note,
501 prompt_tokens,
502 completion_tokens,
503 &accounted_at,
504 )?;
505 self.storage.commit()
506 })();
507 if result.is_err() {
508 let _ = self.storage.rollback();
509 }
510 result
511 }
512
513 pub(super) fn distill_token_period_start(&self, now: &str) -> Result<String> {
514 let window_hours = self
515 .storage
516 .get_meta("evolve.distill_token_window_hours")?
517 .and_then(|value| value.parse::<i64>().ok())
518 .unwrap_or(24)
519 .max(1);
520 Ok(hours_ago(now, window_hours))
521 }
522}