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