From 4d1b1e63be63f45d970b6ddfdca8dd16d6671480 Mon Sep 17 00:00:00 2001 From: Daniele Date: Mon, 10 Aug 2026 22:13:32 +0100 Subject: [PATCH] fix(async-tasks): wake the parent conversation by id, not by source MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit DurableSink resumed the parent through ChatHub::resume(&source), which resolves whatever session the source currently points at. Since one source can now carry several conversations (secondary tabs, a reset since the task started), that pointer is no longer the conversation the result was delivered into: the recovery ran on the wrong one, found nothing pending, and returned silently — the delivered task_completed sat unread until the user's next message drove a normal turn. The sink already knows the parent session id, so resume it directly through resume_for_session, keeping the same in-flight guard. The source lookup and the now-unused pool field go with it. Adds a crate-level regression test: a completed turn, a StoreSink delivery, then a recovery — the result must drive a new round. --- crates/agent-loop/tests/recovery.rs | 28 +++++++++++++++++++ .../src/loop_adapters/async_task.rs | 17 ++++++----- 2 files changed, 36 insertions(+), 9 deletions(-) diff --git a/crates/agent-loop/tests/recovery.rs b/crates/agent-loop/tests/recovery.rs index 6455bbd..9bf3868 100644 --- a/crates/agent-loop/tests/recovery.rs +++ b/crates/agent-loop/tests/recovery.rs @@ -379,6 +379,34 @@ async fn an_interrupted_parallel_batch_is_reaped_and_the_parent_resumes() { assert_eq!(report.frames_resumed, 1, "the root continues with the failures in view"); } +// ── async result wake-up (reproduction) ────────────────────────────────────── + +#[tokio::test] +async fn an_idle_conversation_woken_by_an_async_result_continues() { + use agent_loop::delegate::{AsyncResultSink, CompletedTask, StoreSink}; + use agent_loop::ids::TaskId; + + let h = H::new(vec![Step::message("processing the task result")], vec![]).await; + + // The parent's turn is complete: user message, final assistant reply. + h.store.append(h.root, NewMessage::user("start a task")).await.unwrap(); + h.store.append(h.root, NewMessage::assistant("started, I'll let you know", None)).await.unwrap(); + + // The task finishes: the sink writes the synthetic delivery, then the host + // wakes the conversation with a recovery. + let sink = StoreSink::new(h.store.clone()); + sink.deliver(h.conv.clone(), CompletedTask { + id: TaskId(7), + title: "research".into(), + result: "the answer is 42".into(), + }) + .await + .unwrap(); + + let report = h.recover().await; + assert_eq!(report.frames_resumed, 1, "the delivered result must drive a new round"); +} + // ── resolve_pending ────────────────────────────────────────────────────────── #[tokio::test] diff --git a/crates/skald-core/src/loop_adapters/async_task.rs b/crates/skald-core/src/loop_adapters/async_task.rs index 1a145dd..22cd997 100644 --- a/crates/skald-core/src/loop_adapters/async_task.rs +++ b/crates/skald-core/src/loop_adapters/async_task.rs @@ -99,17 +99,20 @@ impl AsyncExecutor for CronExecutor { /// /// `ChatHub::resume` skips a session with a turn already in flight, which is the /// right rule here too: a live loop reads the store each round and picks the -/// result up on its own. +/// result up on its own. The wake-up addresses the parent **by session id**, +/// never by source: one source may now carry several conversations (secondary +/// tabs) or have moved to a fresh one since the task started, and resuming the +/// source's active session would run the recovery on the wrong conversation — +/// a silent no-op there, while this result sat unread until the next message. pub struct DurableSink { inner: StoreSink, - pool: Arc, hub: Arc, } impl DurableSink { pub fn new(pool: Arc, hub: Arc) -> Self { - let store: Arc = Arc::new(SqliteHistory::new(pool.clone())); - Self { inner: StoreSink::new(store), pool, hub } + let store: Arc = Arc::new(SqliteHistory::new(pool)); + Self { inner: StoreSink::new(store), hub } } } @@ -119,10 +122,6 @@ impl AsyncResultSink for DurableSink { self.inner.deliver(parent.clone(), task).await?; let session_id = SqliteHistory::session_id(&parent)?; - let source = crate::db::chat_sessions::find_by_id(&self.pool, session_id) - .await? - .map(|s| s.source) - .ok_or_else(|| anyhow::anyhow!("deliver: session {session_id} not found"))?; - self.hub.resume(&source).await + self.hub.resume_for_session(session_id).await } }