Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
255a027
feat: unified thread model for web gateway
ilblackdragon Mar 6, 2026
1b030b2
fix: address PR review — TOCTOU races, stale comment, debounce, broad…
ilblackdragon Mar 6, 2026
e5851c6
fix: sort in-memory thread fallback by updated_at descending
ilblackdragon Mar 6, 2026
7afaa59
fix: retry libsql connect() on transient "unable to open database file"
ilblackdragon Mar 6, 2026
092ca85
fix: use ON CONFLICT with index expressions instead of named constraints
ilblackdragon Mar 6, 2026
b5c79df
style: fix rustfmt chain collapse in heartbeat.rs
ilblackdragon Mar 6, 2026
dae234b
fix: skip broadcast when thread_id is None instead of sending empty
ilblackdragon Mar 6, 2026
87e07b9
test: add libsql routine/heartbeat conversation idempotency tests
ilblackdragon Mar 6, 2026
2b9eb75
feat: show "New chat" title for empty threads
ilblackdragon Mar 7, 2026
431a9de
fix: thread sorting, routine isolation, and UI polish
ilblackdragon Mar 7, 2026
d0abde1
Merge remote-tracking branch 'origin/main' into feat/unified-thread-m…
ilblackdragon Mar 7, 2026
bda2b2b
fix: routine history display, orphaned tool_results, duplicate system…
ilblackdragon Mar 7, 2026
78714ac
Merge remote-tracking branch 'origin/main' into feat/unified-thread-m…
ilblackdragon Mar 7, 2026
e16d63a
merge: resolve conflicts with origin/main (tool intent nudges)
ilblackdragon Mar 7, 2026
32afa42
fix: address PR #607 review — RwLock held across await, missing owner…
ilblackdragon Mar 7, 2026
7a6a26a
Merge remote-tracking branch 'origin/main' into feat/unified-thread-m…
ilblackdragon Mar 7, 2026
48ccfda
chore: gitignore trace_*.json files and remove stale traces
ilblackdragon Mar 7, 2026
a282893
chore: remove trace JSON files from repo
ilblackdragon Mar 7, 2026
2bd9a39
Merge remote-tracking branch 'origin/main' into feat/unified-thread-m…
ilblackdragon Mar 7, 2026
db760a1
fix: proper HTTP status codes for routine errors, read-only input gua…
ilblackdragon Mar 7, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -22,3 +22,4 @@ bench-results/
# WASM build artifacts (loaded from disk, not bundled)
*.wasm

trace_*.json
13 changes: 13 additions & 0 deletions migrations/V11__conversation_unique_indexes.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
-- Partial unique indexes to prevent duplicate singleton conversations.
-- These guard against TOCTOU races in get_or_create_routine_conversation
-- and get_or_create_heartbeat_conversation.

-- One routine conversation per user per routine_id.
CREATE UNIQUE INDEX IF NOT EXISTS uq_conv_routine
ON conversations (user_id, (metadata->>'routine_id'))
WHERE metadata->>'routine_id' IS NOT NULL;

-- One heartbeat conversation per user.
CREATE UNIQUE INDEX IF NOT EXISTS uq_conv_heartbeat
ON conversations (user_id)
WHERE metadata->>'thread_type' = 'heartbeat';

Copilot AI Mar 6, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This migration creates uq_conv_routine/uq_conv_heartbeat as unique indexes, but the Postgres Store code uses ON CONFLICT ON CONSTRAINT uq_conv_*, which requires a named unique constraint. Either update the Store SQL to use index inference, or add ALTER TABLE conversations ADD CONSTRAINT uq_conv_* UNIQUE USING INDEX uq_conv_*; here so the constraint exists.

Suggested change
WHERE metadata->>'thread_type' = 'heartbeat';
WHERE metadata->>'thread_type' = 'heartbeat';
ALTER TABLE conversations
ADD CONSTRAINT uq_conv_routine UNIQUE USING INDEX uq_conv_routine;
ALTER TABLE conversations
ADD CONSTRAINT uq_conv_heartbeat UNIQUE USING INDEX uq_conv_heartbeat;

Copilot uses AI. Check for mistakes.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 092ca85 — the migration correctly creates unique indexes; the Store SQL was the problem. Changed to use ON CONFLICT (column_expr) WHERE condition syntax which works with both indexes and constraints.

25 changes: 24 additions & 1 deletion src/agent/agent_loop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,9 @@ pub struct Agent {
pub(super) heartbeat_config: Option<HeartbeatConfig>,
pub(super) hygiene_config: Option<crate::config::HygieneConfig>,
pub(super) routine_config: Option<RoutineConfig>,
/// Optional slot to expose the routine engine to the gateway for manual triggering.
pub(super) routine_engine_slot:
Option<Arc<tokio::sync::RwLock<Option<Arc<crate::agent::routine_engine::RoutineEngine>>>>>,
}

