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