mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-08-25 14:53:34 +00:00
* feat: add inbound attachment support to WASM channel system Add attachment record to WIT interface and implement inbound media parsing across all four channel implementations (Telegram, Slack, WhatsApp, Discord). Attachments flow from WASM channels through EmittedMessage to IncomingMessage with validation (size limits, MIME allowlist, count caps) at the host boundary. - Add `attachment` record to `emitted-message` in wit/channel.wit - Add `IncomingAttachment` struct to channel.rs and re-export - Add host-side validation (20MB total, 10 max, MIME allowlist) - Telegram: parse photo, document, audio, video, voice, sticker - Slack: parse file attachments with url_private - WhatsApp: parse image, audio, video, document with captions - Discord: backward-compatible empty attachments - Update FEATURE_PARITY.md section 7 - Add fixture-based tests per channel and host integration tests [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> * feat: integrate outbound attachment support and reconcile WIT types (#409) Reconcile PR #409's outbound attachment work with our inbound attachment support into a unified design: WIT type split: - `inbound-attachment` in channel-host: metadata-only (id, mime_type, filename, size_bytes, source_url, storage_key, extracted_text) - `attachment` in channel: raw bytes (filename, mime_type, data) on agent-response for outbound sending Outbound features (from PR #409): - `on-broadcast` WIT export for proactive messages without prior inbound - Telegram: multipart sendPhoto/sendDocument with auto photo→document fallback for files >10MB - wrapper.rs: `call_on_broadcast`, `read_attachments` from disk, attachment params threaded through `call_on_respond` - HTTP tool: `save_to` param for binary downloads to /tmp/ (50MB limit, path traversal protection, SSRF-safe redirect following) - Message tool: allow /tmp/ paths for attachments alongside base_dir - Credential env var fallback in inject_channel_credentials Channel updates: - All 4 channels implement on_broadcast (Telegram full, others stub) - Telegram: polling_enabled config, adjusted poll timeout - Inbound attachment types renamed to InboundAttachment in all channels Tests: 1965 passing (9 new), 0 clippy warnings [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> * feat: add audio transcription pipeline and extensible WIT attachment design Add host-side transcription middleware (OpenAI Whisper) that detects audio attachments with inline data on incoming messages and transcribes them automatically. Refactor WIT inbound-attachment to use extras-json and a store-attachment-data host function instead of typed fields, so future attachment properties (dimensions, codec, etc.) don't require WIT changes that invalidate all channel plugins. - Add src/transcription/ module: TranscriptionProvider trait, TranscriptionMiddleware, AudioFormat enum, OpenAI Whisper provider - Add src/config/transcription.rs: TRANSCRIPTION_ENABLED/MODEL/BASE_URL - Wire middleware into agent message loop via AgentDeps - WIT: replace data + duration-secs with extras-json + store-attachment-data - Host: parse extras-json for well-known keys, merge stored binary data - Telegram: download voice files via store-attachment-data, add duration to extras-json, add /file/bot to HTTP allowlist, voice-only placeholder - Add reqwest multipart feature for Whisper API uploads - 5 regression tests for transcription middleware Co-Authored-By: Claude Opus 4.6 <[email protected]> * feat: wire attachment processing into LLM pipeline with multimodal image support Attachments on incoming messages are now augmented into user text via XML tags before entering the turn system, and images with data are passed as multimodal content parts (base64 data URIs) to LLM providers. This enables audio transcripts, document text, and image content to reach the LLM without changes to ChatMessage serialization or provider interfaces. - Add src/agent/attachments.rs with augment_with_attachments() and 9 unit tests - Add ContentPart/ImageUrl types to llm::provider with OpenAI-compatible serde - Carry image_content_parts transiently on Turn (skipped in serialization) - Update nearai_chat and rig_adapter to serialize multimodal content - Add 3 e2e tests verifying attachments flow through the full agent loop Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: CI failures — formatting, version bumps, and Telegram voice test - Fix cargo fmt formatting in attachments.rs, nearai_chat.rs, rig_adapter.rs, e2e_attachments.rs - Bump channel registry versions 0.1.0 → 0.2.0 (discord, slack, telegram, whatsapp) to satisfy version-bump CI check - Fix Telegram test_extract_attachments_voice: add missing required `duration` field to voice fixture JSON Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: bump WIT channel version to 0.3.0, fix Telegram voice test, add pre-commit hook - Bump wit/channel.wit package version 0.2.0 → 0.3.0 (interface changed with store-attachment-data) - Update WIT_CHANNEL_VERSION constant and registry wit_version fields to match - Fix Telegram test_extract_attachments_voice: gate voice download behind #[cfg(target_arch = "wasm32")] so host functions aren't called in native tests, update assertions for generated filename and extras_json duration - Add @0.3.0 linker stubs in wit_compat.rs - Add .githooks/pre-commit hook that runs scripts/check-version-bumps.sh when WIT or extension sources are staged - Symlink commit-msg regression hook into .githooks/ [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> * refactor: extract voice download from extract_attachments into handle_message Move download_voice_file + store_attachment_data calls out of extract_attachments into a separate download_and_store_voice function called from handle_message. This keeps extract_attachments as a pure data-mapping function with no host calls, making it fully testable in native unit tests without #[cfg(target_arch)] gates. [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: address PR review comments — security, correctness, and code quality Security fixes: - Add path validation to read_attachments (restrict to /tmp/) preventing arbitrary file reads from compromised tools - Escape XML special characters in attachment filenames, MIME types, and extracted text to prevent prompt injection via tag spoofing - Percent-encode file_id in Telegram getFile URL to prevent query injection - Clone SecretString directly instead of expose_secret().to_string() Correctness fixes: - Fix store_attachment_data overwrite accounting: subtract old entry size before adding new to prevent inflated totals and false rejections - Use max(reported, stored_size) for attachment size accounting to prevent WASM channels from under-reporting size_bytes to bypass limits - Add application/octet-stream to MIME allowlist (channels default unknown types to this) Code quality: - Extract send_response helper in Telegram, deduplicating on_respond and on_broadcast - Rename misleading Discord test to test_parse_slash_command_interaction - Fix .githooks/commit-msg to use relative symlink (portable across machines) [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> * feat: add tool_upgrade command + fix TOCTOU in save_to path validation Add `tool_upgrade` — a new extension management tool that automatically detects and reinstalls WASM extensions with outdated WIT versions. Preserves authentication secrets during upgrade. Supports upgrading a single extension by name or all installed WASM tools/channels at once. Fix TOCTOU in `validate_save_to_path`: validate the path *before* creating parent directories, so traversal paths like `/tmp/../../etc/` cannot cause filesystem mutations outside /tmp before being rejected. [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: unify WIT package version to 0.3.0 across tool.wit and all capabilities tool.wit and channel.wit share the `near:agent` package namespace, so they must declare the same version. Bumps tool.wit from 0.2.0 to 0.3.0 and updates all capabilities files and registry entries to match. Fixes `cargo component build` failure: "package identifier near:[email protected] does not match previous package name of near:[email protected]" [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: move WIT file comments after package declaration WIT treats `//` comments before `package` as doc comments. When both tool.wit and channel.wit had header comments, the parser rejected them as "doc comments on multiple 'package' items". Move comments after the package declaration in both files. Also bumps tool registry versions to 0.2.0 to match the WIT 0.3.0 bump. [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> * feat: display extension versions in gateway Extensions tab Add version field to InstalledExtension and RegistryEntry types, pipe through the web API (ExtensionInfo, RegistryEntryInfo), and render as a badge in the gateway UI for both installed and available extensions. For installed WASM extensions, version is read from the capabilities file with a fallback to the registry entry when the local file has no version (old installations). Bump all extension Cargo.toml and registry JSON versions from 0.1.0 to 0.2.0 to keep them in sync. [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> * feat: add document text extraction middleware for PDF, Office, and text files Extract text from document attachments (PDF, DOCX, PPTX, XLSX, RTF, plain text, code files) so the LLM can reason about uploaded documents. Uses pdf-extract for PDFs, zip+XML parsing for Office XML formats, and UTF-8 decode for text files. Wired into the agent loop after transcription middleware. Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: download document files in Telegram channel for text extraction The DocumentExtractionMiddleware needs file bytes in the attachment `data` field, but only voice files were being downloaded. Document attachments (PDFs, DOCX, etc.) had empty `data` and a source_url with a credential placeholder that only works inside the WASM host's http_request. Add `download_and_store_documents()` that downloads non-voice, non-image, non-audio attachments via the existing two-step getFile→download flow and stores bytes via `store_attachment_data` for host-side extraction. Also rename `download_voice_file` → `download_telegram_file` since it's generic for any file_id. Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: allow Office MIME types and increase file download limit for Telegram Two issues preventing document extraction from Telegram: 1. PPTX/DOCX/XLSX MIME types (application/vnd.*) were dropped by the WASM host attachment allowlist — add application/vnd., application/msword, and application/rtf prefixes. 2. Telegram file downloads over 10 MB failed with "Response body too large" — set max_response_bytes to 20 MB in Telegram capabilities. [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: report document extraction errors back to user instead of silently skipping - Bump max_response_bytes to 50 MB for Telegram file downloads - When document extraction fails (too large, download error, parse error), set extracted_text to a user-friendly error message instead of leaving it None. This ensures the LLM tells the user what went wrong. - On Telegram download failure, set extracted_text with the error so the user sees feedback even when the file never reaches the extraction middleware. Co-Authored-By: Claude Opus 4.6 <[email protected]> * feat: store extracted document text in workspace memory for search/recall After document extraction succeeds, write the extracted text to workspace memory at `documents/{date}/{filename}`. This enables: - Full-text and semantic search over past uploaded documents - Cross-conversation recall ("what did that PDF say?") - Automatic chunking and embedding via the workspace pipeline Documents are stored with metadata header (uploader, channel, date, MIME type). Error messages (extraction failures) are not stored — only successful extractions. Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: CI failures — formatting, unused assignment warning - Run cargo fmt on document_extraction and agent_loop modules - Suppress unused_assignments warning on trace_llm_ref (used only behind #[cfg(feature = "libsql")]) [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: address PR review comments — security, correctness, and code quality Security fixes: - Remove SSRF-prone download() from DocumentExtractionMiddleware (#13) - Sanitize filenames in workspace path to prevent directory traversal (#11) - Pre-check file size before reading in WASM wrapper to prevent OOM (#2) - Percent-encode file_id in Telegram source URLs (#7) Correctness fixes: - Clear image_content_parts on turn end to prevent memory leak (#1) - Find first *successful* transcription instead of first overall (#3) - Enforce data.len() size limit in document extraction (#10) - Use UTF-8 safe truncation with char_indices() (#12) Robustness & code quality: - Add 120s timeout to OpenAI Whisper HTTP client (#5) - Trim trailing slash from Whisper base_url (#6) - Allow ~/.ironclaw/ paths in WASM wrapper (#8) - Return error from on_broadcast in Slack/Discord/WhatsApp (#9) - Fix doc comment in HTTP tool (#4) Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: formatting — cargo fmt Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: address latest PR review — doc comments, error messages, version bumps - Fix DocumentExtractionMiddleware doc comment (no longer downloads from source_url) - Fix error message: "no inline data" instead of "no download URL" - Log error + fallback instead of silent unwrap_or_default on Whisper HTTP client - Bump all capabilities.json versions from 0.1.0 to 0.2.0 to match Cargo.toml Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: remove unsupported profile: minimal from CI workflows [skip-regression-check] dtolnay/rust-toolchain@stable does not accept the 'profile' input (it was a parameter for the deprecated actions-rs/toolchain action). Co-Authored-By: Claude Opus 4.6 <[email protected]> * fix: merge with latest main — resolve compilation errors and PR review nits - Add version: None to RegistryEntry/InstalledExtension test constructors - Fix MessageContent type mismatches in nearai_chat tests (String → MessageContent::Text) - Fix .contains() calls on MessageContent — use .as_text().unwrap() - Remove redundant trace_llm_ref = None assignment in test_rig - Check data size before clone in document extraction to avoid unnecessary allocation [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <[email protected]> --------- Co-authored-by: Claude Opus 4.6 <[email protected]>
1481 lines
55 KiB
Rust
1481 lines
55 KiB
Rust
//! Thread and session operations for the agent.
|
|
//!
|
|
//! Extracted from `agent_loop.rs` to isolate thread management (user input
|
|
//! processing, undo/redo, approval, auth, persistence) from the core loop.
|
|
|
|
use std::sync::Arc;
|
|
|
|
use tokio::sync::Mutex;
|
|
use tokio::task::JoinSet;
|
|
use uuid::Uuid;
|
|
|
|
use crate::agent::Agent;
|
|
use crate::agent::compaction::ContextCompactor;
|
|
use crate::agent::dispatcher::{
|
|
AgenticLoopResult, check_auth_required, execute_chat_tool_standalone, parse_auth_result,
|
|
};
|
|
use crate::agent::session::{PendingApproval, Session, ThreadState};
|
|
use crate::agent::submission::SubmissionResult;
|
|
use crate::channels::web::util::truncate_preview;
|
|
use crate::channels::{IncomingMessage, StatusUpdate};
|
|
use crate::context::JobContext;
|
|
use crate::error::Error;
|
|
use crate::llm::ChatMessage;
|
|
use crate::tools::redact_params;
|
|
|
|
impl Agent {
|
|
/// Hydrate a historical thread from DB into memory if not already present.
|
|
///
|
|
/// Called before `resolve_thread` so that the session manager finds the
|
|
/// thread on lookup instead of creating a new one.
|
|
///
|
|
/// Creates an in-memory thread with the exact UUID the frontend sent,
|
|
/// even when the conversation has zero messages (e.g. a brand-new
|
|
/// assistant thread). Without this, `resolve_thread` would mint a
|
|
/// fresh UUID and all messages would land in the wrong conversation.
|
|
pub(super) async fn maybe_hydrate_thread(
|
|
&self,
|
|
message: &IncomingMessage,
|
|
external_thread_id: &str,
|
|
) {
|
|
// Only hydrate UUID-shaped thread IDs (web gateway uses UUIDs)
|
|
let thread_uuid = match Uuid::parse_str(external_thread_id) {
|
|
Ok(id) => id,
|
|
Err(_) => return,
|
|
};
|
|
|
|
// Check if already in memory
|
|
let session = self
|
|
.session_manager
|
|
.get_or_create_session(&message.user_id)
|
|
.await;
|
|
{
|
|
let sess = session.lock().await;
|
|
if sess.threads.contains_key(&thread_uuid) {
|
|
return;
|
|
}
|
|
}
|
|
|
|
// Load history from DB (may be empty for a newly created thread).
|
|
let mut chat_messages: Vec<ChatMessage> = Vec::new();
|
|
let msg_count;
|
|
|
|
if let Some(store) = self.store() {
|
|
let db_messages = store
|
|
.list_conversation_messages(thread_uuid)
|
|
.await
|
|
.unwrap_or_default();
|
|
msg_count = db_messages.len();
|
|
chat_messages = db_messages
|
|
.iter()
|
|
.filter_map(|m| match m.role.as_str() {
|
|
"user" => Some(ChatMessage::user(&m.content)),
|
|
"assistant" => Some(ChatMessage::assistant(&m.content)),
|
|
// tool_calls rows are UI metadata (tool name + preview),
|
|
// not part of the LLM conversation context.
|
|
_ => None,
|
|
})
|
|
.collect();
|
|
} else {
|
|
msg_count = 0;
|
|
}
|
|
|
|
// Create thread with the historical ID and restore messages
|
|
let session_id = {
|
|
let sess = session.lock().await;
|
|
sess.id
|
|
};
|
|
|
|
let mut thread = crate::agent::session::Thread::with_id(thread_uuid, session_id);
|
|
if !chat_messages.is_empty() {
|
|
thread.restore_from_messages(chat_messages);
|
|
}
|
|
|
|
// Insert into session and register with session manager
|
|
{
|
|
let mut sess = session.lock().await;
|
|
sess.threads.insert(thread_uuid, thread);
|
|
sess.active_thread = Some(thread_uuid);
|
|
sess.last_active_at = chrono::Utc::now();
|
|
}
|
|
|
|
self.session_manager
|
|
.register_thread(
|
|
&message.user_id,
|
|
&message.channel,
|
|
thread_uuid,
|
|
Arc::clone(&session),
|
|
)
|
|
.await;
|
|
|
|
tracing::debug!(
|
|
"Hydrated thread {} from DB ({} messages)",
|
|
thread_uuid,
|
|
msg_count
|
|
);
|
|
}
|
|
|
|
pub(super) async fn process_user_input(
|
|
&self,
|
|
message: &IncomingMessage,
|
|
session: Arc<Mutex<Session>>,
|
|
thread_id: Uuid,
|
|
content: &str,
|
|
) -> Result<SubmissionResult, Error> {
|
|
// First check thread state without holding lock during I/O
|
|
let thread_state = {
|
|
let sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
thread.state
|
|
};
|
|
|
|
// Check thread state
|
|
match thread_state {
|
|
ThreadState::Processing => {
|
|
return Ok(SubmissionResult::error(
|
|
"Turn in progress. Use /interrupt to cancel.",
|
|
));
|
|
}
|
|
ThreadState::AwaitingApproval => {
|
|
return Ok(SubmissionResult::error(
|
|
"Waiting for approval. Use /interrupt to cancel.",
|
|
));
|
|
}
|
|
ThreadState::Completed => {
|
|
return Ok(SubmissionResult::error(
|
|
"Thread completed. Use /thread new.",
|
|
));
|
|
}
|
|
ThreadState::Idle | ThreadState::Interrupted => {
|
|
// Can proceed
|
|
}
|
|
}
|
|
|
|
// Safety validation for user input
|
|
let validation = self.safety().validate_input(content);
|
|
if !validation.is_valid {
|
|
let details = validation
|
|
.errors
|
|
.iter()
|
|
.map(|e| format!("{}: {}", e.field, e.message))
|
|
.collect::<Vec<_>>()
|
|
.join("; ");
|
|
return Ok(SubmissionResult::error(format!(
|
|
"Input rejected by safety validation: {}",
|
|
details
|
|
)));
|
|
}
|
|
|
|
let violations = self.safety().check_policy(content);
|
|
if violations
|
|
.iter()
|
|
.any(|rule| rule.action == crate::safety::PolicyAction::Block)
|
|
{
|
|
return Ok(SubmissionResult::error("Input rejected by safety policy."));
|
|
}
|
|
|
|
// Scan inbound messages for secrets (API keys, tokens).
|
|
// Catching them here prevents the LLM from echoing them back, which
|
|
// would trigger the outbound leak detector and create error loops.
|
|
if let Some(warning) = self.safety().scan_inbound_for_secrets(content) {
|
|
tracing::warn!(
|
|
user = %message.user_id,
|
|
channel = %message.channel,
|
|
"Inbound message blocked: contains leaked secret"
|
|
);
|
|
return Ok(SubmissionResult::error(warning));
|
|
}
|
|
|
|
// Handle explicit commands (starting with /) directly
|
|
// Everything else goes through the normal agentic loop with tools
|
|
let temp_message = IncomingMessage {
|
|
content: content.to_string(),
|
|
..message.clone()
|
|
};
|
|
|
|
if let Some(intent) = self.router.route_command(&temp_message) {
|
|
// Explicit command like /status, /job, /list - handle directly
|
|
return self.handle_job_or_command(intent, message).await;
|
|
}
|
|
|
|
// Natural language goes through the agentic loop
|
|
// Job tools (create_job, list_jobs, etc.) are in the tool registry
|
|
|
|
// Auto-compact if needed BEFORE adding new turn
|
|
{
|
|
let mut sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get_mut(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
|
|
let messages = thread.messages();
|
|
if let Some(strategy) = self.context_monitor.suggest_compaction(&messages) {
|
|
let pct = self.context_monitor.usage_percent(&messages);
|
|
tracing::info!("Context at {:.1}% capacity, auto-compacting", pct);
|
|
|
|
// Notify the user that compaction is happening
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::Status(format!(
|
|
"Context at {:.0}% capacity, compacting...",
|
|
pct
|
|
)),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
|
|
let compactor = ContextCompactor::new(self.llm().clone(), self.safety().clone());
|
|
if let Err(e) = compactor
|
|
.compact(thread, strategy, self.workspace().map(|w| w.as_ref()))
|
|
.await
|
|
{
|
|
tracing::warn!("Auto-compaction failed: {}", e);
|
|
}
|
|
}
|
|
}
|
|
|
|
// Create checkpoint before turn
|
|
let undo_mgr = self.session_manager.get_undo_manager(thread_id).await;
|
|
{
|
|
let sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
|
|
let mut mgr = undo_mgr.lock().await;
|
|
mgr.checkpoint(
|
|
thread.turn_number(),
|
|
thread.messages(),
|
|
format!("Before turn {}", thread.turn_number()),
|
|
);
|
|
}
|
|
|
|
// Augment content with attachment context (transcripts, metadata, images)
|
|
let augmented =
|
|
crate::agent::attachments::augment_with_attachments(content, &message.attachments);
|
|
let (effective_content, image_parts) = match &augmented {
|
|
Some(result) => (result.text.as_str(), result.image_parts.clone()),
|
|
None => (content, Vec::new()),
|
|
};
|
|
|
|
// Start the turn and get messages
|
|
let turn_messages = {
|
|
let mut sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get_mut(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
let turn = thread.start_turn(effective_content);
|
|
turn.image_content_parts = image_parts;
|
|
thread.messages()
|
|
};
|
|
|
|
// Persist user message to DB immediately so it survives crashes
|
|
self.persist_user_message(thread_id, &message.user_id, effective_content)
|
|
.await;
|
|
|
|
// Send thinking status
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::Thinking("Processing...".into()),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
|
|
// Run the agentic tool execution loop
|
|
let result = self
|
|
.run_agentic_loop(message, session.clone(), thread_id, turn_messages)
|
|
.await;
|
|
|
|
// Re-acquire lock and check if interrupted
|
|
let mut sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get_mut(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
|
|
if thread.state == ThreadState::Interrupted {
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::Status("Interrupted".into()),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
return Ok(SubmissionResult::Interrupted);
|
|
}
|
|
|
|
// Complete, fail, or request approval
|
|
match result {
|
|
Ok(AgenticLoopResult::Response(response)) => {
|
|
// Hook: TransformResponse — allow hooks to modify or reject the final response
|
|
let response = {
|
|
let event = crate::hooks::HookEvent::ResponseTransform {
|
|
user_id: message.user_id.clone(),
|
|
thread_id: thread_id.to_string(),
|
|
response: response.clone(),
|
|
};
|
|
match self.hooks().run(&event).await {
|
|
Err(crate::hooks::HookError::Rejected { reason }) => {
|
|
format!("[Response filtered: {}]", reason)
|
|
}
|
|
Err(err) => {
|
|
format!("[Response blocked by hook policy: {}]", err)
|
|
}
|
|
Ok(crate::hooks::HookOutcome::Continue {
|
|
modified: Some(new_response),
|
|
}) => new_response,
|
|
_ => response, // fail-open: use original
|
|
}
|
|
};
|
|
|
|
thread.complete_turn(&response);
|
|
let tool_calls = thread
|
|
.turns
|
|
.last()
|
|
.map(|t| t.tool_calls.clone())
|
|
.unwrap_or_default();
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::Status("Done".into()),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
|
|
// Persist tool calls then assistant response (user message already persisted at turn start)
|
|
self.persist_tool_calls(thread_id, &message.user_id, &tool_calls)
|
|
.await;
|
|
self.persist_assistant_response(thread_id, &message.user_id, &response)
|
|
.await;
|
|
|
|
Ok(SubmissionResult::response(response))
|
|
}
|
|
Ok(AgenticLoopResult::NeedApproval { pending }) => {
|
|
// Store pending approval in thread and update state
|
|
let request_id = pending.request_id;
|
|
let tool_name = pending.tool_name.clone();
|
|
let description = pending.description.clone();
|
|
let parameters = pending.display_parameters.clone();
|
|
thread.await_approval(pending);
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::Status("Awaiting approval".into()),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
Ok(SubmissionResult::NeedApproval {
|
|
request_id,
|
|
tool_name,
|
|
description,
|
|
parameters,
|
|
})
|
|
}
|
|
Err(e) => {
|
|
thread.fail_turn(e.to_string());
|
|
// User message already persisted at turn start; nothing else to save
|
|
Ok(SubmissionResult::error(e.to_string()))
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Persist the user message to the DB at turn start (before the agentic loop).
|
|
///
|
|
/// This ensures the user message is durable even if the process crashes
|
|
/// mid-response. Call this right after `thread.start_turn()`.
|
|
pub(super) async fn persist_user_message(
|
|
&self,
|
|
thread_id: Uuid,
|
|
user_id: &str,
|
|
user_input: &str,
|
|
) {
|
|
let store = match self.store() {
|
|
Some(s) => Arc::clone(s),
|
|
None => return,
|
|
};
|
|
|
|
if let Err(e) = store
|
|
.ensure_conversation(thread_id, "gateway", user_id, None)
|
|
.await
|
|
{
|
|
tracing::warn!("Failed to ensure conversation {}: {}", thread_id, e);
|
|
return;
|
|
}
|
|
|
|
if let Err(e) = store
|
|
.add_conversation_message(thread_id, "user", user_input)
|
|
.await
|
|
{
|
|
tracing::warn!("Failed to persist user message: {}", e);
|
|
}
|
|
}
|
|
|
|
/// Persist the assistant response to the DB after the agentic loop completes.
|
|
///
|
|
/// Re-ensures the conversation row exists so that assistant responses are
|
|
/// still persisted even if `persist_user_message` failed transiently at
|
|
/// turn start (e.g. a brief DB blip that resolved before response time).
|
|
pub(super) async fn persist_assistant_response(
|
|
&self,
|
|
thread_id: Uuid,
|
|
user_id: &str,
|
|
response: &str,
|
|
) {
|
|
let store = match self.store() {
|
|
Some(s) => Arc::clone(s),
|
|
None => return,
|
|
};
|
|
|
|
if let Err(e) = store
|
|
.ensure_conversation(thread_id, "gateway", user_id, None)
|
|
.await
|
|
{
|
|
tracing::warn!("Failed to ensure conversation {}: {}", thread_id, e);
|
|
return;
|
|
}
|
|
|
|
if let Err(e) = store
|
|
.add_conversation_message(thread_id, "assistant", response)
|
|
.await
|
|
{
|
|
tracing::warn!("Failed to persist assistant message: {}", e);
|
|
}
|
|
}
|
|
|
|
/// Persist tool call summaries to the DB as a `role="tool_calls"` message.
|
|
///
|
|
/// Stored between the user and assistant messages so that
|
|
/// `build_turns_from_db_messages` can reconstruct the tool call history.
|
|
/// Content is a JSON array of tool call summaries.
|
|
pub(super) async fn persist_tool_calls(
|
|
&self,
|
|
thread_id: Uuid,
|
|
user_id: &str,
|
|
tool_calls: &[crate::agent::session::TurnToolCall],
|
|
) {
|
|
if tool_calls.is_empty() {
|
|
return;
|
|
}
|
|
|
|
let store = match self.store() {
|
|
Some(s) => Arc::clone(s),
|
|
None => return,
|
|
};
|
|
|
|
let summaries: Vec<serde_json::Value> = tool_calls
|
|
.iter()
|
|
.map(|tc| {
|
|
let mut obj = serde_json::json!({ "name": tc.name });
|
|
if let Some(ref result) = tc.result {
|
|
let preview = match result {
|
|
serde_json::Value::String(s) => truncate_preview(s, 500),
|
|
other => truncate_preview(&other.to_string(), 500),
|
|
};
|
|
obj["result_preview"] = serde_json::Value::String(preview);
|
|
}
|
|
if let Some(ref error) = tc.error {
|
|
obj["error"] = serde_json::Value::String(truncate_preview(error, 200));
|
|
}
|
|
obj
|
|
})
|
|
.collect();
|
|
|
|
let content = match serde_json::to_string(&summaries) {
|
|
Ok(c) => c,
|
|
Err(e) => {
|
|
tracing::warn!("Failed to serialize tool calls: {}", e);
|
|
return;
|
|
}
|
|
};
|
|
|
|
if let Err(e) = store
|
|
.ensure_conversation(thread_id, "gateway", user_id, None)
|
|
.await
|
|
{
|
|
tracing::warn!("Failed to ensure conversation {}: {}", thread_id, e);
|
|
return;
|
|
}
|
|
|
|
if let Err(e) = store
|
|
.add_conversation_message(thread_id, "tool_calls", &content)
|
|
.await
|
|
{
|
|
tracing::warn!("Failed to persist tool calls: {}", e);
|
|
}
|
|
}
|
|
|
|
pub(super) async fn process_undo(
|
|
&self,
|
|
session: Arc<Mutex<Session>>,
|
|
thread_id: Uuid,
|
|
) -> Result<SubmissionResult, Error> {
|
|
let undo_mgr = self.session_manager.get_undo_manager(thread_id).await;
|
|
let mut mgr = undo_mgr.lock().await;
|
|
|
|
if !mgr.can_undo() {
|
|
return Ok(SubmissionResult::ok_with_message("Nothing to undo."));
|
|
}
|
|
|
|
let mut sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get_mut(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
|
|
// Save current state to redo, get previous checkpoint
|
|
let current_messages = thread.messages();
|
|
let current_turn = thread.turn_number();
|
|
|
|
if let Some(checkpoint) = mgr.undo(current_turn, current_messages) {
|
|
// Extract values before consuming the reference
|
|
let turn_number = checkpoint.turn_number;
|
|
let messages = checkpoint.messages.clone();
|
|
let undo_count = mgr.undo_count();
|
|
// Restore thread from checkpoint
|
|
thread.restore_from_messages(messages);
|
|
Ok(SubmissionResult::ok_with_message(format!(
|
|
"Undone to turn {}. {} undo(s) remaining.",
|
|
turn_number, undo_count
|
|
)))
|
|
} else {
|
|
Ok(SubmissionResult::error("Undo failed."))
|
|
}
|
|
}
|
|
|
|
pub(super) async fn process_redo(
|
|
&self,
|
|
session: Arc<Mutex<Session>>,
|
|
thread_id: Uuid,
|
|
) -> Result<SubmissionResult, Error> {
|
|
let undo_mgr = self.session_manager.get_undo_manager(thread_id).await;
|
|
let mut mgr = undo_mgr.lock().await;
|
|
|
|
if !mgr.can_redo() {
|
|
return Ok(SubmissionResult::ok_with_message("Nothing to redo."));
|
|
}
|
|
|
|
let mut sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get_mut(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
|
|
let current_messages = thread.messages();
|
|
let current_turn = thread.turn_number();
|
|
|
|
if let Some(checkpoint) = mgr.redo(current_turn, current_messages) {
|
|
thread.restore_from_messages(checkpoint.messages);
|
|
Ok(SubmissionResult::ok_with_message(format!(
|
|
"Redone to turn {}.",
|
|
checkpoint.turn_number
|
|
)))
|
|
} else {
|
|
Ok(SubmissionResult::error("Redo failed."))
|
|
}
|
|
}
|
|
|
|
pub(super) async fn process_interrupt(
|
|
&self,
|
|
session: Arc<Mutex<Session>>,
|
|
thread_id: Uuid,
|
|
) -> Result<SubmissionResult, Error> {
|
|
let mut sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get_mut(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
|
|
match thread.state {
|
|
ThreadState::Processing | ThreadState::AwaitingApproval => {
|
|
thread.interrupt();
|
|
Ok(SubmissionResult::ok_with_message("Interrupted."))
|
|
}
|
|
_ => Ok(SubmissionResult::ok_with_message("Nothing to interrupt.")),
|
|
}
|
|
}
|
|
|
|
pub(super) async fn process_compact(
|
|
&self,
|
|
session: Arc<Mutex<Session>>,
|
|
thread_id: Uuid,
|
|
) -> Result<SubmissionResult, Error> {
|
|
let mut sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get_mut(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
|
|
let messages = thread.messages();
|
|
let usage = self.context_monitor.usage_percent(&messages);
|
|
let strategy = self
|
|
.context_monitor
|
|
.suggest_compaction(&messages)
|
|
.unwrap_or(
|
|
crate::agent::context_monitor::CompactionStrategy::Summarize { keep_recent: 5 },
|
|
);
|
|
|
|
let compactor = ContextCompactor::new(self.llm().clone(), self.safety().clone());
|
|
match compactor
|
|
.compact(thread, strategy, self.workspace().map(|w| w.as_ref()))
|
|
.await
|
|
{
|
|
Ok(result) => {
|
|
let mut msg = format!(
|
|
"Compacted: {} turns removed, {} → {} tokens (was {:.1}% full)",
|
|
result.turns_removed, result.tokens_before, result.tokens_after, usage
|
|
);
|
|
if result.summary_written {
|
|
msg.push_str(", summary saved to workspace");
|
|
}
|
|
Ok(SubmissionResult::ok_with_message(msg))
|
|
}
|
|
Err(e) => Ok(SubmissionResult::error(format!("Compaction failed: {}", e))),
|
|
}
|
|
}
|
|
|
|
pub(super) async fn process_clear(
|
|
&self,
|
|
session: Arc<Mutex<Session>>,
|
|
thread_id: Uuid,
|
|
) -> Result<SubmissionResult, Error> {
|
|
let mut sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get_mut(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
thread.turns.clear();
|
|
thread.state = ThreadState::Idle;
|
|
|
|
// Clear undo history too
|
|
let undo_mgr = self.session_manager.get_undo_manager(thread_id).await;
|
|
undo_mgr.lock().await.clear();
|
|
|
|
Ok(SubmissionResult::ok_with_message("Thread cleared."))
|
|
}
|
|
|
|
/// Process an approval or rejection of a pending tool execution.
|
|
pub(super) async fn process_approval(
|
|
&self,
|
|
message: &IncomingMessage,
|
|
session: Arc<Mutex<Session>>,
|
|
thread_id: Uuid,
|
|
request_id: Option<Uuid>,
|
|
approved: bool,
|
|
always: bool,
|
|
) -> Result<SubmissionResult, Error> {
|
|
// Get pending approval for this thread
|
|
let pending = {
|
|
let mut sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get_mut(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
|
|
if thread.state != ThreadState::AwaitingApproval {
|
|
// Stale or duplicate approval (tool already executed) — silently ignore.
|
|
tracing::debug!(
|
|
%thread_id,
|
|
state = ?thread.state,
|
|
"Ignoring stale approval: thread not in AwaitingApproval state"
|
|
);
|
|
return Ok(SubmissionResult::ok_with_message(""));
|
|
}
|
|
|
|
thread.take_pending_approval()
|
|
};
|
|
|
|
let pending = match pending {
|
|
Some(p) => p,
|
|
None => {
|
|
tracing::debug!(
|
|
%thread_id,
|
|
"Ignoring stale approval: no pending approval found"
|
|
);
|
|
return Ok(SubmissionResult::ok_with_message(""));
|
|
}
|
|
};
|
|
|
|
// Verify request ID if provided
|
|
if let Some(req_id) = request_id
|
|
&& req_id != pending.request_id
|
|
{
|
|
// Put it back and return error
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id) {
|
|
thread.await_approval(pending);
|
|
}
|
|
return Ok(SubmissionResult::error(
|
|
"Request ID mismatch. Use the correct request ID.",
|
|
));
|
|
}
|
|
|
|
if approved {
|
|
// If always, add to auto-approved set
|
|
if always {
|
|
let mut sess = session.lock().await;
|
|
sess.auto_approve_tool(&pending.tool_name);
|
|
tracing::info!(
|
|
"Auto-approved tool '{}' for session {}",
|
|
pending.tool_name,
|
|
sess.id
|
|
);
|
|
}
|
|
|
|
// Reset thread state to processing
|
|
{
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id) {
|
|
thread.state = ThreadState::Processing;
|
|
}
|
|
}
|
|
|
|
// Execute the approved tool and continue the loop
|
|
let mut job_ctx =
|
|
JobContext::with_user(&message.user_id, "chat", "Interactive chat session");
|
|
job_ctx.http_interceptor = self.deps.http_interceptor.clone();
|
|
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::ToolStarted {
|
|
name: pending.tool_name.clone(),
|
|
},
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
|
|
let tool_result = self
|
|
.execute_chat_tool(&pending.tool_name, &pending.parameters, &job_ctx)
|
|
.await;
|
|
|
|
let tool_ref = self.tools().get(&pending.tool_name).await;
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::tool_completed(
|
|
pending.tool_name.clone(),
|
|
&tool_result,
|
|
&pending.display_parameters,
|
|
tool_ref.as_deref(),
|
|
),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
|
|
if let Ok(ref output) = tool_result
|
|
&& !output.is_empty()
|
|
{
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::ToolResult {
|
|
name: pending.tool_name.clone(),
|
|
preview: output.clone(),
|
|
},
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
}
|
|
|
|
// Build context including the tool result
|
|
let mut context_messages = pending.context_messages;
|
|
let deferred_tool_calls = pending.deferred_tool_calls;
|
|
|
|
// Record result in thread
|
|
{
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id)
|
|
&& let Some(turn) = thread.last_turn_mut()
|
|
{
|
|
match &tool_result {
|
|
Ok(output) => {
|
|
turn.record_tool_result(serde_json::json!(output));
|
|
}
|
|
Err(e) => {
|
|
turn.record_tool_error(e.to_string());
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// If tool_auth returned awaiting_token, enter auth mode and
|
|
// return instructions directly (skip agentic loop continuation).
|
|
if let Some((ext_name, instructions)) =
|
|
check_auth_required(&pending.tool_name, &tool_result)
|
|
{
|
|
self.handle_auth_intercept(
|
|
&session,
|
|
thread_id,
|
|
message,
|
|
&tool_result,
|
|
ext_name,
|
|
instructions.clone(),
|
|
)
|
|
.await;
|
|
return Ok(SubmissionResult::response(instructions));
|
|
}
|
|
|
|
// Add tool result to context
|
|
let result_content = match tool_result {
|
|
Ok(output) => {
|
|
let sanitized = self
|
|
.safety()
|
|
.sanitize_tool_output(&pending.tool_name, &output);
|
|
self.safety().wrap_for_llm(
|
|
&pending.tool_name,
|
|
&sanitized.content,
|
|
sanitized.was_modified,
|
|
)
|
|
}
|
|
Err(e) => format!("Error: {}", e),
|
|
};
|
|
|
|
context_messages.push(ChatMessage::tool_result(
|
|
&pending.tool_call_id,
|
|
&pending.tool_name,
|
|
result_content,
|
|
));
|
|
|
|
// Replay deferred tool calls from the same assistant message so
|
|
// every tool_use ID gets a matching tool_result before the next
|
|
// LLM call.
|
|
if !deferred_tool_calls.is_empty() {
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::Thinking(format!(
|
|
"Executing {} deferred tool(s)...",
|
|
deferred_tool_calls.len()
|
|
)),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
}
|
|
|
|
// === Phase 1: Preflight (sequential) ===
|
|
// Walk deferred tools checking approval. Collect runnable
|
|
// tools; stop at the first that needs approval.
|
|
let mut runnable: Vec<crate::llm::ToolCall> = Vec::new();
|
|
let mut approval_needed: Option<(
|
|
usize,
|
|
crate::llm::ToolCall,
|
|
Arc<dyn crate::tools::Tool>,
|
|
)> = None;
|
|
|
|
for (idx, tc) in deferred_tool_calls.iter().enumerate() {
|
|
if let Some(tool) = self.tools().get(&tc.name).await {
|
|
use crate::tools::ApprovalRequirement;
|
|
let needs_approval = match tool.requires_approval(&tc.arguments) {
|
|
ApprovalRequirement::Never => false,
|
|
ApprovalRequirement::UnlessAutoApproved => {
|
|
let sess = session.lock().await;
|
|
!sess.is_tool_auto_approved(&tc.name)
|
|
}
|
|
ApprovalRequirement::Always => true,
|
|
};
|
|
|
|
if needs_approval {
|
|
approval_needed = Some((idx, tc.clone(), tool));
|
|
break; // remaining tools stay deferred
|
|
}
|
|
}
|
|
|
|
runnable.push(tc.clone());
|
|
}
|
|
|
|
// === Phase 2: Parallel execution ===
|
|
let exec_results: Vec<(crate::llm::ToolCall, Result<String, Error>)> = if runnable.len()
|
|
<= 1
|
|
{
|
|
// Single tool (or none): execute inline
|
|
let mut results = Vec::new();
|
|
for tc in &runnable {
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::ToolStarted {
|
|
name: tc.name.clone(),
|
|
},
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
|
|
let result = self
|
|
.execute_chat_tool(&tc.name, &tc.arguments, &job_ctx)
|
|
.await;
|
|
|
|
let deferred_tool = self.tools().get(&tc.name).await;
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::tool_completed(
|
|
tc.name.clone(),
|
|
&result,
|
|
&tc.arguments,
|
|
deferred_tool.as_deref(),
|
|
),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
|
|
results.push((tc.clone(), result));
|
|
}
|
|
results
|
|
} else {
|
|
// Multiple tools: execute in parallel via JoinSet
|
|
let mut join_set = JoinSet::new();
|
|
let runnable_count = runnable.len();
|
|
|
|
for (spawn_idx, tc) in runnable.iter().enumerate() {
|
|
let tools = self.tools().clone();
|
|
let safety = self.safety().clone();
|
|
let channels = self.channels.clone();
|
|
let job_ctx = job_ctx.clone();
|
|
let tc = tc.clone();
|
|
let channel = message.channel.clone();
|
|
let metadata = message.metadata.clone();
|
|
|
|
join_set.spawn(async move {
|
|
let _ = channels
|
|
.send_status(
|
|
&channel,
|
|
StatusUpdate::ToolStarted {
|
|
name: tc.name.clone(),
|
|
},
|
|
&metadata,
|
|
)
|
|
.await;
|
|
|
|
let result = execute_chat_tool_standalone(
|
|
&tools,
|
|
&safety,
|
|
&tc.name,
|
|
&tc.arguments,
|
|
&job_ctx,
|
|
)
|
|
.await;
|
|
|
|
let par_tool = tools.get(&tc.name).await;
|
|
let _ = channels
|
|
.send_status(
|
|
&channel,
|
|
StatusUpdate::tool_completed(
|
|
tc.name.clone(),
|
|
&result,
|
|
&tc.arguments,
|
|
par_tool.as_deref(),
|
|
),
|
|
&metadata,
|
|
)
|
|
.await;
|
|
|
|
(spawn_idx, tc, result)
|
|
});
|
|
}
|
|
|
|
// Collect and reorder by original index
|
|
let mut ordered: Vec<Option<(crate::llm::ToolCall, Result<String, Error>)>> =
|
|
(0..runnable_count).map(|_| None).collect();
|
|
while let Some(join_result) = join_set.join_next().await {
|
|
match join_result {
|
|
Ok((idx, tc, result)) => {
|
|
ordered[idx] = Some((tc, result));
|
|
}
|
|
Err(e) => {
|
|
if e.is_panic() {
|
|
tracing::error!("Deferred tool execution task panicked: {}", e);
|
|
} else {
|
|
tracing::error!("Deferred tool execution task cancelled: {}", e);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// Fill panicked slots with error results
|
|
ordered
|
|
.into_iter()
|
|
.enumerate()
|
|
.map(|(i, opt)| {
|
|
opt.unwrap_or_else(|| {
|
|
let tc = runnable[i].clone();
|
|
let err: Error = crate::error::ToolError::ExecutionFailed {
|
|
name: tc.name.clone(),
|
|
reason: "Task failed during execution".to_string(),
|
|
}
|
|
.into();
|
|
(tc, Err(err))
|
|
})
|
|
})
|
|
.collect()
|
|
};
|
|
|
|
// === Phase 3: Post-flight (sequential, in original order) ===
|
|
// Process all results before any conditional return so every
|
|
// tool result is recorded in the session audit trail.
|
|
let mut deferred_auth: Option<String> = None;
|
|
|
|
for (tc, deferred_result) in exec_results {
|
|
if let Ok(ref output) = deferred_result
|
|
&& !output.is_empty()
|
|
{
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::ToolResult {
|
|
name: tc.name.clone(),
|
|
preview: output.clone(),
|
|
},
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
}
|
|
|
|
// Record in thread
|
|
{
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id)
|
|
&& let Some(turn) = thread.last_turn_mut()
|
|
{
|
|
match &deferred_result {
|
|
Ok(output) => turn.record_tool_result(serde_json::json!(output)),
|
|
Err(e) => turn.record_tool_error(e.to_string()),
|
|
}
|
|
}
|
|
}
|
|
|
|
// Auth detection — defer return until all results are recorded
|
|
if deferred_auth.is_none()
|
|
&& let Some((ext_name, instructions)) =
|
|
check_auth_required(&tc.name, &deferred_result)
|
|
{
|
|
self.handle_auth_intercept(
|
|
&session,
|
|
thread_id,
|
|
message,
|
|
&deferred_result,
|
|
ext_name,
|
|
instructions.clone(),
|
|
)
|
|
.await;
|
|
deferred_auth = Some(instructions);
|
|
}
|
|
|
|
let deferred_content = match deferred_result {
|
|
Ok(output) => {
|
|
let sanitized = self.safety().sanitize_tool_output(&tc.name, &output);
|
|
self.safety().wrap_for_llm(
|
|
&tc.name,
|
|
&sanitized.content,
|
|
sanitized.was_modified,
|
|
)
|
|
}
|
|
Err(e) => format!("Error: {}", e),
|
|
};
|
|
|
|
context_messages.push(ChatMessage::tool_result(&tc.id, &tc.name, deferred_content));
|
|
}
|
|
|
|
// Return auth response after all results are recorded
|
|
if let Some(instructions) = deferred_auth {
|
|
return Ok(SubmissionResult::response(instructions));
|
|
}
|
|
|
|
// Handle approval if a tool needed it
|
|
if let Some((approval_idx, tc, tool)) = approval_needed {
|
|
let new_pending = PendingApproval {
|
|
request_id: Uuid::new_v4(),
|
|
tool_name: tc.name.clone(),
|
|
parameters: tc.arguments.clone(),
|
|
display_parameters: redact_params(&tc.arguments, tool.sensitive_params()),
|
|
description: tool.description().to_string(),
|
|
tool_call_id: tc.id.clone(),
|
|
context_messages: context_messages.clone(),
|
|
deferred_tool_calls: deferred_tool_calls[approval_idx + 1..].to_vec(),
|
|
};
|
|
|
|
let request_id = new_pending.request_id;
|
|
let tool_name = new_pending.tool_name.clone();
|
|
let description = new_pending.description.clone();
|
|
let parameters = new_pending.display_parameters.clone();
|
|
|
|
{
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id) {
|
|
thread.await_approval(new_pending);
|
|
}
|
|
}
|
|
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::Status("Awaiting approval".into()),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
|
|
return Ok(SubmissionResult::NeedApproval {
|
|
request_id,
|
|
tool_name,
|
|
description,
|
|
parameters,
|
|
});
|
|
}
|
|
|
|
// Continue the agentic loop (a tool was already executed this turn)
|
|
let result = self
|
|
.run_agentic_loop(message, session.clone(), thread_id, context_messages)
|
|
.await;
|
|
|
|
// Handle the result
|
|
let mut sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get_mut(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
|
|
match result {
|
|
Ok(AgenticLoopResult::Response(response)) => {
|
|
thread.complete_turn(&response);
|
|
let tool_calls = thread
|
|
.turns
|
|
.last()
|
|
.map(|t| t.tool_calls.clone())
|
|
.unwrap_or_default();
|
|
// User message already persisted at turn start; save tool calls then assistant response
|
|
self.persist_tool_calls(thread_id, &message.user_id, &tool_calls)
|
|
.await;
|
|
self.persist_assistant_response(thread_id, &message.user_id, &response)
|
|
.await;
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::Status("Done".into()),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
Ok(SubmissionResult::response(response))
|
|
}
|
|
Ok(AgenticLoopResult::NeedApproval {
|
|
pending: new_pending,
|
|
}) => {
|
|
let request_id = new_pending.request_id;
|
|
let tool_name = new_pending.tool_name.clone();
|
|
let description = new_pending.description.clone();
|
|
let parameters = new_pending.display_parameters.clone();
|
|
thread.await_approval(new_pending);
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::Status("Awaiting approval".into()),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
Ok(SubmissionResult::NeedApproval {
|
|
request_id,
|
|
tool_name,
|
|
description,
|
|
parameters,
|
|
})
|
|
}
|
|
Err(e) => {
|
|
thread.fail_turn(e.to_string());
|
|
// User message already persisted at turn start
|
|
Ok(SubmissionResult::error(e.to_string()))
|
|
}
|
|
}
|
|
} else {
|
|
// Rejected - complete the turn with a rejection message and persist
|
|
let rejection = format!(
|
|
"Tool '{}' was rejected. The agent will not execute this tool.\n\n\
|
|
You can continue the conversation or try a different approach.",
|
|
pending.tool_name
|
|
);
|
|
{
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id) {
|
|
thread.clear_pending_approval();
|
|
thread.complete_turn(&rejection);
|
|
// User message already persisted at turn start; save rejection response
|
|
self.persist_assistant_response(thread_id, &message.user_id, &rejection)
|
|
.await;
|
|
}
|
|
}
|
|
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::Status("Rejected".into()),
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
|
|
Ok(SubmissionResult::response(rejection))
|
|
}
|
|
}
|
|
|
|
/// Handle an auth-required result from a tool execution.
|
|
///
|
|
/// Enters auth mode on the thread, completes + persists the turn,
|
|
/// and sends the AuthRequired status to the channel.
|
|
/// Returns the instructions string for the caller to wrap in a response.
|
|
async fn handle_auth_intercept(
|
|
&self,
|
|
session: &Arc<Mutex<Session>>,
|
|
thread_id: Uuid,
|
|
message: &IncomingMessage,
|
|
tool_result: &Result<String, Error>,
|
|
ext_name: String,
|
|
instructions: String,
|
|
) {
|
|
let auth_data = parse_auth_result(tool_result);
|
|
{
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id) {
|
|
thread.enter_auth_mode(ext_name.clone());
|
|
thread.complete_turn(&instructions);
|
|
// User message already persisted at turn start; save auth instructions
|
|
self.persist_assistant_response(thread_id, &message.user_id, &instructions)
|
|
.await;
|
|
}
|
|
}
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::AuthRequired {
|
|
extension_name: ext_name,
|
|
instructions: Some(instructions.clone()),
|
|
auth_url: auth_data.auth_url,
|
|
setup_url: auth_data.setup_url,
|
|
},
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
}
|
|
|
|
/// Handle an auth token submitted while the thread is in auth mode.
|
|
///
|
|
/// The token goes directly to the extension manager's credential store,
|
|
/// completely bypassing logging, turn creation, history, and compaction.
|
|
pub(super) async fn process_auth_token(
|
|
&self,
|
|
message: &IncomingMessage,
|
|
pending: &crate::agent::session::PendingAuth,
|
|
token: &str,
|
|
session: Arc<Mutex<Session>>,
|
|
thread_id: Uuid,
|
|
) -> Result<Option<String>, Error> {
|
|
let token = token.trim();
|
|
|
|
// Clear auth mode regardless of outcome
|
|
{
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id) {
|
|
thread.pending_auth = None;
|
|
}
|
|
}
|
|
|
|
let ext_mgr = match self.deps.extension_manager.as_ref() {
|
|
Some(mgr) => mgr,
|
|
None => return Ok(Some("Extension manager not available.".to_string())),
|
|
};
|
|
|
|
match ext_mgr.auth(&pending.extension_name, Some(token)).await {
|
|
Ok(result) if result.is_authenticated() => {
|
|
tracing::info!(
|
|
"Extension '{}' authenticated via auth mode",
|
|
pending.extension_name
|
|
);
|
|
|
|
// Auto-activate so tools are available immediately after auth
|
|
match ext_mgr.activate(&pending.extension_name).await {
|
|
Ok(activate_result) => {
|
|
let tool_count = activate_result.tools_loaded.len();
|
|
let tool_list = if activate_result.tools_loaded.is_empty() {
|
|
String::new()
|
|
} else {
|
|
format!("\n\nTools: {}", activate_result.tools_loaded.join(", "))
|
|
};
|
|
let msg = format!(
|
|
"{} authenticated and activated ({} tools loaded).{}",
|
|
pending.extension_name, tool_count, tool_list
|
|
);
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::AuthCompleted {
|
|
extension_name: pending.extension_name.clone(),
|
|
success: true,
|
|
message: msg.clone(),
|
|
},
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
Ok(Some(msg))
|
|
}
|
|
Err(e) => {
|
|
tracing::warn!(
|
|
"Extension '{}' authenticated but activation failed: {}",
|
|
pending.extension_name,
|
|
e
|
|
);
|
|
let msg = format!(
|
|
"{} authenticated successfully, but activation failed: {}. \
|
|
Try activating manually.",
|
|
pending.extension_name, e
|
|
);
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::AuthCompleted {
|
|
extension_name: pending.extension_name.clone(),
|
|
success: true,
|
|
message: msg.clone(),
|
|
},
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
Ok(Some(msg))
|
|
}
|
|
}
|
|
}
|
|
Ok(result) => {
|
|
// Invalid token, re-enter auth mode
|
|
{
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id) {
|
|
thread.enter_auth_mode(pending.extension_name.clone());
|
|
}
|
|
}
|
|
let msg = result
|
|
.instructions()
|
|
.map(String::from)
|
|
.unwrap_or_else(|| "Invalid token. Please try again.".to_string());
|
|
// Re-emit AuthRequired so web UI re-shows the card
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::AuthRequired {
|
|
extension_name: pending.extension_name.clone(),
|
|
instructions: Some(msg.clone()),
|
|
auth_url: result.auth_url().map(String::from),
|
|
setup_url: result.setup_url().map(String::from),
|
|
},
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
Ok(Some(msg))
|
|
}
|
|
Err(e) => {
|
|
let msg = format!(
|
|
"Authentication failed for {}: {}",
|
|
pending.extension_name, e
|
|
);
|
|
let _ = self
|
|
.channels
|
|
.send_status(
|
|
&message.channel,
|
|
StatusUpdate::AuthCompleted {
|
|
extension_name: pending.extension_name.clone(),
|
|
success: false,
|
|
message: msg.clone(),
|
|
},
|
|
&message.metadata,
|
|
)
|
|
.await;
|
|
Ok(Some(msg))
|
|
}
|
|
}
|
|
}
|
|
|
|
pub(super) async fn process_new_thread(
|
|
&self,
|
|
message: &IncomingMessage,
|
|
) -> Result<SubmissionResult, Error> {
|
|
let session = self
|
|
.session_manager
|
|
.get_or_create_session(&message.user_id)
|
|
.await;
|
|
let mut sess = session.lock().await;
|
|
let thread = sess.create_thread();
|
|
let thread_id = thread.id;
|
|
Ok(SubmissionResult::ok_with_message(format!(
|
|
"New thread: {}",
|
|
thread_id
|
|
)))
|
|
}
|
|
|
|
pub(super) async fn process_switch_thread(
|
|
&self,
|
|
message: &IncomingMessage,
|
|
target_thread_id: Uuid,
|
|
) -> Result<SubmissionResult, Error> {
|
|
let session = self
|
|
.session_manager
|
|
.get_or_create_session(&message.user_id)
|
|
.await;
|
|
let mut sess = session.lock().await;
|
|
|
|
if sess.switch_thread(target_thread_id) {
|
|
Ok(SubmissionResult::ok_with_message(format!(
|
|
"Switched to thread {}",
|
|
target_thread_id
|
|
)))
|
|
} else {
|
|
Ok(SubmissionResult::error("Thread not found."))
|
|
}
|
|
}
|
|
|
|
pub(super) async fn process_resume(
|
|
&self,
|
|
session: Arc<Mutex<Session>>,
|
|
thread_id: Uuid,
|
|
checkpoint_id: Uuid,
|
|
) -> Result<SubmissionResult, Error> {
|
|
let undo_mgr = self.session_manager.get_undo_manager(thread_id).await;
|
|
let mut mgr = undo_mgr.lock().await;
|
|
|
|
if let Some(checkpoint) = mgr.restore(checkpoint_id) {
|
|
let mut sess = session.lock().await;
|
|
let thread = sess
|
|
.threads
|
|
.get_mut(&thread_id)
|
|
.ok_or_else(|| Error::from(crate::error::JobError::NotFound { id: thread_id }))?;
|
|
thread.restore_from_messages(checkpoint.messages);
|
|
Ok(SubmissionResult::ok_with_message(format!(
|
|
"Resumed from checkpoint: {}",
|
|
checkpoint.description
|
|
)))
|
|
} else {
|
|
Ok(SubmissionResult::error("Checkpoint not found."))
|
|
}
|
|
}
|
|
}
|