Compare commits

..
Author SHA1 Message Date
8a320ae9db fix(routines): complete full_job execution reliability overhaul (#1650)
* fix(routines): persist full LLM transcript and remove sandbox gate for full_job

Routine execution output was invisible — routine_fire returned a one-liner,
routine_history had no actual output, and the conversation thread contained
only a summary. Full-job routines also hard-failed without Docker.

Three fixes:

1. **Full transcript persistence**: execute_lightweight now persists every
   message (prompt, LLM responses, tool calls with params, tool results) to
   the routine's conversation thread as it executes, not just a summary
   after the fact.

2. **Routine output visibility**: routine_history includes conversation_id
   and recent_output messages. routine_fire tells the user to check
   routine_history. Web detail page has a "View Execution Thread" button
   that navigates to the chat tab. ROUTINE_OK stores "No issues found"
   instead of None. Full-job summary pulls actual job output instead of
   generic "Job X finished".

3. **Remove SandboxReadiness gate**: full_job routines dispatch through the
   scheduler like regular /job commands — no Docker required. The
   SandboxReadiness enum is removed entirely.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* style: apply cargo fmt

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(worker): treat AutonomousUnavailable tool errors as recoverable

The job worker crashed the entire job when a tool was denied for
autonomous execution (e.g. secret_list). The error was already recorded
in reason_ctx for the LLM to see, but process_tool_result_job returned
Err which propagated through the agentic loop and terminated the job.

Now all tool errors (including AutonomousUnavailable) return Ok,
letting the LLM see the denial and try a different approach.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(llm): sanitize tool names for OpenAI Codex Responses API

The Codex API requires tool names to match `^[a-zA-Z0-9_-]+$` but
MCP/extension tools can have dots in their names (e.g. `mcp.server.tool`).
This caused HTTP 400 errors when the job worker sent tool calls back
to the LLM.

Sanitize tool names in both `convert_tool_definition` and
`convert_message` (function_call items) by replacing invalid characters
with underscores.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(routines): inject execution context into full_job description [skip-regression-check]

When a full_job routine dispatches a job, the LLM had no context that
it was already executing inside a routine. It wasted iterations on
infrastructure (discovering tools, creating routines, setting up auth)
instead of doing the actual work.

Prepend a clear directive to the job description telling the LLM that
tools and the routine are already configured, and to execute the task
directly.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(mcp): auto-refresh expired OAuth tokens on access [skip-regression-check]

When IronClaw restarts, MCP servers fail with "Secret has expired"
because get_access_token() checks token expiry locally and returns an
error before any HTTP request is made — so the existing 401-retry
refresh logic never triggers.

Now get_access_token() catches SecretError::Expired and automatically
calls refresh_access_token() using the stored refresh token. If the
refresh succeeds, the new token is returned transparently. If it fails,
the error message includes both the expiry and the refresh failure.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(mcp): align refresh token naming and set expiry on stored tokens

Two bugs prevented MCP OAuth token auto-refresh on restart:

1. Naming mismatch: the hosted OAuth flow stored the refresh token as
   `{token_secret_name}_refresh_token` (e.g. `mcp_notion_access_token_refresh_token`)
   but `McpServerConfig::refresh_token_secret_name()` returned
   `mcp_notion_refresh_token`. The refresh token was there but unfindable.

2. Missing expiry: `store_tokens` in auth.rs never called `with_expiry()`
   even though `AccessToken::expires_in` was available. Combined with the
   fix from the previous commit (auto-refresh on Expired), tokens stored
   via the MCP auth flow will now also trigger refresh correctly.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(web): show activity and transitions for agent jobs in job detail [skip-regression-check]

The job events endpoint only checked sandbox jobs for ownership,
returning 404 for agent jobs dispatched from routines. The detail
handler also returned empty transitions for agent jobs.

- events handler: fall back to agent job ownership check
- detail handler: populate transitions from job's state history

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* feat(routines): expose max_iterations for full_job routines (default 25)

The max_iterations parameter was hardcoded to 10 and not configurable
via routine_create or routine_update, causing complex tasks to hit the
iteration cap.

- Add max_iterations to full_job execution schema (1-200, default 25)
- Thread it through parse → build → RoutineAction
- Support updating via routine_update
- Raise default from 10 to 25

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(routines): break self-dialogue loop after full_job plan execution

After plan execution, the completion-check Q&A ("Is the job complete?" /
"No, not complete...") was left in the message context, causing the
agentic loop to repeat the same analysis instead of calling tools.

Replace the stale dialogue with an action-oriented continuation prompt
that instructs the LLM to use tools for remaining work. Also strip
<suggestions> tags from all job output since they're only meaningful
for interactive chat sessions.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(repl): prevent test hang in single-message mode

In single-message mode, start() stored a clone of the mpsc sender in
self.msg_tx for approval injection. After the thread sent /quit and
exited, the stored clone kept the stream alive, so stream.next()
blocked forever in the test assertion that the stream ends.

Skip storing the sender in single-message mode since interactive
approval is not needed.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(jobs): treat text responses as final answer in agentic loop

When the LLM produces a non-empty text response with no tool intent
(already filtered by the nudge mechanism), it is the job's final
answer. Previously, handle_text_response only exited the loop if the
text matched rigid completion phrases like "job is complete". Natural
summaries like "Weekly review completed and saved to Notion" were
added to context and the loop continued, causing the LLM to restate
the same summary until max_iterations was hit.

Now any non-empty text response marks the job complete and stops the
loop, matching the chat dispatcher behavior.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* perf(tests): reduce skills catalog network failure test from 10s to 1s

The test_search_returns_error_on_network_failure test connects to an
unreachable RFC 5737 TEST-NET IP and waited for the full 10s production
REQUEST_TIMEOUT. Add with_url_and_timeout test helper and use a 1s
timeout instead. [skip-regression-check]

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(tools): accept 'message' as alias for 'content' in message tool

LLMs frequently call the message tool with {"message": "..."} instead
of {"content": "..."}. Fall back to the 'message' key when 'content'
is missing to avoid InvalidParameters errors during autonomous job
execution.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(tools): attach thread_id for gateway broadcast in message tool

When the message tool broadcasts to all channels (channel=null), it
sent an OutgoingResponse without a thread_id. The gateway silently
dropped these messages (returned Ok but never sent the SSE event),
so they appeared in repl but not in the web UI.

The thread_id was only populated when channel was explicitly "gateway".
Now it is always populated from notify_thread_id metadata, so
broadcast_all delivers to the gateway correctly.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(gateway): return error instead of silently dropping messages

Gateway broadcast() and respond() previously returned Ok(()) when
thread_id was missing, silently swallowing the message. Callers
(message tool, agent loop) believed delivery succeeded when it didn't.

Now returns ChannelError::MissingRoutingTarget so callers can detect
and report the failure. Four regression tests verify the contract:
respond/broadcast with and without thread_id.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix: resolve rebase conflicts with staging

Restore sandbox_readiness field removed by pre-rebase commits (staging
still uses it). Update repl test to match staging's single-message
behavior (no longer sends /quit). Add missing reasoning field to
ToolCall in codex test.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix(tools): log error when routine conversation lookup fails

The routine_history tool silently swallowed errors from
get_or_create_routine_conversation, returning empty output without
any diagnostic logging. Add tracing::warn so failures are visible
in logs. [skip-regression-check]

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix: address PR #1650 review comments

- E2E test: accept submitted/accepted as success states in job assertion
- TimeTool: remove operation from required schema (defaults to "now")
- jobs handler: log DB errors server-side, return generic message to client
- routines handler: use read-only find_routine_conversation on GET
- codex provider: reverse-map sanitized tool names so MCP tools resolve

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix: address zmanian review feedback on PR #1650

- MCP refresh token: fall back to legacy secret name (mcp_{name}_refresh_token)
  so existing users don't need to re-authenticate after the naming fix
- Job worker: replace fragile messages.pop() with truncate-to-saved-count
  to avoid maintenance hazard if message flow changes
- Document cost implications of max_iterations 10->25 default bump
- Revert Cargo.toml dist profile change (thin LTO comment, codegen-units=16)
  as it's unrelated to this PR

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

* fix: resolve rebase conflicts and address new Copilot comments

- Fix no_silent_drop tests for updated GatewayConfig (user_id moved to
  GatewayChannel::new second arg, user_tokens removed)
- Fix handle_text_response param name (_reason_ctx -> reason_ctx)
- Fix missing has_text_response field in test JobDelegate
- Propagate row.get errors in find_routine_conversation instead of
  unwrap_or_default
- Only fall back to legacy refresh token name on NotFound/Expired,
  propagate real errors (DB, decryption)

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>

---------

Co-authored-by: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 12:27:43 -07:00
46 changed files with 1270 additions and 1752 deletions
+3 -2
View File
@@ -191,9 +191,10 @@ HEARTBEAT_NOTIFY_CHANNEL=cli
HEARTBEAT_NOTIFY_USER=default
# Memory hygiene settings (automatic cleanup of stale workspace documents)
# Runs on each heartbeat tick; discovers cleanup targets from .config metadata
# Runs on each heartbeat tick; identity files (IDENTITY.md, SOUL.md) are never deleted
# MEMORY_HYGIENE_ENABLED=true
# MEMORY_HYGIENE_VERSION_KEEP_COUNT=50 # max versions to keep per document
# MEMORY_HYGIENE_DAILY_RETENTION_DAYS=30 # delete daily/ docs older than this many days
# MEMORY_HYGIENE_CONVERSATION_RETENTION_DAYS=7 # delete conversations/ docs older than this many days
# MEMORY_HYGIENE_CADENCE_HOURS=12 # minimum hours between cleanup passes
# Docker Sandbox
+1 -2
View File
@@ -252,8 +252,7 @@ strip = true # Remove debug symbols from release binaries
# The profile that 'cargo dist' will build with
[profile.dist]
inherits = "release"
lto = "fat" # Full cross-crate LTO (slow build, better codegen)
codegen-units = 1 # Single codegen unit for maximum optimization
lto = "thin"
# Config for 'dist'
[workspace.metadata.dist]
-23
View File
@@ -1,23 +0,0 @@
-- Document version history for workspace files.
-- Every content update saves the previous content as a version,
-- enabling rollback and audit trails.
CREATE TABLE memory_document_versions (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
document_id UUID NOT NULL REFERENCES memory_documents(id) ON DELETE CASCADE,
version INTEGER NOT NULL,
content TEXT NOT NULL,
content_hash TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
changed_by TEXT,
UNIQUE(document_id, version)
);
CREATE INDEX idx_doc_versions_lookup
ON memory_document_versions(document_id, version DESC);
-- GIN index on metadata for JSON path queries (used by hygiene to find
-- .config documents with hygiene.enabled). The metadata column already
-- exists (V1) but was never indexed.
CREATE INDEX idx_memory_documents_metadata
ON memory_documents USING GIN (metadata jsonb_path_ops);
+20
View File
@@ -1269,6 +1269,14 @@ pub(crate) fn extract_suggestions(text: &str) -> (String, Vec<String>) {
(cleaned, suggestions)
}
/// Remove `<suggestions>` tags from a response, returning only the cleaned text.
///
/// Convenience wrapper around [`extract_suggestions`] for callers that don't
/// need the parsed suggestion list (e.g. job worker, plan completion check).
pub(crate) fn strip_suggestions(text: &str) -> String {
extract_suggestions(text).0
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
@@ -2539,6 +2547,18 @@ mod tests {
assert_eq!(suggestions, vec!["ok"]); // safety: test
}
#[test]
fn test_strip_suggestions_removes_tags() {
let input = "The job is complete.\n<suggestions>[\"Check logs\"]</suggestions>";
assert_eq!(super::strip_suggestions(input), "The job is complete."); // safety: test
}
#[test]
fn test_strip_suggestions_no_tag_passthrough() {
let input = "Plain text without tags.";
assert_eq!(super::strip_suggestions(input), input); // safety: test
}
#[test]
fn test_tool_error_format_includes_tool_name() {
let tool_name = "http";
+4 -21
View File
@@ -276,8 +276,8 @@ impl HeartbeatRunner {
.await;
if report.had_work() {
tracing::info!(
directories_cleaned = ?report.directories_cleaned,
versions_pruned = report.versions_pruned,
daily_logs_deleted = report.daily_logs_deleted,
conversation_docs_deleted = report.conversation_docs_deleted,
"heartbeat: memory hygiene deleted stale documents"
);
}
@@ -590,23 +590,6 @@ pub fn spawn_multi_user_heartbeat(
let workspace = Arc::new(Workspace::new_with_db(user_id, Arc::clone(store.db())));
// Run memory hygiene per user (same as single-user heartbeat).
let hygiene_ws = Arc::clone(&workspace);
let hygiene_cfg = hygiene_config.clone();
let hygiene_user = user_id.clone();
tokio::spawn(async move {
let report =
crate::workspace::hygiene::run_if_due(&hygiene_ws, &hygiene_cfg).await;
if report.had_work() {
tracing::info!(
user_id = hygiene_user,
directories_cleaned = ?report.directories_cleaned,
versions_pruned = report.versions_pruned,
"multi-user heartbeat: memory hygiene deleted stale documents"
);
}
});
// Drain completed tasks to stay within the concurrency cap.
while join_set.len() >= MAX_CONCURRENT_HEARTBEATS {
if let Some(join_result) = join_set.join_next().await {
@@ -633,8 +616,8 @@ pub fn spawn_multi_user_heartbeat(
if report.had_work() {
tracing::info!(
user_id = uid,
directories_cleaned = ?report.directories_cleaned,
versions_pruned = report.versions_pruned,
daily_logs_deleted = report.daily_logs_deleted,
conversation_docs_deleted = report.conversation_docs_deleted,
"multi-user heartbeat: memory hygiene deleted stale documents"
);
}
+1
View File
@@ -36,6 +36,7 @@ pub(crate) use agent_loop::truncate_for_preview;
pub use agent_loop::{Agent, AgentDeps};
pub use compaction::{CompactionResult, ContextCompactor};
pub use context_monitor::{CompactionStrategy, ContextBreakdown, ContextMonitor};
pub(crate) use dispatcher::strip_suggestions;
pub use heartbeat::{
HeartbeatConfig, HeartbeatResult, HeartbeatRunner, spawn_heartbeat, spawn_multi_user_heartbeat,
};
+6 -1
View File
@@ -265,8 +265,13 @@ fn default_max_tokens() -> u32 {
4096
}
/// Default max agentic loop iterations for full_job routines.
///
/// Raised from 10 to 25 to accommodate multi-step tool chains that
/// stalled at the old cap. Worst-case LLM cost is 2.5x higher per run;
/// callers needing tighter budgets should set `max_iterations` explicitly.
fn default_max_iterations() -> u32 {
10
25
}
fn default_max_tool_rounds() -> u32 {
+13 -1
View File
@@ -1292,11 +1292,23 @@ async fn execute_full_job(
}
metadata["notify_user"] = serde_json::json!(&routine.notify.user);
// Prepend execution context so the LLM knows it's already inside a
// routine and should execute the task directly — not set up infrastructure.
let contextualized_description = format!(
"IMPORTANT: You are executing inside routine \"{routine_name}\". \
The routine and its schedule are already configured. \
Tools and credentials are already set up. \
Do NOT create routines, jobs, or try to discover/install/authenticate tools. \
Execute the task directly.\n\n{desc}",
routine_name = routine.name,
desc = execution.description,
);
let job_id = scheduler
.dispatch_job(
&routine.user_id,
execution.title,
execution.description,
&contextualized_description,
Some(metadata),
)
.await
+13 -24
View File
@@ -492,10 +492,12 @@ impl Channel for ReplChannel {
async fn start(&self) -> Result<MessageStream, ChannelError> {
let (tx, rx) = mpsc::channel(32);
// Approval prompts inject responses back through this sender.
// In single-message mode we keep it until the turn finishes, then
// drop it after enqueuing /quit so the receiver stream can close.
if let Ok(mut guard) = self.msg_tx.lock() {
// Store tx so send_status can inject approval responses directly.
// Skip for single-message mode — no interactive approval is needed
// and the extra sender would keep the stream open after /quit.
if self.single_message.is_none()
&& let Ok(mut guard) = self.msg_tx.lock()
{
*guard = Some(tx.clone());
}
let single_message = self.single_message.clone();
@@ -914,8 +916,10 @@ mod tests {
use super::*;
/// Regression: single-message mode must close the stream after the one
/// message so callers (and tests) don't hang forever.
#[tokio::test]
async fn single_message_mode_sends_message_then_quit() {
async fn single_message_mode_sends_message_and_closes_stream() {
let repl = ReplChannel::with_message("hi".to_string());
let mut stream = repl.start().await.expect("repl start should succeed");
@@ -926,30 +930,15 @@ mod tests {
assert_eq!(first.channel, "repl");
assert_eq!(first.content, "hi");
assert!(
timeout(Duration::from_millis(100), stream.next())
.await
.is_err(),
"single-message mode should wait for the turn to finish before quitting"
);
repl.respond(&first, OutgoingResponse::text("done"))
.await
.expect("respond should succeed");
let second = timeout(Duration::from_secs(1), stream.next())
.await
.expect("timed out waiting for quit message")
.expect("quit message missing");
assert_eq!(second.channel, "repl");
assert_eq!(second.content, "/quit");
// The spawned thread sent the message and returned, dropping its
// sender. Because we skip storing a clone in msg_tx for single-
// message mode, the stream should close immediately.
assert!(
timeout(Duration::from_secs(1), stream.next())
.await
.expect("timed out waiting for stream to close")
.is_none(),
"stream should end after /quit"
"stream should end after the single message"
);
}
}
+26 -11
View File
@@ -236,6 +236,18 @@ pub async fn jobs_detail_handler(
(end - start).num_seconds().max(0) as u64
});
// Build transitions from the job's state transition history.
let transitions: Vec<TransitionInfo> = ctx
.transitions
.iter()
.map(|t| TransitionInfo {
from: t.from.to_string(),
to: t.to.to_string(),
timestamp: t.timestamp.to_rfc3339(),
reason: t.reason.clone(),
})
.collect();
// Only show prompt bar for jobs that have a running worker (Pending/InProgress).
// Stuck jobs have no active worker loop, so messages would be silently dropped.
let is_promptable = matches!(
@@ -255,7 +267,7 @@ pub async fn jobs_detail_handler(
project_dir: None,
browse_url: None,
job_mode: None,
transitions: Vec::new(),
transitions,
can_restart: state.scheduler.is_some(),
can_prompt: is_promptable && state.scheduler.is_some(),
job_kind: Some("agent".to_string()),
@@ -643,25 +655,28 @@ pub async fn jobs_events_handler(
.parse()
.map_err(|_| (StatusCode::BAD_REQUEST, "Invalid job ID".to_string()))?;
// Verify ownership before returning events.
match store.get_sandbox_job(job_id).await {
Ok(Some(job)) => {
if job.user_id != user.user_id {
return Err((StatusCode::NOT_FOUND, "Job not found".to_string()));
// Verify ownership before returning events (check both sandbox and agent jobs).
let is_owner = match store.get_sandbox_job(job_id).await {
Ok(Some(job)) => job.user_id == user.user_id,
Ok(None) => {
// Fall back to agent job ownership check.
match store.get_job(job_id).await {
Ok(Some(ctx)) => ctx.user_id == user.user_id,
_ => false,
}
}
Ok(None) => {
return Err((StatusCode::NOT_FOUND, "Job not found".to_string()));
}
Err(e) => {
return Err(db_error("jobs_handler", e));
return Err(db_error("jobs_events_handler", e));
}
};
if !is_owner {
return Err((StatusCode::NOT_FOUND, "Job not found".to_string()));
}
let events = store
.list_job_events(job_id, None)
.await
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
.map_err(|e| db_error("jobs_events_handler", e))?;
let events_json: Vec<serde_json::Value> = events
.into_iter()
+11
View File
@@ -122,6 +122,16 @@ pub async fn routines_detail_handler(
.collect();
let routine_info = RoutineInfo::from_routine(&routine);
// Read-only lookup — do not create a conversation on a GET request.
// The conversation is created lazily when the routine first executes.
let conversation_id = store
.find_routine_conversation(routine.id, &routine.user_id)
.await
.unwrap_or_else(|e| {
tracing::warn!(routine_id = %routine.id, error = %e, "Failed to look up routine conversation");
None
});
Ok(Json(RoutineDetailResponse {
id: routine.id,
name: routine.name.clone(),
@@ -139,6 +149,7 @@ pub async fn routines_detail_handler(
run_count: routine.run_count,
consecutive_failures: routine.consecutive_failures,
created_at: routine.created_at.to_rfc3339(),
conversation_id,
recent_runs,
}))
}
+8 -8
View File
@@ -347,10 +347,10 @@ impl Channel for GatewayChannel {
let thread_id = match &msg.thread_id {
Some(tid) => tid.clone(),
None => {
tracing::warn!(
"Gateway respond with no thread_id — skipping (clients would drop it)"
);
return Ok(());
return Err(ChannelError::MissingRoutingTarget {
name: "gateway".to_string(),
reason: "respond() requires a thread_id on the incoming message".to_string(),
});
}
};
@@ -507,10 +507,10 @@ impl Channel for GatewayChannel {
let thread_id = match response.thread_id {
Some(tid) => tid,
None => {
tracing::warn!(
"Gateway broadcast with no thread_id — skipping (clients would drop it)"
);
return Ok(());
return Err(ChannelError::MissingRoutingTarget {
name: "gateway".to_string(),
reason: "broadcast() requires a thread_id on the response".to_string(),
});
}
};
self.state.sse.broadcast_for_user(
+12
View File
@@ -4265,6 +4265,13 @@ function renderRoutineDetail(routine) {
html += '<div class="job-description"><h3>Action</h3>'
+ '<pre class="action-json">' + escapeHtml(JSON.stringify(routine.action, null, 2)) + '</pre></div>';
// Conversation thread link
if (routine.conversation_id) {
html += '<div class="job-description">'
+ '<a href="#" data-action="view-routine-thread" data-id="' + escapeHtml(routine.conversation_id) + '" class="btn-primary" style="display:inline-block;margin:0.5rem 0">'
+ 'View Execution Thread</a></div>';
}
// Recent runs
if (routine.recent_runs && routine.recent_runs.length > 0) {
html += '<div class="job-timeline-section"><h3>Recent Runs</h3>'
@@ -6190,6 +6197,11 @@ document.addEventListener('click', function(e) {
switchTab('jobs');
openJobDetail(el.dataset.id);
break;
case 'view-routine-thread':
e.preventDefault();
switchTab('chat');
switchThread(el.dataset.id);
break;
case 'copy-tee-report':
copyTeeReport();
break;
+1
View File
@@ -1,3 +1,4 @@
//! Integration tests for the web gateway module.
mod multi_tenant;
mod no_silent_drop;
+93
View File
@@ -0,0 +1,93 @@
//! Regression tests: the gateway channel must never silently drop messages.
//!
//! Previously, `respond()` and `broadcast()` returned `Ok(())` when thread_id
//! was missing, making callers believe the message was delivered when it wasn't.
//! These tests ensure that missing routing info produces an explicit error.
use crate::channels::channel::{Channel, IncomingMessage, OutgoingResponse};
use crate::channels::web::GatewayChannel;
use crate::config::GatewayConfig;
use crate::error::ChannelError;
fn test_gateway() -> GatewayChannel {
GatewayChannel::new(
GatewayConfig {
host: "127.0.0.1".to_string(),
port: 0,
auth_token: Some("test-token".to_string()),
workspace_read_scopes: vec![],
memory_layers: vec![],
},
"test-user".to_string(),
)
}
#[tokio::test]
async fn gateway_respond_without_thread_id_returns_error() {
let gw = test_gateway();
let msg = IncomingMessage::new("gateway", "test-user", "hello");
// msg has no thread_id by default
assert!(msg.thread_id.is_none());
let response = OutgoingResponse::text("reply");
let result = gw.respond(&msg, response).await;
assert!(
result.is_err(),
"respond() must not silently succeed without thread_id"
);
assert!(
matches!(result, Err(ChannelError::MissingRoutingTarget { .. })),
"Expected MissingRoutingTarget, got: {:?}",
result
);
}
#[tokio::test]
async fn gateway_respond_with_thread_id_succeeds() {
let gw = test_gateway();
let mut msg = IncomingMessage::new("gateway", "test-user", "hello");
msg.thread_id = Some("thread-123".to_string());
let response = OutgoingResponse::text("reply");
let result = gw.respond(&msg, response).await;
assert!(
result.is_ok(),
"respond() should succeed with thread_id: {:?}",
result
);
}
#[tokio::test]
async fn gateway_broadcast_without_thread_id_returns_error() {
let gw = test_gateway();
let response = OutgoingResponse::text("notification");
// response has no thread_id by default
let result = gw.broadcast("test-user", response).await;
assert!(
result.is_err(),
"broadcast() must not silently succeed without thread_id"
);
assert!(
matches!(result, Err(ChannelError::MissingRoutingTarget { .. })),
"Expected MissingRoutingTarget, got: {:?}",
result
);
}
#[tokio::test]
async fn gateway_broadcast_with_thread_id_succeeds() {
let gw = test_gateway();
let response = OutgoingResponse::text("notification").in_thread("thread-456".to_string());
let result = gw.broadcast("test-user", response).await;
assert!(
result.is_ok(),
"broadcast() should succeed with thread_id: {:?}",
result
);
}
+1
View File
@@ -768,6 +768,7 @@ pub struct RoutineDetailResponse {
pub run_count: u64,
pub consecutive_failures: u32,
pub created_at: String,
pub conversation_id: Option<Uuid>,
pub recent_runs: Vec<RoutineRunInfo>,
}
+13 -5
View File
@@ -10,8 +10,10 @@ use crate::error::ConfigError;
pub struct HygieneConfig {
/// Whether hygiene is enabled. Env: `MEMORY_HYGIENE_ENABLED` (default: true).
pub enabled: bool,
/// Maximum versions to keep per document. Env: `MEMORY_HYGIENE_VERSION_KEEP_COUNT` (default: 50).
pub version_keep_count: u32,
/// Days before `daily/` documents are deleted. Env: `MEMORY_HYGIENE_DAILY_RETENTION_DAYS` (default: 30).
pub daily_retention_days: u32,
/// Days before `conversations/` documents are deleted. Env: `MEMORY_HYGIENE_CONVERSATION_RETENTION_DAYS` (default: 7).
pub conversation_retention_days: u32,
/// Minimum hours between hygiene passes. Env: `MEMORY_HYGIENE_CADENCE_HOURS` (default: 12).
pub cadence_hours: u32,
}
@@ -20,7 +22,8 @@ impl Default for HygieneConfig {
fn default() -> Self {
Self {
enabled: true,
version_keep_count: 50,
daily_retention_days: 30,
conversation_retention_days: 7,
cadence_hours: 12,
}
}
@@ -30,7 +33,11 @@ impl HygieneConfig {
pub(crate) fn resolve() -> Result<Self, ConfigError> {
Ok(Self {
enabled: parse_bool_env("MEMORY_HYGIENE_ENABLED", true)?,
version_keep_count: parse_optional_env("MEMORY_HYGIENE_VERSION_KEEP_COUNT", 50)?,
daily_retention_days: parse_optional_env("MEMORY_HYGIENE_DAILY_RETENTION_DAYS", 30)?,
conversation_retention_days: parse_optional_env(
"MEMORY_HYGIENE_CONVERSATION_RETENTION_DAYS",
7,
)?,
cadence_hours: parse_optional_env("MEMORY_HYGIENE_CADENCE_HOURS", 12)?,
})
}
@@ -40,7 +47,8 @@ impl HygieneConfig {
pub fn to_workspace_config(&self) -> crate::workspace::hygiene::HygieneConfig {
crate::workspace::hygiene::HygieneConfig {
enabled: self.enabled,
version_keep_count: self.version_keep_count,
daily_retention_days: self.daily_retention_days,
conversation_retention_days: self.conversation_retention_days,
cadence_hours: self.cadence_hours,
state_dir: ironclaw_base_dir(),
}
+35
View File
@@ -290,6 +290,41 @@ impl ConversationStore for LibSqlBackend {
result
}
async fn find_routine_conversation(
&self,
routine_id: Uuid,
user_id: &str,
) -> Result<Option<Uuid>, DatabaseError> {
let conn = self.connect().await?;
let rid = routine_id.to_string();
let mut rows = conn
.query(
r#"
SELECT id FROM conversations
WHERE user_id = ?1 AND json_extract(metadata, '$.routine_id') = ?2
LIMIT 1
"#,
params![user_id, rid],
)
.await
.map_err(|e| DatabaseError::Query(e.to_string()))?;
if let Some(row) = rows
.next()
.await
.map_err(|e| DatabaseError::Query(e.to_string()))?
{
let id_str: String = row.get(0).map_err(|e| {
DatabaseError::Query(format!("Failed to read conversation id: {e}"))
})?;
let id = id_str
.parse()
.map_err(|_| DatabaseError::Serialization("Invalid UUID".to_string()))?;
return Ok(Some(id));
}
Ok(None)
}
/// Uses BEGIN IMMEDIATE to serialize concurrent writers and prevent
/// duplicate heartbeat conversations (TOCTOU race).
async fn get_or_create_heartbeat_conversation(
+2 -326
View File
@@ -13,8 +13,8 @@ use super::{
use crate::db::WorkspaceStore;
use crate::error::{DatabaseError, WorkspaceError};
use crate::workspace::{
DocumentVersion, MemoryChunk, MemoryDocument, RankedResult, SearchConfig, SearchResult,
VersionSummary, WorkspaceEntry, fuse_results,
MemoryChunk, MemoryDocument, RankedResult, SearchConfig, SearchResult, WorkspaceEntry,
fuse_results,
};
use chrono::Utc;
@@ -840,330 +840,6 @@ impl WorkspaceStore for LibSqlBackend {
Ok(fuse_results(fts_results, vector_results, config))
}
// ==================== Metadata ====================
async fn update_document_metadata(
&self,
id: Uuid,
metadata: &serde_json::Value,
) -> Result<(), WorkspaceError> {
let conn = self
.connect()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: e.to_string(),
})?;
let now = fmt_ts(&Utc::now());
let meta_str =
serde_json::to_string(metadata).map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to serialize metadata: {e}"),
})?;
conn.execute(
"UPDATE memory_documents SET metadata = ?2, updated_at = ?3 WHERE id = ?1",
params![id.to_string(), meta_str, now],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to update metadata: {e}"),
})?;
Ok(())
}
async fn find_config_documents(
&self,
user_id: &str,
agent_id: Option<Uuid>,
) -> Result<Vec<MemoryDocument>, WorkspaceError> {
let conn = self
.connect()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: e.to_string(),
})?;
let agent_str = agent_id.map(|a| a.to_string());
let mut rows = conn
.query(
r#"
SELECT id, user_id, agent_id, path, content,
created_at, updated_at, metadata
FROM memory_documents
WHERE user_id = ?1 AND agent_id IS ?2
AND (path LIKE '%/.config' OR path = '.config')
ORDER BY path
"#,
params![user_id, agent_str],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to find config documents: {e}"),
})?;
let mut docs = Vec::new();
while let Some(row) = rows
.next()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to read config document row: {e}"),
})?
{
docs.push(row_to_memory_document(&row));
}
Ok(docs)
}
// ==================== Versioning ====================
async fn save_version(
&self,
document_id: Uuid,
content: &str,
content_hash: &str,
changed_by: Option<&str>,
) -> Result<i32, WorkspaceError> {
let conn = self
.connect()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: e.to_string(),
})?;
let id = Uuid::new_v4().to_string();
let doc_id = document_id.to_string();
let now = fmt_ts(&Utc::now());
// Use a transaction to prevent race conditions: the SELECT and INSERT
// must be atomic so concurrent writers don't allocate the same version.
let tx = conn
.transaction()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to start transaction: {e}"),
})?;
// Get next version number (inside transaction — serializes writers)
let mut rows = tx
.query(
"SELECT COALESCE(MAX(version), 0) + 1 FROM memory_document_versions WHERE document_id = ?1",
params![doc_id.clone()],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to get next version number: {e}"),
})?;
let next_version = if let Some(row) =
rows.next()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to read version number: {e}"),
})? {
get_i64(&row, 0) as i32
} else {
1
};
drop(rows);
tx.execute(
r#"
INSERT INTO memory_document_versions
(id, document_id, version, content, content_hash, created_at, changed_by)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
"#,
params![
id,
doc_id,
next_version as i64,
content,
content_hash,
now,
changed_by
],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to save version: {e}"),
})?;
tx.commit()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to commit version: {e}"),
})?;
Ok(next_version)
}
async fn get_version(
&self,
document_id: Uuid,
version: i32,
) -> Result<DocumentVersion, WorkspaceError> {
let conn = self
.connect()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: e.to_string(),
})?;
let mut rows = conn
.query(
r#"
SELECT id, document_id, version, content, content_hash,
created_at, changed_by
FROM memory_document_versions
WHERE document_id = ?1 AND version = ?2
"#,
params![document_id.to_string(), version as i64],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to get version: {e}"),
})?;
let row = rows
.next()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to read version row: {e}"),
})?
.ok_or(WorkspaceError::VersionNotFound {
document_id,
version,
})?;
Ok(DocumentVersion {
id: get_text(&row, 0)
.parse()
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Invalid version UUID: {e}"),
})?,
document_id: get_text(&row, 1)
.parse()
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Invalid document UUID: {e}"),
})?,
version: get_i64(&row, 2) as i32,
content: get_text(&row, 3),
content_hash: get_text(&row, 4),
created_at: get_ts(&row, 5),
changed_by: get_opt_text(&row, 6),
})
}
async fn list_versions(
&self,
document_id: Uuid,
limit: i64,
) -> Result<Vec<VersionSummary>, WorkspaceError> {
let conn = self
.connect()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: e.to_string(),
})?;
let mut rows = conn
.query(
r#"
SELECT version, content_hash, created_at, changed_by
FROM memory_document_versions
WHERE document_id = ?1
ORDER BY version DESC
LIMIT ?2
"#,
params![document_id.to_string(), limit],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to list versions: {e}"),
})?;
let mut versions = Vec::new();
while let Some(row) = rows
.next()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to read version row: {e}"),
})?
{
versions.push(VersionSummary {
version: get_i64(&row, 0) as i32,
content_hash: get_text(&row, 1),
created_at: get_ts(&row, 2),
changed_by: get_opt_text(&row, 3),
});
}
Ok(versions)
}
async fn get_latest_version_number(
&self,
document_id: Uuid,
) -> Result<Option<i32>, WorkspaceError> {
let conn = self
.connect()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: e.to_string(),
})?;
let mut rows = conn
.query(
"SELECT MAX(version) FROM memory_document_versions WHERE document_id = ?1",
params![document_id.to_string()],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to get latest version number: {e}"),
})?;
if let Some(row) = rows
.next()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to read version number: {e}"),
})?
{
// MAX returns NULL if no rows — libsql returns Null for the value
let val = row.get::<libsql::Value>(0).ok();
match val {
Some(libsql::Value::Integer(v)) => Ok(Some(v as i32)),
_ => Ok(None),
}
} else {
Ok(None)
}
}
async fn prune_versions(
&self,
document_id: Uuid,
keep_count: i32,
) -> Result<u64, WorkspaceError> {
let conn = self
.connect()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: e.to_string(),
})?;
let doc_id = document_id.to_string();
let result = conn
.execute(
r#"
DELETE FROM memory_document_versions
WHERE document_id = ?1
AND version NOT IN (
SELECT version FROM memory_document_versions
WHERE document_id = ?1
ORDER BY version DESC
LIMIT ?2
)
"#,
params![doc_id, keep_count as i64],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to prune versions: {e}"),
})?;
Ok(result)
}
}
#[cfg(test)]
-19
View File
@@ -785,25 +785,6 @@ CREATE TABLE IF NOT EXISTS api_tokens (
);
CREATE INDEX IF NOT EXISTS idx_api_tokens_user ON api_tokens(user_id);
CREATE INDEX IF NOT EXISTS idx_api_tokens_hash ON api_tokens(token_hash);
"#,
),
(
15,
"document_versions",
r#"
CREATE TABLE IF NOT EXISTS memory_document_versions (
id TEXT PRIMARY KEY,
document_id TEXT NOT NULL REFERENCES memory_documents(id) ON DELETE CASCADE,
version INTEGER NOT NULL,
content TEXT NOT NULL,
content_hash TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')),
changed_by TEXT,
UNIQUE(document_id, version)
);
CREATE INDEX IF NOT EXISTS idx_doc_versions_lookup
ON memory_document_versions(document_id, version DESC);
"#,
),
];
+7 -61
View File
@@ -391,6 +391,13 @@ pub trait ConversationStore: Send + Sync {
routine_name: &str,
user_id: &str,
) -> Result<Uuid, DatabaseError>;
/// Read-only lookup for an existing routine conversation. Returns `None`
/// if the routine has never executed (no conversation created yet).
async fn find_routine_conversation(
&self,
routine_id: Uuid,
user_id: &str,
) -> Result<Option<Uuid>, DatabaseError>;
async fn get_or_create_heartbeat_conversation(
&self,
user_id: &str,
@@ -700,67 +707,6 @@ pub trait WorkspaceStore: Send + Sync {
config: &SearchConfig,
) -> Result<Vec<SearchResult>, WorkspaceError>;
// ==================== Metadata ====================
/// Update the metadata JSON field on a document (full replacement).
async fn update_document_metadata(
&self,
id: Uuid,
metadata: &serde_json::Value,
) -> Result<(), WorkspaceError>;
/// Find all `.config` documents in the workspace.
///
/// Returns documents whose path ends with `/.config` or equals `.config`.
/// Used by the hygiene system to discover metadata-driven cleanup targets.
async fn find_config_documents(
&self,
user_id: &str,
agent_id: Option<Uuid>,
) -> Result<Vec<MemoryDocument>, WorkspaceError>;
// ==================== Versioning ====================
/// Save the current content of a document as a new version.
///
/// Returns the new version number (1-based, monotonically increasing).
async fn save_version(
&self,
document_id: Uuid,
content: &str,
content_hash: &str,
changed_by: Option<&str>,
) -> Result<i32, WorkspaceError>;
/// Get a specific version of a document.
async fn get_version(
&self,
document_id: Uuid,
version: i32,
) -> Result<crate::workspace::DocumentVersion, WorkspaceError>;
/// List versions of a document (newest first).
async fn list_versions(
&self,
document_id: Uuid,
limit: i64,
) -> Result<Vec<crate::workspace::VersionSummary>, WorkspaceError>;
/// Get the latest version number for a document, or `None` if no versions exist.
async fn get_latest_version_number(
&self,
document_id: Uuid,
) -> Result<Option<i32>, WorkspaceError>;
/// Delete old versions, keeping only the most recent `keep_count`.
///
/// Returns the number of versions deleted.
async fn prune_versions(
&self,
document_id: Uuid,
keep_count: i32,
) -> Result<u64, WorkspaceError>;
// ==================== Multi-scope read methods ====================
//
// Default implementations loop over user_ids calling single-scope methods,
+11 -65
View File
@@ -25,8 +25,7 @@ use crate::history::{
LlmCallRecord, SandboxJobRecord, SandboxJobSummary, SettingRow, Store,
};
use crate::workspace::{
DocumentVersion, MemoryChunk, MemoryDocument, Repository, SearchConfig, SearchResult,
VersionSummary, WorkspaceEntry,
MemoryChunk, MemoryDocument, Repository, SearchConfig, SearchResult, WorkspaceEntry,
};
/// PostgreSQL database backend.
@@ -138,6 +137,16 @@ impl ConversationStore for PgBackend {
.await
}
async fn find_routine_conversation(
&self,
routine_id: Uuid,
user_id: &str,
) -> Result<Option<Uuid>, DatabaseError> {
self.store
.find_routine_conversation(routine_id, user_id)
.await
}
async fn get_or_create_heartbeat_conversation(
&self,
user_id: &str,
@@ -786,69 +795,6 @@ impl WorkspaceStore for PgBackend {
.list_directory_multi(user_ids, agent_id, directory)
.await
}
// ==================== Metadata ====================
async fn update_document_metadata(
&self,
id: Uuid,
metadata: &serde_json::Value,
) -> Result<(), WorkspaceError> {
self.repo.update_document_metadata(id, metadata).await
}
async fn find_config_documents(
&self,
user_id: &str,
agent_id: Option<Uuid>,
) -> Result<Vec<MemoryDocument>, WorkspaceError> {
self.repo.find_config_documents(user_id, agent_id).await
}
// ==================== Versioning ====================
async fn save_version(
&self,
document_id: Uuid,
content: &str,
content_hash: &str,
changed_by: Option<&str>,
) -> Result<i32, WorkspaceError> {
self.repo
.save_version(document_id, content, content_hash, changed_by)
.await
}
async fn get_version(
&self,
document_id: Uuid,
version: i32,
) -> Result<DocumentVersion, WorkspaceError> {
self.repo.get_version(document_id, version).await
}
async fn list_versions(
&self,
document_id: Uuid,
limit: i64,
) -> Result<Vec<VersionSummary>, WorkspaceError> {
self.repo.list_versions(document_id, limit).await
}
async fn get_latest_version_number(
&self,
document_id: Uuid,
) -> Result<Option<i32>, WorkspaceError> {
self.repo.get_latest_version_number(document_id).await
}
async fn prune_versions(
&self,
document_id: Uuid,
keep_count: i32,
) -> Result<u64, WorkspaceError> {
self.repo.prune_versions(document_id, keep_count).await
}
}
// ==================== UserStore ====================
-6
View File
@@ -315,12 +315,6 @@ pub enum WorkspaceError {
#[error("Write rejected for '{path}': prompt injection detected ({reason})")]
InjectionRejected { path: String, reason: String },
#[error("Version not found: document {document_id} version {version}")]
VersionNotFound { document_id: Uuid, version: i32 },
#[error("Patch failed for '{path}': {reason}")]
PatchFailed { path: String, reason: String },
}
/// Orchestrator errors (internal API, container management).
+21
View File
@@ -1771,6 +1771,27 @@ impl Store {
Ok(row.get("id"))
}
/// Read-only lookup for an existing routine conversation.
pub async fn find_routine_conversation(
&self,
routine_id: Uuid,
user_id: &str,
) -> Result<Option<Uuid>, DatabaseError> {
let conn = self.conn().await?;
let rid = routine_id.to_string();
let row = conn
.query_opt(
r#"
SELECT id FROM conversations
WHERE user_id = $1 AND metadata->>'routine_id' = $2
LIMIT 1
"#,
&[&user_id, &rid],
)
.await?;
Ok(row.map(|r| r.get("id")))
}
/// Get or create the singleton heartbeat conversation for a user.
///
/// Looks for a conversation where `metadata->>'thread_type' = 'heartbeat'`.
+133 -3
View File
@@ -276,8 +276,33 @@ impl LlmProvider for OpenAiCodexProvider {
&self,
request: ToolCompletionRequest,
) -> Result<ToolCompletionResponse, LlmError> {
// Build a reverse map so we can translate sanitized names back to originals.
// Only needed when sanitization actually changes a name (e.g. MCP tools with dots).
let name_map: std::collections::HashMap<String, String> = request
.tools
.iter()
.filter_map(|t| {
let sanitized = sanitize_tool_name(&t.name);
if sanitized != t.name {
Some((sanitized, t.name.clone()))
} else {
None
}
})
.collect();
let body = self.build_request_body(&request.messages, Some(&request.tools));
let parsed = self.send_request(body).await?;
let mut parsed = self.send_request(body).await?;
// Reverse-map sanitized tool names back to originals so the caller
// can look them up in the tool registry.
if !name_map.is_empty() {
for tc in &mut parsed.tool_calls {
if let Some(original) = name_map.get(&tc.name) {
tc.name = original.clone();
}
}
}
let finish_reason = if !parsed.tool_calls.is_empty() {
FinishReason::ToolUse
@@ -421,7 +446,7 @@ fn convert_message(msg: &ChatMessage, index: usize) -> Vec<serde_json::Value> {
serde_json::json!({
"type": "function_call",
"call_id": tc.id,
"name": tc.name,
"name": sanitize_tool_name(&tc.name),
"arguments": args_str,
})
})
@@ -452,6 +477,20 @@ fn convert_message(msg: &ChatMessage, index: usize) -> Vec<serde_json::Value> {
}
}
/// Sanitize a tool name to match the OpenAI Responses API pattern `^[a-zA-Z0-9_-]+$`.
/// Replaces any invalid character (e.g. dots in MCP tool names) with underscores.
fn sanitize_tool_name(name: &str) -> String {
name.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '_' || c == '-' {
c
} else {
'_'
}
})
.collect()
}
/// Convert a `ToolDefinition` to Responses API tool format.
///
/// Applies strict-mode schema normalization (same as OpenAI Chat Completions):
@@ -461,7 +500,7 @@ fn convert_tool_definition(tool: &ToolDefinition) -> serde_json::Value {
serde_json::json!({
"type": "function",
"name": tool.name,
"name": sanitize_tool_name(&tool.name),
"description": tool.description,
"parameters": normalize_schema_strict(&tool.parameters),
})
@@ -1093,4 +1132,95 @@ data: {"type":"response.completed","response":{"status":"completed","usage":{"in
assert_eq!(parsed.tool_calls[1].name, "read_file");
assert_eq!(parsed.finish_reason, FinishReason::ToolUse);
}
/// Regression test: tool names with dots (e.g. MCP tools) must be sanitized
/// to match OpenAI's `^[a-zA-Z0-9_-]+$` pattern.
#[test]
fn test_sanitize_tool_name_replaces_dots() {
assert_eq!(super::sanitize_tool_name("memory_search"), "memory_search");
assert_eq!(
super::sanitize_tool_name("mcp.server.tool"),
"mcp_server_tool"
);
assert_eq!(super::sanitize_tool_name("tool@v2"), "tool_v2");
assert_eq!(super::sanitize_tool_name("my-tool"), "my-tool");
}
/// Regression test: convert_tool_definition sanitizes the name.
#[test]
fn test_convert_tool_definition_sanitizes_name() {
let tool = ToolDefinition {
name: "mcp.server.search".to_string(),
description: "Search".to_string(),
parameters: serde_json::json!({"type": "object", "properties": {}}),
};
let json = super::convert_tool_definition(&tool);
assert_eq!(json["name"], "mcp_server_search");
}
/// Regression test: function_call items sanitize tool names.
#[test]
fn test_convert_message_sanitizes_tool_call_name() {
let tool_calls = vec![ToolCall {
id: "call_1".to_string(),
name: "mcp.server.search".to_string(),
arguments: serde_json::json!({"q": "test"}),
reasoning: None,
}];
let msg = ChatMessage::assistant_with_tool_calls(None, tool_calls);
let items = super::convert_message(&msg, 0);
assert_eq!(items[0]["name"], "mcp_server_search");
}
/// Regression: sanitized tool names in API responses must be reverse-mapped
/// back to original names so the tool registry can look them up.
#[test]
fn test_sanitized_name_reverse_mapping() {
use std::collections::HashMap;
let tools = [
ToolDefinition {
name: "mcp.server.search".to_string(),
description: "Search".to_string(),
parameters: serde_json::json!({"type": "object", "properties": {}}),
},
ToolDefinition {
name: "memory_search".to_string(),
description: "Memory".to_string(),
parameters: serde_json::json!({"type": "object", "properties": {}}),
},
];
// Build name map (same logic as complete_with_tools)
let name_map: HashMap<String, String> = tools
.iter()
.filter_map(|t| {
let sanitized = super::sanitize_tool_name(&t.name);
if sanitized != t.name {
Some((sanitized, t.name.clone()))
} else {
None
}
})
.collect();
// Only the MCP tool should appear (its name changed)
assert_eq!(name_map.len(), 1);
assert_eq!(
name_map.get("mcp_server_search"),
Some(&"mcp.server.search".to_string())
);
// Simulate a tool call coming back with the sanitized name
let mut tc = ToolCall {
id: "call_1".to_string(),
name: "mcp_server_search".to_string(),
arguments: serde_json::json!({}),
reasoning: None,
};
if let Some(original) = name_map.get(&tc.name) {
tc.name = original.clone();
}
assert_eq!(tc.name, "mcp.server.search");
}
}
+10 -2
View File
@@ -182,8 +182,14 @@ impl SkillCatalog {
/// Create a catalog with a custom registry URL (for testing).
#[cfg(test)]
pub fn with_url(url: &str) -> Self {
Self::with_url_and_timeout(url, REQUEST_TIMEOUT)
}
/// Create a catalog with a custom registry URL and timeout (for testing).
#[cfg(test)]
pub fn with_url_and_timeout(url: &str, timeout: Duration) -> Self {
let client = reqwest::Client::builder()
.timeout(REQUEST_TIMEOUT)
.timeout(timeout)
.user_agent(concat!("ironclaw/", env!("CARGO_PKG_VERSION")))
.build()
.unwrap_or_default();
@@ -458,7 +464,9 @@ mod tests {
#[tokio::test]
async fn test_search_returns_error_on_network_failure() {
// Use RFC 5737 TEST-NET-1 (192.0.2.0/24) for reliable failure even behind proxies.
let catalog = SkillCatalog::with_url("http://192.0.2.1:9999");
// Short timeout so the test doesn't block for the full 10s REQUEST_TIMEOUT.
let catalog =
SkillCatalog::with_url_and_timeout("http://192.0.2.1:9999", Duration::from_secs(1));
let outcome = catalog.search("test").await;
assert!(outcome.results.is_empty());
assert!(outcome.error.is_some());
+4 -140
View File
@@ -246,26 +246,9 @@ impl Tool for MemoryWriteTool {
"type": "boolean",
"description": "Skip privacy classification and write directly to the specified layer without redirect. Use when you're certain the content belongs in the target layer.",
"default": false
},
"metadata": {
"type": "object",
"description": "Optional metadata to set on the document (e.g., {\"skip_indexing\": true, \"hygiene\": {\"enabled\": true, \"retention_days\": 7}})"
},
"old_string": {
"type": "string",
"description": "When present, switches to patch mode: finds and replaces this exact string in the document. Requires target to be a path (not 'memory' or 'daily_log')."
},
"new_string": {
"type": "string",
"description": "Replacement string (required when old_string is present)."
},
"replace_all": {
"type": "boolean",
"description": "If true, replace all occurrences of old_string. Default: false.",
"default": false
}
},
"required": []
"required": ["content"]
})
}
@@ -276,9 +259,7 @@ impl Tool for MemoryWriteTool {
) -> Result<ToolOutput, ToolError> {
let start = std::time::Instant::now();
// In patch mode (old_string present), content is not required.
let is_patch_mode = params.get("old_string").and_then(|v| v.as_str()).is_some();
let content = params.get("content").and_then(|v| v.as_str()).unwrap_or("");
let content = require_str(&params, "content")?;
let target = params
.get("target")
@@ -318,9 +299,9 @@ impl Tool for MemoryWriteTool {
return Ok(ToolOutput::success(output, start.elapsed()));
}
if !is_patch_mode && content.trim().is_empty() {
if content.trim().is_empty() {
return Err(ToolError::InvalidParameters(
"content cannot be empty (use old_string/new_string for patch mode)".to_string(),
"content cannot be empty".to_string(),
));
}
@@ -349,46 +330,6 @@ impl Tool for MemoryWriteTool {
path => path.to_string(),
};
// Patch mode: if old_string is provided, do search-and-replace instead of write/append.
let old_string = params.get("old_string").and_then(|v| v.as_str());
if let Some(old_str) = old_string {
let new_str = params
.get("new_string")
.and_then(|v| v.as_str())
.ok_or_else(|| {
ToolError::InvalidParameters(
"new_string is required when old_string is provided".to_string(),
)
})?;
let replace_all = params
.get("replace_all")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let result = workspace
.patch(&resolved_path, old_str, new_str, replace_all)
.await
.map_err(map_write_err)?;
// Apply metadata if provided
if let Some(meta) = params.get("metadata")
&& meta.is_object()
{
workspace
.update_metadata(result.document.id, meta)
.await
.map_err(map_write_err)?;
}
let output = serde_json::json!({
"status": "patched",
"path": resolved_path,
"replacements": result.replacements,
"content_length": result.document.content.len(),
});
return Ok(ToolOutput::success(output, start.elapsed()));
}
// When a layer is specified, route through layer-aware methods for ALL targets.
// Otherwise, use default workspace methods (which include injection scanning).
let layer_result = if let Some(layer_name) = layer {
@@ -492,24 +433,6 @@ impl Tool for MemoryWriteTool {
}
}
// Apply metadata if provided (after write/append, works for all targets).
// We read the document once to get its ID — this is a hot read right
// after the write, so it's effectively free (same DB connection/cache).
if let Some(meta) = params.get("metadata")
&& meta.is_object()
{
match workspace.read(&resolved_path).await {
Ok(doc) => {
if let Err(e) = workspace.update_metadata(doc.id, meta).await {
tracing::warn!(path = %resolved_path, "failed to update metadata: {e}");
}
}
Err(e) => {
tracing::warn!(path = %resolved_path, "failed to read doc for metadata update: {e}");
}
}
}
let mut output = serde_json::json!({
"status": "written",
"path": resolved_path,
@@ -578,15 +501,6 @@ impl Tool for MemoryReadTool {
"path": {
"type": "string",
"description": "Path to the file (e.g., 'MEMORY.md', 'daily/2024-01-15.md', 'projects/alpha/notes.md')"
},
"version": {
"type": "integer",
"description": "Read a specific historical version of the document (omit for current content)"
},
"list_versions": {
"type": "boolean",
"description": "If true, return version history instead of file content",
"default": false
}
},
"required": ["path"]
@@ -611,61 +525,11 @@ impl Tool for MemoryReadTool {
}
let workspace = self.resolver.resolve(&ctx.user_id).await;
let list_versions = params
.get("list_versions")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let version = params
.get("version")
.and_then(|v| v.as_i64())
.map(|v| v as i32);
// Read the document first (needed for document_id in all version operations)
let doc = workspace
.read(path)
.await
.map_err(|e| ToolError::ExecutionFailed(format!("Read failed: {}", e)))?;
// List versions mode
if list_versions {
let versions = workspace
.list_versions(doc.id, 50)
.await
.map_err(|e| ToolError::ExecutionFailed(format!("List versions failed: {}", e)))?;
let output = serde_json::json!({
"path": doc.path,
"versions": versions.iter().map(|v| serde_json::json!({
"version": v.version,
"content_hash": v.content_hash,
"created_at": v.created_at.to_rfc3339(),
"changed_by": v.changed_by,
})).collect::<Vec<_>>(),
"version_count": versions.len(),
});
return Ok(ToolOutput::success(output, start.elapsed()));
}
// Specific version mode
if let Some(ver) = version {
let version_doc = workspace
.get_version(doc.id, ver)
.await
.map_err(|e| ToolError::ExecutionFailed(format!("Get version failed: {}", e)))?;
let output = serde_json::json!({
"path": doc.path,
"version": version_doc.version,
"content": version_doc.content,
"content_hash": version_doc.content_hash,
"created_at": version_doc.created_at.to_rfc3339(),
"changed_by": version_doc.changed_by,
});
return Ok(ToolOutput::success(output, start.elapsed()));
}
// Normal read
let output = serde_json::json!({
"path": doc.path,
"content": doc.content,
+37 -3
View File
@@ -224,7 +224,13 @@ impl Tool for MessageTool {
) -> Result<ToolOutput, ToolError> {
let start = std::time::Instant::now();
let content = require_str(&params, "content")?;
// Accept "message" as an alias for "content" — LLMs frequently use
// the wrong parameter name in autonomous job execution.
let content = require_str(&params, "content").or_else(|_| {
require_str(&params, "message").map_err(|_| {
ToolError::InvalidParameters("missing 'content' parameter".to_string())
})
})?;
let explicit_channel = params
.get("channel")
@@ -323,8 +329,11 @@ impl Tool for MessageTool {
if !attachments.is_empty() {
response = response.with_attachments(attachments);
}
if channel.as_deref() == Some("gateway")
&& response.thread_id.is_none()
// Attach thread_id so the gateway can route the message into the
// correct conversation. Previously this only fired when channel was
// explicitly "gateway", which meant broadcast_all (channel=null) sent
// a response without a thread_id and the gateway silently dropped it.
if response.thread_id.is_none()
&& let Some(thread_id) = metadata_string(&ctx.metadata, "notify_thread_id")
{
response = response.in_thread(thread_id);
@@ -480,6 +489,31 @@ mod tests {
assert!(params.get("attachments").is_some());
}
/// Regression: LLMs frequently pass {"message": "..."} instead of
/// {"content": "..."}. The tool should accept both.
#[tokio::test]
async fn message_param_alias_accepted() {
let tool = MessageTool::new(Arc::new(ChannelManager::new()));
tool.set_context(Some("gateway".to_string()), Some("user".to_string()))
.await;
let ctx = crate::context::JobContext::new("test", "test");
// "message" alias should not produce InvalidParameters
let result = tool
.execute(serde_json::json!({"message": "hello from alias"}), &ctx)
.await;
// Execution may fail for other reasons (no real channel), but
// the error must NOT be about a missing 'content' parameter.
if let Err(ref e) = result {
let msg = e.to_string();
assert!(
!msg.contains("missing 'content'"),
"Should accept 'message' as alias for 'content', got: {msg}"
);
}
}
#[tokio::test]
async fn message_tool_set_context_updates_defaults() {
let tool = MessageTool::new(Arc::new(ChannelManager::new()));
+68 -4
View File
@@ -65,6 +65,7 @@ struct NormalizedExecutionRequest {
context_paths: Vec<String>,
use_tools: bool,
max_tool_rounds: u32,
max_iterations: u32,
}
#[derive(Debug, Clone, PartialEq, Eq)]
@@ -328,6 +329,13 @@ fn full_job_execution_variant() -> Value {
"type": "string",
"enum": ["full_job"],
"description": "Full-job execution mode."
},
"max_iterations": {
"type": "integer",
"description": "Maximum LLM iterations for the job (default: 25). Increase for complex multi-step tasks.",
"default": 25,
"minimum": 1,
"maximum": 200
}
},
"required": ["mode"]
@@ -644,6 +652,12 @@ pub(crate) fn routine_update_parameters_schema() -> Value {
"description": {
"type": "string",
"description": "New description"
},
"max_iterations": {
"type": "integer",
"description": "Maximum LLM iterations for full_job routines (1-200).",
"minimum": 1,
"maximum": 200
}
},
"required": ["name"]
@@ -887,11 +901,16 @@ fn parse_routine_execution(
.clamp(1, crate::agent::routine::MAX_TOOL_ROUNDS_LIMIT as u64)
as u32;
let max_iterations = u64_field(params, "execution", "max_iterations", &["max_iterations"])
.unwrap_or(25)
.clamp(1, 200) as u32;
Ok(NormalizedExecutionRequest {
mode,
context_paths,
use_tools,
max_tool_rounds,
max_iterations,
})
}
@@ -972,7 +991,7 @@ fn build_routine_action(
NormalizedExecutionMode::FullJob => RoutineAction::FullJob {
title: name.to_string(),
description: prompt.to_string(),
max_iterations: 10,
max_iterations: execution.max_iterations,
},
}
}
@@ -1317,6 +1336,12 @@ impl Tool for RoutineUpdateTool {
}
}
if let Some(iters) = params.get("max_iterations").and_then(|v| v.as_u64())
&& let RoutineAction::FullJob { max_iterations, .. } = &mut routine.action
{
*max_iterations = (iters.clamp(1, 200)) as u32;
}
// Validate timezone param if provided
let new_timezone = params
.get("timezone")
@@ -1544,6 +1569,7 @@ impl Tool for RoutineFireTool {
"name": name,
"run_id": run_id.to_string(),
"status": "fired",
"note": "Routine is executing asynchronously. Use routine_history to check the result.",
});
Ok(ToolOutput::success(result, start.elapsed()))
@@ -1642,10 +1668,47 @@ impl Tool for RoutineHistoryTool {
})
.collect();
// Look up the routine's conversation thread and fetch recent messages
// so the user can see the full output of routine runs.
let (conversation_id, recent_output) = match self
.store
.get_or_create_routine_conversation(routine.id, name, &ctx.user_id)
.await
{
Ok(conv_id) => {
let messages = self
.store
.list_conversation_messages_paginated(conv_id, None, limit)
.await
.map(|(msgs, _)| msgs)
.unwrap_or_default();
let msg_list: Vec<serde_json::Value> = messages
.iter()
.map(|m| {
serde_json::json!({
"role": m.role,
"content": m.content,
"timestamp": m.created_at.to_rfc3339(),
})
})
.collect();
(Some(conv_id.to_string()), msg_list)
}
Err(e) => {
tracing::warn!(
routine = %name,
"Failed to fetch routine conversation thread: {e}"
);
(None, Vec::new())
}
};
let result = serde_json::json!({
"routine": name,
"total_runs": routine.run_count,
"conversation_id": conversation_id,
"runs": run_list,
"recent_output": recent_output,
});
Ok(ToolOutput::success(result, start.elapsed()))
@@ -2282,8 +2345,8 @@ mod tests {
.and_then(Value::as_object)
.expect("full_job properties");
assert!(
full_job_props.len() == 1 && full_job_props.contains_key("mode"),
"full_job variant should only expose the execution mode",
full_job_props.contains_key("mode") && full_job_props.contains_key("max_iterations"),
"full_job variant should expose mode and max_iterations",
);
}
@@ -2491,6 +2554,7 @@ mod tests {
context_paths: Vec::new(),
use_tools: false,
max_tool_rounds: 3,
max_iterations: 25,
};
let action = build_routine_action("issue-1316", "Run it", &execution);
@@ -2503,7 +2567,7 @@ mod tests {
max_iterations,
} if title == "issue-1316"
&& description == "Run it"
&& max_iterations == 10
&& max_iterations == 25
));
}
}
+6 -3
View File
@@ -5,7 +5,7 @@ use chrono::{DateTime, LocalResult, NaiveDate, NaiveDateTime, TimeZone, Utc};
use chrono_tz::Tz;
use crate::context::JobContext;
use crate::tools::tool::{Tool, ToolError, ToolOutput, require_str};
use crate::tools::tool::{Tool, ToolError, ToolOutput};
/// Tool for getting current time and date operations.
pub struct TimeTool;
@@ -62,7 +62,7 @@ impl Tool for TimeTool {
"description": "Second timestamp for diff."
}
},
"required": ["operation"]
"required": []
})
}
@@ -73,7 +73,10 @@ impl Tool for TimeTool {
) -> Result<ToolOutput, ToolError> {
let start = std::time::Instant::now();
let operation = require_str(&params, "operation")?;
let operation = params
.get("operation")
.and_then(|v| v.as_str())
.unwrap_or("now");
let result = match operation {
"now" => execute_now(&params, ctx)?,
+28 -7
View File
@@ -954,16 +954,22 @@ pub async fn store_tokens(
server_config: &McpServerConfig,
token: &AccessToken,
) -> Result<(), AuthError> {
// Store access token
let params = CreateSecretParams::new(server_config.token_secret_name(), &token.access_token)
.with_provider(format!("mcp:{}", server_config.name));
// Store access token (with expiry if provided)
let mut params =
CreateSecretParams::new(server_config.token_secret_name(), &token.access_token)
.with_provider(format!("mcp:{}", server_config.name));
if let Some(secs) = token.expires_in {
let expires_at = chrono::Utc::now() + chrono::Duration::seconds(secs as i64);
params = params.with_expiry(expires_at);
}
secrets
.create(user_id, params)
.await
.map_err(|e| AuthError::Secrets(e.to_string()))?;
// Store refresh token if present
// Store refresh token if present (no expiry — long-lived)
if let Some(ref refresh_token) = token.refresh_token {
let params =
CreateSecretParams::new(server_config.refresh_token_secret_name(), refresh_token)
@@ -1064,11 +1070,26 @@ pub async fn refresh_access_token(
// Get client_id (from config or stored DCR)
let client_id = get_client_id(server_config, secrets, user_id).await?;
// Get the refresh token
let refresh_token = secrets
// Get the refresh token (try current name, fall back to legacy name for
// users who authenticated before the naming convention was fixed).
// Only fall back on NotFound/Expired — propagate real errors (DB, decryption).
let refresh_token = match secrets
.get_decrypted(user_id, &server_config.refresh_token_secret_name())
.await
.map_err(|e| AuthError::RefreshFailed(format!("No refresh token: {}", e)))?;
{
Ok(token) => token,
Err(crate::secrets::SecretError::NotFound(_) | crate::secrets::SecretError::Expired) => {
secrets
.get_decrypted(user_id, &server_config.legacy_refresh_token_secret_name())
.await
.map_err(|e| AuthError::RefreshFailed(format!("No refresh token: {}", e)))?
}
Err(e) => {
return Err(AuthError::RefreshFailed(format!(
"Failed to read refresh token: {e}"
)));
}
};
// Discover the token endpoint
let token_url = if let Some(ref oauth) = server_config.oauth {
+30
View File
@@ -259,6 +259,9 @@ impl McpClient {
}
/// Get the access token for this server (if authenticated).
///
/// If the stored token has expired, automatically attempts a refresh using
/// the stored refresh token before failing.
async fn get_access_token(&self) -> Result<Option<String>, ToolError> {
let Some(ref secrets) = self.secrets else {
return Ok(None);
@@ -272,6 +275,33 @@ impl McpClient {
{
Ok(token) => Ok(Some(token.expose().to_string())),
Err(crate::secrets::SecretError::NotFound(_)) => Ok(None),
Err(crate::secrets::SecretError::Expired) => {
// Token expired — attempt refresh before failing.
tracing::info!(
server = %self.server_name,
"Access token expired, attempting refresh"
);
match refresh_access_token(config, secrets, &self.user_id).await {
Ok(new_token) => {
tracing::info!(
server = %self.server_name,
"Access token refreshed successfully"
);
Ok(Some(new_token.access_token))
}
Err(e) => {
tracing::warn!(
server = %self.server_name,
"Token refresh failed: {}", e
);
Err(ToolError::ExternalService(format!(
"Failed to get access token: Secret has expired \
and refresh failed: {}",
e
)))
}
}
}
Err(e) => Err(ToolError::ExternalService(format!(
"Failed to get access token: {}",
e
+19
View File
@@ -250,7 +250,19 @@ impl McpServerConfig {
}
/// Get the secret name used to store the refresh token.
///
/// Matches the convention used by the hosted OAuth flow in
/// `store_oauth_tokens`: `{token_secret_name}_refresh_token`.
pub fn refresh_token_secret_name(&self) -> String {
format!("{}_refresh_token", self.token_secret_name())
}
/// Legacy secret name for refresh tokens (pre-v0.22).
///
/// Earlier versions stored refresh tokens as `mcp_{name}_refresh_token`
/// instead of `{token_secret_name}_refresh_token`. Used as a fallback
/// during lookup to avoid forcing re-auth on existing users.
pub fn legacy_refresh_token_secret_name(&self) -> String {
format!("mcp_{}_refresh_token", self.name)
}
@@ -750,8 +762,15 @@ mod tests {
fn test_token_secret_names() {
let config = McpServerConfig::new("notion", "https://mcp.notion.com");
assert_eq!(config.token_secret_name(), "mcp_notion_access_token");
// Refresh token name follows the hosted OAuth convention:
// {token_secret_name}_refresh_token
assert_eq!(
config.refresh_token_secret_name(),
"mcp_notion_access_token_refresh_token"
);
// Legacy name used before v0.22 — fallback lookup prevents forced re-auth
assert_eq!(
config.legacy_refresh_token_secret_name(),
"mcp_notion_refresh_token"
);
}
+22
View File
@@ -225,4 +225,26 @@ mod tests {
"The tool returned: TASK_COMPLETE signal"
));
}
#[test]
fn signals_completion_after_suggestions_stripped() {
// Regression: after stripping <suggestions> tags, the completion
// signal should still be detected in the cleaned text.
assert!(llm_signals_completion(
"The job is complete. All requested work has been finished."
));
}
#[test]
fn signals_completion_self_dialogue_pattern() {
// Regression: the "not complete" pattern that caused the self-dialogue
// loop when left in job context after plan completion.
assert!(!llm_signals_completion(
"No — the job is **not complete**.\n\n\
What still needs to be done:\n\
1. Fetch actual meeting note contents\n\
2. Create the Notion page\n\
3. Send the completion message"
));
}
}
+123 -23
View File
@@ -828,14 +828,11 @@ Report when the job is complete or if you encounter issues you cannot resolve."#
}),
);
if matches!(
&e,
Error::Tool(crate::error::ToolError::AutonomousUnavailable { .. })
) {
Err(e)
} else {
Ok(())
}
// All tool errors (including AutonomousUnavailable) are
// recoverable — the error message is already recorded in
// reason_ctx so the LLM can see it and try a different
// approach. Returning Err here would kill the entire job.
Ok(())
}
}
}
@@ -930,17 +927,31 @@ Report when the job is complete or if you encounter issues you cannot resolve."#
tokio::time::sleep(Duration::from_millis(100)).await;
}
// Plan completed, check with LLM if job is done
// Plan completed — ask the LLM whether the job is done.
let msg_count_before = reason_ctx.messages.len();
reason_ctx.messages.push(ChatMessage::user(
"All planned actions have been executed. Is the job complete? If not, what else needs to be done?",
"All planned actions have been executed. Assess the results: \
if the job is fully complete, state that the job is complete. \
Otherwise, briefly list what remains.",
));
let response = reasoning.respond(reason_ctx).await?;
reason_ctx.messages.push(ChatMessage::assistant(&response));
let response = crate::agent::strip_suggestions(&response);
if crate::util::llm_signals_completion(&response) {
reason_ctx.messages.push(ChatMessage::assistant(&response));
self.mark_completed().await?;
} else {
// Replace the completion-check exchange with an action-oriented
// continuation prompt. Leaving the "Is the job complete?" / "No"
// dialogue in context causes the agentic loop to repeat the same
// analysis instead of calling tools (self-dialogue loop).
reason_ctx.messages.truncate(msg_count_before);
reason_ctx.messages.push(ChatMessage::user(format!(
"The planned actions are done but the job is not yet complete. \
Remaining work:\n\n{response}\n\n\
Continue executing now use tools to finish the job."
)));
tracing::info!(
"Job {} plan completed but work remains, falling back to direct selection",
self.job_id
@@ -1420,16 +1431,20 @@ impl<'a> LoopDelegate for JobDelegate<'a> {
return TextAction::Continue;
}
// Check for explicit completion
if crate::util::llm_signals_completion(text) {
if let Err(e) = self.worker.mark_completed().await {
tracing::warn!(
"Failed to mark job {} as completed: {}",
self.worker.job_id,
e
);
}
return TextAction::Return(LoopOutcome::Response(text.to_string()));
// Jobs run autonomously — strip <suggestions> tags that are only
// meaningful for interactive chat sessions.
let text = crate::agent::strip_suggestions(text);
// A non-empty text response with no tool intent (already filtered
// by the agentic loop's nudge mechanism) is the LLM's final answer.
// Mark the job complete and stop the loop. Without this, the LLM
// restates its summary every iteration until the cap is hit.
if let Err(e) = self.worker.mark_completed().await {
tracing::warn!(
"Failed to mark job {} as completed: {}",
self.worker.job_id,
e
);
}
// Track that a substantive response has been produced.
@@ -1437,7 +1452,7 @@ impl<'a> LoopDelegate for JobDelegate<'a> {
.store(true, std::sync::atomic::Ordering::Relaxed);
// Add assistant response to context
reason_ctx.messages.push(ChatMessage::assistant(text));
reason_ctx.messages.push(ChatMessage::assistant(&text));
self.worker.log_event(
"message",
@@ -1447,7 +1462,7 @@ impl<'a> LoopDelegate for JobDelegate<'a> {
}),
);
TextAction::Continue
TextAction::Return(LoopOutcome::Response(text))
}
async fn execute_tool_calls(
@@ -1456,6 +1471,9 @@ impl<'a> LoopDelegate for JobDelegate<'a> {
content: Option<String>,
reason_ctx: &mut ReasoningContext,
) -> Result<Option<LoopOutcome>, crate::error::Error> {
// Strip suggestions from accompanying text (not useful in job context).
let content = content.map(|c| crate::agent::strip_suggestions(&c));
if let Some(ref text) = content {
self.worker.log_event(
"message",
@@ -2156,6 +2174,53 @@ mod tests {
);
}
/// Regression: a text response without rigid completion phrases (e.g.
/// "Weekly review completed and saved to Notion") must still terminate the
/// agentic loop and mark the job complete, rather than continuing until
/// max_iterations.
#[tokio::test]
async fn test_text_response_terminates_loop_without_explicit_completion_phrase() {
let worker = make_worker(vec![]).await;
worker
.context_manager()
.update_context(worker.job_id, |ctx| {
ctx.transition_to(JobState::InProgress, None)
})
.await
.unwrap() // safety: test
.unwrap(); // safety: test
let (_, mut rx) = tokio::sync::mpsc::channel(1);
let delegate = JobDelegate {
worker: &worker,
rx: tokio::sync::Mutex::new(&mut rx),
consecutive_rate_limits: std::sync::atomic::AtomicUsize::new(0),
has_text_response: std::sync::atomic::AtomicBool::new(false),
};
let mut reason_ctx = ReasoningContext::new();
// Text that a real LLM would produce but doesn't match llm_signals_completion
let action = delegate
.handle_text_response(
"Weekly review created in Notion and notification sent.",
&mut reason_ctx,
)
.await;
assert!(
matches!(action, TextAction::Return(_)),
"Text response should terminate the loop, got Continue"
); // safety: test
let ctx = worker
.context_manager()
.get_context(worker.job_id)
.await
.unwrap(); // safety: test
assert_eq!(ctx.state, JobState::Completed); // safety: test
}
/// Regression test: selections_to_tool_calls must preserve tool_call_id
/// so that tool_result messages match the assistant_with_tool_calls message
/// and are not treated as orphaned by sanitize_tool_messages.
@@ -2429,4 +2494,39 @@ mod tests {
}
));
}
/// Regression test: AutonomousUnavailable errors must be recoverable.
/// Previously the job worker treated them as fatal, killing the entire
/// job instead of feeding the error back to the LLM.
#[tokio::test]
async fn test_autonomous_unavailable_is_recoverable() {
let worker = make_worker(vec![]).await;
let mut reason_ctx = ReasoningContext::new();
let selection = ToolSelection {
tool_name: "secret_list".to_string(),
parameters: serde_json::json!({}),
reasoning: "list secrets".to_string(),
alternatives: vec![],
tool_call_id: "call_123".to_string(),
};
let err = Error::Tool(crate::error::ToolError::AutonomousUnavailable {
name: "secret_list".to_string(),
reason: "not available in autonomous jobs".to_string(),
});
let result = worker
.process_tool_result_job(&mut reason_ctx, &selection, Err(err))
.await;
assert!(
result.is_ok(),
"AutonomousUnavailable must be recoverable, not fatal: {:?}",
result
);
// The error should be fed back to the LLM as a message.
assert!(
!reason_ctx.messages.is_empty(),
"Error message should be added to reason_ctx for the LLM"
);
}
}
-248
View File
@@ -2,7 +2,6 @@
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use uuid::Uuid;
/// Well-known document paths.
@@ -38,139 +37,6 @@ pub mod paths {
pub const ASSISTANT_DIRECTIVES: &str = "context/assistant-directives.md";
}
/// Name of the folder-level configuration document.
///
/// A document at `{directory}/.config` carries metadata flags that apply
/// as defaults to all documents in that directory (e.g., `skip_indexing`,
/// `hygiene` settings). Individual document metadata overrides folder defaults.
pub const CONFIG_FILE_NAME: &str = ".config";
/// Typed overlay for the `metadata` JSON field on [`MemoryDocument`].
///
/// Fields use `Option` so that only explicitly set flags participate in
/// the merge chain (document metadata → folder `.config` → system defaults).
/// Unknown fields are preserved via `serde(flatten)`.
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
pub struct DocumentMetadata {
/// When `true`, skip chunking and embedding for this document/folder.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub skip_indexing: Option<bool>,
/// When `true`, skip automatic versioning for this document/folder.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub skip_versioning: Option<bool>,
/// Hygiene (auto-cleanup) configuration for this folder.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub hygiene: Option<HygieneMetadata>,
/// Preserve unknown fields for forward compatibility.
#[serde(flatten)]
pub extra: serde_json::Map<String, serde_json::Value>,
}
impl DocumentMetadata {
/// Parse from a raw JSON [`serde_json::Value`].
///
/// Returns [`Default`] if the value is not an object or cannot be parsed.
pub fn from_value(value: &serde_json::Value) -> Self {
serde_json::from_value(value.clone()).unwrap_or_default()
}
/// Convert to a JSON [`serde_json::Value`].
pub fn to_value(&self) -> serde_json::Value {
serde_json::to_value(self).unwrap_or(serde_json::json!({}))
}
/// Merge two metadata values: `overlay` keys win over `base` keys.
///
/// This is a shallow merge at the top-level keys — nested objects are
/// replaced wholesale, not recursively merged. This keeps the semantics
/// simple and predictable across both PostgreSQL and libSQL.
pub fn merge(base: &serde_json::Value, overlay: &serde_json::Value) -> serde_json::Value {
let mut merged = match base {
serde_json::Value::Object(map) => map.clone(),
_ => serde_json::Map::new(),
};
if let serde_json::Value::Object(over) = overlay {
for (k, v) in over {
merged.insert(k.clone(), v.clone());
}
}
serde_json::Value::Object(merged)
}
}
/// Hygiene (auto-cleanup) settings for a folder.
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct HygieneMetadata {
/// Whether this folder is a hygiene target.
pub enabled: bool,
/// Delete documents older than this many days.
#[serde(default = "default_retention_days")]
pub retention_days: u32,
}
fn default_retention_days() -> u32 {
30
}
/// A historical version of a workspace document.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DocumentVersion {
/// Version record ID.
pub id: Uuid,
/// Parent document ID.
pub document_id: Uuid,
/// Version number (1-based, monotonically increasing per document).
pub version: i32,
/// Full document content at this version.
pub content: String,
/// SHA-256 hash of `content` (hex-encoded, prefixed with `sha256:`).
pub content_hash: String,
/// When this version was created.
pub created_at: DateTime<Utc>,
/// Who/what created this version (e.g. `"agent"`, `"user:alice"`).
pub changed_by: Option<String>,
}
/// Summary of a document version (without full content).
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct VersionSummary {
/// Version number.
pub version: i32,
/// SHA-256 hash of the version's content.
pub content_hash: String,
/// When this version was created.
pub created_at: DateTime<Utc>,
/// Who/what created this version.
pub changed_by: Option<String>,
}
/// Result of a workspace patch operation.
#[derive(Debug, Clone)]
pub struct PatchResult {
/// The updated document.
pub document: MemoryDocument,
/// Number of replacements made.
pub replacements: usize,
}
/// Compute a SHA-256 hash of content, returned as `"sha256:{hex}"`.
pub fn content_sha256(content: &str) -> String {
let mut hasher = Sha256::new();
hasher.update(content.as_bytes());
let result = hasher.finalize();
format!("sha256:{:x}", result)
}
/// Check if a path refers to a `.config` document.
pub fn is_config_path(path: &str) -> bool {
let file_name = path.rsplit('/').next().unwrap_or(path);
file_name == CONFIG_FILE_NAME
}
/// Paths treated as identity documents for multi-scope isolation.
///
/// These files are always read from the primary scope only — never from
@@ -494,120 +360,6 @@ mod tests {
assert_eq!(result[0].updated_at, Some(ts));
}
#[test]
fn test_document_metadata_default_is_empty() {
let meta = DocumentMetadata::default();
assert_eq!(meta.skip_indexing, None);
assert_eq!(meta.skip_versioning, None);
assert_eq!(meta.hygiene, None);
assert!(meta.extra.is_empty());
}
#[test]
fn test_document_metadata_from_value_full() {
let value = serde_json::json!({
"skip_indexing": true,
"skip_versioning": false,
"hygiene": { "enabled": true, "retention_days": 7 }
});
let meta = DocumentMetadata::from_value(&value);
assert_eq!(meta.skip_indexing, Some(true));
assert_eq!(meta.skip_versioning, Some(false));
let hygiene = meta.hygiene.unwrap();
assert!(hygiene.enabled);
assert_eq!(hygiene.retention_days, 7);
}
#[test]
fn test_document_metadata_from_value_partial() {
let value = serde_json::json!({"skip_indexing": true});
let meta = DocumentMetadata::from_value(&value);
assert_eq!(meta.skip_indexing, Some(true));
assert_eq!(meta.hygiene, None);
}
#[test]
fn test_document_metadata_from_value_invalid() {
let meta = DocumentMetadata::from_value(&serde_json::json!("not an object"));
assert_eq!(meta, DocumentMetadata::default());
}
#[test]
fn test_document_metadata_preserves_unknown_fields() {
let value = serde_json::json!({
"skip_indexing": true,
"custom_field": "hello"
});
let meta = DocumentMetadata::from_value(&value);
assert_eq!(meta.skip_indexing, Some(true));
assert_eq!(
meta.extra.get("custom_field").and_then(|v| v.as_str()),
Some("hello")
);
// Round-trip preserves the field
let back = meta.to_value();
assert_eq!(
back.get("custom_field").and_then(|v| v.as_str()),
Some("hello")
);
}
#[test]
fn test_document_metadata_merge() {
let base = serde_json::json!({"skip_indexing": false, "hygiene": {"enabled": true, "retention_days": 30}});
let overlay = serde_json::json!({"skip_indexing": true, "skip_versioning": true});
let merged = DocumentMetadata::merge(&base, &overlay);
let meta = DocumentMetadata::from_value(&merged);
// Overlay wins
assert_eq!(meta.skip_indexing, Some(true));
assert_eq!(meta.skip_versioning, Some(true));
// Base preserved when not overridden
assert!(meta.hygiene.is_some());
}
#[test]
fn test_document_metadata_merge_empty_base() {
let base = serde_json::json!({});
let overlay = serde_json::json!({"skip_indexing": true});
let merged = DocumentMetadata::merge(&base, &overlay);
let meta = DocumentMetadata::from_value(&merged);
assert_eq!(meta.skip_indexing, Some(true));
}
#[test]
fn test_hygiene_metadata_default_retention() {
let value = serde_json::json!({"enabled": true});
let hygiene: HygieneMetadata = serde_json::from_value(value).unwrap();
assert!(hygiene.enabled);
assert_eq!(hygiene.retention_days, 30);
}
#[test]
fn test_content_sha256_deterministic() {
let hash1 = content_sha256("hello world");
let hash2 = content_sha256("hello world");
assert_eq!(hash1, hash2);
assert!(hash1.starts_with("sha256:"));
}
#[test]
fn test_content_sha256_different_content() {
let hash1 = content_sha256("hello");
let hash2 = content_sha256("world");
assert_ne!(hash1, hash2);
}
#[test]
fn test_is_config_path() {
assert!(is_config_path(".config"));
assert!(is_config_path("daily/.config"));
assert!(is_config_path("frontend/widgets/.config"));
assert!(!is_config_path("daily/2024-01-15.md"));
assert!(!is_config_path("MEMORY.md"));
assert!(!is_config_path(".config.bak"));
}
#[test]
fn test_merge_workspace_entries_sorted_by_path() {
let entries = vec![
+248 -184
View File
@@ -1,10 +1,8 @@
//! Memory hygiene: automatic cleanup of stale workspace documents.
//!
//! Runs on a configurable cadence and discovers which directories have hygiene
//! enabled by reading `.config` metadata documents. This is a **metadata-driven**
//! approach: instead of hardcoding `daily/` and `conversations/`, the system
//! respects `hygiene.enabled` and `hygiene.retention_days` set on each folder's
//! `.config` document.
//! Runs on a configurable cadence and deletes daily log entries and conversation
//! documents older than their respective retention periods. Identity files
//! (`IDENTITY.md`, `SOUL.md`, etc.) are never touched.
//!
//! A global [`AtomicBool`] guard prevents concurrent hygiene passes, which
//! avoids TOCTOU races on the state file and Windows file-locking errors
@@ -12,16 +10,18 @@
//! pass completes.
//!
//! ```text
//! ┌──────────────────────────────────────────────────
//! │ Hygiene Pass
//! │
//! │ 0. Acquire RUNNING guard (skip if held)
//! │ 1. Check cadence (skip if ran recently)
//! │ 2. Save state (claim the cadence window)
//! │ 3. Discover .config docs with hygiene.enabled
//! │ 4. For each: cleanup_directory(parent, retention)
//! │ 5. Log summary
//! └──────────────────────────────────────────────────┘
//! ┌─────────────────────────────────────────────┐
//! │ Hygiene Pass │
//! │ │
//! │ 0. Acquire RUNNING guard (skip if held) │
//! │ 1. Check cadence (skip if ran recently) │
//! │ 2. Save state (claim the cadence window) │
//! │ 3. List daily/ documents
//! │ 4. Delete those older than daily_retention
//! │ 5. List conversations/ documents
//! │ 6. Delete those older than conversation_ret │
//! │ 7. Log summary │
//! └─────────────────────────────────────────────┘
//! ```
use std::path::PathBuf;
@@ -31,22 +31,46 @@ use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::bootstrap::ironclaw_base_dir;
use crate::workspace::{DocumentMetadata, Workspace, is_config_path};
use crate::workspace::Workspace;
/// Global guard preventing concurrent hygiene passes.
static RUNNING: AtomicBool = AtomicBool::new(false);
/// Paths that must never be deleted by hygiene, regardless of age.
const IDENTITY_PATHS: &[&str] = &[
crate::workspace::document::paths::MEMORY,
crate::workspace::document::paths::IDENTITY,
crate::workspace::document::paths::SOUL,
crate::workspace::document::paths::AGENTS,
crate::workspace::document::paths::USER,
crate::workspace::document::paths::HEARTBEAT,
crate::workspace::document::paths::README,
crate::workspace::document::paths::TOOLS,
crate::workspace::document::paths::BOOTSTRAP,
];
/// Check if a document path is an identity document that must never be deleted.
///
/// Performs case-insensitive comparison to handle case-insensitive filesystems
/// (Windows, macOS) and prevent accidental deletion of identity docs with
/// different casing (e.g., memory.md, MEMORY.MD, Memory.md).
fn is_identity_path(path: &str) -> bool {
let file_name = path.rsplit('/').next().unwrap_or(path);
let file_name_lower = file_name.to_lowercase();
IDENTITY_PATHS
.iter()
.any(|&p| p.to_lowercase() == file_name_lower)
}
/// Configuration for workspace hygiene.
#[derive(Debug, Clone)]
pub struct HygieneConfig {
/// Whether hygiene is enabled at all.
pub enabled: bool,
/// Maximum number of versions to keep per document.
///
/// TODO: Wire up global version pruning once per-document iteration
/// is efficient (e.g., via a dedicated DB query). For now this field
/// is stored in config but not actively enforced during hygiene passes.
pub version_keep_count: u32,
/// Documents in `daily/` older than this many days are deleted.
pub daily_retention_days: u32,
/// Documents in `conversations/` older than this many days are deleted.
pub conversation_retention_days: u32,
/// Minimum hours between hygiene passes.
pub cadence_hours: u32,
/// Directory to store state file (default: `~/.ironclaw`).
@@ -57,7 +81,8 @@ impl Default for HygieneConfig {
fn default() -> Self {
Self {
enabled: true,
version_keep_count: 50,
daily_retention_days: 30,
conversation_retention_days: 7,
cadence_hours: 12,
state_dir: ironclaw_base_dir(),
}
@@ -73,10 +98,10 @@ struct HygieneState {
/// Summary of what a hygiene pass cleaned up.
#[derive(Debug, Default)]
pub struct HygieneReport {
/// Per-directory cleanup results: `(directory_path, deleted_count)`.
pub directories_cleaned: Vec<(String, u32)>,
/// Number of document versions pruned across all documents.
pub versions_pruned: u64,
/// Number of daily log documents deleted.
pub daily_logs_deleted: u32,
/// Number of conversation documents deleted.
pub conversation_docs_deleted: u32,
/// Whether the run was skipped (cadence not yet elapsed).
pub skipped: bool,
}
@@ -84,7 +109,7 @@ pub struct HygieneReport {
impl HygieneReport {
/// True if any cleanup work was done.
pub fn had_work(&self) -> bool {
self.directories_cleaned.iter().any(|(_, n)| *n > 0) || self.versions_pruned > 0
self.daily_logs_deleted > 0 || self.conversation_docs_deleted > 0
}
}
@@ -143,51 +168,30 @@ pub async fn run_if_due(workspace: &Workspace, config: &HygieneConfig) -> Hygien
// TOCTOU races where another task reads stale state.
save_state(&state_file);
tracing::info!("memory hygiene: starting cleanup pass");
tracing::info!(
daily_retention_days = config.daily_retention_days,
conversation_retention_days = config.conversation_retention_days,
"memory hygiene: starting cleanup pass"
);
let mut report = HygieneReport::default();
// Discover directories that have hygiene enabled via .config metadata.
let config_docs = match workspace.find_config_documents().await {
Ok(docs) => docs,
Err(e) => {
tracing::warn!("memory hygiene: failed to discover .config documents: {e}");
return report;
}
};
// Delete old daily logs
match cleanup_daily_logs(workspace, config.daily_retention_days).await {
Ok(count) => report.daily_logs_deleted = count,
Err(e) => tracing::warn!("memory hygiene: failed to clean daily logs: {e}"),
}
for doc in &config_docs {
let meta = DocumentMetadata::from_value(&doc.metadata);
let Some(hygiene) = meta.hygiene else {
continue;
};
if !hygiene.enabled {
continue;
}
// Derive the parent directory from the .config path.
let directory = match doc.path.rsplit_once('/') {
Some((dir, _)) => format!("{dir}/"),
None => continue, // root-level .config — skip
};
match cleanup_directory(workspace, &directory, hygiene.retention_days).await {
Ok(deleted) => {
if deleted > 0 {
tracing::info!(directory, deleted, "memory hygiene: cleaned directory");
}
report.directories_cleaned.push((directory, deleted));
}
Err(e) => {
tracing::warn!(directory, "memory hygiene: failed to clean directory: {e}");
}
}
// Delete old conversation documents
match cleanup_conversation_docs(workspace, config.conversation_retention_days).await {
Ok(count) => report.conversation_docs_deleted = count,
Err(e) => tracing::warn!("memory hygiene: failed to clean conversation docs: {e}"),
}
if report.had_work() {
tracing::info!(
directories_cleaned = ?report.directories_cleaned,
versions_pruned = report.versions_pruned,
daily_logs_deleted = report.daily_logs_deleted,
conversation_docs_deleted = report.conversation_docs_deleted,
"memory hygiene: cleanup complete"
);
} else {
@@ -206,41 +210,88 @@ impl Drop for RunningGuard {
}
}
/// Delete documents in `directory` that are older than `retention_days`.
///
/// Skips directories and `.config` files (which must never be deleted by
/// hygiene). Returns the number of documents deleted.
async fn cleanup_directory(
/// Delete daily log documents older than `retention_days`.
async fn cleanup_daily_logs(
workspace: &Workspace,
directory: &str,
retention_days: u32,
) -> Result<u32, anyhow::Error> {
let cutoff = Utc::now() - chrono::Duration::days(i64::from(retention_days));
let entries = workspace.list(directory).await?;
let entries = workspace.list("daily/").await?;
let mut deleted = 0u32;
for entry in entries {
if entry.is_directory {
continue;
}
if is_config_path(&entry.path) {
// Never delete identity documents
if is_identity_path(&entry.path) {
continue;
}
// Check if the document is old enough to delete
if let Some(updated_at) = entry.updated_at
&& updated_at < cutoff
{
let path = if entry.path.starts_with(directory) {
let path = if entry.path.starts_with("daily/") {
entry.path.clone()
} else {
format!("{}{}", directory, entry.path)
format!("daily/{}", entry.path)
};
if let Err(e) = workspace.delete(&path).await {
tracing::warn!(path, "memory hygiene: failed to delete: {e}");
} else {
tracing::debug!(path, "memory hygiene: deleted stale document");
tracing::debug!(path, "memory hygiene: deleted old daily log");
deleted += 1;
}
}
}
Ok(deleted)
}
/// Delete conversation documents older than `retention_days`.
async fn cleanup_conversation_docs(
workspace: &Workspace,
retention_days: u32,
) -> Result<u32, anyhow::Error> {
let cutoff = Utc::now() - chrono::Duration::days(i64::from(retention_days));
let entries = workspace.list("conversations/").await?;
let mut deleted = 0u32;
for entry in entries {
if entry.is_directory {
continue;
}
// Never delete identity documents
if is_identity_path(&entry.path) {
continue;
}
// Check if the document is old enough to delete
if let Some(updated_at) = entry.updated_at
&& updated_at < cutoff
{
let path = if entry.path.starts_with("conversations/") {
entry.path.clone()
} else {
format!("conversations/{}", entry.path)
};
if let Err(e) = workspace.delete(&path).await {
tracing::warn!(
path,
"memory hygiene: failed to delete conversation doc: {e}"
);
} else {
tracing::debug!(path, "memory hygiene: deleted old conversation doc");
deleted += 1;
}
}
}
Ok(deleted)
}
@@ -298,7 +349,8 @@ mod tests {
fn default_config_is_reasonable() {
let cfg = HygieneConfig::default();
assert!(cfg.enabled);
assert_eq!(cfg.version_keep_count, 50);
assert_eq!(cfg.daily_retention_days, 30);
assert_eq!(cfg.conversation_retention_days, 7);
assert_eq!(cfg.cadence_hours, 12);
}
@@ -310,33 +362,84 @@ mod tests {
}
#[test]
fn report_had_work_when_directories_cleaned() {
fn report_had_work_when_deleted() {
let report = HygieneReport {
directories_cleaned: vec![("daily/".to_string(), 3)],
versions_pruned: 0,
daily_logs_deleted: 3,
conversation_docs_deleted: 0,
skipped: false,
};
assert!(report.had_work());
}
#[test]
fn report_had_work_when_versions_pruned() {
fn report_had_work_when_conversation_deleted() {
let report = HygieneReport {
directories_cleaned: vec![],
versions_pruned: 5,
daily_logs_deleted: 0,
conversation_docs_deleted: 2,
skipped: false,
};
assert!(report.had_work());
}
#[test]
fn report_no_work_when_zero_deletions() {
let report = HygieneReport {
directories_cleaned: vec![("daily/".to_string(), 0)],
versions_pruned: 0,
skipped: false,
};
assert!(!report.had_work());
fn is_identity_path_excludes_sacred_docs() {
for name in [
"MEMORY.md",
"IDENTITY.md",
"SOUL.md",
"AGENTS.md",
"USER.md",
"HEARTBEAT.md",
"README.md",
"TOOLS.md",
"BOOTSTRAP.md",
] {
assert!(is_identity_path(name), "{name} should be excluded");
assert!(
is_identity_path(&format!("conversations/{name}")),
"conversations/{name} should be excluded via path"
);
}
}
#[test]
fn is_identity_path_case_insensitive() {
// Verify case-insensitive matching for case-insensitive filesystems
assert!(
is_identity_path("memory.md"),
"lowercase memory.md should be excluded"
);
assert!(
is_identity_path("Memory.md"),
"mixed case Memory.md should be excluded"
);
assert!(
is_identity_path("MEMORY.MD"),
"uppercase MEMORY.MD should be excluded"
);
assert!(
is_identity_path("identity.md"),
"lowercase identity.md should be excluded"
);
assert!(
is_identity_path("conversations/soul.md"),
"conversations/soul.md should be excluded"
);
assert!(
is_identity_path("conversations/SOUL.MD"),
"conversations/SOUL.MD should be excluded"
);
}
#[test]
fn is_identity_path_allows_normal_docs() {
for path in [
"daily/2024-01-01.md",
"conversations/chat-abc.md",
"notes.md",
] {
assert!(!is_identity_path(path), "{path} should not be excluded");
}
}
#[test]
@@ -449,105 +552,61 @@ mod tests {
Arc::new(Workspace::new_with_db("default", db.clone()))
}
/// Helper to seed a .config document with hygiene metadata on a directory.
async fn seed_hygiene_config(workspace: &Workspace, directory: &str, retention_days: u32) {
let config_path = format!("{}.config", directory);
// Create the .config document with empty content
workspace
.write(&config_path, "")
.await
.expect("write .config");
// Read back to get the document ID
let doc = workspace
.read(&config_path)
.await
.expect("read .config doc");
// Set hygiene metadata
workspace
.update_metadata(
doc.id,
&serde_json::json!({
"hygiene": {"enabled": true, "retention_days": retention_days},
"skip_versioning": true
}),
)
.await
.expect("set metadata");
}
#[tokio::test]
async fn cleanup_directory_skips_config_files() {
async fn cleanup_daily_logs_preserves_identity_documents() {
let (db, _tmp) = create_test_db().await;
let ws = create_workspace(&db);
// Write documents including a .config
// Write several regular documents (non-identity)
ws.write("daily/2024-01-15.md", "Old log")
.await
.expect("write log");
ws.write("daily/.config", "").await.expect("write config");
.expect("write log 1");
ws.write("daily/2024-01-20.md", "Another log")
.await
.expect("write log 2");
// Write an identity document
ws.write("MEMORY.md", "Long-term curated memory")
.await
.expect("write identity");
// List before cleanup
let before = ws.list("daily/").await.expect("list before");
let daily_count_before = before.iter().filter(|e| !e.is_directory).count();
assert!(daily_count_before >= 2, "should have at least 2 daily logs");
// Run cleanup with 0-day retention (deletes everything old)
let deleted = cleanup_directory(&ws, "daily/", 0)
// This tests that even with aggressive cleanup, identity docs survive
let deleted = cleanup_daily_logs(&ws, 0)
.await
.expect("cleanup_directory");
.expect("cleanup_daily_logs");
// Should have deleted the log but not the .config
// Should have deleted some documents (the daily logs)
assert!(deleted > 0, "should have deleted old daily documents");
// Verify .config still exists
let config_doc = db
.get_document_by_path("default", None, "daily/.config")
// Verify identity doc still exists
let identity = db
.get_document_by_path("default", None, "MEMORY.md")
.await
.expect("get .config doc");
assert_eq!(config_doc.path, "daily/.config");
.expect("get identity doc");
assert_eq!(identity.path, "MEMORY.md");
assert_eq!(identity.content, "Long-term curated memory");
}
#[tokio::test]
async fn cleanup_directory_handles_empty_directory() {
async fn cleanup_conversation_docs_handles_empty_directory() {
let (db, _tmp) = create_test_db().await;
let ws = create_workspace(&db);
// Run cleanup on an empty directory
let deleted = cleanup_directory(&ws, "conversations/", 7)
// Run cleanup on an empty directory (conversations/ doesn't exist)
let deleted = cleanup_conversation_docs(&ws, 7)
.await
.expect("cleanup_directory");
.expect("cleanup_conversation_docs");
// Should delete 0 (nothing to delete)
assert_eq!(deleted, 0, "should delete 0 from empty directory");
}
#[tokio::test]
async fn metadata_driven_cleanup_discovers_directories() {
let (db, _tmp) = create_test_db().await;
let ws = create_workspace(&db);
// Seed .config with hygiene enabled on daily/
seed_hygiene_config(&ws, "daily/", 0).await;
// Write some documents
ws.write("daily/log1.md", "content 1")
.await
.expect("write doc 1");
ws.write("daily/log2.md", "content 2")
.await
.expect("write doc 2");
let config = HygieneConfig {
enabled: true,
version_keep_count: 50,
cadence_hours: 12,
state_dir: _tmp.path().to_path_buf(),
};
// First run should discover daily/ and clean it
let report = run_if_due(&ws, &config).await;
assert!(!report.skipped, "first run should not be skipped");
assert!(report.had_work(), "should have cleaned documents");
assert!(
!report.directories_cleaned.is_empty(),
"should have at least one directory cleaned"
);
}
#[tokio::test]
async fn cleanup_respects_cadence_prevents_concurrent_runs() {
let (db, _tmp) = create_test_db().await;
@@ -555,7 +614,8 @@ mod tests {
let config = HygieneConfig {
enabled: true,
version_keep_count: 50,
daily_retention_days: 30,
conversation_retention_days: 7,
cadence_hours: 12,
state_dir: _tmp.path().to_path_buf(),
};
@@ -567,6 +627,13 @@ mod tests {
// Second run immediately should be skipped (cadence not elapsed)
let report2 = run_if_due(&ws, &config).await;
assert!(report2.skipped, "second run should be skipped by cadence");
// Report structure should be correct
assert_eq!(
report1.daily_logs_deleted + report1.conversation_docs_deleted,
0,
"first run should have clean counts"
);
}
#[tokio::test]
@@ -574,10 +641,6 @@ mod tests {
let (db, _tmp) = create_test_db().await;
let ws = create_workspace(&db);
// Seed hygiene on both directories
seed_hygiene_config(&ws, "daily/", 0).await;
seed_hygiene_config(&ws, "conversations/", 0).await;
// Write some documents
ws.write("daily/log1.md", "content 1")
.await
@@ -589,34 +652,35 @@ mod tests {
.await
.expect("write doc 3");
// Run with 0-day retention via direct cleanup_directory calls
let deleted_daily = cleanup_directory(&ws, "daily/", 0)
.await
.expect("cleanup daily");
let deleted_conv = cleanup_directory(&ws, "conversations/", 0)
// Run with 0-day retention to delete everything non-identity
let deleted_daily = cleanup_daily_logs(&ws, 0).await.expect("cleanup daily");
let deleted_conv = cleanup_conversation_docs(&ws, 0)
.await
.expect("cleanup conversations");
// Both should report deletions
assert!(deleted_daily > 0, "should report deleted daily logs");
assert_eq!(deleted_conv, 1, "should report 1 deleted conversation doc");
// Verify HygieneReport aggregation
// Create a HygieneReport and verify aggregation works
let report = HygieneReport {
directories_cleaned: vec![
("daily/".to_string(), deleted_daily),
("conversations/".to_string(), deleted_conv),
],
versions_pruned: 0,
daily_logs_deleted: deleted_daily,
conversation_docs_deleted: deleted_conv,
skipped: false,
};
// Verify HygieneReport structure
assert!(!report.skipped, "should not be skipped");
assert!(report.had_work(), "report should indicate work was done");
assert!(
report.daily_logs_deleted > 0 || report.conversation_docs_deleted > 0,
"report should have at least one deletion count > 0"
);
// Verify had_work() correctly checks directory counts
// Verify had_work() correctly combines both counts
let no_work = HygieneReport {
directories_cleaned: vec![],
versions_pruned: 0,
daily_logs_deleted: 0,
conversation_docs_deleted: 0,
skipped: false,
};
assert!(!no_work.had_work(), "empty report should indicate no work");
+2 -351
View File
@@ -53,9 +53,8 @@ mod search;
pub use chunker::{ChunkConfig, chunk_document};
pub use document::{
CONFIG_FILE_NAME, DocumentMetadata, DocumentVersion, HygieneMetadata, IDENTITY_PATHS,
MemoryChunk, MemoryDocument, PatchResult, VersionSummary, WorkspaceEntry, content_sha256,
is_config_path, is_identity_path, merge_workspace_entries, paths,
IDENTITY_PATHS, MemoryChunk, MemoryDocument, WorkspaceEntry, is_identity_path,
merge_workspace_entries, paths,
};
pub use embedding_cache::{CachedEmbeddingProvider, EmbeddingCacheConfig};
pub use embeddings::{
@@ -367,101 +366,6 @@ impl WorkspaceStorage {
}
}
}
// ==================== Metadata ====================
async fn update_document_metadata(
&self,
id: Uuid,
metadata: &serde_json::Value,
) -> Result<(), WorkspaceError> {
match self {
#[cfg(feature = "postgres")]
Self::Repo(repo) => repo.update_document_metadata(id, metadata).await,
Self::Db(db) => db.update_document_metadata(id, metadata).await,
}
}
async fn find_config_documents(
&self,
user_id: &str,
agent_id: Option<Uuid>,
) -> Result<Vec<MemoryDocument>, WorkspaceError> {
match self {
#[cfg(feature = "postgres")]
Self::Repo(repo) => repo.find_config_documents(user_id, agent_id).await,
Self::Db(db) => db.find_config_documents(user_id, agent_id).await,
}
}
// ==================== Versioning ====================
async fn save_version(
&self,
document_id: Uuid,
content: &str,
content_hash: &str,
changed_by: Option<&str>,
) -> Result<i32, WorkspaceError> {
match self {
#[cfg(feature = "postgres")]
Self::Repo(repo) => {
repo.save_version(document_id, content, content_hash, changed_by)
.await
}
Self::Db(db) => {
db.save_version(document_id, content, content_hash, changed_by)
.await
}
}
}
async fn get_version(
&self,
document_id: Uuid,
version: i32,
) -> Result<DocumentVersion, WorkspaceError> {
match self {
#[cfg(feature = "postgres")]
Self::Repo(repo) => repo.get_version(document_id, version).await,
Self::Db(db) => db.get_version(document_id, version).await,
}
}
async fn list_versions(
&self,
document_id: Uuid,
limit: i64,
) -> Result<Vec<VersionSummary>, WorkspaceError> {
match self {
#[cfg(feature = "postgres")]
Self::Repo(repo) => repo.list_versions(document_id, limit).await,
Self::Db(db) => db.list_versions(document_id, limit).await,
}
}
async fn get_latest_version_number(
&self,
document_id: Uuid,
) -> Result<Option<i32>, WorkspaceError> {
match self {
#[cfg(feature = "postgres")]
Self::Repo(repo) => repo.get_latest_version_number(document_id).await,
Self::Db(db) => db.get_latest_version_number(document_id).await,
}
}
async fn prune_versions(
&self,
document_id: Uuid,
keep_count: i32,
) -> Result<u64, WorkspaceError> {
match self {
#[cfg(feature = "postgres")]
Self::Repo(repo) => repo.prune_versions(document_id, keep_count).await,
Self::Db(db) => db.prune_versions(document_id, keep_count).await,
}
}
}
/// Default template seeded into HEARTBEAT.md on first access.
@@ -792,199 +696,10 @@ impl Workspace {
.await
}
// ==================== Metadata ====================
/// Update the metadata JSON on a document by ID (full replacement).
pub async fn update_metadata(
&self,
id: Uuid,
metadata: &serde_json::Value,
) -> Result<(), WorkspaceError> {
self.storage.update_document_metadata(id, metadata).await
}
/// Prune old versions for a document, keeping only the most recent `keep_count`.
///
/// Returns the number of versions deleted.
pub async fn prune_versions(
&self,
document_id: Uuid,
keep_count: i32,
) -> Result<u64, WorkspaceError> {
self.storage.prune_versions(document_id, keep_count).await
}
/// Find all `.config` documents in this workspace scope.
pub async fn find_config_documents(&self) -> Result<Vec<MemoryDocument>, WorkspaceError> {
self.storage
.find_config_documents(&self.user_id, self.agent_id)
.await
}
/// Resolve effective metadata for a document path.
///
/// Resolution chain: document's own metadata → nearest ancestor `.config` → defaults.
pub async fn resolve_metadata(&self, path: &str) -> DocumentMetadata {
// 1. Document's own metadata
let doc_meta = self
.storage
.get_document_by_path(&self.user_id, self.agent_id, path)
.await
.ok()
.map(|d| d.metadata);
// 2. Walk up parent directories looking for .config
let mut config_meta = None;
let normalized = normalize_path(path);
let mut current = normalized.as_str();
while let Some(slash_pos) = current.rfind('/') {
let parent = &current[..slash_pos];
let config_path = format!("{}/{CONFIG_FILE_NAME}", parent);
if let Ok(doc) = self
.storage
.get_document_by_path(&self.user_id, self.agent_id, &config_path)
.await
{
config_meta = Some(doc.metadata);
break;
}
current = parent;
}
// Also check root-level .config
if config_meta.is_none()
&& let Ok(doc) = self
.storage
.get_document_by_path(&self.user_id, self.agent_id, CONFIG_FILE_NAME)
.await
{
config_meta = Some(doc.metadata);
}
// 3. Merge: config as base, document metadata as overlay
let base = config_meta.unwrap_or(serde_json::json!({}));
let overlay = doc_meta.unwrap_or(serde_json::json!({}));
let merged = DocumentMetadata::merge(&base, &overlay);
DocumentMetadata::from_value(&merged)
}
// ==================== Versioning ====================
/// List versions of a document (newest first).
pub async fn list_versions(
&self,
document_id: Uuid,
limit: i64,
) -> Result<Vec<VersionSummary>, WorkspaceError> {
self.storage.list_versions(document_id, limit).await
}
/// Get a specific version of a document.
pub async fn get_version(
&self,
document_id: Uuid,
version: i32,
) -> Result<DocumentVersion, WorkspaceError> {
self.storage.get_version(document_id, version).await
}
/// Save the current content as a version if it differs from the latest.
///
/// Returns the new version number, or `None` if skipped (empty content,
/// identical hash, or versioning disabled via metadata).
async fn maybe_save_version(
&self,
document_id: Uuid,
current_content: &str,
path: &str,
changed_by: Option<&str>,
) -> Result<Option<i32>, WorkspaceError> {
// Don't version empty documents
if current_content.is_empty() {
return Ok(None);
}
// Check metadata for skip_versioning flag
let metadata = self.resolve_metadata(path).await;
if metadata.skip_versioning == Some(true) {
return Ok(None);
}
let hash = content_sha256(current_content);
// Check if latest version already has this hash (skip duplicate saves)
if let Ok(Some(latest)) = self.storage.get_latest_version_number(document_id).await
&& let Ok(ver) = self.storage.get_version(document_id, latest).await
&& ver.content_hash == hash
{
return Ok(None);
}
let version = self
.storage
.save_version(document_id, current_content, &hash, changed_by)
.await?;
Ok(Some(version))
}
// ==================== Patch ====================
/// Apply a search-and-replace patch to a workspace document.
///
/// Finds `old_string` in the document and replaces it with `new_string`.
/// If `replace_all` is true, replaces all occurrences; otherwise only the first.
/// Auto-versions before applying the patch.
pub async fn patch(
&self,
path: &str,
old_string: &str,
new_string: &str,
replace_all: bool,
) -> Result<PatchResult, WorkspaceError> {
let path = normalize_path(path);
let doc = self
.storage
.get_document_by_path(&self.user_id, self.agent_id, &path)
.await?;
if !doc.content.contains(old_string) {
return Err(WorkspaceError::PatchFailed {
path,
reason: "old_string not found in document".to_string(),
});
}
let (new_content, count) = if replace_all {
let count = doc.content.matches(old_string).count();
(doc.content.replace(old_string, new_string), count)
} else {
(doc.content.replacen(old_string, new_string, 1), 1)
};
// Injection scan for system prompt files
if is_system_prompt_file(&path) && !new_content.is_empty() {
reject_if_injected(&path, &new_content)?;
}
// Auto-version before updating
let _ = self
.maybe_save_version(doc.id, &doc.content, &path, None)
.await;
self.storage.update_document(doc.id, &new_content).await?;
self.reindex_document(doc.id).await?;
let updated = self.storage.get_document_by_id(doc.id).await?;
Ok(PatchResult {
document: updated,
replacements: count,
})
}
/// Write (create or update) a file.
///
/// Creates parent directories implicitly (they're virtual in the DB).
/// Re-indexes the document for search after writing.
/// Auto-versions the previous content before overwriting.
///
/// # Example
/// ```ignore
@@ -1000,12 +715,6 @@ impl Workspace {
.storage
.get_or_create_document_by_path(&self.user_id, self.agent_id, &path)
.await?;
// Auto-version previous content before overwriting
let _ = self
.maybe_save_version(doc.id, &doc.content, &path, None)
.await;
self.storage.update_document(doc.id, content).await?;
self.reindex_document(doc.id).await?;
@@ -1045,11 +754,6 @@ impl Workspace {
reject_if_injected(&path, &new_content)?;
}
// Auto-version previous content before appending
let _ = self
.maybe_save_version(doc.id, &doc.content, &path, None)
.await;
self.storage.update_document(doc.id, &new_content).await?;
self.reindex_document(doc.id).await?;
Ok(())
@@ -1881,14 +1585,6 @@ impl Workspace {
// Get the document
let doc = self.storage.get_document_by_id(document_id).await?;
// Check metadata for skip_indexing flag
let metadata = self.resolve_metadata(&doc.path).await;
if metadata.skip_indexing == Some(true) {
// Delete any existing chunks and skip indexing
self.storage.delete_chunks(document_id).await?;
return Ok(());
}
// Chunk the content
let chunks = chunk_document(&doc.content, ChunkConfig::default());
@@ -1974,51 +1670,6 @@ impl Workspace {
}
}
// Seed folder-level .config documents for hygiene defaults.
let config_seeds: &[(&str, serde_json::Value)] = &[
(
"daily/.config",
serde_json::json!({
"hygiene": {"enabled": true, "retention_days": 30},
"skip_versioning": true
}),
),
(
"conversations/.config",
serde_json::json!({
"hygiene": {"enabled": true, "retention_days": 7},
"skip_versioning": true
}),
),
];
for (config_path, metadata_value) in config_seeds {
match self.read_primary(config_path).await {
Ok(_) => continue, // Already exists, don't overwrite
Err(WorkspaceError::DocumentNotFound { .. }) => {}
Err(e) => {
tracing::debug!("Failed to check {}: {}", config_path, e);
continue;
}
}
// Create empty document with metadata
if let Ok(doc) = self
.storage
.get_or_create_document_by_path(&self.user_id, self.agent_id, config_path)
.await
{
if let Err(e) = self
.storage
.update_document_metadata(doc.id, metadata_value)
.await
{
tracing::debug!("Failed to set metadata on {}: {}", config_path, e);
} else {
count += 1;
}
}
}
// BOOTSTRAP.md is only seeded on truly fresh workspaces (no identity
// files existed before seeding) AND when no profile exists yet (the user
// may already have a profile from a previous install and doesn't need
+1 -191
View File
@@ -11,9 +11,7 @@ use uuid::Uuid;
use crate::error::WorkspaceError;
use crate::workspace::document::{
DocumentVersion, MemoryChunk, MemoryDocument, VersionSummary, WorkspaceEntry,
};
use crate::workspace::document::{MemoryChunk, MemoryDocument, WorkspaceEntry};
use crate::workspace::search::{RankedResult, SearchConfig, SearchResult, fuse_results};
/// Database repository for workspace operations.
@@ -704,192 +702,4 @@ impl Repository {
}
Ok(crate::workspace::merge_workspace_entries(all_entries))
}
// ==================== Metadata ====================
pub async fn update_document_metadata(
&self,
id: Uuid,
metadata: &serde_json::Value,
) -> Result<(), WorkspaceError> {
let conn = self.conn().await?;
conn.execute(
"UPDATE memory_documents SET metadata = $2, updated_at = NOW() WHERE id = $1",
&[&id, &metadata],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to update metadata: {e}"),
})?;
Ok(())
}
pub async fn find_config_documents(
&self,
user_id: &str,
agent_id: Option<Uuid>,
) -> Result<Vec<MemoryDocument>, WorkspaceError> {
let conn = self.conn().await?;
let rows = conn
.query(
r#"
SELECT id, user_id, agent_id, path, content,
created_at, updated_at, metadata
FROM memory_documents
WHERE user_id = $1 AND agent_id IS NOT DISTINCT FROM $2
AND (path LIKE '%/.config' OR path = '.config')
ORDER BY path
"#,
&[&user_id, &agent_id],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to find config documents: {e}"),
})?;
Ok(rows.iter().map(|r| self.row_to_document(r)).collect())
}
// ==================== Versioning ====================
pub async fn save_version(
&self,
document_id: Uuid,
content: &str,
content_hash: &str,
changed_by: Option<&str>,
) -> Result<i32, WorkspaceError> {
let conn = self.conn().await?;
let row = conn
.query_one(
r#"
INSERT INTO memory_document_versions
(id, document_id, version, content, content_hash, changed_by)
VALUES (
gen_random_uuid(),
$1,
(SELECT COALESCE(MAX(version), 0) + 1
FROM memory_document_versions WHERE document_id = $1),
$2, $3, $4
)
RETURNING version
"#,
&[&document_id, &content, &content_hash, &changed_by],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to save version: {e}"),
})?;
Ok(row.get(0))
}
pub async fn get_version(
&self,
document_id: Uuid,
version: i32,
) -> Result<DocumentVersion, WorkspaceError> {
let conn = self.conn().await?;
let row = conn
.query_opt(
r#"
SELECT id, document_id, version, content, content_hash,
created_at, changed_by
FROM memory_document_versions
WHERE document_id = $1 AND version = $2
"#,
&[&document_id, &version],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to get version: {e}"),
})?
.ok_or(WorkspaceError::VersionNotFound {
document_id,
version,
})?;
Ok(DocumentVersion {
id: row.get(0),
document_id: row.get(1),
version: row.get(2),
content: row.get(3),
content_hash: row.get(4),
created_at: row.get(5),
changed_by: row.get(6),
})
}
pub async fn list_versions(
&self,
document_id: Uuid,
limit: i64,
) -> Result<Vec<VersionSummary>, WorkspaceError> {
let conn = self.conn().await?;
let rows = conn
.query(
r#"
SELECT version, content_hash, created_at, changed_by
FROM memory_document_versions
WHERE document_id = $1
ORDER BY version DESC
LIMIT $2
"#,
&[&document_id, &limit],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to list versions: {e}"),
})?;
Ok(rows
.iter()
.map(|row| VersionSummary {
version: row.get(0),
content_hash: row.get(1),
created_at: row.get(2),
changed_by: row.get(3),
})
.collect())
}
pub async fn get_latest_version_number(
&self,
document_id: Uuid,
) -> Result<Option<i32>, WorkspaceError> {
let conn = self.conn().await?;
let row = conn
.query_one(
"SELECT MAX(version) FROM memory_document_versions WHERE document_id = $1",
&[&document_id],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to get latest version number: {e}"),
})?;
Ok(row.get(0))
}
pub async fn prune_versions(
&self,
document_id: Uuid,
keep_count: i32,
) -> Result<u64, WorkspaceError> {
let conn = self.conn().await?;
let result = conn
.execute(
r#"
DELETE FROM memory_document_versions
WHERE document_id = $1
AND version NOT IN (
SELECT version FROM memory_document_versions
WHERE document_id = $1
ORDER BY version DESC
LIMIT $2
)
"#,
&[&document_id, &(keep_count as i64)],
)
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to prune versions: {e}"),
})?;
Ok(result)
}
}
@@ -14,9 +14,12 @@ scenarios/test_extensions.py
scenarios/test_html_injection.py
scenarios/test_mcp_auth_flow.py
scenarios/test_oauth_credential_fallback.py
scenarios/test_oauth_refresh.py
scenarios/test_oauth_url_parameters.py
scenarios/test_owner_scope.py
scenarios/test_pairing.py
scenarios/test_routine_event_batch.py
scenarios/test_routine_full_job.py
scenarios/test_routine_oauth_credential_injection.py
scenarios/test_skills.py
scenarios/test_sse_reconnect.py
+82
View File
@@ -121,6 +121,75 @@ def _last_user_content(messages: list[dict]) -> str:
return ""
def _is_job_mode(messages: list[dict]) -> bool:
"""Detect if this conversation is a background job (not chat)."""
for msg in messages:
if msg.get("role") == "system":
content = msg.get("content", "")
if "autonomous agent working on a job" in content:
return True
return False
def _count_tool_results(messages: list[dict]) -> int:
"""Count how many tool result messages are in the conversation."""
return sum(1 for m in messages if m.get("role") == "tool")
def match_job_response(messages: list[dict], has_tools: bool) -> dict | None:
"""Handle background job conversations.
Returns a dict with either {"text": ...} or {"tool_call": ...},
or None if this isn't a job conversation.
"""
if not _is_job_mode(messages):
return None
last_user = _last_user_content(messages)
tool_result_count = _count_tool_results(messages)
# Planning call (no tools available = complete() not complete_with_tools())
if "create a plan" in last_user.lower():
return {"text": json.dumps({
"goal": "Complete the requested routine job",
"actions": [
{
"tool_name": "echo",
"parameters": {"message": "job-step-1"},
"reasoning": "First step: echo a test message",
"expected_outcome": "Echo returns the message",
},
{
"tool_name": "time",
"parameters": {"operation": "now"},
"reasoning": "Second step: get the current time",
"expected_outcome": "Returns current timestamp",
},
],
"estimated_cost": 0.001,
"estimated_time_secs": 5,
"confidence": 0.95,
})}
# Post-plan completion check: after tool results, say complete
if "planned actions" in last_user.lower() and tool_result_count >= 2:
return {"text": "The job is complete. All tasks are done."}
# Continuation prompt (from our fix): the plan didn't fully complete,
# now the agentic loop should call tools
if "continue executing now" in last_user.lower() and has_tools:
return {"tool_call": {
"tool_name": "echo",
"arguments": {"message": "continuation-step"},
}}
# After a tool result in the agentic loop, signal completion
if tool_result_count > 0 and has_tools:
return {"text": "The job is complete. All requested work has been finished."}
return None
def match_response(messages: list[dict]) -> str:
content = _last_user_content(messages)
for pattern, response in CANNED_RESPONSES:
@@ -193,6 +262,19 @@ async def chat_completions(request: web.Request) -> web.StreamResponse:
has_tools = bool(body.get("tools"))
cid = f"mock-{uuid.uuid4().hex[:8]}"
# Job-mode conversations (background routine/job execution)
job_resp = match_job_response(messages, has_tools)
if job_resp:
if "tool_call" in job_resp:
tc = job_resp["tool_call"]
if not stream:
return _tool_call_response(cid, tc)
return await _stream_tool_call(request, cid, tc)
text = job_resp["text"]
if not stream:
return _text_response(cid, text)
return await _stream_text(request, cid, text)
# Tool result in messages -> text summary
tr = _find_tool_result(messages)
if tr:
@@ -0,0 +1,133 @@
"""E2E tests for full_job routine execution.
Exercises the complete lifecycle: create a full_job routine via the
web UI, trigger it via the API, and verify the job runs tools and
completes without hitting the iteration cap.
Requires Playwright (browser-based tests).
"""
import asyncio
import uuid
from helpers import SEL, api_get, api_post
# -- Helpers ------------------------------------------------------------------
async def _send_chat_message(page, message: str) -> None:
"""Send a chat message and wait for the assistant turn to appear."""
chat_input = page.locator(SEL["chat_input"])
await chat_input.wait_for(state="visible", timeout=5000)
assistant_messages = page.locator(SEL["message_assistant"])
before_count = await assistant_messages.count()
await chat_input.fill(message)
await chat_input.press("Enter")
await page.wait_for_function(
"""({ selector, expectedCount }) => {
return document.querySelectorAll(selector).length >= expectedCount;
}""",
arg={
"selector": SEL["message_assistant"],
"expectedCount": before_count + 1,
},
timeout=30000,
)
async def _wait_for_routine(base_url: str, name: str, timeout: float = 20.0) -> dict:
"""Poll until the named routine exists."""
for _ in range(int(timeout * 2)):
resp = await api_get(base_url, "/api/routines")
resp.raise_for_status()
for routine in resp.json()["routines"]:
if routine["name"] == name:
return routine
await asyncio.sleep(0.5)
raise AssertionError(f"Routine '{name}' not created within {timeout}s")
async def _get_routine_runs(base_url: str, routine_id: str) -> list[dict]:
"""Fetch routine runs."""
resp = await api_get(base_url, f"/api/routines/{routine_id}/runs")
resp.raise_for_status()
return resp.json()["runs"]
async def _wait_for_completed_run(
base_url: str,
routine_id: str,
*,
timeout: float = 60.0,
) -> dict:
"""Poll until the newest run reaches a terminal state."""
for _ in range(int(timeout * 2)):
runs = await _get_routine_runs(base_url, routine_id)
if runs and runs[0]["status"].lower() not in ("running", "pending"):
return runs[0]
await asyncio.sleep(0.5)
raise AssertionError(
f"Routine '{routine_id}' did not complete within {timeout}s"
)
async def _wait_for_job_terminal(
base_url: str,
job_id: str,
*,
timeout: float = 60.0,
) -> dict:
"""Poll until a job reaches a terminal state."""
terminal = {"completed", "failed", "cancelled", "submitted", "accepted"}
for _ in range(int(timeout * 2)):
resp = await api_get(base_url, f"/api/jobs/{job_id}")
resp.raise_for_status()
detail = resp.json()
if detail.get("state", "").lower() in terminal:
return detail
await asyncio.sleep(0.5)
raise AssertionError(f"Job '{job_id}' did not reach terminal state within {timeout}s")
# -- Tests --------------------------------------------------------------------
async def test_full_job_routine_completes_with_tools(page, ironclaw_server):
"""A full_job routine should plan, execute tools, and complete."""
name = f"fjob-{uuid.uuid4().hex[:8]}"
# Step 1: Create full_job routine via chat
await _send_chat_message(page, f"create full-job owner routine {name}")
routine = await _wait_for_routine(ironclaw_server, name)
assert routine["id"]
assert routine["action_type"] == "full_job"
# Step 2: Trigger the routine
resp = await api_post(ironclaw_server, f"/api/routines/{routine['id']}/trigger")
resp.raise_for_status()
trigger_data = resp.json()
assert trigger_data["status"] == "triggered"
# Step 3: Wait for the run to complete
completed_run = await _wait_for_completed_run(
ironclaw_server, routine["id"], timeout=60
)
# The run should have succeeded (not failed)
assert completed_run["status"].lower() != "failed", (
f"Full job routine run failed: {completed_run}"
)
# Step 4: Verify the job reached a success state.
# Jobs may advance past "completed" to "submitted" or "accepted",
# so treat all post-completion states as success.
success_states = {"completed", "submitted", "accepted"}
if completed_run.get("job_id"):
job = await _wait_for_job_terminal(
ironclaw_server, completed_run["job_id"], timeout=30
)
assert job["state"].lower() in success_states, (
f"Expected job state in {success_states}, got '{job['state']}'"
)
+14 -14
View File
@@ -21,9 +21,7 @@ mod tests {
NotifyConfig, Routine, RoutineAction, RoutineGuardrails, RoutineRun, RunStatus, Trigger,
};
use ironclaw::agent::routine_engine::RoutineEngine;
use ironclaw::agent::{
HeartbeatConfig, HeartbeatRunner, SandboxReadiness, Scheduler, SchedulerDeps,
};
use ironclaw::agent::{HeartbeatConfig, HeartbeatRunner, Scheduler, SchedulerDeps};
use ironclaw::channels::IncomingMessage;
use ironclaw::config::{AgentConfig, RoutineConfig, SafetyConfig};
use ironclaw::context::{ContextManager, JobContext};
@@ -352,7 +350,7 @@ mod tests {
extension_manager,
registry,
safety,
SandboxReadiness::Available,
ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
))
}
@@ -456,7 +454,7 @@ mod tests {
None,
tools,
safety,
SandboxReadiness::DisabledByConfig,
ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
));
// Insert a cron routine with next_fire_at in the past.
@@ -535,7 +533,7 @@ mod tests {
None,
tools,
safety,
SandboxReadiness::DisabledByConfig,
ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
));
// Insert an event routine matching "deploy.*production".
@@ -622,7 +620,7 @@ mod tests {
None,
tools,
safety,
SandboxReadiness::DisabledByConfig,
ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
));
let routine = make_routine(
@@ -731,7 +729,7 @@ mod tests {
None,
tools,
safety,
SandboxReadiness::DisabledByConfig,
ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
));
let mut filters = std::collections::HashMap::new();
@@ -874,7 +872,7 @@ mod tests {
None,
tools,
safety,
SandboxReadiness::DisabledByConfig,
ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
));
// Insert an event routine with 1-hour cooldown.
@@ -955,7 +953,8 @@ mod tests {
let hygiene_config = HygieneConfig {
enabled: false,
version_keep_count: 50,
daily_retention_days: 30,
conversation_retention_days: 7,
cadence_hours: 24,
state_dir: _tmp.path().to_path_buf(),
};
@@ -1001,7 +1000,8 @@ mod tests {
let hygiene_config = HygieneConfig {
enabled: false,
version_keep_count: 50,
daily_retention_days: 30,
conversation_retention_days: 7,
cadence_hours: 24,
state_dir: _tmp.path().to_path_buf(),
};
@@ -1055,7 +1055,7 @@ mod tests {
None,
tools,
safety,
SandboxReadiness::DisabledByConfig,
ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
));
(engine, db, dir)
@@ -1177,7 +1177,7 @@ mod tests {
None,
tools,
safety,
SandboxReadiness::DisabledByConfig,
ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
));
// Create a full_job routine with max_concurrent = 1
@@ -1285,7 +1285,7 @@ mod tests {
None,
tools,
safety,
SandboxReadiness::DisabledByConfig,
ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
));
// Insert a due cron routine
+1 -1
View File
@@ -198,7 +198,7 @@ mod tests {
http_interceptor: None,
transcription: None,
document_extraction: None,
sandbox_readiness: ironclaw::agent::SandboxReadiness::DisabledByConfig,
sandbox_readiness: ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
builder: None,
llm_backend: "nearai".to_string(),
tenant_rates: std::sync::Arc::new(ironclaw::tenant::TenantRateRegistry::new(4, 3)),
+2 -1
View File
@@ -264,7 +264,8 @@ impl GatewayWorkflowHarness {
http_interceptor: None,
transcription: None,
document_extraction: None,
sandbox_readiness: ironclaw::agent::SandboxReadiness::DisabledByConfig,
sandbox_readiness:
ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
builder: None,
llm_backend: "nearai".to_string(),
tenant_rates: std::sync::Arc::new(ironclaw::tenant::TenantRateRegistry::new(4, 3)),
+2 -2
View File
@@ -650,7 +650,7 @@ impl TestRigBuilder {
None,
components.tools.clone(),
components.safety.clone(),
ironclaw::agent::SandboxReadiness::Available, // tests don't use real Docker
ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
));
components
.tools
@@ -759,7 +759,7 @@ impl TestRigBuilder {
http_interceptor,
transcription: None,
document_extraction: None,
sandbox_readiness: ironclaw::agent::SandboxReadiness::Available, // tests don't use real Docker
sandbox_readiness: ironclaw::agent::routine_engine::SandboxReadiness::DisabledByConfig,
builder: None,
llm_backend: "nearai".to_string(),
tenant_rates: std::sync::Arc::new(ironclaw::tenant::TenantRateRegistry::new(4, 3)),