fix(cron): deterministic in_channel seed + companion thread-surface seed
Two live failures from the Alice canary (2026-08-19, jobs 28a24afebd81 / 83b93f8be379), both leaving a continuable in_channel cron with amnesia: 1. Seed mirrored via origin heuristics and silently dropped the brief. _seed_cron_channel_session created the flat session row, then mirror_to_session RE-DISCOVERED the target via find_session_by_origin — whose multi-candidate bail-out returns None on a populated chat (flat session + N per-message thread sessions sharing one chat_id, mixed user_ids). Receipt: 'in_channel seed did NOT land on slack:D0BJTDCSR7C'. mirror_to_session now accepts an explicit session_id and both cron seeds pass the exact row they just created; origin-scan remains the fallback for callers that genuinely don't know the target. 2. The brief's OWN THREAD was never seeded. in_channel delivers flat, but a flat Slack message still invites a thread reply (the natural mobile affordance — exactly what the user did). That reply keys to (chat, thread=<brief ts>), which no seed touched. The delivery's message_id now anchors a companion _seed_cron_thread_session so BOTH reply surfaces (plain channel message AND in-thread reply) continue the job. Also: thread-seed failures upgraded debug→WARNING (silent seed failure IS the user-facing bug), and the thread seed reports landed/not-landed. Regression tests drive both against the live failure shapes: exact-session mirror asserted via session_id kwarg; thread companion asserted via the SendResult message_id anchor. Clean-fixture blind spot noted: the E2E harness used a fresh store with one row, which is why heuristic rediscovery looked fine pre-production.
This commit is contained in:
+48
-9
@@ -1752,6 +1752,7 @@ def _seed_cron_thread_session(
|
||||
from gateway.config import Platform
|
||||
from gateway.session import SessionSource
|
||||
|
||||
seeded_session_id: Optional[str] = None
|
||||
session_store = getattr(adapter, "_session_store", None)
|
||||
if session_store is not None:
|
||||
try:
|
||||
@@ -1780,7 +1781,11 @@ def _seed_cron_thread_session(
|
||||
)
|
||||
# Ensure the thread-keyed session row exists so the mirror has
|
||||
# a target and the user's later reply joins the same session.
|
||||
session_store.get_or_create_session(dest_source)
|
||||
# Capture the exact id — the mirror writes into THIS row, not
|
||||
# an origin-heuristic rediscovery (which bails on populated
|
||||
# chats; same class as the flat-seed live failure 2026-08-19).
|
||||
_entry = session_store.get_or_create_session(dest_source)
|
||||
seeded_session_id = getattr(_entry, "session_id", None)
|
||||
|
||||
from gateway.mirror import mirror_to_session
|
||||
|
||||
@@ -1789,7 +1794,7 @@ def _seed_cron_thread_session(
|
||||
# in-thread reply produces assistant→user→... off a phantom assistant
|
||||
# message. Pass the seed user_id so the mirror resolves the exact
|
||||
# thread-keyed session row we just created.
|
||||
mirror_to_session(
|
||||
ok = mirror_to_session(
|
||||
platform_name,
|
||||
str(chat_id),
|
||||
f"[Cron delivery: {job.get('name') or job.get('id', 'cron')}]\n{text}",
|
||||
@@ -1797,13 +1802,23 @@ def _seed_cron_thread_session(
|
||||
thread_id=str(thread_id),
|
||||
user_id="system:cron",
|
||||
role="user",
|
||||
session_id=seeded_session_id,
|
||||
)
|
||||
logger.info(
|
||||
"Job '%s': opened continuable thread %s on %s:%s and seeded the brief",
|
||||
job.get("id", "?"), thread_id, platform_name, chat_id,
|
||||
)
|
||||
if ok:
|
||||
logger.info(
|
||||
"Job '%s': opened continuable thread %s on %s:%s and seeded the brief",
|
||||
job.get("id", "?"), thread_id, platform_name, chat_id,
|
||||
)
|
||||
else:
|
||||
logger.warning(
|
||||
"Job '%s': thread seed did NOT land on %s:%s thread=%s — an "
|
||||
"in-thread reply will not see this brief",
|
||||
job.get("id", "?"), platform_name, chat_id, thread_id,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.debug(
|
||||
# WARNING, not debug: a silent seed failure IS the continuation-
|
||||
# amnesia bug (Alice 2026-08-19) — it must be visible in production.
|
||||
logger.warning(
|
||||
"Job '%s': seeding cron thread session failed for %s:%s:%s: %s",
|
||||
job.get("id", "?"), platform_name, chat_id, thread_id, e,
|
||||
)
|
||||
@@ -1862,6 +1877,7 @@ def _seed_cron_channel_session(
|
||||
|
||||
chat_type = "dm" if is_dm else "group"
|
||||
session_store = getattr(adapter, "_session_store", None)
|
||||
seeded_session_id: Optional[str] = None
|
||||
if session_store is not None:
|
||||
try:
|
||||
platform_enum = Platform(platform_name.lower())
|
||||
@@ -1877,8 +1893,13 @@ def _seed_cron_channel_session(
|
||||
thread_id=None, # flat — the whole-channel/DM session
|
||||
)
|
||||
# Create the flat session row so the mirror has a target and the
|
||||
# user's later plain reply joins the SAME session.
|
||||
session_store.get_or_create_session(dest_source)
|
||||
# user's later plain reply joins the SAME session. Capture the
|
||||
# exact session id: the mirror must write into THIS row, not
|
||||
# re-discover it via origin heuristics (which bail out on
|
||||
# populated chats where the flat session coexists with
|
||||
# per-message thread sessions — live failure, Alice 2026-08-19).
|
||||
_entry = session_store.get_or_create_session(dest_source)
|
||||
seeded_session_id = getattr(_entry, "session_id", None)
|
||||
|
||||
from gateway.mirror import mirror_to_session
|
||||
|
||||
@@ -1889,6 +1910,7 @@ def _seed_cron_channel_session(
|
||||
source_label="cron",
|
||||
thread_id=None,
|
||||
user_id=str(user_id) if user_id else None,
|
||||
session_id=seeded_session_id,
|
||||
role="user",
|
||||
)
|
||||
if ok:
|
||||
@@ -2912,6 +2934,7 @@ def _deliver_result(job: dict, content: str, adapters=None, loop=None) -> Option
|
||||
text_to_send = cleaned_delivery_content.strip()
|
||||
adapter_ok = True
|
||||
timed_out = False
|
||||
delivered_message_id = None
|
||||
if text_to_send:
|
||||
from agent.async_utils import safe_schedule_threadsafe
|
||||
|
||||
@@ -3009,9 +3032,11 @@ def _deliver_result(job: dict, content: str, adapters=None, loop=None) -> Option
|
||||
if isinstance(send_result, dict):
|
||||
send_success = bool(send_result.get("success", False))
|
||||
send_raw_response = send_result.get("raw_response")
|
||||
delivered_message_id = send_result.get("message_id")
|
||||
else:
|
||||
send_success = _confirm_adapter_delivery(send_result)
|
||||
send_raw_response = getattr(send_result, "raw_response", None)
|
||||
delivered_message_id = getattr(send_result, "message_id", None)
|
||||
|
||||
if not send_success:
|
||||
if isinstance(send_result, dict):
|
||||
@@ -3127,6 +3152,20 @@ def _deliver_result(job: dict, content: str, adapters=None, loop=None) -> Option
|
||||
"— a plain reply will not see this brief",
|
||||
job["id"], platform_name, chat_id,
|
||||
)
|
||||
# Companion THREAD-surface seed (live gap, Alice
|
||||
# 2026-08-19): a flat brief is still a Slack message
|
||||
# the user can reply to IN ITS THREAD — the natural
|
||||
# mobile/desktop affordance — and that reply keys to
|
||||
# (chat, thread=<brief ts>), a session the flat seed
|
||||
# never touches. Seed it too so BOTH reply surfaces
|
||||
# continue the job. Uses the delivered message id as
|
||||
# the thread anchor; best-effort like every seed.
|
||||
if delivered_message_id:
|
||||
_seed_cron_thread_session(
|
||||
job, runtime_adapter, platform_name, chat_id,
|
||||
str(delivered_message_id), mirror_text,
|
||||
chat_name=origin.get("chat_name"),
|
||||
)
|
||||
elif in_channel_surface and not origin_target:
|
||||
logger.warning(
|
||||
"Job '%s': in_channel delivery to %s:%s is not the "
|
||||
|
||||
+19
-6
@@ -30,6 +30,7 @@ def mirror_to_session(
|
||||
thread_id: Optional[str] = None,
|
||||
user_id: Optional[str] = None,
|
||||
role: str = "assistant",
|
||||
session_id: Optional[str] = None,
|
||||
) -> bool:
|
||||
"""
|
||||
Append a delivery-mirror message to the target session's transcript.
|
||||
@@ -37,6 +38,17 @@ def mirror_to_session(
|
||||
Finds the gateway session that matches the given platform + chat_id,
|
||||
then writes a mirror entry to both the JSONL transcript and SQLite DB.
|
||||
|
||||
``session_id``: when the caller already KNOWS the exact session (e.g. the
|
||||
cron in_channel seed, which just created the row via
|
||||
``get_or_create_session``), pass it to skip the origin-scan heuristics
|
||||
entirely. ``_find_session_id`` matches by origin (chat_id + user
|
||||
preference with a multi-candidate bail-out), which is correct for
|
||||
"mirror into whatever conversation lives here" callers but WRONG for a
|
||||
caller holding the precise target — on a populated chat (flat session +
|
||||
N per-message thread sessions sharing one chat_id) the scan can refuse
|
||||
to guess and silently drop the mirror (live failure, Alice 2026-08-19:
|
||||
'in_channel seed did NOT land').
|
||||
|
||||
``role`` defaults to ``"assistant"`` — correct for the interactive
|
||||
``send_message`` mirror, where the mirrored text is the agent's own
|
||||
outgoing reply (a genuine assistant turn). Callers mirroring text that is
|
||||
@@ -52,12 +64,13 @@ def mirror_to_session(
|
||||
All errors are caught -- this is never fatal.
|
||||
"""
|
||||
try:
|
||||
session_id = _find_session_id(
|
||||
platform,
|
||||
str(chat_id),
|
||||
thread_id=thread_id,
|
||||
user_id=user_id,
|
||||
)
|
||||
if not session_id:
|
||||
session_id = _find_session_id(
|
||||
platform,
|
||||
str(chat_id),
|
||||
thread_id=thread_id,
|
||||
user_id=user_id,
|
||||
)
|
||||
if not session_id:
|
||||
logger.debug(
|
||||
"Mirror: no session found for %s:%s:%s:%s",
|
||||
|
||||
@@ -2385,7 +2385,9 @@ class TestCronContinuableSurfaceInChannel:
|
||||
|
||||
def _slack_adapter(self, supports_inchannel=True, with_store=True):
|
||||
adapter = AsyncMock()
|
||||
adapter.send.return_value = MagicMock(success=True)
|
||||
adapter.send.return_value = MagicMock(
|
||||
success=True, message_id="msg_1", raw_response=None,
|
||||
)
|
||||
# Capability flag read via getattr in the scheduler.
|
||||
adapter.supports_inchannel_continuable = supports_inchannel
|
||||
# A live session store so the in_channel seed can CREATE the flat row
|
||||
@@ -2466,6 +2468,50 @@ class TestCronContinuableSurfaceInChannel:
|
||||
# session key includes it on per-user-isolated chats.
|
||||
assert seed_mock.call_args.kwargs.get("user_id") == "U_HUMAN"
|
||||
|
||||
def test_seed_mirrors_into_exact_created_session_on_populated_chat(self):
|
||||
"""REGRESSION (live, Alice 2026-08-19 19:17): 'in_channel seed did NOT
|
||||
land'. The seed created the flat session row, then mirror_to_session
|
||||
re-discovered the target via origin heuristics — and on a populated
|
||||
chat (flat session + N per-message thread sessions sharing chat_id,
|
||||
mixed user_ids) find_session_by_origin's multi-candidate bail-out
|
||||
returned None, silently dropping the brief. The seed must mirror into
|
||||
the EXACT session row it just created, no rediscovery."""
|
||||
from cron.scheduler import _seed_cron_channel_session
|
||||
|
||||
store = MagicMock()
|
||||
created = MagicMock()
|
||||
created.session_id = "sess-flat-exact"
|
||||
store.get_or_create_session.return_value = created
|
||||
adapter = MagicMock()
|
||||
adapter._session_store = store
|
||||
|
||||
with patch("gateway.mirror.mirror_to_session", return_value=True) as mirror_mock:
|
||||
ok = _seed_cron_channel_session(
|
||||
{"id": "j1", "name": "Brief"}, adapter, "slack", "D0BJTDCSR7C",
|
||||
"Daily brief", is_dm=True, user_id="U_HUMAN", chat_name=None,
|
||||
)
|
||||
assert ok is True
|
||||
# The mirror received the exact created session id — origin-scan
|
||||
# heuristics (and their populated-chat bail-out) are out of the path.
|
||||
assert mirror_mock.call_args.kwargs.get("session_id") == "sess-flat-exact"
|
||||
|
||||
def test_in_channel_also_seeds_thread_surface_of_delivered_brief(self):
|
||||
"""REGRESSION (live, Alice 2026-08-19 19:19): the user replied IN THE
|
||||
BRIEF'S THREAD (Slack's natural affordance on a flat message). That
|
||||
reply keys to (chat, thread=<brief ts>) — a session in_channel mode
|
||||
never seeded, so the agent had no idea about its own brief. The flat
|
||||
delivery's message_id must anchor a companion thread-surface seed."""
|
||||
adapter = self._slack_adapter(supports_inchannel=True)
|
||||
with patch("cron.scheduler._seed_cron_channel_session", return_value=True), \
|
||||
patch("cron.scheduler._seed_cron_thread_session") as thread_seed_mock:
|
||||
self._run_inchannel_delivery(
|
||||
{"slack": {"cron_continuable_surface": "in_channel"}}, adapter,
|
||||
attach_to_session=False,
|
||||
)
|
||||
thread_seed_mock.assert_called_once()
|
||||
# Anchored on the delivered message id (the router's SendResult).
|
||||
assert thread_seed_mock.call_args.args[4] == "msg_1"
|
||||
|
||||
|
||||
class TestMultiTargetDeliveryContinuesOnFailure:
|
||||
"""When delivery to one target fails inside the standalone thread-pool
|
||||
|
||||
Reference in New Issue
Block a user