Compare commits

...
Author SHA1 Message Date
[email protected]andClaude Opus 4.6 3bae87b04e feat(frontend): hash-based URL navigation for page refresh persistence
Navigation state is now encoded in window.location.hash so refreshing
the page (or sharing a URL) restores the current view:

  #/chat                   → chat tab, assistant thread
  #/chat/{threadId}        → specific conversation
  #/memory/{path/to/file}  → memory browser with file open
  #/jobs/{jobId}           → job detail view
  #/routines/{id}          → routine detail view
  #/settings/{subtab}      → settings sub-tab (extensions, etc.)
  #/logs                   → logs tab

Hooked into all navigation functions: switchTab, switchThread,
switchToAssistant, createNewThread, readMemoryFile, openJobDetail,
closeJobDetail, openRoutineDetail, closeRoutineDetail,
switchSettingsSubtab.

Thread restore is deferred until loadThreads() completes (async),
then the pending thread ID is matched against the loaded thread list.

Browser back/forward buttons work via hashchange listener.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 14:31:16 -07:00
[email protected]andClaude Opus 4.6 c47c82bf79 feat(frontend): structured data cards + chat renderer API for rich message rendering
Agent responses containing JSON/structured data (like mission results,
status objects) now render as styled cards with labeled fields, status
badges, and monospaced IDs instead of raw text.

Built-in rendering:
- Detects inline JSON objects (including Python-style single quotes)
- Renders as data cards with key-value rows
- Status/state fields get colored badges (success/error/pending)
- UUIDs rendered in monospace

Extensible via widgets:
- IronClaw.registerChatRenderer({ id, match, render, priority })
- First matching renderer wins (priority ordering)
- Renderer gets the content element to mutate in place

