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