feat: unified thread model for web gateway - #607
Conversation
Every piece of activity (user chat, routine run, heartbeat alert, external channel message) now lives in its own thread, properly isolated, with meaningful titles and visual distinction. Key changes: - Add `channel` field to ConversationSummary and ThreadInfo so the gateway can distinguish thread origins (gateway, telegram, routine, heartbeat). - Add `list_conversations_all_channels` to Database trait (both postgres and libsql) so chat_threads_handler shows cross-channel threads. - Routine runs get a persistent conversation per routine via `get_or_create_routine_conversation`; notifications carry thread_id. - Heartbeat gets a persistent conversation via `get_or_create_heartbeat_conversation`; HeartbeatRunner accepts an optional Database store and binds notifications to the thread. - Fix broadcast() in web gateway to propagate response.thread_id instead of hardcoding empty string. - Fix isCurrentThread(null) returning true (the core notification leak bug) — now returns false so events without a thread_id don't leak into the active thread. - Rewrite frontend thread sidebar: meaningful titles with channel-specific fallbacks, relative timestamps instead of turn counts, channel badges for non-gateway threads, unread notification dots, read-only indicator for external channel threads. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Summary of ChangesHello, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request introduces a comprehensive unified thread model for the web gateway, significantly enhancing how conversations are managed and presented. The changes ensure that all types of interactions, whether from internal routines, heartbeat alerts, or external messaging platforms, are properly organized into distinct, persistent threads. This not only prevents notification cross-talk but also provides users with a much clearer and more informative overview of their conversations through an enriched sidebar interface, improving overall user experience and system clarity. Highlights
Changelog
Activity
Using Gemini Code AssistThe full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips. Invoking Gemini You can request assistance from Gemini at any point by creating a comment using either
Customization To customize Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a Limitations & Feedback Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counter productive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for Github and other Google products, sign up here. You can also get AI-powered code generation, chat, as well as code reviews directly in the IDE at no cost with the Gemini Code Assist IDE Extension. Footnotes
|
There was a problem hiding this comment.
Code Review
This pull request introduces a unified thread model, allowing conversations from different channels like Telegram and Signal to be visible in the web gateway. It also ensures that routine and heartbeat notifications are persisted in their own dedicated threads, preventing them from leaking into active conversations. The frontend has been updated to support these changes with features like channel badges, unread indicators, and relative timestamps.
The implementation is solid, but I've identified a couple of race conditions in the database logic for creating routine and heartbeat conversations. These could lead to duplicate conversation records under concurrent load. I've provided detailed comments and suggestions on how to address these to ensure data consistency, aligning with best practices for preventing Time-of-Check-to-Time-of-Use (TOCTOU) race conditions.
| async fn get_or_create_routine_conversation( | ||
| &self, | ||
| routine_id: Uuid, | ||
| routine_name: &str, | ||
| user_id: &str, | ||
| ) -> Result<Uuid, DatabaseError> { | ||
| let conn = self.connect().await?; | ||
| let rid = routine_id.to_string(); | ||
|
|
||
| let mut rows = conn | ||
| .query( | ||
| r#" | ||
| SELECT id FROM conversations | ||
| WHERE user_id = ?1 AND json_extract(metadata, '$.routine_id') = ?2 | ||
| LIMIT 1 | ||
| "#, | ||
| params![user_id, rid], | ||
| ) | ||
| .await | ||
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | ||
|
|
||
| if let Some(row) = rows | ||
| .next() | ||
| .await | ||
| .map_err(|e| DatabaseError::Query(e.to_string()))? | ||
| { | ||
| let id_str: String = row.get(0).unwrap_or_default(); | ||
| return id_str | ||
| .parse() | ||
| .map_err(|_| DatabaseError::Serialization("Invalid UUID".to_string())); | ||
| } | ||
|
|
||
| let id = Uuid::new_v4(); | ||
| let metadata = serde_json::json!({ | ||
| "thread_type": "routine", | ||
| "routine_id": routine_id.to_string(), | ||
| "routine_name": routine_name, | ||
| }); | ||
| conn.execute( | ||
| "INSERT INTO conversations (id, channel, user_id, metadata) VALUES (?1, ?2, ?3, ?4)", | ||
| params![id.to_string(), "routine", user_id, metadata.to_string()], | ||
| ) | ||
| .await | ||
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | ||
| Ok(id) | ||
| } |
There was a problem hiding this comment.
This get_or_create_routine_conversation function has a race condition. If two concurrent routine runs for a new routine execute this function, both could check for an existing conversation, find none, and then both attempt to insert a new one. This would result in duplicate conversations for the same routine, as there doesn't appear to be a unique constraint to prevent this.
The same issue exists in get_or_create_heartbeat_conversation.
To fix this, you should ensure atomicity. A robust way to handle this is to add a unique index on the identifying fields and then use an "upsert" pattern. For SQLite/libSQL, this would involve:
- Adding a unique index to the database schema on the expression that identifies a routine conversation.
CREATE UNIQUE INDEX uq_routine_conversation ON conversations(user_id, json_extract(metadata, '$.routine_id')) WHERE json_extract(metadata, '$.routine_id') IS NOT NULL;
- Change the logic to first attempt an
INSERTwithON CONFLICT DO NOTHING, and thenSELECTthe ID. This ensures that even with concurrent requests, only one row is created.
References
- To prevent a time-of-check-to-time-of-use (TOCTOU) race condition when checking a condition under a read lock before performing a write, re-verify the condition after acquiring the write lock but before executing the write. This is known as the double-checked locking pattern. The suggested database
ON CONFLICTmechanism provides an atomic way to achieve this.
There was a problem hiding this comment.
Fixed in 1b030b2. libsql now uses BEGIN IMMEDIATE transactions to serialize concurrent writers for both get_or_create_routine_conversation and get_or_create_heartbeat_conversation. Also added matching unique indexes (uq_conv_routine, uq_conv_heartbeat) to the libsql schema as a defense-in-depth measure.
| pub async fn get_or_create_routine_conversation( | ||
| &self, | ||
| routine_id: Uuid, | ||
| routine_name: &str, | ||
| user_id: &str, | ||
| ) -> Result<Uuid, DatabaseError> { | ||
| let conn = self.conn().await?; | ||
| let rid = routine_id.to_string(); | ||
|
|
||
| let row = conn | ||
| .query_opt( | ||
| r#" | ||
| SELECT id FROM conversations | ||
| WHERE user_id = $1 AND metadata->>'routine_id' = $2 | ||
| LIMIT 1 | ||
| "#, | ||
| &[&user_id, &rid], | ||
| ) | ||
| .await?; | ||
|
|
||
| if let Some(row) = row { | ||
| return Ok(row.get("id")); | ||
| } | ||
|
|
||
| let id = Uuid::new_v4(); | ||
| let metadata = serde_json::json!({ | ||
| "thread_type": "routine", | ||
| "routine_id": routine_id.to_string(), | ||
| "routine_name": routine_name, | ||
| }); | ||
| conn.execute( | ||
| "INSERT INTO conversations (id, channel, user_id, metadata) VALUES ($1, $2, $3, $4)", | ||
| &[&id, &"routine", &user_id, &metadata], | ||
| ) | ||
| .await?; | ||
|
|
||
| Ok(id) | ||
| } |
There was a problem hiding this comment.
This get_or_create_routine_conversation function has a race condition. If two concurrent routine runs for a new routine execute this function, both could check for an existing conversation, find none, and then both attempt to insert a new one. This would result in duplicate conversations for the same routine.
The same issue exists in get_or_create_heartbeat_conversation.
For PostgreSQL, you can fix this atomically and efficiently by using an INSERT ... ON CONFLICT statement. This requires a unique index on the fields that identify a routine conversation.
- Add a partial unique index to your database schema:
CREATE UNIQUE INDEX uq_routine_conversation ON conversations (user_id, (metadata->>'routine_id')) WHERE (metadata->>'routine_id') IS NOT NULL;
- Then, you can refactor this function to perform an "upsert" followed by a select, which is atomic. This pattern is robust against race conditions and ensures a single conversation is created per routine. Note that the
ON CONFLICTclause needs to match the unique index definition exactly.
pub async fn get_or_create_routine_conversation(
&self,
routine_id: Uuid,
routine_name: &str,
user_id: &str,
) -> Result<Uuid, DatabaseError> {
let conn = self.conn().await?;
let rid = routine_id.to_string();
// Use an ON CONFLICT clause to handle concurrent creation atomically.
// This requires a unique index on (user_id, (metadata->>'routine_id')).
let new_id = Uuid::new_v4();
let metadata = serde_json::json!({
"thread_type": "routine",
"routine_id": routine_id.to_string(),
"routine_name": routine_name,
});
conn.execute(
r#"
INSERT INTO conversations (id, channel, user_id, metadata)
VALUES ($1, 'routine', $2, $3)
ON CONFLICT (user_id, (metadata->>'routine_id')) WHERE (metadata->>'routine_id') IS NOT NULL
DO NOTHING
"#,
&[&new_id, &user_id, &metadata],
)
.await?;
// Now, select the ID. It will either be the one we just inserted
// or the one that already existed.
let row = conn
.query_one(
r#"
SELECT id FROM conversations
WHERE user_id = $1 AND metadata->>'routine_id' = $2
LIMIT 1
"#,
&[&user_id, &rid],
)
.await?;
Ok(row.get("id"))
}References
- To prevent a time-of-check-to-time-of-use (TOCTOU) race condition when checking a condition under a read lock before performing a write, re-verify the condition after acquiring the write lock but before executing the write. This is known as the double-checked locking pattern. The suggested database
ON CONFLICTmechanism provides an atomic way to achieve this.
There was a problem hiding this comment.
Fixed in 1b030b2. Postgres now uses INSERT ON CONFLICT on a new partial unique index (uq_conv_routine) + SELECT-back, making the operation atomic. Same pattern applied to get_or_create_heartbeat_conversation. Migration V11 adds the indexes.
There was a problem hiding this comment.
Pull request overview
This PR unifies the web gateway’s thread model so notifications and cross-channel conversations are routed into persistent per-thread conversations and can be displayed coherently in the web UI sidebar.
Changes:
- Added conversation
channelto thread list DTOs and expanded thread listing to include all channels. - Introduced persistent “routine” and “heartbeat” conversations and attached
thread_idto their notifications. - Updated web UI thread filtering and sidebar rendering (badges, relative timestamps, unread indicators, read-only behavior).
Reviewed changes
Copilot reviewed 13 out of 13 changed files in this pull request and generated 7 comments.
Show a summary per file
| File | Description |
|---|---|
src/history/store.rs |
Extends ConversationSummary with channel and adds store methods for listing all channels + routine/heartbeat conversation creation. |
src/db/mod.rs |
Extends ConversationStore trait with new cross-channel listing and routine/heartbeat helpers. |
src/db/postgres.rs |
Delegates new ConversationStore methods to the existing Store. |
src/db/libsql/conversations.rs |
Implements new conversation listing and routine/heartbeat get-or-create logic for libsql. |
src/agent/routine_engine.rs |
Persists routine run results into a dedicated conversation and forwards thread_id on notifications. |
src/agent/heartbeat.rs |
Adds optional DB store support to persist heartbeat alerts and sets thread_id on outgoing notifications. |
src/agent/agent_loop.rs |
Passes the optional database store into spawn_heartbeat. |
src/channels/web/server.rs |
Uses all-channel conversation listing and includes channel in ThreadInfo. |
src/channels/web/handlers/chat.rs |
Mirrors server handler updates to include channel in ThreadInfo. |
src/channels/web/types.rs |
Adds optional channel field to ThreadInfo and tests for serialization behavior. |
src/channels/web/mod.rs |
Updates gateway broadcast() to propagate response.thread_id. |
src/channels/web/static/app.js |
Fixes thread filtering, adds unread tracking, and rewrites sidebar rendering with channel-aware UI. |
src/channels/web/static/style.css |
Adds styles for channel badges and unread indicators in the thread list. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| let mut rows = conn | ||
| .query( | ||
| r#" | ||
| SELECT id FROM conversations | ||
| WHERE user_id = ?1 AND json_extract(metadata, '$.routine_id') = ?2 | ||
| LIMIT 1 | ||
| "#, | ||
| params![user_id, rid], | ||
| ) | ||
| .await | ||
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | ||
|
|
||
| if let Some(row) = rows | ||
| .next() | ||
| .await | ||
| .map_err(|e| DatabaseError::Query(e.to_string()))? | ||
| { | ||
| let id_str: String = row.get(0).unwrap_or_default(); | ||
| return id_str | ||
| .parse() | ||
| .map_err(|_| DatabaseError::Serialization("Invalid UUID".to_string())); | ||
| } | ||
|
|
||
| let id = Uuid::new_v4(); | ||
| let metadata = serde_json::json!({ | ||
| "thread_type": "routine", | ||
| "routine_id": routine_id.to_string(), | ||
| "routine_name": routine_name, | ||
| }); | ||
| conn.execute( | ||
| "INSERT INTO conversations (id, channel, user_id, metadata) VALUES (?1, ?2, ?3, ?4)", | ||
| params![id.to_string(), "routine", user_id, metadata.to_string()], | ||
| ) | ||
| .await | ||
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | ||
| Ok(id) |
There was a problem hiding this comment.
get_or_create_routine_conversation is implemented as a query-then-insert without any uniqueness constraint. Under concurrent routine executions this can create duplicate routine conversations for the same routine_id. Consider introducing a unique key (schema-level) and using an UPSERT/transactional approach to make this atomic.
There was a problem hiding this comment.
Fixed in 1b030b2 — same issue as Gemini flagged. libsql now uses BEGIN IMMEDIATE transactions; unique indexes added to schema.
| let mut rows = conn | ||
| .query( | ||
| r#" | ||
| SELECT id FROM conversations | ||
| WHERE user_id = ?1 AND json_extract(metadata, '$.thread_type') = 'heartbeat' | ||
| LIMIT 1 | ||
| "#, | ||
| params![user_id], | ||
| ) | ||
| .await | ||
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | ||
|
|
||
| if let Some(row) = rows | ||
| .next() | ||
| .await | ||
| .map_err(|e| DatabaseError::Query(e.to_string()))? | ||
| { | ||
| let id_str: String = row.get(0).unwrap_or_default(); | ||
| return id_str | ||
| .parse() | ||
| .map_err(|_| DatabaseError::Serialization("Invalid UUID".to_string())); | ||
| } | ||
|
|
||
| let id = Uuid::new_v4(); | ||
| let metadata = serde_json::json!({ "thread_type": "heartbeat" }); | ||
| conn.execute( | ||
| "INSERT INTO conversations (id, channel, user_id, metadata) VALUES (?1, ?2, ?3, ?4)", | ||
| params![id.to_string(), "heartbeat", user_id, metadata.to_string()], | ||
| ) | ||
| .await | ||
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | ||
| Ok(id) |
There was a problem hiding this comment.
get_or_create_heartbeat_conversation uses query-then-insert, which can produce multiple heartbeat conversations for the same user under concurrent sends. If heartbeat is meant to be a singleton thread, enforce it with a unique constraint and implement this as an UPSERT (or wrap in a transaction with appropriate locking).
| let mut rows = conn | |
| .query( | |
| r#" | |
| SELECT id FROM conversations | |
| WHERE user_id = ?1 AND json_extract(metadata, '$.thread_type') = 'heartbeat' | |
| LIMIT 1 | |
| "#, | |
| params![user_id], | |
| ) | |
| .await | |
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | |
| if let Some(row) = rows | |
| .next() | |
| .await | |
| .map_err(|e| DatabaseError::Query(e.to_string()))? | |
| { | |
| let id_str: String = row.get(0).unwrap_or_default(); | |
| return id_str | |
| .parse() | |
| .map_err(|_| DatabaseError::Serialization("Invalid UUID".to_string())); | |
| } | |
| let id = Uuid::new_v4(); | |
| let metadata = serde_json::json!({ "thread_type": "heartbeat" }); | |
| conn.execute( | |
| "INSERT INTO conversations (id, channel, user_id, metadata) VALUES (?1, ?2, ?3, ?4)", | |
| params![id.to_string(), "heartbeat", user_id, metadata.to_string()], | |
| ) | |
| .await | |
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | |
| Ok(id) | |
| // Wrap the lookup-and-insert logic in a transaction to avoid races under concurrency. | |
| let tx_result: Result<Uuid, DatabaseError> = async { | |
| // Acquire a write lock so only one concurrent writer can execute this block at a time. | |
| conn.execute("BEGIN IMMEDIATE", params![]) | |
| .await | |
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | |
| let mut rows = conn | |
| .query( | |
| r#" | |
| SELECT id FROM conversations | |
| WHERE user_id = ?1 AND json_extract(metadata, '$.thread_type') = 'heartbeat' | |
| LIMIT 1 | |
| "#, | |
| params![user_id], | |
| ) | |
| .await | |
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | |
| let id: Uuid = if let Some(row) = rows | |
| .next() | |
| .await | |
| .map_err(|e| DatabaseError::Query(e.to_string()))? | |
| { | |
| let id_str: String = row.get(0).unwrap_or_default(); | |
| id_str | |
| .parse() | |
| .map_err(|_| { | |
| DatabaseError::Serialization("Invalid UUID".to_string()) | |
| })? | |
| } else { | |
| let new_id = Uuid::new_v4(); | |
| let metadata = serde_json::json!({ "thread_type": "heartbeat" }); | |
| conn.execute( | |
| "INSERT INTO conversations (id, channel, user_id, metadata) VALUES (?1, ?2, ?3, ?4)", | |
| params![new_id.to_string(), "heartbeat", user_id, metadata.to_string()], | |
| ) | |
| .await | |
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | |
| new_id | |
| }; | |
| conn.execute("COMMIT", params![]) | |
| .await | |
| .map_err(|e| DatabaseError::Query(e.to_string()))?; | |
| Ok(id) | |
| } | |
| .await; | |
| match tx_result { | |
| Ok(id) => Ok(id), | |
| Err(e) => { | |
| // Best-effort rollback if the transaction failed. | |
| let _ = conn.execute("ROLLBACK", params![]).await; | |
| Err(e) | |
| } | |
| } |
There was a problem hiding this comment.
Fixed in 1b030b2. libsql heartbeat get-or-create now uses BEGIN IMMEDIATE transaction. Unique index uq_conv_heartbeat also added to schema.
| // Events without a thread_id (legacy) are always shown. | ||
| function isCurrentThread(threadId) { | ||
| if (!threadId) return true; | ||
| if (!threadId) return false; |
There was a problem hiding this comment.
The comment above isCurrentThread says events without a thread_id (legacy) are always shown, but the new logic returns false when threadId is falsy. Please update the comment (or the logic) so they match; otherwise future changes may reintroduce notification leaking or drop legacy events unexpectedly.
| if (!threadId) return false; | |
| if (!threadId) return true; |
There was a problem hiding this comment.
Fixed in 1b030b2 — updated the comment to match the new behavior: // Events without a thread_id are dropped (must be scoped to a thread).
| if (!isCurrentThread(data.thread_id)) { | ||
| if (data.thread_id) { | ||
| unreadThreads.set(data.thread_id, (unreadThreads.get(data.thread_id) || 0) + 1); | ||
| loadThreads(); | ||
| } |
There was a problem hiding this comment.
On every SSE response event for a non-active thread, the client calls loadThreads(), which triggers a full /api/chat/threads fetch and DOM rebuild. Under frequent off-thread notifications (routines/heartbeat), this can cause a request storm and UI jank. Consider debouncing loadThreads() (e.g., timer-based coalescing) or updating only the unread counter in-place.
There was a problem hiding this comment.
Fixed in 1b030b2 — added a 500ms debounce via debouncedLoadThreads() so rapid SSE events coalesce into a single API call.
| let thread_id = response.thread_id.unwrap_or_default(); | ||
| self.state.sse.broadcast(SseEvent::Response { | ||
| content: response.content, | ||
| thread_id: String::new(), | ||
| thread_id, | ||
| }); |
There was a problem hiding this comment.
broadcast() converts a missing response.thread_id into an empty string. With the updated client-side filtering (isCurrentThread treats falsy thread IDs as not current), this will silently drop any broadcasted responses that don't set thread_id (e.g. OutgoingResponse::text(...) notifications). Prefer keeping thread_id optional in SseEvent::Response (so JSON can omit/null it) or avoid broadcasting when thread_id is None and log a warning.
There was a problem hiding this comment.
Fixed in 1b030b2 — added a tracing::warn! log when thread_id is None so the empty-string fallback is no longer silent.
| let row = conn | ||
| .query_opt( | ||
| r#" | ||
| SELECT id FROM conversations | ||
| WHERE user_id = $1 AND metadata->>'routine_id' = $2 | ||
| LIMIT 1 | ||
| "#, | ||
| &[&user_id, &rid], | ||
| ) | ||
| .await?; | ||
|
|
||
| if let Some(row) = row { | ||
| return Ok(row.get("id")); | ||
| } | ||
|
|
||
| let id = Uuid::new_v4(); | ||
| let metadata = serde_json::json!({ | ||
| "thread_type": "routine", | ||
| "routine_id": routine_id.to_string(), | ||
| "routine_name": routine_name, | ||
| }); | ||
| conn.execute( | ||
| "INSERT INTO conversations (id, channel, user_id, metadata) VALUES ($1, $2, $3, $4)", | ||
| &[&id, &"routine", &user_id, &metadata], | ||
| ) |
There was a problem hiding this comment.
get_or_create_routine_conversation does a SELECT then INSERT without any uniqueness constraint/transactional protection. If two routine runs for the same routine_id overlap, this can create duplicate conversations and split history/unread counts. Consider adding a deterministic key (e.g. routine_id column) with a unique index, or use an UPSERT pattern (INSERT .. ON CONFLICT .. RETURNING) guarded by a unique constraint.
There was a problem hiding this comment.
Fixed in 1b030b2. Postgres routine get-or-create now uses INSERT ON CONFLICT on uq_conv_routine unique index + SELECT-back. V11 migration adds the index.
| let row = conn | ||
| .query_opt( | ||
| r#" | ||
| SELECT id FROM conversations | ||
| WHERE user_id = $1 AND metadata->>'thread_type' = 'heartbeat' | ||
| LIMIT 1 | ||
| "#, | ||
| &[&user_id], | ||
| ) | ||
| .await?; | ||
|
|
||
| if let Some(row) = row { | ||
| return Ok(row.get("id")); | ||
| } | ||
|
|
||
| let id = Uuid::new_v4(); | ||
| let metadata = serde_json::json!({ | ||
| "thread_type": "heartbeat", | ||
| }); | ||
| conn.execute( | ||
| "INSERT INTO conversations (id, channel, user_id, metadata) VALUES ($1, $2, $3, $4)", | ||
| &[&id, &"heartbeat", &user_id, &metadata], | ||
| ) | ||
| .await?; |
There was a problem hiding this comment.
get_or_create_heartbeat_conversation uses a SELECT-then-INSERT without a uniqueness constraint, so concurrent heartbeat notifications can create multiple “singleton” heartbeat conversations for the same user. To guarantee singleton semantics, add a unique constraint (e.g. on (user_id, channel) for channel='heartbeat' or a dedicated column) and use INSERT .. ON CONFLICT to atomically fetch/create.
There was a problem hiding this comment.
Fixed in 1b030b2. Postgres heartbeat get-or-create now uses INSERT ON CONFLICT on uq_conv_heartbeat unique index + SELECT-back. V11 migration adds the index.
…cast warning - Fix TOCTOU race in get_or_create_routine_conversation (postgres): use INSERT ON CONFLICT on new uq_conv_routine unique index + SELECT-back. - Fix TOCTOU race in get_or_create_heartbeat_conversation (postgres): use INSERT ON CONFLICT on new uq_conv_heartbeat unique index + SELECT-back. - Fix TOCTOU race in get_or_create_routine_conversation (libsql): use BEGIN IMMEDIATE transaction to serialize concurrent writers. - Fix TOCTOU race in get_or_create_heartbeat_conversation (libsql): use BEGIN IMMEDIATE transaction to serialize concurrent writers. - Add V11 migration with partial unique indexes for postgres. - Add matching unique indexes to libsql schema. - Update stale comment on isCurrentThread (said "always shown" but logic now returns false for missing thread_id). - Debounce loadThreads() on off-thread SSE events to prevent request storms. - Log warning in broadcast() when thread_id is None (clients will drop it). Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
The in-memory thread list fallback (when no DB is available) used HashMap::values() which has no guaranteed ordering. Sort by updated_at descending to match the SQL query ordering. [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 15 out of 15 changed files in this pull request and generated 4 comments.
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| r#" | ||
| INSERT INTO conversations (id, channel, user_id, metadata) | ||
| VALUES ($1, 'routine', $2, $3) | ||
| ON CONFLICT ON CONSTRAINT uq_conv_routine DO NOTHING |
There was a problem hiding this comment.
ON CONFLICT ON CONSTRAINT uq_conv_routine will fail at runtime because V11 creates uq_conv_routine as a unique index, not a unique constraint (Postgres only allows ON CONFLICT ON CONSTRAINT for named constraints). Fix by either (a) switching this INSERT to use index inference (ON CONFLICT (user_id, (metadata->>'routine_id')) WHERE metadata->>'routine_id' IS NOT NULL DO NOTHING) or (b) updating the migration to add a named unique constraint (e.g., ALTER TABLE ... ADD CONSTRAINT ... UNIQUE USING INDEX ...) and then keeping ON CONFLICT ON CONSTRAINT.
| ON CONFLICT ON CONSTRAINT uq_conv_routine DO NOTHING | |
| ON CONFLICT (user_id, (metadata->>'routine_id')) | |
| WHERE metadata->>'routine_id' IS NOT NULL | |
| DO NOTHING |
There was a problem hiding this comment.
Fixed in 092ca85 — switched from ON CONFLICT ON CONSTRAINT uq_conv_routine to ON CONFLICT (user_id, (metadata->>'routine_id')) WHERE metadata->>'routine_id' IS NOT NULL, which works with unique indexes.
| r#" | ||
| INSERT INTO conversations (id, channel, user_id, metadata) | ||
| VALUES ($1, 'heartbeat', $2, $3) | ||
| ON CONFLICT ON CONSTRAINT uq_conv_heartbeat DO NOTHING |
There was a problem hiding this comment.
Same issue as the routine insert: ON CONFLICT ON CONSTRAINT uq_conv_heartbeat assumes uq_conv_heartbeat is a named constraint, but the migration creates it as a unique index. This will error at runtime on Postgres. Use index inference (ON CONFLICT (user_id) WHERE metadata->>'thread_type' = 'heartbeat' DO NOTHING) or change the migration to create/attach a unique constraint and keep ON CONFLICT ON CONSTRAINT.
| ON CONFLICT ON CONSTRAINT uq_conv_heartbeat DO NOTHING | |
| ON CONFLICT (user_id) WHERE metadata->>'thread_type' = 'heartbeat' DO NOTHING |
There was a problem hiding this comment.
Fixed in 092ca85 — switched from ON CONFLICT ON CONSTRAINT uq_conv_heartbeat to ON CONFLICT (user_id) WHERE metadata->>'thread_type' = 'heartbeat'.
| -- One heartbeat conversation per user. | ||
| CREATE UNIQUE INDEX IF NOT EXISTS uq_conv_heartbeat | ||
| ON conversations (user_id) | ||
| WHERE metadata->>'thread_type' = 'heartbeat'; |
There was a problem hiding this comment.
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.
| 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; |
There was a problem hiding this comment.
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.
| if (thread.title) return thread.title; | ||
| const ch = thread.channel || 'gateway'; | ||
| if (thread.thread_type === 'heartbeat') return 'Heartbeat Alerts'; | ||
| if (thread.thread_type === 'routine') return 'Routine: ' + (thread.title || thread.id.substring(0, 8)); |
There was a problem hiding this comment.
In threadTitle(), the routine branch uses (thread.title || thread.id.substring(0, 8)), but thread.title has already been checked and returned at the top of the function, so it will always be falsy here. Consider simplifying this to avoid misleading/unused logic (and to make it clearer what routine threads will display).
| if (thread.thread_type === 'routine') return 'Routine: ' + (thread.title || thread.id.substring(0, 8)); | |
| if (thread.thread_type === 'routine') return 'Routine: ' + thread.id.substring(0, 8); |
There was a problem hiding this comment.
Fixed in 092ca85 — removed the dead thread.title || since thread.title is always falsy at that point (already returned on line 1268).
The cron ticker's background task occasionally fails with "unable to open database file" when creating a new SQLite connection concurrently with the main thread. Add retry with exponential backoff (50ms, 100ms, 200ms) to handle transient VFS/locking issues in libsql's local mode. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
PostgreSQL ON CONFLICT ON CONSTRAINT requires a named table constraint, but V11 migration creates unique indexes. Switch to the expression form (ON CONFLICT (columns) WHERE condition) which works with unique indexes. Also fix dead code in threadTitle() where thread.title was already checked on the previous line. [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 16 out of 16 changed files in this pull request and generated 2 comments.
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| function isCurrentThread(threadId) { | ||
| if (!threadId) return true; | ||
| if (!threadId) return false; | ||
| if (!currentThreadId) return true; |
There was a problem hiding this comment.
isCurrentThread() returns true when currentThreadId is unset, but connectSSE() is started before the initial loadThreads() call sets a default thread. This allows off-thread notifications (routine/heartbeat/etc.) that arrive during startup to still be rendered into the chat view, partially reintroducing the “notification leaking” this PR is trying to prevent. Consider either (1) treating !currentThreadId as not-current (return false), or (2) delaying SSE connection until after the initial thread selection is established (e.g., loadThreads() resolves and sets currentThreadId).
| if (!currentThreadId) return true; | |
| if (!currentThreadId) return false; |
There was a problem hiding this comment.
The !currentThreadId returning true on line 428 is intentional — during the brief startup window before loadThreads() resolves, showing incoming messages is better UX than a blank screen. Once loadThreads() sets currentThreadId (typically < 100ms), proper filtering kicks in. Returning false here would mean the user sees nothing until they click a thread. The startup race window is negligibly short.
| let thread_id = match response.thread_id { | ||
| Some(tid) => tid, | ||
| None => { | ||
| tracing::warn!( | ||
| "Gateway broadcast with no thread_id — event will be dropped by clients" | ||
| ); | ||
| String::new() | ||
| } | ||
| }; | ||
| self.state.sse.broadcast(SseEvent::Response { | ||
| content: response.content, | ||
| thread_id: String::new(), | ||
| thread_id, | ||
| }); |
There was a problem hiding this comment.
When response.thread_id is None, broadcast() logs a warning but still broadcasts an SSE Response with an empty thread_id. Clients will drop this event (per isCurrentThread), so this becomes avoidable SSE traffic and can still create log spam if any legacy broadcasts remain. Consider returning early after the warning (skip broadcasting), or routing these into a dedicated “system” thread once that exists.
There was a problem hiding this comment.
Fixed in dae234b — broadcast() now returns early when thread_id is None instead of sending an SSE event with an empty thread_id that clients would drop.
[skip-regression-check] Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Clients drop SSE events with empty thread_id anyway, so avoid the unnecessary network traffic by returning early. [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 16 out of 16 changed files in this pull request and generated 5 comments.
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| created_at: chrono::Utc::now().to_rfc3339(), | ||
| updated_at: chrono::Utc::now().to_rfc3339(), | ||
| title: None, | ||
| thread_type: Some("assistant".to_string()), | ||
| channel: Some("gateway".to_string()), | ||
| }); |
There was a problem hiding this comment.
In the DB-backed path, the assistant conversation is only discovered if it happens to be included in list_conversations_all_channels(..., 50). If the user has >50 conversations and the assistant thread is older than the cutoff, assistant_thread stays None and the handler synthesizes an assistant thread with updated_at = now, which misrepresents last activity and can cause confusing UI (shows as just updated). Consider always fetching the assistant conversation explicitly (e.g., separate query by assistant_id, or increase limit / UNION the assistant row) instead of fabricating timestamps.
| @@ -483,8 +485,10 @@ | |||
| updated_at: t.updated_at.to_rfc3339(), | |||
| title: None, | |||
| thread_type: None, | |||
| channel: Some("gateway".to_string()), | |||
| }) | |||
| .collect(); | |||
| threads.sort_by(|a, b| b.updated_at.cmp(&a.updated_at)); | |||
There was a problem hiding this comment.
The fallback in-memory thread list is sorted using updated_at string comparison. RFC3339 strings produced by to_rfc3339() can vary in formatting (e.g., presence/precision of fractional seconds), so lexicographic ordering is not guaranteed to match chronological ordering. Prefer sorting by the underlying DateTime values (e.g., sort sess.threads by t.updated_at before converting to strings).
| 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, |
There was a problem hiding this comment.
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).
| created_at: chrono::Utc::now().to_rfc3339(), | ||
| updated_at: chrono::Utc::now().to_rfc3339(), | ||
| title: None, | ||
| thread_type: Some("assistant".to_string()), | ||
| channel: Some("gateway".to_string()), | ||
| }); |
There was a problem hiding this comment.
In the DB-backed path, the assistant conversation is only discovered if it happens to be included in list_conversations_all_channels(..., 50). If the user has >50 conversations and the assistant thread is older than the cutoff, assistant_thread stays None and the handler synthesizes an assistant thread with updated_at = now, which misrepresents last activity and can cause confusing UI (shows as just updated). Consider always fetching the assistant conversation explicitly (e.g., separate query by assistant_id, or increase limit / UNION the assistant row) instead of fabricating timestamps.
| @@ -1094,8 +1096,10 @@ | |||
| updated_at: t.updated_at.to_rfc3339(), | |||
| title: None, | |||
| thread_type: None, | |||
| channel: Some("gateway".to_string()), | |||
| }) | |||
| .collect(); | |||
| threads.sort_by(|a, b| b.updated_at.cmp(&a.updated_at)); | |||
There was a problem hiding this comment.
The fallback in-memory thread list is sorted using updated_at string comparison. RFC3339 strings produced by to_rfc3339() can vary in formatting (e.g., presence/precision of fractional seconds), so lexicographic ordering is not guaranteed to match chronological ordering. Prefer sorting by the underlying DateTime values (e.g., sort sess.threads by t.updated_at before converting to strings).
Add tests proving get_or_create_routine_conversation returns the same conversation ID across multiple invocations with the same routine_id. Add debug logging to routine engine to track conversation resolution. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
… messages Three independent fixes with regression tests: 1. Routine conversations now display in the web UI. build_turns_from_db_messages() handles standalone assistant messages (no preceding user message) by creating turns with empty user_input. Frontend skips empty user bubbles. 2. Worker select_tools and execute_plan paths now push an assistant_with_tool_calls message before tool execution, preventing sanitize_tool_messages from rewriting tool_results as orphaned user messages. 3. Reasoning::plan() and respond_with_tools() merge system messages from context into a single system prompt instead of creating [system, system, ...] sequences that strict LLM providers (Qwen) reject. Also: sidebar padding/spacing improvements, wider thread panel (240px). Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 26 out of 26 changed files in this pull request and generated 5 comments.
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| let engine_guard = state.routine_engine.read().await; | ||
| let engine = engine_guard.as_ref().ok_or(( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| "Database not available".to_string(), | ||
| "Routine engine not available".to_string(), | ||
| ))?; | ||
|
|
||
| let routine_id = Uuid::parse_str(&id) | ||
| .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid routine ID".to_string()))?; | ||
|
|
||
| let routine = store | ||
| .get_routine(routine_id) | ||
| let run_id = engine | ||
| .fire_manual(routine_id) | ||
| .await |
There was a problem hiding this comment.
The routine_engine slot read-lock is held across the .await of engine.fire_manual(...). Consider cloning the engine Arc out of the lock and dropping the guard before awaiting to prevent unnecessary lock retention.
There was a problem hiding this comment.
Fixed in 32afa42 — routines_trigger_handler now clones the Arc<RoutineEngine> out of the RwLock before awaiting fire_manual().
| self.safety().clone(), | ||
| Some(notify_tx), | ||
| self.store().map(Arc::clone), | ||
| )) |
There was a problem hiding this comment.
Heartbeat persistence uses HeartbeatRunner's config.notify_user_id to choose which user's heartbeat conversation to write to, but the agent currently constructs AgentHeartbeatConfig without calling .with_notify(...) (only sets interval). That means persisted heartbeat messages will default to user_id "default" even when hb_config.notify_user is configured, and the heartbeat thread may not appear for the intended user. Populate the runner config with the resolved notify user/channel when wiring heartbeat in the agent loop.
There was a problem hiding this comment.
Fixed in 32afa42 — agent_loop.rs now wires hb_config.notify_user and hb_config.notify_channel into AgentHeartbeatConfig via .with_notify().
| let engine_guard = state.routine_engine.read().await; | ||
| let engine = engine_guard.as_ref().ok_or(( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| "Database not available".to_string(), | ||
| "Routine engine not available".to_string(), | ||
| ))?; | ||
|
|
||
| let routine_id = Uuid::parse_str(&id) | ||
| .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid routine ID".to_string()))?; | ||
|
|
||
| let routine = store | ||
| .get_routine(routine_id) | ||
| let run_id = engine | ||
| .fire_manual(routine_id) | ||
| .await | ||
| .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))? | ||
| .ok_or((StatusCode::NOT_FOUND, "Routine not found".to_string()))?; | ||
|
|
||
| // Send the routine prompt through the message pipeline as a manual trigger. | ||
| let prompt = match &routine.action { | ||
| crate::agent::routine::RoutineAction::Lightweight { prompt, .. } => prompt.clone(), | ||
| crate::agent::routine::RoutineAction::FullJob { | ||
| title, description, .. | ||
| } => format!("{}: {}", title, description), | ||
| }; | ||
|
|
||
| let content = format!("[routine:{}] {}", routine.name, prompt); | ||
| let msg = IncomingMessage::new("gateway", &state.user_id, content); | ||
|
|
||
| let tx_guard = state.msg_tx.read().await; | ||
| let tx = tx_guard.as_ref().ok_or(( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| "Channel not started".to_string(), | ||
| ))?; | ||
|
|
||
| tx.send(msg).await.map_err(|_| { | ||
| ( | ||
| StatusCode::INTERNAL_SERVER_ERROR, | ||
| "Channel closed".to_string(), | ||
| ) | ||
| })?; | ||
| .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; |
There was a problem hiding this comment.
routines_trigger_handler no longer verifies that the routine being triggered belongs to the current gateway user_id. Since RoutineStore::get_routine(id) is not user-scoped, a caller who can guess a UUID could trigger another user's routine. Reintroduce an ownership check (load the routine and compare routine.user_id to state.user_id) or change RoutineEngine::fire_manual to take user_id and enforce it internally.
There was a problem hiding this comment.
Fixed in 32afa42 — fire_manual() now takes user_id: Option<&str> and returns RoutineError::NotAuthorized if the routine's user_id doesn't match. Gateway handler passes Some(&state.user_id).
| let engine_guard = state.routine_engine.read().await; | ||
| let engine = engine_guard.as_ref().ok_or(( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| "Database not available".to_string(), | ||
| "Routine engine not available".to_string(), | ||
| ))?; | ||
|
|
||
| let routine_id = Uuid::parse_str(&id) | ||
| .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid routine ID".to_string()))?; | ||
|
|
||
| let routine = store | ||
| .get_routine(routine_id) | ||
| let run_id = engine | ||
| .fire_manual(routine_id) | ||
| .await |
There was a problem hiding this comment.
This handler holds the async routine_engine slot read-lock across the .await of engine.fire_manual(...). To avoid lock contention and potential deadlock patterns, clone the Arc<RoutineEngine> out of the slot (or otherwise drop the guard) before awaiting the manual fire.
There was a problem hiding this comment.
Fixed in 32afa42 — same fix as handlers/routines.rs: clones Arc<RoutineEngine> out of the RwLock before awaiting.
| let engine_guard = state.routine_engine.read().await; | ||
| let engine = engine_guard.as_ref().ok_or(( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| "Database not available".to_string(), | ||
| "Routine engine not available".to_string(), | ||
| ))?; | ||
|
|
||
| let routine_id = Uuid::parse_str(&id) | ||
| .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid routine ID".to_string()))?; | ||
|
|
||
| let routine = store | ||
| .get_routine(routine_id) | ||
| let run_id = engine | ||
| .fire_manual(routine_id) | ||
| .await | ||
| .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))? | ||
| .ok_or((StatusCode::NOT_FOUND, "Routine not found".to_string()))?; | ||
|
|
||
| if routine.user_id != state.user_id { | ||
| return Err((StatusCode::FORBIDDEN, "Access denied".to_string())); | ||
| } | ||
|
|
||
| // Send the routine prompt through the message pipeline as a manual trigger. | ||
| let prompt = match &routine.action { | ||
| crate::agent::routine::RoutineAction::Lightweight { prompt, .. } => prompt.clone(), | ||
| crate::agent::routine::RoutineAction::FullJob { | ||
| title, description, .. | ||
| } => format!("{}: {}", title, description), | ||
| }; | ||
|
|
||
| let content = format!("[routine:{}] {}", routine.name, prompt); | ||
| let thread_id = format!( | ||
| "routine-{}-{}", | ||
| routine_id, | ||
| chrono::Utc::now().timestamp_millis() | ||
| ); | ||
| let msg = IncomingMessage::new("gateway", &state.user_id, content).with_thread(thread_id); | ||
|
|
||
| let tx_guard = state.msg_tx.read().await; | ||
| let tx = tx_guard.as_ref().ok_or(( | ||
| StatusCode::SERVICE_UNAVAILABLE, | ||
| "Channel not started".to_string(), | ||
| ))?; | ||
|
|
||
| tx.send(msg).await.map_err(|_| { | ||
| ( | ||
| StatusCode::INTERNAL_SERVER_ERROR, | ||
| "Channel closed".to_string(), | ||
| ) | ||
| })?; | ||
| .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; |
There was a problem hiding this comment.
routines_trigger_handler no longer verifies that the routine belongs to state.user_id. Because get_routine(id) is not user-scoped, this can allow triggering routines owned by other users. Add an ownership check (load routine and compare routine.user_id to state.user_id) or enforce it inside RoutineEngine::fire_manual.
There was a problem hiding this comment.
Fixed in 32afa42 — same ownership check fix. fire_manual(routine_id, Some(&state.user_id)) now enforces that the routine belongs to the requesting user.
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
…ship check, heartbeat config - Clone Arc<RoutineEngine> out of RwLock before .await in trigger handler - Add user_id ownership check to fire_manual() with NotAuthorized error - Wire heartbeat notify_user/notify_channel from config to AgentHeartbeatConfig Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 37 out of 45 changed files in this pull request and generated 10 comments.
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| "Channel closed".to_string(), | ||
| ) | ||
| })?; | ||
| .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; |
There was a problem hiding this comment.
routines_trigger_handler now maps any RoutineEngine::fire_manual error to HTTP 500. This will return 500 for expected conditions like unknown routine IDs (RoutineError::NotFound) and unauthorized triggers (RoutineError::NotAuthorized). Please match on RoutineError and translate to appropriate status codes (e.g. 404 for NotFound, 403 for NotAuthorized, 409/400 for Disabled/MaxConcurrent), preserving 500 only for unexpected failures.
| .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; | |
| .map_err(|e| { | |
| match e { | |
| crate::routines::engine::RoutineError::NotFound => ( | |
| StatusCode::NOT_FOUND, | |
| "Routine not found".to_string(), | |
| ), | |
| crate::routines::engine::RoutineError::NotAuthorized => ( | |
| StatusCode::FORBIDDEN, | |
| "Not authorized to trigger this routine".to_string(), | |
| ), | |
| crate::routines::engine::RoutineError::Disabled => ( | |
| StatusCode::CONFLICT, | |
| "Routine is disabled".to_string(), | |
| ), | |
| crate::routines::engine::RoutineError::MaxConcurrent => ( | |
| StatusCode::CONFLICT, | |
| "Maximum concurrent runs reached".to_string(), | |
| ), | |
| other => ( | |
| StatusCode::INTERNAL_SERVER_ERROR, | |
| other.to_string(), | |
| ), | |
| } | |
| })?; |
There was a problem hiding this comment.
Fixed in db760a1 — fire_manual errors now map to proper HTTP status codes: NotFound → 404, NotAuthorized → 403, Disabled/MaxConcurrent → 409, others → 500.
| "Channel closed".to_string(), | ||
| ) | ||
| })?; | ||
| .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; |
There was a problem hiding this comment.
routines_trigger_handler now returns HTTP 500 for all fire_manual errors. This breaks previous behavior where NotFound/AccessDenied were surfaced as 404/403, and makes client error handling harder. Consider mapping RoutineError variants to specific status codes (404 NotFound, 403 NotAuthorized, etc.) and leaving 500 for unexpected errors only.
| .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; | |
| .map_err(|e| { | |
| let status = match e { | |
| RoutineError::NotFound(..) => StatusCode::NOT_FOUND, | |
| RoutineError::AccessDenied(..) | RoutineError::NotAuthorized(..) => StatusCode::FORBIDDEN, | |
| _ => StatusCode::INTERNAL_SERVER_ERROR, | |
| }; | |
| (status, e.to_string()) | |
| })?; |
There was a problem hiding this comment.
Fixed in db760a1 — same fix applied to server.rs handler. RoutineError variants now map to 404/403/409/500.
| { | ||
| "model_name": "recorded-zai-org/GLM-5-FP8", | ||
| "memory_snapshot": [ | ||
| { | ||
| "path": "AGENTS.md", | ||
| "content": "# Agent Instructions\n\nYou are a personal AI assistant with access to tools and persistent memory.\n\n## Guidelines\n\n- Always search memory before answering questions about prior conversations\n- Write important facts and decisions to memory for future reference\n- Use the daily log for session-level notes\n- Be concise but thorough" | ||
| }, | ||
| { | ||
| "path": "HEARTBEAT.md", | ||
| "content": "# Heartbeat Checklist\n\n<!-- Keep this file empty to skip heartbeat API calls.\n Add tasks below when you want the agent to check something periodically.\n\n Example:\n - [ ] Check for unread emails needing a reply\n - [ ] Review today's calendar for upcoming meetings\n - [ ] Check CI build status for main branch\n-->" | ||
| }, | ||
| { | ||
| "path": "IDENTITY.md", | ||
| "content": "# Identity\n\nName: Alfred\nNature: A secure personal AI assistant\n\nEdit this file to give your agent a custom name and personality." | ||
| }, |
There was a problem hiding this comment.
This trace JSON file appears to be a generated recording (model name, memory snapshots, user prompts, OAuth URLs/state, etc.). Committing it to the repo risks leaking sensitive/user data and adds significant noise/size to the PR. Please remove it from the pull request (and consider adding a .gitignore rule or storing sanitized fixtures under a dedicated test/ directory if needed).
| { | ||
| "model_name": "recorded-zai-org/GLM-5-FP8", | ||
| "memory_snapshot": [ | ||
| { | ||
| "path": "AGENTS.md", | ||
| "content": "# Agent Instructions\n\nYou are a personal AI assistant with access to tools and persistent memory.\n\n## Guidelines\n\n- Always search memory before answering questions about prior conversations\n- Write important facts and decisions to memory for future reference\n- Use the daily log for session-level notes\n- Be concise but thorough" | ||
| }, | ||
| { | ||
| "path": "HEARTBEAT.md", | ||
| "content": "# Heartbeat Checklist\n\n<!-- Keep this file empty to skip heartbeat API calls.\n Add tasks below when you want the agent to check something periodically.\n\n Example:\n - [ ] Check for unread emails needing a reply\n - [ ] Review today's calendar for upcoming meetings\n - [ ] Check CI build status for main branch\n-->" | ||
| }, | ||
| { | ||
| "path": "IDENTITY.md", | ||
| "content": "# Identity\n\nName: Alfred\nNature: A secure personal AI assistant\n\nEdit this file to give your agent a custom name and personality." | ||
| }, |
There was a problem hiding this comment.
This trace JSON file contains recorded conversation content and workspace snapshots. Keeping it in the repository is a privacy/security risk and likely not intended as a source artifact. Please drop it from the PR (or replace with a sanitized, minimal fixture under a test-only path).
| { | ||
| "model_name": "recorded-zai-org/GLM-5-FP8", | ||
| "memory_snapshot": [ | ||
| { | ||
| "path": "AGENTS.md", | ||
| "content": "# Agent Instructions\n\nYou are a personal AI assistant with access to tools and persistent memory.\n\n## Guidelines\n\n- Always search memory before answering questions about prior conversations\n- Write important facts and decisions to memory for future reference\n- Use the daily log for session-level notes\n- Be concise but thorough" | ||
| }, | ||
| { | ||
| "path": "HEARTBEAT.md", | ||
| "content": "# Heartbeat Checklist\n\n<!-- Keep this file empty to skip heartbeat API calls.\n Add tasks below when you want the agent to check something periodically.\n\n Example:\n - [ ] Check for unread emails needing a reply\n - [ ] Review today's calendar for upcoming meetings\n - [ ] Check CI build status for main branch\n-->" | ||
| }, | ||
| { | ||
| "path": "IDENTITY.md", | ||
| "content": "# Identity\n\nName: Alfred\nNature: A secure personal AI assistant\n\nEdit this file to give your agent a custom name and personality." | ||
| }, |
There was a problem hiding this comment.
This trace JSON file looks like an agent recording/debug artifact and includes user content. Committing it increases repo size and can leak private data. Please remove it from the PR and prevent future accidental commits (e.g., via .gitignore).
| { | ||
| "model_name": "recorded-zai-org/GLM-5-FP8", | ||
| "memory_snapshot": [ | ||
| { | ||
| "path": "AGENTS.md", | ||
| "content": "# Agent Instructions\n\nYou are a personal AI assistant with access to tools and persistent memory.\n\n## Guidelines\n\n- Always search memory before answering questions about prior conversations\n- Write important facts and decisions to memory for future reference\n- Use the daily log for session-level notes\n- Be concise but thorough" | ||
| }, | ||
| { | ||
| "path": "HEARTBEAT.md", | ||
| "content": "# Heartbeat Checklist\n\n<!-- Keep this file empty to skip heartbeat API calls.\n Add tasks below when you want the agent to check something periodically.\n\n Example:\n - [ ] Check for unread emails needing a reply\n - [ ] Review today's calendar for upcoming meetings\n - [ ] Check CI build status for main branch\n-->" | ||
| }, | ||
| { | ||
| "path": "IDENTITY.md", | ||
| "content": "# Identity\n\nName: Alfred\nNature: A secure personal AI assistant\n\nEdit this file to give your agent a custom name and personality." | ||
| }, |
There was a problem hiding this comment.
This trace JSON file appears to be a recorded interaction log. It includes conversation content and other runtime details that shouldn’t be committed to the repository. Please remove it from this PR (or provide a sanitized, minimal test fixture if required).
| { | ||
| "model_name": "recorded-zai-org/GLM-5-FP8", | ||
| "memory_snapshot": [ | ||
| { | ||
| "path": "AGENTS.md", | ||
| "content": "# Agent Instructions\n\nYou are a personal AI assistant with access to tools and persistent memory.\n\n## Guidelines\n\n- Always search memory before answering questions about prior conversations\n- Write important facts and decisions to memory for future reference\n- Use the daily log for session-level notes\n- Be concise but thorough" | ||
| }, | ||
| { | ||
| "path": "HEARTBEAT.md", | ||
| "content": "# Heartbeat Checklist\n\n<!-- Keep this file empty to skip heartbeat API calls.\n Add tasks below when you want the agent to check something periodically.\n\n Example:\n - [ ] Check for unread emails needing a reply\n - [ ] Review today's calendar for upcoming meetings\n - [ ] Check CI build status for main branch\n-->" | ||
| }, | ||
| { | ||
| "path": "IDENTITY.md", | ||
| "content": "# Identity\n\nName: Alfred\nNature: A secure personal AI assistant\n\nEdit this file to give your agent a custom name and personality." | ||
| }, |
There was a problem hiding this comment.
This trace JSON file looks like a recorded session log and includes potentially sensitive content. It should not be committed as part of this feature change. Please remove it from the PR and add a .gitignore rule if these traces are generated locally.
| { | ||
| "model_name": "recorded-zai-org/GLM-5-FP8", | ||
| "memory_snapshot": [ | ||
| { | ||
| "path": "AGENTS.md", | ||
| "content": "# Agent Instructions\n\nYou are a personal AI assistant with access to tools and persistent memory.\n\n## Guidelines\n\n- Always search memory before answering questions about prior conversations\n- Write important facts and decisions to memory for future reference\n- Use the daily log for session-level notes\n- Be concise but thorough" | ||
| }, | ||
| { | ||
| "path": "HEARTBEAT.md", | ||
| "content": "# Heartbeat Checklist\n\n<!-- Keep this file empty to skip heartbeat API calls.\n Add tasks below when you want the agent to check something periodically.\n\n Example:\n - [ ] Check for unread emails needing a reply\n - [ ] Review today's calendar for upcoming meetings\n - [ ] Check CI build status for main branch\n-->" | ||
| }, | ||
| { | ||
| "path": "IDENTITY.md", | ||
| "content": "# Identity\n\nName: Alfred\nNature: A secure personal AI assistant\n\nEdit this file to give your agent a custom name and personality." | ||
| }, |
There was a problem hiding this comment.
This trace JSON file includes user inputs and potentially sensitive runtime details. It looks like a debug/recording artifact rather than production source. Please remove it from version control (or sanitize and relocate to test fixtures if it’s needed for regression reproduction).
| { | ||
| "model_name": "recorded-zai-org/GLM-5-FP8", | ||
| "memory_snapshot": [ | ||
| { | ||
| "path": "AGENTS.md", | ||
| "content": "# Agent Instructions\n\nYou are a personal AI assistant with access to tools and persistent memory.\n\n## Guidelines\n\n- Always search memory before answering questions about prior conversations\n- Write important facts and decisions to memory for future reference\n- Use the daily log for session-level notes\n- Be concise but thorough" | ||
| }, | ||
| { | ||
| "path": "HEARTBEAT.md", | ||
| "content": "# Heartbeat Checklist\n\n<!-- Keep this file empty to skip heartbeat API calls.\n Add tasks below when you want the agent to check something periodically.\n\n Example:\n - [ ] Check for unread emails needing a reply\n - [ ] Review today's calendar for upcoming meetings\n - [ ] Check CI build status for main branch\n-->" | ||
| }, | ||
| { | ||
| "path": "IDENTITY.md", | ||
| "content": "# Identity\n\nName: Alfred\nNature: A secure personal AI assistant\n\nEdit this file to give your agent a custom name and personality." | ||
| }, |
There was a problem hiding this comment.
This trace JSON file contains a full recorded interaction (including tool outputs and auth URLs). This is sensitive and should not be committed to the repo. Please remove it from the PR and consider adding safeguards (gitignore/pre-commit) to prevent trace artifacts from being checked in.
| { | ||
| "model_name": "recorded-zai-org/GLM-5-FP8", | ||
| "memory_snapshot": [ | ||
| { | ||
| "path": "AGENTS.md", | ||
| "content": "# Agent Instructions\n\nYou are a personal AI assistant with access to tools and persistent memory.\n\n## Guidelines\n\n- Always search memory before answering questions about prior conversations\n- Write important facts and decisions to memory for future reference\n- Use the daily log for session-level notes\n- Be concise but thorough" | ||
| }, | ||
| { | ||
| "path": "HEARTBEAT.md", | ||
| "content": "# Heartbeat Checklist\n\n<!-- Keep this file empty to skip heartbeat API calls.\n Add tasks below when you want the agent to check something periodically.\n\n Example:\n - [ ] Check for unread emails needing a reply\n - [ ] Review today's calendar for upcoming meetings\n - [ ] Check CI build status for main branch\n-->" | ||
| }, | ||
| { | ||
| "path": "IDENTITY.md", | ||
| "content": "# Identity\n\nName: Alfred\nNature: A secure personal AI assistant\n\nEdit this file to give your agent a custom name and personality." | ||
| }, |
There was a problem hiding this comment.
This trace JSON file appears to be an execution recording/debug artifact. It includes conversation content and increases the repository size. Please remove it from this PR (or sanitize and move to a dedicated test-fixture location if it’s required).
…odel # Conflicts: # tests/openai_compat_integration.rs # tests/ws_gateway_integration.rs
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
…odel # Conflicts: # src/llm/reasoning.rs
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 29 out of 30 changed files in this pull request and generated 2 comments.
Comments suppressed due to low confidence (1)
src/channels/web/static/app.js:289
connectSSE’sresponsehandler callsenableChatInput()unconditionally before refreshing threads. When the active thread is an external/read-only channel, this can briefly re-enable input untilloadThreads()completes (and other SSE handlers may re-enable it without a subsequent refresh). Consider enabling/disabling input based on the current thread’s channel/type rather than always enabling here.
finalizeActivityGroup();
addMessage('assistant', data.content);
enableChatInput();
// Refresh thread list so new titles appear after first message
loadThreads();
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| function enableChatInput() { | ||
| // no-op: input and send button are always enabled | ||
| const input = document.getElementById('chat-input'); | ||
| const btn = document.getElementById('send-btn'); | ||
| if (input) { | ||
| input.disabled = false; | ||
| input.placeholder = 'Message or / for commands...'; | ||
| } | ||
| if (btn) btn.disabled = false; | ||
| } |
There was a problem hiding this comment.
enableChatInput() always re-enables the input/button. With the new concept of read-only external threads, any call site that invokes enableChatInput() (e.g. SSE status/error/auth events) can accidentally re-enable input while a read-only thread is active. Consider folding the read-only check into this function (or adding a separate syncChatInputForCurrentThread() helper) so the read-only rule is enforced centrally.
There was a problem hiding this comment.
Fixed in db760a1 — enableChatInput() now checks a currentThreadIsReadOnly flag and returns early if the current thread is read-only. The flag is set during thread switching in loadThreads() and switchToAssistant().
| let thread_id = match response.thread_id { | ||
| Some(tid) => tid, | ||
| None => { | ||
| tracing::warn!( | ||
| "Gateway broadcast with no thread_id — skipping (clients would drop it)" | ||
| ); | ||
| return Ok(()); | ||
| } | ||
| }; | ||
| self.state.sse.broadcast(SseEvent::Response { | ||
| content: response.content, | ||
| thread_id: String::new(), | ||
| thread_id, | ||
| }); |
There was a problem hiding this comment.
broadcast() now correctly skips responses without a thread_id, but respond() still uses msg.thread_id.clone().unwrap_or_default() (sending an empty string). With the frontend now dropping falsy thread_ids, any response to a message without a thread id will be silently lost. Consider applying the same guard/logging in respond() (or otherwise enforcing that all gateway IncomingMessages have a thread id).
There was a problem hiding this comment.
Fixed in db760a1 — respond() now matches broadcast() behavior: returns early with a warning when thread_id is None instead of sending an empty string.
…rd, respond thread_id - Map RoutineError::NotFound → 404, NotAuthorized → 403, Disabled → 409 - Guard enableChatInput() against re-enabling on read-only threads - Skip respond() when thread_id is None (matches broadcast() behavior) [skip-regression-check] Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* feat: unified thread model for web gateway
Every piece of activity (user chat, routine run, heartbeat alert, external
channel message) now lives in its own thread, properly isolated, with
meaningful titles and visual distinction.
Key changes:
- Add `channel` field to ConversationSummary and ThreadInfo so the gateway
can distinguish thread origins (gateway, telegram, routine, heartbeat).
- Add `list_conversations_all_channels` to Database trait (both postgres
and libsql) so chat_threads_handler shows cross-channel threads.
- Routine runs get a persistent conversation per routine via
`get_or_create_routine_conversation`; notifications carry thread_id.
- Heartbeat gets a persistent conversation via
`get_or_create_heartbeat_conversation`; HeartbeatRunner accepts an
optional Database store and binds notifications to the thread.
- Fix broadcast() in web gateway to propagate response.thread_id instead
of hardcoding empty string.
- Fix isCurrentThread(null) returning true (the core notification leak
bug) — now returns false so events without a thread_id don't leak into
the active thread.
- Rewrite frontend thread sidebar: meaningful titles with channel-specific
fallbacks, relative timestamps instead of turn counts, channel badges
for non-gateway threads, unread notification dots, read-only indicator
for external channel threads.
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: address PR review — TOCTOU races, stale comment, debounce, broadcast warning
- Fix TOCTOU race in get_or_create_routine_conversation (postgres):
use INSERT ON CONFLICT on new uq_conv_routine unique index + SELECT-back.
- Fix TOCTOU race in get_or_create_heartbeat_conversation (postgres):
use INSERT ON CONFLICT on new uq_conv_heartbeat unique index + SELECT-back.
- Fix TOCTOU race in get_or_create_routine_conversation (libsql):
use BEGIN IMMEDIATE transaction to serialize concurrent writers.
- Fix TOCTOU race in get_or_create_heartbeat_conversation (libsql):
use BEGIN IMMEDIATE transaction to serialize concurrent writers.
- Add V11 migration with partial unique indexes for postgres.
- Add matching unique indexes to libsql schema.
- Update stale comment on isCurrentThread (said "always shown" but logic
now returns false for missing thread_id).
- Debounce loadThreads() on off-thread SSE events to prevent request storms.
- Log warning in broadcast() when thread_id is None (clients will drop it).
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: sort in-memory thread fallback by updated_at descending
The in-memory thread list fallback (when no DB is available) used
HashMap::values() which has no guaranteed ordering. Sort by
updated_at descending to match the SQL query ordering.
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: retry libsql connect() on transient "unable to open database file"
The cron ticker's background task occasionally fails with "unable to
open database file" when creating a new SQLite connection concurrently
with the main thread. Add retry with exponential backoff (50ms, 100ms,
200ms) to handle transient VFS/locking issues in libsql's local mode.
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: use ON CONFLICT with index expressions instead of named constraints
PostgreSQL ON CONFLICT ON CONSTRAINT requires a named table constraint,
but V11 migration creates unique indexes. Switch to the expression form
(ON CONFLICT (columns) WHERE condition) which works with unique indexes.
Also fix dead code in threadTitle() where thread.title was already
checked on the previous line.
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* style: fix rustfmt chain collapse in heartbeat.rs
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: skip broadcast when thread_id is None instead of sending empty
Clients drop SSE events with empty thread_id anyway, so avoid the
unnecessary network traffic by returning early.
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* test: add libsql routine/heartbeat conversation idempotency tests
Add tests proving get_or_create_routine_conversation returns the same
conversation ID across multiple invocations with the same routine_id.
Add debug logging to routine engine to track conversation resolution.
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* feat: show "New chat" title for empty threads
- threadTitle() returns "New chat" when turn_count is 0
- Assistant thread label updates dynamically from API data
- Default HTML label changed from "Assistant" to "New chat"
- New threads naturally sort to top via last_activity DESC
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: thread sorting, routine isolation, and UI polish
- Fix libsql timestamp format mismatch causing broken thread sort order.
SQLite defaults used `datetime('now')` (space-separated) while Rust code
used RFC3339 (T-separated), breaking string-based ORDER BY. All INSERTs
now use RFC3339, and queries use `datetime()` to normalize comparison.
- Route manual routine triggers through RoutineEngine.fire_manual() instead
of injecting as regular chat messages, so routines always run in their
dedicated conversation thread.
- Add RoutineEngineSlot to GatewayState for gateway<->engine communication.
- Derive routine thread titles from conversation metadata (routine_name)
instead of showing truncated UUID hashes.
- Make chat_new_thread_handler persist to DB synchronously so loadThreads()
sees newly created threads immediately.
- Fix enableChatInput() no-op and wrong element ID in disableChatInputReadOnly().
- Fix handlers/chat.rs stale gateway-only query (use list_conversations_all_channels).
- Sort in-memory threads by DateTime before converting to RFC3339 strings.
- Trigger debouncedLoadThreads() on thinking/status SSE events for non-current
threads so routine/heartbeat threads appear in sidebar promptly.
- Remove "Threads" text from sidebar header.
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: routine history display, orphaned tool_results, duplicate system messages
Three independent fixes with regression tests:
1. Routine conversations now display in the web UI. build_turns_from_db_messages()
handles standalone assistant messages (no preceding user message) by creating
turns with empty user_input. Frontend skips empty user bubbles.
2. Worker select_tools and execute_plan paths now push an
assistant_with_tool_calls message before tool execution, preventing
sanitize_tool_messages from rewriting tool_results as orphaned user messages.
3. Reasoning::plan() and respond_with_tools() merge system messages from
context into a single system prompt instead of creating [system, system, ...]
sequences that strict LLM providers (Qwen) reject.
Also: sidebar padding/spacing improvements, wider thread panel (240px).
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: address PR nearai#607 review — RwLock held across await, missing ownership check, heartbeat config
- Clone Arc<RoutineEngine> out of RwLock before .await in trigger handler
- Add user_id ownership check to fire_manual() with NotAuthorized error
- Wire heartbeat notify_user/notify_channel from config to AgentHeartbeatConfig
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* chore: gitignore trace_*.json files and remove stale traces
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* chore: remove trace JSON files from repo
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: proper HTTP status codes for routine errors, read-only input guard, respond thread_id
- Map RoutineError::NotFound → 404, NotAuthorized → 403, Disabled → 409
- Guard enableChatInput() against re-enabling on read-only threads
- Skip respond() when thread_id is None (matches broadcast() behavior)
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
* feat: unified thread model for web gateway
Every piece of activity (user chat, routine run, heartbeat alert, external
channel message) now lives in its own thread, properly isolated, with
meaningful titles and visual distinction.
Key changes:
- Add `channel` field to ConversationSummary and ThreadInfo so the gateway
can distinguish thread origins (gateway, telegram, routine, heartbeat).
- Add `list_conversations_all_channels` to Database trait (both postgres
and libsql) so chat_threads_handler shows cross-channel threads.
- Routine runs get a persistent conversation per routine via
`get_or_create_routine_conversation`; notifications carry thread_id.
- Heartbeat gets a persistent conversation via
`get_or_create_heartbeat_conversation`; HeartbeatRunner accepts an
optional Database store and binds notifications to the thread.
- Fix broadcast() in web gateway to propagate response.thread_id instead
of hardcoding empty string.
- Fix isCurrentThread(null) returning true (the core notification leak
bug) — now returns false so events without a thread_id don't leak into
the active thread.
- Rewrite frontend thread sidebar: meaningful titles with channel-specific
fallbacks, relative timestamps instead of turn counts, channel badges
for non-gateway threads, unread notification dots, read-only indicator
for external channel threads.
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: address PR review — TOCTOU races, stale comment, debounce, broadcast warning
- Fix TOCTOU race in get_or_create_routine_conversation (postgres):
use INSERT ON CONFLICT on new uq_conv_routine unique index + SELECT-back.
- Fix TOCTOU race in get_or_create_heartbeat_conversation (postgres):
use INSERT ON CONFLICT on new uq_conv_heartbeat unique index + SELECT-back.
- Fix TOCTOU race in get_or_create_routine_conversation (libsql):
use BEGIN IMMEDIATE transaction to serialize concurrent writers.
- Fix TOCTOU race in get_or_create_heartbeat_conversation (libsql):
use BEGIN IMMEDIATE transaction to serialize concurrent writers.
- Add V11 migration with partial unique indexes for postgres.
- Add matching unique indexes to libsql schema.
- Update stale comment on isCurrentThread (said "always shown" but logic
now returns false for missing thread_id).
- Debounce loadThreads() on off-thread SSE events to prevent request storms.
- Log warning in broadcast() when thread_id is None (clients will drop it).
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: sort in-memory thread fallback by updated_at descending
The in-memory thread list fallback (when no DB is available) used
HashMap::values() which has no guaranteed ordering. Sort by
updated_at descending to match the SQL query ordering.
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: retry libsql connect() on transient "unable to open database file"
The cron ticker's background task occasionally fails with "unable to
open database file" when creating a new SQLite connection concurrently
with the main thread. Add retry with exponential backoff (50ms, 100ms,
200ms) to handle transient VFS/locking issues in libsql's local mode.
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: use ON CONFLICT with index expressions instead of named constraints
PostgreSQL ON CONFLICT ON CONSTRAINT requires a named table constraint,
but V11 migration creates unique indexes. Switch to the expression form
(ON CONFLICT (columns) WHERE condition) which works with unique indexes.
Also fix dead code in threadTitle() where thread.title was already
checked on the previous line.
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* style: fix rustfmt chain collapse in heartbeat.rs
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: skip broadcast when thread_id is None instead of sending empty
Clients drop SSE events with empty thread_id anyway, so avoid the
unnecessary network traffic by returning early.
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* test: add libsql routine/heartbeat conversation idempotency tests
Add tests proving get_or_create_routine_conversation returns the same
conversation ID across multiple invocations with the same routine_id.
Add debug logging to routine engine to track conversation resolution.
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* feat: show "New chat" title for empty threads
- threadTitle() returns "New chat" when turn_count is 0
- Assistant thread label updates dynamically from API data
- Default HTML label changed from "Assistant" to "New chat"
- New threads naturally sort to top via last_activity DESC
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: thread sorting, routine isolation, and UI polish
- Fix libsql timestamp format mismatch causing broken thread sort order.
SQLite defaults used `datetime('now')` (space-separated) while Rust code
used RFC3339 (T-separated), breaking string-based ORDER BY. All INSERTs
now use RFC3339, and queries use `datetime()` to normalize comparison.
- Route manual routine triggers through RoutineEngine.fire_manual() instead
of injecting as regular chat messages, so routines always run in their
dedicated conversation thread.
- Add RoutineEngineSlot to GatewayState for gateway<->engine communication.
- Derive routine thread titles from conversation metadata (routine_name)
instead of showing truncated UUID hashes.
- Make chat_new_thread_handler persist to DB synchronously so loadThreads()
sees newly created threads immediately.
- Fix enableChatInput() no-op and wrong element ID in disableChatInputReadOnly().
- Fix handlers/chat.rs stale gateway-only query (use list_conversations_all_channels).
- Sort in-memory threads by DateTime before converting to RFC3339 strings.
- Trigger debouncedLoadThreads() on thinking/status SSE events for non-current
threads so routine/heartbeat threads appear in sidebar promptly.
- Remove "Threads" text from sidebar header.
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: routine history display, orphaned tool_results, duplicate system messages
Three independent fixes with regression tests:
1. Routine conversations now display in the web UI. build_turns_from_db_messages()
handles standalone assistant messages (no preceding user message) by creating
turns with empty user_input. Frontend skips empty user bubbles.
2. Worker select_tools and execute_plan paths now push an
assistant_with_tool_calls message before tool execution, preventing
sanitize_tool_messages from rewriting tool_results as orphaned user messages.
3. Reasoning::plan() and respond_with_tools() merge system messages from
context into a single system prompt instead of creating [system, system, ...]
sequences that strict LLM providers (Qwen) reject.
Also: sidebar padding/spacing improvements, wider thread panel (240px).
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: address PR nearai#607 review — RwLock held across await, missing ownership check, heartbeat config
- Clone Arc<RoutineEngine> out of RwLock before .await in trigger handler
- Add user_id ownership check to fire_manual() with NotAuthorized error
- Wire heartbeat notify_user/notify_channel from config to AgentHeartbeatConfig
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* chore: gitignore trace_*.json files and remove stale traces
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* chore: remove trace JSON files from repo
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
* fix: proper HTTP status codes for routine errors, read-only input guard, respond thread_id
- Map RoutineError::NotFound → 404, NotAuthorized → 403, Disabled → 409
- Guard enableChatInput() against re-enabling on read-only threads
- Skip respond() when thread_id is None (matches broadcast() behavior)
[skip-regression-check]
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
Summary
isCurrentThread(null)now returnsfalseinstead oftrue, so routine/heartbeat notifications no longer leak into the active threadchat_threads_handlernow queries all channels (not justgateway), so Telegram/Signal conversations appear in the sidebarthread_idon notificationsChanges
src/history/store.rschanneltoConversationSummary, add 3 new Store methodssrc/db/mod.rslist_conversations_all_channels,get_or_create_routine_conversation,get_or_create_heartbeat_conversationtoConversationStoretraitsrc/db/postgres.rssrc/db/libsql/conversations.rssrc/agent/routine_engine.rsthread_idto notificationssrc/agent/heartbeat.rsDatabasestore, persist heartbeat alerts to dedicated conversationsrc/agent/agent_loop.rsspawn_heartbeatsrc/channels/web/server.rslist_conversations_all_channels, populatechannelonThreadInfosrc/channels/web/types.rschanneltoThreadInfosrc/channels/web/mod.rsbroadcast()to propagateresponse.thread_idsrc/channels/web/static/app.jsisCurrentThread, rewrite thread sidebar, add badges/unread dots, read-only for externalsrc/channels/web/static/style.cssTest plan
cargo check --all-features— compilescargo check --no-default-features --features libsql— compilescargo clippy --all --benches --tests --examples --all-features— zero warningscargo test— all existing + 5 new tests passtest_conversation_summary_has_channel_field— ConversationSummary includes channeltest_conversation_summary_channel_various_values— channel works for all channel typestest_thread_info_channel_serialized— ThreadInfo serializes channeltest_thread_info_channel_omitted_when_none— channel omitted when Nonetest_spawn_heartbeat_accepts_store_param— spawn_heartbeat signature includes store🤖 Generated with Claude Code