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