Also adds ChatRenderer variant to WidgetSlot enum.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 14:04:11 -07:00
[email protected]andClaude Opus 4.6 fc32a31515 fix: address CI failures — license, rust-version, formatting, manifest warnings
- Add license = "MIT OR Apache-2.0" to ironclaw_frontend Cargo.toml (cargo-deny)
- Fix rust-version to 1.92 to match other crates
- Log warning for invalid widget manifests instead of silent skip
- Run cargo fmt across all files

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 12:27:41 -07:00
[email protected]andClaude Opus 4.6 c54c8e0fb5 feat(frontend): extract frontend into ironclaw_frontend crate with widget extension system
Moves all frontend static assets (app.js, style.css, index.html, i18n/*,
theme-init.js, favicon.ico) from src/channels/web/static/ into a dedicated
ironclaw_frontend crate. The crate also adds:

- Layout configuration types (branding, tab order, chat features, per-widget config)
- Widget manifest types with named slot system (tab, chat_header, sidebar, etc.)
- CSS scoping utility (auto-prefixes selectors with [data-widget="id"])
- Bundle assembly (injects layout config, widgets, and custom CSS into HTML)
- Frontend API endpoints (GET/PUT layout, list widgets, serve widget files)
- Browser-side IronClaw.registerWidget() API with authenticated fetch,
  event subscription, theme access, and i18n

Widgets are stored in workspace at frontend/widgets/{id}/ and served via
the API. Layout config is stored at frontend/layout.json. The agent can
create/edit both using existing memory_write/memory_read tools.

Gateway handlers now reference ironclaw_frontend::assets constants instead
of include_str!() with local paths, completing the separation.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 12:27:41 -07:00
[email protected] ebdae267a0 Merge remote-tracking branch 'origin/staging' into feat/workspace-metadata-versioning-patch
# Conflicts:
#	src/agent/heartbeat.rs
#	src/db/libsql_migrations.rs
2026-03-28 12:27:02 -07:00
[email protected]andClaude Opus 4.6 cb4cf123c6 fix: address review feedback — transaction safety, patch mode, formatting
- Wrap libSQL save_version in a transaction to prevent race condition
  where concurrent writers could allocate the same version number
- Make content optional in memory_write when in patch mode (old_string
  present) — LLM no longer forced to provide unused content param
- Improve metadata update error handling with explicit match arms
- Run cargo fmt across all files

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 12:12:41 -07:00
fd41bdf4be fix(worker): treat empty LLM response after text output as completion (#1677)
* fix(worker): treat empty LLM response after text output as completion

When a job's LLM produces a substantive text response (e.g., formatted
results from a routine) and the next LLM call returns empty or errors,
the worker now treats this as successful completion instead of
continuing the loop until failure.

Previously, empty responses always triggered TextAction::Continue,
causing the loop to re-call the LLM. The LLM had nothing more to say,
so the provider returned "Response contained no message or tool call
(empty)". This made routine jobs that successfully produced results
report as "failed".

The fix adds a `has_text_response` flag to JobDelegate:
- After any non-empty text response: flag is set
- Empty text after flag is set: treated as completion
- LLM errors (select_tools/respond_with_tools) after flag: treated
  as completion instead of propagating
- Empty text before any output: still retries (rate-limit backoff)

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

* fix(worker): restrict error swallowing to EmptyResponse variant only

- Add LlmError::EmptyResponse variant for when LLM returns no content
- Update nearai_chat and github_copilot providers to emit EmptyResponse
  instead of InvalidResponse for empty/no-choice responses
- try_complete_on_error now only swallows EmptyResponse (not AuthFailed,
  ContextLengthExceeded, Http, Io, etc.)
- Extract is_completion_eligible_error as testable pure function
- Log mark_completed errors at warn level instead of silently dropping
- Add EmptyResponse to retry and circuit breaker transient classifications
- Rewrite test to exercise real classification logic against all variants

Addresses review feedback from zmanian and gemini-code-assist.

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

* refactor(worker): extract mark_completed_or_warn helper to DRY completion logic

Extract shared mark-completed + warn-on-failure pattern into a single
helper method used by both try_complete_on_error and handle_text_response.

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

---------

Co-authored-by: j-bloggs <[email protected]>
Co-authored-by: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 18:46:08 +01:00
de5a1c7b0d fix(worker): replace script -qfc with pty-process for injection-safe PTY (#1678)
- Add pty-process crate (MIT, tokio async support) for PTY allocation
- Spawn claude CLI with pty-process::Command::arg() chaining instead of
  building a shell string for script -qfc
- Eliminates all shell injection surfaces: prompt, model, session_id
  are passed via execve, never interpreted by a shell
- Keep stderr on separate pipe to prevent NDJSON parse breakage
  (pty-process attaches PTY to all fds by default)
- Gate PTY behind #[cfg(unix)] with direct-spawn fallback for Windows CI
- Read stdout from PTY master (implements tokio::io::AsyncRead)
- Add regression tests: arg vector construction + PTY allocation

Addresses review feedback from zmanian and gemini-code-assist.

Co-authored-by: j-bloggs <[email protected]>
Co-authored-by: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 16:31:49 +01:00
9ce3a9fc53 feat(discord): implement on_broadcast via DM channel creation (#1693)
- Implement broadcast_dm() that creates a DM channel with the target
  user (POST /users/@me/channels, cached by Discord) and sends the
  message to it
- Extract DISCORD_API_BASE constant for all Discord REST API URLs
- Extract send_channel_message() shared helper to deduplicate message
  posting between on_respond and broadcast_dm
- Add snowflake validation on user_id before API calls
- Fix pre-existing clippy redundant_closure warning
- Use typed DmChannelResponse struct instead of serde_json::Value

Closes no specific issue — completes the previously stubbed on_broadcast.

Co-authored-by: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 16:31:27 +01:00
AchieveandGitHub 9bb19a98f7 fix(web): redact database error details from API responses (#1711) 2026-03-28 15:13:28 +01:00
AchieveandGitHub 0b33ca9926 fix(oauth): tighten legacy state validation and fallback handling (#1701)
* fix(oauth): tighten legacy state validation and fallback handling

* style: fix formatting

* refactor: separate validation checks for clearer error messages
2026-03-28 15:10:39 +01:00
AchieveandGitHub 9ba10eac35 fix(db): add tracing warn for naive timestamp fallback and improve parse_timestamp tests (#1700)
* fix(db): add tracing warn for naive timestamp fallback and improve parse_timestamp tests

* style: fix formatting
2026-03-28 15:08:25 +01:00
AchieveandGitHub 27e8d6f8dd fix(wasm): use typed WASM schema as advertised schema when available (#1699) 2026-03-28 15:07:16 +01:00
Henry ParkandGitHub f49f368355 Clean up extension credentials on uninstall (#1718)
* Clean up extension credentials on uninstall

* Address PR review feedback

* Cover channel webhook secrets on uninstall

* Harden tool secret cleanup detection
2026-03-28 14:46:45 +01:00
[email protected]andClaude Opus 4.6 7532065590 feat(workspace): metadata-driven indexing/hygiene, document versioning, and patch support
Foundation for the extensible frontend system. Workspace documents now
support metadata flags (skip_indexing, skip_versioning, hygiene config)
via folder-level .config documents and per-file overrides, replacing
hardcoded hygiene targets and indexing behavior.

Key changes:
- DocumentMetadata type with resolution chain (doc → folder .config → defaults)
- Document versioning: auto-saves previous content on write/append/patch
- Workspace patch: search-and-replace editing via memory_write tool
- Hygiene rewrite: discovers cleanup targets from .config metadata
  instead of hardcoded daily/ and conversations/ directories
- memory_read gains version/list_versions params
- memory_write gains metadata/old_string/new_string/replace_all params
- V14 migration adds memory_document_versions table (both PG + libSQL)

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 00:28:32 -07:00
53 changed files with 4778 additions and 438 deletions
+2 -3
View File
@@ -191,10 +191,9 @@ HEARTBEAT_NOTIFY_CHANNEL=cli
HEARTBEAT_NOTIFY_USER=default
# Memory hygiene settings (automatic cleanup of stale workspace documents)
# Runs on each heartbeat tick; identity files (IDENTITY.md, SOUL.md) are never deleted
# Runs on each heartbeat tick; discovers cleanup targets from .config metadata
# MEMORY_HYGIENE_ENABLED=true
# 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_VERSION_KEEP_COUNT=50 # max versions to keep per document
# MEMORY_HYGIENE_CADENCE_HOURS=12 # minimum hours between cleanup passes
# Docker Sandbox
Generated
+26 -5
View File
@@ -3150,7 +3150,7 @@ dependencies = [
"libc",
"percent-encoding",
"pin-project-lite",
"socket2 0.5.10",
"socket2 0.6.3",
"system-configuration",
"tokio",
"tower-service",
@@ -3429,6 +3429,7 @@ dependencies = [
"iana-time-zone",
"insta",
"ironclaw_common",
"ironclaw_frontend",
"ironclaw_safety",
"json5",
"libsql",
@@ -3439,6 +3440,7 @@ dependencies = [
"pgvector",
"postgres-types",
"pretty_assertions",
"pty-process",
"rand 0.8.5",
"readabilityrs",
"refinery",
@@ -3495,6 +3497,15 @@ dependencies = [
"serde_json",
]
[[package]]
name = "ironclaw_frontend"
version = "0.1.0"
dependencies = [
"serde",
"serde_json",
"thiserror 2.0.18",
]
[[package]]
name = "ironclaw_safety"
version = "0.2.0"
@@ -3524,7 +3535,7 @@ checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46"
dependencies = [
"hermit-abi",
"libc",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -4906,6 +4917,16 @@ dependencies = [
"syn 1.0.109",
]
[[package]]
name = "pty-process"
version = "0.5.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "71cec9e2670207c5ebb9e477763c74436af3b9091dd550b9fb3c1bec7f3ea266"
dependencies = [
"rustix 1.1.4",
"tokio",
]
[[package]]
name = "pulley-interpreter"
version = "28.0.1"
@@ -4930,7 +4951,7 @@ dependencies = [
"quinn-udp",
"rustc-hash 2.1.1",
"rustls 0.23.37",
"socket2 0.5.10",
"socket2 0.6.3",
"thiserror 2.0.18",
"tokio",
"tracing",
@@ -4967,9 +4988,9 @@ dependencies = [
"cfg_aliases",
"libc",
"once_cell",
"socket2 0.5.10",
"socket2 0.6.3",
"tracing",
"windows-sys 0.59.0",
"windows-sys 0.60.2",
]
[[package]]
+6 -1
View File
@@ -1,5 +1,5 @@
[workspace]
members = [".", "crates/ironclaw_common", "crates/ironclaw_safety"]
members = [".", "crates/ironclaw_common", "crates/ironclaw_safety", "crates/ironclaw_frontend"]
exclude = [
"channels-src/discord",
"channels-src/telegram",
@@ -105,6 +105,7 @@ cron = "0.13"
ironclaw_common = { path = "crates/ironclaw_common", version = "0.1.0" }
# Safety/sanitization
ironclaw_frontend = { path = "crates/ironclaw_frontend", version = "0.1.0" }
ironclaw_safety = { path = "crates/ironclaw_safety", version = "0.2.0" }
regex = "1"
aho-corasick = "1"
@@ -189,6 +190,10 @@ json5 = { version = "0.4", optional = true }
[target.'cfg(target_os = "macos")'.dependencies]
security-framework = "3"
# PTY allocation for Claude CLI stdout buffering fix (Unix only)
[target.'cfg(unix)'.dependencies]
pty-process = { version = "0.5", features = ["async"] }
# Linux secret-service (GNOME Keyring, KWallet)
[target.'cfg(target_os = "linux")'.dependencies]
secret-service = { version = "4", features = ["rt-tokio-crypto-rust"] }
+118 -23
View File
@@ -28,6 +28,9 @@ use std::{cmp::Ordering, collections::HashMap};
use ed25519_dalek::{Signature, Verifier, VerifyingKey};
use serde::{Deserialize, Serialize};
/// Discord REST API v10 base URL.
const DISCORD_API_BASE: &str = "https://discord.com/api/v10";
use exports::near::agent::channel::{
AgentResponse, ChannelConfig, Guest, HttpEndpointConfig, IncomingHttpRequest,
OutgoingHttpResponse, PollConfig, StatusUpdate,
@@ -427,7 +430,7 @@ impl Guest for DiscordChannel {
(
"PATCH",
format!(
"https://discord.com/api/v10/webhooks/{}/{}/messages/@original",
"{DISCORD_API_BASE}/webhooks/{}/{}/messages/@original",
application_id, token
),
)
@@ -438,20 +441,7 @@ impl Guest for DiscordChannel {
payload["allowed_mentions"] = serde_json::json!({
"replied_user": true
});
let mention_payload = serde_json::to_vec(&payload)
.map_err(|e| format!("Failed to serialize mention payload: {}", e))?;
let mention_url = format!(
"https://discord.com/api/v10/channels/{}/messages",
metadata.channel_id
);
let result = channel_host::http_request(
"POST",
&mention_url,
&discord_auth_headers_json(true),
Some(&mention_payload),
None,
);
return map_discord_response(result);
return send_channel_message(&metadata.channel_id, payload);
} else {
return Err("Unsupported Discord response metadata".to_string());
};
@@ -469,8 +459,8 @@ impl Guest for DiscordChannel {
fn on_status(_update: StatusUpdate) {}
fn on_broadcast(_user_id: String, _response: AgentResponse) -> Result<(), String> {
Err("broadcast not yet implemented for Discord channel".to_string())
fn on_broadcast(user_id: String, response: AgentResponse) -> Result<(), String> {
broadcast_dm(&user_id, &response.content)
}
fn on_shutdown() {
@@ -501,6 +491,21 @@ fn map_discord_response(
}
}
/// Post a JSON payload to a Discord channel as a new message.
fn send_channel_message(channel_id: &str, payload: serde_json::Value) -> Result<(), String> {
let payload_bytes = serde_json::to_vec(&payload)
.map_err(|e| format!("Failed to serialize message: {}", e))?;
let url = format!("{DISCORD_API_BASE}/channels/{}/messages", channel_id);
let result = channel_host::http_request(
"POST",
&url,
&discord_auth_headers_json(true),
Some(&payload_bytes),
None,
);
map_discord_response(result)
}
fn load_runtime_config() -> DiscordRuntimeConfig {
channel_host::workspace_read("config.json")
.and_then(|raw| serde_json::from_str::<DiscordRuntimeConfig>(&raw).ok())
@@ -539,7 +544,7 @@ fn get_or_fetch_bot_id() -> Option<String> {
let response = channel_host::http_request(
"GET",
"https://discord.com/api/v10/users/@me",
&format!("{DISCORD_API_BASE}/users/@me"),
&discord_auth_headers_json(false),
None,
Some(10_000),
@@ -659,7 +664,7 @@ fn poll_channel_mentions(channel_id: &str, bot_id: &str) {
fn fetch_latest_message_id(channel_id: &str) -> Option<String> {
let url = format!(
"https://discord.com/api/v10/channels/{}/messages?limit=1",
"{DISCORD_API_BASE}/channels/{}/messages?limit=1",
channel_id
);
let response = channel_host::http_request(
@@ -697,7 +702,7 @@ fn fetch_messages_after_cursor(
for page in 0..MAX_PAGES {
let url = format!(
"https://discord.com/api/v10/channels/{}/messages?limit={}&after={}",
"{DISCORD_API_BASE}/channels/{}/messages?limit={}&after={}",
channel_id, PAGE_LIMIT, after
);
let response = match channel_host::http_request(
@@ -986,7 +991,7 @@ fn handle_slash_command(interaction: &DiscordInteraction) -> bool {
);
// Attempt to notify user of internal error
let url = format!(
"https://discord.com/api/v10/webhooks/{}/{}",
"{DISCORD_API_BASE}/webhooks/{}/{}",
interaction.application_id, interaction.token
);
let payload = serde_json::json!({
@@ -1106,7 +1111,7 @@ fn check_sender_permission(
}
let dm_policy =
channel_host::workspace_read(DM_POLICY_PATH).unwrap_or_else(|| default_dm_policy());
channel_host::workspace_read(DM_POLICY_PATH).unwrap_or_else(default_dm_policy);
if dm_policy == "open" {
return true;
}
@@ -1161,7 +1166,7 @@ fn check_sender_permission(
/// Send a pairing code as an ephemeral Discord followup message.
fn send_pairing_reply(ctx: &PairingReplyCtx, code: &str) -> Result<(), String> {
let url = format!(
"https://discord.com/api/v10/webhooks/{}/{}",
"{DISCORD_API_BASE}/webhooks/{}/{}",
ctx.application_id, ctx.token
);
let payload = serde_json::json!({
@@ -1194,6 +1199,57 @@ fn send_pairing_reply(ctx: &PairingReplyCtx, code: &str) -> Result<(), String> {
}
}
/// Send a broadcast message to a Discord user via DM.
///
/// Creates a DM channel with the user (Discord caches this, so repeated calls
/// for the same user reuse the existing channel) and then posts the message.
fn broadcast_dm(user_id: &str, content: &str) -> Result<(), String> {
// Validate user_id is a plausible Discord snowflake (numeric, 17-20 digits)
// to avoid injecting arbitrary strings into API URLs.
if user_id.is_empty()
|| !user_id.chars().all(|c| c.is_ascii_digit())
|| user_id.len() < 17
|| user_id.len() > 20
{
return Err(format!("Invalid Discord user ID: '{}'", user_id));
}
// Step 1: Open (or reuse) a DM channel with the target user.
let create_dm_payload = serde_json::json!({ "recipient_id": user_id });
let create_dm_bytes = serde_json::to_vec(&create_dm_payload)
.map_err(|e| format!("Failed to serialize DM channel request: {}", e))?;
let dm_response = channel_host::http_request(
"POST",
&format!("{DISCORD_API_BASE}/users/@me/channels"),
&discord_auth_headers_json(true),
Some(&create_dm_bytes),
Some(10_000),
)
.map_err(|e| format!("Failed to create DM channel: {}", e))?;
if !(200..300).contains(&dm_response.status) {
let body = String::from_utf8_lossy(&dm_response.body);
return Err(format!(
"Discord create-DM failed: {} - {}",
dm_response.status, body
));
}
#[derive(Deserialize)]
struct DmChannelResponse {
id: String,
}
let dm_channel: DmChannelResponse = serde_json::from_slice(&dm_response.body)
.map_err(|e| format!("Failed to parse DM channel response: {}", e))?;
let channel_id = &dm_channel.id;
// Step 2: Send the message to the DM channel.
let truncated = truncate_message(content);
let payload = serde_json::json!({ "content": truncated });
send_channel_message(channel_id, payload)
}
fn json_response(status: u16, value: serde_json::Value) -> OutgoingHttpResponse {
let body = serde_json::to_vec(&value).unwrap_or_default();
let headers = serde_json::json!({"Content-Type": "application/json"});
@@ -1593,4 +1649,43 @@ mod tests {
assert_eq!(interaction.interaction_type, 2);
assert!(interaction.data.is_some());
}
#[test]
fn test_broadcast_dm_payload_format() {
// Verify the DM channel creation payload is well-formed JSON that
// Discord's API expects.
let user_id = "123456789012345678";
let payload = serde_json::json!({ "recipient_id": user_id });
let serialized = serde_json::to_vec(&payload).unwrap();
let parsed: serde_json::Value = serde_json::from_slice(&serialized).unwrap();
assert_eq!(
parsed.get("recipient_id").and_then(|v| v.as_str()),
Some(user_id)
);
}
#[test]
fn test_broadcast_message_truncation() {
// Broadcast uses truncate_message, verify it handles content within
// Discord's 2000-char limit for DMs.
let short = "Hello from broadcast";
assert_eq!(truncate_message(short), short);
let long = "x".repeat(2500);
let result = truncate_message(&long);
assert!(result.len() <= 2006); // 1990 content + 16 suffix
assert!(result.ends_with("\n... (truncated)"));
}
#[test]
fn test_broadcast_dm_validates_snowflake() {
// broadcast_dm rejects invalid Discord snowflake IDs before making
// any API calls. We can call it directly since invalid IDs are
// rejected before any host function is invoked.
assert!(broadcast_dm("", "hi").is_err());
assert!(broadcast_dm("abc", "hi").is_err());
assert!(broadcast_dm("12345", "hi").is_err()); // too short
assert!(broadcast_dm("123456789012345678901", "hi").is_err()); // too long
assert!(broadcast_dm("12345678901234567x", "hi").is_err()); // non-digit
}
}
+15
View File
@@ -0,0 +1,15 @@
[package]
name = "ironclaw_frontend"
version = "0.1.0"
edition = "2024"
rust-version = "1.92"
description = "Frontend assets, layout configuration, and widget extension system for IronClaw"
license = "MIT OR Apache-2.0"
[package.metadata.dist]
dist = false
[dependencies]
serde = { version = "1", features = ["derive"] }
serde_json = "1"
thiserror = "2"
+37
View File
@@ -0,0 +1,37 @@
//! Embedded static assets for the IronClaw web gateway.
//!
//! All frontend files are compiled into the binary via `include_str!()` /
//! `include_bytes!()`. The web gateway serves these as the default baseline;
//! workspace-stored customizations (layout config, widgets, CSS overrides)
//! are layered on top at runtime.
// ==================== Core Files ====================
/// Main HTML page (SPA shell).
pub const INDEX_HTML: &str = include_str!("../static/index.html");
/// Main application JavaScript.
pub const APP_JS: &str = include_str!("../static/app.js");
/// Base stylesheet.
pub const STYLE_CSS: &str = include_str!("../static/style.css");
/// Theme initialization script (runs synchronously in `<head>` to prevent FOUC).
pub const THEME_INIT_JS: &str = include_str!("../static/theme-init.js");
/// Favicon.
pub const FAVICON_ICO: &[u8] = include_bytes!("../static/favicon.ico");
// ==================== Internationalization ====================
/// i18n core library.
pub const I18N_INDEX_JS: &str = include_str!("../static/i18n/index.js");
/// English translations.
pub const I18N_EN_JS: &str = include_str!("../static/i18n/en.js");
/// Chinese (Simplified) translations.
pub const I18N_ZH_CN_JS: &str = include_str!("../static/i18n/zh-CN.js");
/// i18n integration with the app.
pub const I18N_APP_JS: &str = include_str!("../static/i18n-app.js");
+237
View File
@@ -0,0 +1,237 @@
//! Frontend bundle assembly.
//!
//! Combines the embedded base HTML with workspace customizations (layout
//! config, widgets, CSS overrides) into the final served page.
use crate::layout::LayoutConfig;
use crate::widget::{WidgetManifest, scope_css};
/// A resolved frontend bundle ready for serving.
///
/// Contains the layout configuration, resolved widgets (with their JS/CSS
/// content loaded), and any custom CSS overrides.
#[derive(Debug, Clone, Default)]
pub struct FrontendBundle {
/// Layout configuration (branding, tabs, chat settings).
pub layout: LayoutConfig,
/// Resolved widgets with their source code loaded.
pub widgets: Vec<ResolvedWidget>,
/// Custom CSS to append after the base stylesheet.
pub custom_css: Option<String>,
}
/// A widget with its manifest and source files loaded.
#[derive(Debug, Clone)]
pub struct ResolvedWidget {
/// Widget metadata.
pub manifest: WidgetManifest,
/// JavaScript source code (`index.js`).
pub js: String,
/// Optional CSS source code (`style.css`), auto-scoped.
pub css: Option<String>,
}
/// Inject frontend customizations into the base HTML template.
///
/// Modifications:
///
/// **Before `</head>`:**
/// - Branding CSS custom property overrides
/// - Title override (replaces `<title>` content)
///
/// **Before `</body>`:**
/// - Layout config as `window.__IRONCLAW_LAYOUT__`
/// - Scoped widget `<style>` blocks
/// - Widget `<script type="module">` tags
/// - Custom CSS `<style>` block
pub fn assemble_index(base_html: &str, bundle: &FrontendBundle) -> String {
let mut head_injections = Vec::new();
let mut body_injections = Vec::new();
// --- Head injections ---
// Branding CSS variables
let css_vars = bundle.layout.branding.to_css_vars();
if !css_vars.is_empty() {
head_injections.push(format!("<style>{}</style>", css_vars));
}
// --- Body injections ---
// Layout config as global variable
if let Ok(layout_json) = serde_json::to_string(&bundle.layout) {
body_injections.push(format!(
"<script>window.__IRONCLAW_LAYOUT__ = {};</script>",
layout_json
));
}
// Widget CSS (scoped) and JS
for widget in &bundle.widgets {
if let Some(ref css) = widget.css {
let scoped = scope_css(css, &widget.manifest.id);
if !scoped.trim().is_empty() {
body_injections.push(format!(
"<style data-widget=\"{}\">{}</style>",
widget.manifest.id, scoped
));
}
}
// Widget JS as module script served from API
body_injections.push(format!(
"<script type=\"module\" src=\"/api/frontend/widget/{}/index.js\"></script>",
widget.manifest.id
));
}
// Custom CSS
if let Some(ref custom_css) = bundle.custom_css {
if !custom_css.trim().is_empty() {
body_injections.push(format!("<style data-custom-css>{}</style>", custom_css));
}
}
// --- Assemble ---
let mut result = base_html.to_string();
// Inject before </head>
if !head_injections.is_empty() {
let head_block = head_injections.join("\n");
if let Some(pos) = result.rfind("</head>") {
result.insert_str(pos, &format!("\n{}\n", head_block));
}
}
// Override <title> if branding title is set
if let Some(ref title) = bundle.layout.branding.title {
if let Some(start) = result.find("<title>") {
if let Some(end) = result[start..].find("</title>") {
let end = start + end + "</title>".len();
result.replace_range(start..end, &format!("<title>{}</title>", title));
}
}
}
// Inject before </body>
if !body_injections.is_empty() {
let body_block = body_injections.join("\n");
if let Some(pos) = result.rfind("</body>") {
result.insert_str(pos, &format!("\n{}\n", body_block));
}
}
result
}
#[cfg(test)]
mod tests {
use super::*;
use crate::layout::*;
use crate::widget::*;
const MINIMAL_HTML: &str =
"<!DOCTYPE html><html><head><title>IronClaw</title></head><body></body></html>";
#[test]
fn test_assemble_index_no_customizations() {
let bundle = FrontendBundle::default();
let result = assemble_index(MINIMAL_HTML, &bundle);
// Layout config is always injected (even when default/empty)
assert!(result.contains("window.__IRONCLAW_LAYOUT__"));
// No branding overrides or custom CSS
assert!(!result.contains("--color-primary"));
assert!(!result.contains("data-custom-css"));
}
#[test]
fn test_assemble_index_branding_title() {
let bundle = FrontendBundle {
layout: LayoutConfig {
branding: BrandingConfig {
title: Some("Acme AI".to_string()),
..Default::default()
},
..Default::default()
},
..Default::default()
};
let result = assemble_index(MINIMAL_HTML, &bundle);
assert!(result.contains("<title>Acme AI</title>"));
assert!(!result.contains("<title>IronClaw</title>"));
}
#[test]
fn test_assemble_index_branding_colors() {
let bundle = FrontendBundle {
layout: LayoutConfig {
branding: BrandingConfig {
colors: Some(BrandingColors {
primary: Some("#0066cc".to_string()),
accent: None,
}),
..Default::default()
},
..Default::default()
},
..Default::default()
};
let result = assemble_index(MINIMAL_HTML, &bundle);
assert!(result.contains("--color-primary: #0066cc;"));
}
#[test]
fn test_assemble_index_layout_config_injected() {
let bundle = FrontendBundle {
layout: LayoutConfig {
tabs: TabConfig {
hidden: Some(vec!["routines".to_string()]),
..Default::default()
},
..Default::default()
},
..Default::default()
};
let result = assemble_index(MINIMAL_HTML, &bundle);
assert!(result.contains("window.__IRONCLAW_LAYOUT__"));
assert!(result.contains("routines"));
}
#[test]
fn test_assemble_index_widget_script() {
let bundle = FrontendBundle {
widgets: vec![ResolvedWidget {
manifest: WidgetManifest {
id: "dashboard".to_string(),
name: "Dashboard".to_string(),
slot: WidgetSlot::Tab,
icon: None,
position: None,
},
js: "console.log('hello');".to_string(),
css: Some(".panel { color: red; }".to_string()),
}],
..Default::default()
};
let result = assemble_index(MINIMAL_HTML, &bundle);
assert!(result.contains("src=\"/api/frontend/widget/dashboard/index.js\""));
assert!(result.contains("data-widget=\"dashboard\""));
assert!(result.contains("[data-widget=\"dashboard\"] .panel"));
}
#[test]
fn test_assemble_index_custom_css() {
let bundle = FrontendBundle {
custom_css: Some("body { background: #111; }".to_string()),
..Default::default()
};
let result = assemble_index(MINIMAL_HTML, &bundle);
assert!(result.contains("data-custom-css"));
assert!(result.contains("background: #111;"));
}
}
+179
View File
@@ -0,0 +1,179 @@
//! Layout configuration types for frontend customization.
//!
//! A [`LayoutConfig`] is stored as `frontend/layout.json` in the workspace.
//! It controls branding, tab visibility/order, chat features, and per-widget
//! configuration. All fields are optional with sensible defaults.
use std::collections::HashMap;
use serde::{Deserialize, Serialize};
/// Top-level layout configuration.
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct LayoutConfig {
/// Branding overrides (title, logo, colors).
#[serde(default)]
pub branding: BrandingConfig,
/// Tab bar configuration.
#[serde(default)]
pub tabs: TabConfig,
/// Chat panel configuration.
#[serde(default)]
pub chat: ChatConfig,
/// Per-widget instance configuration (keyed by widget ID).
#[serde(default)]
pub widgets: HashMap<String, WidgetInstanceConfig>,
}
/// Branding overrides for the gateway UI.
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct BrandingConfig {
/// Page title (replaces default "IronClaw").
#[serde(default, skip_serializing_if = "Option::is_none")]
pub title: Option<String>,
/// Subtitle shown below the title.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub subtitle: Option<String>,
/// URL to a logo image.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub logo_url: Option<String>,
/// URL to a custom favicon.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub favicon_url: Option<String>,
/// Color overrides (injected as CSS custom properties on `:root`).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub colors: Option<BrandingColors>,
}
/// Color overrides for the UI theme.
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct BrandingColors {
/// Primary brand color (e.g., `"#0066cc"`).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub primary: Option<String>,
/// Accent color.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub accent: Option<String>,
}
/// Tab bar layout configuration.
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct TabConfig {
/// Ordered list of tab IDs to display (built-in + widget tabs).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub order: Option<Vec<String>>,
/// Tab IDs to hide from the tab bar.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub hidden: Option<Vec<String>>,
/// Default tab to show on load.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub default_tab: Option<String>,
}
/// Chat panel feature flags.
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct ChatConfig {
/// Show suggestion chips below the input.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub suggestions: Option<bool>,
/// Enable image upload in the chat input.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub image_upload: Option<bool>,
}
/// Per-widget instance configuration.
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct WidgetInstanceConfig {
/// Whether this widget is enabled.
#[serde(default)]
pub enabled: bool,
/// Arbitrary widget-specific configuration passed to `widget.init()`.
#[serde(default)]
pub config: serde_json::Value,
}
impl BrandingConfig {
/// Generate CSS custom property overrides for injection into `:root`.
pub fn to_css_vars(&self) -> String {
let mut vars = Vec::new();
if let Some(ref colors) = self.colors {
if let Some(ref primary) = colors.primary {
vars.push(format!("--color-primary: {};", primary));
}
if let Some(ref accent) = colors.accent {
vars.push(format!("--color-accent: {};", accent));
}
}
if vars.is_empty() {
String::new()
} else {
format!(":root {{ {} }}", vars.join(" "))
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_layout_config_default_is_empty() {
let config = LayoutConfig::default();
assert!(config.branding.title.is_none());
assert!(config.tabs.order.is_none());
assert!(config.widgets.is_empty());
}
#[test]
fn test_layout_config_roundtrip() {
let json = serde_json::json!({
"branding": { "title": "Acme AI", "colors": { "primary": "#0066cc" } },
"tabs": { "order": ["chat", "memory"], "hidden": ["routines"] },
"widgets": { "dashboard": { "enabled": true, "config": { "refresh": 30 } } }
});
let config: LayoutConfig = serde_json::from_value(json).unwrap();
assert_eq!(config.branding.title.as_deref(), Some("Acme AI"));
assert_eq!(config.tabs.hidden.as_ref().map(|h| h.len()), Some(1));
assert!(config.widgets.get("dashboard").is_some_and(|w| w.enabled));
}
#[test]
fn test_branding_css_vars_empty() {
let branding = BrandingConfig::default();
assert!(branding.to_css_vars().is_empty());
}
#[test]
fn test_branding_css_vars_with_colors() {
let branding = BrandingConfig {
colors: Some(BrandingColors {
primary: Some("#0066cc".to_string()),
accent: Some("#ff6b00".to_string()),
}),
..Default::default()
};
let css = branding.to_css_vars();
assert!(css.contains("--color-primary: #0066cc;"));
assert!(css.contains("--color-accent: #ff6b00;"));
}
#[test]
fn test_partial_deserialization() {
let json = serde_json::json!({"branding": {"title": "Test"}});
let config: LayoutConfig = serde_json::from_value(json).unwrap();
assert_eq!(config.branding.title.as_deref(), Some("Test"));
assert!(config.chat.suggestions.is_none());
}
}
+36
View File
@@ -0,0 +1,36 @@
//! IronClaw Frontend — assets, layout configuration, and widget extension system.
//!
//! This crate owns the complete frontend for the IronClaw web gateway:
//!
//! - **Embedded assets** (`assets` module): HTML, JS, CSS, i18n files compiled
//! into the binary for zero-dependency serving.
//! - **Layout configuration** (`layout` module): Branding, tab order, feature
//! flags — customizable per-tenant via workspace.
//! - **Widget system** (`widget` module): Self-contained frontend components
//! that plug into named slots in the UI.
//! - **Bundle assembly** (`bundle` module): Combines base assets with workspace
//! customizations into the final served HTML.
pub mod assets;
mod bundle;
mod layout;
mod widget;
pub use bundle::{FrontendBundle, ResolvedWidget, assemble_index};
pub use layout::{
BrandingColors, BrandingConfig, ChatConfig, LayoutConfig, TabConfig, WidgetInstanceConfig,
};
pub use widget::{WidgetManifest, WidgetSlot, scope_css};
/// Errors from frontend operations.
#[derive(Debug, thiserror::Error)]
pub enum FrontendError {
#[error("Layout configuration is invalid: {reason}")]
InvalidLayout { reason: String },
#[error("Widget '{id}' not found")]
WidgetNotFound { id: String },
#[error("Widget manifest is invalid: {reason}")]
InvalidManifest { reason: String },
}
+178
View File
@@ -0,0 +1,178 @@
//! Widget system types and utilities.
//!
//! Widgets are self-contained frontend components that plug into named
//! [`WidgetSlot`]s in the UI. Each widget has a manifest (`widget.json`)
//! and implementation files (`index.js`, optional `style.css`).
use serde::{Deserialize, Serialize};
/// Widget manifest — metadata about a widget component.
///
/// Stored as `frontend/widgets/{id}/manifest.json` in the workspace.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WidgetManifest {
/// Unique widget identifier (must be a valid HTML attribute value).
pub id: String,
/// Human-readable widget name.
pub name: String,
/// Where this widget is rendered in the UI.
pub slot: WidgetSlot,
/// Optional icon identifier (CSS class or emoji).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub icon: Option<String>,
/// Positioning hint (e.g., `"after:memory"`, `"before:jobs"`).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub position: Option<String>,
}
/// Named insertion points in the UI where widgets can be rendered.
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "snake_case")]
pub enum WidgetSlot {
/// Full tab panel (adds a new tab to the tab bar).
Tab,
/// Banner area above the chat message list.
ChatHeader,
/// Area below the chat input.
ChatFooter,
/// Extra action buttons next to the send button.
ChatActions,
/// Right sidebar panel.
Sidebar,
/// Left side of the status bar.
StatusLeft,
/// Right side of the status bar.
StatusRight,
/// Additional section in the Settings tab.
SettingsSection,
/// Custom inline renderer for structured data in chat messages.
/// Registered via `IronClaw.registerChatRenderer()` on the browser side.
ChatRenderer,
}
/// Prefix every CSS selector with `[data-widget="{widget_id}"]` for style isolation.
///
/// This prevents widget styles from bleeding into the main app or other widgets.
/// The widget container element gets `data-widget="{id}"` set by the runtime.
///
/// # Example
///
/// ```
/// use ironclaw_frontend::scope_css;
///
/// let scoped = scope_css(".title { color: red; }", "my-widget");
/// assert!(scoped.contains("[data-widget=\"my-widget\"] .title"));
/// ```
pub fn scope_css(css: &str, widget_id: &str) -> String {
let prefix = format!("[data-widget=\"{}\"]", widget_id);
let mut result = String::with_capacity(css.len() + css.len() / 4);
let mut chars = css.chars().peekable();
let mut in_block = false;
let mut current_selector = String::new();
while let Some(ch) = chars.next() {
match ch {
'{' if !in_block => {
// Scope each comma-separated selector
let selectors: Vec<&str> = current_selector.split(',').collect();
let scoped: Vec<String> = selectors
.iter()
.map(|s| {
let s = s.trim();
if s.is_empty() || s.starts_with('@') {
s.to_string()
} else {
format!("{} {}", prefix, s)
}
})
.collect();
result.push_str(&scoped.join(", "));
result.push_str(" {");
current_selector.clear();
in_block = true;
}
'}' if in_block => {
result.push('}');
in_block = false;
}
_ if in_block => {
result.push(ch);
}
_ => {
current_selector.push(ch);
}
}
}
// Append any trailing content
if !current_selector.is_empty() {
result.push_str(&current_selector);
}
result
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_widget_manifest_roundtrip() {
let json = serde_json::json!({
"id": "dashboard",
"name": "Analytics Dashboard",
"slot": "tab",
"icon": "chart-bar",
"position": "after:memory"
});
let manifest: WidgetManifest = serde_json::from_value(json).unwrap();
assert_eq!(manifest.id, "dashboard");
assert_eq!(manifest.slot, WidgetSlot::Tab);
assert_eq!(manifest.icon.as_deref(), Some("chart-bar"));
}
#[test]
fn test_widget_slot_serialization() {
assert_eq!(
serde_json::to_string(&WidgetSlot::ChatHeader).unwrap(),
"\"chat_header\""
);
assert_eq!(
serde_json::to_string(&WidgetSlot::SettingsSection).unwrap(),
"\"settings_section\""
);
}
#[test]
fn test_scope_css_basic() {
let input = ".title { color: red; }";
let result = scope_css(input, "my-widget");
assert!(result.contains("[data-widget=\"my-widget\"] .title"));
assert!(result.contains("color: red;"));
}
#[test]
fn test_scope_css_multiple_selectors() {
let input = ".a, .b { margin: 0; }";
let result = scope_css(input, "w");
assert!(result.contains("[data-widget=\"w\"] .a"));
assert!(result.contains("[data-widget=\"w\"] .b"));
}
#[test]
fn test_scope_css_multiple_rules() {
let input = ".a { color: red; } .b { color: blue; }";
let result = scope_css(input, "w");
assert!(result.contains("[data-widget=\"w\"] .a"));
assert!(result.contains("[data-widget=\"w\"] .b"));
}
#[test]
fn test_scope_css_empty() {
assert_eq!(scope_css("", "w"), "");
}
}
@@ -95,6 +95,118 @@ let authFlowPending = false;
let _ghostSuggestion = '';
let currentSettingsSubtab = 'inference';
// --- Hash-based URL Navigation ---
//
// Encodes navigation state in window.location.hash so refreshing
// the page restores the current tab, thread, memory file, job detail, etc.
//
// Hash format: #/{tab}[/{detail}[/{subtab}]]
// #/chat → chat tab, assistant thread
// #/chat/{threadId} → chat tab, specific thread
// #/memory → memory tab, tree root
// #/memory/{path/to/file} → memory tab, specific file
// #/jobs → jobs list
// #/jobs/{jobId} → job detail
// #/routines → routines list
// #/routines/{id} → routine detail
// #/settings/{subtab} → settings tab with specific sub-tab
// #/logs → logs tab
/** Suppress hash-change handling while we're programmatically updating. */
let _suppressHashChange = false;
/** Update the URL hash to reflect current navigation state. */
function updateHash() {
var parts = [currentTab];
switch (currentTab) {
case 'chat':
if (currentThreadId && currentThreadId !== assistantThreadId) {
parts.push(currentThreadId);
}
break;
case 'memory':
if (typeof currentMemoryPath === 'string' && currentMemoryPath) {
parts.push(currentMemoryPath);
}
break;
case 'jobs':
if (typeof currentJobId !== 'undefined' && currentJobId) {
parts.push(currentJobId);
}
break;
case 'routines':
if (typeof currentRoutineId !== 'undefined' && currentRoutineId) {
parts.push(currentRoutineId);
}
break;
case 'settings':
if (currentSettingsSubtab && currentSettingsSubtab !== 'inference') {
parts.push(currentSettingsSubtab);
}
break;
}
var hash = '#/' + parts.join('/');
_suppressHashChange = true;
if (window.location.hash !== hash) {
window.history.replaceState(null, '', hash);
}
_suppressHashChange = false;
}
/** Parse the current URL hash into navigation state. */
function parseHash() {
var hash = window.location.hash || '';
if (!hash.startsWith('#/')) return null;
var parts = hash.substring(2).split('/');
return {
tab: parts[0] || 'chat',
detail: parts.slice(1).join('/') || null,
};
}
/**
* Restore navigation state from the URL hash.
* Called once after authentication and on hashchange events.
*/
function restoreFromHash() {
var state = parseHash();
if (!state) return;
// Switch tab (without recursively updating hash)
if (state.tab && state.tab !== currentTab) {
switchTab(state.tab);
}
// Restore detail state within the tab
if (state.detail) {
switch (state.tab) {
case 'chat':
// Defer thread switch until threads are loaded
window._pendingThreadRestore = state.detail;
break;
case 'memory':
readMemoryFile(state.detail);
break;
case 'jobs':
openJobDetail(state.detail);
break;
case 'routines':
openRoutineDetail(state.detail);
break;
case 'settings':
switchSettingsSubtab(state.detail);
break;
}
}
}
window.addEventListener('hashchange', function() {
if (_suppressHashChange) return;
restoreFromHash();
});
// --- Streaming Debounce State ---
let _streamBuffer = '';
let _streamDebounceTimer = null;
@@ -197,6 +309,8 @@ function authenticate() {
loadThreads();
loadMemoryTree();
loadJobs();
// Restore navigation state from URL hash (tab, thread, memory file, etc.)
restoreFromHash();
// Apply URL log_level param if present, otherwise just sync the dropdown
if (urlLogLevel) {
setServerLogLevel(urlLogLevel);
@@ -1020,6 +1134,128 @@ function sanitizeRenderedHtml(html) {
return '';
}
// ==================== Structured Data Rendering ====================
//
// Detects JSON objects and key-value data in assistant messages and
// renders them as styled cards instead of raw text. Also supports
// extensible chat renderers via IronClaw.registerChatRenderer().
/**
* Post-process a .message-content element to upgrade structured data into cards.
* Runs registered chat renderers first, then falls back to built-in JSON detection.
*/
function upgradeStructuredData(contentEl) {
// 1. Run registered chat renderers
var renderers = (window.IronClaw && IronClaw._chatRenderers) || [];
for (var i = 0; i < renderers.length; i++) {
try {
if (renderers[i].match(contentEl.textContent, contentEl)) {
renderers[i].render(contentEl, contentEl.textContent);
return; // First matching renderer wins
}
} catch (e) {
console.error('[IronClaw] Chat renderer "' + renderers[i].id + '" failed:', e);
}
}
// 2. Built-in: detect and upgrade inline JSON objects
upgradeInlineJson(contentEl);
}
/**
* Find JSON-like objects in text nodes and replace them with styled cards.
*/
function upgradeInlineJson(contentEl) {
// Walk text content looking for JSON objects: {...} patterns
// Only process <p> and top-level text, not code blocks
var paragraphs = contentEl.querySelectorAll('p');
if (paragraphs.length === 0) {
// No <p> tags — markdown might have produced bare text
paragraphs = [contentEl];
}
paragraphs.forEach(function(p) {
// Skip code blocks
if (p.closest('pre') || p.closest('code')) return;
var html = p.innerHTML;
// Match JSON-like objects: {...} (including Python-style single quotes)
var jsonRegex = /(\{[^{}]*(?:\{[^{}]*\}[^{}]*)*\})/g;
var match;
var replaced = false;
while ((match = jsonRegex.exec(html)) !== null) {
var raw = match[1];
// Normalize Python-style single quotes to double quotes for parsing
var normalized = raw.replace(/'/g, '"');
try {
var obj = JSON.parse(normalized);
if (typeof obj === 'object' && obj !== null && !Array.isArray(obj)) {
var card = buildDataCard(obj);
html = html.substring(0, match.index) + card + html.substring(match.index + match[0].length);
replaced = true;
// Reset regex since we modified the string
jsonRegex.lastIndex = match.index + card.length;
}
} catch (e) {
// Not valid JSON — leave as text
}
}
if (replaced) {
p.innerHTML = html;
}
});
}
/**
* Build an HTML data card from a plain object.
*/
function buildDataCard(obj) {
var keys = Object.keys(obj);
if (keys.length === 0) return '';
var rows = '';
for (var i = 0; i < keys.length; i++) {
var key = keys[i];
var value = obj[key];
var displayKey = key.replace(/_/g, ' ');
var valueClass = 'data-card-value';
var valueHtml;
// Special rendering for known value types
if (key === 'status' || key === 'state') {
var badgeClass = 'status-badge';
var sv = String(value).toLowerCase();
if (sv === 'created' || sv === 'active' || sv === 'success' || sv === 'completed' || sv === 'ok' || sv === 'running') {
badgeClass += ' status-success';
} else if (sv === 'failed' || sv === 'error' || sv === 'cancelled' || sv === 'rejected') {
badgeClass += ' status-error';
} else if (sv === 'pending' || sv === 'waiting' || sv === 'queued') {
badgeClass += ' status-pending';
}
valueHtml = '<span class="' + badgeClass + '">' + escapeHtml(String(value)) + '</span>';
} else if (typeof value === 'object' && value !== null) {
valueHtml = '<code>' + escapeHtml(JSON.stringify(value)) + '</code>';
} else {
// Check if value looks like a UUID or ID
var strVal = String(value);
if (/^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i.test(strVal)) {
valueHtml = '<code class="data-card-id">' + escapeHtml(strVal) + '</code>';
} else {
valueHtml = '<span>' + escapeHtml(strVal) + '</span>';
}
}
rows += '<div class="data-card-row">' +
'<span class="data-card-label">' + escapeHtml(displayKey) + '</span>' +
'<span class="' + valueClass + '">' + valueHtml + '</span>' +
'</div>';
}
return '<div class="data-card">' + rows + '</div>';
}
function copyCodeBlock(btn) {
const pre = btn.parentElement;
const code = pre.querySelector('code');
@@ -1824,6 +2060,8 @@ function createMessageElement(role, content) {
} else {
div.setAttribute('data-raw', content);
contentEl.innerHTML = renderMarkdown(content);
// Upgrade structured data (JSON objects, etc.) into styled cards
upgradeStructuredData(contentEl);
// Syntax highlighting for code blocks
if (typeof hljs !== 'undefined') {
requestAnimationFrame(() => {
@@ -2012,6 +2250,19 @@ function loadThreads() {
list.appendChild(item);
}
// Restore thread from URL hash if pending (deferred from restoreFromHash)
if (window._pendingThreadRestore) {
var pendingId = window._pendingThreadRestore;
window._pendingThreadRestore = null;
// Verify the thread exists in the loaded list
var found = (pendingId === assistantThreadId) ||
threads.some(function(t) { return t.id === pendingId; });
if (found) {
switchThread(pendingId);
return;
}
}
// Default to assistant thread on first load if no thread selected
if (!currentThreadId && assistantThreadId) {
switchToAssistant();
@@ -2051,6 +2302,7 @@ function switchToAssistant() {
oldestTimestamp = null;
loadHistory();
loadThreads();
updateHash();
if (window.innerWidth <= 768) {
const sidebar = document.getElementById('thread-sidebar');
sidebar.classList.remove('expanded-mobile');
@@ -2067,6 +2319,7 @@ function switchThread(threadId) {
oldestTimestamp = null;
loadHistory();
loadThreads();
updateHash();
if (window.innerWidth <= 768) {
const sidebar = document.getElementById('thread-sidebar');
sidebar.classList.remove('expanded-mobile');
@@ -2080,6 +2333,7 @@ function createNewThread() {
document.getElementById('chat-messages').innerHTML = '';
showWelcomeCard();
loadThreads();
updateHash();
}).catch((err) => {
showToast('Failed to create thread: ' + err.message, 'error');
});
@@ -2227,6 +2481,7 @@ function switchTab(tab) {
stopPairingPoll();
}
updateTabIndicator();
updateHash();
}
function updateTabIndicator() {
@@ -2372,6 +2627,7 @@ function toggleExpand(node) {
function readMemoryFile(path) {
currentMemoryPath = path;
updateHash();
// Update breadcrumb
document.getElementById('memory-breadcrumb-path').innerHTML = buildBreadcrumb(path);
document.getElementById('memory-edit-btn').style.display = 'inline-block';
@@ -3685,6 +3941,7 @@ function restartJob(jobId) {
function openJobDetail(jobId) {
currentJobId = jobId;
currentJobSubTab = 'activity';
updateHash();
apiFetch('/api/jobs/' + jobId).then((job) => {
renderJobDetail(job);
}).catch((err) => {
@@ -3697,6 +3954,7 @@ function closeJobDetail() {
currentJobId = null;
jobFilesTreeState = null;
loadJobs();
updateHash();
}
function renderJobDetail(job) {
@@ -4193,6 +4451,7 @@ function renderRoutinesList(routines) {
function openRoutineDetail(id) {
currentRoutineId = id;
updateHash();
apiFetch('/api/routines/' + id).then((routine) => {
renderRoutineDetail(routine);
}).catch((err) => {
@@ -4203,6 +4462,7 @@ function openRoutineDetail(id) {
function closeRoutineDetail() {
currentRoutineId = null;
loadRoutines();
updateHash();
}
function renderRoutineDetail(routine) {
@@ -5180,6 +5440,7 @@ function switchSettingsSubtab(subtab) {
document.querySelector('.settings-layout').classList.add('settings-detail-active');
}
loadSettingsSubtab(subtab);
updateHash();
}
function settingsBack() {
@@ -6324,3 +6585,177 @@ document.getElementById('settings-search-input').addEventListener('input', funct
activePanel.appendChild(empty);
}
});
// ==================== Widget Extension System ====================
//
// Provides a registration API for frontend widgets. Widgets are self-contained
// components that plug into named slots in the UI (tabs, sidebar, status bar, etc.).
//
// Widget authors call IronClaw.registerWidget({ id, name, slot, init, ... })
// from their module script. The init() function receives a container DOM element
// and the IronClaw.api object for authenticated fetch, event subscription, etc.
window.IronClaw = window.IronClaw || {};
IronClaw.widgets = new Map();
IronClaw._widgetInitQueue = [];
IronClaw._chatRenderers = [];
/**
* Register a widget component.
* @param {Object} def - Widget definition
* @param {string} def.id - Unique widget identifier
* @param {string} def.name - Display name
* @param {string} def.slot - Target slot ('tab', 'chat_header', etc.)
* @param {string} [def.icon] - Icon identifier
* @param {Function} def.init - Called with (container, api) when widget activates
* @param {Function} [def.activate] - Called when widget becomes visible
* @param {Function} [def.deactivate] - Called when widget is hidden
* @param {Function} [def.destroy] - Called when widget is removed
*/
IronClaw.registerWidget = function(def) {
if (!def.id || !def.init) {
console.error('[IronClaw] Widget registration requires id and init:', def);
return;
}
IronClaw.widgets.set(def.id, def);
if (def.slot === 'tab') {
_addWidgetTab(def);
}
};
/**
* Register a chat renderer for custom inline rendering of structured data.
*
* Chat renderers run against each assistant message. The first renderer
* whose `match()` returns true gets to transform the content.
*
* @param {Object} def - Renderer definition
* @param {string} def.id - Unique identifier
* @param {Function} def.match - (textContent, element) => boolean
* @param {Function} def.render - (element, textContent) => void (mutate element in place)
* @param {number} [def.priority=0] - Higher priority runs first
*/
IronClaw.registerChatRenderer = function(def) {
if (!def.id || !def.match || !def.render) {
console.error('[IronClaw] Chat renderer requires id, match, and render:', def);
return;
}
IronClaw._chatRenderers.push(def);
// Sort by priority (higher first)
IronClaw._chatRenderers.sort(function(a, b) {
return (b.priority || 0) - (a.priority || 0);
});
};
/**
* API object exposed to widgets for safe interaction with the app.
*/
IronClaw.api = {
/** Authenticated fetch wrapper — injects the session token. */
fetch: function(path, opts) {
opts = opts || {};
opts.headers = Object.assign({}, opts.headers || {}, {
'Authorization': 'Bearer ' + token
});
return fetch(path, opts);
},
/** Subscribe to an SSE/WebSocket event type. Returns an unsubscribe function. */
subscribe: function(eventType, handler) {
if (!window._widgetEventHandlers) window._widgetEventHandlers = {};
if (!window._widgetEventHandlers[eventType]) window._widgetEventHandlers[eventType] = [];
window._widgetEventHandlers[eventType].push(handler);
return function() {
var handlers = window._widgetEventHandlers[eventType];
if (handlers) {
var idx = handlers.indexOf(handler);
if (idx !== -1) handlers.splice(idx, 1);
}
};
},
/** Current theme information. */
theme: {
get current() { return document.documentElement.dataset.theme || 'dark'; }
},
/** Internationalization helper. */
i18n: {
t: function(key) { return (window.I18n && window.I18n.t) ? window.I18n.t(key) : key; }
},
/** Navigate to a tab by ID. */
navigate: function(tabId) {
if (typeof switchTab === 'function') switchTab(tabId);
}
};
/**
* Add a widget as a new tab in the tab bar.
* @private
*/
function _addWidgetTab(def) {
var tabBar = document.querySelector('.tab-bar');
var tabContent = document.querySelector('.tab-content') || document.getElementById('tab-content');
if (!tabBar || !tabContent) {
// DOM not ready yet — queue for later
IronClaw._widgetInitQueue.push(def);
return;
}
// Create tab button
var btn = document.createElement('button');
btn.className = 'tab-btn';
btn.dataset.tab = def.id;
btn.textContent = def.name;
if (def.icon) {
btn.dataset.icon = def.icon;
}
btn.addEventListener('click', function() {
if (typeof switchTab === 'function') switchTab(def.id);
});
// Insert before the settings tab (last built-in tab) or at the end
var settingsBtn = tabBar.querySelector('[data-tab="settings"]');
if (settingsBtn) {
tabBar.insertBefore(btn, settingsBtn);
} else {
tabBar.appendChild(btn);
}
// Create container panel
var panel = document.createElement('div');
panel.className = 'tab-panel';
panel.dataset.tab = def.id;
panel.dataset.widget = def.id;
panel.style.display = 'none';
tabContent.appendChild(panel);
// Initialize the widget
try {
def.init(panel, IronClaw.api);
} catch (e) {
console.error('[IronClaw] Widget "' + def.id + '" init failed:', e);
panel.innerHTML = '<div style="padding:2rem;color:var(--color-error,red);">Widget "' +
def.id + '" failed to load: ' + (e.message || e) + '</div>';
}
}
// Apply layout config if injected by the server
if (window.__IRONCLAW_LAYOUT__) {
(function() {
var layout = window.__IRONCLAW_LAYOUT__;
// Apply branding title
if (layout.branding && layout.branding.title) {
var titleEl = document.querySelector('.app-title');
if (titleEl) titleEl.textContent = layout.branding.title;
}
// Apply tab visibility
if (layout.tabs && layout.tabs.hidden) {
layout.tabs.hidden.forEach(function(tabId) {
var btn = document.querySelector('.tab-btn[data-tab="' + tabId + '"]');
if (btn) btn.style.display = 'none';
});
}
})();
}

Before

Width:  |  Height:  |  Size: 3.8 KiB

After

Width:  |  Height:  |  Size: 3.8 KiB

@@ -2106,6 +2106,90 @@ body {
background: var(--accent-subtle);
}
/* --- Data Cards (inline structured data) --- */
.data-card {
display: flex;
flex-direction: column;
gap: 0;
margin: 8px 0;
background: var(--bg-tertiary);
border: 1px solid var(--border);
border-radius: var(--radius-lg);
border-left: 3px solid var(--accent);
overflow: hidden;
}
.data-card-row {
display: flex;
align-items: baseline;
gap: var(--space-3);
padding: 6px 14px;
border-bottom: 1px solid var(--border);
}
.data-card-row:last-child {
border-bottom: none;
}
.data-card-label {
font-size: 12px;
font-weight: 500;
color: var(--text-secondary);
text-transform: capitalize;
min-width: 80px;
flex-shrink: 0;
}
.data-card-value {
font-size: var(--text-sm);
color: var(--text-primary);
word-break: break-word;
}
.data-card-value code {
font-family: var(--font-mono);
font-size: 12px;
padding: 1px 5px;
background: var(--bg-secondary);
border-radius: var(--radius-sm);
}
.data-card-id {
font-family: var(--font-mono);
font-size: 11px;
color: var(--text-secondary);
padding: 1px 5px;
background: var(--bg-secondary);
border-radius: var(--radius-sm);
}
/* Status badges */
.status-badge {
display: inline-block;
padding: 2px 10px;
border-radius: 10px;
font-size: 12px;
font-weight: 600;
text-transform: capitalize;
}
.status-success {
background: rgba(52, 211, 153, 0.15);
color: var(--success, #34d399);
}
.status-error {
background: rgba(248, 113, 113, 0.15);
color: var(--error, #f87171);
}
.status-pending {
background: rgba(251, 191, 36, 0.15);
color: var(--warning, #fbbf24);
}
/* Clickable job rows */
.job-row {
cursor: pointer;
+23
View File
@@ -0,0 +1,23 @@
-- 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);
+21 -4
View File
@@ -276,8 +276,8 @@ impl HeartbeatRunner {
.await;
if report.had_work() {
tracing::info!(
daily_logs_deleted = report.daily_logs_deleted,
conversation_docs_deleted = report.conversation_docs_deleted,
directories_cleaned = ?report.directories_cleaned,
versions_pruned = report.versions_pruned,
"heartbeat: memory hygiene deleted stale documents"
);
}
@@ -590,6 +590,23 @@ 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 {
@@ -616,8 +633,8 @@ pub fn spawn_multi_user_heartbeat(
if report.had_work() {
tracing::info!(
user_id = uid,
daily_logs_deleted = report.daily_logs_deleted,
conversation_docs_deleted = report.conversation_docs_deleted,
directories_cleaned = ?report.directories_cleaned,
versions_pruned = report.versions_pruned,
"multi-user heartbeat: memory hygiene deleted stale documents"
);
}
+149
View File
@@ -0,0 +1,149 @@
//! Frontend extension API handlers.
//!
//! Provides endpoints for reading/writing layout configuration and
//! discovering/serving widget files from the workspace.
use std::sync::Arc;
use axum::{
Json,
extract::{Path, State},
http::{StatusCode, header},
response::IntoResponse,
};
use ironclaw_frontend::{LayoutConfig, WidgetManifest};
use crate::channels::web::auth::AuthenticatedUser;
use crate::channels::web::handlers::memory::resolve_workspace;
use crate::channels::web::server::GatewayState;
/// `GET /api/frontend/layout` — return the current layout configuration.
///
/// Reads `frontend/layout.json` from the workspace. Returns an empty
/// default config if the file doesn't exist.
pub async fn frontend_layout_handler(
State(state): State<Arc<GatewayState>>,
AuthenticatedUser(user): AuthenticatedUser,
) -> Result<Json<LayoutConfig>, (StatusCode, String)> {
let workspace = resolve_workspace(&state, &user).await?;
let layout = match workspace.read("frontend/layout.json").await {
Ok(doc) => serde_json::from_str(&doc.content).unwrap_or_default(),
Err(_) => LayoutConfig::default(),
};
Ok(Json(layout))
}
/// `PUT /api/frontend/layout` — update the layout configuration.
///
/// Writes the provided layout config to `frontend/layout.json` in workspace.
pub async fn frontend_layout_update_handler(
State(state): State<Arc<GatewayState>>,
AuthenticatedUser(user): AuthenticatedUser,
Json(layout): Json<LayoutConfig>,
) -> Result<StatusCode, (StatusCode, String)> {
let workspace = resolve_workspace(&state, &user).await?;
let content = serde_json::to_string_pretty(&layout).map_err(|e| {
(
StatusCode::BAD_REQUEST,
format!("Invalid layout config: {e}"),
)
})?;
workspace
.write("frontend/layout.json", &content)
.await
.map_err(|e| {
tracing::error!("Failed to write layout config: {e}");
(
StatusCode::INTERNAL_SERVER_ERROR,
"Failed to write layout config".to_string(),
)
})?;
Ok(StatusCode::OK)
}
/// `GET /api/frontend/widgets` — list all widget manifests.
///
/// Scans `frontend/widgets/` in workspace for directories containing
/// `manifest.json` and returns their parsed manifests.
pub async fn frontend_widgets_handler(
State(state): State<Arc<GatewayState>>,
AuthenticatedUser(user): AuthenticatedUser,
) -> Result<Json<Vec<WidgetManifest>>, (StatusCode, String)> {
let workspace = resolve_workspace(&state, &user).await?;
let entries = workspace
.list("frontend/widgets/")
.await
.unwrap_or_default();
let mut manifests = Vec::new();
for entry in entries {
if !entry.is_directory {
continue;
}
let manifest_path = format!("frontend/widgets/{}/manifest.json", entry.name());
if let Ok(doc) = workspace.read(&manifest_path).await {
match serde_json::from_str::<WidgetManifest>(&doc.content) {
Ok(manifest) => manifests.push(manifest),
Err(e) => {
tracing::warn!(
path = %manifest_path,
"skipping widget with invalid manifest: {e}"
);
}
}
}
}
Ok(Json(manifests))
}
/// `GET /api/frontend/widget/{id}/{*file}` — serve a widget file.
///
/// Serves JS/CSS files from `frontend/widgets/{id}/{file}` in workspace
/// with appropriate MIME types.
pub async fn frontend_widget_file_handler(
State(state): State<Arc<GatewayState>>,
AuthenticatedUser(user): AuthenticatedUser,
Path((id, file)): Path<(String, String)>,
) -> Result<impl IntoResponse, (StatusCode, String)> {
// Reject path traversal
if id.contains("..") || file.contains("..") {
return Err((StatusCode::BAD_REQUEST, "Invalid path".to_string()));
}
let workspace = resolve_workspace(&state, &user).await?;
let path = format!("frontend/widgets/{}/{}", id, file);
let doc = workspace.read(&path).await.map_err(|_| {
(
StatusCode::NOT_FOUND,
format!("Widget file not found: {path}"),
)
})?;
// Determine MIME type from extension
let content_type = if file.ends_with(".js") {
"application/javascript"
} else if file.ends_with(".css") {
"text/css"
} else if file.ends_with(".json") {
"application/json"
} else {
"text/plain"
};
Ok((
[
(header::CONTENT_TYPE, content_type),
(header::CACHE_CONTROL, "no-cache"),
],
doc.content,
))
}
+30 -32
View File
@@ -15,6 +15,14 @@ use crate::channels::web::auth::AuthenticatedUser;
use crate::channels::web::server::GatewayState;
use crate::channels::web::types::*;
fn db_error(context: &str, e: impl std::fmt::Display) -> (StatusCode, String) {
tracing::error!(%e, context, "Database error in jobs handler");
(
StatusCode::INTERNAL_SERVER_ERROR,
"Internal database error".to_string(),
)
}
pub async fn jobs_list_handler(
State(state): State<Arc<GatewayState>>,
AuthenticatedUser(user): AuthenticatedUser,
@@ -213,10 +221,7 @@ pub async fn jobs_detail_handler(
}
Ok(None) => {}
Err(e) => {
return Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("Database error: {}", e),
));
return Err(db_error("jobs_handler", e));
}
}
@@ -257,10 +262,7 @@ pub async fn jobs_detail_handler(
}))
}
Ok(None) => Err((StatusCode::NOT_FOUND, "Job not found".to_string())),
Err(e) => Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("Database error: {}", e),
)),
Err(e) => Err(db_error("jobs_handler", e)),
}
}
@@ -304,10 +306,7 @@ pub async fn jobs_cancel_handler(
}
Ok(None) => {}
Err(e) => {
return Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("Database error: {}", e),
));
return Err(db_error("jobs_handler", e));
}
}
}
@@ -350,10 +349,7 @@ pub async fn jobs_cancel_handler(
}
Ok(None) => {}
Err(e) => {
return Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("Database error: {}", e),
));
return Err(db_error("jobs_handler", e));
}
}
}
@@ -471,10 +467,7 @@ pub async fn jobs_restart_handler(
}
Ok(None) => {}
Err(e) => {
return Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("Database error: {}", e),
));
return Err(db_error("jobs_handler", e));
}
}
@@ -530,10 +523,7 @@ pub async fn jobs_restart_handler(
})))
}
Ok(None) => Err((StatusCode::NOT_FOUND, "Job not found".to_string())),
Err(e) => Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("Database error: {}", e),
)),
Err(e) => Err(db_error("jobs_handler", e)),
}
}
@@ -609,10 +599,7 @@ pub async fn jobs_prompt_handler(
return Err((StatusCode::NOT_FOUND, "Job not found".to_string()));
}
Err(e) => {
return Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("Database error: {}", e),
));
return Err(db_error("jobs_handler", e));
}
}
}
@@ -667,10 +654,7 @@ pub async fn jobs_events_handler(
return Err((StatusCode::NOT_FOUND, "Job not found".to_string()));
}
Err(e) => {
return Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("Database error: {}", e),
));
return Err(db_error("jobs_handler", e));
}
}
@@ -823,3 +807,17 @@ pub async fn job_files_read_handler(
content,
}))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_db_error_does_not_leak_details() {
let (status, body) = db_error("test_context", "relation \"jobs\" does not exist");
assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR);
assert_eq!(body, "Internal database error");
assert!(!body.contains("relation"));
assert!(!body.contains("does not exist"));
}
}
+1
View File
@@ -16,6 +16,7 @@ pub mod users;
pub mod chat;
#[allow(dead_code)]
pub mod extensions;
pub mod frontend;
#[allow(dead_code)]
pub mod settings;
#[allow(dead_code)]
+3 -3
View File
@@ -13,20 +13,20 @@ use crate::channels::web::types::*;
// --- Static file handlers ---
pub async fn index_handler() -> Html<&'static str> {
Html(include_str!("../static/index.html"))
Html(ironclaw_frontend::assets::INDEX_HTML)
}
pub async fn css_handler() -> impl IntoResponse {
(
[(header::CONTENT_TYPE, "text/css")],
include_str!("../static/style.css"),
ironclaw_frontend::assets::STYLE_CSS,
)
}
pub async fn js_handler() -> impl IntoResponse {
(
[(header::CONTENT_TYPE, "application/javascript")],
include_str!("../static/app.js"),
ironclaw_frontend::assets::APP_JS,
)
}
+28 -9
View File
@@ -33,6 +33,10 @@ use crate::channels::relay::DEFAULT_RELAY_NAME;
use crate::channels::web::auth::{
AuthenticatedUser, CombinedAuthState, UserIdentity, auth_middleware,
};
use crate::channels::web::handlers::frontend::{
frontend_layout_handler, frontend_layout_update_handler, frontend_widget_file_handler,
frontend_widgets_handler,
};
use crate::channels::web::handlers::jobs::{
job_files_list_handler, job_files_read_handler, jobs_cancel_handler, jobs_detail_handler,
jobs_events_handler, jobs_list_handler, jobs_prompt_handler, jobs_restart_handler,
@@ -580,6 +584,16 @@ pub async fn start_server(
"/api/tokens/{id}",
axum::routing::delete(super::handlers::tokens::tokens_revoke_handler),
)
// Frontend extension API
.route(
"/api/frontend/layout",
get(frontend_layout_handler).put(frontend_layout_update_handler),
)
.route("/api/frontend/widgets", get(frontend_widgets_handler))
.route(
"/api/frontend/widget/{id}/{*file}",
get(frontend_widget_file_handler),
)
// Gateway control plane
.route("/api/gateway/status", get(gateway_status_handler))
// OpenAI-compatible API
@@ -726,6 +740,11 @@ pub async fn start_server(
}
// --- Static file handlers ---
//
// All frontend assets are embedded in the `ironclaw_frontend` crate.
// These handlers serve them with appropriate MIME types and cache headers.
use ironclaw_frontend::assets;
async fn index_handler() -> impl IntoResponse {
(
@@ -733,7 +752,7 @@ async fn index_handler() -> impl IntoResponse {
(header::CONTENT_TYPE, "text/html; charset=utf-8"),
(header::CACHE_CONTROL, "no-cache"),
],
include_str!("static/index.html"),
assets::INDEX_HTML,
)
}
@@ -743,7 +762,7 @@ async fn css_handler() -> impl IntoResponse {
(header::CONTENT_TYPE, "text/css"),
(header::CACHE_CONTROL, "no-cache"),
],
include_str!("static/style.css"),
assets::STYLE_CSS,
)
}
@@ -753,7 +772,7 @@ async fn js_handler() -> impl IntoResponse {
(header::CONTENT_TYPE, "application/javascript"),
(header::CACHE_CONTROL, "no-cache"),
],
include_str!("static/app.js"),
assets::APP_JS,
)
}
@@ -763,7 +782,7 @@ async fn theme_init_handler() -> impl IntoResponse {
(header::CONTENT_TYPE, "application/javascript"),
(header::CACHE_CONTROL, "no-cache"),
],
include_str!("static/theme-init.js"),
assets::THEME_INIT_JS,
)
}
@@ -773,7 +792,7 @@ async fn favicon_handler() -> impl IntoResponse {
(header::CONTENT_TYPE, "image/x-icon"),
(header::CACHE_CONTROL, "public, max-age=86400"),
],
include_bytes!("static/favicon.ico").as_slice(),
assets::FAVICON_ICO,
)
}
@@ -783,7 +802,7 @@ async fn i18n_index_handler() -> impl IntoResponse {
(header::CONTENT_TYPE, "application/javascript"),
(header::CACHE_CONTROL, "no-cache"),
],
include_str!("static/i18n/index.js"),
assets::I18N_INDEX_JS,
)
}
@@ -793,7 +812,7 @@ async fn i18n_en_handler() -> impl IntoResponse {
(header::CONTENT_TYPE, "application/javascript"),
(header::CACHE_CONTROL, "no-cache"),
],
include_str!("static/i18n/en.js"),
assets::I18N_EN_JS,
)
}
@@ -803,7 +822,7 @@ async fn i18n_zh_handler() -> impl IntoResponse {
(header::CONTENT_TYPE, "application/javascript"),
(header::CACHE_CONTROL, "no-cache"),
],
include_str!("static/i18n/zh-CN.js"),
assets::I18N_ZH_CN_JS,
)
}
@@ -813,7 +832,7 @@ async fn i18n_app_handler() -> impl IntoResponse {
(header::CONTENT_TYPE, "application/javascript"),
(header::CACHE_CONTROL, "no-cache"),
],
include_str!("static/i18n-app.js"),
assets::I18N_APP_JS,
)
}
+122 -4
View File
@@ -569,6 +569,42 @@ pub async fn sweep_expired_flows(registry: &PendingOAuthRegistry) {
const HOSTED_STATE_PREFIX: &str = "ic2";
const HOSTED_STATE_CHECKSUM_BYTES: usize = 12;
/// Maximum length for a legacy flow ID or instance name.
const LEGACY_STATE_MAX_LEN: usize = 128;
/// Minimum length for a legacy flow ID.
const LEGACY_STATE_MIN_LEN: usize = 8;
/// Validate that a legacy state component (flow_id or instance_name) contains
/// only safe characters: alphanumeric, dash, underscore.
fn is_valid_legacy_state_component(s: &str) -> bool {
!s.is_empty()
&& s.len() <= LEGACY_STATE_MAX_LEN
&& s.bytes()
.all(|b| b.is_ascii_alphanumeric() || b == b'-' || b == b'_')
}
fn validate_legacy_flow_id(flow_id: &str) -> Result<(), String> {
if flow_id.len() < LEGACY_STATE_MIN_LEN {
return Err(format!(
"Legacy OAuth flow_id too short ({} chars, minimum {LEGACY_STATE_MIN_LEN})",
flow_id.len()
));
}
if flow_id.len() > LEGACY_STATE_MAX_LEN {
return Err(format!(
"Legacy OAuth flow_id too long ({} chars, maximum {LEGACY_STATE_MAX_LEN})",
flow_id.len()
));
}
if !flow_id
.bytes()
.all(|b| b.is_ascii_alphanumeric() || b == b'-' || b == b'_')
{
return Err("Legacy OAuth flow_id contains invalid characters".to_string());
}
Ok(())
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DecodedHostedOAuthState {
pub flow_id: String,
@@ -653,6 +689,17 @@ pub fn decode_hosted_oauth_state(state: &str) -> Result<DecodedHostedOAuthState,
if flow_id.is_empty() {
return Err("Hosted OAuth legacy state is missing flow_id".to_string());
}
validate_legacy_flow_id(flow_id)?;
if !instance_name.is_empty() && !is_valid_legacy_state_component(instance_name) {
return Err(format!(
"Legacy OAuth instance name contains invalid characters or exceeds max length ({LEGACY_STATE_MAX_LEN})"
));
}
tracing::debug!(
flow_id,
instance_name,
"Decoded legacy prefixed OAuth state"
);
return Ok(DecodedHostedOAuthState {
flow_id: flow_id.to_string(),
instance_name: if instance_name.is_empty() {
@@ -668,6 +715,9 @@ pub fn decode_hosted_oauth_state(state: &str) -> Result<DecodedHostedOAuthState,
return Err("Hosted OAuth state is empty".to_string());
}
validate_legacy_flow_id(state)?;
tracing::debug!(flow_id = state, "Decoded legacy raw OAuth state");
Ok(DecodedHostedOAuthState {
flow_id: state.to_string(),
instance_name: None,
@@ -1734,13 +1784,13 @@ mod tests {
fn test_decode_hosted_oauth_state_accepts_legacy_formats() {
use crate::cli::oauth_defaults::decode_hosted_oauth_state;
let decoded = decode_hosted_oauth_state("kind-deer:abc123").expect("legacy prefixed");
assert_eq!(decoded.flow_id, "abc123");
let decoded = decode_hosted_oauth_state("kind-deer:abc12345").expect("legacy prefixed");
assert_eq!(decoded.flow_id, "abc12345");
assert_eq!(decoded.instance_name.as_deref(), Some("kind-deer"));
assert!(decoded.is_legacy);
let decoded = decode_hosted_oauth_state("abc123").expect("legacy raw");
assert_eq!(decoded.flow_id, "abc123");
let decoded = decode_hosted_oauth_state("abc12345").expect("legacy raw");
assert_eq!(decoded.flow_id, "abc12345");
assert_eq!(decoded.instance_name, None);
assert!(decoded.is_legacy);
}
@@ -1864,4 +1914,72 @@ mod tests {
assert_eq!(decoded_no_instance.instance_name, None);
assert!(!decoded_no_instance.is_legacy);
}
/// Legacy flow IDs that are too short must be rejected (#1443).
#[test]
fn test_legacy_state_rejects_short_flow_id() {
use crate::cli::oauth_defaults::decode_hosted_oauth_state;
let err = decode_hosted_oauth_state("abc").expect_err("short raw flow_id");
assert!(err.contains("too short"), "unexpected error: {err}");
let err = decode_hosted_oauth_state("inst:abc").expect_err("short prefixed flow_id");
assert!(err.contains("too short"), "unexpected error: {err}");
}
/// Legacy flow IDs with invalid characters must be rejected (#1443).
#[test]
fn test_legacy_state_rejects_invalid_characters() {
use crate::cli::oauth_defaults::decode_hosted_oauth_state;
let err = decode_hosted_oauth_state("flow id with spaces!").expect_err("spaces in flow_id");
assert!(
err.contains("invalid characters"),
"unexpected error: {err}"
);
let err = decode_hosted_oauth_state("inst:flow/id?bad=yes")
.expect_err("special chars in prefixed flow_id");
assert!(
err.contains("invalid characters"),
"unexpected error: {err}"
);
}
/// Legacy instance names with invalid characters must be rejected (#1444).
#[test]
fn test_legacy_state_rejects_invalid_instance_name() {
use crate::cli::oauth_defaults::decode_hosted_oauth_state;
let err = decode_hosted_oauth_state("bad instance!:valid-flow-id-12345")
.expect_err("invalid instance name");
assert!(err.contains("instance name"), "unexpected error: {err}");
}
/// Excessively long legacy flow IDs must be rejected (#1443).
#[test]
fn test_legacy_state_rejects_oversized_flow_id() {
use crate::cli::oauth_defaults::decode_hosted_oauth_state;
let long_id = "a".repeat(200);
let err = decode_hosted_oauth_state(&long_id).expect_err("oversized flow_id");
assert!(err.contains("too long"), "unexpected error: {err}");
}
/// Valid legacy flow IDs at boundary lengths are accepted.
#[test]
fn test_legacy_state_accepts_boundary_lengths() {
use crate::cli::oauth_defaults::decode_hosted_oauth_state;
// Exactly 8 chars (minimum)
let decoded = decode_hosted_oauth_state("abcd1234").expect("8-char flow_id");
assert_eq!(decoded.flow_id, "abcd1234");
assert!(decoded.is_legacy);
// Exactly 128 chars (maximum)
let max_id = "a".repeat(128);
let decoded = decode_hosted_oauth_state(&max_id).expect("128-char flow_id");
assert_eq!(decoded.flow_id, max_id);
assert!(decoded.is_legacy);
}
}
+5 -13
View File
@@ -10,10 +10,8 @@ use crate::error::ConfigError;
pub struct HygieneConfig {
/// Whether hygiene is enabled. Env: `MEMORY_HYGIENE_ENABLED` (default: true).
pub enabled: bool,
/// 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,
/// Maximum versions to keep per document. Env: `MEMORY_HYGIENE_VERSION_KEEP_COUNT` (default: 50).
pub version_keep_count: u32,
/// Minimum hours between hygiene passes. Env: `MEMORY_HYGIENE_CADENCE_HOURS` (default: 12).
pub cadence_hours: u32,
}
@@ -22,8 +20,7 @@ impl Default for HygieneConfig {
fn default() -> Self {
Self {
enabled: true,
daily_retention_days: 30,
conversation_retention_days: 7,
version_keep_count: 50,
cadence_hours: 12,
}
}
@@ -33,11 +30,7 @@ impl HygieneConfig {
pub(crate) fn resolve() -> Result<Self, ConfigError> {
Ok(Self {
enabled: parse_bool_env("MEMORY_HYGIENE_ENABLED", true)?,
daily_retention_days: parse_optional_env("MEMORY_HYGIENE_DAILY_RETENTION_DAYS", 30)?,
conversation_retention_days: parse_optional_env(
"MEMORY_HYGIENE_CONVERSATION_RETENTION_DAYS",
7,
)?,
version_keep_count: parse_optional_env("MEMORY_HYGIENE_VERSION_KEEP_COUNT", 50)?,
cadence_hours: parse_optional_env("MEMORY_HYGIENE_CADENCE_HOURS", 12)?,
})
}
@@ -47,8 +40,7 @@ impl HygieneConfig {
pub fn to_workspace_config(&self) -> crate::workspace::hygiene::HygieneConfig {
crate::workspace::hygiene::HygieneConfig {
enabled: self.enabled,
daily_retention_days: self.daily_retention_days,
conversation_retention_days: self.conversation_retention_days,
version_keep_count: self.version_keep_count,
cadence_hours: self.cadence_hours,
state_dir: ironclaw_base_dir(),
}
+17 -13
View File
@@ -17,7 +17,6 @@ mod workspace;
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use async_trait::async_trait;
use chrono::{DateTime, NaiveDateTime, Utc};
@@ -34,8 +33,6 @@ use crate::workspace::MemoryDocument;
use crate::db::libsql_migrations;
static NAIVE_TIMESTAMP_LOGGED: AtomicBool = AtomicBool::new(false);
/// Explicit column list for routines table (matches positional access in `row_to_routine_libsql`).
pub(crate) const ROUTINE_COLUMNS: &str = "\
id, name, description, user_id, enabled, \
@@ -167,13 +164,11 @@ impl LibSqlBackend {
///
/// Returns an error if none of the formats match.
pub(crate) fn parse_timestamp(s: &str) -> Result<DateTime<Utc>, String> {
let log_naive_timestamp_once = || {
if !NAIVE_TIMESTAMP_LOGGED.swap(true, Ordering::Relaxed) {
tracing::debug!(
timestamp = %s,
"parsed naive timestamp without timezone; assuming UTC for backward compatibility"
);
}
let log_naive_timestamp = || {
tracing::warn!(
timestamp = %s,
"parsed naive timestamp, assuming UTC — consider migrating to RFC 3339"
);
};
// RFC 3339 (our canonical write format)
@@ -182,12 +177,12 @@ pub(crate) fn parse_timestamp(s: &str) -> Result<DateTime<Utc>, String> {
}
// Naive with fractional seconds (legacy or SQLite datetime() output)
if let Ok(ndt) = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S%.f") {
log_naive_timestamp_once();
log_naive_timestamp();
return Ok(ndt.and_utc());
}
// Naive without fractional seconds (legacy format)
if let Ok(ndt) = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S") {
log_naive_timestamp_once();
log_naive_timestamp();
return Ok(ndt.and_utc());
}
Err(format!("unparseable timestamp: {:?}", s))
@@ -439,7 +434,7 @@ mod tests {
use chrono::{TimeZone, Utc};
use crate::db::Database;
use crate::db::libsql::{LibSqlBackend, normalize_notify_user, parse_timestamp};
use crate::db::libsql::{LibSqlBackend, fmt_ts, normalize_notify_user, parse_timestamp};
#[test]
fn test_normalize_notify_user_treats_legacy_default_as_missing() {
@@ -468,6 +463,15 @@ mod tests {
assert_eq!(naive_without_millis, expected);
}
#[test]
fn test_fmt_ts_roundtrips_through_parse_timestamp() {
let original = Utc.with_ymd_and_hms(2026, 6, 15, 8, 30, 45).unwrap()
+ chrono::Duration::milliseconds(123);
let formatted = fmt_ts(&original);
let parsed = parse_timestamp(&formatted).unwrap();
assert_eq!(parsed, original);
}
#[tokio::test]
async fn test_libsql_now_format_is_rfc3339_and_parseable() {
let backend = LibSqlBackend::new_memory().await.unwrap();
+326 -2
View File
@@ -13,8 +13,8 @@ use super::{
use crate::db::WorkspaceStore;
use crate::error::{DatabaseError, WorkspaceError};
use crate::workspace::{
MemoryChunk, MemoryDocument, RankedResult, SearchConfig, SearchResult, WorkspaceEntry,
fuse_results,
DocumentVersion, MemoryChunk, MemoryDocument, RankedResult, SearchConfig, SearchResult,
VersionSummary, WorkspaceEntry, fuse_results,
};
use chrono::Utc;
@@ -840,6 +840,330 @@ 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,6 +785,25 @@ 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);
"#,
),
];
+61
View File
@@ -700,6 +700,67 @@ 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,
+65 -1
View File
@@ -25,7 +25,8 @@ use crate::history::{
LlmCallRecord, SandboxJobRecord, SandboxJobSummary, SettingRow, Store,
};
use crate::workspace::{
MemoryChunk, MemoryDocument, Repository, SearchConfig, SearchResult, WorkspaceEntry,
DocumentVersion, MemoryChunk, MemoryDocument, Repository, SearchConfig, SearchResult,
VersionSummary, WorkspaceEntry,
};
/// PostgreSQL database backend.
@@ -785,6 +786,69 @@ 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,6 +315,12 @@ 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).
+649 -1
View File
@@ -53,6 +53,38 @@ struct HostedOAuthFlowStart {
flow: crate::cli::oauth_defaults::PendingOAuthFlow,
}
#[derive(Debug, Default)]
struct SecretCleanupPlan {
base_secrets: HashSet<String>,
companion_secrets: HashMap<String, HashSet<String>>,
}
impl SecretCleanupPlan {
fn add_base_secret(&mut self, secret_name: impl AsRef<str>) {
self.base_secrets
.insert(secret_name.as_ref().to_lowercase());
}
fn add_companion_secret(
&mut self,
base_secret_name: impl AsRef<str>,
companion_secret_name: impl AsRef<str>,
) {
self.companion_secrets
.entry(base_secret_name.as_ref().to_lowercase())
.or_default()
.insert(companion_secret_name.as_ref().to_lowercase());
}
}
fn oauth_refresh_secret_name(secret_name: &str) -> String {
format!("{}_refresh_token", secret_name.to_lowercase())
}
fn oauth_scopes_secret_name(secret_name: &str) -> String {
format!("{}_scopes", secret_name.to_lowercase())
}
fn normalize_oauth_callback_path(path: &str) -> String {
let trimmed_path = path.trim_end_matches('/');
if trimmed_path.is_empty() {
@@ -1602,6 +1634,10 @@ impl ExtensionManager {
match kind {
ExtensionKind::McpServer => {
let cleanup_plan = self
.collect_secret_cleanup_plan(name, kind, user_id)
.await?;
// Unregister tools with this server's prefix
let tool_names: Vec<String> = self
.tool_registry
@@ -1623,6 +1659,9 @@ impl ExtensionManager {
.await
.map_err(|e| ExtensionError::Config(e.to_string()))?;
self.cleanup_uninstalled_extension_secrets(cleanup_plan, user_id)
.await;
Ok(format!(
"Removed MCP server '{}' and {} tool(s)",
name,
@@ -1630,6 +1669,10 @@ impl ExtensionManager {
))
}
ExtensionKind::WasmTool => {
let cleanup_plan = self
.collect_secret_cleanup_plan(name, kind, user_id)
.await?;
// Unregister from tool registry
self.tool_registry.unregister(name).await;
@@ -1674,9 +1717,16 @@ impl ExtensionManager {
let _ = tokio::fs::remove_file(&cap_path).await;
}
self.cleanup_uninstalled_extension_secrets(cleanup_plan, user_id)
.await;
Ok(format!("Removed WASM tool '{}'", name))
}
ExtensionKind::WasmChannel => {
let cleanup_plan = self
.collect_secret_cleanup_plan(name, kind, user_id)
.await?;
// Remove from active set and persist
self.active_channel_names.write().await.remove(name);
self.persist_active_channels(user_id).await;
@@ -1702,6 +1752,9 @@ impl ExtensionManager {
let _ = tokio::fs::remove_file(&cap_path).await;
}
self.cleanup_uninstalled_extension_secrets(cleanup_plan, user_id)
.await;
Ok(format!(
"Removed channel '{}'. Restart IronClaw for the change to take effect.",
name
@@ -2999,6 +3052,258 @@ impl ExtensionManager {
crate::tools::wasm::CapabilitiesFile::from_bytes(&cap_bytes).ok()
}
async fn load_channel_capabilities(
&self,
name: &str,
) -> Option<crate::channels::wasm::ChannelCapabilitiesFile> {
let cap_path = self
.wasm_channels_dir
.join(format!("{}.capabilities.json", name));
let cap_bytes = tokio::fs::read(&cap_path).await.ok()?;
crate::channels::wasm::ChannelCapabilitiesFile::from_bytes(&cap_bytes).ok()
}
async fn collect_secret_cleanup_plan(
&self,
name: &str,
kind: ExtensionKind,
user_id: &str,
) -> Result<SecretCleanupPlan, ExtensionError> {
let mut plan = SecretCleanupPlan::default();
match kind {
ExtensionKind::WasmTool => {
if let Some(cap) = self.load_tool_capabilities(name).await {
for secret_name in Self::tool_secret_names(&cap) {
plan.add_base_secret(secret_name);
}
if let Some(auth) = cap.auth {
plan.add_base_secret(&auth.secret_name);
plan.add_companion_secret(
&auth.secret_name,
oauth_refresh_secret_name(&auth.secret_name),
);
plan.add_companion_secret(
&auth.secret_name,
oauth_scopes_secret_name(&auth.secret_name),
);
}
}
}
ExtensionKind::WasmChannel => {
if let Some(cap) = self.load_channel_capabilities(name).await {
for secret_name in Self::channel_secret_names(&cap) {
plan.add_base_secret(secret_name);
}
}
}
ExtensionKind::McpServer => {
let server = self
.get_mcp_server(name, user_id)
.await
.map_err(|e| ExtensionError::Config(e.to_string()))?;
let token_secret_name = server.token_secret_name();
plan.add_base_secret(&token_secret_name);
plan.add_base_secret(server.client_id_secret_name());
// MCP OAuth can persist companion secrets through two paths:
// the MCP auth helper uses `mcp_<name>_refresh_token`, while the
// hosted gateway callback stores companions alongside the access
// token secret (`<token_secret>_refresh_token` / `_scopes`).
plan.add_companion_secret(&token_secret_name, server.refresh_token_secret_name());
plan.add_companion_secret(
&token_secret_name,
oauth_refresh_secret_name(&token_secret_name),
);
plan.add_companion_secret(
&token_secret_name,
oauth_scopes_secret_name(&token_secret_name),
);
}
ExtensionKind::ChannelRelay => {}
}
Ok(plan)
}
async fn cleanup_uninstalled_extension_secrets(&self, plan: SecretCleanupPlan, user_id: &str) {
let referenced_secrets = match self.collect_referenced_secret_names(user_id).await {
Ok(secret_names) => secret_names,
Err(error) => {
tracing::warn!(
user_id,
error,
"Failed to determine which secrets are still referenced; keeping secrets"
);
return;
}
};
for base_secret in &plan.base_secrets {
if referenced_secrets.contains(base_secret) {
continue;
}
self.delete_secret_best_effort(user_id, base_secret).await;
if let Some(companion_secrets) = plan.companion_secrets.get(base_secret) {
for companion_secret in companion_secrets {
if !referenced_secrets.contains(companion_secret) {
self.delete_secret_best_effort(user_id, companion_secret)
.await;
}
}
}
}
}
async fn delete_secret_best_effort(&self, user_id: &str, secret_name: &str) {
if let Err(error) = self.secrets.delete(user_id, secret_name).await {
tracing::warn!(
user_id,
secret_name,
error = %error,
"Failed to delete secret while uninstalling extension"
);
}
}
async fn collect_referenced_secret_names(
&self,
user_id: &str,
) -> Result<HashSet<String>, String> {
let mut referenced_secret_names = HashSet::new();
let tools = discover_tools(&self.wasm_tools_dir)
.await
.map_err(|e| format!("discover tools: {e}"))?;
for (tool_name, discovered_tool) in &tools {
let cap = self
.load_tool_capabilities(tool_name)
.await
.ok_or_else(|| {
let path = discovered_tool
.capabilities_path
.as_ref()
.map(|path| path.display().to_string())
.unwrap_or_else(|| format!("{} (missing)", tool_name));
format!("load tool capabilities for {tool_name}: {path}")
})?;
referenced_secret_names.extend(Self::tool_secret_names(&cap));
}
let channels = crate::channels::wasm::discover_channels(&self.wasm_channels_dir)
.await
.map_err(|e| format!("discover channels: {e}"))?;
for (channel_name, discovered_channel) in &channels {
let cap = self
.load_channel_capabilities(channel_name)
.await
.ok_or_else(|| {
let path = discovered_channel
.capabilities_path
.as_ref()
.map(|path| path.display().to_string())
.unwrap_or_else(|| format!("{} (missing)", channel_name));
format!("load channel capabilities for {channel_name}: {path}")
})?;
referenced_secret_names.extend(Self::channel_secret_names(&cap));
}
let mcp_servers = self
.load_mcp_servers(user_id)
.await
.map_err(|e| format!("load MCP servers: {e}"))?;
for server in &mcp_servers.servers {
referenced_secret_names.extend(Self::mcp_server_secret_names(server));
}
Ok(referenced_secret_names)
}
fn tool_secret_names(cap: &crate::tools::wasm::CapabilitiesFile) -> HashSet<String> {
let mut names = HashSet::new();
if let Some(auth) = &cap.auth {
names.insert(auth.secret_name.to_lowercase());
}
if let Some(setup) = &cap.setup {
names.extend(
setup
.required_secrets
.iter()
.map(|secret| secret.name.to_lowercase()),
);
}
if let Some(http) = &cap.http {
names.extend(
http.credentials
.values()
.map(|credential| credential.secret_name.to_lowercase()),
);
}
if let Some(webhook) = &cap.webhook {
if let Some(secret_name) = &webhook.secret_name {
names.insert(secret_name.to_lowercase());
}
if let Some(secret_name) = &webhook.signature_key_secret_name {
names.insert(secret_name.to_lowercase());
}
if let Some(secret_name) = &webhook.hmac_secret_name {
names.insert(secret_name.to_lowercase());
}
}
names
}
fn channel_secret_names(
cap: &crate::channels::wasm::ChannelCapabilitiesFile,
) -> HashSet<String> {
let mut names: HashSet<String> = cap
.setup
.required_secrets
.iter()
.map(|secret| secret.name.to_lowercase())
.collect();
if let Some(http) = cap.capabilities.tool.http.as_ref() {
names.extend(
http.credentials
.values()
.map(|credential| credential.secret_name.to_lowercase()),
);
}
if let Some(webhook) = cap
.capabilities
.channel
.as_ref()
.and_then(|channel| channel.webhook.as_ref())
{
if webhook.secret_header.is_some() || webhook.secret_name.is_some() {
names.insert(cap.webhook_secret_name().to_lowercase());
}
if let Some(secret_name) = cap.signature_key_secret_name() {
names.insert(secret_name.to_lowercase());
}
if let Some(secret_name) = cap.hmac_secret_name() {
names.insert(secret_name.to_lowercase());
}
}
names
}
fn mcp_server_secret_names(server: &McpServerConfig) -> HashSet<String> {
[
server.token_secret_name().to_lowercase(),
server.client_id_secret_name().to_lowercase(),
]
.into_iter()
.collect()
}
/// Collect merged OAuth scopes from all installed tools sharing the same secret_name.
///
/// When multiple tools share an OAuth provider (e.g., google-calendar and google-drive
@@ -6033,6 +6338,8 @@ mod tests {
ExtensionError, ExtensionKind, ExtensionSource, InstallResult, VerificationChallenge,
};
use crate::pairing::PairingStore;
use crate::secrets::CreateSecretParams;
use crate::tools::mcp::McpServerConfig;
fn require(condition: bool, message: impl Into<String>) -> Result<(), String> {
if condition {
@@ -6353,6 +6660,38 @@ mod tests {
tools_dir
}
fn write_test_channel(
dir: &std::path::Path,
name: &str,
capabilities_json: &str,
) -> std::path::PathBuf {
let channels_dir = dir.join("channels");
std::fs::create_dir_all(&channels_dir).expect("channels dir");
std::fs::write(
channels_dir.join(format!("{name}.wasm")),
b"not-a-real-wasm",
)
.expect("wasm");
std::fs::write(
channels_dir.join(format!("{name}.capabilities.json")),
capabilities_json,
)
.expect("capabilities");
channels_dir
}
async fn store_test_secret(
manager: &crate::extensions::manager::ExtensionManager,
name: &str,
value: &str,
) {
manager
.secrets
.create("test", CreateSecretParams::new(name, value))
.await
.expect("store secret");
}
#[test]
fn test_setting_value_is_present() {
assert!(
@@ -7423,7 +7762,13 @@ mod tests {
// Regression: remove() only checked channel_runtime for shutdown, missing
// relay-only mode where only relay_channel_manager is set.
let dir = tempfile::tempdir().expect("temp dir");
let mgr = make_test_manager(None, dir.path().to_path_buf());
let (store, _db_dir) = make_test_store().await;
let mgr = make_test_manager_with_dirs(
None,
dir.path().join("tools"),
dir.path().join("channels"),
Some(store),
);
// Set up relay channel manager with a stub channel
let cm = Arc::new(crate::channels::ChannelManager::new());
@@ -7450,6 +7795,8 @@ mod tests {
.await
.expect("store team_id");
}
store_test_secret(&mgr, "relay:slack-relay:oauth_state", "nonce").await;
store_test_secret(&mgr, "relay:slack-relay:stream_token", "legacy-token").await;
// Verify channel exists before removal
assert!(cm.get_channel("slack-relay").await.is_some());
@@ -7478,6 +7825,30 @@ mod tests {
cm.get_channel("slack-relay").await.is_none(),
"relay channel should be removed from the channel manager"
);
assert!(
!mgr.secrets
.exists("test", "relay:slack-relay:oauth_state")
.await
.expect("oauth state exists query"),
"relay oauth_state secret should be removed"
);
assert!(
!mgr.secrets
.exists("test", "relay:slack-relay:stream_token")
.await
.expect("stream token exists query"),
"relay legacy stream token should be removed"
);
assert_eq!(
mgr.store
.as_ref()
.expect("store")
.get_setting("test", "relay:slack-relay:team_id")
.await
.expect("team_id query"),
None,
"relay team_id setting should be removed"
);
}
#[tokio::test]
@@ -7585,6 +7956,185 @@ mod tests {
);
}
#[tokio::test]
async fn test_remove_wasm_tool_deletes_unique_secrets() {
let dir = tempfile::tempdir().expect("temp dir");
let tools_dir = write_test_tool(
dir.path(),
"github",
r#"{
"name": "github",
"auth": { "secret_name": "github_token" },
"setup": {
"required_secrets": [
{ "name": "github_client_secret", "prompt": "GitHub client secret for testing cleanup behavior." }
]
},
"http": {
"credentials": {
"service_token": {
"secret_name": "github_service_token",
"location": { "type": "bearer" }
}
}
},
"webhook": {
"hmac_secret_name": "github_webhook_secret"
}
}"#,
);
let mgr = make_test_manager_with_dirs(None, tools_dir, dir.path().join("channels"), None);
store_test_secret(&mgr, "github_token", "access-token").await;
store_test_secret(&mgr, "github_token_refresh_token", "refresh-token").await;
store_test_secret(&mgr, "github_token_scopes", "repo workflow").await;
store_test_secret(&mgr, "github_client_secret", "client-secret").await;
store_test_secret(&mgr, "github_service_token", "service-token").await;
store_test_secret(&mgr, "github_webhook_secret", "webhook-secret").await;
mgr.remove("github", "test")
.await
.expect("remove should succeed");
for secret_name in [
"github_token",
"github_token_refresh_token",
"github_token_scopes",
"github_client_secret",
"github_service_token",
"github_webhook_secret",
] {
assert!(
!mgr.secrets
.exists("test", secret_name)
.await
.expect("exists query"),
"secret {secret_name} should be deleted"
);
}
}
#[tokio::test]
async fn test_remove_wasm_tool_keeps_secrets_when_other_tool_capabilities_missing() {
let dir = tempfile::tempdir().expect("temp dir");
let tools_dir = write_test_tool(
dir.path(),
"github",
r#"{
"name": "github",
"auth": { "secret_name": "shared_token" }
}"#,
);
std::fs::write(tools_dir.join("broken.wasm"), b"fake-tool").expect("write tool");
let mgr = make_test_manager_with_dirs(None, tools_dir, dir.path().join("channels"), None);
store_test_secret(&mgr, "shared_token", "access-token").await;
store_test_secret(&mgr, "shared_token_refresh_token", "refresh-token").await;
store_test_secret(&mgr, "shared_token_scopes", "repo").await;
mgr.remove("github", "test")
.await
.expect("remove should succeed");
for secret_name in [
"shared_token",
"shared_token_refresh_token",
"shared_token_scopes",
] {
assert!(
mgr.secrets
.exists("test", secret_name)
.await
.expect("exists query"),
"secret {secret_name} should be retained when reference detection is uncertain"
);
}
}
#[tokio::test]
async fn test_remove_wasm_tool_keeps_shared_secrets_until_last_extension() {
let dir = tempfile::tempdir().expect("temp dir");
write_test_tool(
dir.path(),
"google-calendar",
r#"{
"name": "google-calendar",
"auth": { "secret_name": "google_oauth_token" },
"setup": {
"required_secrets": [
{ "name": "google_oauth_client_id", "prompt": "Google OAuth client id for cleanup testing." },
{ "name": "google_oauth_client_secret", "prompt": "Google OAuth client secret for cleanup testing." }
]
}
}"#,
);
let tools_dir = write_test_tool(
dir.path(),
"google-drive",
r#"{
"name": "google-drive",
"auth": { "secret_name": "google_oauth_token" },
"setup": {
"required_secrets": [
{ "name": "google_oauth_client_id", "prompt": "Google OAuth client id for cleanup testing." },
{ "name": "google_oauth_client_secret", "prompt": "Google OAuth client secret for cleanup testing." }
]
}
}"#,
);
let mgr = make_test_manager_with_dirs(None, tools_dir, dir.path().join("channels"), None);
for (secret_name, value) in [
("google_oauth_token", "access-token"),
("google_oauth_token_refresh_token", "refresh-token"),
("google_oauth_token_scopes", "calendar drive"),
("google_oauth_client_id", "client-id"),
("google_oauth_client_secret", "client-secret"),
] {
store_test_secret(&mgr, secret_name, value).await;
}
mgr.remove("google-calendar", "test")
.await
.expect("first remove should succeed");
for secret_name in [
"google_oauth_token",
"google_oauth_token_refresh_token",
"google_oauth_token_scopes",
"google_oauth_client_id",
"google_oauth_client_secret",
] {
assert!(
mgr.secrets
.exists("test", secret_name)
.await
.expect("exists query"),
"shared secret {secret_name} should remain while google-drive is still installed"
);
}
mgr.remove("google-drive", "test")
.await
.expect("second remove should succeed");
for secret_name in [
"google_oauth_token",
"google_oauth_token_refresh_token",
"google_oauth_token_scopes",
"google_oauth_client_id",
"google_oauth_client_secret",
] {
assert!(
!mgr.secrets
.exists("test", secret_name)
.await
.expect("exists query"),
"shared secret {secret_name} should be deleted after the last tool is removed"
);
}
}
#[tokio::test]
async fn test_remove_wasm_channel_clears_activation_error_and_deletes_files() {
let dir = tempfile::tempdir().expect("temp dir");
@@ -7619,6 +8169,104 @@ mod tests {
);
}
#[tokio::test]
async fn test_remove_wasm_channel_deletes_setup_secrets() {
let dir = tempfile::tempdir().expect("temp dir");
let channels_dir = write_test_channel(
dir.path(),
"telegram",
r#"{
"type": "channel",
"name": "telegram",
"setup": {
"required_secrets": [
{
"name": "telegram_bot_token",
"prompt": "Telegram bot token used to verify uninstall cleanup behavior."
}
]
},
"capabilities": {
"http": {
"credentials": {
"tenant_token": {
"secret_name": "telegram_service_token",
"location": { "type": "bearer" }
}
}
},
"channel": {
"webhook": {
"secret_header": "X-Telegram-Bot-Api-Secret-Token",
"secret_name": "telegram_webhook_secret"
}
}
}
}"#,
);
let mgr = make_test_manager_with_dirs(None, dir.path().join("tools"), channels_dir, None);
store_test_secret(&mgr, "telegram_bot_token", "123:telegram-token").await;
store_test_secret(&mgr, "telegram_service_token", "tenant-service-token").await;
store_test_secret(&mgr, "telegram_webhook_secret", "webhook-secret").await;
mgr.remove("telegram", "test")
.await
.expect("remove should succeed");
for secret_name in [
"telegram_bot_token",
"telegram_service_token",
"telegram_webhook_secret",
] {
assert!(
!mgr.secrets
.exists("test", secret_name)
.await
.expect("exists query"),
"channel secret {secret_name} should be deleted"
);
}
}
#[tokio::test]
async fn test_remove_mcp_server_deletes_stored_secrets() {
let dir = tempfile::tempdir().expect("temp dir");
let (store, _db_dir) = make_test_store().await;
let mgr = make_test_manager_with_dirs(
None,
dir.path().join("tools"),
dir.path().join("channels"),
Some(Arc::clone(&store)),
);
let server = McpServerConfig::new("notion", "https://example.com/mcp");
mgr.add_mcp_server(server.clone(), "test")
.await
.expect("add mcp server");
store_test_secret(&mgr, &server.token_secret_name(), "access-token").await;
store_test_secret(&mgr, &server.refresh_token_secret_name(), "refresh-token").await;
store_test_secret(&mgr, &server.client_id_secret_name(), "client-id").await;
mgr.remove("notion", "test")
.await
.expect("remove should succeed");
for secret_name in [
server.token_secret_name(),
server.refresh_token_secret_name(),
server.client_id_secret_name(),
] {
assert!(
!mgr.secrets
.exists("test", &secret_name)
.await
.expect("exists query"),
"MCP secret {secret_name} should be deleted"
);
}
}
#[test]
fn test_sanitize_url_with_query_params() {
let url = "https://api.example.com/path?api_key=secret123&token=abc";
+1
View File
@@ -234,6 +234,7 @@ fn is_transient(err: &LlmError) -> bool {
LlmError::RequestFailed { .. }
| LlmError::RateLimited { .. }
| LlmError::InvalidResponse { .. }
| LlmError::EmptyResponse { .. }
| LlmError::SessionExpired { .. }
| LlmError::SessionRenewalFailed { .. }
| LlmError::Http(_)
+3
View File
@@ -17,6 +17,9 @@ pub enum LlmError {
#[error("Invalid response from {provider}: {reason}")]
InvalidResponse { provider: String, reason: String },
#[error("Empty response from {provider}: no content returned")]
EmptyResponse { provider: String },
#[error("Context length exceeded: {used} tokens used, {limit} allowed")]
ContextLengthExceeded { used: usize, limit: usize },
+2 -4
View File
@@ -231,9 +231,8 @@ impl LlmProvider for GithubCopilotProvider {
.choices
.into_iter()
.next()
.ok_or_else(|| LlmError::InvalidResponse {
.ok_or_else(|| LlmError::EmptyResponse {
provider: "github_copilot".to_string(),
reason: "No choices in response".to_string(),
})?;
let (content, _tool_calls) = extract_choice_content(&choice);
@@ -309,9 +308,8 @@ impl LlmProvider for GithubCopilotProvider {
.choices
.into_iter()
.next()
.ok_or_else(|| LlmError::InvalidResponse {
.ok_or_else(|| LlmError::EmptyResponse {
provider: "github_copilot".to_string(),
reason: "No choices in response".to_string(),
})?;
let (content, tool_calls) = extract_choice_content(&choice);
+2 -4
View File
@@ -490,9 +490,8 @@ impl LlmProvider for NearAiChatProvider {
.choices
.into_iter()
.next()
.ok_or_else(|| LlmError::InvalidResponse {
.ok_or_else(|| LlmError::EmptyResponse {
provider: "nearai_chat".to_string(),
reason: "No choices in response".to_string(),
})?;
// Fall back to reasoning_content when content is null (same as
@@ -570,9 +569,8 @@ impl LlmProvider for NearAiChatProvider {
.choices
.into_iter()
.next()
.ok_or_else(|| LlmError::InvalidResponse {
.ok_or_else(|| LlmError::EmptyResponse {
provider: "nearai_chat".to_string(),
reason: "No choices in response".to_string(),
})?;
let tool_calls: Vec<ToolCall> = choice
+1
View File
@@ -48,6 +48,7 @@ pub(crate) fn is_retryable(err: &LlmError) -> bool {
LlmError::RequestFailed { .. }
| LlmError::RateLimited { .. }
| LlmError::InvalidResponse { .. }
| LlmError::EmptyResponse { .. }
| LlmError::SessionRenewalFailed { .. }
| LlmError::Http(_)
| LlmError::Io(_)
+140 -4
View File
@@ -246,9 +246,26 @@ 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": ["content"]
"required": []
})
}
@@ -259,7 +276,9 @@ impl Tool for MemoryWriteTool {
) -> Result<ToolOutput, ToolError> {
let start = std::time::Instant::now();
let content = require_str(&params, "content")?;
// 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 target = params
.get("target")
@@ -299,9 +318,9 @@ impl Tool for MemoryWriteTool {
return Ok(ToolOutput::success(output, start.elapsed()));
}
if content.trim().is_empty() {
if !is_patch_mode && content.trim().is_empty() {
return Err(ToolError::InvalidParameters(
"content cannot be empty".to_string(),
"content cannot be empty (use old_string/new_string for patch mode)".to_string(),
));
}
@@ -330,6 +349,46 @@ 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 {
@@ -433,6 +492,24 @@ 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,
@@ -501,6 +578,15 @@ 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"]
@@ -525,11 +611,61 @@ 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,
+4 -2
View File
@@ -124,8 +124,10 @@ impl WasmToolLoader {
let wasm_bytes = fs::read(wasm_path).await?;
// Read capabilities (optional) and extract OAuth refresh config
// and tool description. Parameter schema is auto-derived from the
// WASM module's schema() export (see WasmToolSchemas::compact_schema).
// and tool description. Parameter schema is NOT read from the
// capabilities file — it is auto-derived from the WASM module's
// schema() export at prepare time (see WasmToolSchemas::compact_schema),
// so no schema override is needed here.
let (capabilities, oauth_refresh, description) = if let Some(cap_path) = capabilities_path {
if cap_path.exists() {
let cap_bytes = fs::read(cap_path).await?;
+92 -5
View File
@@ -759,14 +759,31 @@ impl WasmToolSchemas {
}
let kept: serde_json::Map<String, serde_json::Value> = all_properties
.into_iter()
.iter()
.filter(|(name, prop)| {
required.contains(name) || prop.get("enum").is_some() || prop.get("const").is_some()
required.contains(name.as_str())
|| prop.get("enum").is_some()
|| prop.get("const").is_some()
})
.map(|(k, v)| (k.clone(), v.clone()))
.collect();
if kept.is_empty() {
return Self::permissive_schema();
// When the schema has typed properties but none survived the
// required/enum filter, include all typed properties so the LLM
// sees meaningful parameter hints instead of permissive `{}`.
let typed: serde_json::Map<String, serde_json::Value> = all_properties
.into_iter()
.filter(|(_, prop)| schema_is_typed_property(prop))
.collect();
if typed.is_empty() {
return Self::permissive_schema();
}
return serde_json::json!({
"type": "object",
"properties": typed,
"additionalProperties": true,
});
}
let kept_required: Vec<serde_json::Value> = required
@@ -1991,6 +2008,58 @@ mod tests {
);
}
#[tokio::test]
async fn test_typed_schema_without_required_is_advertised() {
// Regression test for #1303: when a WASM tool exports a typed schema
// with no required/enum fields, the advertised schema should still
// contain the typed properties instead of falling back to permissive {}.
let discovery_schema = serde_json::json!({
"type": "object",
"properties": {
"query": { "type": "string" },
"limit": { "type": "integer" }
}
});
let runtime = Arc::new(WasmToolRuntime::new(WasmRuntimeConfig::for_testing()).unwrap());
let prepared = runtime
.prepare("typed_search", b"\0asm\x0d\0\x01\0", None)
.await
.unwrap();
let mut wrapper =
super::WasmToolWrapper::new(Arc::clone(&runtime), prepared, Capabilities::default());
wrapper.schemas = super::WasmToolSchemas::new(discovery_schema.clone());
wrapper.description = "Typed search tool".to_string();
let advertised = wrapper.parameters_schema();
let props = advertised["properties"].as_object().unwrap();
// Both typed properties should be preserved in the advertised schema
assert!(
props.contains_key("query"),
"advertised schema should contain 'query' property"
);
assert!(
props.contains_key("limit"),
"advertised schema should contain 'limit' property"
);
assert_eq!(props.len(), 2);
// The schema should NOT be permissive
assert!(
!super::WasmToolSchemas::is_permissive_schema(&advertised),
"advertised schema should not be permissive when typed properties exist"
);
// No tool_info hint needed since typed properties are visible
let schema = wrapper.schema();
assert!(
!schema.description.contains("tool_info"),
"description should not contain tool_info hint: {}",
schema.description
);
}
#[test]
fn test_compact_schema_keeps_required_and_enum_properties() {
let schema = serde_json::json!({
@@ -2028,8 +2097,8 @@ mod tests {
}
#[test]
fn test_compact_schema_falls_back_to_permissive_when_empty() {
// No required, no enum → permissive fallback
fn test_compact_schema_preserves_typed_properties_when_no_required() {
// No required, no enum, but typed properties → keep all typed props
let schema = serde_json::json!({
"type": "object",
"properties": {
@@ -2038,6 +2107,24 @@ mod tests {
}
});
let compacted = super::WasmToolSchemas::compact_schema(&schema);
let props = compacted["properties"].as_object().unwrap();
assert_eq!(props.len(), 2);
assert!(props.contains_key("query"));
assert!(props.contains_key("limit"));
assert_eq!(compacted["additionalProperties"], true);
}
#[test]
fn test_compact_schema_falls_back_to_permissive_when_no_typed_properties() {
// Properties with no type info → permissive fallback
let schema = serde_json::json!({
"type": "object",
"properties": {
"data": {}
}
});
let compacted = super::WasmToolSchemas::compact_schema(&schema);
assert!(compacted["properties"].as_object().unwrap().is_empty());
}
+144 -36
View File
@@ -31,6 +31,7 @@ use std::time::Duration;
use serde::{Deserialize, Serialize};
use tokio::io::{AsyncBufReadExt, BufReader};
#[cfg(not(unix))]
use tokio::process::Command;
use uuid::Uuid;
@@ -340,6 +341,11 @@ impl ClaudeBridgeRuntime {
/// Spawn a `claude` CLI process and stream its output.
///
/// Uses a PTY on Unix so Node.js line-buffers stdout instead of
/// full-buffering (which causes the bridge to hang on non-TTY pipes).
/// Arguments are passed via `execve` (no shell) — injection-safe by
/// construction.
///
/// Returns the session_id if captured from the `system` init message.
async fn run_claude_session(
&self,
@@ -347,47 +353,102 @@ impl ClaudeBridgeRuntime {
resume_session_id: Option<&str>,
extra_env: &std::collections::HashMap<String, String>,
) -> Result<Option<String>, WorkerError> {
let mut cmd = Command::new("claude");
cmd.arg("-p")
.arg(prompt)
.arg("--output-format")
.arg("stream-json")
.arg("--verbose")
.arg("--max-turns")
.arg(self.config.max_turns.to_string())
.arg("--model")
.arg(&self.config.model);
let max_turns_str = self.config.max_turns.to_string();
if let Some(sid) = resume_session_id {
cmd.arg("--resume").arg(sid);
}
// Inject credentials into the child process environment without
// mutating the global process env (which is unsafe in multi-threaded programs).
cmd.envs(extra_env);
cmd.current_dir("/workspace")
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped());
let mut child = cmd.spawn().map_err(|e| WorkerError::ExecutionFailed {
reason: format!("failed to spawn claude: {}", e),
})?;
let stdout = child
.stdout
.take()
.ok_or_else(|| WorkerError::ExecutionFailed {
reason: "failed to capture claude stdout".to_string(),
// Spawn with PTY on Unix to fix Node.js stdout buffering.
// All arguments are passed individually via execve — never through
// a shell interpreter. This eliminates shell injection by construction.
#[cfg(unix)]
let (mut child, stdout, stderr) = {
let (pty, pts) = pty_process::open().map_err(|e| WorkerError::ExecutionFailed {
reason: format!("failed to allocate PTY: {}", e),
})?;
let stderr = child
.stderr
.take()
.ok_or_else(|| WorkerError::ExecutionFailed {
reason: "failed to capture claude stderr".to_string(),
let mut cmd = pty_process::Command::new("claude");
cmd = cmd
.arg("-p")
.arg(prompt)
.arg("--output-format")
.arg("stream-json")
.arg("--verbose")
.arg("--max-turns")
.arg(&max_turns_str)
.arg("--model")
.arg(&self.config.model);
if let Some(sid) = resume_session_id {
cmd = cmd.arg("--resume").arg(sid);
}
cmd = cmd.envs(extra_env.iter());
cmd = cmd.current_dir("/workspace");
// Keep stderr on a separate pipe — pty-process attaches the PTY
// to all fds by default, which would merge stderr into the PTY
// stream and break NDJSON parsing.
cmd = cmd.stderr(std::process::Stdio::piped());
let mut child = cmd.spawn(pts).map_err(|e| WorkerError::ExecutionFailed {
reason: format!("failed to spawn claude with PTY: {}", e),
})?;
let stderr = child
.stderr
.take()
.ok_or_else(|| WorkerError::ExecutionFailed {
reason: "failed to capture claude stderr".to_string(),
})?;
// stdout comes from the PTY master, which implements AsyncRead
let stdout: Box<dyn tokio::io::AsyncRead + Unpin + Send> = Box::new(pty);
(child, stdout, stderr)
};
// Non-Unix fallback (Windows CI) — no PTY, direct spawn.
// Claude bridge only runs in Linux Docker containers, so this path
// exists solely for compilation on Windows targets.
#[cfg(not(unix))]
let (mut child, stdout, stderr) = {
let mut cmd = Command::new("claude");
cmd.arg("-p")
.arg(prompt)
.arg("--output-format")
.arg("stream-json")
.arg("--verbose")
.arg("--max-turns")
.arg(&max_turns_str)
.arg("--model")
.arg(&self.config.model);
if let Some(sid) = resume_session_id {
cmd.arg("--resume").arg(sid);
}
cmd.envs(extra_env);
cmd.current_dir("/workspace")
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped());
let mut child = cmd.spawn().map_err(|e| WorkerError::ExecutionFailed {
reason: format!("failed to spawn claude: {}", e),
})?;
let stdout_pipe = child
.stdout
.take()
.ok_or_else(|| WorkerError::ExecutionFailed {
reason: "failed to capture claude stdout".to_string(),
})?;
let stderr = child
.stderr
.take()
.ok_or_else(|| WorkerError::ExecutionFailed {
reason: "failed to capture claude stderr".to_string(),
})?;
let stdout: Box<dyn tokio::io::AsyncRead + Unpin + Send> = Box::new(stdout_pipe);
(child, stdout, stderr)
};
// Spawn stderr reader that forwards lines as log events
let client_for_stderr = Arc::clone(&self.client);
let job_id = self.config.job_id;
@@ -1027,4 +1088,51 @@ mod tests {
let copied = copy_dir_recursive(nonexistent, dst.path()).unwrap();
assert_eq!(copied, 0);
}
/// Regression test: arguments are passed individually (not via shell string),
/// so shell metacharacters in prompt/model/session_id are harmless.
#[test]
fn command_args_no_shell_interpretation() {
// Prompt, model, and session_id may contain shell metacharacters from
// user-supplied task descriptions or LLM output. Since we use
// Command::arg() (execve), these are passed as literal strings.
let prompt = "Fix the user's bug; echo $HOME && rm -rf /";
let model = "claude-3-opus-20240229";
let session_id = "'; DROP TABLE jobs; --";
let max_turns = 10u32;
let max_turns_str = max_turns.to_string();
let args: Vec<&str> = vec![
"-p",
prompt,
"--output-format",
"stream-json",
"--verbose",
"--max-turns",
&max_turns_str,
"--model",
model,
"--resume",
session_id,
];
// All values present as literal strings — no shell interpretation
// ["-p", prompt, "--output-format", "stream-json", "--verbose",
// "--max-turns", "10", "--model", model, "--resume", session_id]
assert_eq!(args[1], prompt);
assert_eq!(args[8], model);
assert_eq!(args[10], session_id);
// Shell metacharacters preserved, not expanded
assert!(args[1].contains("$HOME"));
assert!(args[1].contains("&&"));
assert!(args[10].contains("'; DROP TABLE"));
}
/// Verify PTY is available on Unix platforms.
#[cfg(unix)]
#[tokio::test]
async fn pty_opens_successfully() {
let result = pty_process::open();
assert!(result.is_ok(), "PTY allocation should succeed on Unix");
}
}
+148 -4
View File
@@ -391,6 +391,7 @@ Report when the job is complete or if you encounter issues you cannot resolve."#
worker: self,
rx: tokio::sync::Mutex::new(rx),
consecutive_rate_limits: std::sync::atomic::AtomicUsize::new(0),
has_text_response: std::sync::atomic::AtomicBool::new(false),
};
let config = AgenticLoopConfig {
@@ -1101,6 +1102,15 @@ fn store_fallback_in_metadata(
}
/// Job delegate: implements `LoopDelegate` for the background job context.
/// Whether an LLM error represents a completion-eligible empty response.
///
/// Only `EmptyResponse` (provider returned no choices/content) qualifies.
/// Infrastructure errors (`AuthFailed`, `Http`, `Io`, etc.) never qualify —
/// they must propagate even if prior text output was produced.
fn is_completion_eligible_error(error: &crate::error::LlmError) -> bool {
matches!(error, crate::error::LlmError::EmptyResponse { .. })
}
///
/// Handles: signal channel (stop/ping/user messages), cancellation checks,
/// rate-limit retry, parallel tool execution, DB persistence, SSE broadcasting.
@@ -1109,6 +1119,10 @@ struct JobDelegate<'a> {
rx: tokio::sync::Mutex<&'a mut mpsc::Receiver<WorkerMessage>>,
/// Tracks consecutive rate-limit errors to fail fast instead of burning iterations.
consecutive_rate_limits: std::sync::atomic::AtomicUsize,
/// Whether a substantive (non-empty) text response has been produced.
/// When true, an empty follow-up response is treated as job completion
/// rather than a retry signal (prevents spurious failures in routines).
has_text_response: std::sync::atomic::AtomicBool,
}
impl<'a> JobDelegate<'a> {
@@ -1161,6 +1175,53 @@ impl<'a> JobDelegate<'a> {
finish_reason: crate::llm::FinishReason::Stop,
})
}
/// Mark the job as completed, logging a warning on failure.
async fn mark_completed_or_warn(&self, context: &str) {
if let Err(e) = self.worker.mark_completed().await {
tracing::warn!(
job_id = %self.worker.job_id,
error = %e,
"Failed to mark job completed ({context})"
);
}
}
/// If a substantive text response was already produced and the error
/// indicates the LLM simply returned nothing, treat it as successful
/// completion rather than a fatal failure.
///
/// Only swallows `EmptyResponse` — infrastructure errors (`AuthFailed`,
/// `ContextLengthExceeded`, `Http`, `Io`, etc.) always propagate.
///
/// Returns `Some(empty RespondOutput)` when the error should be swallowed,
/// `None` when it should propagate normally.
async fn try_complete_on_error(
&self,
context: &str,
error: &crate::error::LlmError,
) -> Option<crate::llm::RespondOutput> {
if !is_completion_eligible_error(error) {
return None;
}
if !self
.has_text_response
.load(std::sync::atomic::Ordering::Relaxed)
{
return None;
}
tracing::info!(
job_id = %self.worker.job_id,
error = %error,
"{context} empty response after text output — treating as completion"
);
self.mark_completed_or_warn(context).await;
Some(crate::llm::RespondOutput {
result: RespondResult::Text(String::new()),
usage: crate::llm::TokenUsage::default(),
finish_reason: crate::llm::FinishReason::Stop,
})
}
}
#[async_trait]
@@ -1291,7 +1352,12 @@ impl<'a> LoopDelegate for JobDelegate<'a> {
Err(crate::error::LlmError::RateLimited { retry_after, .. }) => {
return self.handle_rate_limit(retry_after, "tool selection").await;
}
Err(e) => return Err(e.into()),
Err(e) => {
if let Some(output) = self.try_complete_on_error("select_tools", &e).await {
return Ok(output);
}
return Err(e.into());
}
};
// Fall back to respond_with_tools
@@ -1321,7 +1387,12 @@ impl<'a> LoopDelegate for JobDelegate<'a> {
self.handle_rate_limit(retry_after, "respond_with_tools")
.await
}
Err(e) => Err(e.into()),
Err(e) => {
if let Some(output) = self.try_complete_on_error("respond_with_tools", &e).await {
return Ok(output);
}
Err(e.into())
}
}
}
@@ -1330,9 +1401,22 @@ impl<'a> LoopDelegate for JobDelegate<'a> {
text: &str,
reason_ctx: &mut ReasoningContext,
) -> TextAction {
// Empty text from rate-limit backoff retry — skip processing and let the
// loop proceed to the next iteration which will re-call the LLM.
// Empty text after a substantive response means the LLM has finished.
// Treat as successful completion rather than continuing the loop (which
// would produce "Response contained no message or tool call (empty)").
if text.is_empty() {
if self
.has_text_response
.load(std::sync::atomic::Ordering::Relaxed)
{
tracing::debug!(
job_id = %self.worker.job_id,
"Empty response after text output — treating as completion"
);
self.mark_completed_or_warn("empty text response").await;
return TextAction::Return(LoopOutcome::Response(String::new()));
}
// No prior text response — this is likely a rate-limit backoff retry.
return TextAction::Continue;
}
@@ -1348,6 +1432,10 @@ impl<'a> LoopDelegate for JobDelegate<'a> {
return TextAction::Return(LoopOutcome::Response(text.to_string()));
}
// Track that a substantive response has been produced.
self.has_text_response
.store(true, std::sync::atomic::Ordering::Relaxed);
// Add assistant response to context
reason_ctx.messages.push(ChatMessage::assistant(text));
@@ -2285,4 +2373,60 @@ mod tests {
assert_eq!(telegram[0].0, "owner-scope");
assert_eq!(telegram[0].1.content, "hello from routine");
}
/// Regression test: only `EmptyResponse` errors are eligible for
/// completion-swallowing. Infrastructure errors must always propagate.
#[test]
fn is_completion_eligible_only_matches_empty_response() {
use crate::error::LlmError;
// EmptyResponse is eligible
assert!(super::is_completion_eligible_error(
&LlmError::EmptyResponse {
provider: "test".to_string(),
}
));
// All other variants are NOT eligible
assert!(!super::is_completion_eligible_error(
&LlmError::InvalidResponse {
provider: "test".to_string(),
reason: "parse error".to_string(),
}
));
assert!(!super::is_completion_eligible_error(
&LlmError::AuthFailed {
provider: "test".to_string(),
}
));
assert!(!super::is_completion_eligible_error(
&LlmError::ContextLengthExceeded {
used: 100_000,
limit: 50_000,
}
));
assert!(!super::is_completion_eligible_error(
&LlmError::ModelNotAvailable {
provider: "test".to_string(),
model: "gpt-4".to_string(),
}
));
assert!(!super::is_completion_eligible_error(
&LlmError::RequestFailed {
provider: "test".to_string(),
reason: "timeout".to_string(),
}
));
assert!(!super::is_completion_eligible_error(
&LlmError::SessionExpired {
provider: "test".to_string(),
}
));
assert!(!super::is_completion_eligible_error(
&LlmError::SessionRenewalFailed {
provider: "test".to_string(),
reason: "timeout".to_string(),
}
));
}
}
+248
View File
@@ -2,6 +2,7 @@
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use uuid::Uuid;
/// Well-known document paths.
@@ -37,6 +38,139 @@ 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
@@ -360,6 +494,120 @@ 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![
+194 -258
View File
@@ -1,8 +1,10 @@
//! Memory hygiene: automatic cleanup of stale workspace documents.
//!
//! 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.
//! 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.
//!
//! A global [`AtomicBool`] guard prevents concurrent hygiene passes, which
//! avoids TOCTOU races on the state file and Windows file-locking errors
@@ -10,18 +12,16 @@
//! 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. List daily/ documents
//! │ 4. Delete those older than daily_retention
//! │ 5. List conversations/ documents
//! │ 6. Delete those older than conversation_ret │
//! │ 7. 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. Discover .config docs with hygiene.enabled
//! │ 4. For each: cleanup_directory(parent, retention)
//! │ 5. Log summary
//! └──────────────────────────────────────────────────┘
//! ```
use std::path::PathBuf;
@@ -31,46 +31,22 @@ use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::bootstrap::ironclaw_base_dir;
use crate::workspace::Workspace;
use crate::workspace::{DocumentMetadata, Workspace, is_config_path};
/// 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,
/// 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,
/// 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,
/// Minimum hours between hygiene passes.
pub cadence_hours: u32,
/// Directory to store state file (default: `~/.ironclaw`).
@@ -81,8 +57,7 @@ impl Default for HygieneConfig {
fn default() -> Self {
Self {
enabled: true,
daily_retention_days: 30,
conversation_retention_days: 7,
version_keep_count: 50,
cadence_hours: 12,
state_dir: ironclaw_base_dir(),
}
@@ -98,10 +73,10 @@ struct HygieneState {
/// Summary of what a hygiene pass cleaned up.
#[derive(Debug, Default)]
pub struct HygieneReport {
/// Number of daily log documents deleted.
pub daily_logs_deleted: u32,
/// Number of conversation documents deleted.
pub conversation_docs_deleted: u32,
/// 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,
/// Whether the run was skipped (cadence not yet elapsed).
pub skipped: bool,
}
@@ -109,7 +84,7 @@ pub struct HygieneReport {
impl HygieneReport {
/// True if any cleanup work was done.
pub fn had_work(&self) -> bool {
self.daily_logs_deleted > 0 || self.conversation_docs_deleted > 0
self.directories_cleaned.iter().any(|(_, n)| *n > 0) || self.versions_pruned > 0
}
}
@@ -168,30 +143,51 @@ 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!(
daily_retention_days = config.daily_retention_days,
conversation_retention_days = config.conversation_retention_days,
"memory hygiene: starting cleanup pass"
);
tracing::info!("memory hygiene: starting cleanup pass");
let mut report = HygieneReport::default();
// 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}"),
}
// 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 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}"),
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}");
}
}
}
if report.had_work() {
tracing::info!(
daily_logs_deleted = report.daily_logs_deleted,
conversation_docs_deleted = report.conversation_docs_deleted,
directories_cleaned = ?report.directories_cleaned,
versions_pruned = report.versions_pruned,
"memory hygiene: cleanup complete"
);
} else {
@@ -210,88 +206,41 @@ impl Drop for RunningGuard {
}
}
/// Delete daily log documents older than `retention_days`.
async fn cleanup_daily_logs(
/// 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(
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("daily/").await?;
let entries = workspace.list(directory).await?;
let mut deleted = 0u32;
for entry in entries {
if entry.is_directory {
continue;
}
// Never delete identity documents
if is_identity_path(&entry.path) {
if is_config_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("daily/") {
let path = if entry.path.starts_with(directory) {
entry.path.clone()
} else {
format!("daily/{}", entry.path)
format!("{}{}", directory, 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 old daily log");
tracing::debug!(path, "memory hygiene: deleted stale document");
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)
}
@@ -349,8 +298,7 @@ mod tests {
fn default_config_is_reasonable() {
let cfg = HygieneConfig::default();
assert!(cfg.enabled);
assert_eq!(cfg.daily_retention_days, 30);
assert_eq!(cfg.conversation_retention_days, 7);
assert_eq!(cfg.version_keep_count, 50);
assert_eq!(cfg.cadence_hours, 12);
}
@@ -362,84 +310,33 @@ mod tests {
}
#[test]
fn report_had_work_when_deleted() {
fn report_had_work_when_directories_cleaned() {
let report = HygieneReport {
daily_logs_deleted: 3,
conversation_docs_deleted: 0,
directories_cleaned: vec![("daily/".to_string(), 3)],
versions_pruned: 0,
skipped: false,
};
assert!(report.had_work());
}
#[test]
fn report_had_work_when_conversation_deleted() {
fn report_had_work_when_versions_pruned() {
let report = HygieneReport {
daily_logs_deleted: 0,
conversation_docs_deleted: 2,
directories_cleaned: vec![],
versions_pruned: 5,
skipped: false,
};
assert!(report.had_work());
}
#[test]
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");
}
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());
}
#[test]
@@ -552,61 +449,105 @@ mod tests {
Arc::new(Workspace::new_with_db("default", db.clone()))
}
#[tokio::test]
async fn cleanup_daily_logs_preserves_identity_documents() {
let (db, _tmp) = create_test_db().await;
let ws = create_workspace(&db);
// Write several regular documents (non-identity)
ws.write("daily/2024-01-15.md", "Old log")
/// 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 log 1");
ws.write("daily/2024-01-20.md", "Another log")
.expect("write .config");
// Read back to get the document ID
let doc = workspace
.read(&config_path)
.await
.expect("write log 2");
// Write an identity document
ws.write("MEMORY.md", "Long-term curated memory")
.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("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)
// This tests that even with aggressive cleanup, identity docs survive
let deleted = cleanup_daily_logs(&ws, 0)
.await
.expect("cleanup_daily_logs");
// Should have deleted some documents (the daily logs)
assert!(deleted > 0, "should have deleted old daily documents");
// Verify identity doc still exists
let identity = db
.get_document_by_path("default", None, "MEMORY.md")
.await
.expect("get identity doc");
assert_eq!(identity.path, "MEMORY.md");
assert_eq!(identity.content, "Long-term curated memory");
.expect("set metadata");
}
#[tokio::test]
async fn cleanup_conversation_docs_handles_empty_directory() {
async fn cleanup_directory_skips_config_files() {
let (db, _tmp) = create_test_db().await;
let ws = create_workspace(&db);
// Run cleanup on an empty directory (conversations/ doesn't exist)
let deleted = cleanup_conversation_docs(&ws, 7)
// Write documents including a .config
ws.write("daily/2024-01-15.md", "Old log")
.await
.expect("cleanup_conversation_docs");
.expect("write log");
ws.write("daily/.config", "").await.expect("write config");
// Run cleanup with 0-day retention (deletes everything old)
let deleted = cleanup_directory(&ws, "daily/", 0)
.await
.expect("cleanup_directory");
// Should have deleted the log but not the .config
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")
.await
.expect("get .config doc");
assert_eq!(config_doc.path, "daily/.config");
}
#[tokio::test]
async fn cleanup_directory_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)
.await
.expect("cleanup_directory");
// 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;
@@ -614,8 +555,7 @@ mod tests {
let config = HygieneConfig {
enabled: true,
daily_retention_days: 30,
conversation_retention_days: 7,
version_keep_count: 50,
cadence_hours: 12,
state_dir: _tmp.path().to_path_buf(),
};
@@ -627,13 +567,6 @@ 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]
@@ -641,6 +574,10 @@ 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
@@ -652,35 +589,34 @@ mod tests {
.await
.expect("write doc 3");
// 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)
// 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)
.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");
// Create a HygieneReport and verify aggregation works
// Verify HygieneReport aggregation
let report = HygieneReport {
daily_logs_deleted: deleted_daily,
conversation_docs_deleted: deleted_conv,
directories_cleaned: vec![
("daily/".to_string(), deleted_daily),
("conversations/".to_string(), deleted_conv),
],
versions_pruned: 0,
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 combines both counts
// Verify had_work() correctly checks directory counts
let no_work = HygieneReport {
daily_logs_deleted: 0,
conversation_docs_deleted: 0,
directories_cleaned: vec![],
versions_pruned: 0,
skipped: false,
};
assert!(!no_work.had_work(), "empty report should indicate no work");
+351 -2
View File
@@ -53,8 +53,9 @@ mod search;
pub use chunker::{ChunkConfig, chunk_document};
pub use document::{
IDENTITY_PATHS, MemoryChunk, MemoryDocument, WorkspaceEntry, is_identity_path,
merge_workspace_entries, paths,
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,
};
pub use embedding_cache::{CachedEmbeddingProvider, EmbeddingCacheConfig};
pub use embeddings::{
@@ -366,6 +367,101 @@ 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.
@@ -696,10 +792,199 @@ 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
@@ -715,6 +1000,12 @@ 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?;
@@ -754,6 +1045,11 @@ 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(())
@@ -1585,6 +1881,14 @@ 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());
@@ -1670,6 +1974,51 @@ 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
+191 -1
View File
@@ -11,7 +11,9 @@ use uuid::Uuid;
use crate::error::WorkspaceError;
use crate::workspace::document::{MemoryChunk, MemoryDocument, WorkspaceEntry};
use crate::workspace::document::{
DocumentVersion, MemoryChunk, MemoryDocument, VersionSummary, WorkspaceEntry,
};
use crate::workspace::search::{RankedResult, SearchConfig, SearchResult, fuse_results};
/// Database repository for workspace operations.
@@ -702,4 +704,192 @@ 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)
}
}
+2
View File
@@ -53,6 +53,7 @@ HEADED=1 pytest scenarios/
| `test_skills.py` | Skills tab UI visibility, ClawHub search (skipped if registry unreachable), install + remove lifecycle |
| `test_sse_reconnect.py` | SSE reconnects after programmatic `eventSource.close()` + `connectSSE()`; history is reloaded after reconnect |
| `test_tool_approval.py` | Approval card appears, buttons disable on approve/deny, parameters toggle via `page.evaluate("showApproval(...)")`; the waiting-approval regression uses a real HTTP tool call |
| `test_extension_uninstall_cleanup.py` | Real install/setup/remove coverage for WASM tools, WASM channels, OAuth-backed shared Google tools, and MCP servers; verifies uninstall deletes stored secrets from the libSQL `secrets` table while preserving shared credentials until the last referencing extension is removed |
| `test_oauth_refresh.py` | Hosted Gmail OAuth regression: complete setup via `/oauth/callback`, expire the stored access token in libSQL, trigger a real `gmail` tool call through `/api/chat/send`, and verify refresh goes through the mock `/oauth/refresh` proxy without forwarding `client_secret` |
## `helpers.py`
@@ -77,6 +78,7 @@ All fixtures are defined in `tests/e2e/conftest.py`. Running `pytest scenarios/`
| `mock_llm_server` | Starts `mock_llm.py --port 0`, reads the assigned port from stdout, waits for `/v1/models` to return 200. Yields the base URL. |
| `ironclaw_server` | Starts the ironclaw binary with a minimal env (see below), waits for `/api/health` (timeout 60s). Yields the base URL. On teardown sends **SIGINT** (not SIGTERM) so the tokio ctrl_c handler triggers a graceful shutdown and LLVM coverage data is flushed. |
| `hosted_oauth_refresh_server` | Starts a second ironclaw instance with a dedicated libSQL DB and `GOOGLE_OAUTH_CLIENT_ID=hosted-google-client-id`, while still pointing `IRONCLAW_OAUTH_EXCHANGE_URL` at `mock_llm.py`. Yields a dict with `base_url`, `db_path`, `gateway_user_id`, and `mock_llm_url` for the hosted refresh regression scenario. |
| `extension_cleanup_server` | Starts an isolated ironclaw instance with its own temp DB/home/WASM dirs, `SECRETS_MASTER_KEY`, and hosted-style OAuth env so uninstall-cleanup scenarios can inspect the `secrets` table without interfering with the shared E2E server state. |
| `browser` | Launches a single Chromium instance (headless by default; set `HEADED=1` for headed). Shared across all tests. |
### Function-scoped fixtures
+109
View File
@@ -443,6 +443,115 @@ async def hosted_oauth_refresh_server(
home_tmpdir.cleanup()
@pytest.fixture(scope="session")
async def extension_cleanup_server(
ironclaw_binary,
mock_llm_server,
):
"""Start an isolated ironclaw instance for uninstall secret cleanup E2E tests."""
reserved = _reserve_loopback_sockets(2)
db_tmpdir = tempfile.TemporaryDirectory(prefix="ironclaw-e2e-cleanup-db-")
home_tmpdir = tempfile.TemporaryDirectory(prefix="ironclaw-e2e-cleanup-home-")
tools_tmpdir = tempfile.TemporaryDirectory(prefix="ironclaw-e2e-cleanup-tools-")
channels_tmpdir = tempfile.TemporaryDirectory(prefix="ironclaw-e2e-cleanup-channels-")
try:
gateway_port = reserved[0].getsockname()[1]
http_port = reserved[1].getsockname()[1]
for sock in reserved:
if sock.fileno() != -1:
sock.close()
db_path = os.path.join(db_tmpdir.name, "extension-cleanup.db")
home_dir = home_tmpdir.name
env = {
"PATH": os.environ.get("PATH", "/usr/bin:/bin"),
"HOME": home_dir,
"IRONCLAW_BASE_DIR": os.path.join(home_dir, ".ironclaw"),
"RUST_LOG": "ironclaw=info",
"RUST_BACKTRACE": "1",
"IRONCLAW_OWNER_ID": OWNER_SCOPE_ID,
"GATEWAY_ENABLED": "true",
"GATEWAY_HOST": "127.0.0.1",
"GATEWAY_PORT": str(gateway_port),
"GATEWAY_AUTH_TOKEN": AUTH_TOKEN,
"GATEWAY_USER_ID": OWNER_SCOPE_ID,
"HTTP_HOST": "127.0.0.1",
"HTTP_PORT": str(http_port),
"HTTP_WEBHOOK_SECRET": HTTP_WEBHOOK_SECRET,
"CLI_ENABLED": "false",
"LLM_BACKEND": "openai_compatible",
"LLM_BASE_URL": mock_llm_server,
"LLM_MODEL": "mock-model",
"DATABASE_BACKEND": "libsql",
"LIBSQL_PATH": db_path,
"SECRETS_MASTER_KEY": "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef",
"SANDBOX_ENABLED": "false",
"SKILLS_ENABLED": "true",
"ROUTINES_ENABLED": "true",
"HEARTBEAT_ENABLED": "false",
"EMBEDDING_ENABLED": "false",
"WASM_ENABLED": "true",
"WASM_TOOLS_DIR": tools_tmpdir.name,
"WASM_CHANNELS_DIR": channels_tmpdir.name,
"ONBOARD_COMPLETED": "true",
"IRONCLAW_OAUTH_CALLBACK_URL": "https://oauth.test.example/oauth/callback",
"IRONCLAW_OAUTH_EXCHANGE_URL": mock_llm_server,
"GOOGLE_OAUTH_CLIENT_ID": "hosted-google-client-id",
}
_forward_coverage_env(env)
proc = await asyncio.create_subprocess_exec(
ironclaw_binary, "--no-onboard",
stdin=asyncio.subprocess.DEVNULL,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
env=env,
)
startup_kill_attempted = False
base_url = f"http://127.0.0.1:{gateway_port}"
try:
await wait_for_ready(f"{base_url}/api/health", timeout=60)
yield {
"base_url": base_url,
"db_path": db_path,
"gateway_user_id": OWNER_SCOPE_ID,
"mock_llm_url": mock_llm_server,
}
except TimeoutError:
if proc.returncode is None:
startup_kill_attempted = True
await _stop_process(proc, timeout=2)
returncode = proc.returncode
stderr_bytes = b""
if proc.stderr:
try:
stderr_bytes = await asyncio.wait_for(proc.stderr.read(8192), timeout=2)
except asyncio.TimeoutError:
pass
stderr_text = stderr_bytes.decode("utf-8", errors="replace")
pytest.fail(
f"extension cleanup server failed to start on port {gateway_port} "
f"(returncode={returncode}).\nstderr:\n{stderr_text}"
)
finally:
if proc.returncode is None:
if startup_kill_attempted:
await _stop_process(proc, timeout=2)
else:
await _stop_process(proc, sig=signal.SIGINT, timeout=10)
if proc.returncode is None:
await _stop_process(proc, timeout=2)
finally:
for sock in reserved:
if sock.fileno() != -1:
sock.close()
db_tmpdir.cleanup()
home_tmpdir.cleanup()
tools_tmpdir.cleanup()
channels_tmpdir.cleanup()
@pytest.fixture(scope="session")
async def http_channel_server(ironclaw_server, server_ports):
"""HTTP webhook channel base URL."""
@@ -0,0 +1,266 @@
"""Extension uninstall secret cleanup E2E tests.
Exercises real install/setup/auth/remove flows and verifies the backing
secrets table is cleaned up when extensions are uninstalled.
"""
import sqlite3
from urllib.parse import parse_qs, urlparse
import httpx
from helpers import api_get, api_post
def _extract_state(auth_url: str) -> str:
parsed = urlparse(auth_url)
state = parse_qs(parsed.query).get("state", [None])[0]
assert state, f"auth_url should include state: {auth_url}"
return state
def _secret_exists(db_path: str, user_id: str, name: str) -> bool:
with sqlite3.connect(db_path) as conn:
row = conn.execute(
"SELECT 1 FROM secrets WHERE user_id = ?1 AND name = ?2 LIMIT 1",
(user_id, name),
).fetchone()
return row is not None
def _secret_names(db_path: str, user_id: str) -> set[str]:
with sqlite3.connect(db_path) as conn:
rows = conn.execute(
"SELECT name FROM secrets WHERE user_id = ?1",
(user_id,),
).fetchall()
return {row[0] for row in rows}
async def _get_extension(base_url: str, name: str) -> dict | None:
response = await api_get(base_url, "/api/extensions", timeout=15)
response.raise_for_status()
for extension in response.json().get("extensions", []):
if extension["name"] == name:
return extension
return None
async def _ensure_removed(base_url: str, name: str) -> None:
extension = await _get_extension(base_url, name)
if extension is not None:
response = await api_post(base_url, f"/api/extensions/{name}/remove", timeout=30)
assert response.status_code == 200, response.text
assert response.json().get("success") is True, response.text
async def _install_extension(
base_url: str,
name: str,
*,
kind: str | None = None,
url: str | None = None,
) -> None:
payload = {"name": name}
if kind is not None:
payload["kind"] = kind
if url is not None:
payload["url"] = url
response = await api_post(
base_url,
"/api/extensions/install",
json=payload,
timeout=180,
)
assert response.status_code == 200, response.text
assert response.json().get("success") is True, response.text
async def test_remove_wasm_tool_deletes_unique_secret(extension_cleanup_server):
server = extension_cleanup_server["base_url"]
db_path = extension_cleanup_server["db_path"]
user_id = extension_cleanup_server["gateway_user_id"]
await _ensure_removed(server, "web-search")
await _install_extension(server, "web-search")
setup_response = await api_post(
server,
"/api/extensions/web-search/setup",
json={"secrets": {"brave_api_key": "cleanup-test-key"}},
timeout=30,
)
assert setup_response.status_code == 200, setup_response.text
assert setup_response.json().get("success") is True, setup_response.text
assert _secret_exists(db_path, user_id, "brave_api_key")
remove_response = await api_post(
server,
"/api/extensions/web-search/remove",
timeout=30,
)
assert remove_response.status_code == 200, remove_response.text
assert remove_response.json().get("success") is True, remove_response.text
assert not _secret_exists(db_path, user_id, "brave_api_key")
async def test_remove_wasm_channel_deletes_setup_secrets(extension_cleanup_server):
server = extension_cleanup_server["base_url"]
db_path = extension_cleanup_server["db_path"]
user_id = extension_cleanup_server["gateway_user_id"]
await _ensure_removed(server, "discord")
await _install_extension(server, "discord", kind="wasm_channel")
setup_response = await api_post(
server,
"/api/extensions/discord/setup",
json={
"secrets": {
"discord_bot_token": "cleanup-discord-bot-token",
"discord_public_key": "cleanup-discord-public-key",
}
},
timeout=30,
)
assert setup_response.status_code == 200, setup_response.text
assert setup_response.json().get("success") is True, setup_response.text
assert _secret_exists(db_path, user_id, "discord_bot_token")
assert _secret_exists(db_path, user_id, "discord_public_key")
remove_response = await api_post(
server,
"/api/extensions/discord/remove",
timeout=30,
)
assert remove_response.status_code == 200, remove_response.text
assert remove_response.json().get("success") is True, remove_response.text
assert not _secret_exists(db_path, user_id, "discord_bot_token")
assert not _secret_exists(db_path, user_id, "discord_public_key")
async def test_remove_shared_google_oauth_secrets_after_last_tool(extension_cleanup_server):
server = extension_cleanup_server["base_url"]
db_path = extension_cleanup_server["db_path"]
user_id = extension_cleanup_server["gateway_user_id"]
await _ensure_removed(server, "gmail")
await _ensure_removed(server, "google-drive")
await _install_extension(server, "gmail")
await _install_extension(server, "google-drive")
setup_response = await api_post(
server,
"/api/extensions/gmail/setup",
json={"secrets": {}},
timeout=30,
)
assert setup_response.status_code == 200, setup_response.text
auth_url = setup_response.json().get("auth_url")
assert auth_url, setup_response.text
async with httpx.AsyncClient() as client:
callback_response = await client.get(
f"{server}/oauth/callback",
params={"code": "mock_auth_code", "state": _extract_state(auth_url)},
timeout=30,
follow_redirects=True,
)
assert callback_response.status_code == 200, callback_response.text[:400]
shared_secrets = [
"google_oauth_token",
"google_oauth_token_refresh_token",
"google_oauth_token_scopes",
]
for secret_name in shared_secrets:
assert _secret_exists(db_path, user_id, secret_name), f"expected {secret_name} to exist"
gmail_remove_response = await api_post(
server,
"/api/extensions/gmail/remove",
timeout=30,
)
assert gmail_remove_response.status_code == 200, gmail_remove_response.text
assert gmail_remove_response.json().get("success") is True, gmail_remove_response.text
for secret_name in shared_secrets:
assert _secret_exists(db_path, user_id, secret_name), (
f"{secret_name} should remain while google-drive is still installed"
)
drive_remove_response = await api_post(
server,
"/api/extensions/google-drive/remove",
timeout=30,
)
assert drive_remove_response.status_code == 200, drive_remove_response.text
assert drive_remove_response.json().get("success") is True, drive_remove_response.text
for secret_name in shared_secrets:
assert not _secret_exists(db_path, user_id, secret_name), (
f"{secret_name} should be deleted after the last Google tool is removed"
)
async def test_remove_mcp_server_deletes_stored_secrets(extension_cleanup_server):
server = extension_cleanup_server["base_url"]
db_path = extension_cleanup_server["db_path"]
user_id = extension_cleanup_server["gateway_user_id"]
mcp_url = f"{extension_cleanup_server['mock_llm_url']}/mcp"
await _ensure_removed(server, "mock-mcp")
await _install_extension(server, "mock-mcp", kind="mcp_server", url=mcp_url)
setup_response = await api_post(
server,
"/api/extensions/mock-mcp/setup",
json={"secrets": {}},
timeout=30,
)
assert setup_response.status_code == 200, setup_response.text
auth_url = setup_response.json().get("auth_url")
if auth_url is None:
activate_response = await api_post(
server,
"/api/extensions/mock-mcp/activate",
timeout=30,
)
assert activate_response.status_code == 200, activate_response.text
auth_url = activate_response.json().get("auth_url")
assert auth_url, "mock-mcp should require OAuth in E2E"
async with httpx.AsyncClient() as client:
callback_response = await client.get(
f"{server}/oauth/callback",
params={"code": "mock_mcp_code", "state": _extract_state(auth_url)},
timeout=30,
follow_redirects=True,
)
assert callback_response.status_code == 200, callback_response.text[:400]
expected_mcp_secrets = [
"mcp_mock-mcp_access_token",
"mcp_mock-mcp_client_id",
]
stored_secret_names = _secret_names(db_path, user_id)
for secret_name in expected_mcp_secrets:
assert secret_name in stored_secret_names, (
f"expected {secret_name} to exist; stored secrets were {sorted(stored_secret_names)}"
)
remove_response = await api_post(
server,
"/api/extensions/mock-mcp/remove",
timeout=30,
)
assert remove_response.status_code == 200, remove_response.text
assert remove_response.json().get("success") is True, remove_response.text
remaining_secret_names = _secret_names(db_path, user_id)
assert not any(name.startswith("mcp_mock-mcp_") for name in remaining_secret_names), (
f"mock-mcp secrets should be deleted on remove; remaining secrets were "
f"{sorted(remaining_secret_names)}"
)
+2 -4
View File
@@ -955,8 +955,7 @@ mod tests {
let hygiene_config = HygieneConfig {
enabled: false,
daily_retention_days: 30,
conversation_retention_days: 7,
version_keep_count: 50,
cadence_hours: 24,
state_dir: _tmp.path().to_path_buf(),
};
@@ -1002,8 +1001,7 @@ mod tests {
let hygiene_config = HygieneConfig {
enabled: false,
daily_retention_days: 30,
conversation_retention_days: 7,
version_keep_count: 50,
cadence_hours: 24,
state_dir: _tmp.path().to_path_buf(),
};