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