1#![forbid(unsafe_code)]
4
5mod services;
6
7pub use kcode_kennedy_session_objects::ResolvedObject;
8pub use kcode_telegram_session_coordinator::validate_file_name as validate_delivery_file_name;
9pub use services::{Api as Service, LocalServices as Capabilities};
10
11use std::{
12 collections::{BTreeMap, HashMap, HashSet},
13 future::Future,
14 time::Duration,
15};
16
17use anyhow::Context as _;
18use chrono::{DateTime, Utc};
19use kcode_commit_session::{CommitReceipt, CommitRequest, PlannedNode};
20use kcode_dev_tools::{
21 ATTACH_OBJECT_WEB_LIB_TOOL, CALL_RUST_BIN_TOOL, RUST_BIN_TOOLS, RUST_LIB_TOOLS, WEB_LIB_TOOLS,
22 WRITE_RUST_BIN_TOOL, WRITE_RUST_LIB_TOOL, WRITE_WEB_LIB_TOOL, proposed_write_snapshot,
23};
24use kcode_dev_tools_chatend::{
25 FreeformWrite, SourceSnapshot, apply_snapshot, decode_freeform_write, prepare_freeform_write,
26};
27use kcode_history_ingress_context::{
28 Outcome as HistoryIngressContextOutcome, RecoveryOutcome as ContextRecoveryOutcome,
29};
30use kcode_kennedy_kweb_loader::{load_durable_batch, node_from_value};
31use kcode_kennedy_session_presentation::{RenderRequest, render};
32use kcode_kennedy_session_tool_contracts::{
33 DecodedTool, ManagedObjectArguments, ValidationRequest, decode, decode_managed_objects,
34 validate,
35};
36use kcode_kennedy_subagent_context::Context as SubagentContext;
37use kcode_kweb_context::{
38 Context as KwebContext, Node as KwebNode, NodeDraft, StagedCreate as KwebStagedCreate,
39};
40use kcode_kweb_db::NodeId;
41use kcode_server_object_envelopes::encode_file;
42use kcode_session_history::{
43 NewSession, Session as HistorySession,
44 chatend::{BoxContent, BoxId, BoxOwner, EventId, EventKind, SessionKind, SessionMetadata},
45};
46use kcode_speaker_system::KTOOLS as SPEECH_CLASSIFICATION_TOOLS;
47use serde::{Deserialize, Serialize};
48use serde_json::{Value, json};
49use sha2::{Digest, Sha256};
50use uuid::Uuid;
51
52const AGENT_LOOP_ROUND_LIMIT: u64 = 100;
53const BROWSER_CONVERSATION_REQUEST_TIMEOUT: Duration = Duration::from_secs(90 * 60);
54const HISTORY_INGRESS_REQUEST_TIMEOUT: Duration = Duration::from_secs(90 * 60);
55const WAKEUP_REQUEST_TIMEOUT: Duration = Duration::from_secs(90 * 60);
56const MAX_MEDIA_ENRICHMENT_BYTES: u64 = 20 * 1024 * 1024;
57const KWEB_TOOL_INSTANCE: &str = "kweb";
58const TASK_BOARD_TOOLS: [&str; 8] = [
59 "CreateTaskCategory",
60 "GetTaskCategory",
61 "RemoveTaskCategory",
62 "CreateTask",
63 "GetTask",
64 "UpdateTask",
65 "RemoveTask",
66 "GetTopTaskOrphan",
67];
68const CONTEXT_OVERFLOW_WARNING_BOX_NAME: &str = "Context overflow warning";
69const CONTEXT_OVERFLOW_WARNING: &str = "Context size was exceeded, some context has been dehydrated. The session is now at risk of destabilizing, please perform any cleanup tasks and end the session";
70const INGRESS_FORCE_COMMIT_NOTE: &str = "ingress_force_commit";
71
72#[derive(Clone, Debug)]
74pub struct RuntimeModel {
75 pub model: String,
76 pub reasoning_effort: String,
77 pub context_window_tokens: u64,
78}
79
80impl RuntimeModel {
81 pub fn from_intelligence(runtime: kcode_intelligence_router::RuntimeModel) -> Self {
82 Self {
83 model: runtime.model,
84 reasoning_effort: runtime.reasoning_effort,
85 context_window_tokens: runtime.context_window_tokens,
86 }
87 }
88
89 fn attribution(&self) -> String {
90 format!("{}-{}", self.model, self.reasoning_effort)
91 }
92}
93
94#[derive(Clone, Debug, PartialEq, Eq)]
95pub enum AgentMode {
96 Conversation,
97 FreeTime,
98 Wakeup,
99 Ingress { record_id: Option<String> },
100}
101
102#[derive(Clone, Debug)]
103pub struct SessionOptions {
104 pub session_type: String,
105 pub root_node_ids: Vec<String>,
106 pub reference_root_node_ids: Vec<String>,
107 pub channel: Value,
108 pub free_time: Value,
109 pub orchestration: Value,
110 pub provenance_id: Option<String>,
111 pub mode: AgentMode,
112 pub source_session_type: Option<String>,
113 pub group_context: Value,
114 pub rust_lib_session_id: Option<String>,
115}
116
117impl SessionOptions {
118 pub fn conversation(session_type: impl Into<String>, roots: Vec<String>) -> Self {
119 Self {
120 session_type: session_type.into(),
121 root_node_ids: roots,
122 reference_root_node_ids: Vec::new(),
123 channel: Value::Null,
124 free_time: Value::Null,
125 orchestration: json!({"owner":"backend","status":"idle"}),
126 provenance_id: None,
127 mode: AgentMode::Conversation,
128 source_session_type: None,
129 group_context: Value::Null,
130 rust_lib_session_id: None,
131 }
132 }
133}
134
135fn restore_session_type(options: &mut SessionOptions, state: &Value) {
136 if !matches!(&options.mode, AgentMode::Ingress { .. }) {
137 options.session_type = state
138 .get("sessionType")
139 .and_then(Value::as_str)
140 .unwrap_or(&options.session_type)
141 .to_owned();
142 }
143}
144
145fn restore_commit_receipt(restored: Option<&Value>) -> anyhow::Result<Option<CommitReceipt>> {
146 restored
147 .and_then(|state| state.get("commitReceipt"))
148 .filter(|receipt| !receipt.is_null())
149 .cloned()
150 .map(serde_json::from_value)
151 .transpose()
152 .context("decoding the stored session commit receipt")
153}
154
155#[derive(Clone, Debug, Default, Deserialize, Serialize)]
156#[serde(rename_all = "camelCase")]
157struct KwebPlan {
158 creates: Vec<StagedNodeCreate>,
159 updates: BTreeMap<String, PlannedNode>,
160}
161
162#[derive(Clone, Debug, Deserialize, Serialize)]
163#[serde(rename_all = "camelCase")]
164struct StagedNodeCreate {
165 pending_id: String,
166 data: PlannedNode,
167}
168
169struct CreateNodeArguments {
170 parents: Vec<String>,
171 owner: String,
172 short_name: String,
173 short_description: String,
174 long_description: String,
175}
176
177impl KwebPlan {
178 fn restore(restored: Option<&Value>, journal: &HistorySession) -> anyhow::Result<Self> {
179 if let Some(plan) = restored.and_then(|state| state.get("kwebPlan")) {
180 return serde_json::from_value(plan.clone()).context("decoding the staged Kweb plan");
181 }
182 let latest = journal.state().events.iter().rev().find_map(|event| {
183 let EventKind::KwebPlanChanged { operation } = &event.kind else {
184 return None;
185 };
186 operation.get("plan")
187 });
188 latest
189 .cloned()
190 .map(serde_json::from_value)
191 .transpose()
192 .context("decoding the staged Kweb plan")
193 .map(Option::unwrap_or_default)
194 }
195
196 fn created(&self, id: &str) -> Option<&PlannedNode> {
197 self.creates
198 .iter()
199 .find(|create| create.pending_id == id)
200 .map(|create| &create.data)
201 }
202
203 fn created_mut(&mut self, id: &str) -> Option<&mut PlannedNode> {
204 self.creates
205 .iter_mut()
206 .find(|create| create.pending_id == id)
207 .map(|create| &mut create.data)
208 }
209}
210
211pub struct Session {
212 api: Service,
213 runtime: RuntimeModel,
214 journal: HistorySession,
215 plan: KwebPlan,
216 pub session_type: String,
217 pub channel: Value,
218 pub free_time: Value,
219 pub orchestration: Value,
220 pub provenance_id: Option<String>,
221 pub rust_lib_session_id: String,
222 pub root_node_ids: Vec<String>,
223 pub reference_root_node_ids: Vec<String>,
224 pub started_at: String,
225 pub transcript: Vec<Value>,
226 pub pending_turn: bool,
227 pub pending_external_event_id: Option<String>,
228 pub completed: bool,
229 pub rounds_used: u64,
230 commit_receipt: Option<CommitReceipt>,
231 commit_author: String,
232 mode: AgentMode,
233 source_session_type: Option<String>,
234 group_context: Value,
235 context: KwebContext,
236 free_time_end_reason: Option<String>,
237 fatal_persistence_error: Option<String>,
238}
239
240#[derive(Clone, Copy, Debug, Eq, PartialEq)]
241enum InputStage {
242 Accepted,
243}
244
245#[derive(Clone, Copy, Debug, Eq, PartialEq)]
246enum ContextRecovery {
247 NotNeeded,
248 Recovered,
249 Irreducible,
250}
251
252fn render_load_nodes_result(
253 journal: &HistorySession,
254 changed_box_ids: &[BoxId],
255) -> anyhow::Result<String> {
256 let changed_box_ids = changed_box_ids
257 .iter()
258 .map(ToString::to_string)
259 .collect::<Vec<_>>();
260 let projected_boxes = journal
261 .state()
262 .projection()
263 .items
264 .into_iter()
265 .filter(|item| !item.marker)
266 .map(|item| (item.box_id.to_string(), item.text))
267 .collect::<Vec<_>>();
268 render(RenderRequest::LoadNodes {
269 changed_box_ids: &changed_box_ids,
270 projected_boxes: &projected_boxes,
271 })
272}
273
274fn provider_tool_result_with_context_footer(journal: &HistorySession, result: &str) -> String {
275 let footer = journal.state().projection().footer;
276 render(RenderRequest::ProviderFooter {
277 result,
278 footer: &footer,
279 })
280 .expect("provider-footer rendering is infallible")
281}
282
283fn append_slow_tool_duration(text: &mut String, elapsed: Duration) {
284 *text = render(RenderRequest::SlowTool { text, elapsed })
285 .expect("slow-tool rendering is infallible");
286}
287
288fn render_web_search_result(
289 result: &kcode_intelligence_router::SearchResponse,
290) -> anyhow::Result<String> {
291 let sources = result
292 .sources
293 .iter()
294 .map(|source| (source.title.clone(), source.url.clone()))
295 .collect::<Vec<_>>();
296 render(RenderRequest::WebSearch {
297 answer: &result.answer,
298 sources: &sources,
299 })
300}
301
302fn render_web_fetch_result(
303 result: &kcode_intelligence_router::FetchResponse,
304) -> anyhow::Result<String> {
305 render(RenderRequest::WebFetch {
306 url: &result.url,
307 title: result.title.as_deref(),
308 content_type: &result.content_type,
309 truncated: result.truncated,
310 content: &result.content,
311 })
312}
313
314fn render_media_annotation_result(
315 object_id: &str,
316 file_name: &str,
317 content_type: &str,
318 result: &kcode_intelligence_router::AnnotationResponse,
319) -> anyhow::Result<String> {
320 render(RenderRequest::MediaAnnotation {
321 object_id,
322 file_name,
323 content_type,
324 model: &result.model,
325 complete: result.complete,
326 incomplete_reason: result.incomplete_reason.as_deref(),
327 text: &result.text,
328 })
329}
330
331fn render_audio_transcription_result(
332 object_id: &str,
333 file_name: &str,
334 content_type: &str,
335 result: &kcode_intelligence_router::TranscriptionResponse,
336) -> anyhow::Result<String> {
337 render(RenderRequest::AudioTranscription {
338 object_id,
339 file_name,
340 content_type,
341 model: &result.model,
342 text: &result.text,
343 })
344}
345
346fn render_document_extraction_result(
347 object_id: &str,
348 file_name: &str,
349 result: &kcode_intelligence_router::DocumentExtraction,
350) -> anyhow::Result<String> {
351 render(RenderRequest::DocumentExtraction {
352 object_id,
353 file_name,
354 format: &result.format,
355 characters: result.characters,
356 truncated: result.truncated,
357 text: &result.text,
358 })
359}
360
361struct ToolCall {
362 name: String,
363 arguments: Value,
364}
365
366#[derive(Deserialize)]
367#[serde(rename_all = "camelCase", deny_unknown_fields)]
368struct TaskId {
369 task_id: String,
370}
371
372#[derive(Deserialize)]
373#[serde(rename_all = "camelCase", deny_unknown_fields)]
374struct CategoryId {
375 category_id: String,
376}
377
378#[derive(Deserialize)]
379#[serde(rename_all = "camelCase", deny_unknown_fields)]
380struct CategoryCall {
381 category_id: String,
382 #[serde(default)]
383 offset: u64,
384 #[serde(default = "task_page_limit")]
385 limit: u32,
386}
387
388#[derive(Deserialize)]
389#[serde(deny_unknown_fields)]
390struct EmptyCall {}
391
392fn task_page_limit() -> u32 {
393 50
394}
395
396struct RecordedToolInvocation {
397 invocation_id: String,
398 tool_instance: String,
399 tool_name: String,
400}
401
402struct PendingFreeformWrite {
403 request: FreeformWrite,
404 call_box_id: BoxId,
405}
406
407fn tool_call_box_content(call: &ToolCall) -> anyhow::Result<BoxContent> {
408 if matches!(
409 call.name.as_str(),
410 WRITE_RUST_LIB_TOOL | WRITE_WEB_LIB_TOOL | WRITE_RUST_BIN_TOOL
411 ) {
412 let name = call
413 .arguments
414 .get("name")
415 .and_then(Value::as_str)
416 .map(|name| name.chars().take(255).collect::<String>());
417 let file_count = call
418 .arguments
419 .get("files")
420 .and_then(Value::as_array)
421 .map(Vec::len);
422 return Ok(BoxContent {
423 text: serde_json::to_string_pretty(&json!({
424 "name":call.name,
425 "arguments":{
426 "name":name,
427 "fileCount":file_count,
428 "completeFileContents":"omitted from active context; retained in the durable tool invocation"
429 }
430 }))?,
431 objects: Vec::new(),
432 metadata: json!({
433 "compactedToolInvocation":true,
434 "toolName":call.name,
435 }),
436 });
437 }
438 Ok(BoxContent::text(serde_json::to_string_pretty(
439 &json!({"name":call.name,"arguments":call.arguments}),
440 )?))
441}
442
443struct ToolOutcome {
444 text: String,
445 store_result: bool,
446 ok: bool,
447 end_session: bool,
448 freeform_write: Option<FreeformWrite>,
449 managed_source_snapshot: Option<SourceSnapshot>,
450}
451
452fn result_displays_snapshot(result: &str, snapshot: &SourceSnapshot) -> bool {
453 result == snapshot.text
454}
455
456fn subagent_managed_write_fits(
457 context: &SubagentContext,
458 call: &ToolCall,
459 budget: &kcode_agent_runtime::ContextBudget,
460) -> bool {
461 let Some(snapshot) = proposed_write_snapshot(&call.name, &call.arguments) else {
462 return true;
463 };
464 let state = context.source_state(&snapshot);
465 budget.fits_state(state.key, state.text)
466}
467
468struct KennedySubagentHost<'a> {
469 session: &'a mut Session,
470 context: SubagentContext,
471 captures: HashMap<String, FreeformWrite>,
472}
473
474struct KennedySessionHost<'a, C> {
475 session: &'a mut Session,
476 checkpoint: &'a mut C,
477 accounting: Option<kcode_intelligence_chatend::TopLevelCall>,
478 pending_freeform_write: Option<PendingFreeformWrite>,
479 deadline_after_response: bool,
480}
481
482fn is_kweb_mutation(name: &str) -> bool {
483 matches!(
484 name,
485 "ConnectNodes" | "ConsolidateFanout" | "SetFixedConnection" | "CreateNode" | "UpdateNode"
486 )
487}
488
489fn referenced_pending_nodes(decoded: &DecodedTool) -> Vec<String> {
490 let mut ids = Vec::new();
491 match decoded {
492 DecodedTool::ConnectNodes(nodes) => ids.extend(nodes.iter()),
493 DecodedTool::ConsolidateFanout {
494 parent,
495 fanout,
496 aggregator,
497 } => {
498 ids.push(parent);
499 ids.extend(fanout.iter());
500 ids.push(aggregator);
501 }
502 DecodedTool::SetFixedConnection {
503 parent,
504 child,
505 slot: _,
506 } => {
507 ids.push(parent);
508 ids.extend(child.iter());
509 }
510 DecodedTool::CreateNode { parents, owner, .. } => {
511 ids.extend(parents.iter());
512 ids.push(owner);
513 }
514 DecodedTool::UpdateNode { id, owner, .. } => {
515 ids.push(id);
516 ids.push(owner);
517 }
518 _ => {}
519 }
520 ids.into_iter()
521 .filter(|id| id.starts_with("pending:"))
522 .cloned()
523 .collect()
524}
525
526fn subagent_unavailable_reason(name: &str) -> Option<&'static str> {
527 match name {
528 "RunSubagent" => {
529 Some("RunSubagent is unavailable inside a subagent. Only Kennedy may launch subagents.")
530 }
531 "EndSession" => Some(
532 "EndSession is unavailable inside a subagent. A child cannot control the parent session lifecycle.",
533 ),
534 "DehydrateBoxes" | "SummarizeBox" | "HydrateBox" | "BoxesIntoObjects" => {
535 Some("Parent box controls are unavailable inside a box-free subagent context.")
536 }
537 _ => None,
538 }
539}
540
541fn execute_kweb_mutation(
542 name: &str,
543 decoded: DecodedTool,
544 context: &KwebContext,
545 plan: &mut KwebPlan,
546 journal: &mut HistorySession,
547) -> anyhow::Result<String> {
548 match (name, decoded) {
549 ("ConnectNodes", DecodedTool::ConnectNodes(ids)) => {
550 Session::connect_nodes(plan, context, ids)
551 }
552 (
553 "ConsolidateFanout",
554 DecodedTool::ConsolidateFanout {
555 parent,
556 fanout,
557 aggregator,
558 },
559 ) => Session::consolidate_fanout(plan, context, parent, fanout, aggregator),
560 (
561 "SetFixedConnection",
562 DecodedTool::SetFixedConnection {
563 parent,
564 child,
565 slot,
566 },
567 ) => Session::set_fixed_connection(plan, context, parent, child, slot),
568 (
569 "CreateNode",
570 DecodedTool::CreateNode {
571 parents,
572 owner,
573 short_name,
574 short_description,
575 long_description,
576 },
577 ) => Session::create_node(
578 plan,
579 context,
580 journal,
581 CreateNodeArguments {
582 parents,
583 owner,
584 short_name,
585 short_description,
586 long_description,
587 },
588 ),
589 (
590 "UpdateNode",
591 DecodedTool::UpdateNode {
592 id,
593 owner,
594 short_name,
595 short_description,
596 long_description,
597 },
598 ) => Session::update_node(
599 plan,
600 context,
601 id,
602 owner,
603 short_name,
604 short_description,
605 long_description,
606 ),
607 _ => anyhow::bail!("decoded contract for {name} did not match its Kweb mutation"),
608 }
609}
610
611fn kweb_plan_projection(plan: &KwebPlan) -> (BTreeMap<String, NodeDraft>, Vec<KwebStagedCreate>) {
612 let updates = plan
613 .updates
614 .iter()
615 .map(|(id, node)| (id.clone(), kweb_node_draft(node)))
616 .collect();
617 let creates = plan
618 .creates
619 .iter()
620 .map(|create| KwebStagedCreate {
621 pending_id: create.pending_id.clone(),
622 data: kweb_node_draft(&create.data),
623 })
624 .collect();
625 (updates, creates)
626}
627
628impl Session {
629 pub async fn new(
630 api: Service,
631 system_prompt: String,
632 runtime: RuntimeModel,
633 mut options: SessionOptions,
634 restored: Option<&Value>,
635 ) -> anyhow::Result<Self> {
636 if let Some(state) = restored {
637 restore_session_type(&mut options, state);
638 options.channel = state.get("channel").cloned().unwrap_or(options.channel);
639 options.free_time = state.get("freeTime").cloned().unwrap_or(options.free_time);
640 options.orchestration = state
641 .get("orchestration")
642 .cloned()
643 .unwrap_or(options.orchestration);
644 }
645 if options.group_context.is_null() {
646 options.group_context = options
647 .channel
648 .get("groupContext")
649 .cloned()
650 .unwrap_or(Value::Null);
651 }
652 options
653 .reference_root_node_ids
654 .retain(|id| !options.root_node_ids.contains(id));
655 options.reference_root_node_ids.sort();
656 options.reference_root_node_ids.dedup();
657
658 let started_at = restored
659 .and_then(|state| state.get("startedAt"))
660 .and_then(Value::as_str)
661 .map(str::to_owned)
662 .unwrap_or_else(|| Utc::now().to_rfc3339());
663 let rust_lib_session_id = restored
664 .and_then(|state| state.get("rustLibSessionId"))
665 .and_then(Value::as_str)
666 .map(str::to_owned)
667 .or(options.rust_lib_session_id.clone())
668 .unwrap_or_else(|| format!("kennedy:{}", Uuid::new_v4()));
669 let history_session_id = restored
670 .and_then(|state| state.get("sessionId"))
671 .and_then(Value::as_str)
672 .map(str::to_owned);
673 let source_session_type = options.source_session_type.clone().or_else(|| {
674 restored
675 .and_then(|state| state.get("sourceSessionType"))
676 .and_then(Value::as_str)
677 .map(str::to_owned)
678 });
679 let session_id = history_session_id
680 .clone()
681 .unwrap_or_else(|| Uuid::new_v4().to_string());
682 let metadata = SessionMetadata {
683 session_id: session_id.clone(),
684 kind: session_kind(&options.session_type, &options.mode),
685 created_at: started_at.clone(),
686 effective_context_tokens: runtime.context_window_tokens,
687 channel: options.channel.clone(),
688 };
689 let mut journal = if history_session_id.is_some() {
690 api.history_session(metadata, &runtime.model)
691 .with_context(|| {
692 format!(
693 "opening authoritative session {session_id} (legacy snapshots are intentionally unsupported)"
694 )
695 })?
696 } else {
697 api.create_history_session(NewSession {
698 kind: metadata.kind,
699 created_at: metadata.created_at,
700 effective_context_tokens: metadata.effective_context_tokens,
701 channel: metadata.channel,
702 })?
703 };
704 let mut context =
705 KwebContext::new(options.root_node_ids.clone()).map_err(anyhow::Error::new)?;
706 restore_kweb_context(&journal, &mut context)?;
707 let plan = KwebPlan::restore(restored, &journal)?;
708 let transcript = kcode_kennedy_session_ingress::transcript_from_journal(&journal);
709 let (pending_turn, pending_external_event_id) = restore_pending_turn(restored, &transcript);
710
711 let needs_initialization = !journal
712 .state()
713 .boxes
714 .values()
715 .any(|state| matches!(state.owner, BoxOwner::System));
716 let commit_receipt = restore_commit_receipt(restored)?;
717 let commit_author = restored
718 .and_then(|state| state.get("commitAuthor"))
719 .and_then(Value::as_str)
720 .map(str::to_owned)
721 .unwrap_or_else(|| runtime.attribution());
722 if let Some(receipt) = &commit_receipt {
723 journal.mark_completed(receipt.session_object_id.to_string());
724 }
725 let completed =
726 journal.state().completed_session_object.is_some() || commit_receipt.is_some();
727 let mut session = Self {
728 api,
729 runtime,
730 journal,
731 plan,
732 session_type: options.session_type,
733 channel: options.channel,
734 free_time: options.free_time,
735 orchestration: options.orchestration,
736 provenance_id: options.provenance_id,
737 rust_lib_session_id,
738 root_node_ids: options.root_node_ids,
739 reference_root_node_ids: options.reference_root_node_ids,
740 started_at,
741 transcript,
742 pending_turn,
743 pending_external_event_id,
744 completed,
745 rounds_used: restored
746 .and_then(|state| state.get("roundsUsed"))
747 .and_then(Value::as_u64)
748 .unwrap_or_default(),
749 commit_receipt,
750 commit_author,
751 mode: options.mode,
752 source_session_type,
753 group_context: options.group_context,
754 context,
755 free_time_end_reason: None,
756 fatal_persistence_error: None,
757 };
758
759 if matches!(session.mode, AgentMode::Ingress { .. }) && !session.journal.is_sealed() {
760 session.journal.repair_unfinished_tools(now())?;
761 }
762 if session.journal.is_sealed() {
763 anyhow::ensure!(
764 !matches!(session.mode, AgentMode::Conversation),
765 "a read-only conversation has an unexpectedly sealed session log"
766 );
767 if session.commit_receipt.is_none() {
768 session.finalize_kweb_session()?;
769 }
770 session.completed = true;
771 return Ok(session);
772 }
773
774 if needs_initialization {
775 session.journal.create_box(
776 now(),
777 "Kennedy system prompt",
778 BoxOwner::System,
779 BoxContent::text(&system_prompt),
780 )?;
781 if session.session_type == "telegram-group" && !session.group_context.is_null() {
782 session.journal.create_box(
783 now(),
784 "Telegram group context",
785 BoxOwner::Controller,
786 BoxContent::text(kcode_telegram_session_coordinator::format_group_context(
787 &session.group_context,
788 )),
789 )?;
790 }
791 let roots = session.root_node_ids.clone();
792 let invocation =
793 session.record_tool_invocation("LoadNodes", json!({"identifiers":&roots}))?;
794 let result = load_durable_batch(session.api.kmap(), &mut session.context, &roots)?;
795 session.sync_kweb_boxes()?;
796 session.record_tool_completion(
797 Some(&invocation),
798 json!({"ok":true,"automatic":true,"identifiers":roots,"result":result}),
799 )?;
800 } else {
801 session.sync_kweb_boxes()?;
802 }
803 if matches!(session.mode, AgentMode::Ingress { .. })
804 && !session.completed
805 && !session.journal.state().history_ingress_started
806 {
807 session.prepare_history_ingress(&system_prompt).await?;
808 }
809 Ok(session)
810 }
811
812 async fn prepare_history_ingress(&mut self, prompt: &str) -> anyhow::Result<()> {
813 let cost_at_ingress = self.journal.state().projection().status;
814 if !self.journal.state().source_terminated {
815 self.journal.record(
816 now(),
817 EventKind::SourceTerminated {
818 reason: "history_ingress".into(),
819 },
820 )?;
821 }
822 let system_box = self
823 .journal
824 .state()
825 .boxes
826 .values()
827 .find(|state| matches!(state.owner, BoxOwner::System))
828 .map(|state| state.id)
829 .context("session has no system-prompt box")?;
830 self.journal
831 .update_box(now(), system_box, BoxContent::text(prompt))?;
832 let ingress_kind = session_kind(&self.session_type, &self.mode);
833 if self.journal.state().metadata.effective_context_tokens
834 != self.runtime.context_window_tokens
835 || self.journal.state().metadata.kind != ingress_kind
836 {
837 self.journal
838 .configure_context(ingress_kind, self.runtime.context_window_tokens);
839 }
840 self.journal.create_box(
841 now(),
842 "Session cost at ingress",
843 BoxOwner::Controller,
844 BoxContent::text(cost_summary(
845 "session cost before history ingress",
846 cost_at_ingress.estimated_cost_usd_nanos,
847 cost_at_ingress.unpriced_provider_calls,
848 )),
849 )?;
850 self.revalidate_loaded_nodes().await?;
851 match kcode_history_ingress_context::prepare(&mut self.journal, now())? {
852 HistoryIngressContextOutcome::Ready => {}
853 HistoryIngressContextOutcome::OverCapacity {
854 estimated_tokens,
855 target_tokens,
856 } => {
857 self.journal.record(
858 now(),
859 EventKind::Note {
860 label: INGRESS_FORCE_COMMIT_NOTE.into(),
861 value: json!({
862 "reason":"fully_dehydrated_context_above_initial_target",
863 "estimatedTokens":estimated_tokens,
864 "initialTargetTokens":target_tokens,
865 }),
866 },
867 )?;
868 self.pending_turn = false;
869 self.finalize_kweb_session()?;
870 self.completed = true;
871 return Ok(());
872 }
873 }
874 self.journal
875 .record(now(), EventKind::HistoryIngressStarted)?;
876 self.pending_turn = true;
877 Ok(())
878 }
879
880 async fn revalidate_loaded_nodes(&mut self) -> anyhow::Result<()> {
881 let direct = self.context.loaded_node_ids().to_vec();
882 load_durable_batch(self.api.kmap(), &mut self.context, &direct)?;
883 self.sync_kweb_boxes()?;
884 Ok(())
885 }
886
887 fn stage_user_input(&mut self, text: &str, metadata: &Value) -> Option<InputStage> {
888 let recorded_at = now();
889 let result = (|| -> anyhow::Result<Option<InputStage>> {
890 let Some(staged) = kcode_kennedy_session_ingress::stage_user_input(
891 &mut self.journal,
892 text,
893 metadata,
894 &recorded_at,
895 )?
896 else {
897 return Ok(None);
898 };
899 self.transcript.push(staged.transcript);
900 self.recover_context_overflow(staged.external_event_id.as_deref(), &[])?;
901 Ok(Some(InputStage::Accepted))
902 })();
903 match result {
904 Ok(stage) => stage,
905 Err(error) => {
906 self.fatal_persistence_error = Some(error.to_string());
907 tracing::error!(error=%error, "Could not durably stage session input");
908 Some(InputStage::Accepted)
909 }
910 }
911 }
912
913 pub fn append_final_user_message(&mut self, text: &str, metadata: &Value) -> bool {
914 self.stage_user_input(text, metadata).is_some()
915 }
916
917 pub fn stage_source_message(
918 &mut self,
919 kennedy: bool,
920 text: &str,
921 metadata: Value,
922 ) -> anyhow::Result<()> {
923 let staged = kcode_kennedy_session_ingress::stage_source_input(
924 &mut self.journal,
925 kennedy,
926 text,
927 metadata,
928 &now(),
929 )?;
930 self.transcript.push(staged.transcript);
931 self.recover_context_overflow(staged.external_event_id.as_deref(), &[])?;
932 Ok(())
933 }
934
935 pub fn answer_for_external_event(&self, id: &str) -> Option<&Value> {
936 self.transcript.iter().rev().find(|entry| {
937 is_terminal_external_response(entry)
938 && entry.get("externalEventId").and_then(Value::as_str) == Some(id)
939 })
940 }
941
942 pub fn responses_for_external_event(&self, id: &str) -> Vec<&Value> {
943 self.transcript
944 .iter()
945 .filter(|entry| {
946 matches!(
947 entry.get("role").and_then(Value::as_str),
948 Some("kennedy" | "system")
949 ) && entry.get("externalEventId").and_then(Value::as_str) == Some(id)
950 })
951 .collect()
952 }
953
954 pub fn resolve_object(&mut self, object_id: &str) -> anyhow::Result<ResolvedObject> {
955 let api = self.api.clone();
956 kcode_kennedy_session_objects::resolve_object(
957 &mut self.journal,
958 object_id,
959 move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
960 )
961 }
962
963 fn resolve_media_object(&mut self, object_id: &str) -> anyhow::Result<ResolvedObject> {
964 let api = self.api.clone();
965 kcode_kennedy_session_objects::resolve_media_object(
966 &mut self.journal,
967 object_id,
968 MAX_MEDIA_ENRICHMENT_BYTES,
969 move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
970 )
971 }
972
973 fn resolve_image_object(
974 &mut self,
975 object_id: &str,
976 ) -> anyhow::Result<(Vec<u8>, String, String)> {
977 let resolved = self.resolve_media_object(object_id)?;
978 anyhow::ensure!(
979 resolved.media_type.starts_with("image/"),
980 "GenerateImage reference {object_id} is not an image"
981 );
982 Ok((resolved.bytes, resolved.file_name, resolved.media_type))
983 }
984
985 fn recover_context_overflow(
986 &mut self,
987 external_event_id: Option<&str>,
988 pinned_box_ids: &[BoxId],
989 ) -> anyhow::Result<ContextRecovery> {
990 let projection = self.journal.state().projection();
991 let target_tokens = self.journal.state().active_context_limit();
992 if projection.estimated_tokens <= target_tokens {
993 return Ok(ContextRecovery::NotNeeded);
994 }
995 let projection_hash = hex::encode(Sha256::digest(projection.render().as_bytes()));
996 let already_irreducible = self
997 .journal
998 .state()
999 .events
1000 .iter()
1001 .rev()
1002 .find_map(|event| match &event.kind {
1003 EventKind::Note { label, value } if label == "context_overflow_recovery" => {
1004 Some(value)
1005 }
1006 _ => None,
1007 })
1008 .is_some_and(|value| {
1009 value.get("irreducible").and_then(Value::as_bool) == Some(true)
1010 && value.get("limitTokens").and_then(Value::as_u64) == Some(target_tokens)
1011 && value.get("projectionHash").and_then(Value::as_str)
1012 == Some(projection_hash.as_str())
1013 });
1014 if already_irreducible {
1015 return Ok(ContextRecovery::Irreducible);
1016 }
1017
1018 let before_tokens = projection.estimated_tokens;
1019 let mut metadata = json!({
1020 "transcriptRole":"system",
1021 "contextOverflowWarning":true,
1022 "projectedTokens":before_tokens,
1023 "limitTokens":target_tokens,
1024 });
1025 if let Some(id) = external_event_id {
1026 metadata["externalEventId"] = json!(id);
1027 }
1028 let warning_box_id = self.journal.create_box(
1029 now(),
1030 CONTEXT_OVERFLOW_WARNING_BOX_NAME,
1031 BoxOwner::Controller,
1032 BoxContent {
1033 text: CONTEXT_OVERFLOW_WARNING.into(),
1034 objects: Vec::new(),
1035 metadata,
1036 },
1037 )?;
1038 let mut transcript = json!({
1039 "role":"system",
1040 "content":CONTEXT_OVERFLOW_WARNING,
1041 "contextOverflowWarning":true,
1042 });
1043 if let Some(id) = external_event_id {
1044 transcript["externalEventId"] = json!(id);
1045 }
1046 self.transcript.push(transcript);
1047
1048 let mut pins = pinned_box_ids.to_vec();
1049 if !pins.contains(&warning_box_id) {
1050 pins.push(warning_box_id);
1051 }
1052 let outcome = kcode_history_ingress_context::recover(&mut self.journal, now(), &pins)?;
1053 let (dehydrated_box_ids, estimated_tokens, target_tokens, irreducible) = match outcome {
1054 ContextRecoveryOutcome::Recovered {
1055 dehydrated_box_ids,
1056 estimated_tokens,
1057 target_tokens,
1058 } => (dehydrated_box_ids, estimated_tokens, target_tokens, false),
1059 ContextRecoveryOutcome::OverCapacity {
1060 dehydrated_box_ids,
1061 estimated_tokens,
1062 target_tokens,
1063 } => (dehydrated_box_ids, estimated_tokens, target_tokens, true),
1064 };
1065 let final_projection_hash = hex::encode(Sha256::digest(
1066 self.journal.state().projection().render().as_bytes(),
1067 ));
1068 self.journal.record(
1069 now(),
1070 EventKind::Note {
1071 label: "context_overflow_recovery".into(),
1072 value: json!({
1073 "beforeTokens":before_tokens,
1074 "estimatedTokens":estimated_tokens,
1075 "limitTokens":target_tokens,
1076 "dehydratedBoxIds":dehydrated_box_ids,
1077 "irreducible":irreducible,
1078 "projectionHash":final_projection_hash,
1079 }),
1080 },
1081 )?;
1082 if irreducible {
1083 if matches!(self.mode, AgentMode::Ingress { .. }) {
1084 self.request_ingress_force_commit(
1085 "irreducible_context_overflow",
1086 estimated_tokens,
1087 )?;
1088 } else if !self.journal.state().source_terminated {
1089 self.journal.record(
1090 now(),
1091 EventKind::SourceTerminated {
1092 reason: "irreducible_context_overflow".into(),
1093 },
1094 )?;
1095 }
1096 Ok(ContextRecovery::Irreducible)
1097 } else {
1098 Ok(ContextRecovery::Recovered)
1099 }
1100 }
1101
1102 fn request_ingress_force_commit(
1103 &mut self,
1104 reason: &str,
1105 projected_tokens: u64,
1106 ) -> anyhow::Result<()> {
1107 if self.ingress_force_commit_requested() {
1108 return Ok(());
1109 }
1110 self.journal.record(
1111 now(),
1112 EventKind::Note {
1113 label: INGRESS_FORCE_COMMIT_NOTE.into(),
1114 value: json!({
1115 "reason":reason,
1116 "projectedTokens":projected_tokens,
1117 "limitTokens":self.journal.state().ingress_context_limit(),
1118 }),
1119 },
1120 )?;
1121 Ok(())
1122 }
1123
1124 fn ingress_force_commit_requested(&self) -> bool {
1125 self.journal.state().events.iter().rev().any(|event| {
1126 matches!(
1127 &event.kind,
1128 EventKind::Note { label, .. } if label == INGRESS_FORCE_COMMIT_NOTE
1129 )
1130 })
1131 }
1132
1133 pub fn requires_history_ingress(&self) -> bool {
1134 matches!(self.mode, AgentMode::Conversation) && self.journal.state().source_terminated
1135 }
1136
1137 pub fn stage_free_time_opening(&mut self) -> bool {
1138 if self.pending_turn {
1139 return false;
1140 }
1141 let mut blocks = vec![
1142 render(RenderRequest::FreeTimeOpening {
1143 free_time: &self.free_time,
1144 })
1145 .expect("free-time opening rendering is infallible"),
1146 ];
1147 if let Some(message) = self
1148 .free_time
1149 .get("handoffMessage")
1150 .and_then(Value::as_str)
1151 .filter(|message| !message.trim().is_empty())
1152 {
1153 blocks.push(format!(
1154 "Message from the previous self-time session:\n\n{message}"
1155 ));
1156 }
1157 let Some(stage) = self.stage_user_input(&blocks.join("\n\n"), &json!({"kind":"self-time"}))
1158 else {
1159 return false;
1160 };
1161 self.pending_turn = matches!(stage, InputStage::Accepted);
1162 true
1163 }
1164
1165 pub fn stage_wakeup_opening(&mut self) -> anyhow::Result<bool> {
1166 if self.pending_turn {
1167 return Ok(false);
1168 }
1169 let marker = self
1170 .channel
1171 .get("wakeupMarker")
1172 .and_then(Value::as_str)
1173 .context("wakeup session is missing its acquired time marker")?;
1174 let marker = DateTime::parse_from_rfc3339(marker)
1175 .context("wakeup session has an invalid acquired time marker")?
1176 .with_timezone(&Utc);
1177 let text = render(RenderRequest::WakeupOpening { marker })?;
1178 let Some(stage) = self.stage_user_input(
1179 &text,
1180 &json!({"kind":"wakeup","wakeupMarker":marker.to_rfc3339()}),
1181 ) else {
1182 return Ok(false);
1183 };
1184 self.pending_turn = matches!(stage, InputStage::Accepted);
1185 Ok(true)
1186 }
1187
1188 pub fn begin_user_turn(&mut self, text: &str, metadata: &Value) -> bool {
1189 if self.pending_turn {
1190 return false;
1191 }
1192 let Some(stage) = self.stage_user_input(text, metadata) else {
1193 return false;
1194 };
1195 debug_assert_eq!(stage, InputStage::Accepted);
1196 self.rounds_used = 0;
1197 self.pending_turn = true;
1198 self.pending_external_event_id = metadata
1199 .get("externalEventId")
1200 .and_then(Value::as_str)
1201 .map(str::to_owned);
1202 true
1203 }
1204
1205 pub fn reset_exhausted_turn_rounds_for_retry(&mut self) {
1206 if matches!(self.mode, AgentMode::Conversation)
1207 && self.rounds_used >= AGENT_LOOP_ROUND_LIMIT
1208 {
1209 self.rounds_used = 0;
1210 }
1211 }
1212
1213 pub fn interrupt_current_turn(&mut self) -> anyhow::Result<()> {
1214 self.journal.repair_unfinished_tools(now())?;
1215 let notice = "The user stopped this agent turn.";
1216 let mut metadata = json!({"transcriptRole":"system","userStopped":true});
1217 let mut transcript_entry = json!({
1218 "role":"system",
1219 "content":notice,
1220 "userStopped":true,
1221 });
1222 if let Some(external_event_id) = &self.pending_external_event_id {
1223 metadata["externalEventId"] = json!(external_event_id);
1224 transcript_entry["externalEventId"] = json!(external_event_id);
1225 }
1226 self.journal.create_box(
1227 now(),
1228 "Turn stopped",
1229 BoxOwner::Controller,
1230 BoxContent {
1231 text: notice.into(),
1232 objects: Vec::new(),
1233 metadata,
1234 },
1235 )?;
1236 self.transcript.push(transcript_entry);
1237 self.pending_turn = false;
1238 self.pending_external_event_id = None;
1239 self.orchestration =
1240 json!({"owner":"backend","status":"idle","lastOutcome":"user-stopped"});
1241 Ok(())
1242 }
1243
1244 pub async fn run_pending_turn<C, F>(
1245 &mut self,
1246 operation_id: Uuid,
1247 mut checkpoint: C,
1248 ) -> anyhow::Result<Option<String>>
1249 where
1250 C: FnMut(Value) -> F + Send,
1251 F: Future<Output = anyhow::Result<()>> + Send,
1252 {
1253 if let Some(error) = self.fatal_persistence_error.take() {
1254 anyhow::bail!("session journal write failed: {error}");
1255 }
1256 if !self.pending_turn {
1257 return Ok(None);
1258 }
1259 let runtime = self.api.agent_runtime();
1260 let user_id = self
1261 .root_node_ids
1262 .first()
1263 .context("session has no user root for intelligence accounting")?
1264 .clone();
1265 let request = kcode_agent_runtime::SessionRunRequest {
1266 user_id,
1267 operation_id,
1268 rounds_used: self.rounds_used,
1269 round_limit: AGENT_LOOP_ROUND_LIMIT,
1270 };
1271 let mut host = KennedySessionHost {
1272 session: self,
1273 checkpoint: &mut checkpoint,
1274 accounting: None,
1275 pending_freeform_write: None,
1276 deadline_after_response: false,
1277 };
1278 let result = runtime.run_session(request, &mut host).await;
1279 let round_limit = result
1280 .as_ref()
1281 .is_err_and(kcode_agent_runtime::is_session_round_limit);
1282 let result = if round_limit && matches!(host.session.mode, AgentMode::Ingress { .. }) {
1283 host.session.request_ingress_force_commit(
1284 "agent_loop_round_limit",
1285 host.session.journal.state().projection().estimated_tokens,
1286 )?;
1287 (host.checkpoint)(host.session.snapshot()?).await?;
1288 None
1289 } else {
1290 result?
1291 };
1292 drop(host);
1293 match self.mode {
1294 AgentMode::Conversation => {
1295 if self.journal.state().source_terminated {
1296 self.pending_turn = false;
1297 self.pending_external_event_id = None;
1298 checkpoint(self.snapshot()?).await?;
1299 return Ok(None);
1300 }
1301 let Some(answer) = result else {
1302 if self
1303 .pending_external_event_id
1304 .as_deref()
1305 .and_then(|id| self.answer_for_external_event(id))
1306 .is_some()
1307 {
1308 self.pending_turn = false;
1309 self.pending_external_event_id = None;
1310 checkpoint(self.snapshot()?).await?;
1311 return Ok(None);
1312 }
1313 anyhow::bail!(
1314 "Kennedy ended a conversational turn without an assistant response"
1315 );
1316 };
1317 self.pending_turn = false;
1318 self.pending_external_event_id = None;
1319 checkpoint(self.snapshot()?).await?;
1320 Ok(Some(answer))
1321 }
1322 AgentMode::FreeTime | AgentMode::Wakeup | AgentMode::Ingress { .. } => {
1323 self.pending_turn = false;
1324 self.pending_external_event_id = None;
1325 self.finalize_kweb_session()?;
1326 self.completed = true;
1327 checkpoint(self.snapshot()?).await?;
1328 Ok(None)
1329 }
1330 }
1331 }
1332
1333 fn project_descendant<T>(
1334 &mut self,
1335 outcome: Result<kcode_intelligence_router::Accounted<T>, services::ApiError>,
1336 ) -> anyhow::Result<T> {
1337 match outcome {
1338 Ok(accounted) => {
1339 kcode_intelligence_chatend::record_descendant_receipt(
1340 &mut self.journal,
1341 &accounted.receipt,
1342 )?;
1343 Ok(accounted.value)
1344 }
1345 Err(error) => {
1346 if let Some(receipt) = &error.receipt {
1347 kcode_intelligence_chatend::record_descendant_receipt(
1348 &mut self.journal,
1349 receipt,
1350 )?;
1351 }
1352 Err(error.into())
1353 }
1354 }
1355 }
1356
1357 async fn run_subagent(
1358 &mut self,
1359 model: String,
1360 reasoning_effort: Option<String>,
1361 context_node_ids: Vec<String>,
1362 task: String,
1363 parent_operation_id: Uuid,
1364 ) -> anyhow::Result<String> {
1365 let reasoning_effort =
1366 reasoning_effort.unwrap_or_else(|| self.runtime.reasoning_effort.clone());
1367 let mut context = Vec::with_capacity(context_node_ids.len());
1368 for node_id in &context_node_ids {
1369 context.push(self.api.kmap_node(node_id)?.data.long_description);
1370 }
1371 let user_id = self
1372 .root_node_ids
1373 .first()
1374 .context("session has no user root for subagent intelligence accounting")?
1375 .clone();
1376 let timeout = self.agent_request_timeout();
1377 let runtime = self.api.agent_runtime();
1378 let first_event = self.journal.state().events.len();
1379 let cost_before = self.journal.state().projection().status;
1380 let subagent_context = SubagentContext::new(self.root_node_ids.clone())?;
1381 let result = {
1382 let mut host = KennedySubagentHost {
1383 session: self,
1384 context: subagent_context,
1385 captures: HashMap::new(),
1386 };
1387 runtime
1388 .run(
1389 kcode_agent_runtime::RunRequest {
1390 user_id,
1391 parent_operation_id,
1392 model,
1393 reasoning_effort,
1394 context,
1395 task,
1396 timeout,
1397 start_metadata: json!({"contextNodeIds":context_node_ids}),
1398 },
1399 &mut host,
1400 )
1401 .await
1402 };
1403 match result {
1404 Ok(result) => {
1405 let cost_after = self.journal.state().projection().status;
1406 Ok(format!(
1407 "{}\n\n[{}]",
1408 result.answer,
1409 cost_summary(
1410 "subagent cost",
1411 cost_after
1412 .estimated_cost_usd_nanos
1413 .saturating_sub(cost_before.estimated_cost_usd_nanos),
1414 cost_after
1415 .unpriced_provider_calls
1416 .saturating_sub(cost_before.unpriced_provider_calls),
1417 )
1418 ))
1419 }
1420 Err(error) => {
1421 let may_have_effects =
1422 self.journal.state().events[first_event..]
1423 .iter()
1424 .any(|event| {
1425 matches!(
1426 &event.kind,
1427 EventKind::Note { label, .. } if label == "subagent_tool_call"
1428 )
1429 });
1430 if may_have_effects {
1431 Err(error.context(
1432 "the subagent failed after making Ktool calls; some tool effects may already have occurred",
1433 ))
1434 } else {
1435 Err(error)
1436 }
1437 }
1438 }
1439 }
1440
1441 async fn complete_subagent_freeform_write(
1442 &mut self,
1443 context: &mut SubagentContext,
1444 request: FreeformWrite,
1445 contents: String,
1446 budget: &kcode_agent_runtime::ContextBudget,
1447 ) -> anyhow::Result<kcode_agent_runtime::ToolOutcome> {
1448 let kind = request.kind();
1449 let freeform_tool = request.write_tool();
1450 anyhow::ensure!(
1451 context.source_is_open(kind, request.name()),
1452 "{} {:?} is not open in this subagent context. Call {} first.",
1453 kind.label(),
1454 request.name(),
1455 kind.open_tool()
1456 );
1457 let backend_arguments = request.capture_subagent(&mut self.journal, &now(), contents)?;
1458 let preview = self
1459 .api
1460 .managed_source_execute(
1461 &self.rust_lib_session_id,
1462 request.preview_tool(),
1463 backend_arguments.clone(),
1464 Vec::new(),
1465 )
1466 .await?;
1467 let preview = preview
1468 .snapshot
1469 .context("subagent freeform write preview omitted its source snapshot")?;
1470 let preview_state = context.source_state(&preview);
1471 anyhow::ensure!(
1472 budget.fits_state(preview_state.key, preview_state.text),
1473 "{freeform_tool} was not run because its resulting source state would exceed the subagent context limit"
1474 );
1475 let execution = self
1476 .api
1477 .managed_source_execute(
1478 &self.rust_lib_session_id,
1479 freeform_tool,
1480 backend_arguments,
1481 Vec::new(),
1482 )
1483 .await?;
1484 let snapshot = execution
1485 .snapshot
1486 .context("subagent freeform write omitted its resulting source snapshot")?;
1487 let state = context.apply_source_snapshot(snapshot);
1488 Ok(kcode_agent_runtime::ToolOutcome {
1489 text: execution.text,
1490 ok: true,
1491 state_updates: state.update.into_iter().collect(),
1492 displayed_state_keys: Vec::new(),
1493 capture: None,
1494 })
1495 }
1496
1497 async fn complete_freeform_write(
1498 &mut self,
1499 pending: PendingFreeformWrite,
1500 contents: String,
1501 ) -> anyhow::Result<ToolOutcome> {
1502 let request = pending.request;
1503 let freeform_tool = request.write_tool();
1504 let backend_arguments =
1505 request.capture(&mut self.journal, &now(), pending.call_box_id, contents)?;
1506 let preview_result = self
1507 .api
1508 .managed_source_execute(
1509 &self.rust_lib_session_id,
1510 request.preview_tool(),
1511 backend_arguments.clone(),
1512 Vec::new(),
1513 )
1514 .await;
1515 let preview = match preview_result {
1516 Ok(preview) => preview,
1517 Err(error) => {
1518 return Ok(ToolOutcome {
1519 text: format!("{freeform_tool} failed: {error}"),
1520 store_result: true,
1521 ok: false,
1522 end_session: false,
1523 freeform_write: None,
1524 managed_source_snapshot: None,
1525 });
1526 }
1527 };
1528 let _preview = preview
1529 .snapshot
1530 .context("freeform write preview omitted the resulting source snapshot")?;
1531 request.source_box_id(&self.journal)?;
1532
1533 let execution_result = self
1534 .api
1535 .managed_source_execute(
1536 &self.rust_lib_session_id,
1537 freeform_tool,
1538 backend_arguments,
1539 Vec::new(),
1540 )
1541 .await;
1542 let execution = match execution_result {
1543 Ok(execution) => execution,
1544 Err(error) => {
1545 return Ok(ToolOutcome {
1546 text: format!("{freeform_tool} failed: {error}"),
1547 store_result: true,
1548 ok: false,
1549 end_session: false,
1550 freeform_write: None,
1551 managed_source_snapshot: None,
1552 });
1553 }
1554 };
1555 let snapshot = execution
1556 .snapshot
1557 .context("freeform write omitted the resulting source snapshot")?;
1558 apply_snapshot(&mut self.journal, &now(), snapshot)?;
1559 Ok(ToolOutcome {
1560 text: execution.text,
1561 store_result: false,
1562 ok: true,
1563 end_session: false,
1564 freeform_write: None,
1565 managed_source_snapshot: None,
1566 })
1567 }
1568
1569 async fn send_telegram_dm(&mut self, arguments: &Value) -> anyhow::Result<String> {
1570 let request = kcode_telegram_session_coordinator::parse_private_request(arguments)?;
1571 let attachments = self.telegram_delivery_attachments(request.attachments)?;
1572 let caller_holds_user_lock = self.session_type == "telegram"
1573 && self.channel.get("telegramUserId").and_then(Value::as_i64)
1574 == Some(request.telegram_user_id);
1575 self.api
1576 .telegram()
1577 .send_private(kcode_telegram_session_coordinator::PrivateDelivery {
1578 telegram_user_id: request.telegram_user_id,
1579 message: request.message,
1580 attachments,
1581 caller_holds_user_lock,
1582 })
1583 .await
1584 }
1585
1586 async fn send_telegram_group_message(&mut self, arguments: &Value) -> anyhow::Result<String> {
1587 let request = kcode_telegram_session_coordinator::parse_group_request(arguments)?;
1588 let attachments = self.telegram_delivery_attachments(request.attachments)?;
1589 self.api
1590 .telegram()
1591 .send_group(kcode_telegram_session_coordinator::GroupDelivery {
1592 root_node_id: request.root_node_id,
1593 message: request.message,
1594 attachments,
1595 })
1596 .await
1597 }
1598
1599 fn telegram_delivery_attachments(
1600 &mut self,
1601 requests: Vec<kcode_telegram_session_coordinator::AttachmentRequest>,
1602 ) -> anyhow::Result<Vec<kcode_telegram_session_coordinator::Attachment>> {
1603 let api = self.api.clone();
1604 kcode_kennedy_session_objects::delivery_attachments(
1605 &mut self.journal,
1606 requests,
1607 move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
1608 )
1609 }
1610
1611 async fn execute_tool(
1612 &mut self,
1613 call: &ToolCall,
1614 operation_id: Uuid,
1615 ) -> anyhow::Result<ToolOutcome> {
1616 self.assert_tool_allowed(&call.name)?;
1617 let decoded = decode(&call.name, &call.arguments)?;
1618 let mut end_session = false;
1619 let mut store_result = true;
1620 let mut freeform_write = None;
1621 let mut managed_source_snapshot = None;
1622 let text = match (call.name.as_str(), decoded) {
1623 ("SendTelegramDM", _) => self.send_telegram_dm(&call.arguments).await?,
1624 ("SendTelegramGroupMessage", _) => {
1625 self.send_telegram_group_message(&call.arguments).await?
1626 }
1627 (
1628 "RunSubagent",
1629 Some(DecodedTool::RunSubagent {
1630 model,
1631 reasoning_effort,
1632 context_node_ids,
1633 task,
1634 }),
1635 ) => {
1636 let first_event = self.journal.state().events.len();
1637 match self
1638 .run_subagent(
1639 model,
1640 reasoning_effort,
1641 context_node_ids,
1642 task,
1643 operation_id,
1644 )
1645 .await
1646 {
1647 Ok(response) => response,
1648 Err(error) => {
1649 let may_have_effects = self.journal.state().events[first_event..]
1650 .iter()
1651 .any(|event| {
1652 matches!(
1653 &event.kind,
1654 EventKind::Note { label, .. }
1655 if label == "subagent_tool_call"
1656 )
1657 });
1658 if may_have_effects {
1659 return Err(error.context(
1660 "the subagent failed after making Ktool calls; some tool effects may already have occurred",
1661 ));
1662 }
1663 return Err(error);
1664 }
1665 }
1666 }
1667 ("EndSession", Some(DecodedTool::EndSession { message })) => {
1668 anyhow::ensure!(
1669 !matches!(self.mode, AgentMode::Conversation),
1670 "EndSession is only available during an autonomous or history-ingress session"
1671 );
1672 end_session = true;
1673 if matches!(self.mode, AgentMode::FreeTime)
1674 && let Some(message) = message.filter(|message| !message.trim().is_empty())
1675 {
1676 self.free_time["nextSessionMessage"] = json!(message);
1677 }
1678 "Session ending.".into()
1679 }
1680 ("DehydrateBoxes", Some(DecodedTool::BoxIds(ids))) => {
1681 self.journal.dehydrate_boxes(now(), &ids)?;
1682 format!(
1683 "Dehydrated boxes {}.",
1684 ids.iter()
1685 .map(ToString::to_string)
1686 .collect::<Vec<_>>()
1687 .join(", ")
1688 )
1689 }
1690 ("SummarizeBox", Some(DecodedTool::SummarizeBox { box_id, summary })) => {
1691 self.journal.summarize_box(now(), box_id, summary)?;
1692 format!("Summarized box {box_id}.")
1693 }
1694 ("HydrateBox", Some(DecodedTool::BoxId(id))) => {
1695 self.journal.rehydrate_box(now(), id)?;
1696 let external_event_id = self.pending_external_event_id.clone();
1697 match self.recover_context_overflow(external_event_id.as_deref(), &[id])? {
1698 ContextRecovery::NotNeeded => format!("Hydrated box {id}."),
1699 ContextRecovery::Recovered => {
1700 format!("Hydrated box {id}.\n\n{CONTEXT_OVERFLOW_WARNING}")
1701 }
1702 ContextRecovery::Irreducible => anyhow::bail!(CONTEXT_OVERFLOW_WARNING),
1703 }
1704 }
1705 ("BoxesIntoObjects", Some(DecodedTool::BoxIds(ids))) => {
1706 kcode_kennedy_box_text_objects::stage_box_text_objects(
1707 &mut self.journal,
1708 &ids,
1709 &now(),
1710 )?
1711 }
1712 ("LoadNodes", Some(DecodedTool::LoadNodes(identifiers))) => {
1713 load_durable_batch(self.api.kmap(), &mut self.context, &identifiers)?;
1714 let changed = self.sync_kweb_boxes()?;
1715 store_result = false;
1716 render_load_nodes_result(&self.journal, &changed)?
1717 }
1718 (
1719 "EmitObject",
1720 Some(DecodedTool::EmitObject {
1721 object_id,
1722 file_name,
1723 }),
1724 ) => {
1725 anyhow::ensure!(
1726 matches!(self.mode, AgentMode::Conversation),
1727 "EmitObject is only available in a conversation"
1728 );
1729 let object = self.resolve_object(&object_id)?;
1730 let file_name = file_name.unwrap_or_else(|| object.file_name.clone());
1731 if let Some(maximum) = self.channel.get("maxObjectBytes").and_then(Value::as_u64) {
1732 anyhow::ensure!(
1733 !object.bytes.is_empty(),
1734 "object {object_id} is empty and cannot be sent through this channel"
1735 );
1736 anyhow::ensure!(
1737 object.bytes.len() as u64 <= maximum,
1738 "object {object_id} is {} bytes, over this channel's {maximum}-byte limit",
1739 object.bytes.len()
1740 );
1741 }
1742 let descriptor = json!({
1743 "objectId":object_id,
1744 "fileName":file_name,
1745 "mediaType":object.media_type,
1746 "byteLength":object.bytes.len(),
1747 });
1748 let mut metadata = json!({
1749 "outputKind":"object",
1750 "attachments":[descriptor.clone()],
1751 });
1752 if let Some(external_event_id) = &self.pending_external_event_id {
1753 metadata["externalEventId"] = json!(external_event_id);
1754 }
1755 let content = BoxContent {
1756 text: String::new(),
1757 objects: vec![object_id.clone()],
1758 metadata,
1759 };
1760 self.journal
1761 .create_box(now(), "Kennedy message", BoxOwner::Kennedy, content)?;
1762 let mut transcript = json!({
1763 "role":"kennedy",
1764 "content":"",
1765 "objects":[object_id],
1766 "attachments":[descriptor],
1767 });
1768 if let Some(external_event_id) = &self.pending_external_event_id {
1769 transcript["externalEventId"] = json!(external_event_id);
1770 }
1771 self.transcript.push(transcript);
1772 store_result = false;
1773 "Object emitted to the user.".into()
1774 }
1775 ("WebSearch", Some(DecodedTool::WebSearch { question, model })) => {
1776 let user_id = self
1777 .root_node_ids
1778 .first()
1779 .context("session has no user root for intelligence accounting")?
1780 .clone();
1781 let outcome = self
1782 .api
1783 .search(
1784 &user_id,
1785 kcode_intelligence_router::SearchRequest {
1786 question,
1787 model,
1788 operation_id: Uuid::new_v4(),
1789 parent_operation_id: Some(operation_id),
1790 },
1791 )
1792 .await;
1793 let result = self.project_descendant(outcome)?;
1794 render_web_search_result(&result)?
1795 }
1796 ("WebFetch", Some(DecodedTool::WebFetch(url))) => {
1797 let user_id = self
1798 .root_node_ids
1799 .first()
1800 .context("session has no user root for intelligence accounting")?;
1801 let result = self
1802 .api
1803 .fetch(
1804 user_id,
1805 kcode_intelligence_router::FetchRequest {
1806 url,
1807 operation_id: Uuid::new_v4(),
1808 parent_operation_id: Some(operation_id),
1809 },
1810 )
1811 .await?;
1812 render_web_fetch_result(&result)?
1813 }
1814 ("StageTelegramGroupMedia", Some(DecodedTool::StageTelegramGroupMedia(message_id))) => {
1815 let media_ref = kcode_telegram_session_coordinator::group_media_reference(
1816 &self.group_context,
1817 message_id,
1818 )?;
1819 let chat_id = media_ref.chat_id;
1820 let api = self.api.clone();
1821 let staged = kcode_kennedy_session_objects::stage_telegram_group_media(
1822 &mut self.journal,
1823 kcode_kennedy_session_objects::TelegramStageRequest {
1824 chat_id,
1825 message_id,
1826 maximum_bytes: MAX_MEDIA_ENRICHMENT_BYTES,
1827 transport_metadata: media_ref.transport_metadata(),
1828 recorded_at: now(),
1829 },
1830 || api.telegram().group_message_media(chat_id, message_id),
1831 |media_type| {
1832 kcode_telegram_session_coordinator::group_media_file_name(
1833 &media_ref, media_type,
1834 )
1835 },
1836 )?;
1837 render(RenderRequest::StagedTelegramMedia {
1838 pending_id: &staged.descriptor.pending_id,
1839 kind: &staged.kind,
1840 file_name: &staged.descriptor.file_name,
1841 media_type: &staged.descriptor.media_type,
1842 size_bytes: staged.descriptor.size_bytes,
1843 message_id,
1844 reused: staged.reused,
1845 })?
1846 }
1847 (
1848 "TranscribeAudio",
1849 Some(DecodedTool::MediaEnrichment {
1850 object_id,
1851 model,
1852 prompt,
1853 }),
1854 ) => {
1855 let object = self.resolve_media_object(&object_id)?;
1856 validate(ValidationRequest::TranscribableAudio(&object.media_type))?;
1857 validate(ValidationRequest::TranscriptionModel(&model))?;
1858 let user_id = self
1859 .root_node_ids
1860 .first()
1861 .context("session has no user root for intelligence accounting")?
1862 .clone();
1863 let outcome = self
1864 .api
1865 .transcribe_audio(
1866 &user_id,
1867 &model,
1868 &prompt,
1869 object.bytes,
1870 object.file_name.clone(),
1871 &object.media_type,
1872 None,
1873 operation_id,
1874 )
1875 .await;
1876 let result = self.project_descendant(outcome)?;
1877 render_audio_transcription_result(
1878 &object.object_id,
1879 &object.file_name,
1880 &object.media_type,
1881 &result,
1882 )?
1883 }
1884 (
1885 "AnnotateMedia",
1886 Some(DecodedTool::MediaEnrichment {
1887 object_id,
1888 model,
1889 prompt,
1890 }),
1891 ) => {
1892 let media = self.resolve_media_object(&object_id)?;
1893 validate(ValidationRequest::Annotation {
1894 model: &model,
1895 media_type: &media.media_type,
1896 })?;
1897 let user_id = self
1898 .root_node_ids
1899 .first()
1900 .context("session has no user root for intelligence accounting")?
1901 .clone();
1902 let outcome = self
1903 .api
1904 .annotate_media(
1905 &user_id,
1906 &model,
1907 &prompt,
1908 media.bytes,
1909 media.file_name.clone(),
1910 &media.media_type,
1911 operation_id,
1912 )
1913 .await;
1914 let result = self.project_descendant(outcome)?;
1915 render_media_annotation_result(
1916 &media.object_id,
1917 &media.file_name,
1918 &media.media_type,
1919 &result,
1920 )?
1921 }
1922 (
1923 "GenerateImage",
1924 Some(DecodedTool::GenerateImage {
1925 model,
1926 prompt,
1927 reference_object_ids,
1928 }),
1929 ) => {
1930 let mut references = Vec::with_capacity(reference_object_ids.len());
1931 for object_id in &reference_object_ids {
1932 references.push(self.resolve_image_object(object_id)?);
1933 }
1934 let user_id = self
1935 .root_node_ids
1936 .first()
1937 .context("session has no user root for intelligence accounting")?
1938 .clone();
1939 let outcome = self
1940 .api
1941 .generate_image(&user_id, &model, &prompt, references, operation_id)
1942 .await;
1943 let result = self.project_descendant(outcome)?;
1944 let size = result.bytes.len();
1945 let file_name =
1946 format!("generated-image.{}", image_extension(&result.content_type));
1947 let object_id = self.api.save_generated_image(
1948 result.bytes,
1949 &file_name,
1950 &result.content_type,
1951 &result.model,
1952 )?;
1953 format!(
1954 "Generated image.\nObject: {object_id}\nFile: {file_name}\nContent type: {}\nSize: {size} bytes\nModel: {}\nUse EmitObject with {object_id} to deliver it.",
1955 result.content_type, result.model
1956 )
1957 }
1958 ("ExtractDocumentText", Some(DecodedTool::ObjectId(object_id))) => {
1959 let object = self.resolve_media_object(&object_id)?;
1960 validate(ValidationRequest::ExtractableDocument {
1961 media_type: &object.media_type,
1962 file_name: &object.file_name,
1963 })?;
1964 let result = self
1965 .api
1966 .extract_document(object.bytes, object.file_name.clone(), &object.media_type)
1967 .await?;
1968 render_document_extraction_result(&object.object_id, &object.file_name, &result)?
1969 }
1970 (name, None) if SPEECH_CLASSIFICATION_TOOLS.contains(&name) => {
1971 self.api
1972 .execute_speech_classification_tool(name, call.arguments.clone())
1973 .await?
1974 }
1975 (name, None) if TASK_BOARD_TOOLS.contains(&name) => {
1976 self.execute_task_board_tool(name, &call.arguments).await?
1977 }
1978 (name, Some(decoded)) if is_kweb_mutation(name) => {
1979 let text = execute_kweb_mutation(
1980 name,
1981 decoded,
1982 &self.context,
1983 &mut self.plan,
1984 &mut self.journal,
1985 )?;
1986 self.sync_kweb_boxes()?;
1987 text
1988 }
1989 (name, None)
1990 if RUST_LIB_TOOLS.contains(&name)
1991 || WEB_LIB_TOOLS.contains(&name)
1992 || RUST_BIN_TOOLS.contains(&name) =>
1993 {
1994 if let Some(request) = prepare_freeform_write(&self.journal, name, &call.arguments)?
1995 {
1996 store_result = false;
1997 let acknowledgement = request.acknowledgement();
1998 freeform_write = Some(request);
1999 acknowledgement
2000 } else {
2001 let object_ids = if name == CALL_RUST_BIN_TOOL {
2002 decode_managed_objects(ManagedObjectArguments::RustBinary(&call.arguments))?
2003 } else if name == ATTACH_OBJECT_WEB_LIB_TOOL {
2004 decode_managed_objects(ManagedObjectArguments::WebLibraryAttachment(
2005 &call.arguments,
2006 ))?
2007 } else {
2008 Vec::new()
2009 };
2010 let mut objects = Vec::with_capacity(object_ids.len());
2011 for object_id in object_ids {
2012 objects.push(self.resolve_object(&object_id)?.bytes);
2013 }
2014 let execution = self
2015 .api
2016 .managed_source_execute(
2017 &self.rust_lib_session_id,
2018 name,
2019 call.arguments.clone(),
2020 objects,
2021 )
2022 .await?;
2023 if let Some(snapshot) = execution.snapshot {
2024 managed_source_snapshot = Some(snapshot);
2025 store_result = false;
2026 }
2027 execution.text
2028 }
2029 }
2030 (name, Some(_)) => {
2031 anyhow::bail!("decoded contract for {name} did not match its dispatch lane")
2032 }
2033 (name, None) => anyhow::bail!("Tool {name} is not available"),
2034 };
2035 Ok(ToolOutcome {
2036 text,
2037 store_result,
2038 ok: true,
2039 end_session,
2040 freeform_write,
2041 managed_source_snapshot,
2042 })
2043 }
2044
2045 async fn execute_task_board_tool(
2046 &self,
2047 name: &str,
2048 arguments: &Value,
2049 ) -> anyhow::Result<String> {
2050 let board = self
2051 .api
2052 .task_board()
2053 .context("task board is not configured")?
2054 .clone();
2055 let name = name.to_owned();
2056 let arguments = arguments.clone();
2057 let user_id = self
2058 .root_node_ids
2059 .first()
2060 .context("session has no user root for task-category lookup")?
2061 .clone();
2062 tokio::task::spawn_blocking(move || -> anyhow::Result<String> {
2063 let output = match name.as_str() {
2064 "CreateTaskCategory" => serde_json::to_string_pretty(
2065 &board.create_category(serde_json::from_value(arguments)?)?,
2066 )?,
2067 "GetTaskCategory" => {
2068 let call: CategoryCall = serde_json::from_value(arguments)?;
2069 serde_json::to_string_pretty(&board.category(
2070 &call.category_id,
2071 kcode_task_board::BrowsePage {
2072 user_id,
2073 offset: call.offset,
2074 limit: call.limit,
2075 },
2076 )?)?
2077 }
2078 "RemoveTaskCategory" => {
2079 let call: CategoryId = serde_json::from_value(arguments)?;
2080 board.remove_category(&call.category_id)?;
2081 format!("Removed category {}.", call.category_id)
2082 }
2083 "CreateTask" => serde_json::to_string_pretty(
2084 &board.create_task(serde_json::from_value(arguments)?)?,
2085 )?,
2086 "GetTask" => {
2087 let call: TaskId = serde_json::from_value(arguments)?;
2088 serde_json::to_string_pretty(&board.task(&call.task_id)?)?
2089 }
2090 "UpdateTask" => serde_json::to_string_pretty(
2091 &board.update_task(serde_json::from_value(arguments)?)?,
2092 )?,
2093 "RemoveTask" => {
2094 let call: TaskId = serde_json::from_value(arguments)?;
2095 board.remove_task(&call.task_id)?;
2096 format!("Removed task {}.", call.task_id)
2097 }
2098 "GetTopTaskOrphan" => {
2099 let _: EmptyCall = serde_json::from_value(arguments)?;
2100 serde_json::to_string_pretty(&board.top_orphan()?)?
2101 }
2102 _ => anyhow::bail!("Tool {name} is not a task-board operation"),
2103 };
2104 Ok(output)
2105 })
2106 .await
2107 .context("task-board worker stopped")?
2108 }
2109
2110 fn assert_tool_allowed(&self, name: &str) -> anyhow::Result<()> {
2111 let write = matches!(
2112 name,
2113 "ConnectNodes"
2114 | "ConsolidateFanout"
2115 | "SetFixedConnection"
2116 | "CreateNode"
2117 | "UpdateNode"
2118 );
2119 anyhow::ensure!(
2120 !write || !matches!(self.mode, AgentMode::Conversation),
2121 "{name} requires the global Kweb write lane and is unavailable in a read-only conversation"
2122 );
2123 if name == "EndSession" {
2124 anyhow::ensure!(
2125 !matches!(self.mode, AgentMode::Conversation),
2126 "EndSession is unavailable in a conversation"
2127 );
2128 }
2129 Ok(())
2130 }
2131
2132 fn sync_kweb_boxes(&mut self) -> anyhow::Result<Vec<BoxId>> {
2133 let (updates, creates) = kweb_plan_projection(&self.plan);
2134 self.context
2135 .sync_chatend(&mut self.journal, now(), &updates, &creates)
2136 .map_err(anyhow::Error::new)
2137 }
2138
2139 fn record_tool_invocation(
2140 &mut self,
2141 name: &str,
2142 arguments: Value,
2143 ) -> anyhow::Result<RecordedToolInvocation> {
2144 let invocation = RecordedToolInvocation {
2145 invocation_id: Uuid::new_v4().to_string(),
2146 tool_instance: tool_instance(name),
2147 tool_name: name.into(),
2148 };
2149 self.journal.record(
2150 now(),
2151 EventKind::ToolInvoked {
2152 tool_instance: invocation.tool_instance.clone(),
2153 tool_name: invocation.tool_name.clone(),
2154 arguments,
2155 invocation_id: Some(invocation.invocation_id.clone()),
2156 },
2157 )?;
2158 Ok(invocation)
2159 }
2160
2161 fn record_tool_completion(
2162 &mut self,
2163 invocation: Option<&RecordedToolInvocation>,
2164 outcome: Value,
2165 ) -> anyhow::Result<EventId> {
2166 let (tool_instance, tool_name, invocation_id) = invocation
2167 .map(|invocation| {
2168 (
2169 invocation.tool_instance.clone(),
2170 invocation.tool_name.clone(),
2171 Some(invocation.invocation_id.clone()),
2172 )
2173 })
2174 .unwrap_or_else(|| ("call_ktool".into(), "call_ktool".into(), None));
2175 self.journal.record(
2176 now(),
2177 EventKind::ToolCompleted {
2178 tool_instance,
2179 tool_name,
2180 outcome,
2181 invocation_id,
2182 },
2183 )
2184 }
2185
2186 fn node_data(plan: &KwebPlan, context: &KwebContext, id: &str) -> anyhow::Result<PlannedNode> {
2187 if let Some(data) = plan.created(id) {
2188 return Ok(data.clone());
2189 }
2190 if let Some(data) = plan.updates.get(id) {
2191 return Ok(data.clone());
2192 }
2193 let node = context
2194 .node(id)
2195 .with_context(|| format!("Kweb context does not contain node {id}"))?;
2196 Ok(planned_node(node))
2197 }
2198
2199 fn put_node_data(plan: &mut KwebPlan, id: &str, data: PlannedNode) -> anyhow::Result<()> {
2200 if let Some(created) = plan.created_mut(id) {
2201 *created = data;
2202 } else {
2203 canonical_id(id)?;
2204 plan.updates.insert(id.to_owned(), data);
2205 }
2206 Ok(())
2207 }
2208
2209 fn connect_nodes(
2210 plan: &mut KwebPlan,
2211 context: &KwebContext,
2212 ids: Vec<String>,
2213 ) -> anyhow::Result<String> {
2214 for id in &ids {
2215 Self::ensure_known_node(plan, context, id)?;
2216 }
2217 for id in &ids {
2218 let mut data = Self::node_data(plan, context, id)?;
2219 let mut recent = ids
2220 .iter()
2221 .filter(|other| *other != id)
2222 .cloned()
2223 .collect::<Vec<_>>();
2224 for other in data.recent_connections {
2225 if &other != id && !recent.contains(&other) {
2226 recent.push(other);
2227 }
2228 }
2229 data.recent_connections = recent;
2230 Self::put_node_data(plan, id, data)?;
2231 }
2232 Ok(format!(
2233 "Staged connections among nodes {}.",
2234 ids.join(", ")
2235 ))
2236 }
2237
2238 fn consolidate_fanout(
2239 plan: &mut KwebPlan,
2240 context: &KwebContext,
2241 parent: String,
2242 fanout: Vec<String>,
2243 aggregator: String,
2244 ) -> anyhow::Result<String> {
2245 for id in std::iter::once(&parent)
2246 .chain(std::iter::once(&aggregator))
2247 .chain(fanout.iter())
2248 {
2249 Self::ensure_known_node(plan, context, id)?;
2250 }
2251 let mut parent_data = Self::node_data(plan, context, &parent)?;
2252 parent_data
2253 .recent_connections
2254 .retain(|id| !fanout.contains(id));
2255 if !parent_data.recent_connections.contains(&aggregator) {
2256 parent_data.recent_connections.push(aggregator.clone());
2257 }
2258 let mut aggregator_data = Self::node_data(plan, context, &aggregator)?;
2259 for id in fanout {
2260 if !aggregator_data.recent_connections.contains(&id) {
2261 aggregator_data.recent_connections.push(id);
2262 }
2263 }
2264 Self::put_node_data(plan, &parent, parent_data)?;
2265 Self::put_node_data(plan, &aggregator, aggregator_data)?;
2266 Ok(format!(
2267 "Staged fanout consolidation from node {parent} into node {aggregator}."
2268 ))
2269 }
2270
2271 fn set_fixed_connection(
2272 plan: &mut KwebPlan,
2273 context: &KwebContext,
2274 parent: String,
2275 child: Option<String>,
2276 slot: usize,
2277 ) -> anyhow::Result<String> {
2278 Self::ensure_known_node(plan, context, &parent)?;
2279 if let Some(child) = &child {
2280 Self::ensure_known_node(plan, context, child)?;
2281 anyhow::ensure!(child != &parent, "a node cannot connect to itself");
2282 }
2283 let mut data = Self::node_data(plan, context, &parent)?;
2284 if let Some(child) = child.clone() {
2285 anyhow::ensure!(
2286 slot <= data.fixed_connections.len() + 1,
2287 "fixed connection positions must remain contiguous"
2288 );
2289 data.fixed_connections.retain(|id| id != &child);
2290 if slot - 1 < data.fixed_connections.len() {
2291 data.fixed_connections[slot - 1] = child;
2292 } else {
2293 data.fixed_connections.push(child);
2294 }
2295 } else if slot > 0 && slot - 1 < data.fixed_connections.len() {
2296 data.fixed_connections.remove(slot - 1);
2297 }
2298 Self::put_node_data(plan, &parent, data)?;
2299 Ok(match child {
2300 Some(child) => {
2301 format!("Staged node {child} in fixed slot {slot} of node {parent}.")
2302 }
2303 None => format!("Cleared fixed slot {slot} of node {parent} in the staged plan."),
2304 })
2305 }
2306
2307 fn create_node(
2308 plan: &mut KwebPlan,
2309 context: &KwebContext,
2310 journal: &mut HistorySession,
2311 arguments: CreateNodeArguments,
2312 ) -> anyhow::Result<String> {
2313 let CreateNodeArguments {
2314 parents,
2315 owner,
2316 short_name,
2317 short_description,
2318 long_description,
2319 } = arguments;
2320 for id in parents.iter().chain(std::iter::once(&owner)) {
2321 if id != "self" && id != "unowned" {
2322 Self::ensure_known_node(plan, context, id)?;
2323 }
2324 }
2325 let pending = journal.allocate_pending_node(now())?.to_string();
2326 plan.creates.push(StagedNodeCreate {
2327 pending_id: pending.clone(),
2328 data: PlannedNode {
2329 short_name,
2330 short_description,
2331 long_description,
2332 owner,
2333 fixed_connections: Vec::new(),
2334 recent_connections: parents.clone(),
2335 objects: Vec::new(),
2336 attach_session_archive: true,
2337 },
2338 });
2339 for parent in parents {
2340 let mut data = Self::node_data(plan, context, &parent)?;
2341 data.recent_connections.retain(|id| id != &pending);
2342 data.recent_connections.insert(0, pending.clone());
2343 Self::put_node_data(plan, &parent, data)?;
2344 }
2345 Ok(format!("Created staged node {pending}."))
2346 }
2347
2348 fn update_node(
2349 plan: &mut KwebPlan,
2350 context: &KwebContext,
2351 id: String,
2352 owner: String,
2353 short_name: String,
2354 short_description: String,
2355 long_description: String,
2356 ) -> anyhow::Result<String> {
2357 Self::ensure_known_node(plan, context, &id)?;
2358 if owner != "self" && owner != "unowned" {
2359 Self::ensure_known_node(plan, context, &owner)?;
2360 }
2361 let mut data = Self::node_data(plan, context, &id)?;
2362 data.owner = owner;
2363 data.short_name = short_name;
2364 data.short_description = short_description;
2365 data.long_description = long_description;
2366 data.attach_session_archive = true;
2367 Self::put_node_data(plan, &id, data)?;
2368 Ok(format!("Staged the update to node {id}."))
2369 }
2370
2371 fn ensure_known_node(plan: &KwebPlan, context: &KwebContext, id: &str) -> anyhow::Result<()> {
2372 if id.starts_with("pending:") {
2373 anyhow::ensure!(
2374 plan.created(id).is_some(),
2375 "pending node {id} is not part of this session"
2376 );
2377 } else {
2378 canonical_id(id)?;
2379 anyhow::ensure!(
2380 context.contains_full_node(id) || plan.updates.contains_key(id),
2381 "node {id} is not loaded; call LoadNodes first"
2382 );
2383 }
2384 Ok(())
2385 }
2386
2387 fn finalize_kweb_session(&mut self) -> anyhow::Result<()> {
2388 if self.commit_receipt.is_some() {
2389 return Ok(());
2390 }
2391 self.journal.repair_unfinished_tools(now())?;
2392 self.journal.seal()?;
2393 let archive = self.journal.archive_bytes()?;
2394 let object_locations = self
2395 .journal
2396 .objects()
2397 .iter()
2398 .map(|(id, location)| (id.clone(), location.clone()))
2399 .collect::<Vec<_>>();
2400 let mut objects = BTreeMap::new();
2401 for (id, location) in object_locations {
2402 let pending_id = id.to_string();
2403 let transport_kind =
2404 kcode_kennedy_session_objects::staged_descriptor(&self.journal, &id)?
2405 .transport_kind;
2406 let bytes = encode_file(
2407 &pending_id,
2408 location.metadata.file_name.as_deref(),
2409 &location.metadata.media_type,
2410 transport_kind.as_deref(),
2411 self.journal.read_object(&id)?,
2412 )
2413 .with_context(|| format!("encoding staged object {pending_id}"))?;
2414 anyhow::ensure!(
2415 objects.insert(pending_id.clone(), bytes).is_none(),
2416 "duplicate staged object {pending_id}"
2417 );
2418 }
2419 let mut creates = BTreeMap::new();
2420 for create in &self.plan.creates {
2421 anyhow::ensure!(
2422 creates
2423 .insert(create.pending_id.clone(), create.data.clone())
2424 .is_none(),
2425 "duplicate staged node {}",
2426 create.pending_id
2427 );
2428 }
2429 let updates = self
2430 .plan
2431 .updates
2432 .iter()
2433 .map(|(node_id, data)| {
2434 node_id
2435 .parse::<NodeId>()
2436 .with_context(|| format!("{node_id:?} is not a canonical node ID"))
2437 .map(|node_id| (node_id, data.clone()))
2438 })
2439 .collect::<anyhow::Result<BTreeMap<_, _>>>()?;
2440 let result = self.api.commit_kweb_session(CommitRequest {
2441 idempotency_key: self.journal.state().metadata.session_id.clone(),
2442 author: self.commit_author.clone(),
2443 source_created_at: DateTime::parse_from_rfc3339(&self.started_at)
2444 .context("session start timestamp is invalid")?
2445 .with_timezone(&Utc),
2446 archive,
2447 objects,
2448 creates,
2449 updates,
2450 })?;
2451 self.journal
2452 .mark_completed(result.session_object_id.to_string());
2453 self.commit_receipt = Some(result);
2454 Ok(())
2455 }
2456
2457 fn prepare_free_time_round(&mut self) -> anyhow::Result<bool> {
2458 if !matches!(self.mode, AgentMode::FreeTime) {
2459 return Ok(false);
2460 }
2461 let Some(deadline) = deadline(&self.free_time) else {
2462 return Ok(false);
2463 };
2464 if Utc::now() >= deadline {
2465 self.free_time_end_reason = Some("deadline".into());
2466 self.journal.create_box(
2467 now(),
2468 "Self-time timer",
2469 BoxOwner::Controller,
2470 BoxContent::text(
2471 "The self-time deadline has arrived. Finish without starting more tool work.",
2472 ),
2473 )?;
2474 return Ok(true);
2475 }
2476 Ok(false)
2477 }
2478
2479 fn refresh_runtime_prompt(&mut self) -> anyhow::Result<()> {
2480 let Some(system_box) = self
2481 .journal
2482 .state()
2483 .boxes
2484 .values()
2485 .find(|state| matches!(state.owner, BoxOwner::System))
2486 .map(|state| state.id)
2487 else {
2488 return Ok(());
2489 };
2490 let current = self
2491 .journal
2492 .state()
2493 .boxes
2494 .get(&system_box)
2495 .context("system-prompt box disappeared")?
2496 .canonical
2497 .content
2498 .text
2499 .clone();
2500 let marker = "\n\nCurrent runtime\n\n";
2501 let Some((prefix, _)) = current.rsplit_once(marker) else {
2502 return Ok(());
2503 };
2504 let refreshed = format!(
2505 "{prefix}{marker}{}",
2506 render(RenderRequest::RuntimeDescription {
2507 model: &self.runtime.model,
2508 reasoning_effort: &self.runtime.reasoning_effort,
2509 current_time: Utc::now(),
2510 })?
2511 );
2512 if refreshed != current {
2513 self.journal
2514 .update_box(now(), system_box, BoxContent::text(refreshed))?;
2515 }
2516 Ok(())
2517 }
2518
2519 fn agent_request_timeout(&self) -> Option<Duration> {
2520 if matches!(self.mode, AgentMode::Conversation) && self.session_type == "conversation" {
2521 return Some(BROWSER_CONVERSATION_REQUEST_TIMEOUT);
2522 }
2523 if matches!(self.mode, AgentMode::Ingress { .. }) {
2524 return Some(HISTORY_INGRESS_REQUEST_TIMEOUT);
2525 }
2526 if matches!(self.mode, AgentMode::Wakeup) {
2527 return Some(WAKEUP_REQUEST_TIMEOUT);
2528 }
2529 if matches!(self.mode, AgentMode::FreeTime) {
2530 let deadline = deadline(&self.free_time)?;
2531 return Some(Duration::from_secs(
2532 (deadline - Utc::now()).num_seconds().max(1) as u64 + 360,
2533 ));
2534 }
2535 None
2536 }
2537
2538 pub fn refresh_telegram_group_context(
2539 &mut self,
2540 group_context: &Value,
2541 current_message_id: Option<&str>,
2542 ) -> anyhow::Result<()> {
2543 if self.session_type != "telegram-group" {
2544 return Ok(());
2545 }
2546 self.channel["groupContext"] = group_context.clone();
2547 self.group_context = group_context.clone();
2548 self.journal.create_box(
2549 now(),
2550 "Telegram group update",
2551 BoxOwner::Controller,
2552 BoxContent::text(kcode_telegram_session_coordinator::format_group_context(
2553 group_context,
2554 )),
2555 )?;
2556 self.recover_context_overflow(current_message_id, &[])?;
2557 Ok(())
2558 }
2559
2560 pub fn finalize_free_time(&mut self, reason: &str) -> anyhow::Result<()> {
2561 anyhow::ensure!(
2562 matches!(reason, "tool" | "deadline" | "hard-stop" | "user-stop"),
2563 "invalid self-time completion reason"
2564 );
2565 self.free_time["sliceEndedReason"] = json!(reason);
2566 self.free_time["sliceEndedAt"] = json!(now());
2567 self.pending_turn = false;
2568 self.pending_external_event_id = None;
2569 Ok(())
2570 }
2571
2572 pub fn commit_current_write_session(&mut self) -> anyhow::Result<()> {
2573 anyhow::ensure!(
2574 matches!(
2575 self.mode,
2576 AgentMode::FreeTime | AgentMode::Wakeup | AgentMode::Ingress { .. }
2577 ),
2578 "a read-only conversation cannot be committed as a Kweb write session"
2579 );
2580 self.finalize_kweb_session()?;
2581 self.completed = true;
2582 Ok(())
2583 }
2584
2585 pub fn snapshot(&self) -> anyhow::Result<Value> {
2586 let projection = self.journal.state().projection();
2587 let chatend_text = projection.render();
2588 let session_status = projection.status.clone();
2589 Ok(json!({
2590 "format":"kennedy-chatend",
2591 "version":1,
2592 "stateVersion":3,
2593 "sessionId":self.journal.state().metadata.session_id,
2594 "chatendMetadata":self.journal.state().metadata,
2595 "sessionType":self.session_type,
2596 "sourceSessionType":self.source_session_type,
2597 "channel":self.channel,
2598 "freeTime":self.free_time,
2599 "orchestration":self.orchestration,
2600 "provenanceId":self.provenance_id,
2601 "rustLibSessionId":self.rust_lib_session_id,
2602 "rootNodeIds":self.root_node_ids,
2603 "referenceRootNodeIds":self.reference_root_node_ids,
2604 "startedAt":self.started_at,
2605 "transcript":self.transcript,
2606 "pendingTurn":self.pending_turn,
2607 "pendingExternalEventId":self.pending_external_event_id,
2608 "roundsUsed":self.rounds_used,
2609 "completed":self.completed,
2610 "sessionObjectId":self.journal.state().completed_session_object,
2611 "commitReceipt":self.commit_receipt,
2612 "commitAuthor":self.commit_author,
2613 "providerModel":self.runtime.model,
2614 "kwebPlan":self.plan,
2615 "boxCount":self.journal.state().boxes.len(),
2616 "eventCount":self.journal.state().events.len(),
2617 "boxes":self.journal.state().boxes,
2618 "events":self.journal.state().events,
2619 "context":projection,
2620 "sessionStatus":session_status,
2621 "chatendText":chatend_text,
2622 }))
2623 }
2624
2625 pub async fn release_managed_sources(&self) {
2626 self.api
2627 .release_managed_sources(&self.rust_lib_session_id)
2628 .await;
2629 }
2630}
2631
2632impl<C, F> kcode_agent_runtime::SessionHost for KennedySessionHost<'_, C>
2633where
2634 C: FnMut(Value) -> F + Send,
2635 F: Future<Output = anyhow::Result<()>> + Send,
2636{
2637 fn prepare_round<'a>(
2638 &'a mut self,
2639 round: u64,
2640 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::RoundPreparation> {
2641 Box::pin(async move {
2642 self.session.rounds_used = round;
2643 self.session.refresh_runtime_prompt()?;
2644 self.deadline_after_response = self.session.prepare_free_time_round()?;
2645 let external_event_id = self.session.pending_external_event_id.clone();
2646 if self
2647 .session
2648 .recover_context_overflow(external_event_id.as_deref(), &[])?
2649 == ContextRecovery::Irreducible
2650 || (matches!(self.session.mode, AgentMode::Ingress { .. })
2651 && self.session.ingress_force_commit_requested())
2652 {
2653 return Ok(kcode_agent_runtime::RoundPreparation::Complete(None));
2654 }
2655 let input = self.session.journal.state().projection().render();
2656 Ok(kcode_agent_runtime::RoundPreparation::Run(
2657 kcode_agent_runtime::PreparedRound {
2658 input,
2659 model: self.session.runtime.model.clone(),
2660 reasoning_effort: self.session.runtime.reasoning_effort.clone(),
2661 tool_description: call_ktool_description(),
2662 timeout: self.session.agent_request_timeout(),
2663 },
2664 ))
2665 })
2666 }
2667
2668 fn record<'a>(
2669 &'a mut self,
2670 event: kcode_agent_runtime::SessionEvent,
2671 ) -> kcode_agent_runtime::HostFuture<'a, ()> {
2672 Box::pin(async move {
2673 match event {
2674 kcode_agent_runtime::SessionEvent::InferenceSubmitted {
2675 manifest_hash,
2676 model,
2677 ..
2678 } => {
2679 let projection = self.session.journal.state().projection();
2680 self.accounting = Some(kcode_intelligence_chatend::TopLevelCall::new(
2681 manifest_hash.clone(),
2682 model,
2683 ));
2684 self.session.journal.record(
2685 now(),
2686 EventKind::InferenceSubmitted {
2687 manifest_hash,
2688 estimated_input_tokens: projection.estimated_tokens,
2689 raw_estimated_input_tokens: Some(projection.raw_estimated_tokens),
2690 },
2691 )?;
2692 }
2693 kcode_agent_runtime::SessionEvent::ProviderInput { input, .. } => {
2694 self.session.journal.record(
2695 now(),
2696 EventKind::Note {
2697 label: "provider_input".into(),
2698 value: Value::String(input),
2699 },
2700 )?;
2701 }
2702 kcode_agent_runtime::SessionEvent::UsageUpdated { usage, .. } => {
2703 self.accounting
2704 .as_mut()
2705 .context("provider usage arrived before inference submission")?
2706 .usage_updated(&mut self.session.journal, &now(), &usage)?;
2707 }
2708 kcode_agent_runtime::SessionEvent::ProviderReceipt { usage, .. } => {
2709 self.accounting
2710 .take()
2711 .context("provider receipt arrived before inference submission")?
2712 .completed(&mut self.session.journal, &now(), usage.as_ref())?;
2713 }
2714 }
2715 let snapshot = self.session.snapshot()?;
2716 (self.checkpoint)(snapshot).await
2717 })
2718 }
2719
2720 fn execute_tool<'a>(
2721 &'a mut self,
2722 call: anyhow::Result<kcode_agent_runtime::ToolCall>,
2723 operation_id: Uuid,
2724 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionToolOutcome> {
2725 Box::pin(async move {
2726 if let Some(pending) = &self.pending_freeform_write {
2727 let text = format!(
2728 "{} is awaiting the complete file contents; no other Ktool can run before that output.",
2729 pending.request.write_tool()
2730 );
2731 self.session
2732 .record_tool_completion(None, json!({"ok":false,"result":text}))?;
2733 let text = provider_tool_result_with_context_footer(&self.session.journal, &text);
2734 (self.checkpoint)(self.session.snapshot()?).await?;
2735 return Ok(kcode_agent_runtime::SessionToolOutcome {
2736 text,
2737 ok: false,
2738 capture: Some(json!(true)),
2739 stop: false,
2740 finish_after_round: false,
2741 emitted_response: false,
2742 });
2743 }
2744 let tool_started_at = std::time::Instant::now();
2745 let mut created_call_box_id = None;
2746 let mut recorded_invocation = None;
2747 let transcript_start = self.session.transcript.len();
2748 let mut emitted_response = false;
2749 let mut outcome = match call {
2750 Ok(call) => {
2751 let call = ToolCall {
2752 name: call.name,
2753 arguments: call.arguments,
2754 };
2755 let call_name = format!("Kennedy tool call: {}", call.name);
2756 let call_content = tool_call_box_content(&call)?;
2757 recorded_invocation = Some(
2758 self.session
2759 .record_tool_invocation(&call.name, call.arguments.clone())?,
2760 );
2761 created_call_box_id = Some(self.session.journal.create_box(
2762 now(),
2763 call_name,
2764 BoxOwner::Kennedy,
2765 call_content,
2766 )?);
2767 let external_event_id = self.session.pending_external_event_id.clone();
2768 if self
2769 .session
2770 .recover_context_overflow(external_event_id.as_deref(), &[])?
2771 == ContextRecovery::Irreducible
2772 {
2773 ToolOutcome {
2774 text: CONTEXT_OVERFLOW_WARNING.into(),
2775 store_result: false,
2776 ok: false,
2777 end_session: false,
2778 freeform_write: None,
2779 managed_source_snapshot: None,
2780 }
2781 } else {
2782 match self.session.execute_tool(&call, operation_id).await {
2783 Ok(outcome) => {
2784 emitted_response = call.name == "EmitObject" && outcome.ok;
2785 outcome
2786 }
2787 Err(error) => ToolOutcome {
2788 text: format!("{} failed: {error}", call.name),
2789 store_result: call.name != "LoadNodes",
2790 ok: false,
2791 end_session: false,
2792 freeform_write: None,
2793 managed_source_snapshot: None,
2794 },
2795 }
2796 }
2797 }
2798 Err(error) => ToolOutcome {
2799 text: error.to_string(),
2800 store_result: true,
2801 ok: false,
2802 end_session: false,
2803 freeform_write: None,
2804 managed_source_snapshot: None,
2805 },
2806 };
2807 if let Some(snapshot) = outcome.managed_source_snapshot.take() {
2808 apply_snapshot(&mut self.session.journal, &now(), snapshot)?;
2809 outcome.store_result = false;
2810 }
2811 append_slow_tool_duration(&mut outcome.text, tool_started_at.elapsed());
2812 let capture = if let Some(request) = outcome.freeform_write.take() {
2813 self.pending_freeform_write = Some(PendingFreeformWrite {
2814 request,
2815 call_box_id: created_call_box_id
2816 .context("freeform write call box was not created")?,
2817 });
2818 Some(json!(true))
2819 } else {
2820 None
2821 };
2822 if outcome.store_result {
2823 self.session.journal.create_box(
2824 now(),
2825 "Kennedy tool result",
2826 BoxOwner::Controller,
2827 BoxContent::text(&outcome.text),
2828 )?;
2829 }
2830 let external_event_id = self.session.pending_external_event_id.clone();
2831 let recovery = self
2832 .session
2833 .recover_context_overflow(external_event_id.as_deref(), &[])?;
2834 let context_warning_added =
2835 self.session.transcript[transcript_start..]
2836 .iter()
2837 .any(|entry| {
2838 entry.get("contextOverflowWarning").and_then(Value::as_bool) == Some(true)
2839 });
2840 let mut provider_text = outcome.text.clone();
2841 if context_warning_added && !provider_text.contains(CONTEXT_OVERFLOW_WARNING) {
2842 if !provider_text.is_empty() {
2843 provider_text.push_str("\n\n");
2844 }
2845 provider_text.push_str(CONTEXT_OVERFLOW_WARNING);
2846 }
2847 let provider_text =
2848 provider_tool_result_with_context_footer(&self.session.journal, &provider_text);
2849 self.session.record_tool_completion(
2850 recorded_invocation.as_ref(),
2851 json!({"ok":outcome.ok,"result":outcome.text}),
2852 )?;
2853 let stop = recovery == ContextRecovery::Irreducible
2854 || (matches!(self.session.mode, AgentMode::Ingress { .. })
2855 && self.session.ingress_force_commit_requested())
2856 || (!matches!(self.session.mode, AgentMode::Ingress { .. })
2857 && self.session.journal.state().source_terminated);
2858 let snapshot = self.session.snapshot()?;
2859 (self.checkpoint)(snapshot).await?;
2860 Ok(kcode_agent_runtime::SessionToolOutcome {
2861 text: provider_text,
2862 ok: outcome.ok,
2863 capture,
2864 stop,
2865 finish_after_round: outcome.end_session,
2866 emitted_response,
2867 })
2868 })
2869 }
2870
2871 fn complete_capture<'a>(
2872 &'a mut self,
2873 _capture: Value,
2874 contents: String,
2875 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionControl> {
2876 Box::pin(async move {
2877 let pending = self
2878 .pending_freeform_write
2879 .take()
2880 .context("provider completed without a pending freeform write")?;
2881 let result_metadata = pending.request.clone();
2882 let outcome = self
2883 .session
2884 .complete_freeform_write(pending, contents)
2885 .await?;
2886 if outcome.store_result {
2887 self.session.journal.create_box(
2888 now(),
2889 "Kennedy tool result",
2890 BoxOwner::Controller,
2891 BoxContent::text(&outcome.text),
2892 )?;
2893 }
2894 self.session.journal.record(
2895 now(),
2896 EventKind::Note {
2897 label: "write_file_freeform_result".into(),
2898 value: result_metadata.result_record(outcome.ok, &outcome.text),
2899 },
2900 )?;
2901 let external_event_id = self.session.pending_external_event_id.clone();
2902 let recovery = self
2903 .session
2904 .recover_context_overflow(external_event_id.as_deref(), &[])?;
2905 let snapshot = self.session.snapshot()?;
2906 (self.checkpoint)(snapshot).await?;
2907 if recovery == ContextRecovery::Irreducible
2908 || (matches!(self.session.mode, AgentMode::Ingress { .. })
2909 && self.session.ingress_force_commit_requested())
2910 || (!matches!(self.session.mode, AgentMode::Ingress { .. })
2911 && self.session.journal.state().source_terminated)
2912 || self.deadline_after_response
2913 {
2914 return Ok(kcode_agent_runtime::SessionControl::Complete(None));
2915 }
2916 self.session.journal.create_box(
2917 now(),
2918 controller_box_name(&self.session.mode),
2919 BoxOwner::Controller,
2920 BoxContent::text(controller_message(
2921 &self.session.mode,
2922 &self.session.free_time,
2923 )),
2924 )?;
2925 Ok(kcode_agent_runtime::SessionControl::Continue)
2926 })
2927 }
2928
2929 fn complete_round<'a>(
2930 &'a mut self,
2931 completion: kcode_agent_runtime::RoundCompletion,
2932 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionControl> {
2933 Box::pin(async move {
2934 let answer = completion.answer.trim().to_owned();
2935 let mut completion_recovery = ContextRecovery::NotNeeded;
2936 if !answer.is_empty() {
2937 let mut content = BoxContent::text(answer.clone());
2938 if let Some(id) = &self.session.pending_external_event_id {
2939 content.metadata["externalEventId"] = json!(id);
2940 }
2941 self.session.journal.create_box(
2942 now(),
2943 "Kennedy message",
2944 BoxOwner::Kennedy,
2945 content,
2946 )?;
2947 let mut transcript = json!({"role":"kennedy","content":answer});
2948 if let Some(id) = &self.session.pending_external_event_id {
2949 transcript["externalEventId"] = json!(id);
2950 }
2951 self.session.transcript.push(transcript);
2952 let external_event_id = self.session.pending_external_event_id.clone();
2953 completion_recovery = self
2954 .session
2955 .recover_context_overflow(external_event_id.as_deref(), &[])?;
2956 }
2957 let snapshot = self.session.snapshot()?;
2958 (self.checkpoint)(snapshot).await?;
2959 if completion_recovery == ContextRecovery::Irreducible
2960 || (matches!(self.session.mode, AgentMode::Ingress { .. })
2961 && self.session.ingress_force_commit_requested())
2962 || (!matches!(self.session.mode, AgentMode::Ingress { .. })
2963 && self.session.journal.state().source_terminated)
2964 {
2965 return Ok(kcode_agent_runtime::SessionControl::Complete(None));
2966 }
2967 if completion.finish_requested || self.deadline_after_response {
2968 return Ok(kcode_agent_runtime::SessionControl::Complete(
2969 (!answer.is_empty()).then_some(answer),
2970 ));
2971 }
2972 if matches!(self.session.mode, AgentMode::Conversation) && !answer.is_empty() {
2973 return Ok(kcode_agent_runtime::SessionControl::Complete(Some(answer)));
2974 }
2975 if matches!(self.session.mode, AgentMode::Conversation) && completion.emitted_response {
2976 return Ok(kcode_agent_runtime::SessionControl::Complete(None));
2977 }
2978 let solo_ingress_response =
2979 matches!(self.session.mode, AgentMode::Ingress { .. }) && !answer.is_empty();
2980 anyhow::ensure!(
2981 completion.used_tool || solo_ingress_response,
2982 "provider completed without a response or tool call"
2983 );
2984 self.session.journal.create_box(
2985 now(),
2986 controller_box_name(&self.session.mode),
2987 BoxOwner::Controller,
2988 BoxContent::text(controller_message(
2989 &self.session.mode,
2990 &self.session.free_time,
2991 )),
2992 )?;
2993 Ok(kcode_agent_runtime::SessionControl::Continue)
2994 })
2995 }
2996}
2997
2998impl kcode_agent_runtime::Host for KennedySubagentHost<'_> {
2999 fn render_tool_call(&mut self, call: &kcode_agent_runtime::ToolCall) -> anyhow::Result<String> {
3000 Ok(tool_call_box_content(&ToolCall {
3001 name: call.name.clone(),
3002 arguments: call.arguments.clone(),
3003 })?
3004 .text)
3005 }
3006
3007 fn execute_tool<'a>(
3008 &'a mut self,
3009 call: kcode_agent_runtime::ToolCall,
3010 operation_id: Uuid,
3011 budget: kcode_agent_runtime::ContextBudget,
3012 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::ToolOutcome> {
3013 Box::pin(async move {
3014 let call = ToolCall {
3015 name: call.name,
3016 arguments: call.arguments,
3017 };
3018 if let Some(reason) = subagent_unavailable_reason(&call.name) {
3019 return Ok(kcode_agent_runtime::ToolOutcome::failure(reason));
3020 }
3021 if budget.estimated_tokens() > budget.max_input_tokens() {
3022 return Ok(kcode_agent_runtime::ToolOutcome::failure(
3023 "The Ktool call was not run because its retained invocation would exceed the subagent context limit.",
3024 ));
3025 }
3026 if !subagent_managed_write_fits(&self.context, &call, &budget) {
3027 return Ok(kcode_agent_runtime::ToolOutcome::failure(
3028 "The managed-source write was not run because its resulting current state would exceed the subagent context limit.",
3029 ));
3030 }
3031
3032 let tool_started_at = std::time::Instant::now();
3033
3034 if call.name == "LoadNodes" {
3035 let Some(DecodedTool::LoadNodes(identifiers)) =
3036 decode(&call.name, &call.arguments)?
3037 else {
3038 return Ok(kcode_agent_runtime::ToolOutcome::failure(
3039 "LoadNodes did not match its tool contract.",
3040 ));
3041 };
3042 load_durable_batch(
3043 self.session.api.kmap(),
3044 self.context.kweb_mut(),
3045 &identifiers,
3046 )?;
3047 let (updates, creates) = kweb_plan_projection(&self.session.plan);
3048 let changes = self.context.reconcile_kweb(&updates, &creates)?;
3049 let displayed_state_keys = changes.displayed_state_keys();
3050 let mut text = if changes.is_empty() {
3051 "LoadNodes completed. The subagent Kweb projection was already current.".into()
3052 } else {
3053 changes.display_text()
3054 };
3055 append_slow_tool_duration(&mut text, tool_started_at.elapsed());
3056 return Ok(kcode_agent_runtime::ToolOutcome {
3057 text,
3058 ok: true,
3059 state_updates: changes.updates,
3060 displayed_state_keys,
3061 capture: None,
3062 });
3063 }
3064
3065 if is_kweb_mutation(&call.name) {
3066 self.session.assert_tool_allowed(&call.name)?;
3067 let decoded = decode(&call.name, &call.arguments)?
3068 .with_context(|| format!("{} did not match its tool contract", call.name))?;
3069 let referenced_pending = referenced_pending_nodes(&decoded);
3070 let prior_create_count = self.session.plan.creates.len();
3071 let mut text = execute_kweb_mutation(
3072 &call.name,
3073 decoded,
3074 self.context.kweb(),
3075 &mut self.session.plan,
3076 &mut self.session.journal,
3077 )?;
3078 self.context.include_staged_nodes(
3079 referenced_pending.into_iter().chain(
3080 self.session.plan.creates[prior_create_count..]
3081 .iter()
3082 .map(|create| create.pending_id.clone()),
3083 ),
3084 );
3085 let (updates, creates) = kweb_plan_projection(&self.session.plan);
3086 let changes = self.context.reconcile_kweb(&updates, &creates)?;
3087 append_slow_tool_duration(&mut text, tool_started_at.elapsed());
3088 return Ok(kcode_agent_runtime::ToolOutcome {
3089 text,
3090 ok: true,
3091 state_updates: changes.updates,
3092 displayed_state_keys: Vec::new(),
3093 capture: None,
3094 });
3095 }
3096
3097 if let Some(request) = decode_freeform_write(&call.name, &call.arguments)? {
3098 if !self.context.source_is_open(request.kind(), request.name()) {
3099 return Ok(kcode_agent_runtime::ToolOutcome::failure(format!(
3100 "{} {:?} is not open in this subagent context. Call {} first.",
3101 request.kind().label(),
3102 request.name(),
3103 request.kind().open_tool()
3104 )));
3105 }
3106 let acknowledgement = request.acknowledgement();
3107 let id = Uuid::new_v4().to_string();
3108 self.captures.insert(id.clone(), request);
3109 return Ok(kcode_agent_runtime::ToolOutcome {
3110 text: acknowledgement,
3111 ok: true,
3112 state_updates: Vec::new(),
3113 displayed_state_keys: Vec::new(),
3114 capture: Some(Value::String(id)),
3115 });
3116 }
3117
3118 let mut outcome = match self.session.execute_tool(&call, operation_id).await {
3119 Ok(outcome) => outcome,
3120 Err(error) => {
3121 let mut text = format!("{} failed: {error}", call.name);
3122 append_slow_tool_duration(&mut text, tool_started_at.elapsed());
3123 return Ok(kcode_agent_runtime::ToolOutcome::failure(text));
3124 }
3125 };
3126 let displays_managed_snapshot = outcome
3127 .managed_source_snapshot
3128 .as_ref()
3129 .is_some_and(|snapshot| result_displays_snapshot(&outcome.text, snapshot));
3130 append_slow_tool_duration(&mut outcome.text, tool_started_at.elapsed());
3131 let (state_updates, displayed_state_keys) =
3132 if let Some(snapshot) = outcome.managed_source_snapshot.take() {
3133 let state = self.context.apply_source_snapshot(snapshot);
3134 let displayed = displays_managed_snapshot.then_some(state.key);
3135 (
3136 state.update.into_iter().collect(),
3137 displayed.into_iter().collect(),
3138 )
3139 } else {
3140 (Vec::new(), Vec::new())
3141 };
3142 let capture = outcome.freeform_write.take().map(|request| {
3143 let id = Uuid::new_v4().to_string();
3144 self.captures.insert(id.clone(), request);
3145 Value::String(id)
3146 });
3147 Ok(kcode_agent_runtime::ToolOutcome {
3148 text: outcome.text,
3149 ok: outcome.ok,
3150 state_updates,
3151 displayed_state_keys,
3152 capture,
3153 })
3154 })
3155 }
3156
3157 fn complete_capture<'a>(
3158 &'a mut self,
3159 capture: Value,
3160 contents: String,
3161 budget: kcode_agent_runtime::ContextBudget,
3162 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::ToolOutcome> {
3163 Box::pin(async move {
3164 let id = capture
3165 .as_str()
3166 .context("subagent freeform capture token is invalid")?;
3167 let request = self
3168 .captures
3169 .remove(id)
3170 .context("subagent freeform capture token is unknown")?;
3171 self.session
3172 .complete_subagent_freeform_write(&mut self.context, request, contents, &budget)
3173 .await
3174 })
3175 }
3176
3177 fn record(&mut self, event: kcode_agent_runtime::AuditEvent) -> anyhow::Result<()> {
3178 kcode_intelligence_chatend::record_subagent_event(&mut self.session.journal, &now(), &event)
3179 }
3180}
3181
3182fn cost_summary(label: &str, estimated_cost_usd_nanos: u64, unpriced_calls: u64) -> String {
3183 render(RenderRequest::CostSummary {
3184 label,
3185 estimated_cost_usd_nanos,
3186 unpriced_calls,
3187 })
3188 .expect("cost-summary rendering is infallible")
3189}
3190
3191fn restore_kweb_context(journal: &HistorySession, context: &mut KwebContext) -> anyhow::Result<()> {
3192 let Some(tool) = journal.state().tools.get(KWEB_TOOL_INSTANCE) else {
3193 return Ok(());
3194 };
3195 let mut nodes = BTreeMap::new();
3196 for slot in &tool.slots {
3197 let state = journal
3198 .state()
3199 .box_state(slot.box_id)
3200 .context("Kweb slot references a missing box")?;
3201 if let Some(node) = state.canonical.content.metadata.get("storedNode") {
3202 let node = match serde_json::from_value::<KwebNode>(node.clone()) {
3203 Ok(node) => node,
3204 Err(_) => node_from_value(node).context("decoding a stored Kweb context node")?,
3205 };
3206 nodes.insert(node.id.clone(), node);
3207 }
3208 }
3209 let mut direct = journal
3210 .state()
3211 .events
3212 .iter()
3213 .flat_map(|event| {
3214 let EventKind::ToolInvoked {
3215 tool_name,
3216 arguments,
3217 ..
3218 } = &event.kind
3219 else {
3220 return Vec::new();
3221 };
3222 match tool_name.as_str() {
3223 "LoadNodes" => arguments
3224 .get("identifiers")
3225 .and_then(Value::as_array)
3226 .into_iter()
3227 .flatten()
3228 .filter_map(Value::as_str)
3229 .map(str::to_owned)
3230 .collect(),
3231 "LoadNode" => arguments
3232 .get("identifier")
3233 .and_then(Value::as_str)
3234 .map(str::to_owned)
3235 .into_iter()
3236 .collect(),
3237 _ => Vec::new(),
3238 }
3239 })
3240 .collect::<Vec<_>>();
3241 if direct.is_empty() {
3242 direct = context.root_node_ids().to_vec();
3243 }
3244 context
3245 .restore(nodes.into_values(), direct)
3246 .map_err(anyhow::Error::new)
3247}
3248
3249fn is_terminal_external_response(entry: &Value) -> bool {
3250 match entry.get("role").and_then(Value::as_str) {
3251 Some("kennedy") => true,
3252 Some("system") => {
3253 entry.get("contextOverflowWarning").and_then(Value::as_bool) != Some(true)
3254 }
3255 _ => false,
3256 }
3257}
3258
3259fn restore_pending_turn(restored: Option<&Value>, transcript: &[Value]) -> (bool, Option<String>) {
3260 let mut unanswered_external = Vec::<String>::new();
3261 let mut answered_external = HashSet::<String>::new();
3262 for entry in transcript {
3263 if entry.get("role").and_then(Value::as_str) == Some("user") {
3264 if let Some(id) = entry.get("externalEventId").and_then(Value::as_str) {
3265 answered_external.remove(id);
3266 unanswered_external.retain(|candidate| candidate != id);
3267 unanswered_external.push(id.to_owned());
3268 }
3269 continue;
3270 }
3271 if !is_terminal_external_response(entry) {
3272 continue;
3273 }
3274 if let Some(id) = entry.get("externalEventId").and_then(Value::as_str) {
3275 answered_external.insert(id.to_owned());
3276 unanswered_external.retain(|candidate| candidate != id);
3277 } else if entry.get("userStopped").and_then(Value::as_bool) == Some(true)
3278 && let Some(id) = unanswered_external.pop()
3279 {
3280 answered_external.insert(id);
3281 }
3282 }
3283 let journal_pending_external = unanswered_external.last().cloned();
3284 let restored_pending = restored
3285 .and_then(|state| state.get("pendingTurn"))
3286 .and_then(Value::as_bool)
3287 .unwrap_or(false);
3288 let restored_external = restored
3289 .and_then(|state| state.get("pendingExternalEventId"))
3290 .and_then(Value::as_str);
3291 let restored_external_answered =
3292 restored_external.is_some_and(|id| answered_external.contains(id));
3293 let pending_external_event_id = journal_pending_external.or_else(|| {
3294 restored_external
3295 .filter(|id| !answered_external.contains(*id))
3296 .map(str::to_owned)
3297 });
3298 let pending_turn =
3299 pending_external_event_id.is_some() || (restored_pending && !restored_external_answered);
3300 (pending_turn, pending_external_event_id)
3301}
3302
3303fn kweb_node_draft(node: &PlannedNode) -> NodeDraft {
3304 NodeDraft {
3305 short_name: node.short_name.clone(),
3306 short_description: node.short_description.clone(),
3307 long_description: node.long_description.clone(),
3308 owner: node.owner.clone(),
3309 fixed_connections: node.fixed_connections.clone(),
3310 recent_connections: node.recent_connections.clone(),
3311 objects: node.objects.clone(),
3312 }
3313}
3314
3315fn planned_node(node: &KwebNode) -> PlannedNode {
3316 PlannedNode {
3317 short_name: node.short_name.clone(),
3318 short_description: node.short_description.clone(),
3319 long_description: node.long_description.clone(),
3320 owner: node.owner.clone(),
3321 fixed_connections: node
3322 .fixed_connections
3323 .iter()
3324 .map(|connection| connection.id.clone())
3325 .collect(),
3326 recent_connections: node
3327 .recent_connections
3328 .iter()
3329 .map(|connection| connection.id.clone())
3330 .collect(),
3331 objects: node.objects.clone(),
3332 attach_session_archive: true,
3333 }
3334}
3335
3336fn session_kind(session_type: &str, mode: &AgentMode) -> SessionKind {
3337 if matches!(mode, AgentMode::Ingress { .. }) {
3338 return SessionKind::HistoryIngress;
3339 }
3340 match session_type {
3341 "conversation" => SessionKind::Conversation,
3342 "telegram" => SessionKind::Telegram,
3343 "telegram-group" => SessionKind::TelegramGroup,
3344 "free-time" => SessionKind::SelfTime,
3345 "wakeup" => SessionKind::Other("wakeup".into()),
3346 "audio" => SessionKind::AudioIngress,
3347 other => SessionKind::Other(other.into()),
3348 }
3349}
3350
3351fn tool_instance(name: &str) -> String {
3352 if name == "LoadNodes" {
3353 return KWEB_TOOL_INSTANCE.into();
3354 }
3355 format!("{name}:{}", Uuid::new_v4())
3356}
3357
3358fn canonical_id(value: &str) -> anyhow::Result<String> {
3359 value
3360 .parse::<NodeId>()
3361 .with_context(|| format!("{value:?} is not a canonical node ID"))?;
3362 Ok(value.into())
3363}
3364
3365fn image_extension(media_type: &str) -> &'static str {
3366 match media_type
3367 .split(';')
3368 .next()
3369 .unwrap_or(media_type)
3370 .trim()
3371 .to_ascii_lowercase()
3372 .as_str()
3373 {
3374 "image/jpeg" => "jpg",
3375 "image/webp" => "webp",
3376 _ => "png",
3377 }
3378}
3379
3380fn call_ktool_description() -> String {
3381 render(RenderRequest::CallKtoolDescription).expect("Ktool-description rendering is infallible")
3382}
3383
3384fn now() -> String {
3385 Utc::now().to_rfc3339()
3386}
3387
3388fn deadline(value: &Value) -> Option<DateTime<Utc>> {
3389 value
3390 .get("deadlineAt")
3391 .and_then(Value::as_str)
3392 .and_then(|value| DateTime::parse_from_rfc3339(value).ok())
3393 .map(|value| value.with_timezone(&Utc))
3394}
3395
3396fn controller_box_name(mode: &AgentMode) -> &'static str {
3397 match mode {
3398 AgentMode::Conversation => "Turn continuation",
3399 AgentMode::FreeTime => "Self-time continuation",
3400 AgentMode::Wakeup => "Wakeup continuation",
3401 AgentMode::Ingress { .. } => "History-ingress continuation",
3402 }
3403}
3404
3405fn controller_message(mode: &AgentMode, free_time: &Value) -> String {
3406 let mode = match mode {
3407 AgentMode::Conversation => "conversation",
3408 AgentMode::FreeTime => "free-time",
3409 AgentMode::Wakeup => "wakeup",
3410 AgentMode::Ingress { .. } => "ingress",
3411 };
3412 render(RenderRequest::ControllerMessage { mode, free_time })
3413 .expect("known controller modes render successfully")
3414}
3415
3416#[cfg(test)]
3417mod tests {
3418 use std::time::{SystemTime, UNIX_EPOCH};
3419
3420 use kcode_session_history::{Config as HistoryConfig, NewSession, SessionHistory};
3421
3422 use super::*;
3423
3424 fn test_journal(label: &str) -> (std::path::PathBuf, HistorySession) {
3425 let root = std::env::temp_dir().join(format!(
3426 "kcode-kennedy-sessions-{label}-{}-{}",
3427 std::process::id(),
3428 SystemTime::now()
3429 .duration_since(UNIX_EPOCH)
3430 .unwrap()
3431 .as_nanos()
3432 ));
3433 let history = SessionHistory::open(HistoryConfig {
3434 directory: root.join("sessions"),
3435 completed_list: root.join("completed.jsonl"),
3436 provider_cost_compatibility: None,
3437 })
3438 .unwrap();
3439 let journal = history
3440 .create_session(NewSession {
3441 kind: SessionKind::SelfTime,
3442 created_at: "2026-08-05T00:00:00Z".into(),
3443 effective_context_tokens: 10_000,
3444 channel: Value::Null,
3445 })
3446 .unwrap();
3447 (root, journal)
3448 }
3449
3450 #[test]
3451 fn subagents_reject_parent_controls_but_allow_delivery_effects() {
3452 for unavailable in [
3453 "RunSubagent",
3454 "EndSession",
3455 "DehydrateBoxes",
3456 "SummarizeBox",
3457 "HydrateBox",
3458 "BoxesIntoObjects",
3459 ] {
3460 assert!(subagent_unavailable_reason(unavailable).is_some());
3461 }
3462 for delegated in [
3463 "EmitObject",
3464 "SendTelegramDM",
3465 "SendTelegramGroupMessage",
3466 "LoadNodes",
3467 "ExtractDocumentText",
3468 ] {
3469 assert_eq!(subagent_unavailable_reason(delegated), None);
3470 }
3471 }
3472
3473 #[test]
3474 fn only_complete_snapshot_results_claim_to_display_managed_state() {
3475 let snapshot = SourceSnapshot {
3476 kind: kcode_dev_tools::ManagedSourceKind::RustLibrary,
3477 name: "example".into(),
3478 text: "complete source".into(),
3479 };
3480
3481 assert!(result_displays_snapshot("complete source", &snapshot));
3482 assert!(!result_displays_snapshot(
3483 "Wrote file src/lib.rs in Rust library example.",
3484 &snapshot
3485 ));
3486 }
3487
3488 #[test]
3489 fn child_kweb_mutation_changes_plan_without_creating_parent_boxes() {
3490 let (root, mut journal) = test_journal("subagent-kweb");
3491 let node_id = "AAAAAAAE".to_owned();
3492 let mut context = KwebContext::new(vec![node_id.clone()]).unwrap();
3493 context
3494 .apply_load(
3495 KwebNode {
3496 id: node_id.clone(),
3497 short_name: "Old".into(),
3498 short_description: "Old summary".into(),
3499 long_description: "Old details".into(),
3500 owner: "self".into(),
3501 fixed_connections: Vec::new(),
3502 recent_connections: Vec::new(),
3503 objects: Vec::new(),
3504 last_modified_by: "test".into(),
3505 last_modified_at: None,
3506 },
3507 Vec::new(),
3508 )
3509 .unwrap();
3510 let mut plan = KwebPlan::default();
3511
3512 let result = execute_kweb_mutation(
3513 "UpdateNode",
3514 DecodedTool::UpdateNode {
3515 id: node_id.clone(),
3516 owner: "self".into(),
3517 short_name: "New".into(),
3518 short_description: "New summary".into(),
3519 long_description: "New details".into(),
3520 },
3521 &context,
3522 &mut plan,
3523 &mut journal,
3524 )
3525 .unwrap();
3526
3527 assert_eq!(result, format!("Staged the update to node {node_id}."));
3528 assert_eq!(plan.updates[&node_id].long_description, "New details");
3529 assert!(journal.state().boxes.is_empty());
3530 std::fs::remove_dir_all(root).unwrap();
3531 }
3532}