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