#![forbid(unsafe_code)]
mod services;
pub use kcode_kennedy_session_objects::ResolvedObject;
pub use kcode_telegram_session_coordinator::validate_file_name as validate_delivery_file_name;
pub use services::{Api as Service, LocalServices as Capabilities};
use std::{
collections::{BTreeMap, HashMap, HashSet},
future::Future,
time::Duration,
};
use anyhow::Context as _;
use chrono::{DateTime, Utc};
use kcode_commit_session::{CommitReceipt, CommitRequest, PlannedNode};
use kcode_dev_tools::{
ATTACH_OBJECT_WEB_LIB_TOOL, CALL_RUST_BIN_TOOL, RUST_BIN_TOOLS, RUST_LIB_TOOLS, WEB_LIB_TOOLS,
WRITE_RUST_BIN_TOOL, WRITE_RUST_LIB_TOOL, WRITE_WEB_LIB_TOOL, proposed_write_snapshot,
};
use kcode_dev_tools_chatend::{
FreeformWrite, SourceSnapshot, apply_snapshot, prepare_freeform_write, source_box_id,
};
use kcode_history_ingress_context::{
Outcome as HistoryIngressContextOutcome, RecoveryOutcome as ContextRecoveryOutcome,
};
use kcode_kennedy_kweb_loader::{load_durable_batch, node_from_value};
use kcode_kennedy_session_presentation::{RenderRequest, render};
use kcode_kennedy_session_tool_contracts::{
DecodedTool, ManagedObjectArguments, ValidationRequest, decode, decode_managed_objects,
validate,
};
use kcode_kweb_context::{
Context as KwebContext, Node as KwebNode, NodeDraft, StagedCreate as KwebStagedCreate,
};
use kcode_kweb_db::NodeId;
use kcode_server_object_envelopes::encode_file;
use kcode_session_history::{
NewSession, Session as HistorySession,
chatend::{
BoxContent, BoxId, BoxOwner, EventId, EventKind, Representation, SessionKind,
SessionMetadata,
},
};
use kcode_speech_classification::KTOOLS as SPEECH_CLASSIFICATION_TOOLS;
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use sha2::{Digest, Sha256};
use uuid::Uuid;
const AGENT_LOOP_ROUND_LIMIT: u64 = 100;
const BROWSER_CONVERSATION_REQUEST_TIMEOUT: Duration = Duration::from_secs(90 * 60);
const HISTORY_INGRESS_REQUEST_TIMEOUT: Duration = Duration::from_secs(90 * 60);
const WAKEUP_REQUEST_TIMEOUT: Duration = Duration::from_secs(90 * 60);
const MAX_MEDIA_ENRICHMENT_BYTES: u64 = 20 * 1024 * 1024;
const KWEB_TOOL_INSTANCE: &str = "kweb";
const CONTEXT_OVERFLOW_WARNING_BOX_NAME: &str = "Context overflow warning";
const 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";
const INGRESS_FORCE_COMMIT_NOTE: &str = "ingress_force_commit";
#[derive(Clone, Debug)]
pub struct RuntimeModel {
pub model: String,
pub reasoning_effort: String,
pub context_window_tokens: u64,
}
impl RuntimeModel {
pub fn from_intelligence(runtime: kcode_intelligence_router::RuntimeModel) -> Self {
Self {
model: runtime.model,
reasoning_effort: runtime.reasoning_effort,
context_window_tokens: runtime.context_window_tokens,
}
}
fn attribution(&self) -> String {
format!("{}-{}", self.model, self.reasoning_effort)
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum AgentMode {
Conversation,
FreeTime,
Wakeup,
Ingress { record_id: Option<String> },
}
#[derive(Clone, Debug)]
pub struct SessionOptions {
pub session_type: String,
pub root_node_ids: Vec<String>,
pub reference_root_node_ids: Vec<String>,
pub channel: Value,
pub free_time: Value,
pub orchestration: Value,
pub provenance_id: Option<String>,
pub mode: AgentMode,
pub source_session_type: Option<String>,
pub group_context: Value,
pub rust_lib_session_id: Option<String>,
}
impl SessionOptions {
pub fn conversation(session_type: impl Into<String>, roots: Vec<String>) -> Self {
Self {
session_type: session_type.into(),
root_node_ids: roots,
reference_root_node_ids: Vec::new(),
channel: Value::Null,
free_time: Value::Null,
orchestration: json!({"owner":"backend","status":"idle"}),
provenance_id: None,
mode: AgentMode::Conversation,
source_session_type: None,
group_context: Value::Null,
rust_lib_session_id: None,
}
}
}
fn restore_session_type(options: &mut SessionOptions, state: &Value) {
if !matches!(&options.mode, AgentMode::Ingress { .. }) {
options.session_type = state
.get("sessionType")
.and_then(Value::as_str)
.unwrap_or(&options.session_type)
.to_owned();
}
}
fn restore_commit_receipt(restored: Option<&Value>) -> anyhow::Result<Option<CommitReceipt>> {
restored
.and_then(|state| state.get("commitReceipt"))
.filter(|receipt| !receipt.is_null())
.cloned()
.map(serde_json::from_value)
.transpose()
.context("decoding the stored session commit receipt")
}
#[derive(Clone, Debug, Default, Deserialize, Serialize)]
#[serde(rename_all = "camelCase")]
struct KwebPlan {
creates: Vec<StagedNodeCreate>,
updates: BTreeMap<String, PlannedNode>,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase")]
struct StagedNodeCreate {
pending_id: String,
data: PlannedNode,
}
impl KwebPlan {
fn restore(restored: Option<&Value>, journal: &HistorySession) -> anyhow::Result<Self> {
if let Some(plan) = restored.and_then(|state| state.get("kwebPlan")) {
return serde_json::from_value(plan.clone()).context("decoding the staged Kweb plan");
}
let latest = journal.state().events.iter().rev().find_map(|event| {
let EventKind::KwebPlanChanged { operation } = &event.kind else {
return None;
};
operation.get("plan")
});
latest
.cloned()
.map(serde_json::from_value)
.transpose()
.context("decoding the staged Kweb plan")
.map(Option::unwrap_or_default)
}
fn created(&self, id: &str) -> Option<&PlannedNode> {
self.creates
.iter()
.find(|create| create.pending_id == id)
.map(|create| &create.data)
}
fn created_mut(&mut self, id: &str) -> Option<&mut PlannedNode> {
self.creates
.iter_mut()
.find(|create| create.pending_id == id)
.map(|create| &mut create.data)
}
}
pub struct Session {
api: Service,
runtime: RuntimeModel,
journal: HistorySession,
plan: KwebPlan,
pub session_type: String,
pub channel: Value,
pub free_time: Value,
pub orchestration: Value,
pub provenance_id: Option<String>,
pub rust_lib_session_id: String,
pub root_node_ids: Vec<String>,
pub reference_root_node_ids: Vec<String>,
pub started_at: String,
pub transcript: Vec<Value>,
pub pending_turn: bool,
pub pending_external_event_id: Option<String>,
pub completed: bool,
pub rounds_used: u64,
commit_receipt: Option<CommitReceipt>,
commit_author: String,
mode: AgentMode,
source_session_type: Option<String>,
group_context: Value,
context: KwebContext,
free_time_end_reason: Option<String>,
fatal_persistence_error: Option<String>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum InputStage {
Accepted,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum ContextRecovery {
NotNeeded,
Recovered,
Irreducible,
}
fn render_load_nodes_result(
journal: &HistorySession,
changed_box_ids: &[BoxId],
) -> anyhow::Result<String> {
let changed_box_ids = changed_box_ids
.iter()
.map(ToString::to_string)
.collect::<Vec<_>>();
let projected_boxes = journal
.state()
.projection()
.items
.into_iter()
.filter(|item| !item.marker)
.map(|item| (item.box_id.to_string(), item.text))
.collect::<Vec<_>>();
render(RenderRequest::LoadNodes {
changed_box_ids: &changed_box_ids,
projected_boxes: &projected_boxes,
})
}
fn provider_tool_result_with_context_footer(journal: &HistorySession, result: &str) -> String {
let footer = journal.state().projection().footer;
render(RenderRequest::ProviderFooter {
result,
footer: &footer,
})
.expect("provider-footer rendering is infallible")
}
fn append_slow_tool_duration(text: &mut String, elapsed: Duration) {
*text = render(RenderRequest::SlowTool { text, elapsed })
.expect("slow-tool rendering is infallible");
}
fn render_web_search_result(
result: &kcode_intelligence_router::SearchResponse,
) -> anyhow::Result<String> {
let sources = result
.sources
.iter()
.map(|source| (source.title.clone(), source.url.clone()))
.collect::<Vec<_>>();
render(RenderRequest::WebSearch {
answer: &result.answer,
sources: &sources,
})
}
fn render_web_fetch_result(
result: &kcode_intelligence_router::FetchResponse,
) -> anyhow::Result<String> {
render(RenderRequest::WebFetch {
url: &result.url,
title: result.title.as_deref(),
content_type: &result.content_type,
truncated: result.truncated,
content: &result.content,
})
}
fn render_media_annotation_result(
object_id: &str,
file_name: &str,
content_type: &str,
result: &kcode_intelligence_router::AnnotationResponse,
) -> anyhow::Result<String> {
render(RenderRequest::MediaAnnotation {
object_id,
file_name,
content_type,
model: &result.model,
complete: result.complete,
incomplete_reason: result.incomplete_reason.as_deref(),
text: &result.text,
})
}
fn render_audio_transcription_result(
object_id: &str,
file_name: &str,
content_type: &str,
result: &kcode_intelligence_router::TranscriptionResponse,
) -> anyhow::Result<String> {
render(RenderRequest::AudioTranscription {
object_id,
file_name,
content_type,
model: &result.model,
text: &result.text,
})
}
fn render_document_extraction_result(
object_id: &str,
file_name: &str,
result: &kcode_intelligence_router::DocumentExtraction,
) -> anyhow::Result<String> {
render(RenderRequest::DocumentExtraction {
object_id,
file_name,
format: &result.format,
characters: result.characters,
truncated: result.truncated,
text: &result.text,
})
}
struct ToolCall {
name: String,
arguments: Value,
}
struct RecordedToolInvocation {
invocation_id: String,
tool_instance: String,
tool_name: String,
}
struct PendingFreeformWrite {
request: FreeformWrite,
call_box_id: BoxId,
}
fn tool_call_box_content(call: &ToolCall) -> anyhow::Result<BoxContent> {
if matches!(
call.name.as_str(),
WRITE_RUST_LIB_TOOL | WRITE_WEB_LIB_TOOL | WRITE_RUST_BIN_TOOL
) {
let name = call
.arguments
.get("name")
.and_then(Value::as_str)
.map(|name| name.chars().take(255).collect::<String>());
let file_count = call
.arguments
.get("files")
.and_then(Value::as_array)
.map(Vec::len);
return Ok(BoxContent {
text: serde_json::to_string_pretty(&json!({
"name":call.name,
"arguments":{
"name":name,
"fileCount":file_count,
"completeFileContents":"omitted from active context; retained in the durable tool invocation"
}
}))?,
objects: Vec::new(),
metadata: json!({
"compactedToolInvocation":true,
"toolName":call.name,
}),
});
}
Ok(BoxContent::text(serde_json::to_string_pretty(
&json!({"name":call.name,"arguments":call.arguments}),
)?))
}
struct ToolOutcome {
text: String,
store_result: bool,
ok: bool,
end_session: bool,
freeform_write: Option<FreeformWrite>,
managed_source_snapshot: Option<SourceSnapshot>,
}
#[derive(Clone)]
struct ChangedSubagentState {
box_id: BoxId,
name: String,
text: Option<String>,
hide_from_parent: bool,
}
#[derive(Clone, Copy)]
struct CanonicalBoxVersion {
event_id: EventId,
active: bool,
tool_owned: bool,
dehydrated: bool,
}
type CanonicalBoxVersions = BTreeMap<BoxId, CanonicalBoxVersion>;
fn subagent_state_key(box_id: BoxId) -> String {
format!("tool-state:{box_id}")
}
fn result_displays_snapshot(result: &str, snapshot: &SourceSnapshot) -> bool {
result == snapshot.text
}
fn canonical_box_versions(journal: &HistorySession) -> CanonicalBoxVersions {
journal
.state()
.boxes
.iter()
.map(|(box_id, state)| {
(
*box_id,
CanonicalBoxVersion {
event_id: state.canonical.event_id,
active: state.active,
tool_owned: matches!(state.owner, BoxOwner::Tool { .. }),
dehydrated: matches!(state.representation, Representation::Dehydrated { .. }),
},
)
})
.collect()
}
fn changed_subagent_tool_states(
journal: &HistorySession,
previous: &CanonicalBoxVersions,
) -> Vec<ChangedSubagentState> {
journal
.state()
.boxes
.iter()
.filter_map(|(box_id, state)| {
let current_tool_owned = matches!(state.owner, BoxOwner::Tool { .. });
let prior = previous.get(box_id);
if !current_tool_owned && !prior.is_some_and(|prior| prior.tool_owned) {
return None;
}
let changed = prior.is_none_or(|prior| {
prior.event_id != state.canonical.event_id
|| prior.active != state.active
|| prior.tool_owned != current_tool_owned
});
changed.then(|| ChangedSubagentState {
box_id: *box_id,
name: state.name.clone(),
text: (current_tool_owned && !state.canonical.content.text.is_empty())
.then(|| state.canonical.content.text.clone()),
hide_from_parent: prior.is_none_or(|prior| prior.dehydrated || !prior.active),
})
})
.collect()
}
fn subagent_managed_write_fits(
journal: &HistorySession,
call: &ToolCall,
budget: &kcode_agent_runtime::ContextBudget,
) -> bool {
let Some(snapshot) = proposed_write_snapshot(&call.name, &call.arguments) else {
return true;
};
let kind = snapshot.kind;
let key = source_box_id(journal, kind, &snapshot.name)
.map(subagent_state_key)
.unwrap_or_else(|| format!("prospective-managed-state:{:?}:{}", kind, snapshot.name));
budget.fits_state(
key,
format!(
"Current Managed {} {}:\n{}",
kind.label(),
snapshot.name,
snapshot.text
),
)
}
fn subagent_state_updates(
states: &[ChangedSubagentState],
) -> Vec<kcode_agent_runtime::StateUpdate> {
states
.iter()
.map(|state| kcode_agent_runtime::StateUpdate {
key: subagent_state_key(state.box_id),
text: state
.text
.as_ref()
.map(|text| format!("Current {}:\n{text}", state.name)),
})
.collect()
}
struct KennedySubagentHost<'a> {
session: &'a mut Session,
captures: HashMap<String, FreeformWrite>,
}
struct KennedySessionHost<'a, C> {
session: &'a mut Session,
checkpoint: &'a mut C,
accounting: Option<kcode_intelligence_chatend::TopLevelCall>,
pending_freeform_write: Option<PendingFreeformWrite>,
deadline_after_response: bool,
}
impl Session {
pub async fn new(
api: Service,
system_prompt: String,
runtime: RuntimeModel,
mut options: SessionOptions,
restored: Option<&Value>,
) -> anyhow::Result<Self> {
if let Some(state) = restored {
restore_session_type(&mut options, state);
options.channel = state.get("channel").cloned().unwrap_or(options.channel);
options.free_time = state.get("freeTime").cloned().unwrap_or(options.free_time);
options.orchestration = state
.get("orchestration")
.cloned()
.unwrap_or(options.orchestration);
}
if options.group_context.is_null() {
options.group_context = options
.channel
.get("groupContext")
.cloned()
.unwrap_or(Value::Null);
}
options
.reference_root_node_ids
.retain(|id| !options.root_node_ids.contains(id));
options.reference_root_node_ids.sort();
options.reference_root_node_ids.dedup();
let started_at = restored
.and_then(|state| state.get("startedAt"))
.and_then(Value::as_str)
.map(str::to_owned)
.unwrap_or_else(|| Utc::now().to_rfc3339());
let rust_lib_session_id = restored
.and_then(|state| state.get("rustLibSessionId"))
.and_then(Value::as_str)
.map(str::to_owned)
.or(options.rust_lib_session_id.clone())
.unwrap_or_else(|| format!("kennedy:{}", Uuid::new_v4()));
let history_session_id = restored
.and_then(|state| state.get("sessionId"))
.and_then(Value::as_str)
.map(str::to_owned);
let source_session_type = options.source_session_type.clone().or_else(|| {
restored
.and_then(|state| state.get("sourceSessionType"))
.and_then(Value::as_str)
.map(str::to_owned)
});
let session_id = history_session_id
.clone()
.unwrap_or_else(|| Uuid::new_v4().to_string());
let metadata = SessionMetadata {
session_id: session_id.clone(),
kind: session_kind(&options.session_type, &options.mode),
created_at: started_at.clone(),
effective_context_tokens: runtime.context_window_tokens,
channel: options.channel.clone(),
};
let mut journal = if history_session_id.is_some() {
api.history_session(metadata, &runtime.model)
.with_context(|| {
format!(
"opening authoritative session {session_id} (legacy snapshots are intentionally unsupported)"
)
})?
} else {
api.create_history_session(NewSession {
kind: metadata.kind,
created_at: metadata.created_at,
effective_context_tokens: metadata.effective_context_tokens,
channel: metadata.channel,
})?
};
let mut context =
KwebContext::new(options.root_node_ids.clone()).map_err(anyhow::Error::new)?;
restore_kweb_context(&journal, &mut context)?;
let plan = KwebPlan::restore(restored, &journal)?;
let transcript = kcode_kennedy_session_ingress::transcript_from_journal(&journal);
let (pending_turn, pending_external_event_id) = restore_pending_turn(restored, &transcript);
let needs_initialization = !journal
.state()
.boxes
.values()
.any(|state| matches!(state.owner, BoxOwner::System));
let commit_receipt = restore_commit_receipt(restored)?;
let commit_author = restored
.and_then(|state| state.get("commitAuthor"))
.and_then(Value::as_str)
.map(str::to_owned)
.unwrap_or_else(|| runtime.attribution());
if let Some(receipt) = &commit_receipt {
journal.mark_completed(receipt.session_object_id.to_string());
}
let completed =
journal.state().completed_session_object.is_some() || commit_receipt.is_some();
let mut session = Self {
api,
runtime,
journal,
plan,
session_type: options.session_type,
channel: options.channel,
free_time: options.free_time,
orchestration: options.orchestration,
provenance_id: options.provenance_id,
rust_lib_session_id,
root_node_ids: options.root_node_ids,
reference_root_node_ids: options.reference_root_node_ids,
started_at,
transcript,
pending_turn,
pending_external_event_id,
completed,
rounds_used: restored
.and_then(|state| state.get("roundsUsed"))
.and_then(Value::as_u64)
.unwrap_or_default(),
commit_receipt,
commit_author,
mode: options.mode,
source_session_type,
group_context: options.group_context,
context,
free_time_end_reason: None,
fatal_persistence_error: None,
};
if matches!(session.mode, AgentMode::Ingress { .. }) && !session.journal.is_sealed() {
session.journal.repair_unfinished_tools(now())?;
}
if session.journal.is_sealed() {
anyhow::ensure!(
!matches!(session.mode, AgentMode::Conversation),
"a read-only conversation has an unexpectedly sealed session log"
);
if session.commit_receipt.is_none() {
session.finalize_kweb_session()?;
}
session.completed = true;
return Ok(session);
}
if needs_initialization {
session.journal.create_box(
now(),
"Kennedy system prompt",
BoxOwner::System,
BoxContent::text(&system_prompt),
)?;
if session.session_type == "telegram-group" && !session.group_context.is_null() {
session.journal.create_box(
now(),
"Telegram group context",
BoxOwner::Controller,
BoxContent::text(kcode_telegram_session_coordinator::format_group_context(
&session.group_context,
)),
)?;
}
let roots = session.root_node_ids.clone();
let invocation =
session.record_tool_invocation("LoadNodes", json!({"identifiers":&roots}))?;
let result = load_durable_batch(session.api.kmap(), &mut session.context, &roots)?;
session.sync_kweb_boxes()?;
session.record_tool_completion(
Some(&invocation),
json!({"ok":true,"automatic":true,"identifiers":roots,"result":result}),
)?;
} else {
session.sync_kweb_boxes()?;
}
if matches!(session.mode, AgentMode::Ingress { .. })
&& !session.completed
&& !session.journal.state().history_ingress_started
{
session.prepare_history_ingress(&system_prompt).await?;
}
Ok(session)
}
async fn prepare_history_ingress(&mut self, prompt: &str) -> anyhow::Result<()> {
let cost_at_ingress = self.journal.state().projection().status;
if !self.journal.state().source_terminated {
self.journal.record(
now(),
EventKind::SourceTerminated {
reason: "history_ingress".into(),
},
)?;
}
let system_box = self
.journal
.state()
.boxes
.values()
.find(|state| matches!(state.owner, BoxOwner::System))
.map(|state| state.id)
.context("session has no system-prompt box")?;
self.journal
.update_box(now(), system_box, BoxContent::text(prompt))?;
let ingress_kind = session_kind(&self.session_type, &self.mode);
if self.journal.state().metadata.effective_context_tokens
!= self.runtime.context_window_tokens
|| self.journal.state().metadata.kind != ingress_kind
{
self.journal
.configure_context(ingress_kind, self.runtime.context_window_tokens);
}
self.journal.create_box(
now(),
"Session cost at ingress",
BoxOwner::Controller,
BoxContent::text(cost_summary(
"session cost before history ingress",
cost_at_ingress.estimated_cost_usd_nanos,
cost_at_ingress.unpriced_provider_calls,
)),
)?;
self.revalidate_loaded_nodes().await?;
match kcode_history_ingress_context::prepare(&mut self.journal, now())? {
HistoryIngressContextOutcome::Ready => {}
HistoryIngressContextOutcome::OverCapacity {
estimated_tokens,
target_tokens,
} => {
self.journal.record(
now(),
EventKind::Note {
label: INGRESS_FORCE_COMMIT_NOTE.into(),
value: json!({
"reason":"fully_dehydrated_context_above_initial_target",
"estimatedTokens":estimated_tokens,
"initialTargetTokens":target_tokens,
}),
},
)?;
self.pending_turn = false;
self.finalize_kweb_session()?;
self.completed = true;
return Ok(());
}
}
self.journal
.record(now(), EventKind::HistoryIngressStarted)?;
self.pending_turn = true;
Ok(())
}
async fn revalidate_loaded_nodes(&mut self) -> anyhow::Result<()> {
let direct = self.context.loaded_node_ids().to_vec();
load_durable_batch(self.api.kmap(), &mut self.context, &direct)?;
self.sync_kweb_boxes()?;
Ok(())
}
fn stage_user_input(&mut self, text: &str, metadata: &Value) -> Option<InputStage> {
let recorded_at = now();
let result = (|| -> anyhow::Result<Option<InputStage>> {
let Some(staged) = kcode_kennedy_session_ingress::stage_user_input(
&mut self.journal,
text,
metadata,
&recorded_at,
)?
else {
return Ok(None);
};
self.transcript.push(staged.transcript);
self.recover_context_overflow(staged.external_event_id.as_deref(), &[])?;
Ok(Some(InputStage::Accepted))
})();
match result {
Ok(stage) => stage,
Err(error) => {
self.fatal_persistence_error = Some(error.to_string());
tracing::error!(error=%error, "Could not durably stage session input");
Some(InputStage::Accepted)
}
}
}
pub fn append_final_user_message(&mut self, text: &str, metadata: &Value) -> bool {
self.stage_user_input(text, metadata).is_some()
}
pub fn stage_source_message(
&mut self,
kennedy: bool,
text: &str,
metadata: Value,
) -> anyhow::Result<()> {
let staged = kcode_kennedy_session_ingress::stage_source_input(
&mut self.journal,
kennedy,
text,
metadata,
&now(),
)?;
self.transcript.push(staged.transcript);
self.recover_context_overflow(staged.external_event_id.as_deref(), &[])?;
Ok(())
}
pub fn answer_for_external_event(&self, id: &str) -> Option<&Value> {
self.transcript.iter().rev().find(|entry| {
is_terminal_external_response(entry)
&& entry.get("externalEventId").and_then(Value::as_str) == Some(id)
})
}
pub fn responses_for_external_event(&self, id: &str) -> Vec<&Value> {
self.transcript
.iter()
.filter(|entry| {
matches!(
entry.get("role").and_then(Value::as_str),
Some("kennedy" | "system")
) && entry.get("externalEventId").and_then(Value::as_str) == Some(id)
})
.collect()
}
pub fn resolve_object(&mut self, object_id: &str) -> anyhow::Result<ResolvedObject> {
let api = self.api.clone();
kcode_kennedy_session_objects::resolve_object(
&mut self.journal,
object_id,
move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
)
}
fn resolve_media_object(&mut self, object_id: &str) -> anyhow::Result<ResolvedObject> {
let api = self.api.clone();
kcode_kennedy_session_objects::resolve_media_object(
&mut self.journal,
object_id,
MAX_MEDIA_ENRICHMENT_BYTES,
move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
)
}
fn resolve_image_object(
&mut self,
object_id: &str,
) -> anyhow::Result<(Vec<u8>, String, String)> {
let resolved = self.resolve_media_object(object_id)?;
anyhow::ensure!(
resolved.media_type.starts_with("image/"),
"GenerateImage reference {object_id} is not an image"
);
Ok((resolved.bytes, resolved.file_name, resolved.media_type))
}
fn recover_context_overflow(
&mut self,
external_event_id: Option<&str>,
pinned_box_ids: &[BoxId],
) -> anyhow::Result<ContextRecovery> {
let projection = self.journal.state().projection();
let target_tokens = self.journal.state().active_context_limit();
if projection.estimated_tokens <= target_tokens {
return Ok(ContextRecovery::NotNeeded);
}
let projection_hash = hex::encode(Sha256::digest(projection.render().as_bytes()));
let already_irreducible = self
.journal
.state()
.events
.iter()
.rev()
.find_map(|event| match &event.kind {
EventKind::Note { label, value } if label == "context_overflow_recovery" => {
Some(value)
}
_ => None,
})
.is_some_and(|value| {
value.get("irreducible").and_then(Value::as_bool) == Some(true)
&& value.get("limitTokens").and_then(Value::as_u64) == Some(target_tokens)
&& value.get("projectionHash").and_then(Value::as_str)
== Some(projection_hash.as_str())
});
if already_irreducible {
return Ok(ContextRecovery::Irreducible);
}
let before_tokens = projection.estimated_tokens;
let mut metadata = json!({
"transcriptRole":"system",
"contextOverflowWarning":true,
"projectedTokens":before_tokens,
"limitTokens":target_tokens,
});
if let Some(id) = external_event_id {
metadata["externalEventId"] = json!(id);
}
let warning_box_id = self.journal.create_box(
now(),
CONTEXT_OVERFLOW_WARNING_BOX_NAME,
BoxOwner::Controller,
BoxContent {
text: CONTEXT_OVERFLOW_WARNING.into(),
objects: Vec::new(),
metadata,
},
)?;
let mut transcript = json!({
"role":"system",
"content":CONTEXT_OVERFLOW_WARNING,
"contextOverflowWarning":true,
});
if let Some(id) = external_event_id {
transcript["externalEventId"] = json!(id);
}
self.transcript.push(transcript);
let mut pins = pinned_box_ids.to_vec();
if !pins.contains(&warning_box_id) {
pins.push(warning_box_id);
}
let outcome = kcode_history_ingress_context::recover(&mut self.journal, now(), &pins)?;
let (dehydrated_box_ids, estimated_tokens, target_tokens, irreducible) = match outcome {
ContextRecoveryOutcome::Recovered {
dehydrated_box_ids,
estimated_tokens,
target_tokens,
} => (dehydrated_box_ids, estimated_tokens, target_tokens, false),
ContextRecoveryOutcome::OverCapacity {
dehydrated_box_ids,
estimated_tokens,
target_tokens,
} => (dehydrated_box_ids, estimated_tokens, target_tokens, true),
};
let final_projection_hash = hex::encode(Sha256::digest(
self.journal.state().projection().render().as_bytes(),
));
self.journal.record(
now(),
EventKind::Note {
label: "context_overflow_recovery".into(),
value: json!({
"beforeTokens":before_tokens,
"estimatedTokens":estimated_tokens,
"limitTokens":target_tokens,
"dehydratedBoxIds":dehydrated_box_ids,
"irreducible":irreducible,
"projectionHash":final_projection_hash,
}),
},
)?;
if irreducible {
if matches!(self.mode, AgentMode::Ingress { .. }) {
self.request_ingress_force_commit(
"irreducible_context_overflow",
estimated_tokens,
)?;
} else if !self.journal.state().source_terminated {
self.journal.record(
now(),
EventKind::SourceTerminated {
reason: "irreducible_context_overflow".into(),
},
)?;
}
Ok(ContextRecovery::Irreducible)
} else {
Ok(ContextRecovery::Recovered)
}
}
fn request_ingress_force_commit(
&mut self,
reason: &str,
projected_tokens: u64,
) -> anyhow::Result<()> {
if self.ingress_force_commit_requested() {
return Ok(());
}
self.journal.record(
now(),
EventKind::Note {
label: INGRESS_FORCE_COMMIT_NOTE.into(),
value: json!({
"reason":reason,
"projectedTokens":projected_tokens,
"limitTokens":self.journal.state().ingress_context_limit(),
}),
},
)?;
Ok(())
}
fn ingress_force_commit_requested(&self) -> bool {
self.journal.state().events.iter().rev().any(|event| {
matches!(
&event.kind,
EventKind::Note { label, .. } if label == INGRESS_FORCE_COMMIT_NOTE
)
})
}
pub fn requires_history_ingress(&self) -> bool {
matches!(self.mode, AgentMode::Conversation) && self.journal.state().source_terminated
}
pub fn stage_free_time_opening(&mut self) -> bool {
if self.pending_turn {
return false;
}
let mut blocks = vec![
render(RenderRequest::FreeTimeOpening {
free_time: &self.free_time,
})
.expect("free-time opening rendering is infallible"),
];
if let Some(message) = self
.free_time
.get("handoffMessage")
.and_then(Value::as_str)
.filter(|message| !message.trim().is_empty())
{
blocks.push(format!(
"Message from the previous self-time session:\n\n{message}"
));
}
let Some(stage) = self.stage_user_input(&blocks.join("\n\n"), &json!({"kind":"self-time"}))
else {
return false;
};
self.pending_turn = matches!(stage, InputStage::Accepted);
true
}
pub fn stage_wakeup_opening(&mut self) -> anyhow::Result<bool> {
if self.pending_turn {
return Ok(false);
}
let marker = self
.channel
.get("wakeupMarker")
.and_then(Value::as_str)
.context("wakeup session is missing its acquired time marker")?;
let marker = DateTime::parse_from_rfc3339(marker)
.context("wakeup session has an invalid acquired time marker")?
.with_timezone(&Utc);
let text = render(RenderRequest::WakeupOpening { marker })?;
let Some(stage) = self.stage_user_input(
&text,
&json!({"kind":"wakeup","wakeupMarker":marker.to_rfc3339()}),
) else {
return Ok(false);
};
self.pending_turn = matches!(stage, InputStage::Accepted);
Ok(true)
}
pub fn begin_user_turn(&mut self, text: &str, metadata: &Value) -> bool {
if self.pending_turn {
return false;
}
let Some(stage) = self.stage_user_input(text, metadata) else {
return false;
};
debug_assert_eq!(stage, InputStage::Accepted);
self.rounds_used = 0;
self.pending_turn = true;
self.pending_external_event_id = metadata
.get("externalEventId")
.and_then(Value::as_str)
.map(str::to_owned);
true
}
pub fn reset_exhausted_turn_rounds_for_retry(&mut self) {
if matches!(self.mode, AgentMode::Conversation)
&& self.rounds_used >= AGENT_LOOP_ROUND_LIMIT
{
self.rounds_used = 0;
}
}
pub fn interrupt_current_turn(&mut self) -> anyhow::Result<()> {
self.journal.repair_unfinished_tools(now())?;
let notice = "The user stopped this agent turn.";
let mut metadata = json!({"transcriptRole":"system","userStopped":true});
let mut transcript_entry = json!({
"role":"system",
"content":notice,
"userStopped":true,
});
if let Some(external_event_id) = &self.pending_external_event_id {
metadata["externalEventId"] = json!(external_event_id);
transcript_entry["externalEventId"] = json!(external_event_id);
}
self.journal.create_box(
now(),
"Turn stopped",
BoxOwner::Controller,
BoxContent {
text: notice.into(),
objects: Vec::new(),
metadata,
},
)?;
self.transcript.push(transcript_entry);
self.pending_turn = false;
self.pending_external_event_id = None;
self.orchestration =
json!({"owner":"backend","status":"idle","lastOutcome":"user-stopped"});
Ok(())
}
pub async fn run_pending_turn<C, F>(
&mut self,
operation_id: Uuid,
mut checkpoint: C,
) -> anyhow::Result<Option<String>>
where
C: FnMut(Value) -> F + Send,
F: Future<Output = anyhow::Result<()>> + Send,
{
if let Some(error) = self.fatal_persistence_error.take() {
anyhow::bail!("session journal write failed: {error}");
}
if !self.pending_turn {
return Ok(None);
}
let runtime = self.api.agent_runtime();
let user_id = self
.root_node_ids
.first()
.context("session has no user root for intelligence accounting")?
.clone();
let request = kcode_agent_runtime::SessionRunRequest {
user_id,
operation_id,
rounds_used: self.rounds_used,
round_limit: AGENT_LOOP_ROUND_LIMIT,
};
let mut host = KennedySessionHost {
session: self,
checkpoint: &mut checkpoint,
accounting: None,
pending_freeform_write: None,
deadline_after_response: false,
};
let result = runtime.run_session(request, &mut host).await;
let round_limit = result
.as_ref()
.is_err_and(kcode_agent_runtime::is_session_round_limit);
let result = if round_limit && matches!(host.session.mode, AgentMode::Ingress { .. }) {
host.session.request_ingress_force_commit(
"agent_loop_round_limit",
host.session.journal.state().projection().estimated_tokens,
)?;
(host.checkpoint)(host.session.snapshot()?).await?;
None
} else {
result?
};
drop(host);
match self.mode {
AgentMode::Conversation => {
if self.journal.state().source_terminated {
self.pending_turn = false;
self.pending_external_event_id = None;
checkpoint(self.snapshot()?).await?;
return Ok(None);
}
let Some(answer) = result else {
if self
.pending_external_event_id
.as_deref()
.and_then(|id| self.answer_for_external_event(id))
.is_some()
{
self.pending_turn = false;
self.pending_external_event_id = None;
checkpoint(self.snapshot()?).await?;
return Ok(None);
}
anyhow::bail!(
"Kennedy ended a conversational turn without an assistant response"
);
};
self.pending_turn = false;
self.pending_external_event_id = None;
checkpoint(self.snapshot()?).await?;
Ok(Some(answer))
}
AgentMode::FreeTime | AgentMode::Wakeup | AgentMode::Ingress { .. } => {
self.pending_turn = false;
self.pending_external_event_id = None;
self.finalize_kweb_session()?;
self.completed = true;
checkpoint(self.snapshot()?).await?;
Ok(None)
}
}
}
fn project_descendant<T>(
&mut self,
outcome: Result<kcode_intelligence_router::Accounted<T>, services::ApiError>,
) -> anyhow::Result<T> {
match outcome {
Ok(accounted) => {
kcode_intelligence_chatend::record_descendant_receipt(
&mut self.journal,
&accounted.receipt,
)?;
Ok(accounted.value)
}
Err(error) => {
if let Some(receipt) = &error.receipt {
kcode_intelligence_chatend::record_descendant_receipt(
&mut self.journal,
receipt,
)?;
}
Err(error.into())
}
}
}
async fn run_subagent(
&mut self,
model: String,
reasoning_effort: Option<String>,
context_node_ids: Vec<String>,
task: String,
parent_operation_id: Uuid,
) -> anyhow::Result<String> {
let reasoning_effort =
reasoning_effort.unwrap_or_else(|| self.runtime.reasoning_effort.clone());
let mut context = Vec::with_capacity(context_node_ids.len());
for node_id in &context_node_ids {
context.push(self.api.kmap_node(node_id)?.data.long_description);
}
let user_id = self
.root_node_ids
.first()
.context("session has no user root for subagent intelligence accounting")?
.clone();
let timeout = self.agent_request_timeout();
let runtime = self.api.agent_runtime();
let first_event = self.journal.state().events.len();
let cost_before = self.journal.state().projection().status;
let result = {
let mut host = KennedySubagentHost {
session: self,
captures: HashMap::new(),
};
runtime
.run(
kcode_agent_runtime::RunRequest {
user_id,
parent_operation_id,
model,
reasoning_effort,
context,
task,
timeout,
start_metadata: json!({"contextNodeIds":context_node_ids}),
},
&mut host,
)
.await
};
match result {
Ok(result) => {
let cost_after = self.journal.state().projection().status;
Ok(format!(
"{}\n\n[{}]",
result.answer,
cost_summary(
"subagent cost",
cost_after
.estimated_cost_usd_nanos
.saturating_sub(cost_before.estimated_cost_usd_nanos),
cost_after
.unpriced_provider_calls
.saturating_sub(cost_before.unpriced_provider_calls),
)
))
}
Err(error) => {
let may_have_effects =
self.journal.state().events[first_event..]
.iter()
.any(|event| {
matches!(
&event.kind,
EventKind::Note { label, .. } if label == "subagent_tool_call"
)
});
if may_have_effects {
Err(error.context(
"the subagent failed after making Ktool calls; some tool effects may already have occurred",
))
} else {
Err(error)
}
}
}
}
async fn complete_subagent_freeform_write(
&mut self,
request: FreeformWrite,
contents: String,
budget: &kcode_agent_runtime::ContextBudget,
) -> anyhow::Result<(ToolOutcome, Vec<ChangedSubagentState>)> {
let kind = request.kind();
let freeform_tool = request.write_tool();
let backend_arguments = request.capture_subagent(&mut self.journal, &now(), contents)?;
let preview = self
.api
.managed_source_execute(
&self.rust_lib_session_id,
request.preview_tool(),
backend_arguments.clone(),
Vec::new(),
)
.await?;
let preview = preview
.snapshot
.context("subagent freeform write preview omitted its source snapshot")?;
let source_box_id = request.source_box_id(&self.journal)?;
anyhow::ensure!(
budget.fits_state(
format!("tool-state:{source_box_id}"),
format!(
"Current Managed {} {}:\n{}",
kind.label(),
preview.name,
preview.text
),
),
"{freeform_tool} was not run because its resulting source state would exceed the subagent context limit"
);
let previous = canonical_box_versions(&self.journal);
let execution = self
.api
.managed_source_execute(
&self.rust_lib_session_id,
freeform_tool,
backend_arguments,
Vec::new(),
)
.await?;
let snapshot = execution
.snapshot
.context("subagent freeform write omitted its resulting source snapshot")?;
apply_snapshot(&mut self.journal, &now(), snapshot)?;
let states = changed_subagent_tool_states(&self.journal, &previous);
Ok((
ToolOutcome {
text: execution.text,
store_result: false,
ok: true,
end_session: false,
freeform_write: None,
managed_source_snapshot: None,
},
states,
))
}
async fn complete_freeform_write(
&mut self,
pending: PendingFreeformWrite,
contents: String,
) -> anyhow::Result<ToolOutcome> {
let request = pending.request;
let freeform_tool = request.write_tool();
let backend_arguments =
request.capture(&mut self.journal, &now(), pending.call_box_id, contents)?;
let preview_result = self
.api
.managed_source_execute(
&self.rust_lib_session_id,
request.preview_tool(),
backend_arguments.clone(),
Vec::new(),
)
.await;
let preview = match preview_result {
Ok(preview) => preview,
Err(error) => {
return Ok(ToolOutcome {
text: format!("{freeform_tool} failed: {error}"),
store_result: true,
ok: false,
end_session: false,
freeform_write: None,
managed_source_snapshot: None,
});
}
};
let _preview = preview
.snapshot
.context("freeform write preview omitted the resulting source snapshot")?;
request.source_box_id(&self.journal)?;
let execution_result = self
.api
.managed_source_execute(
&self.rust_lib_session_id,
freeform_tool,
backend_arguments,
Vec::new(),
)
.await;
let execution = match execution_result {
Ok(execution) => execution,
Err(error) => {
return Ok(ToolOutcome {
text: format!("{freeform_tool} failed: {error}"),
store_result: true,
ok: false,
end_session: false,
freeform_write: None,
managed_source_snapshot: None,
});
}
};
let snapshot = execution
.snapshot
.context("freeform write omitted the resulting source snapshot")?;
apply_snapshot(&mut self.journal, &now(), snapshot)?;
Ok(ToolOutcome {
text: execution.text,
store_result: false,
ok: true,
end_session: false,
freeform_write: None,
managed_source_snapshot: None,
})
}
async fn send_telegram_dm(&mut self, arguments: &Value) -> anyhow::Result<String> {
let request = kcode_telegram_session_coordinator::parse_private_request(arguments)?;
let attachments = self.telegram_delivery_attachments(request.attachments)?;
let caller_holds_user_lock = self.session_type == "telegram"
&& self.channel.get("telegramUserId").and_then(Value::as_i64)
== Some(request.telegram_user_id);
self.api
.telegram()
.send_private(kcode_telegram_session_coordinator::PrivateDelivery {
telegram_user_id: request.telegram_user_id,
message: request.message,
attachments,
caller_holds_user_lock,
})
.await
}
async fn send_telegram_group_message(&mut self, arguments: &Value) -> anyhow::Result<String> {
let request = kcode_telegram_session_coordinator::parse_group_request(arguments)?;
let attachments = self.telegram_delivery_attachments(request.attachments)?;
self.api
.telegram()
.send_group(kcode_telegram_session_coordinator::GroupDelivery {
root_node_id: request.root_node_id,
message: request.message,
attachments,
})
.await
}
fn telegram_delivery_attachments(
&mut self,
requests: Vec<kcode_telegram_session_coordinator::AttachmentRequest>,
) -> anyhow::Result<Vec<kcode_telegram_session_coordinator::Attachment>> {
let api = self.api.clone();
kcode_kennedy_session_objects::delivery_attachments(
&mut self.journal,
requests,
move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
)
}
async fn execute_tool(
&mut self,
call: &ToolCall,
operation_id: Uuid,
) -> anyhow::Result<ToolOutcome> {
self.assert_tool_allowed(&call.name)?;
let decoded = decode(&call.name, &call.arguments)?;
let mut end_session = false;
let mut store_result = true;
let mut freeform_write = None;
let mut managed_source_snapshot = None;
let text = match (call.name.as_str(), decoded) {
("SendTelegramDM", _) => self.send_telegram_dm(&call.arguments).await?,
("SendTelegramGroupMessage", _) => {
self.send_telegram_group_message(&call.arguments).await?
}
(
"RunSubagent",
Some(DecodedTool::RunSubagent {
model,
reasoning_effort,
context_node_ids,
task,
}),
) => {
let first_event = self.journal.state().events.len();
match self
.run_subagent(
model,
reasoning_effort,
context_node_ids,
task,
operation_id,
)
.await
{
Ok(response) => response,
Err(error) => {
let may_have_effects = self.journal.state().events[first_event..]
.iter()
.any(|event| {
matches!(
&event.kind,
EventKind::Note { label, .. }
if label == "subagent_tool_call"
)
});
if may_have_effects {
return Err(error.context(
"the subagent failed after making Ktool calls; some tool effects may already have occurred",
));
}
return Err(error);
}
}
}
("EndSession", Some(DecodedTool::EndSession { message })) => {
anyhow::ensure!(
!matches!(self.mode, AgentMode::Conversation),
"EndSession is only available during an autonomous or history-ingress session"
);
end_session = true;
if matches!(self.mode, AgentMode::FreeTime)
&& let Some(message) = message.filter(|message| !message.trim().is_empty())
{
self.free_time["nextSessionMessage"] = json!(message);
}
"Session ending.".into()
}
("DehydrateBoxes", Some(DecodedTool::BoxIds(ids))) => {
self.journal.dehydrate_boxes(now(), &ids)?;
format!(
"Dehydrated boxes {}.",
ids.iter()
.map(ToString::to_string)
.collect::<Vec<_>>()
.join(", ")
)
}
("SummarizeBox", Some(DecodedTool::SummarizeBox { box_id, summary })) => {
self.journal.summarize_box(now(), box_id, summary)?;
format!("Summarized box {box_id}.")
}
("HydrateBox", Some(DecodedTool::BoxId(id))) => {
self.journal.rehydrate_box(now(), id)?;
let external_event_id = self.pending_external_event_id.clone();
match self.recover_context_overflow(external_event_id.as_deref(), &[id])? {
ContextRecovery::NotNeeded => format!("Hydrated box {id}."),
ContextRecovery::Recovered => {
format!("Hydrated box {id}.\n\n{CONTEXT_OVERFLOW_WARNING}")
}
ContextRecovery::Irreducible => anyhow::bail!(CONTEXT_OVERFLOW_WARNING),
}
}
("BoxesIntoObjects", Some(DecodedTool::BoxIds(ids))) => {
kcode_kennedy_box_text_objects::stage_box_text_objects(
&mut self.journal,
&ids,
&now(),
)?
}
("LoadNodes", Some(DecodedTool::LoadNodes(identifiers))) => {
load_durable_batch(self.api.kmap(), &mut self.context, &identifiers)?;
let changed = self.sync_kweb_boxes()?;
store_result = false;
render_load_nodes_result(&self.journal, &changed)?
}
(
"EmitObject",
Some(DecodedTool::EmitObject {
object_id,
file_name,
}),
) => {
anyhow::ensure!(
matches!(self.mode, AgentMode::Conversation),
"EmitObject is only available in a conversation"
);
let object = self.resolve_object(&object_id)?;
let file_name = file_name.unwrap_or_else(|| object.file_name.clone());
if let Some(maximum) = self.channel.get("maxObjectBytes").and_then(Value::as_u64) {
anyhow::ensure!(
!object.bytes.is_empty(),
"object {object_id} is empty and cannot be sent through this channel"
);
anyhow::ensure!(
object.bytes.len() as u64 <= maximum,
"object {object_id} is {} bytes, over this channel's {maximum}-byte limit",
object.bytes.len()
);
}
let descriptor = json!({
"objectId":object_id,
"fileName":file_name,
"mediaType":object.media_type,
"byteLength":object.bytes.len(),
});
let mut metadata = json!({
"outputKind":"object",
"attachments":[descriptor.clone()],
});
if let Some(external_event_id) = &self.pending_external_event_id {
metadata["externalEventId"] = json!(external_event_id);
}
let content = BoxContent {
text: String::new(),
objects: vec![object_id.clone()],
metadata,
};
self.journal
.create_box(now(), "Kennedy message", BoxOwner::Kennedy, content)?;
let mut transcript = json!({
"role":"kennedy",
"content":"",
"objects":[object_id],
"attachments":[descriptor],
});
if let Some(external_event_id) = &self.pending_external_event_id {
transcript["externalEventId"] = json!(external_event_id);
}
self.transcript.push(transcript);
store_result = false;
"Object emitted to the user.".into()
}
("WebSearch", Some(DecodedTool::WebSearch { question, model })) => {
let user_id = self
.root_node_ids
.first()
.context("session has no user root for intelligence accounting")?
.clone();
let outcome = self
.api
.search(
&user_id,
kcode_intelligence_router::SearchRequest {
question,
model,
operation_id: Uuid::new_v4(),
parent_operation_id: Some(operation_id),
},
)
.await;
let result = self.project_descendant(outcome)?;
render_web_search_result(&result)?
}
("WebFetch", Some(DecodedTool::WebFetch(url))) => {
let user_id = self
.root_node_ids
.first()
.context("session has no user root for intelligence accounting")?;
let result = self
.api
.fetch(
user_id,
kcode_intelligence_router::FetchRequest {
url,
operation_id: Uuid::new_v4(),
parent_operation_id: Some(operation_id),
},
)
.await?;
render_web_fetch_result(&result)?
}
("StageTelegramGroupMedia", Some(DecodedTool::StageTelegramGroupMedia(message_id))) => {
let media_ref = kcode_telegram_session_coordinator::group_media_reference(
&self.group_context,
message_id,
)?;
let chat_id = media_ref.chat_id;
let api = self.api.clone();
let staged = kcode_kennedy_session_objects::stage_telegram_group_media(
&mut self.journal,
kcode_kennedy_session_objects::TelegramStageRequest {
chat_id,
message_id,
maximum_bytes: MAX_MEDIA_ENRICHMENT_BYTES,
transport_metadata: media_ref.transport_metadata(),
recorded_at: now(),
},
|| api.telegram().group_message_media(chat_id, message_id),
|media_type| {
kcode_telegram_session_coordinator::group_media_file_name(
&media_ref, media_type,
)
},
)?;
render(RenderRequest::StagedTelegramMedia {
pending_id: &staged.descriptor.pending_id,
kind: &staged.kind,
file_name: &staged.descriptor.file_name,
media_type: &staged.descriptor.media_type,
size_bytes: staged.descriptor.size_bytes,
message_id,
reused: staged.reused,
})?
}
(
"TranscribeAudio",
Some(DecodedTool::MediaEnrichment {
object_id,
model,
prompt,
}),
) => {
let object = self.resolve_media_object(&object_id)?;
validate(ValidationRequest::TranscribableAudio(&object.media_type))?;
validate(ValidationRequest::TranscriptionModel(&model))?;
let user_id = self
.root_node_ids
.first()
.context("session has no user root for intelligence accounting")?
.clone();
let outcome = self
.api
.transcribe_audio(
&user_id,
&model,
&prompt,
object.bytes,
object.file_name.clone(),
&object.media_type,
None,
operation_id,
)
.await;
let result = self.project_descendant(outcome)?;
render_audio_transcription_result(
&object.object_id,
&object.file_name,
&object.media_type,
&result,
)?
}
(
"AnnotateMedia",
Some(DecodedTool::MediaEnrichment {
object_id,
model,
prompt,
}),
) => {
let media = self.resolve_media_object(&object_id)?;
validate(ValidationRequest::Annotation {
model: &model,
media_type: &media.media_type,
})?;
let user_id = self
.root_node_ids
.first()
.context("session has no user root for intelligence accounting")?
.clone();
let outcome = self
.api
.annotate_media(
&user_id,
&model,
&prompt,
media.bytes,
media.file_name.clone(),
&media.media_type,
operation_id,
)
.await;
let result = self.project_descendant(outcome)?;
render_media_annotation_result(
&media.object_id,
&media.file_name,
&media.media_type,
&result,
)?
}
(
"GenerateImage",
Some(DecodedTool::GenerateImage {
model,
prompt,
reference_object_ids,
}),
) => {
let mut references = Vec::with_capacity(reference_object_ids.len());
for object_id in &reference_object_ids {
references.push(self.resolve_image_object(object_id)?);
}
let user_id = self
.root_node_ids
.first()
.context("session has no user root for intelligence accounting")?
.clone();
let outcome = self
.api
.generate_image(&user_id, &model, &prompt, references, operation_id)
.await;
let result = self.project_descendant(outcome)?;
let size = result.bytes.len();
let file_name =
format!("generated-image.{}", image_extension(&result.content_type));
let object_id = self.api.save_generated_image(
result.bytes,
&file_name,
&result.content_type,
&result.model,
)?;
format!(
"Generated image.\nObject: {object_id}\nFile: {file_name}\nContent type: {}\nSize: {size} bytes\nModel: {}\nUse EmitObject with {object_id} to deliver it.",
result.content_type, result.model
)
}
("ExtractDocumentText", Some(DecodedTool::ObjectId(object_id))) => {
let object = self.resolve_media_object(&object_id)?;
validate(ValidationRequest::ExtractableDocument {
media_type: &object.media_type,
file_name: &object.file_name,
})?;
let result = self
.api
.extract_document(object.bytes, object.file_name.clone(), &object.media_type)
.await?;
render_document_extraction_result(&object.object_id, &object.file_name, &result)?
}
(name, None) if SPEECH_CLASSIFICATION_TOOLS.contains(&name) => {
self.api
.execute_speech_classification_tool(name, call.arguments.clone())
.await?
}
("ConnectNodes", Some(DecodedTool::ConnectNodes(ids))) => self.connect_nodes(ids)?,
(
"ConsolidateFanout",
Some(DecodedTool::ConsolidateFanout {
parent,
fanout,
aggregator,
}),
) => self.consolidate_fanout(parent, fanout, aggregator)?,
(
"SetFixedConnection",
Some(DecodedTool::SetFixedConnection {
parent,
child,
slot,
}),
) => self.set_fixed_connection(parent, child, slot)?,
(
"CreateNode",
Some(DecodedTool::CreateNode {
parents,
owner,
short_name,
short_description,
long_description,
}),
) => self.create_node(
parents,
owner,
short_name,
short_description,
long_description,
)?,
(
"UpdateNode",
Some(DecodedTool::UpdateNode {
id,
owner,
short_name,
short_description,
long_description,
}),
) => self.update_node(id, owner, short_name, short_description, long_description)?,
(name, None)
if RUST_LIB_TOOLS.contains(&name)
|| WEB_LIB_TOOLS.contains(&name)
|| RUST_BIN_TOOLS.contains(&name) =>
{
if let Some(request) = prepare_freeform_write(&self.journal, name, &call.arguments)?
{
store_result = false;
let acknowledgement = request.acknowledgement();
freeform_write = Some(request);
acknowledgement
} else {
let object_ids = if name == CALL_RUST_BIN_TOOL {
decode_managed_objects(ManagedObjectArguments::RustBinary(&call.arguments))?
} else if name == ATTACH_OBJECT_WEB_LIB_TOOL {
decode_managed_objects(ManagedObjectArguments::WebLibraryAttachment(
&call.arguments,
))?
} else {
Vec::new()
};
let mut objects = Vec::with_capacity(object_ids.len());
for object_id in object_ids {
objects.push(self.resolve_object(&object_id)?.bytes);
}
let execution = self
.api
.managed_source_execute(
&self.rust_lib_session_id,
name,
call.arguments.clone(),
objects,
)
.await?;
if let Some(snapshot) = execution.snapshot {
managed_source_snapshot = Some(snapshot);
store_result = false;
}
execution.text
}
}
(name, Some(_)) => {
anyhow::bail!("decoded contract for {name} did not match its dispatch lane")
}
(name, None) => anyhow::bail!("Tool {name} is not available"),
};
Ok(ToolOutcome {
text,
store_result,
ok: true,
end_session,
freeform_write,
managed_source_snapshot,
})
}
fn assert_tool_allowed(&self, name: &str) -> anyhow::Result<()> {
let write = matches!(
name,
"ConnectNodes"
| "ConsolidateFanout"
| "SetFixedConnection"
| "CreateNode"
| "UpdateNode"
);
anyhow::ensure!(
!write || !matches!(self.mode, AgentMode::Conversation),
"{name} requires the global Kweb write lane and is unavailable in a read-only conversation"
);
if name == "EndSession" {
anyhow::ensure!(
!matches!(self.mode, AgentMode::Conversation),
"EndSession is unavailable in a conversation"
);
}
Ok(())
}
fn sync_kweb_boxes(&mut self) -> anyhow::Result<Vec<BoxId>> {
let updates = self
.plan
.updates
.iter()
.map(|(id, node)| (id.clone(), kweb_node_draft(node)))
.collect::<BTreeMap<_, _>>();
let creates = self
.plan
.creates
.iter()
.map(|create| KwebStagedCreate {
pending_id: create.pending_id.clone(),
data: kweb_node_draft(&create.data),
})
.collect::<Vec<_>>();
self.context
.sync_chatend(&mut self.journal, now(), &updates, &creates)
.map_err(anyhow::Error::new)
}
fn record_tool_invocation(
&mut self,
name: &str,
arguments: Value,
) -> anyhow::Result<RecordedToolInvocation> {
let invocation = RecordedToolInvocation {
invocation_id: Uuid::new_v4().to_string(),
tool_instance: tool_instance(name),
tool_name: name.into(),
};
self.journal.record(
now(),
EventKind::ToolInvoked {
tool_instance: invocation.tool_instance.clone(),
tool_name: invocation.tool_name.clone(),
arguments,
invocation_id: Some(invocation.invocation_id.clone()),
},
)?;
Ok(invocation)
}
fn record_tool_completion(
&mut self,
invocation: Option<&RecordedToolInvocation>,
outcome: Value,
) -> anyhow::Result<EventId> {
let (tool_instance, tool_name, invocation_id) = invocation
.map(|invocation| {
(
invocation.tool_instance.clone(),
invocation.tool_name.clone(),
Some(invocation.invocation_id.clone()),
)
})
.unwrap_or_else(|| ("call_ktool".into(), "call_ktool".into(), None));
self.journal.record(
now(),
EventKind::ToolCompleted {
tool_instance,
tool_name,
outcome,
invocation_id,
},
)
}
fn stage_plan(&mut self) -> anyhow::Result<()> {
self.sync_kweb_boxes()?;
Ok(())
}
fn node_data(&self, id: &str) -> anyhow::Result<PlannedNode> {
if let Some(data) = self.plan.created(id) {
return Ok(data.clone());
}
if let Some(data) = self.plan.updates.get(id) {
return Ok(data.clone());
}
let node = self
.context
.node(id)
.with_context(|| format!("Kweb context does not contain node {id}"))?;
Ok(planned_node(node))
}
fn put_node_data(&mut self, id: &str, data: PlannedNode) -> anyhow::Result<()> {
if let Some(created) = self.plan.created_mut(id) {
*created = data;
} else {
canonical_id(id)?;
self.plan.updates.insert(id.to_owned(), data);
}
Ok(())
}
fn connect_nodes(&mut self, ids: Vec<String>) -> anyhow::Result<String> {
for id in &ids {
self.ensure_known_node(id)?;
}
for id in &ids {
let mut data = self.node_data(id)?;
let mut recent = ids
.iter()
.filter(|other| *other != id)
.cloned()
.collect::<Vec<_>>();
for other in data.recent_connections {
if &other != id && !recent.contains(&other) {
recent.push(other);
}
}
data.recent_connections = recent;
self.put_node_data(id, data)?;
}
self.stage_plan()?;
Ok(format!(
"Staged connections among nodes {}.",
ids.join(", ")
))
}
fn consolidate_fanout(
&mut self,
parent: String,
fanout: Vec<String>,
aggregator: String,
) -> anyhow::Result<String> {
for id in std::iter::once(&parent)
.chain(std::iter::once(&aggregator))
.chain(fanout.iter())
{
self.ensure_known_node(id)?;
}
let mut parent_data = self.node_data(&parent)?;
parent_data
.recent_connections
.retain(|id| !fanout.contains(id));
if !parent_data.recent_connections.contains(&aggregator) {
parent_data.recent_connections.push(aggregator.clone());
}
let mut aggregator_data = self.node_data(&aggregator)?;
for id in fanout {
if !aggregator_data.recent_connections.contains(&id) {
aggregator_data.recent_connections.push(id);
}
}
self.put_node_data(&parent, parent_data)?;
self.put_node_data(&aggregator, aggregator_data)?;
self.stage_plan()?;
Ok(format!(
"Staged fanout consolidation from node {parent} into node {aggregator}."
))
}
fn set_fixed_connection(
&mut self,
parent: String,
child: Option<String>,
slot: usize,
) -> anyhow::Result<String> {
self.ensure_known_node(&parent)?;
if let Some(child) = &child {
self.ensure_known_node(child)?;
anyhow::ensure!(child != &parent, "a node cannot connect to itself");
}
let mut data = self.node_data(&parent)?;
if let Some(child) = child.clone() {
anyhow::ensure!(
slot <= data.fixed_connections.len() + 1,
"fixed connection positions must remain contiguous"
);
data.fixed_connections.retain(|id| id != &child);
if slot - 1 < data.fixed_connections.len() {
data.fixed_connections[slot - 1] = child;
} else {
data.fixed_connections.push(child);
}
} else if slot > 0 && slot - 1 < data.fixed_connections.len() {
data.fixed_connections.remove(slot - 1);
}
self.put_node_data(&parent, data)?;
self.stage_plan()?;
Ok(match child {
Some(child) => {
format!("Staged node {child} in fixed slot {slot} of node {parent}.")
}
None => format!("Cleared fixed slot {slot} of node {parent} in the staged plan."),
})
}
fn create_node(
&mut self,
parents: Vec<String>,
owner: String,
short_name: String,
short_description: String,
long_description: String,
) -> anyhow::Result<String> {
for id in parents.iter().chain(std::iter::once(&owner)) {
if id != "self" && id != "unowned" {
self.ensure_known_node(id)?;
}
}
let pending = self.journal.allocate_pending_node(now())?.to_string();
self.plan.creates.push(StagedNodeCreate {
pending_id: pending.clone(),
data: PlannedNode {
short_name,
short_description,
long_description,
owner,
fixed_connections: Vec::new(),
recent_connections: parents.clone(),
objects: Vec::new(),
attach_session_archive: true,
},
});
for parent in parents {
let mut data = self.node_data(&parent)?;
data.recent_connections.retain(|id| id != &pending);
data.recent_connections.insert(0, pending.clone());
self.put_node_data(&parent, data)?;
}
self.stage_plan()?;
Ok(format!("Created staged node {pending}."))
}
fn update_node(
&mut self,
id: String,
owner: String,
short_name: String,
short_description: String,
long_description: String,
) -> anyhow::Result<String> {
self.ensure_known_node(&id)?;
if owner != "self" && owner != "unowned" {
self.ensure_known_node(&owner)?;
}
let mut data = self.node_data(&id)?;
data.owner = owner;
data.short_name = short_name;
data.short_description = short_description;
data.long_description = long_description;
data.attach_session_archive = true;
self.put_node_data(&id, data)?;
self.stage_plan()?;
Ok(format!("Staged the update to node {id}."))
}
fn ensure_known_node(&self, id: &str) -> anyhow::Result<()> {
if id.starts_with("pending:") {
anyhow::ensure!(
self.plan.created(id).is_some(),
"pending node {id} is not part of this session"
);
} else {
canonical_id(id)?;
anyhow::ensure!(
self.context.contains_full_node(id) || self.plan.updates.contains_key(id),
"node {id} is not loaded; call LoadNodes first"
);
}
Ok(())
}
fn finalize_kweb_session(&mut self) -> anyhow::Result<()> {
if self.commit_receipt.is_some() {
return Ok(());
}
self.journal.repair_unfinished_tools(now())?;
self.journal.seal()?;
let archive = self.journal.archive_bytes()?;
let object_locations = self
.journal
.objects()
.iter()
.map(|(id, location)| (id.clone(), location.clone()))
.collect::<Vec<_>>();
let mut objects = BTreeMap::new();
for (id, location) in object_locations {
let pending_id = id.to_string();
let transport_kind =
kcode_kennedy_session_objects::staged_descriptor(&self.journal, &id)?
.transport_kind;
let bytes = encode_file(
&pending_id,
location.metadata.file_name.as_deref(),
&location.metadata.media_type,
transport_kind.as_deref(),
self.journal.read_object(&id)?,
)
.with_context(|| format!("encoding staged object {pending_id}"))?;
anyhow::ensure!(
objects.insert(pending_id.clone(), bytes).is_none(),
"duplicate staged object {pending_id}"
);
}
let mut creates = BTreeMap::new();
for create in &self.plan.creates {
anyhow::ensure!(
creates
.insert(create.pending_id.clone(), create.data.clone())
.is_none(),
"duplicate staged node {}",
create.pending_id
);
}
let updates = self
.plan
.updates
.iter()
.map(|(node_id, data)| {
node_id
.parse::<NodeId>()
.with_context(|| format!("{node_id:?} is not a canonical node ID"))
.map(|node_id| (node_id, data.clone()))
})
.collect::<anyhow::Result<BTreeMap<_, _>>>()?;
let result = self.api.commit_kweb_session(CommitRequest {
idempotency_key: self.journal.state().metadata.session_id.clone(),
author: self.commit_author.clone(),
source_created_at: DateTime::parse_from_rfc3339(&self.started_at)
.context("session start timestamp is invalid")?
.with_timezone(&Utc),
archive,
objects,
creates,
updates,
})?;
self.journal
.mark_completed(result.session_object_id.to_string());
self.commit_receipt = Some(result);
Ok(())
}
fn prepare_free_time_round(&mut self) -> anyhow::Result<bool> {
if !matches!(self.mode, AgentMode::FreeTime) {
return Ok(false);
}
let Some(deadline) = deadline(&self.free_time) else {
return Ok(false);
};
if Utc::now() >= deadline {
self.free_time_end_reason = Some("deadline".into());
self.journal.create_box(
now(),
"Self-time timer",
BoxOwner::Controller,
BoxContent::text(
"The self-time deadline has arrived. Finish without starting more tool work.",
),
)?;
return Ok(true);
}
Ok(false)
}
fn refresh_runtime_prompt(&mut self) -> anyhow::Result<()> {
let Some(system_box) = self
.journal
.state()
.boxes
.values()
.find(|state| matches!(state.owner, BoxOwner::System))
.map(|state| state.id)
else {
return Ok(());
};
let current = self
.journal
.state()
.boxes
.get(&system_box)
.context("system-prompt box disappeared")?
.canonical
.content
.text
.clone();
let marker = "\n\nCurrent runtime\n\n";
let Some((prefix, _)) = current.rsplit_once(marker) else {
return Ok(());
};
let refreshed = format!(
"{prefix}{marker}{}",
render(RenderRequest::RuntimeDescription {
model: &self.runtime.model,
reasoning_effort: &self.runtime.reasoning_effort,
current_time: Utc::now(),
})?
);
if refreshed != current {
self.journal
.update_box(now(), system_box, BoxContent::text(refreshed))?;
}
Ok(())
}
fn agent_request_timeout(&self) -> Option<Duration> {
if matches!(self.mode, AgentMode::Conversation) && self.session_type == "conversation" {
return Some(BROWSER_CONVERSATION_REQUEST_TIMEOUT);
}
if matches!(self.mode, AgentMode::Ingress { .. }) {
return Some(HISTORY_INGRESS_REQUEST_TIMEOUT);
}
if matches!(self.mode, AgentMode::Wakeup) {
return Some(WAKEUP_REQUEST_TIMEOUT);
}
if matches!(self.mode, AgentMode::FreeTime) {
let deadline = deadline(&self.free_time)?;
return Some(Duration::from_secs(
(deadline - Utc::now()).num_seconds().max(1) as u64 + 360,
));
}
None
}
pub fn refresh_telegram_group_context(
&mut self,
group_context: &Value,
current_message_id: Option<&str>,
) -> anyhow::Result<()> {
if self.session_type != "telegram-group" {
return Ok(());
}
self.channel["groupContext"] = group_context.clone();
self.group_context = group_context.clone();
self.journal.create_box(
now(),
"Telegram group update",
BoxOwner::Controller,
BoxContent::text(kcode_telegram_session_coordinator::format_group_context(
group_context,
)),
)?;
self.recover_context_overflow(current_message_id, &[])?;
Ok(())
}
pub fn finalize_free_time(&mut self, reason: &str) -> anyhow::Result<()> {
anyhow::ensure!(
matches!(reason, "tool" | "deadline" | "hard-stop" | "user-stop"),
"invalid self-time completion reason"
);
self.free_time["sliceEndedReason"] = json!(reason);
self.free_time["sliceEndedAt"] = json!(now());
self.pending_turn = false;
self.pending_external_event_id = None;
Ok(())
}
pub fn commit_current_write_session(&mut self) -> anyhow::Result<()> {
anyhow::ensure!(
matches!(
self.mode,
AgentMode::FreeTime | AgentMode::Wakeup | AgentMode::Ingress { .. }
),
"a read-only conversation cannot be committed as a Kweb write session"
);
self.finalize_kweb_session()?;
self.completed = true;
Ok(())
}
pub fn snapshot(&self) -> anyhow::Result<Value> {
let projection = self.journal.state().projection();
let chatend_text = projection.render();
let session_status = projection.status.clone();
Ok(json!({
"format":"kennedy-chatend",
"version":1,
"stateVersion":3,
"sessionId":self.journal.state().metadata.session_id,
"chatendMetadata":self.journal.state().metadata,
"sessionType":self.session_type,
"sourceSessionType":self.source_session_type,
"channel":self.channel,
"freeTime":self.free_time,
"orchestration":self.orchestration,
"provenanceId":self.provenance_id,
"rustLibSessionId":self.rust_lib_session_id,
"rootNodeIds":self.root_node_ids,
"referenceRootNodeIds":self.reference_root_node_ids,
"startedAt":self.started_at,
"transcript":self.transcript,
"pendingTurn":self.pending_turn,
"pendingExternalEventId":self.pending_external_event_id,
"roundsUsed":self.rounds_used,
"completed":self.completed,
"sessionObjectId":self.journal.state().completed_session_object,
"commitReceipt":self.commit_receipt,
"commitAuthor":self.commit_author,
"providerModel":self.runtime.model,
"kwebPlan":self.plan,
"boxCount":self.journal.state().boxes.len(),
"eventCount":self.journal.state().events.len(),
"boxes":self.journal.state().boxes,
"events":self.journal.state().events,
"context":projection,
"sessionStatus":session_status,
"chatendText":chatend_text,
}))
}
pub async fn release_managed_sources(&self) {
self.api
.release_managed_sources(&self.rust_lib_session_id)
.await;
}
}
impl<C, F> kcode_agent_runtime::SessionHost for KennedySessionHost<'_, C>
where
C: FnMut(Value) -> F + Send,
F: Future<Output = anyhow::Result<()>> + Send,
{
fn prepare_round<'a>(
&'a mut self,
round: u64,
) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::RoundPreparation> {
Box::pin(async move {
self.session.rounds_used = round;
self.session.refresh_runtime_prompt()?;
self.deadline_after_response = self.session.prepare_free_time_round()?;
let external_event_id = self.session.pending_external_event_id.clone();
if self
.session
.recover_context_overflow(external_event_id.as_deref(), &[])?
== ContextRecovery::Irreducible
|| (matches!(self.session.mode, AgentMode::Ingress { .. })
&& self.session.ingress_force_commit_requested())
{
return Ok(kcode_agent_runtime::RoundPreparation::Complete(None));
}
let input = self.session.journal.state().projection().render();
Ok(kcode_agent_runtime::RoundPreparation::Run(
kcode_agent_runtime::PreparedRound {
input,
model: self.session.runtime.model.clone(),
reasoning_effort: self.session.runtime.reasoning_effort.clone(),
tool_description: call_ktool_description(),
timeout: self.session.agent_request_timeout(),
},
))
})
}
fn record<'a>(
&'a mut self,
event: kcode_agent_runtime::SessionEvent,
) -> kcode_agent_runtime::HostFuture<'a, ()> {
Box::pin(async move {
match event {
kcode_agent_runtime::SessionEvent::InferenceSubmitted {
manifest_hash,
model,
..
} => {
let projection = self.session.journal.state().projection();
self.accounting = Some(kcode_intelligence_chatend::TopLevelCall::new(
manifest_hash.clone(),
model,
));
self.session.journal.record(
now(),
EventKind::InferenceSubmitted {
manifest_hash,
estimated_input_tokens: projection.estimated_tokens,
raw_estimated_input_tokens: Some(projection.raw_estimated_tokens),
},
)?;
}
kcode_agent_runtime::SessionEvent::ProviderInput { input, .. } => {
self.session.journal.record(
now(),
EventKind::Note {
label: "provider_input".into(),
value: Value::String(input),
},
)?;
}
kcode_agent_runtime::SessionEvent::UsageUpdated { usage, .. } => {
self.accounting
.as_mut()
.context("provider usage arrived before inference submission")?
.usage_updated(&mut self.session.journal, &now(), &usage)?;
}
kcode_agent_runtime::SessionEvent::ProviderReceipt { usage, .. } => {
self.accounting
.take()
.context("provider receipt arrived before inference submission")?
.completed(&mut self.session.journal, &now(), usage.as_ref())?;
}
}
let snapshot = self.session.snapshot()?;
(self.checkpoint)(snapshot).await
})
}
fn execute_tool<'a>(
&'a mut self,
call: anyhow::Result<kcode_agent_runtime::ToolCall>,
operation_id: Uuid,
) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionToolOutcome> {
Box::pin(async move {
if let Some(pending) = &self.pending_freeform_write {
let text = format!(
"{} is awaiting the complete file contents; no other Ktool can run before that output.",
pending.request.write_tool()
);
self.session
.record_tool_completion(None, json!({"ok":false,"result":text}))?;
let text = provider_tool_result_with_context_footer(&self.session.journal, &text);
(self.checkpoint)(self.session.snapshot()?).await?;
return Ok(kcode_agent_runtime::SessionToolOutcome {
text,
ok: false,
capture: Some(json!(true)),
stop: false,
finish_after_round: false,
emitted_response: false,
});
}
let tool_started_at = std::time::Instant::now();
let mut created_call_box_id = None;
let mut recorded_invocation = None;
let transcript_start = self.session.transcript.len();
let mut emitted_response = false;
let mut outcome = match call {
Ok(call) => {
let call = ToolCall {
name: call.name,
arguments: call.arguments,
};
let call_name = format!("Kennedy tool call: {}", call.name);
let call_content = tool_call_box_content(&call)?;
recorded_invocation = Some(
self.session
.record_tool_invocation(&call.name, call.arguments.clone())?,
);
created_call_box_id = Some(self.session.journal.create_box(
now(),
call_name,
BoxOwner::Kennedy,
call_content,
)?);
let external_event_id = self.session.pending_external_event_id.clone();
if self
.session
.recover_context_overflow(external_event_id.as_deref(), &[])?
== ContextRecovery::Irreducible
{
ToolOutcome {
text: CONTEXT_OVERFLOW_WARNING.into(),
store_result: false,
ok: false,
end_session: false,
freeform_write: None,
managed_source_snapshot: None,
}
} else {
match self.session.execute_tool(&call, operation_id).await {
Ok(outcome) => {
emitted_response = call.name == "EmitObject" && outcome.ok;
outcome
}
Err(error) => ToolOutcome {
text: format!("{} failed: {error}", call.name),
store_result: call.name != "LoadNodes",
ok: false,
end_session: false,
freeform_write: None,
managed_source_snapshot: None,
},
}
}
}
Err(error) => ToolOutcome {
text: error.to_string(),
store_result: true,
ok: false,
end_session: false,
freeform_write: None,
managed_source_snapshot: None,
},
};
if let Some(snapshot) = outcome.managed_source_snapshot.take() {
apply_snapshot(&mut self.session.journal, &now(), snapshot)?;
outcome.store_result = false;
}
append_slow_tool_duration(&mut outcome.text, tool_started_at.elapsed());
let capture = if let Some(request) = outcome.freeform_write.take() {
self.pending_freeform_write = Some(PendingFreeformWrite {
request,
call_box_id: created_call_box_id
.context("freeform write call box was not created")?,
});
Some(json!(true))
} else {
None
};
if outcome.store_result {
self.session.journal.create_box(
now(),
"Kennedy tool result",
BoxOwner::Controller,
BoxContent::text(&outcome.text),
)?;
}
let external_event_id = self.session.pending_external_event_id.clone();
let recovery = self
.session
.recover_context_overflow(external_event_id.as_deref(), &[])?;
let context_warning_added =
self.session.transcript[transcript_start..]
.iter()
.any(|entry| {
entry.get("contextOverflowWarning").and_then(Value::as_bool) == Some(true)
});
let mut provider_text = outcome.text.clone();
if context_warning_added && !provider_text.contains(CONTEXT_OVERFLOW_WARNING) {
if !provider_text.is_empty() {
provider_text.push_str("\n\n");
}
provider_text.push_str(CONTEXT_OVERFLOW_WARNING);
}
let provider_text =
provider_tool_result_with_context_footer(&self.session.journal, &provider_text);
self.session.record_tool_completion(
recorded_invocation.as_ref(),
json!({"ok":outcome.ok,"result":outcome.text}),
)?;
let stop = recovery == ContextRecovery::Irreducible
|| (matches!(self.session.mode, AgentMode::Ingress { .. })
&& self.session.ingress_force_commit_requested())
|| (!matches!(self.session.mode, AgentMode::Ingress { .. })
&& self.session.journal.state().source_terminated);
let snapshot = self.session.snapshot()?;
(self.checkpoint)(snapshot).await?;
Ok(kcode_agent_runtime::SessionToolOutcome {
text: provider_text,
ok: outcome.ok,
capture,
stop,
finish_after_round: outcome.end_session,
emitted_response,
})
})
}
fn complete_capture<'a>(
&'a mut self,
_capture: Value,
contents: String,
) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionControl> {
Box::pin(async move {
let pending = self
.pending_freeform_write
.take()
.context("provider completed without a pending freeform write")?;
let result_metadata = pending.request.clone();
let outcome = self
.session
.complete_freeform_write(pending, contents)
.await?;
if outcome.store_result {
self.session.journal.create_box(
now(),
"Kennedy tool result",
BoxOwner::Controller,
BoxContent::text(&outcome.text),
)?;
}
self.session.journal.record(
now(),
EventKind::Note {
label: "write_file_freeform_result".into(),
value: result_metadata.result_record(outcome.ok, &outcome.text),
},
)?;
let external_event_id = self.session.pending_external_event_id.clone();
let recovery = self
.session
.recover_context_overflow(external_event_id.as_deref(), &[])?;
let snapshot = self.session.snapshot()?;
(self.checkpoint)(snapshot).await?;
if recovery == ContextRecovery::Irreducible
|| (matches!(self.session.mode, AgentMode::Ingress { .. })
&& self.session.ingress_force_commit_requested())
|| (!matches!(self.session.mode, AgentMode::Ingress { .. })
&& self.session.journal.state().source_terminated)
|| self.deadline_after_response
{
return Ok(kcode_agent_runtime::SessionControl::Complete(None));
}
self.session.journal.create_box(
now(),
controller_box_name(&self.session.mode),
BoxOwner::Controller,
BoxContent::text(controller_message(
&self.session.mode,
&self.session.free_time,
)),
)?;
Ok(kcode_agent_runtime::SessionControl::Continue)
})
}
fn complete_round<'a>(
&'a mut self,
completion: kcode_agent_runtime::RoundCompletion,
) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionControl> {
Box::pin(async move {
let answer = completion.answer.trim().to_owned();
let mut completion_recovery = ContextRecovery::NotNeeded;
if !answer.is_empty() {
let mut content = BoxContent::text(answer.clone());
if let Some(id) = &self.session.pending_external_event_id {
content.metadata["externalEventId"] = json!(id);
}
self.session.journal.create_box(
now(),
"Kennedy message",
BoxOwner::Kennedy,
content,
)?;
let mut transcript = json!({"role":"kennedy","content":answer});
if let Some(id) = &self.session.pending_external_event_id {
transcript["externalEventId"] = json!(id);
}
self.session.transcript.push(transcript);
let external_event_id = self.session.pending_external_event_id.clone();
completion_recovery = self
.session
.recover_context_overflow(external_event_id.as_deref(), &[])?;
}
let snapshot = self.session.snapshot()?;
(self.checkpoint)(snapshot).await?;
if completion_recovery == ContextRecovery::Irreducible
|| (matches!(self.session.mode, AgentMode::Ingress { .. })
&& self.session.ingress_force_commit_requested())
|| (!matches!(self.session.mode, AgentMode::Ingress { .. })
&& self.session.journal.state().source_terminated)
{
return Ok(kcode_agent_runtime::SessionControl::Complete(None));
}
if completion.finish_requested || self.deadline_after_response {
return Ok(kcode_agent_runtime::SessionControl::Complete(
(!answer.is_empty()).then_some(answer),
));
}
if matches!(self.session.mode, AgentMode::Conversation) && !answer.is_empty() {
return Ok(kcode_agent_runtime::SessionControl::Complete(Some(answer)));
}
if matches!(self.session.mode, AgentMode::Conversation) && completion.emitted_response {
return Ok(kcode_agent_runtime::SessionControl::Complete(None));
}
let solo_ingress_response =
matches!(self.session.mode, AgentMode::Ingress { .. }) && !answer.is_empty();
anyhow::ensure!(
completion.used_tool || solo_ingress_response,
"provider completed without a response or tool call"
);
self.session.journal.create_box(
now(),
controller_box_name(&self.session.mode),
BoxOwner::Controller,
BoxContent::text(controller_message(
&self.session.mode,
&self.session.free_time,
)),
)?;
Ok(kcode_agent_runtime::SessionControl::Continue)
})
}
}
impl kcode_agent_runtime::Host for KennedySubagentHost<'_> {
fn render_tool_call(&mut self, call: &kcode_agent_runtime::ToolCall) -> anyhow::Result<String> {
Ok(tool_call_box_content(&ToolCall {
name: call.name.clone(),
arguments: call.arguments.clone(),
})?
.text)
}
fn execute_tool<'a>(
&'a mut self,
call: kcode_agent_runtime::ToolCall,
operation_id: Uuid,
budget: kcode_agent_runtime::ContextBudget,
) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::ToolOutcome> {
Box::pin(async move {
let call = ToolCall {
name: call.name,
arguments: call.arguments,
};
if call.name == "RunSubagent" {
return Ok(kcode_agent_runtime::ToolOutcome::failure(
"RunSubagent is unavailable inside a subagent. Only Kennedy may launch subagents.",
));
}
if matches!(
call.name.as_str(),
"SendTelegramDM" | "SendTelegramGroupMessage"
) {
return Ok(kcode_agent_runtime::ToolOutcome::failure(
"Telegram cold delivery is unavailable inside a subagent. Only Kennedy may initiate a Telegram message.",
));
}
if budget.estimated_tokens() > budget.max_input_tokens() {
return Ok(kcode_agent_runtime::ToolOutcome::failure(
"The Ktool call was not run because its retained invocation would exceed the subagent context limit.",
));
}
if !subagent_managed_write_fits(&self.session.journal, &call, &budget) {
return Ok(kcode_agent_runtime::ToolOutcome::failure(
"The managed-source write was not run because its resulting current state would exceed the subagent context limit.",
));
}
let tool_started_at = std::time::Instant::now();
let previous = canonical_box_versions(&self.session.journal);
let mut outcome = match self.session.execute_tool(&call, operation_id).await {
Ok(outcome) => outcome,
Err(error) => {
let mut text = format!("{} failed: {error}", call.name);
append_slow_tool_duration(&mut text, tool_started_at.elapsed());
return Ok(kcode_agent_runtime::ToolOutcome::failure(text));
}
};
let displays_managed_snapshot = outcome
.managed_source_snapshot
.as_ref()
.is_some_and(|snapshot| result_displays_snapshot(&outcome.text, snapshot));
append_slow_tool_duration(&mut outcome.text, tool_started_at.elapsed());
let displayed_state_keys =
if let Some(snapshot) = outcome.managed_source_snapshot.take() {
let box_id = apply_snapshot(&mut self.session.journal, &now(), snapshot)?;
outcome.store_result = false;
if displays_managed_snapshot {
vec![subagent_state_key(box_id)]
} else {
Default::default()
}
} else {
Vec::new()
};
let states = if outcome.ok {
changed_subagent_tool_states(&self.session.journal, &previous)
} else {
Vec::new()
};
let hidden = states
.iter()
.filter(|state| state.hide_from_parent && state.text.is_some())
.map(|state| state.box_id)
.collect::<Vec<_>>();
if !hidden.is_empty() {
self.session.journal.dehydrate_boxes(now(), &hidden)?;
}
let capture = outcome.freeform_write.take().map(|request| {
let id = Uuid::new_v4().to_string();
self.captures.insert(id.clone(), request);
Value::String(id)
});
Ok(kcode_agent_runtime::ToolOutcome {
text: outcome.text,
ok: outcome.ok,
state_updates: subagent_state_updates(&states),
displayed_state_keys,
capture,
})
})
}
fn complete_capture<'a>(
&'a mut self,
capture: Value,
contents: String,
budget: kcode_agent_runtime::ContextBudget,
) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::ToolOutcome> {
Box::pin(async move {
let id = capture
.as_str()
.context("subagent freeform capture token is invalid")?;
let request = self
.captures
.remove(id)
.context("subagent freeform capture token is unknown")?;
let (outcome, states) = self
.session
.complete_subagent_freeform_write(request, contents, &budget)
.await?;
let hidden = states
.iter()
.filter(|state| state.hide_from_parent && state.text.is_some())
.map(|state| state.box_id)
.collect::<Vec<_>>();
if !hidden.is_empty() {
self.session.journal.dehydrate_boxes(now(), &hidden)?;
}
Ok(kcode_agent_runtime::ToolOutcome {
text: outcome.text,
ok: outcome.ok,
state_updates: subagent_state_updates(&states),
displayed_state_keys: Vec::new(),
capture: None,
})
})
}
fn record(&mut self, event: kcode_agent_runtime::AuditEvent) -> anyhow::Result<()> {
kcode_intelligence_chatend::record_subagent_event(&mut self.session.journal, &now(), &event)
}
}
fn cost_summary(label: &str, estimated_cost_usd_nanos: u64, unpriced_calls: u64) -> String {
render(RenderRequest::CostSummary {
label,
estimated_cost_usd_nanos,
unpriced_calls,
})
.expect("cost-summary rendering is infallible")
}
fn restore_kweb_context(journal: &HistorySession, context: &mut KwebContext) -> anyhow::Result<()> {
let Some(tool) = journal.state().tools.get(KWEB_TOOL_INSTANCE) else {
return Ok(());
};
let mut nodes = BTreeMap::new();
for slot in &tool.slots {
let state = journal
.state()
.box_state(slot.box_id)
.context("Kweb slot references a missing box")?;
if let Some(node) = state.canonical.content.metadata.get("storedNode") {
let node = match serde_json::from_value::<KwebNode>(node.clone()) {
Ok(node) => node,
Err(_) => node_from_value(node).context("decoding a stored Kweb context node")?,
};
nodes.insert(node.id.clone(), node);
}
}
let mut direct = journal
.state()
.events
.iter()
.flat_map(|event| {
let EventKind::ToolInvoked {
tool_name,
arguments,
..
} = &event.kind
else {
return Vec::new();
};
match tool_name.as_str() {
"LoadNodes" => arguments
.get("identifiers")
.and_then(Value::as_array)
.into_iter()
.flatten()
.filter_map(Value::as_str)
.map(str::to_owned)
.collect(),
"LoadNode" => arguments
.get("identifier")
.and_then(Value::as_str)
.map(str::to_owned)
.into_iter()
.collect(),
_ => Vec::new(),
}
})
.collect::<Vec<_>>();
if direct.is_empty() {
direct = context.root_node_ids().to_vec();
}
context
.restore(nodes.into_values(), direct)
.map_err(anyhow::Error::new)
}
fn is_terminal_external_response(entry: &Value) -> bool {
match entry.get("role").and_then(Value::as_str) {
Some("kennedy") => true,
Some("system") => {
entry.get("contextOverflowWarning").and_then(Value::as_bool) != Some(true)
}
_ => false,
}
}
fn restore_pending_turn(restored: Option<&Value>, transcript: &[Value]) -> (bool, Option<String>) {
let mut unanswered_external = Vec::<String>::new();
let mut answered_external = HashSet::<String>::new();
for entry in transcript {
if entry.get("role").and_then(Value::as_str) == Some("user") {
if let Some(id) = entry.get("externalEventId").and_then(Value::as_str) {
answered_external.remove(id);
unanswered_external.retain(|candidate| candidate != id);
unanswered_external.push(id.to_owned());
}
continue;
}
if !is_terminal_external_response(entry) {
continue;
}
if let Some(id) = entry.get("externalEventId").and_then(Value::as_str) {
answered_external.insert(id.to_owned());
unanswered_external.retain(|candidate| candidate != id);
} else if entry.get("userStopped").and_then(Value::as_bool) == Some(true)
&& let Some(id) = unanswered_external.pop()
{
answered_external.insert(id);
}
}
let journal_pending_external = unanswered_external.last().cloned();
let restored_pending = restored
.and_then(|state| state.get("pendingTurn"))
.and_then(Value::as_bool)
.unwrap_or(false);
let restored_external = restored
.and_then(|state| state.get("pendingExternalEventId"))
.and_then(Value::as_str);
let restored_external_answered =
restored_external.is_some_and(|id| answered_external.contains(id));
let pending_external_event_id = journal_pending_external.or_else(|| {
restored_external
.filter(|id| !answered_external.contains(*id))
.map(str::to_owned)
});
let pending_turn =
pending_external_event_id.is_some() || (restored_pending && !restored_external_answered);
(pending_turn, pending_external_event_id)
}
fn kweb_node_draft(node: &PlannedNode) -> NodeDraft {
NodeDraft {
short_name: node.short_name.clone(),
short_description: node.short_description.clone(),
long_description: node.long_description.clone(),
owner: node.owner.clone(),
fixed_connections: node.fixed_connections.clone(),
recent_connections: node.recent_connections.clone(),
objects: node.objects.clone(),
}
}
fn planned_node(node: &KwebNode) -> PlannedNode {
PlannedNode {
short_name: node.short_name.clone(),
short_description: node.short_description.clone(),
long_description: node.long_description.clone(),
owner: node.owner.clone(),
fixed_connections: node
.fixed_connections
.iter()
.map(|connection| connection.id.clone())
.collect(),
recent_connections: node
.recent_connections
.iter()
.map(|connection| connection.id.clone())
.collect(),
objects: node.objects.clone(),
attach_session_archive: true,
}
}
fn session_kind(session_type: &str, mode: &AgentMode) -> SessionKind {
if matches!(mode, AgentMode::Ingress { .. }) {
return SessionKind::HistoryIngress;
}
match session_type {
"conversation" => SessionKind::Conversation,
"telegram" => SessionKind::Telegram,
"telegram-group" => SessionKind::TelegramGroup,
"free-time" => SessionKind::SelfTime,
"wakeup" => SessionKind::Other("wakeup".into()),
"audio" => SessionKind::AudioIngress,
other => SessionKind::Other(other.into()),
}
}
fn tool_instance(name: &str) -> String {
if name == "LoadNodes" {
return KWEB_TOOL_INSTANCE.into();
}
format!("{name}:{}", Uuid::new_v4())
}
fn canonical_id(value: &str) -> anyhow::Result<String> {
value
.parse::<NodeId>()
.with_context(|| format!("{value:?} is not a canonical node ID"))?;
Ok(value.into())
}
fn image_extension(media_type: &str) -> &'static str {
match media_type
.split(';')
.next()
.unwrap_or(media_type)
.trim()
.to_ascii_lowercase()
.as_str()
{
"image/jpeg" => "jpg",
"image/webp" => "webp",
_ => "png",
}
}
fn call_ktool_description() -> String {
render(RenderRequest::CallKtoolDescription).expect("Ktool-description rendering is infallible")
}
fn now() -> String {
Utc::now().to_rfc3339()
}
fn deadline(value: &Value) -> Option<DateTime<Utc>> {
value
.get("deadlineAt")
.and_then(Value::as_str)
.and_then(|value| DateTime::parse_from_rfc3339(value).ok())
.map(|value| value.with_timezone(&Utc))
}
fn controller_box_name(mode: &AgentMode) -> &'static str {
match mode {
AgentMode::Conversation => "Turn continuation",
AgentMode::FreeTime => "Self-time continuation",
AgentMode::Wakeup => "Wakeup continuation",
AgentMode::Ingress { .. } => "History-ingress continuation",
}
}
fn controller_message(mode: &AgentMode, free_time: &Value) -> String {
let mode = match mode {
AgentMode::Conversation => "conversation",
AgentMode::FreeTime => "free-time",
AgentMode::Wakeup => "wakeup",
AgentMode::Ingress { .. } => "ingress",
};
render(RenderRequest::ControllerMessage { mode, free_time })
.expect("known controller modes render successfully")
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn subagent_source_displays_and_state_updates_share_a_box_identity() {
let box_id = BoxId(7);
let updates = subagent_state_updates(&[ChangedSubagentState {
box_id,
name: "Rust Library example".into(),
text: Some("updated source".into()),
hide_from_parent: false,
}]);
assert_eq!(subagent_state_key(box_id), updates[0].key);
}
#[test]
fn only_complete_snapshot_results_claim_to_display_managed_state() {
let snapshot = SourceSnapshot {
kind: kcode_dev_tools::ManagedSourceKind::RustLibrary,
name: "example".into(),
text: "complete source".into(),
};
assert!(result_displays_snapshot("complete source", &snapshot));
assert!(!result_displays_snapshot(
"Wrote file src/lib.rs in Rust library example.",
&snapshot
));
}
}