impl Agent {
Expand Down Expand Up @@ -144,9 +147,18 @@ impl Agent {
heartbeat_config,
hygiene_config,
routine_config,
routine_engine_slot: None,
}
}

/// Set the routine engine slot for exposing the engine to the gateway.
pub fn set_routine_engine_slot(
&mut self,
slot: Arc<tokio::sync::RwLock<Option<Arc<crate::agent::routine_engine::RoutineEngine>>>>,
) {
self.routine_engine_slot = Some(slot);
}

// Convenience accessors

/// Get the scheduler (for external wiring, e.g. CreateJobTool).
Expand Down Expand Up @@ -338,8 +350,13 @@ impl Agent {
let heartbeat_handle = if let Some(ref hb_config) = self.heartbeat_config {
if hb_config.enabled {
if let Some(workspace) = self.workspace() {
let config = AgentHeartbeatConfig::default()
let mut config = AgentHeartbeatConfig::default()
.with_interval(std::time::Duration::from_secs(hb_config.interval_secs));
if let (Some(user), Some(channel)) =
(&hb_config.notify_user, &hb_config.notify_channel)
{
config = config.with_notify(user, channel);
}

// Set up notification channel
let (notify_tx, mut notify_rx) =
Expand Down Expand Up @@ -392,6 +409,7 @@ impl Agent {
self.cheap_llm().clone(),
self.safety().clone(),
Some(notify_tx),
self.store().map(Arc::clone),
))
} else {
tracing::warn!("Heartbeat enabled but no workspace available");
Expand Down Expand Up @@ -482,6 +500,11 @@ impl Agent {
// SAFETY: self is consumed by run(), we can smuggle the engine in
// via a local to use in the message loop below.

// Expose engine to gateway for manual triggering
if let Some(ref slot) = self.routine_engine_slot {
*slot.write().await = Some(Arc::clone(&engine));
}

tracing::info!(
"Routines enabled: cron ticker every {}s, max {} concurrent",
rt_config.cron_check_interval_secs,
Expand Down
56 changes: 55 additions & 1 deletion src/agent/heartbeat.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ use std::time::Duration;
use tokio::sync::mpsc;

use crate::channels::OutgoingResponse;
use crate::db::Database;
use crate::llm::{ChatMessage, CompletionRequest, LlmProvider, Reasoning};
use crate::safety::SafetyLayer;
use crate::workspace::Workspace;
Expand Down Expand Up @@ -103,6 +104,7 @@ pub struct HeartbeatRunner {
llm: Arc<dyn LlmProvider>,
safety: Arc<SafetyLayer>,
response_tx: Option<mpsc::Sender<OutgoingResponse>>,
store: Option<Arc<dyn Database>>,
consecutive_failures: u32,
}

Expand All @@ -122,6 +124,7 @@ impl HeartbeatRunner {
llm,
safety,
response_tx: None,
store: None,
consecutive_failures: 0,
}
}
Expand All @@ -132,6 +135,12 @@ impl HeartbeatRunner {
self
}

/// Set the database store for persistent heartbeat conversations.
pub fn with_store(mut self, store: Arc<dyn Database>) -> Self {
self.store = Some(store);
self
}

/// Run the heartbeat loop.
///
/// This runs forever, checking periodically based on the configured interval.
Expand Down Expand Up @@ -292,9 +301,32 @@ impl HeartbeatRunner {
return;
};

let user_id = self.config.notify_user_id.as_deref().unwrap_or("default");

// Persist to heartbeat conversation and get thread_id
let thread_id = if let Some(ref store) = self.store {
match store.get_or_create_heartbeat_conversation(user_id).await {
Ok(conv_id) => {
if let Err(e) = store
.add_conversation_message(conv_id, "assistant", message)
.await
{
tracing::error!("Failed to persist heartbeat message: {}", e);
}
Some(conv_id.to_string())
}
Err(e) => {
tracing::error!("Failed to get heartbeat conversation: {}", e);
None
}
}
} else {
None
};

let response = OutgoingResponse {
content: format!("🔔 *Heartbeat Alert*\n\n{}", message),
thread_id: None,
thread_id,
Comment on lines +304 to +329

Copilot AI Mar 6, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

When store is None, heartbeat notifications are emitted with thread_id = None. The Gateway channel now skips broadcasting responses without a thread_id, so in “no DB / in-memory only” deployments heartbeat alerts won’t appear in the web UI at all. Consider either (a) not spawning heartbeat when running without a store and the notify target includes the gateway, or (b) ensuring a thread_id is always set for gateway delivery (e.g., route to the assistant conversation when a store exists, and otherwise log a clear warning + disable web delivery).

Copilot uses AI. Check for mistakes.
attachments: Vec::new(),
metadata: serde_json::json!({
"source": "heartbeat",
Expand Down Expand Up @@ -356,11 +388,15 @@ pub fn spawn_heartbeat(
llm: Arc<dyn LlmProvider>,
safety: Arc<SafetyLayer>,
response_tx: Option<mpsc::Sender<OutgoingResponse>>,
store: Option<Arc<dyn Database>>,
) -> tokio::task::JoinHandle<()> {
let mut runner = HeartbeatRunner::new(config, hygiene_config, workspace, llm, safety);
if let Some(tx) = response_tx {
runner = runner.with_response_channel(tx);
}
if let Some(s) = store {
runner = runner.with_store(s);
}

tokio::spawn(async move {
runner.run().await;
Expand Down Expand Up @@ -495,4 +531,22 @@ mod tests {
let content = "<!-- comment -->\nActual task here";
assert!(!is_effectively_empty(content));
}

#[test]
fn test_spawn_heartbeat_accepts_store_param() {
// Regression: spawn_heartbeat must accept an optional Database store
// for persisting heartbeat notifications to a dedicated conversation.
// Compile-time check: the 7th parameter is `Option<Arc<dyn Database>>`.
#[allow(clippy::type_complexity)]
let _fn_ptr: fn(
HeartbeatConfig,
HygieneConfig,
Arc<crate::workspace::Workspace>,
Arc<dyn crate::llm::LlmProvider>,
Arc<crate::safety::SafetyLayer>,
Option<tokio::sync::mpsc::Sender<crate::channels::OutgoingResponse>>,
Option<Arc<dyn crate::db::Database>>,
) -> tokio::task::JoinHandle<()> = spawn_heartbeat;
let _ = _fn_ptr;
}
}
50 changes: 48 additions & 2 deletions src/agent/routine_engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -184,7 +184,11 @@ impl RoutineEngine {
///
/// Bypasses cooldown checks (those only apply to cron/event triggers).
/// Still enforces enabled check and concurrent run limit.
pub async fn fire_manual(&self, routine_id: Uuid) -> Result<Uuid, RoutineError> {
pub async fn fire_manual(
&self,
routine_id: Uuid,
user_id: Option<&str>,
) -> Result<Uuid, RoutineError> {
let routine = self
.store
.get_routine(routine_id)
Expand All @@ -194,6 +198,13 @@ impl RoutineEngine {
})?
.ok_or(RoutineError::NotFound { id: routine_id })?;

// Enforce ownership when a user_id is provided (gateway calls).
if let Some(uid) = user_id
&& routine.user_id != uid
{
return Err(RoutineError::NotAuthorized { id: routine_id });
}

if !routine.enabled {
return Err(RoutineError::Disabled {
name: routine.name.clone(),
Expand Down Expand Up @@ -396,13 +407,47 @@ async fn execute_routine(ctx: EngineContext, routine: Routine, run: RoutineRun)
tracing::error!(routine = %routine.name, "Failed to update runtime state: {}", e);
}

// Persist routine result to its dedicated conversation thread
let thread_id = match ctx
.store
.get_or_create_routine_conversation(routine.id, &routine.name, &routine.user_id)
.await
{
Ok(conv_id) => {
tracing::debug!(
routine = %routine.name,
routine_id = %routine.id,
conversation_id = %conv_id,
"Resolved routine conversation thread"
);
// Record the run result as a conversation message
let msg = match (&summary, status) {
(Some(s), _) => format!("[{}] {}: {}", run.trigger_type, status, s),
(None, _) => format!("[{}] {}", run.trigger_type, status),
};
if let Err(e) = ctx
.store
.add_conversation_message(conv_id, "assistant", &msg)
.await
{
tracing::error!(routine = %routine.name, "Failed to persist routine message: {}", e);
}
Some(conv_id.to_string())
}
Err(e) => {
tracing::error!(routine = %routine.name, "Failed to get routine conversation: {}", e);
None
}
};

// Send notifications based on config
send_notification(
&ctx.notify_tx,
&routine.notify,
&routine.name,
status,
summary.as_deref(),
thread_id.as_deref(),
)
.await;
}
Expand Down Expand Up @@ -611,6 +656,7 @@ async fn send_notification(
routine_name: &str,
status: RunStatus,
summary: Option<&str>,
thread_id: Option<&str>,
) {
let should_notify = match status {
RunStatus::Ok => notify.on_success,
Expand All @@ -637,7 +683,7 @@ async fn send_notification(

let response = OutgoingResponse {
content: message,
thread_id: None,
thread_id: thread_id.map(String::from),
attachments: Vec::new(),
metadata: serde_json::json!({
"source": "routine",
Expand Down
Loading
Loading