ROB-3670 - Add M2 Conversation Worker for async conversation processing - #1903
Conversation
Implements M2 of the Conversation History Consolidation project. Holmes now actively claims Conversation rows from Supabase and publishes results back to ConversationEvents in real-time, instead of passively responding over HTTP/SSE. New module holmes/core/conversations_worker/ contains: - worker.py: main orchestrator. Claim loop + per-conversation threads that run the existing call_stream() pipeline and publish events back. - event_publisher.py: batches StreamMessage events into ConversationEvents rows. Flushes immediately on terminal events (ai_answer_end, approval_required, error) and passes compact=True on history compaction. - realtime_manager.py: optional background asyncio thread that joins the per-cluster Presence channel and subscribes to Postgres Changes on the Conversations table. Failure to connect is non-fatal; the worker falls back to polling. - models.py: ConversationTask, ConversationReassignedError. DAL additions (supabase_dal.py): - claim_conversations(holmes_id) - post_conversation_events(...) - complete_conversation(...) - get_conversation_events(...) - get_pending_conversations() (polling fallback helper) New env vars (env_vars.py): - ENABLE_CONVERSATION_WORKER (default False) - CONVERSATION_WORKER_MAX_CONCURRENT - CONVERSATION_WORKER_POLL_INTERVAL_SECONDS - CONVERSATION_WORKER_HEARTBEAT_INTERVAL_SECONDS - CONVERSATION_WORKER_EVENT_BATCH_INTERVAL_SECONDS - CONVERSATION_WORKER_REALTIME_RECONNECT_MAX_SECONDS - CONVERSATION_WORKER_REALTIME_ENABLED Wiring in server.py: the worker is instantiated alongside the existing scheduled prompts executor and gated behind ENABLE_CONVERSATION_WORKER. The legacy /api/chat endpoint is untouched and continues to work. Hydration logic rebuilds conversation_history for follow-up turns by scanning prior ConversationEvents for the last ai_answer_end or approval_required event and reusing its messages array. Verified end-to-end against the staging Supabase instance with the Elasticsearch toolset: single-turn and multi-turn (post_conversation_followup) conversations claim, stream events, and complete correctly. Unit tests cover event batching, terminal-event flushing, compaction flag, and hydration from prior events. Signed-off-by: Claude <noreply@anthropic.com>
The AsyncRealtimeClient from realtime-py calls websockets.connect() directly, which ignores https_proxy env vars and attempts direct DNS/TCP. In environments where all egress must go through an HTTP CONNECT proxy (e.g. Claude Code's sandbox), this causes a DNS resolution failure while the regular HTTPS REST calls to the same host succeed. Install a module-level monkey-patch on realtime._async.client.connect when https_proxy is set: the patch opens an HTTP CONNECT tunnel via python-socks, wraps in SSL for wss://, and hands the pre-connected socket to websockets.connect(..., sock=...). Also: the initial WebSocket connection must use the anon apikey, not the user JWT. Supabase Realtime uses the URL ?apikey=... for server auth and then set_auth() for RLS. The previous code used the user JWT for both, causing HTTP 401. Verified end-to-end: - RealtimeManager connected and subscribed to cluster channel - New conversation insert triggered a Postgres change notification within 142ms of creation - Holmes claimed and processed the conversation via the realtime path (poll interval was 300s; claim happened in well under 1s) If python-socks isn't installed, a warning is logged and the manager falls back to a direct connection attempt (which the outer worker already handles gracefully via polling). Signed-off-by: Claude <noreply@anthropic.com>
…ypes Followup turns that carry only tool_decisions / frontend_tool_results (no new user ask) previously failed with "has no user question". The /api/chat endpoint requires ask because it always appends a new user message; but on resume call_stream consumes tool_decisions against the existing history without needing a new ask. Fix: - Detect resume-only followups (tool_decisions/frontend_tool_results present without ask) and skip build_chat_messages; pass conversation_history through as-is. - Add _extract_last_user_ask helper to pull a placeholder ask from history for ChatRequest validation (ask: str is required by the Pydantic model). End-to-end verification against staging Supabase with bash toolset + enable_tool_approval=true: - Turn 1 events: user_message, token_count, ai_message, start_tool_calling, tool_calling_result (status=approval_required), token_count, approval_required (terminal) - Turn 2 events (post_conversation_followup with tool_decisions): user_message, tool_calling_result (approved execution), token_count, ai_answer_end (terminal) - Final answer correctly echoes the bash command output Also verified stop_conversation triggers ConversationReassignedError in the publisher, which the worker catches and logs — leaving the row in 'stopped' state without overwriting it. New unit tests: - test_publisher_flushes_on_error_event: ERROR is a terminal event - test_publisher_covers_all_stream_event_types: smoke test that every StreamEvents value round-trips through the publisher - test_extract_last_user_ask: helper handles plain + vision messages Signed-off-by: Claude <noreply@anthropic.com>
The original Presence implementation had two bugs that prevented observers
from seeing Holmes:
1. The channel was created without a presence config. Supabase Realtime
only broadcasts presence state to observers when the joining client
includes {config: {presence: {enabled: true, key: ...}}} in the channel
options. Without it, track() calls were accepted silently but the state
never reached other subscribers.
2. track() was being called before the channel finished SUBSCRIBING. The
message was sent on the socket but dropped by the server.
3. _leave_conversation_channel looked up channels by bare topic, but
realtime-py stores them under the "realtime:<topic>" key, so leave was
a no-op.
4. The worker never actually called join_conversation_presence /
leave_conversation_presence — step 4 of the M2 spec was unimplemented.
Changes:
- Enable presence config on both the cluster and per-conversation
channels and wait for SUBSCRIBE ack before calling track().
- Lookup channels by the "realtime:" prefix in _leave_conversation_channel.
- Use untrack() + remove_channel() for a clean LEAVE.
- Worker._process_conversation calls join_conversation_presence at the
start of a claim; _process_conversation_safe's finally block calls
leave_conversation_presence on any exit path (success, failure,
reassignment).
Verified end-to-end against staging Supabase with a separate observer:
Cluster channel (holmes:cluster:{account}:{cluster}):
state = {
"<holmes_id>": [{
"presence_ref": "...",
"holmes_id": "<pid>",
"version": "dev-unknown",
"started_at": "2026-04-13T07:...",
"active_conversations": 0
}]
}
Per-conversation channel (holmes:conversation:{id}):
JOIN (on claim): full metadata + conversation_id field
LEAVE (on done): same presence_ref removed
Signed-off-by: Claude <noreply@anthropic.com>
There was a problem hiding this comment.
Claude Code Review
This repository is configured for manual code reviews. Comment @claude review to trigger a review and subscribe this PR to future pushes, or @claude review once for a one-time review.
Tip: disable this comment in your organization's Code Review settings.
📂 Previous Runs📜 #5 · Run @ __8755b3d__ (#25051477957) — Apr 28, 12:07 UTC✅ Results of HolmesGPT evalsAutomatically triggered by commit 8755b3d on branch Results of HolmesGPT evals
Benchmark Comparison DetailsBaseline: latest ci-benchmark experiment on master Status: Success - 74 test/model combinations loaded Benchmark experiment:
No benchmark data available for comparison. Benchmark has no cost, total tokens, cached tokens data. Will appear after the next weekly benchmark run. Comparison indicators:
📜 #4 · Run @ __52a1aca__ (#24996390540) — Apr 27, 13:05 UTC✅ Results of HolmesGPT evalsAutomatically triggered by commit 52a1aca on branch Results of HolmesGPT evals
Benchmark Comparison DetailsBaseline: latest ci-benchmark experiment on master Status: Success - 74 test/model combinations loaded Benchmark experiment:
No benchmark data available for comparison. Benchmark has no cost, total tokens, cached tokens data. Will appear after the next weekly benchmark run. Comparison indicators:
📜 #3 · Run @ __910fa57__ (#24840511591) — Apr 23, 14:46 UTC✅ Results of HolmesGPT evalsAutomatically triggered by commit 910fa57 on branch Results of HolmesGPT evals
Benchmark Comparison DetailsBaseline: latest ci-benchmark experiment on master Status: Success - 34 test/model combinations loaded Benchmark experiment:
No benchmark data available for comparison. Benchmark has no cost, total tokens, cached tokens data. Will appear after the next weekly benchmark run. Comparison indicators:
📜 #2 · Run @ __afd11cd__ (#24839631399) — Apr 23, 14:11 UTC✅ Results of HolmesGPT evalsAutomatically triggered by commit afd11cd on branch Results of HolmesGPT evals
Benchmark Comparison DetailsBaseline: latest ci-benchmark experiment on master Status: Success - 34 test/model combinations loaded Benchmark experiment:
No benchmark data available for comparison. Benchmark has no cost, total tokens, cached tokens data. Will appear after the next weekly benchmark run. Comparison indicators:
📜 #1 · Run @ __58d7fe6__ (#24833811409) — Apr 23, 12:03 UTC✅ Results of HolmesGPT evalsAutomatically triggered by commit 58d7fe6 on branch Results of HolmesGPT evals
Benchmark Comparison DetailsBaseline: latest ci-benchmark experiment on master Status: Success - 34 test/model combinations loaded Benchmark experiment:
No benchmark data available for comparison. Benchmark has no cost, total tokens, cached tokens data. Will appear after the next weekly benchmark run. Comparison indicators:
✅ Results of HolmesGPT evalsAutomatically triggered by commit 0ea081c on branch Results of HolmesGPT evals
Benchmark Comparison DetailsBaseline: latest ci-benchmark experiment on master Status: Success - 74 test/model combinations loaded Benchmark experiment:
No benchmark data available for comparison. Benchmark has no cost, total tokens, cached tokens data. Will appear after the next weekly benchmark run. Comparison indicators:
📖 Legend
🔄 Re-run evals manually
Option 1: Comment on this PR with Or with more options (one per line): Run evals on a different branch (e.g., master) for comparison:
Quick re-run: Use Option 2: Trigger via GitHub Actions UI → "Run workflow" Option 3: Add PR labels to include extra evals (applies to both automatic runs and
Examples: 🏷️ Valid tags
🤖 Valid models
Commands: CLI: |
|
✅ Docker images ready for
Use these tags to pull the images for testing. 📋 Copy commandsgcloud auth configure-docker us-central1-docker.pkg.dev
docker pull us-central1-docker.pkg.dev/robusta-development/temporary-builds/holmes:14f8fed1
docker tag us-central1-docker.pkg.dev/robusta-development/temporary-builds/holmes:14f8fed1 me-west1-docker.pkg.dev/robusta-development/development/holmes-dev:14f8fed1
docker push me-west1-docker.pkg.dev/robusta-development/development/holmes-dev:14f8fed1
docker pull us-central1-docker.pkg.dev/robusta-development/temporary-builds/holmes-operator:14f8fed1
docker tag us-central1-docker.pkg.dev/robusta-development/temporary-builds/holmes-operator:14f8fed1 me-west1-docker.pkg.dev/robusta-development/development/holmes-operator-dev:14f8fed1
docker push me-west1-docker.pkg.dev/robusta-development/development/holmes-operator-dev:14f8fed1Patch Helm values in one line (choose the chart you use): HolmesGPT chart: helm upgrade --install holmesgpt ./helm/holmes \
--set registry=me-west1-docker.pkg.dev/robusta-development/development \
--set image=holmes-dev:14f8fed1 \
--set operator.registry=me-west1-docker.pkg.dev/robusta-development/development \
--set operator.image=holmes-operator-dev:14f8fed1Robusta wrapper chart: helm upgrade --install robusta robusta/robusta \
--reuse-values \
--set holmes.registry=me-west1-docker.pkg.dev/robusta-development/development \
--set holmes.image=holmes-dev:14f8fed1 \
--set holmes.operator.registry=me-west1-docker.pkg.dev/robusta-development/development \
--set holmes.operator.image=holmes-operator-dev:14f8fed1 |
|
Note Reviews pausedIt looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
WalkthroughAdds a Conversation Worker subsystem: env var configuration, Pydantic models, Supabase DAL RPCs, a realtime subscription manager, an event publisher, a threaded worker integrated into server startup, Holmes status metadata flags, unit and integration tests, and a pytest marker. Changes
Sequence Diagram(s)sequenceDiagram
participant RM as RealtimeManager
participant W as ConversationWorker
participant DAL as SupabaseDal
participant Chat as ChatPipeline
participant EP as EventPublisher
RM->>W: on_new_pending()
W->>DAL: claim_conversations(holmes_id)
DAL-->>W: claimed rows
loop per claimed conversation
W->>DAL: get_conversation_events(conversation_id)
DAL-->>W: prior events
W->>W: hydrate ConversationTask
W->>DAL: update_conversation_status(..., "running")
DAL-->>W: success / raises mismatch
W->>Chat: call_stream(chat_request) (streaming)
Chat-->>EP: stream StreamMessage events
EP->>DAL: post_conversation_events(batch, compact)
DAL-->>EP: seq int / None / raises mismatch
EP-->>W: terminal StreamEvents or None
W->>DAL: update_conversation_status(..., terminal_status)
DAL-->>W: success / raises mismatch
end
Estimated code review effort🎯 5 (Critical) | ⏱️ ~120 minutes Possibly related PRs
Suggested reviewers
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
✅ Deploy Preview for holmes-docs ready!
To edit notification comments on pull requests, go to your Netlify project configuration. |
🔬 CLI Performance Benchmark🔴 Startup Time (no LLM)Measures
🟡 Full CLI with LLMMeasures
PR: |
There was a problem hiding this comment.
Actionable comments posted: 7
🧹 Nitpick comments (4)
tests/core/conversations_worker/test_worker_hydration.py (1)
126-137: Remove redundant import inside test function.
ConversationWorkeris already imported at line 3, so the import at line 128 is redundant. This also aligns with the coding guidelines requiring imports at the top of the file.♻️ Proposed fix
def test_extract_last_user_ask(): """_extract_last_user_ask walks a message history and returns the last user text.""" - from holmes.core.conversations_worker.worker import ConversationWorker - history = [As per coding guidelines: "ALWAYS place Python imports at the top of the file, not inside functions or methods"
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/core/conversations_worker/test_worker_hydration.py` around lines 126 - 137, Remove the redundant local import inside test_extract_last_user_ask: delete the line importing ConversationWorker from holmes.core.conversations_worker.worker and let the test use the ConversationWorker already imported at the top of the file; the target symbol to update is the test function test_extract_last_user_ask and the reference to ConversationWorker._extract_last_user_ask so the test still asserts the same behavior after the inner import is removed.server.py (1)
547-553: Move import to top of file per coding guidelines.The import of
ConversationWorkeris placed inside the conditional block. While this is a common pattern for lazy loading, it violates the coding guidelines that require imports to be placed at the top of the file.♻️ Proposed fix
Move the import to the top of the file with other imports (after line 63):
from holmes.core.conversations_worker import ConversationWorkerThen simplify the conditional instantiation:
conversation_worker = None if ENABLE_CONVERSATION_WORKER: - from holmes.core.conversations_worker import ConversationWorker - conversation_worker = ConversationWorker( dal=dal, config=config, chat_function=chat )As per coding guidelines: "ALWAYS place Python imports at the top of the file, not inside functions or methods"
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@server.py` around lines 547 - 553, The import of ConversationWorker must be moved out of the conditional and placed with the other top-of-file imports (add "from holmes.core.conversations_worker import ConversationWorker" alongside other imports), then simplify the block that instantiates conversation_worker so it only checks ENABLE_CONVERSATION_WORKER and creates ConversationWorker(dal=dal, config=config, chat_function=chat) when enabled (leave conversation_worker = None otherwise); reference ConversationWorker, ENABLE_CONVERSATION_WORKER, and the conversation_worker variable when making the change.holmes/core/conversations_worker/event_publisher.py (1)
122-130: String-based error detection is fragile.Detecting reassignment errors by parsing exception message strings (e.g.,
"assignee mismatch","request sequence mismatch") is brittle and could break if the RPC error messages change. Consider defining a custom exception type in the DAL layer that wraps these specific error conditions, or at minimum, document this coupling so future DAL changes don't silently break this detection logic.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/event_publisher.py` around lines 122 - 130, The current catch block in event_publisher that inspects exception message strings (the variable e) and looks for "assignee mismatch"/"request sequence mismatch"/"is not running" is brittle; update the DAL/RPC layer to raise a specific custom exception (e.g., ReassignmentError or AssigneeMismatchError) for these conditions and then change this handler to catch that concrete exception type (use isinstance(e, ReassignmentError) or add an except ReassignmentError branch) and re-raise ConversationReassignedError accordingly; if changing the DAL isn't possible immediately, add a short comment documenting the coupling and create an internal wrapper exception in the RPC client that translates those message patterns into the new custom exception so event_publisher can rely on type-based detection instead of string matching.holmes/core/conversations_worker/realtime_manager.py (1)
69-84: Star-arg unpacking after keyword argument (B026).The
ws_connect(url, sock=sock, *args, **kwargs)pattern at line 84 is flagged by static analysis. While it works, it's discouraged because the behavior is confusing whenargscontains positional arguments that could conflict with keyword arguments.♻️ Proposed fix - remove *args since kwargs should capture all needed parameters
- async def _proxied_connect(url: str, *args, **kwargs): + async def _proxied_connect(url: str, **kwargs): parsed = urllib.parse.urlparse(url) if parsed.scheme not in ("ws", "wss"): - return await ws_connect(url, *args, **kwargs) + return await ws_connect(url, **kwargs) # skip proxy for localhost targets if parsed.hostname in ("localhost", "127.0.0.1"): - return await ws_connect(url, *args, **kwargs) + return await ws_connect(url, **kwargs) port = parsed.port or (443 if parsed.scheme == "wss" else 80) proxy = Proxy.from_url(proxy_connect_url) sock = await proxy.connect(dest_host=parsed.hostname, dest_port=port) kwargs.setdefault("server_hostname", parsed.hostname) if parsed.scheme == "wss" and "ssl" not in kwargs: kwargs["ssl"] = ssl.create_default_context() - return await ws_connect(url, sock=sock, *args, **kwargs) + return await ws_connect(url, sock=sock, **kwargs)🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 69 - 84, The call pattern using ws_connect(url, sock=sock, *args, **kwargs) in _proxied_connect risks positional/keyword conflicts; remove the star-arg unpacking after keyword arguments by passing only keyword args to ws_connect (e.g., call ws_connect(url, sock=sock, **kwargs)) and likewise replace other ws_connect(url, *args, **kwargs) calls in _proxied_connect with ws_connect(url, **kwargs), ensuring any needed positional parameters are moved into kwargs before calling; see _proxied_connect, ws_connect usage and the proxy setup (proxy_connect_url / Proxy.from_url) to locate and update the calls.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@holmes/core/conversations_worker/event_publisher.py`:
- Around line 48-77: The return type of consume does not allow None even though
the method (and docstring) can return None via self._last_terminal_event; update
the signature of consume to return Optional[StreamEvents] (import Optional from
typing), remove the trailing "# type: ignore[return-value]" and ensure
_last_terminal_event is typed (or annotated where set) as Optional[StreamEvents]
so static checkers accept returning None from consume; keep behavior unchanged.
In `@holmes/core/conversations_worker/models.py`:
- Line 29: Update the conversation_history field to use a concrete typed generic
instead of bare list: change conversation_history: Optional[list] = None to
conversation_history: Optional[List[Dict[str, Any]]] (or the appropriate
dict-like message type) and ensure typing imports (List, Dict, Any) are present
in models.py so type hints match usage in worker.py; keep the default None
unchanged.
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 251-255: The try/except around self.on_new_pending() in
RealtimeManager currently swallows all exceptions; change it to catch Exception
as e and log the error (including stack trace) before continuing so subscription
setup isn't disrupted but failures are visible; use the instance logger (e.g.,
self.logger.exception(...) or self.logger.error(..., exc_info=True)) in the
except block and keep suppressing the exception after logging.
In `@holmes/core/conversations_worker/worker.py`:
- Around line 190-192: The current logging.exception call logs the entire conv
payload (variable conv) which may expose sensitive data; change the log to only
include stable identifiers (e.g., conv.get("id") or conv.get("conversation_id")
and any non-sensitive, masked user id) and remove the raw payload from the
message—so in the exception handler around logging.exception replace the conv
interpolation with those identifier values and keep exc_info=True to preserve
the stack trace; locate the call to logging.exception in this function and
construct a small context dict or formatted string with only the safe
identifiers instead of the full conv object.
- Around line 378-380: The loop that does for msg in reversed(history): assumes
each entry is a dict and directly calls msg.get(...), which can raise if a
malformed non-dict entry appears; modify the loop in the worker to first verify
the item is a dict (e.g., if not isinstance(msg, dict): continue) before
accessing msg.get("role") and msg.get("content"), and optionally guard that
content is a str/non-empty before using it so malformed history entries are
skipped rather than causing an exception.
- Around line 154-175: The _try_claim_and_dispatch method can over-claim because
dal.claim_conversations currently requests all pending conversations; update the
DAL method signature claim_conversations(self, holmes_id, limit: Optional[int] =
None) to accept an optional limit and have it pass that limit into the Supabase
RPC, then in _try_claim_and_dispatch compute remaining_slots =
CONVERSATION_WORKER_MAX_CONCURRENT - len(self._active_conversation_ids) and call
self.dal.claim_conversations(self.holmes_id, limit=remaining_slots) so the RPC
returns at most the worker’s available capacity before building tasks and adding
ids to _active_conversation_ids.
In `@holmes/core/supabase_dal.py`:
- Around line 1007-1031: get_conversation_events is missing an account_id filter
which allows queries to be scoped only by conversation_id; update the signature
of get_conversation_events to accept account_id (e.g., def
get_conversation_events(self, conversation_id: str, account_id: str, min_seq:
int = 0) or add account_id as a required/optional param consistent with callers)
and add an .eq("account_id", account_id) clause to the Supabase query chain
before executing; mirror the approach used in get_pending_conversations for
consistency and ensure callers are updated to pass account_id.
---
Nitpick comments:
In `@holmes/core/conversations_worker/event_publisher.py`:
- Around line 122-130: The current catch block in event_publisher that inspects
exception message strings (the variable e) and looks for "assignee
mismatch"/"request sequence mismatch"/"is not running" is brittle; update the
DAL/RPC layer to raise a specific custom exception (e.g., ReassignmentError or
AssigneeMismatchError) for these conditions and then change this handler to
catch that concrete exception type (use isinstance(e, ReassignmentError) or add
an except ReassignmentError branch) and re-raise ConversationReassignedError
accordingly; if changing the DAL isn't possible immediately, add a short comment
documenting the coupling and create an internal wrapper exception in the RPC
client that translates those message patterns into the new custom exception so
event_publisher can rely on type-based detection instead of string matching.
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 69-84: The call pattern using ws_connect(url, sock=sock, *args,
**kwargs) in _proxied_connect risks positional/keyword conflicts; remove the
star-arg unpacking after keyword arguments by passing only keyword args to
ws_connect (e.g., call ws_connect(url, sock=sock, **kwargs)) and likewise
replace other ws_connect(url, *args, **kwargs) calls in _proxied_connect with
ws_connect(url, **kwargs), ensuring any needed positional parameters are moved
into kwargs before calling; see _proxied_connect, ws_connect usage and the proxy
setup (proxy_connect_url / Proxy.from_url) to locate and update the calls.
In `@server.py`:
- Around line 547-553: The import of ConversationWorker must be moved out of the
conditional and placed with the other top-of-file imports (add "from
holmes.core.conversations_worker import ConversationWorker" alongside other
imports), then simplify the block that instantiates conversation_worker so it
only checks ENABLE_CONVERSATION_WORKER and creates ConversationWorker(dal=dal,
config=config, chat_function=chat) when enabled (leave conversation_worker =
None otherwise); reference ConversationWorker, ENABLE_CONVERSATION_WORKER, and
the conversation_worker variable when making the change.
In `@tests/core/conversations_worker/test_worker_hydration.py`:
- Around line 126-137: Remove the redundant local import inside
test_extract_last_user_ask: delete the line importing ConversationWorker from
holmes.core.conversations_worker.worker and let the test use the
ConversationWorker already imported at the top of the file; the target symbol to
update is the test function test_extract_last_user_ask and the reference to
ConversationWorker._extract_last_user_ask so the test still asserts the same
behavior after the inner import is removed.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 2e0772cd-6ada-425b-9c02-ee39da42c7f4
📒 Files selected for processing (11)
holmes/common/env_vars.pyholmes/core/conversations_worker/__init__.pyholmes/core/conversations_worker/event_publisher.pyholmes/core/conversations_worker/models.pyholmes/core/conversations_worker/realtime_manager.pyholmes/core/conversations_worker/worker.pyholmes/core/supabase_dal.pyserver.pytests/core/conversations_worker/__init__.pytests/core/conversations_worker/test_event_publisher.pytests/core/conversations_worker/test_worker_hydration.py
Addresses review feedback:
1. Batch interval default 0.5s → 1s
2. Heartbeat interval default 10s → 5s
3. Rename CONVERSATION_WORKER_POLL_INTERVAL_SECONDS to
CONVERSATION_WORKER_POLL_INTERVAL_SECONDS_WITHOUT_REALTIME (default 60s).
The worker only polls when realtime is disabled or disconnected; when the
WebSocket is SUBSCRIBED, the claim loop sleeps for an hour and relies on
Postgres Changes notifications to drive claims. RealtimeManager tracks the
connection state via the subscribe callback status and exposes is_connected().
4. Remove unused dal.get_pending_conversations (only claim_conversations is used).
5. dal.get_conversation_events: drop unused min_seq parameter, add
include_compacted=False filter. By default compacted rows are filtered out
since the compaction event's messages array supersedes them.
Compaction verification:
- Unit tests cover the full publisher contract:
• CONVERSATION_HISTORY_COMPACTED event → post_conversation_events with
_compact=True (exactly once per compaction)
• CONVERSATION_HISTORY_COMPACTION_START alone → never sets _compact=True
• Normal flows (no compaction) → _compact=False on every call
• All StreamEvents values round-trip through the publisher
- DAL contract test: compact=True is forwarded to the RPC with the right arg
name; default is False.
- DAL contract test: get_conversation_events filters compacted by default and
returns all rows when include_compacted=True.
- End-to-end: a conversation triggered compaction 5 times against live Holmes.
The conversation completed successfully with ai_answer_end. The publisher
correctly set _compact=True on each compacted event batch. A separate direct
RPC test (/tmp/test_compaction.py) confirmed the staging M1 post_conversation_events
RPC does NOT perform the compacted-row UPDATE — this is an M1 server-side
bug, not an M2 code issue. M2's responsibility is to pass the flag, which
is verified.
Coverage:
__init__.py 100%
models.py 100%
event_publisher.py 88%
worker.py 47%
realtime_manager.py 26% (async WS code — unit tests cover helpers +
state machine; integration tests cover the
full subscribe + track + on-change pipeline)
New test files:
test_dal_contract.py — RPC argument forwarding, compacted filter
test_worker_polling.py — realtime-gated polling gate
test_worker_lifecycle.py — claim dispatch, error handling, presence leave
test_realtime_manager.py — initial state, loop-absent no-ops, idempotent proxy patch
Integration tests that were re-run and passing against staging:
single-turn — 17s, correct ES answer
multi-turn followup — 9-message history, correct contextual answer
tool-approval (bash) — JOIN on turn 1, approve via followup, turn 2
emits approved tool_calling_result + ai_answer_end
stop_conversation — Holmes publishes event → RPC rejects with
"is not running (status=stopped)" → worker
catches ConversationReassignedError → row stays
stopped (not flipped to completed/failed)
cluster Presence — observer sees full holmes_id + version metadata
per-conversation Presence — observer sees JOIN on claim + LEAVE on finish
polling-only mode — realtime disabled, works via 5s poll
heavy compaction — 5 compactions + final ai_answer_end, success
Signed-off-by: Claude <noreply@anthropic.com>
|
Caution Failed to replace (edit) comment. This is likely due to insufficient permissions or the comment being deleted. Error details |
There was a problem hiding this comment.
Actionable comments posted: 3
♻️ Duplicate comments (5)
holmes/core/supabase_dal.py (1)
1015-1038:⚠️ Potential issue | 🟡 MinorScope
get_conversation_events()byaccount_id.This query still loads rows by
conversation_idalone. If a bad ID is passed here, the worker can hydrate history from the wrong tenant as soon as RLS is loosened or misconfigured. Add.eq("account_id", self.account_id)to keep this path tenant-scoped by default.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/supabase_dal.py` around lines 1015 - 1038, The get_conversation_events method is currently filtering only by conversation_id which can return rows from the wrong tenant; update the Supabase query in get_conversation_events to add tenant scoping by calling .eq("account_id", self.account_id) on the query (before applying the compacted filter and ordering) so the query becomes scoped by both conversation_id and account_id; ensure self.account_id is used (and available) when building the query in this method.holmes/core/conversations_worker/realtime_manager.py (1)
183-186:⚠️ Potential issue | 🟡 MinorLog callback failures instead of swallowing them.
These
except Exception: passblocks hide missed wake-up paths and make reconnect issues much harder to diagnose. Logging and continuing preserves the fallback behavior without making callback failures invisible.Also applies to: 269-281
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 183 - 186, The code is silently swallowing exceptions from callback invocations (e.g. the call to self.on_new_pending()), which hides failures; change these try/except blocks to catch Exception and log the full exception (use logger.exception or self.logger.exception with a clear message like "callback failed in realtime_manager: on_new_pending") then continue, and apply the same change to the other callback invocation block referenced (the block around lines 269–281) so all callback failures are logged rather than ignored.holmes/core/conversations_worker/worker.py (3)
400-409:⚠️ Potential issue | 🟡 MinorSkip malformed history entries before dict access.
msg.get(...)assumes every history item is a dict. One bad persisted event will turn resume hydration into a hard failure here; guard withisinstance(msg, dict)first and ignore malformed entries.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 400 - 409, The loop iterating "for msg in reversed(history):" assumes each history item is a dict and will crash on malformed entries; add a guard at the top of that loop (e.g., if not isinstance(msg, dict): continue) so non-dict entries are skipped before calling msg.get; keep the existing checks for content being str or list and the inner part dict checks intact to safely handle vision parts.
211-214:⚠️ Potential issue | 🟠 MajorAvoid logging the full conversation row on parse failures.
convcan include user asks, history, and other metadata. Logging the whole payload here leaks tenant data into logs; keep this to stable identifiers such asconversation_idandrequest_sequence.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 211 - 214, The exception handler currently logs the entire conv object; change it to log only stable identifiers (e.g., conversation_id and request_sequence) instead of the full payload: extract these from conv (e.g., conv.get("conversation_id") and conv.get("request_sequence") or the equivalent keys used in this code) and pass them into the logging.exception message (do not include conv itself), keeping exc_info=True so the stacktrace is preserved; update the logging call in the except block around the conversation task build (where conv is referenced) to only emit those identifiers and no user data.
176-196:⚠️ Potential issue | 🟠 MajorClaim only the remaining worker slots.
The current guard only skips when
active >= CONVERSATION_WORKER_MAX_CONCURRENT. If one slot is free and the RPC returns 10 rows, this path will still submit all 10 and overshoot the configured concurrency. Compute the remaining capacity first and pass it intoclaim_conversations(...)so the RPC can cap the claim set.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 176 - 196, The _try_claim_and_dispatch method can overshoot CONVERSATION_WORKER_MAX_CONCURRENT because it calls dal.claim_conversations(self.holmes_id) unbounded; change it to compute remaining capacity = CONVERSATION_WORKER_MAX_CONCURRENT - active (after reading active under self._active_lock) and pass that capacity into claim_conversations (e.g., claim_conversations(self.holmes_id, capacity)) so the RPC returns at most the number of slots available, then proceed to build tasks and submit as before; update any claim_conversations signature and call sites accordingly to accept and enforce the limit.
🧹 Nitpick comments (2)
tests/core/conversations_worker/test_realtime_manager.py (1)
48-67: Hoistrealtime._async.clientto module scope.Importing the module inside each test violates the repo’s Python import rule and makes the patched module state harder to reason about across tests. Prefer a top-level import and patch/reset its attributes inside the test bodies.
As per coding guidelines,
ALWAYS place Python imports at the top of the file, not inside functions or methods.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/core/conversations_worker/test_realtime_manager.py` around lines 48 - 67, Move the import of realtime._async.client to module scope (top of the test file) instead of importing inside each test; then update tests to operate on that module object by resetting/inspecting its attributes (reset rt._holmes_proxy_patched, preserve/restore rt.connect around _install_proxy_patch_if_needed calls) so each test patches and unpatches explicitly and remains idempotent; reference rt (the hoisted realtime._async.client), its _holmes_proxy_patched flag, connect function, and the _install_proxy_patch_if_needed helper when making these changes.tests/core/conversations_worker/test_dal_contract.py (1)
1-10: Keep this contract test alongsidesupabase_dal.py.This module is asserting the contract of
holmes/core/supabase_dal.py, but it lives undertests/core/conversations_worker/. Moving it to the matching Supabase DAL test path will make the ownership and failure surface much clearer.As per coding guidelines,
Test files should match source structure under tests/ directory.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/core/conversations_worker/test_dal_contract.py` around lines 1 - 10, This contract test for SupabaseDal is in the wrong tree; move tests/core/conversations_worker/test_dal_contract.py so it mirrors the source module holmes/core/supabase_dal.py (e.g., tests/core/supabase_dal/test_dal_contract.py or tests/core/supabase_dal.py), update any relative imports if needed, and ensure the test still imports SupabaseDal and any test fakes correctly so ownership and failures surface next to the SupabaseDal implementation.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 231-252: The pg change callback currently wakes for all
account-wide Conversation writes; modify the subscription or callback to only
react to this worker's cluster pending rows: add cluster_id and pending-status
to the subscription filter used in self._cluster_channel.on_postgres_changes
(e.g. extend account_id_filter with cluster_id=eq.{self.dal.cluster_id} and
status=eq.pending) or, if unable to express that in the filter, update the
_on_pg_change(payload: Dict[str, Any]) function to inspect payload.get("data",
{}) for matching "cluster_id" == self.dal.cluster_id and the new row's status ==
"pending" before calling self.on_new_pending(); keep the logging and exception
handling intact.
In `@holmes/core/conversations_worker/worker.py`:
- Around line 219-236: In _process_conversation_safe, add a specific except
branch for ConversationReassignedError raised by _process_conversation so we
don't run the generic failure path; catch ConversationReassignedError before the
broad Exception handler and simply log or return without calling
dal.complete_conversation (leave the existing generic except Exception to mark
status "failed" and call complete_conversation); reference
_process_conversation_safe, _process_conversation, ConversationReassignedError,
and dal.complete_conversation to locate and implement the change.
In `@tests/core/conversations_worker/test_worker_lifecycle.py`:
- Around line 145-164: The test
test_process_conversation_safe_always_leaves_presence_on_error should also
assert that a reassignment does not mark the conversation complete; update the
test to patch ConversationWorker._process_conversation to raise
ConversationReassignedError (as done) and then add an assertion that
w.dal.complete_conversation.assert_not_called() after calling
ConversationWorker._process_conversation_safe(task), keeping the existing check
that rt.leave_conversation_presence.assert_called_once_with("c1") to ensure
leave runs in the finally block.
---
Duplicate comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 183-186: The code is silently swallowing exceptions from callback
invocations (e.g. the call to self.on_new_pending()), which hides failures;
change these try/except blocks to catch Exception and log the full exception
(use logger.exception or self.logger.exception with a clear message like
"callback failed in realtime_manager: on_new_pending") then continue, and apply
the same change to the other callback invocation block referenced (the block
around lines 269–281) so all callback failures are logged rather than ignored.
In `@holmes/core/conversations_worker/worker.py`:
- Around line 400-409: The loop iterating "for msg in reversed(history):"
assumes each history item is a dict and will crash on malformed entries; add a
guard at the top of that loop (e.g., if not isinstance(msg, dict): continue) so
non-dict entries are skipped before calling msg.get; keep the existing checks
for content being str or list and the inner part dict checks intact to safely
handle vision parts.
- Around line 211-214: The exception handler currently logs the entire conv
object; change it to log only stable identifiers (e.g., conversation_id and
request_sequence) instead of the full payload: extract these from conv (e.g.,
conv.get("conversation_id") and conv.get("request_sequence") or the equivalent
keys used in this code) and pass them into the logging.exception message (do not
include conv itself), keeping exc_info=True so the stacktrace is preserved;
update the logging call in the except block around the conversation task build
(where conv is referenced) to only emit those identifiers and no user data.
- Around line 176-196: The _try_claim_and_dispatch method can overshoot
CONVERSATION_WORKER_MAX_CONCURRENT because it calls
dal.claim_conversations(self.holmes_id) unbounded; change it to compute
remaining capacity = CONVERSATION_WORKER_MAX_CONCURRENT - active (after reading
active under self._active_lock) and pass that capacity into claim_conversations
(e.g., claim_conversations(self.holmes_id, capacity)) so the RPC returns at most
the number of slots available, then proceed to build tasks and submit as before;
update any claim_conversations signature and call sites accordingly to accept
and enforce the limit.
In `@holmes/core/supabase_dal.py`:
- Around line 1015-1038: The get_conversation_events method is currently
filtering only by conversation_id which can return rows from the wrong tenant;
update the Supabase query in get_conversation_events to add tenant scoping by
calling .eq("account_id", self.account_id) on the query (before applying the
compacted filter and ordering) so the query becomes scoped by both
conversation_id and account_id; ensure self.account_id is used (and available)
when building the query in this method.
---
Nitpick comments:
In `@tests/core/conversations_worker/test_dal_contract.py`:
- Around line 1-10: This contract test for SupabaseDal is in the wrong tree;
move tests/core/conversations_worker/test_dal_contract.py so it mirrors the
source module holmes/core/supabase_dal.py (e.g.,
tests/core/supabase_dal/test_dal_contract.py or tests/core/supabase_dal.py),
update any relative imports if needed, and ensure the test still imports
SupabaseDal and any test fakes correctly so ownership and failures surface next
to the SupabaseDal implementation.
In `@tests/core/conversations_worker/test_realtime_manager.py`:
- Around line 48-67: Move the import of realtime._async.client to module scope
(top of the test file) instead of importing inside each test; then update tests
to operate on that module object by resetting/inspecting its attributes (reset
rt._holmes_proxy_patched, preserve/restore rt.connect around
_install_proxy_patch_if_needed calls) so each test patches and unpatches
explicitly and remains idempotent; reference rt (the hoisted
realtime._async.client), its _holmes_proxy_patched flag, connect function, and
the _install_proxy_patch_if_needed helper when making these changes.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 73afdda1-db8a-472f-9b5e-26e749577f79
📒 Files selected for processing (9)
holmes/common/env_vars.pyholmes/core/conversations_worker/realtime_manager.pyholmes/core/conversations_worker/worker.pyholmes/core/supabase_dal.pytests/core/conversations_worker/test_dal_contract.pytests/core/conversations_worker/test_event_publisher.pytests/core/conversations_worker/test_realtime_manager.pytests/core/conversations_worker/test_worker_lifecycle.pytests/core/conversations_worker/test_worker_polling.py
✅ Files skipped from review due to trivial changes (1)
- tests/core/conversations_worker/test_event_publisher.py
🚧 Files skipped from review as they are similar to previous changes (1)
- holmes/common/env_vars.py
…C reads)
The M1 plan and DB schema were updated. This commit aligns Holmes M2 with
those changes:
1. RPC parameter rename: ``_holmes_id`` → ``_assignee`` for
``post_conversation_events`` and ``complete_conversation``. The DAL
method signatures use ``assignee=`` as the kwarg; the worker's
internal ``self.holmes_id`` attribute is unchanged (matches existing
ScheduledPromptsExecutor pattern). ``claim_conversations`` already
used ``_assignee``.
2. ConversationEvents schema: ``request_sequence`` is no longer a column
on event rows. ``seq`` is now globally monotonic per conversation
(UNIQUE on ``(conversation_id, seq)``). Turn boundaries are detected
by the ``user_message`` event itself, not a per-turn sequence number.
``Conversations.request_sequence`` (parent row) still tracks the
current turn for optimistic concurrency in RPC writes — that is
unchanged.
3. ``get_conversation_events`` now goes through a SECURITY DEFINER RPC
instead of a direct table SELECT — Holmes does not need direct read
access to ``ConversationEvents`` under RLS. The RPC flattens all
matching rows' events arrays into a single chronological list
ordered by ``(seq, ord)``, and accepts ``_include_compacted`` and
``_min_seq`` parameters. The DAL signature is::
get_conversation_events(conversation_id, include_compacted=False,
min_seq=1) -> List[event_dict]
Each element is a flat ``{"event": ..., "data": ..., "ts": ...}``
dict with no row metadata (seq/compacted not exposed at this layer).
4. ``_hydrate_task_from_events`` rewritten to consume the flat event
list. Algorithm: find the last ``user_message`` index → that is the
current turn's request; among events with index < that, find the
last ``ai_answer_end`` / ``approval_required`` → its ``messages``
array is the conversation history for the LLM. Cleaner than the
previous row-of-events scan.
Unit tests (44 passing):
- test_dal_contract.py: covers all four DAL methods' RPC argument
forwarding, including ``_assignee`` rename, ``_include_compacted``,
``_min_seq``, and the flat-event return shape.
- test_worker_hydration.py: fixtures rewritten as flat event lists
(no row/seq nesting). Covers first-turn, follow-up reconstruction,
approval-required history, the "stale terminal after current
user_message" guard, and model/extra-fields extraction.
- test_event_publisher.py + test_worker_polling.py +
test_realtime_manager.py + test_worker_lifecycle.py: all updated
for the assignee rename.
Integration tests (re-run against staging, all green):
- single-turn (15s, correct answer)
- multi-turn followup (9-msg history reconstruction)
- tool-approval flow (turn boundaries detected correctly via
user_message events; full pause/resume cycle)
- stop_conversation → ConversationReassignedError
- per-conversation Presence JOIN+LEAVE
- direct RPC compaction (5 events, all flags correct, DAL flat-list
return verified by event-type assertions)
- end-to-end compaction strict (15 events, watermark seq=13, every
compacted flag matches expected)
- cross-turn compaction verified directly: 48-row conversation, all
5 turn-1 rows correctly compacted=True after turn-2 compaction
- parallel stress: 15 concurrent conversations across single-turn,
multi-turn, and tool-approval suites — all passed
Signed-off-by: Claude <noreply@anthropic.com>
There was a problem hiding this comment.
♻️ Duplicate comments (5)
holmes/core/conversations_worker/worker.py (4)
402-417:⚠️ Potential issue | 🟡 MinorGuard against malformed history entries before dict access.
msg.get(...)assumes every entry is a dict. If a non-dict entry slips through, this raises anAttributeErrorand aborts extraction.🔧 Proposed fix
`@staticmethod` def _extract_last_user_ask(history: Optional[list]) -> Optional[str]: """Pull the most recent user message text from an OpenAI-format history.""" if not history: return None for msg in reversed(history): + if not isinstance(msg, dict): + continue if msg.get("role") == "user": content = msg.get("content") if isinstance(content, str): return content if isinstance(content, list): # Vision message: find the first text part for part in content: if isinstance(part, dict) and part.get("type") == "text": return part.get("text") return None🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 402 - 417, The _extract_last_user_ask function assumes each history entry is a dict and directly calls msg.get, which can raise if msg is not a dict; update _extract_last_user_ask to first check isinstance(msg, dict) before accessing msg.get and skip any non-dict entries, and likewise validate that content is either a str or a list before iterating (also ensure list elements are dicts before .get) so malformed history entries are ignored instead of raising.
176-196:⚠️ Potential issue | 🟠 MajorClaiming should be capped by remaining worker slots.
The worker checks capacity at lines 178-184 but then claims all available conversations at line 186 without passing a limit. If the RPC returns more conversations than available slots, they all get dispatched, potentially exceeding
CONVERSATION_WORKER_MAX_CONCURRENT.🔧 Proposed fix
def _try_claim_and_dispatch(self) -> None: - # Respect max concurrency: if we're at capacity, skip claiming + # Respect max concurrency: claim only remaining slots with self._active_lock: active = len(self._active_conversation_ids) - if active >= CONVERSATION_WORKER_MAX_CONCURRENT: + available_slots = CONVERSATION_WORKER_MAX_CONCURRENT - active + if available_slots <= 0: logging.debug( "At max concurrency (%d), skipping claim", active ) return - claimed = self.dal.claim_conversations(self.holmes_id) + claimed = self.dal.claim_conversations(self.holmes_id, limit=available_slots)Also update
SupabaseDal.claim_conversations()to accept and forward thelimitparameter to the RPC.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 176 - 196, The _try_claim_and_dispatch method currently checks capacity but then calls claim_conversations(self.holmes_id) with no limit, allowing more conversations to be claimed than available slots; change _try_claim_and_dispatch to compute remaining_slots = CONVERSATION_WORKER_MAX_CONCURRENT - len(self._active_conversation_ids) (use the same locking around _active_conversation_ids), pass that remaining_slots as a limit argument to dal.claim_conversations(self.holmes_id, limit=remaining_slots), and only proceed if remaining_slots > 0; also update SupabaseDal.claim_conversations to accept a limit parameter and forward it to the RPC so the DB returns at most that many rows.
211-215:⚠️ Potential issue | 🟡 MinorAvoid logging full conversation rows on parse failure.
Line 213 logs the entire
convpayload which may contain sensitive metadata. Log only stable identifiers.🔧 Proposed fix
except Exception: logging.exception( - "Failed to build conversation task from row: %s", conv, exc_info=True + "Failed to build conversation task from row (conversation_id=%s, request_sequence=%s)", + conv.get("conversation_id"), + conv.get("request_sequence"), + exc_info=True, ) return None🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 211 - 215, The except block currently logs the full conv payload which may contain sensitive data; update the logging.exception call in the failing parse handler (the except around the conv parsing where variable conv is available) to log only stable non-sensitive identifiers (e.g., conv.get("conversation_id") or conv.get("id") and conv.get("user_id") / conv.get("tenant_id")) instead of the full conv object, and keep the exc_info=True behavior and the existing return None. Ensure you reference the same logging.exception call (the one that currently includes conv) and replace the conv argument with the chosen identifier fields.
219-242:⚠️ Potential issue | 🟠 MajorHandle
ConversationReassignedErrorbefore the generic failure path.If
ConversationReassignedErrorescapes_process_conversation(), this block marks the conversation as"failed"even though another worker may now own it. CatchConversationReassignedErrorseparately and skipcomplete_conversation(...)in that branch.🔧 Proposed fix
def _process_conversation_safe(self, task: ConversationTask) -> None: try: self._process_conversation(task) + except ConversationReassignedError: + logging.info( + "Conversation %s was reassigned mid-processing, skipping completion", + task.conversation_id, + ) except Exception as e: logging.exception( "Error processing conversation %s: %s", task.conversation_id, e, exc_info=True, ) # Attempt to mark as failed try: self.dal.complete_conversation( conversation_id=task.conversation_id, request_sequence=task.request_sequence, assignee=self.holmes_id, status="failed", ) except Exception: logging.exception( "Failed to mark conversation %s as failed", task.conversation_id, exc_info=True, )🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 219 - 242, The except block in _process_conversation_safe currently treats all exceptions the same and will call self.dal.complete_conversation even when a ConversationReassignedError indicates another worker owns the conversation; update _process_conversation_safe to add a separate except ConversationReassignedError branch that only logs (or returns) and does NOT call self.dal.complete_conversation (leave the existing generic Exception handler to log and call self.dal.complete_conversation with conversation_id=task.conversation_id, request_sequence=task.request_sequence, assignee=self.holmes_id, status="failed"); ensure references to _process_conversation, ConversationReassignedError, and self.dal.complete_conversation are used to locate the change.tests/core/conversations_worker/test_worker_lifecycle.py (1)
145-165: 🛠️ Refactor suggestion | 🟠 MajorAssert that reassignment does not call
complete_conversation().This test raises
ConversationReassignedErrorbut only verifies thatleave_conversation_presenceis called. It should also assert thatcomplete_conversationis not called, since marking a reassigned conversation as failed is incorrect behavior.🔧 Proposed fix
with patch.object(ConversationWorker, "_process_conversation", boom): w._process_conversation_safe(task) # leave must run in the finally even after an error rt.leave_conversation_presence.assert_called_once_with("c1") + # Reassignment should NOT mark the conversation as failed + w.dal.complete_conversation.assert_not_called()Note: This assertion will currently fail because
_process_conversation_safecatchesConversationReassignedErrorin the genericExceptionhandler. Once the worker is fixed to handleConversationReassignedErrorseparately (per the comment onworker.py), this test will pass.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/core/conversations_worker/test_worker_lifecycle.py` around lines 145 - 165, The test should also assert that a ConversationReassignedError does not trigger finishing the conversation: after calling w._process_conversation_safe(task) when _process_conversation raises ConversationReassignedError, add an assertion that complete_conversation was not invoked (e.g., assert that w.complete_conversation or the worker's complete_conversation mock was not called). This ensures ConversationWorker._process_conversation_safe treats ConversationReassignedError specially (leaves presence via rt.leave_conversation_presence but does not call complete_conversation).
🧹 Nitpick comments (3)
holmes/core/conversations_worker/realtime_manager.py (2)
183-186: Log exception in_run()finally block.Similar to the subscription callbacks, this
try-except-passsilently discards errors when waking the worker for fallback polling.🔧 Proposed fix
finally: self._connected = False # Wake the worker so it falls back to polling try: self.on_new_pending() except Exception: - pass + logging.debug("Error in on_new_pending callback during shutdown", exc_info=True)🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 183 - 186, The finally block in _run() currently swallows exceptions from self.on_new_pending() with an empty except; change it to catch Exception as e and log the error (e.g., self.logger.exception(...) or self.logger.error(..., exc_info=True)) so failures in on_new_pending() are recorded; update the except in the _run() method to reference the exception variable and emit a clear log message including the traceback.
82-84: Star-arg unpacking after keyword argument.Line 84 uses
*argsaftersock=sock, which is discouraged (Ruff B026). Reorder to place positional args before keyword args.🔧 Proposed fix
- return await ws_connect(url, sock=sock, *args, **kwargs) + return await ws_connect(url, *args, sock=sock, **kwargs)🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 82 - 84, The ws_connect call uses positional star-arg unpacking after a keyword (sock=sock) which triggers Ruff B026; reorder the arguments so positional args come first. Update the call in realtime_manager.py to pass url and *args before keyword arguments (e.g., await ws_connect(url, *args, sock=sock, **kwargs)), keeping the parsed.scheme check and ssl.create_default_context() usage intact.tests/core/conversations_worker/test_realtime_manager.py (1)
62-79: Consider using a pytest fixture for cleanup to ensure restoration on test failure.The cleanup at lines 77-79 won't run if the test fails before reaching that point. A
finallyblock or fixture would be more robust.🔧 Proposed fix using try/finally
def test_install_proxy_patch_is_idempotent(monkeypatch): """Calling install twice should not double-patch.""" monkeypatch.setenv( "https_proxy", "http://user:pass@proxy.internal:8888" ) import realtime._async.client as rt - rt._holmes_proxy_patched = False original_connect = rt.connect - _install_proxy_patch_if_needed() - first_patched = rt.connect - _install_proxy_patch_if_needed() - second_patched = rt.connect - assert first_patched is second_patched, "patch was reinstalled unexpectedly" - - # Cleanup: restore the original connect fn - rt.connect = original_connect - rt._holmes_proxy_patched = False + rt._holmes_proxy_patched = False + try: + _install_proxy_patch_if_needed() + first_patched = rt.connect + _install_proxy_patch_if_needed() + second_patched = rt.connect + assert first_patched is second_patched, "patch was reinstalled unexpectedly" + finally: + # Cleanup: restore the original connect fn + rt.connect = original_connect + rt._holmes_proxy_patched = False🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/core/conversations_worker/test_realtime_manager.py` around lines 62 - 79, The test test_install_proxy_patch_is_idempotent should ensure cleanup even on failure: capture the original state (rt.connect and rt._holmes_proxy_patched) before calling _install_proxy_patch_if_needed and wrap the test body in a try/finally so the finally block always restores rt.connect and rt._holmes_proxy_patched; reference the module variables rt.connect, rt._holmes_proxy_patched and the helper _install_proxy_patch_if_needed when making the changes.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Duplicate comments:
In `@holmes/core/conversations_worker/worker.py`:
- Around line 402-417: The _extract_last_user_ask function assumes each history
entry is a dict and directly calls msg.get, which can raise if msg is not a
dict; update _extract_last_user_ask to first check isinstance(msg, dict) before
accessing msg.get and skip any non-dict entries, and likewise validate that
content is either a str or a list before iterating (also ensure list elements
are dicts before .get) so malformed history entries are ignored instead of
raising.
- Around line 176-196: The _try_claim_and_dispatch method currently checks
capacity but then calls claim_conversations(self.holmes_id) with no limit,
allowing more conversations to be claimed than available slots; change
_try_claim_and_dispatch to compute remaining_slots =
CONVERSATION_WORKER_MAX_CONCURRENT - len(self._active_conversation_ids) (use the
same locking around _active_conversation_ids), pass that remaining_slots as a
limit argument to dal.claim_conversations(self.holmes_id,
limit=remaining_slots), and only proceed if remaining_slots > 0; also update
SupabaseDal.claim_conversations to accept a limit parameter and forward it to
the RPC so the DB returns at most that many rows.
- Around line 211-215: The except block currently logs the full conv payload
which may contain sensitive data; update the logging.exception call in the
failing parse handler (the except around the conv parsing where variable conv is
available) to log only stable non-sensitive identifiers (e.g.,
conv.get("conversation_id") or conv.get("id") and conv.get("user_id") /
conv.get("tenant_id")) instead of the full conv object, and keep the
exc_info=True behavior and the existing return None. Ensure you reference the
same logging.exception call (the one that currently includes conv) and replace
the conv argument with the chosen identifier fields.
- Around line 219-242: The except block in _process_conversation_safe currently
treats all exceptions the same and will call self.dal.complete_conversation even
when a ConversationReassignedError indicates another worker owns the
conversation; update _process_conversation_safe to add a separate except
ConversationReassignedError branch that only logs (or returns) and does NOT call
self.dal.complete_conversation (leave the existing generic Exception handler to
log and call self.dal.complete_conversation with
conversation_id=task.conversation_id, request_sequence=task.request_sequence,
assignee=self.holmes_id, status="failed"); ensure references to
_process_conversation, ConversationReassignedError, and
self.dal.complete_conversation are used to locate the change.
In `@tests/core/conversations_worker/test_worker_lifecycle.py`:
- Around line 145-165: The test should also assert that a
ConversationReassignedError does not trigger finishing the conversation: after
calling w._process_conversation_safe(task) when _process_conversation raises
ConversationReassignedError, add an assertion that complete_conversation was not
invoked (e.g., assert that w.complete_conversation or the worker's
complete_conversation mock was not called). This ensures
ConversationWorker._process_conversation_safe treats ConversationReassignedError
specially (leaves presence via rt.leave_conversation_presence but does not call
complete_conversation).
---
Nitpick comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 183-186: The finally block in _run() currently swallows exceptions
from self.on_new_pending() with an empty except; change it to catch Exception as
e and log the error (e.g., self.logger.exception(...) or self.logger.error(...,
exc_info=True)) so failures in on_new_pending() are recorded; update the except
in the _run() method to reference the exception variable and emit a clear log
message including the traceback.
- Around line 82-84: The ws_connect call uses positional star-arg unpacking
after a keyword (sock=sock) which triggers Ruff B026; reorder the arguments so
positional args come first. Update the call in realtime_manager.py to pass url
and *args before keyword arguments (e.g., await ws_connect(url, *args,
sock=sock, **kwargs)), keeping the parsed.scheme check and
ssl.create_default_context() usage intact.
In `@tests/core/conversations_worker/test_realtime_manager.py`:
- Around line 62-79: The test test_install_proxy_patch_is_idempotent should
ensure cleanup even on failure: capture the original state (rt.connect and
rt._holmes_proxy_patched) before calling _install_proxy_patch_if_needed and wrap
the test body in a try/finally so the finally block always restores rt.connect
and rt._holmes_proxy_patched; reference the module variables rt.connect,
rt._holmes_proxy_patched and the helper _install_proxy_patch_if_needed when
making the changes.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: c9ffc244-7910-40ee-8a01-634710b5befc
📒 Files selected for processing (11)
holmes/common/env_vars.pyholmes/core/conversations_worker/event_publisher.pyholmes/core/conversations_worker/realtime_manager.pyholmes/core/conversations_worker/worker.pyholmes/core/supabase_dal.pytests/core/conversations_worker/test_dal_contract.pytests/core/conversations_worker/test_event_publisher.pytests/core/conversations_worker/test_realtime_manager.pytests/core/conversations_worker/test_worker_hydration.pytests/core/conversations_worker/test_worker_lifecycle.pytests/core/conversations_worker/test_worker_polling.py
✅ Files skipped from review due to trivial changes (2)
- tests/core/conversations_worker/test_event_publisher.py
- holmes/core/conversations_worker/event_publisher.py
🚧 Files skipped from review as they are similar to previous changes (3)
- holmes/common/env_vars.py
- tests/core/conversations_worker/test_worker_hydration.py
- holmes/core/supabase_dal.py
The M1 plan was updated so every conversation RPC takes _account_id as
its first argument (consistent ordering, ergonomic for SECURITY DEFINER
auth checks). claim_conversations and get_conversation_events already
had it. This commit adds _account_id to the two remaining DAL RPCs:
- post_conversation_events
- complete_conversation
Both pass self.account_id from the SupabaseDal instance. The DAL method
signatures are unchanged (callers don't need to pass account_id).
Verified live against staging (which now requires _account_id on all
four conversation RPCs):
$ poetry run python -c "...probe each rpc..."
post_conversation_events without _account_id → PGRST202 not found
post_conversation_events with _account_id → P0001 (works)
complete_conversation without _account_id → PGRST202 not found
complete_conversation with _account_id → P0001 (works)
Test results after the change:
Unit tests:
44 passed in 46s
Integration (4 in parallel against staging):
✓ single-turn (17.7s)
✓ multi-turn followup (16s, 9-msg history reconstructed)
✓ tool approval flow (turn boundaries via user_message)
✓ per-conversation Presence JOIN+LEAVE
✓ stop_conversation → ConversationReassignedError
Stress (3 suites running concurrently, 15 conversations):
batch-parallel: 8/8 in 60s
multi-turn: 4/4 in 48s
approval: 3/3 in 23s
Compaction:
✓ direct RPC test: all 5 rows have correct compacted flag
✓ e2e with small context (40k): 40 rows, 19 compactions,
watermark seq=39, every compacted flag matches expected
✓ DAL get_conversation_events RPC returns 9 events (default,
uncompacted only) vs 187 events (include_compacted=True)
Signed-off-by: Claude <noreply@anthropic.com>
Five targeted fixes found during review: 1. models.py: typed ConversationTask.conversation_history as ``Optional[List[Dict[str, Any]]]`` (was bare ``Optional[list]``). Matches its usage in worker.py and better aligns with the OpenAI chat-message structure the field holds. 2. realtime_manager.py: the Postgres Changes callback previously woke the claim loop on every Conversations INSERT/UPDATE for the account — including rows for other clusters and non-pending status changes. Supabase Realtime's filter syntax only supports one operator per subscription, so we keep the account_id filter on the server and do cluster/status filtering client-side in ``_on_pg_change``: skip the notification unless the row's ``cluster_id == self.dal.cluster`` AND ``status == 'pending'``. Also tolerate both payload shapes (``record`` and legacy ``new``) for the row body. 3. worker.py ``_process_conversation_safe``: added a specific ``ConversationReassignedError`` branch before the generic ``Exception`` handler. On reassignment we must NOT call ``complete_conversation(status='failed')`` — the conversation's state is already being driven by whoever reassigned it (stop_conversation bumped request_sequence, or another Holmes took over). A stale complete_conversation call would either fail the RPC's status guard or race with the new owner. Presence leave still runs in the finally block regardless. 4. worker.py ``_extract_last_user_ask``: guarded the loop against malformed history entries. Now skips non-dict items, and for vision content requires the text field to be a non-empty string. Defensive change — under normal operation every entry is a dict, but a hand-crafted / backfilled history could contain surprises. 5. test_worker_lifecycle.py: strengthened the reassignment test to assert ``dal.complete_conversation.assert_not_called()`` after ``ConversationReassignedError``, covering fix #3. Also expanded the docstring to explain why the assertion matters. Verification: 44 unit tests pass (now covering the stronger reassignment invariant) Integration against staging: - single-turn / multi-turn / tool approval / per-conv presence - stop_conversation → ConversationReassignedError (row stays 'stopped'; complete_conversation NOT called — verified by new test) - 15-conversation concurrent stress (8 batch + 4 multi + 3 approval) - direct RPC compaction (all compacted flags correct) Signed-off-by: Claude <noreply@anthropic.com>
There was a problem hiding this comment.
Actionable comments posted: 4
♻️ Duplicate comments (3)
holmes/core/conversations_worker/worker.py (2)
176-196:⚠️ Potential issue | 🟠 MajorClaim only the remaining worker capacity.
This still asks the DAL for pending conversations without a limit once
active < max. A worker with one free slot can assign itself the entire backlog and exceedCONVERSATION_WORKER_MAX_CONCURRENT.🛠 Suggested direction
with self._active_lock: active = len(self._active_conversation_ids) - if active >= CONVERSATION_WORKER_MAX_CONCURRENT: + available_slots = CONVERSATION_WORKER_MAX_CONCURRENT - active + if available_slots <= 0: logging.debug( "At max concurrency (%d), skipping claim", active ) return - claimed = self.dal.claim_conversations(self.holmes_id) + claimed = self.dal.claim_conversations( + self.holmes_id, + limit=available_slots, + )The DAL/RPC needs the matching
limitparameter as well.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 176 - 196, The _try_claim_and_dispatch method can over-claim because it asks dal.claim_conversations for all pending items even when there are only a few free slots; change the call to dal.claim_conversations(self.holmes_id) to pass a limit equal to the remaining capacity (CONVERSATION_WORKER_MAX_CONCURRENT - len(self._active_conversation_ids)) so the DAL only returns up to the number of slots available; update any call sites or DAL signature of claim_conversations to accept and honor that limit and keep the subsequent logic that adds to _active_conversation_ids and submits tasks to _executor unchanged.
190-193:⚠️ Potential issue | 🟠 MajorDon't drop a claimed conversation on task-build failure.
By the time this path runs,
claim_conversations()has already assigned the row. ReturningNonehere just abandons it, so the conversation can stay stuck inrunning; the exception path also logs the full row while doing it. Mark the claim failed/released before continuing and log only stable identifiers.Also applies to: 198-214
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 190 - 193, When _build_task_from_conversation_row(conv) returns None you must not leave the DB row stuck as claimed; before continuing, call the appropriate claim-release/fail helper (the same mechanism used by claim_conversations()/release or mark_claim_failed) to mark the claim as released/failed for that conversation (use the conversation's stable id fields such as conversation_id/row id and owner), then continue; also replace any heavy-row logging in these branches with logging of only stable identifiers. Apply this change in the loop handling claimed (the block around _build_task_from_conversation_row) and the other similar block referenced (lines ~198-214).holmes/core/conversations_worker/realtime_manager.py (1)
183-186:⚠️ Potential issue | 🟡 MinorDon't silently swallow
on_new_pending()failures.These
passblocks hide missed wake-ups if the callback ever stops being a bareEvent.set, which makes realtime fallback bugs much harder to diagnose.Also applies to: 285-297
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 183 - 186, Replace the bare "try: self.on_new_pending() except Exception: pass" swallowing behavior with explicit error handling: catch Exception around self.on_new_pending(), call the module/class logger.exception (or process_logger.exception) to record the stacktrace and context, and then either re-raise the exception or propagate a specific error so failures are visible (do the same for the other identical block around self.on_new_pending()). Ensure you reference the same self.on_new_pending() call sites and mirror the handling in both locations rather than silently passing.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 44-59: Move the runtime/local imports in realtime_manager.py
(e.g., realtime._async.client as rt_client, python_socks.async_.asyncio.Proxy,
and websockets.asyncio.client.connect) to module scope; if those imports cause
cycles or are optional, create a small helper module to encapsulate the lazy
import logic and expose safe helpers (preserving the _holmes_proxy_patched flag
behavior) instead of importing inside functions. Also add explicit type
annotations to the subscribe callbacks _on_subscribe_cb and _on_sub for the
parameters status and err (and return types) so they comply with the codebase
typing rules; apply the same import/typing cleanup to the other occurrences
noted (lines referenced in the review).
- Around line 120-141: start() currently doesn't clear the lifecycle events so
restarting uses stale state; before creating/starting the new thread clear the
events (call _stop_event.clear() and _started.clear()) so the new thread's
_thread_entry can run normally, and ensure any leftover _loop/_thread references
are reset if needed (use _thread or _loop checks already present) to allow a
clean restart; update start() to clear these events prior to creating the
Thread.
In `@holmes/core/conversations_worker/worker.py`:
- Around line 18-24: The file has imports after executable code and inside
functions and a partially-typed helper; move all imports (e.g.,
ConversationEventPublisher and any other imports currently below executable
statements or inside methods) to module top-level to satisfy the import-order
rule, and if those imports create an import cycle, extract the dependent code
into a new helper module and import that helper instead; complete the missing
type annotations on the helper function _terminal_to_status() (add precise
parameter and return types consistent with your Status/Terminal enums or classes
used in this module) and remove any in-function lazy imports (replace them with
top-level imports or references to the new helper module).
- Line 61: The current assignment to self.holmes_id (self.holmes_id =
os.environ.get("HOSTNAME") or str(os.getpid())) produces non-unique IDs across
processes; update the initialization in the worker class/constructor where
self.holmes_id is set so it includes a per-process unique component (e.g.,
append or replace with a uuid4 and/or include os.getpid()) — for example, build
holmes_id from HOSTNAME (if present) plus pid and a uuid4 string so each process
yields a globally unique identifier used for presence and assignee keys.
---
Duplicate comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 183-186: Replace the bare "try: self.on_new_pending() except
Exception: pass" swallowing behavior with explicit error handling: catch
Exception around self.on_new_pending(), call the module/class logger.exception
(or process_logger.exception) to record the stacktrace and context, and then
either re-raise the exception or propagate a specific error so failures are
visible (do the same for the other identical block around
self.on_new_pending()). Ensure you reference the same self.on_new_pending() call
sites and mirror the handling in both locations rather than silently passing.
In `@holmes/core/conversations_worker/worker.py`:
- Around line 176-196: The _try_claim_and_dispatch method can over-claim because
it asks dal.claim_conversations for all pending items even when there are only a
few free slots; change the call to dal.claim_conversations(self.holmes_id) to
pass a limit equal to the remaining capacity (CONVERSATION_WORKER_MAX_CONCURRENT
- len(self._active_conversation_ids)) so the DAL only returns up to the number
of slots available; update any call sites or DAL signature of
claim_conversations to accept and honor that limit and keep the subsequent logic
that adds to _active_conversation_ids and submits tasks to _executor unchanged.
- Around line 190-193: When _build_task_from_conversation_row(conv) returns None
you must not leave the DB row stuck as claimed; before continuing, call the
appropriate claim-release/fail helper (the same mechanism used by
claim_conversations()/release or mark_claim_failed) to mark the claim as
released/failed for that conversation (use the conversation's stable id fields
such as conversation_id/row id and owner), then continue; also replace any
heavy-row logging in these branches with logging of only stable identifiers.
Apply this change in the loop handling claimed (the block around
_build_task_from_conversation_row) and the other similar block referenced (lines
~198-214).
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 465e8161-f68d-474f-a508-d32fb9b75df1
📒 Files selected for processing (4)
holmes/core/conversations_worker/models.pyholmes/core/conversations_worker/realtime_manager.pyholmes/core/conversations_worker/worker.pytests/core/conversations_worker/test_worker_lifecycle.py
✅ Files skipped from review due to trivial changes (2)
- tests/core/conversations_worker/test_worker_lifecycle.py
- holmes/core/conversations_worker/models.py
Seven targeted findings from a code review:
1. event_publisher.consume: return type is Optional[StreamEvents] (was
StreamEvents with a "# type: ignore[return-value]" because the docstring
already documented it could return None). The ignore is gone.
2. event_publisher._flush: collapsed the two early-return branches
("no events" and "no events but compact requested") into a single
``if not events_to_flush: return``. Compact with no events is a
no-op either way.
3. event_publisher._flush: the M1 RPCs now prefix mismatch errors with
"MISMATCH ", so the client-side detection is a single
``"mismatch" in str(e).lower()`` instead of three specific substring
checks. Simpler and more robust as the set of mismatch conditions
evolves on the server.
4. realtime_manager: moved the runtime imports (
``realtime._async.client``, ``python_socks.async_.asyncio.Proxy``,
``websockets.asyncio.client.connect``) to module scope. The optional
python-socks is gated with a try/except at import time — if it's
missing the proxy patch becomes a no-op (with a warning) exactly as
before. Also typed the two subscribe callbacks
``_on_subscribe_cb(status: Any, err: Optional[Exception])`` and
``_on_sub(status: Any, err: Optional[Exception])``.
5. realtime_manager.start(): clear ``_stop_event`` and ``_started`` and
reset ``_loop``/``_client``/``_cluster_channel``/``_connected`` before
spawning the thread, so stop() → start() restart cycles don't run
against stale state.
6. worker: moved the in-function imports (``PromptComponent``,
``ToolsetTag``, ``tool_result_storage``, ``build_chat_messages``,
``TracingFactory``, ``StreamEvents``, ``RealtimeManager``) to
module top-level. Verified by launching server.py — no circular
import. Also annotated ``_terminal_to_status`` with
``(terminal: Optional[StreamEvents]) -> str`` and dropped its
in-function StreamEvents import.
7. worker.holmes_id: was ``os.environ.get("HOSTNAME") or str(os.getpid())``
which can collide across pod restarts (HOSTNAME reused) or between
tests. Now composed as ``{hostname}-{pid}-{uuid4[:8]}`` for global
uniqueness. This is the presence key AND the assignee value in
Conversations, so uniqueness matters for correctness under reclaims.
Example value in server log: ``local-36204-432d12a0``.
Verification:
- 44 unit tests pass
- Direct RPC compaction test: all 5 rows have correct compacted flag
- stop_conversation: worker catches the mismatch (now matched by
the simpler 'mismatch' check) and does NOT call complete_conversation
- per-conversation Presence JOIN+LEAVE observed by an external observer
- server starts cleanly, imports module-top, holmes_id format is
hostname-pid-uuid4 as designed
Integration tests (single/multi/approval) against staging succeed when
the LLM uses tool-calling. Currently on this Opus 4.6 session the model
is occasionally emitting ``<tool_call>`` tags as plain text instead of
function-calling — confirmed by hitting ``/api/chat`` directly with the
same prompt and getting the same output, so this is an LLM behavior
(not M2) issue and is outside the scope of this change.
Signed-off-by: Claude <noreply@anthropic.com>
There was a problem hiding this comment.
Actionable comments posted: 2
♻️ Duplicate comments (3)
holmes/core/conversations_worker/worker.py (2)
187-207:⚠️ Potential issue | 🟠 MajorCap claims to the remaining worker capacity.
holmes/core/supabase_dal.py:906-932currently claims all pending rows for the cluster, so once this method seesactive < CONVERSATION_WORKER_MAX_CONCURRENT, it can still dispatch far more work than the remaining slots allow. Computeavailable_slotshere and pass that limit through the DAL/RPC before submitting tasks.Suggested fix
def _try_claim_and_dispatch(self) -> None: - # Respect max concurrency: if we're at capacity, skip claiming + # Respect max concurrency: only claim the remaining available slots with self._active_lock: active = len(self._active_conversation_ids) - if active >= CONVERSATION_WORKER_MAX_CONCURRENT: + available_slots = CONVERSATION_WORKER_MAX_CONCURRENT - active + if available_slots <= 0: logging.debug( "At max concurrency (%d), skipping claim", active ) return - claimed = self.dal.claim_conversations(self.holmes_id) + claimed = self.dal.claim_conversations( + self.holmes_id, limit=available_slots + )🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 187 - 207, The worker currently calls self.dal.claim_conversations() without limiting how many rows are returned, so _try_claim_and_dispatch can over-dispatch; compute available_slots = CONVERSATION_WORKER_MAX_CONCURRENT - len(self._active_conversation_ids) after grabbing the lock, and pass that value into self.dal.claim_conversations(holmes_id, limit=available_slots) (and update any DAL/RPC signature if needed), then only iterate and submit up to available_slots tasks while still adding each task.conversation_id to self._active_conversation_ids before calling self._executor.submit(self._process_conversation_safe, task).
222-225:⚠️ Potential issue | 🟠 MajorAvoid logging the full conversation row on parse failure.
Dumping
convinto the exception message can leak conversation metadata into logs. Log stable identifiers only and keep the stack trace.Suggested fix
except Exception: logging.exception( - "Failed to build conversation task from row: %s", conv, exc_info=True + "Failed to build conversation task from row (conversation_id=%s, request_sequence=%s)", + conv.get("conversation_id"), + conv.get("request_sequence"), + exc_info=True, ) return None🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 222 - 225, The exception handler in worker.py currently logs the entire conversation row object `conv`; change it to log only stable identifier(s) instead (e.g., `conv.get("id")`, `conv.get("conversation_id")`, or whichever unique key exists) while preserving the stack trace via exc_info=True; locate the except block that calls logging.exception("Failed to build conversation task from row: %s", conv, exc_info=True) and replace the message/argument to include only the identifier(s) and any minimal context, removing the full `conv` payload to avoid leaking metadata.holmes/core/conversations_worker/realtime_manager.py (1)
196-199:⚠️ Potential issue | 🟡 MinorLog suppressed
on_new_pending()callback failures.These blocks still swallow callback exceptions entirely, which makes missed wake-ups and reconnect problems much harder to diagnose. Keep suppressing the callback failure if needed, but log it with context before returning.
Also applies to: 297-300, 306-309
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 196 - 199, The try/except blocks that call callback methods (e.g., self.on_new_pending() in RealtimeManager) are swallowing all exceptions; change each to catch Exception as e and log the error with context before returning (for example use self.logger.exception or process_logger.exception with a message like "on_new_pending callback failed" plus the exception), keeping the suppression behavior but ensuring failures are recorded; apply the same change to the other similar blocks noted (the other try/excepts around callback invocations).
🧹 Nitpick comments (1)
holmes/core/conversations_worker/worker.py (1)
61-66: Add the missing return annotation on__init__.This constructor is the odd one out in the file and currently misses an explicit
-> None.As per coding guidelines, "Type hints are required throughout the codebase".
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 61 - 66, The __init__ method is missing its return type annotation; update the constructor signature for the method named __init__ (the one taking parameters dal, config, chat_function in this file) to include an explicit "-> None" return annotation so it follows the project's type-hinting guidelines.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 214-229: The realtime connection only calls
AsyncRealtimeClient.set_auth(user_jwt) once, so when the Supabase session
refreshes the JWT becomes stale; update the logic to monitor token/session
changes and re-call set_auth with the new access token on refresh or reconnect.
Specifically, add a watcher/subscription that gets the latest token from
self.dal.client.auth.get_session().access_token (or the DAL's session change
hook) and invoke self._client.set_auth(new_jwt) whenever the session updates or
after auto-reconnect; ensure errors from set_auth are handled and retried, and
tie this behavior to the existing AsyncRealtimeClient instance stored in
self._client so it persists across reconnects.
In `@holmes/core/conversations_worker/worker.py`:
- Around line 141-145: Add inline comments in the ConversationWorker
stop/teardown section explaining the shutdown strategy: next to
self._executor.shutdown(wait=False) note that we intentionally call shutdown
with wait=False to perform a non-blocking, graceful shutdown that prevents
blocking the main thread while allowing already-submitted conversation tasks to
run to completion (or be killed by process exit), and next to
self._claim_thread.join(timeout=5) explain this gives a short grace period for
the claim loop to exit cleanly without blocking indefinitely; also add a final
note that daemon threads or process exit will terminate any remaining in-flight
conversations if they haven’t finished.
---
Duplicate comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 196-199: The try/except blocks that call callback methods (e.g.,
self.on_new_pending() in RealtimeManager) are swallowing all exceptions; change
each to catch Exception as e and log the error with context before returning
(for example use self.logger.exception or process_logger.exception with a
message like "on_new_pending callback failed" plus the exception), keeping the
suppression behavior but ensuring failures are recorded; apply the same change
to the other similar blocks noted (the other try/excepts around callback
invocations).
In `@holmes/core/conversations_worker/worker.py`:
- Around line 187-207: The worker currently calls self.dal.claim_conversations()
without limiting how many rows are returned, so _try_claim_and_dispatch can
over-dispatch; compute available_slots = CONVERSATION_WORKER_MAX_CONCURRENT -
len(self._active_conversation_ids) after grabbing the lock, and pass that value
into self.dal.claim_conversations(holmes_id, limit=available_slots) (and update
any DAL/RPC signature if needed), then only iterate and submit up to
available_slots tasks while still adding each task.conversation_id to
self._active_conversation_ids before calling
self._executor.submit(self._process_conversation_safe, task).
- Around line 222-225: The exception handler in worker.py currently logs the
entire conversation row object `conv`; change it to log only stable
identifier(s) instead (e.g., `conv.get("id")`, `conv.get("conversation_id")`, or
whichever unique key exists) while preserving the stack trace via exc_info=True;
locate the except block that calls logging.exception("Failed to build
conversation task from row: %s", conv, exc_info=True) and replace the
message/argument to include only the identifier(s) and any minimal context,
removing the full `conv` payload to avoid leaking metadata.
---
Nitpick comments:
In `@holmes/core/conversations_worker/worker.py`:
- Around line 61-66: The __init__ method is missing its return type annotation;
update the constructor signature for the method named __init__ (the one taking
parameters dal, config, chat_function in this file) to include an explicit "->
None" return annotation so it follows the project's type-hinting guidelines.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 95bbd573-414a-42ba-b647-e9c57f5e70f9
📒 Files selected for processing (3)
holmes/core/conversations_worker/event_publisher.pyholmes/core/conversations_worker/realtime_manager.pyholmes/core/conversations_worker/worker.py
🚧 Files skipped from review as they are similar to previous changes (1)
- holmes/core/conversations_worker/event_publisher.py
Supabase access tokens default to 1h validity and the Python client auto-refreshes its stored session in the background, but the realtime WebSocket doesn't pick up rotated tokens automatically. Without an explicit push, RLS-scoped Postgres Changes subscriptions silently stop delivering events once the original JWT expires, leaving the worker stuck in a connected-but-silent state. The RealtimeManager main loop now polls the DAL's session once a minute and re-calls set_auth on the underlying client only when the JWT has actually rotated. The last-pushed JWT is tracked on the instance and reset on start() for clean restarts. Also add inline comments to ConversationWorker.stop() explaining why executor shutdown is non-blocking and why the claim-thread join is bounded — avoiding shutdown hangs on long-running LLM streams. Signed-off-by: Claude <noreply@anthropic.com>
Implement the queued→running two-phase lifecycle for conversations: - claim_conversations now transitions to queued (not running) - Presence is joined immediately for queued conversations with status="queued" in the payload, updated to "running" on dispatch - Worker claims ALL pending conversations eagerly (no capacity gate) and queues them locally; only dispatches up to MAX_CONCURRENT at a time, transitioning each to running via update_conversation_status - When a conversation finishes, the next queued task is dispatched - MISMATCH errors from update_conversation_status are promoted to ConversationReassignedError in the DAL, caught in _dispatch_queued so conversations stopped/reassigned while queued are skipped cleanly - complete_conversation replaced with update_conversation_status RPC which accepts queued/running/completed/failed and only clears assignee on terminal states - Default CONVERSATION_WORKER_MAX_CONCURRENT lowered from 10 to 5 - QUEUED added to ConversationStatus enum 55 unit tests pass (11 new). Signed-off-by: Claude <noreply@anthropic.com>
Add supports_realtime_conversations field to the HolmesMetadata dataclass, set from ENABLE_CONVERSATION_WORKER. This lets initiators (Relay, Frontend) query the HolmesStatus table to decide whether to route conversations through Supabase or the existing /api/chat endpoint. Signed-off-by: Claude <noreply@anthropic.com>
There was a problem hiding this comment.
Actionable comments posted: 2
♻️ Duplicate comments (1)
holmes/core/conversations_worker/worker.py (1)
323-327:⚠️ Potential issue | 🟡 MinorAvoid logging full conversation row on parse failure.
The logging at line 325 includes the entire
convpayload (%s", conv), which may contain sensitive metadata. Log only stable identifiers to avoid leaking sensitive data.🔧 Proposed fix
except Exception: logging.exception( - "Failed to build conversation task from row: %s", conv, exc_info=True + "Failed to build conversation task from row (conversation_id=%s, request_sequence=%s)", + conv.get("conversation_id"), + conv.get("request_sequence"), + exc_info=True, ) return None🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 323 - 327, The exception handler in the conversation task builder is logging the entire conv row (the variable conv) which can leak sensitive data; update the except block in the function that "builds conversation task from row" (the code around the logging.exception call) to log only stable identifiers (e.g., conv.get("id") or conv.get("conversation_id") and conv.get("tenant_id") if present) instead of the full conv payload, preserving logging.exception(..., exc_info=True) and the error context.
🧹 Nitpick comments (6)
holmes/core/conversations_worker/realtime_manager.py (3)
366-370: Log the exception instead of silently passing.Another silent exception swallow that would benefit from debug logging.
🔧 Proposed fix
try: self.on_new_pending() except Exception: - pass + logging.debug("Error in on_new_pending callback on channel error", exc_info=True)🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 366 - 370, The try/except around the call to self.on_new_pending() in realtime_manager.py currently swallows all exceptions; change it to catch Exception as e and log the error (e.g., using self.logger.exception or self.logger.error(..., exc_info=True)) so the exception and stack trace are recorded instead of being silently ignored; keep behavior otherwise the same so the claim loop still falls back to polling.
358-361: Log the exception instead of silently passing.Similar to above, the callback error is silently swallowed. Consider logging at debug level.
🔧 Proposed fix
try: self.on_new_pending() except Exception: - pass + logging.debug("Error in on_new_pending callback after subscribe", exc_info=True)🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 358 - 361, The call to self.on_new_pending() in realtime_manager is swallowing exceptions; change the except block to log the exception at debug level instead of pass — catch Exception as e and call the class logger (e.g., self.logger.debug) with a short message, the exception, and exc_info=True; if self.logger is not available use the module logger (logging.getLogger(__name__)) so failures in on_new_pending are recorded for debugging while preserving the current behavior.
230-234: Log the exception instead of silently passing.This
try-except-passblock in thefinallyclause silently swallows exceptions fromon_new_pending(). While the intent is to not let callback errors disrupt shutdown, logging aids debugging.🔧 Proposed fix
try: self.on_new_pending() except Exception: - pass + logging.debug("Error in on_new_pending callback during shutdown", exc_info=True)🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 230 - 234, The finally clause currently swallows exceptions from on_new_pending() (try: self.on_new_pending() except Exception: pass); replace the silent pass with logging the exception so shutdown isn't disrupted but failures are visible—catch the exception as e and call a logger (e.g., self.logger.exception("Error running on_new_pending during shutdown") or self.logger.error(..., exc_info=True)) so the stacktrace is recorded while preserving the swallow-on-shutdown behavior.tests/core/conversations_worker/test_dal_contract.py (1)
15-28: Consider handlingrpc_data=Nonedefault more explicitly.When
rpc_dataisNone(the default), the mock'sexecute()method is not configured, which means any test accidentally callingdal.client.rpc(...).execute()without setting uprpc_datawill get an unconfiguredMagicMock. This works because the tests that userpc_data=Nonealso setdal.enabled = False, but it's fragile.♻️ Suggested improvement for clarity
def _build_dal(rpc_data: Any = None) -> SupabaseDal: """Build a DAL with a mocked supabase client whose rpc().execute() returns a result with the given data payload.""" dal = SupabaseDal.__new__(SupabaseDal) dal.enabled = True dal.account_id = "acc-1" dal.cluster = "cluster-1" dal.client = MagicMock() dal.client.rpc = MagicMock() - if rpc_data is not None: - dal.client.rpc.return_value = MagicMock( - execute=MagicMock(return_value=MagicMock(data=rpc_data)) - ) + # Always configure the mock chain; rpc_data=None simulates "no data" from RPC + dal.client.rpc.return_value = MagicMock( + execute=MagicMock(return_value=MagicMock(data=rpc_data)) + ) return dal🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/core/conversations_worker/test_dal_contract.py` around lines 15 - 28, The _build_dal helper currently leaves dal.client.rpc(...).execute() unconfigured when rpc_data is None which is fragile; update _build_dal (the SupabaseDal construction) so that dal.client.rpc.return_value is always a mock with an execute attribute — e.g., when rpc_data is None still set dal.client.rpc.return_value to a MagicMock whose execute returns a MagicMock(data=None) (or an explicit empty result), ensuring dal.client.rpc(...).execute() is deterministic regardless of dal.enabled.holmes/core/supabase_dal.py (1)
1021-1025: Move import to module scope.The
ConversationReassignedErrorimport is inside the exception handler. This violates the coding guideline requiring imports at the top of the file.♻️ Proposed fix
Add import at top of file (near line 40):
from holmes.core.conversations_worker.models import ConversationReassignedErrorThen simplify the exception handler:
except Exception as e: # The RPC raises MISMATCH errors when assignee, request_sequence, # or status guards fail — propagate these so the worker can exit # cleanly rather than retrying a stale transition. if "mismatch" in str(e).lower(): - from holmes.core.conversations_worker.models import ( - ConversationReassignedError, - ) - raise ConversationReassignedError(str(e)) from eAs per coding guidelines, "ALWAYS place Python imports at the top of the file, not inside functions or methods."
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/supabase_dal.py` around lines 1021 - 1025, The exception handler currently imports ConversationReassignedError locally and then re-raises it; move the import to module scope by adding "from holmes.core.conversations_worker.models import ConversationReassignedError" near the top of the file and remove the inline import inside the except block so the handler simply does "raise ConversationReassignedError(str(e)) from e" (locate the raise in the same except block where ConversationReassignedError is referenced).holmes/core/conversations_worker/worker.py (1)
302-308: Consider logging at debug level instead of silently passing.While the docstring explicitly states "swallowing any errors," logging at debug level would aid troubleshooting without disrupting the flow.
🔧 Proposed fix
def _leave_presence_quietly(self, conversation_id: str) -> None: """Leave conversation presence, swallowing any errors.""" if self._realtime_manager is not None: try: self._realtime_manager.leave_conversation_presence(conversation_id) except Exception: - pass + logging.debug( + "Error leaving presence for %s (ignored)", conversation_id, exc_info=True + )🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 302 - 308, The _leave_presence_quietly method currently swallows all exceptions; update it to log any caught exceptions at debug level so failures are visible for troubleshooting: inside _leave_presence_quietly (which calls self._realtime_manager.leave_conversation_presence(conversation_id)), catch Exception as e and if a logger exists (e.g. self._logger or self.logger) call its debug method with a clear message including conversation_id and include the exception info (using exc_info=True or formatting e) before continuing to suppress the exception.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Line 97: The call currently passes *args after a keyword argument
(ws_connect(url, sock=sock, *args, **kwargs)), which is unsafe; change the
argument order so positional args are expanded before any keyword arguments
(e.g., invoke ws_connect(url, *args, sock=sock, **kwargs)) or eliminate
positional args entirely by forwarding only keywords; update the call site
referencing ws_connect, url, sock, *args and **kwargs accordingly so *args is
placed before the sock= keyword.
- Around line 276-279: The code reads user_jwt =
self.dal.client.auth.get_session().access_token which can raise AttributeError
if get_session() returns None; update the connection path in realtime_manager.py
to guard against a None session by calling self._maybe_refresh_auth() (or
otherwise refreshing auth) before accessing .access_token, re-fetching session
after the refresh and only accessing .access_token when session is not None;
ensure apikey/user_jwt handling and any downstream logic in the initial
connection path cope with a missing token (e.g., log and proceed or raise a
clear error) and reference the existing methods/attributes:
self.dal.client.auth.get_session(), _maybe_refresh_auth, set_auth, apikey, and
user_jwt.
---
Duplicate comments:
In `@holmes/core/conversations_worker/worker.py`:
- Around line 323-327: The exception handler in the conversation task builder is
logging the entire conv row (the variable conv) which can leak sensitive data;
update the except block in the function that "builds conversation task from row"
(the code around the logging.exception call) to log only stable identifiers
(e.g., conv.get("id") or conv.get("conversation_id") and conv.get("tenant_id")
if present) instead of the full conv payload, preserving logging.exception(...,
exc_info=True) and the error context.
---
Nitpick comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 366-370: The try/except around the call to self.on_new_pending()
in realtime_manager.py currently swallows all exceptions; change it to catch
Exception as e and log the error (e.g., using self.logger.exception or
self.logger.error(..., exc_info=True)) so the exception and stack trace are
recorded instead of being silently ignored; keep behavior otherwise the same so
the claim loop still falls back to polling.
- Around line 358-361: The call to self.on_new_pending() in realtime_manager is
swallowing exceptions; change the except block to log the exception at debug
level instead of pass — catch Exception as e and call the class logger (e.g.,
self.logger.debug) with a short message, the exception, and exc_info=True; if
self.logger is not available use the module logger (logging.getLogger(__name__))
so failures in on_new_pending are recorded for debugging while preserving the
current behavior.
- Around line 230-234: The finally clause currently swallows exceptions from
on_new_pending() (try: self.on_new_pending() except Exception: pass); replace
the silent pass with logging the exception so shutdown isn't disrupted but
failures are visible—catch the exception as e and call a logger (e.g.,
self.logger.exception("Error running on_new_pending during shutdown") or
self.logger.error(..., exc_info=True)) so the stacktrace is recorded while
preserving the swallow-on-shutdown behavior.
In `@holmes/core/conversations_worker/worker.py`:
- Around line 302-308: The _leave_presence_quietly method currently swallows all
exceptions; update it to log any caught exceptions at debug level so failures
are visible for troubleshooting: inside _leave_presence_quietly (which calls
self._realtime_manager.leave_conversation_presence(conversation_id)), catch
Exception as e and if a logger exists (e.g. self._logger or self.logger) call
its debug method with a clear message including conversation_id and include the
exception info (using exc_info=True or formatting e) before continuing to
suppress the exception.
In `@holmes/core/supabase_dal.py`:
- Around line 1021-1025: The exception handler currently imports
ConversationReassignedError locally and then re-raises it; move the import to
module scope by adding "from holmes.core.conversations_worker.models import
ConversationReassignedError" near the top of the file and remove the inline
import inside the except block so the handler simply does "raise
ConversationReassignedError(str(e)) from e" (locate the raise in the same except
block where ConversationReassignedError is referenced).
In `@tests/core/conversations_worker/test_dal_contract.py`:
- Around line 15-28: The _build_dal helper currently leaves
dal.client.rpc(...).execute() unconfigured when rpc_data is None which is
fragile; update _build_dal (the SupabaseDal construction) so that
dal.client.rpc.return_value is always a mock with an execute attribute — e.g.,
when rpc_data is None still set dal.client.rpc.return_value to a MagicMock whose
execute returns a MagicMock(data=None) (or an explicit empty result), ensuring
dal.client.rpc(...).execute() is deterministic regardless of dal.enabled.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 4465a789-34b0-4dd4-af34-1c892581b2aa
📒 Files selected for processing (8)
holmes/common/env_vars.pyholmes/core/conversations_worker/event_publisher.pyholmes/core/conversations_worker/models.pyholmes/core/conversations_worker/realtime_manager.pyholmes/core/conversations_worker/worker.pyholmes/core/supabase_dal.pytests/core/conversations_worker/test_dal_contract.pytests/core/conversations_worker/test_worker_lifecycle.py
🚧 Files skipped from review as they are similar to previous changes (4)
- holmes/core/conversations_worker/event_publisher.py
- holmes/common/env_vars.py
- tests/core/conversations_worker/test_worker_lifecycle.py
- holmes/core/conversations_worker/models.py
Two fixes in realtime_manager.py: 1. The initial connection path called get_session().access_token without checking for a None session, which would raise AttributeError if the Supabase auth session wasn't established yet. Now guards with an explicit None check and logs a warning instead of crashing. 2. The proxied ws_connect call had *args after a keyword argument (sock=sock), which is against convention. Reordered to place *args before keyword args. Signed-off-by: Claude <noreply@anthropic.com>
Every failure path now posts an error event to ConversationEvents before
marking the conversation as failed, so subscribers (Frontend, Relay) can
see the failure reason in the event stream:
- _process_conversation_safe: unexpected exceptions → error event with
the exception message, then update_conversation_status("failed")
- _process_conversation: no user question → error event, then failed
- _run_chat_and_publish: stream ended without terminal → error event
Extracted _post_error_event and _fail_conversation helpers to avoid
duplication. ConversationReassignedError still does NOT post an error
event (the conversation state is owned by whoever reassigned it).
Signed-off-by: Claude <noreply@anthropic.com>
There was a problem hiding this comment.
♻️ Duplicate comments (2)
holmes/core/conversations_worker/realtime_manager.py (2)
230-233:⚠️ Potential issue | 🟡 MinorDon’t silently swallow callback errors in shutdown wake-up path.
on_new_pending()failures are currently hidden, which makes reconnect/fallback debugging harder.Proposed fix
try: self.on_new_pending() - except Exception: - pass + except Exception: + logging.exception( + "Error in on_new_pending callback during realtime loop shutdown", + exc_info=True, + )🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 230 - 233, The shutdown wake-up path currently swallows all exceptions from the callback call to self.on_new_pending(), hiding failures; change the except block to catch Exception as e and log the full exception (use self.logger.exception("on_new_pending failed during shutdown wake-up") if a logger exists on the class, otherwise use logging.getLogger(__name__).exception) so the stacktrace and error are preserved for debugging, but keep execution flowing (do not re-raise) to maintain shutdown semantics.
365-368:⚠️ Potential issue | 🟡 MinorLog callback failures in subscribe status handler instead of ignoring them.
Both callback invocations suppress exceptions with
pass, which can hide production signaling issues.Proposed fix
try: self.on_new_pending() - except Exception: - pass + except Exception: + logging.exception( + "Error in on_new_pending callback after SUBSCRIBED status", + exc_info=True, + ) @@ try: self.on_new_pending() - except Exception: - pass + except Exception: + logging.exception( + "Error in on_new_pending callback after channel failure status", + exc_info=True, + )Also applies to: 374-377
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 365 - 368, The subscribe-status handler currently swallows exceptions from callback invocations (e.g., the call to self.on_new_pending()) using bare except: pass; change these to catch Exception as e and log the failure with full details (using self.logger.exception(...) or self.logger.error(..., exc_info=True)) so callback errors are visible; apply the same change to the other similar callback invocation in that handler to ensure all callback failures are recorded with context.
🧹 Nitpick comments (1)
holmes/core/conversations_worker/realtime_manager.py (1)
122-123: Add explicit attribute types for realtime client/channel fields.These instance attributes are initialized as
Nonewithout annotations, which weakens mypy checks in a typed module.Proposed fix
-from typing import Any, Callable, Dict, Optional, TYPE_CHECKING +from typing import Any, Callable, Dict, Optional, TYPE_CHECKING @@ if TYPE_CHECKING: from holmes.core.supabase_dal import SupabaseDal + from realtime._async.client import AsyncRealtimeChannel @@ - self._client = None - self._cluster_channel = None + self._client: Optional[AsyncRealtimeClient] = None + self._cluster_channel: Optional["AsyncRealtimeChannel"] = NoneAs per coding guidelines, "Type hints are required throughout the codebase (mypy configuration in pyproject.toml)".
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 122 - 123, The two instance attributes self._client and self._cluster_channel are left un-annotated which weakens mypy; update the class __init__ to give them explicit Optional types (e.g. self._client: Optional[RealtimeClient] = None and self._cluster_channel: Optional[ClusterChannel] = None), add the necessary typing import (from typing import Optional) and import or forward-reference the concrete types (RealtimeClient, ClusterChannel) used by this module; if imports would cause circular dependencies, use TYPE_CHECKING guards or string type annotations to avoid runtime imports.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Duplicate comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 230-233: The shutdown wake-up path currently swallows all
exceptions from the callback call to self.on_new_pending(), hiding failures;
change the except block to catch Exception as e and log the full exception (use
self.logger.exception("on_new_pending failed during shutdown wake-up") if a
logger exists on the class, otherwise use logging.getLogger(__name__).exception)
so the stacktrace and error are preserved for debugging, but keep execution
flowing (do not re-raise) to maintain shutdown semantics.
- Around line 365-368: The subscribe-status handler currently swallows
exceptions from callback invocations (e.g., the call to self.on_new_pending())
using bare except: pass; change these to catch Exception as e and log the
failure with full details (using self.logger.exception(...) or
self.logger.error(..., exc_info=True)) so callback errors are visible; apply the
same change to the other similar callback invocation in that handler to ensure
all callback failures are recorded with context.
---
Nitpick comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 122-123: The two instance attributes self._client and
self._cluster_channel are left un-annotated which weakens mypy; update the class
__init__ to give them explicit Optional types (e.g. self._client:
Optional[RealtimeClient] = None and self._cluster_channel:
Optional[ClusterChannel] = None), add the necessary typing import (from typing
import Optional) and import or forward-reference the concrete types
(RealtimeClient, ClusterChannel) used by this module; if imports would cause
circular dependencies, use TYPE_CHECKING guards or string type annotations to
avoid runtime imports.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 3a73185c-0e68-4828-9379-8f0e50da9016
📒 Files selected for processing (1)
holmes/core/conversations_worker/realtime_manager.py
- _REALTIME_CONNECTED_POLL_SECONDS (hardcoded 120) → CONVERSATION_WORKER_POLL_INTERVAL_SECONDS_WITH_REALTIME env var (default 300s / 5 minutes) - Hardcoded max_backoff=120 in realtime_manager.py → uses CONVERSATION_WORKER_REALTIME_RECONNECT_MAX_SECONDS env var (default 120s, already existed but was unused) - Removed unused env vars: CONVERSATION_WORKER_HEARTBEAT_INTERVAL_SECONDS (never imported) https://claude.ai/code/session_013q1t5sReMaC3SboNv3LZ1m Signed-off-by: Claude <noreply@anthropic.com>
…1855/git/HolmesGPT/holmesgpt into claude/implement-m2-holmes-bzoqT
…ssage Three fields from the user_message spec were not being extracted from ConversationEvents or passed to ChatRequest: - frontend_tools: tool definitions from the frontend that Holmes exposes to the LLM but delegates execution of - response_format: optional JSON schema for structured output - behavior_controls: generic prompt component overrides (replaces the individual bash_enabled/fast_mode fields) Added all three to ConversationTask, extraction in _hydrate_task_from_events, and pass-through to ChatRequest. _run_chat_and_publish now uses chat_request.behavior_controls directly instead of reconstructing from individual task fields. https://claude.ai/code/session_013q1t5sReMaC3SboNv3LZ1m Signed-off-by: Claude <noreply@anthropic.com>
…roring ChatRequest fields ConversationTask no longer duplicates every ChatRequest field (ask, images, model, tool_decisions, etc.). Instead it stores the raw user_message data dict and ChatRequest is constructed directly from it in _process_conversation. This means new user_message fields (frontend_tools, response_format, behavior_controls, or future additions) are automatically available without updating ConversationTask. Hydration now just stores user_message_data + conversation_history. The field-by-field extraction and mapping are eliminated. https://claude.ai/code/session_013q1t5sReMaC3SboNv3LZ1m Signed-off-by: Claude <noreply@anthropic.com>
Previously, if a conversation was "pending" but had a terminal event (ai_answer_end / approval_required) after the latest user_message — meaning that user_message was already processed — the worker would silently re-run the stale question. This could happen if the DB state or a retry path got out of sync. Now hydration detects this state: if a terminal event appears after the latest user_message, user_message_data is left empty. The existing "no ask" guard in _process_conversation then fails the conversation cleanly with a diagnostic error event, matching /api/chat's defensive posture (Pydantic would reject a similar request with 422). New test file test_worker_edge_cases.py covers: - Empty events list - user_message with empty-string ask - user_message with no ask / no tool_decisions / no frontend_tool_results - user_message already answered by ai_answer_end - user_message already answered by approval_required - Positive control: resume-only followup with tool_decisions https://claude.ai/code/session_013q1t5sReMaC3SboNv3LZ1m Signed-off-by: Claude <noreply@anthropic.com>
There was a problem hiding this comment.
Actionable comments posted: 1
♻️ Duplicate comments (6)
holmes/core/conversations_worker/worker.py (4)
396-398:⚠️ Potential issue | 🟠 MajorScope hydration to
task.request_sequence.Line 397 loads a flat event stream for the whole conversation, and
_hydrate_task_from_events()then picks the latestuser_messageglobally. Sinceget_conversation_events()flattens events across request sequences, a newer follow-up can cause this worker to hydrate and process the wrong turn. Filter by exactrequest_sequencein the DAL/RPC result before hydration.Suggested direction
- events = self.dal.get_conversation_events(task.conversation_id) + events = self.dal.get_conversation_events( + task.conversation_id, + request_sequence=task.request_sequence, + ) self._hydrate_task_from_events(task, events)If the RPC cannot return one sequence today, extend it to include
seqper event or accept an exact sequence parameter;min_seqalone is not enough because it can still include later turns.Also applies to: 451-483
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 396 - 398, The worker currently loads all conversation events via get_conversation_events and then calls _hydrate_task_from_events, which lets a later request_sequence override the intended turn; change the data retrieval so get_conversation_events (or its RPC) is called/scoped with the exact task.request_sequence (or returns per-event seq values and the call filters by seq) so only events matching task.request_sequence are returned to _hydrate_task_from_events; update the DAL/RPC signature to accept an exact sequence parameter or include seq on each event and filter client-side before hydration so _process_conversation and the _hydrate_task_from_events call operate only on the intended sequence.
212-217:⚠️ Potential issue | 🟠 MajorAvoid claiming more conversations than this worker can actually run.
Line 217 claims all pending rows and parks the overflow in this process. In multi-worker deployments, one instance can monopolize work, and if it exits after claiming, those rows remain assigned/queued instead of being available to other workers. Prefer passing a remaining-capacity or bounded-prefetch limit into
claim_conversations().Suggested direction
- # Claim ALL pending conversations — they transition to queued state. - # There is no capacity check here: we claim eagerly so that no other - # Holmes instance can grab them, and queue them locally until executor - # slots open up. - claimed = self.dal.claim_conversations(self.holmes_id) + with self._active_lock: + active = len(self._active_conversation_ids) + with self._queued_lock: + queued = len(self._queued_tasks) + available = CONVERSATION_WORKER_MAX_CONCURRENT - active - queued + if available <= 0: + return + claimed = self.dal.claim_conversations(self.holmes_id, limit=available)🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 212 - 217, The worker currently calls claim_conversations(self.holmes_id) in _try_claim_and_dispatch which eagerly claims all pending conversations; change this to compute the worker's remaining capacity (e.g., max_workers or executor_slots minus current active/queued tasks) and pass that limit into claim_conversations (e.g., claim_conversations(self.holmes_id, limit)) so only up to remaining-capacity are claimed; update the DAL method signature (claim_conversations) and any callers accordingly to support a bounded-prefetch/remaining_capacity parameter to avoid monopolizing work across workers.
253-265:⚠️ Potential issue | 🟠 MajorRequeue the task when the queued→running transition returns
False.Line 265 drops the popped task on a
Falseresult. Per the DAL contract,Falsecan come from non-mismatch failures, leaving the row queued/assigned in Supabase but no longer present in_queued_tasks. Requeue and stop this dispatch pass so it can retry later.Suggested fix
if not ok: logging.warning( - "Failed to transition conversation %s to running — skipping", + "Failed to transition conversation %s to running — requeuing", task.conversation_id, ) - continue + with self._queued_lock: + self._queued_tasks.appendleft(task) + break🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 253 - 265, The block that calls self.dal.update_conversation_status (used to transition queued→running) currently drops the popped task when it returns False; instead, if update_conversation_status returns False you must requeue the task back into self._queued_tasks and stop the current dispatch pass so it can be retried later. Concretely: when self.dal.update_conversation_status(conversation_id=task.conversation_id, request_sequence=task.request_sequence, assignee=self.holmes_id, status="running") yields False, push the same task back onto self._queued_tasks (or equivalent queue structure) and break/return from the dispatch loop rather than continue, preserving the queued/assigned row in the DAL for a later retry.
63-91: 🛠️ Refactor suggestion | 🟠 MajorComplete the missing type annotations.
Line 68 is missing
-> None, and Lines 86, 91, and 495 still use bare collection types. Please parameterize these so mypy can validate worker state and history handling. As per coding guidelines, "Type hints are required throughout the codebase (mypy configuration in pyproject.toml)."Suggested typing cleanup
def __init__( self, dal: "SupabaseDal", config: "Config", chat_function: ChatFunction, - ): + ) -> None: @@ - self._active_conversation_ids: set = set() + self._active_conversation_ids: set[str] = set() @@ - self._queued_tasks: deque = deque() + self._queued_tasks: deque[ConversationTask] = deque() @@ - def _extract_last_user_ask(history: Optional[list]) -> Optional[str]: + def _extract_last_user_ask( + history: Optional[List[Dict[str, Any]]], + ) -> Optional[str]:Also applies to: 494-495
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/worker.py` around lines 63 - 91, The __init__ method needs an explicit return type and several bare collection types must be parameterized: add "-> None" to the def __init__(...) signature, import typing aliases (from typing import Optional, Set, Deque, Any, List) and change the attributes to typed collections—e.g. _active_conversation_ids: Set[str] = set(), _queued_tasks: Deque[Any] = deque(), and keep _executor: Optional[ThreadPoolExecutor] and _active_lock: threading.Lock as-is; also locate the history variables referenced around lines 494-495 and replace their bare list/dict types with typed counterparts (e.g. List[Any] or List[dict] or a concrete HistoryRecord type) so mypy can validate them.holmes/core/conversations_worker/realtime_manager.py (2)
487-521:⚠️ Potential issue | 🟠 MajorFail subscribe setup unless a real
SUBSCRIBEDack arrives.The error branches set
subscribedand timeout only logs, so_connect_and_subscribe()can return successfully afterCHANNEL_ERROR,CLOSED,TIMED_OUT, or no ack. That makes_full_reconnect()reset backoff even though realtime is not subscribed. Track a success flag and raise on timeout/error so the existing reconnect backoff handles it.Suggested pattern for each subscribe method
subscribed = asyncio.Event() +subscribe_ok = False def _on_subscribe(status: Any, err: Optional[Exception] = None) -> None: + nonlocal subscribe_ok logging.info("Broadcast subscribe status=%s err=%s", status, err) status_str = str(status).upper() if "SUBSCRIBED" in status_str: + subscribe_ok = True self._connected = True subscribed.set() ... elif any( s in status_str for s in ("CHANNEL_ERROR", "CLOSED", "TIMED_OUT") ): self._connected = False subscribed.set() ... await self._channel.subscribe(_on_subscribe) try: await asyncio.wait_for(subscribed.wait(), timeout=5) except asyncio.TimeoutError: logging.warning("Timed out waiting for broadcast subscribe ack") + raise RuntimeError("Timed out waiting for broadcast subscribe ack") +if not subscribe_ok: + raise RuntimeError("Realtime channel subscribe failed")Also applies to: 555-589
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 487 - 521, The subscribe handler currently marks the local Event and returns on any status including errors, letting _connect_and_subscribe() succeed even when not truly subscribed; modify the _on_subscribe callback and surrounding logic (the asyncio.Event "subscribed", the _on_subscribe function, and the await self._channel.subscribe(...) / wait_for(...) block) to track a success flag (e.g., subscribed_success) and only set it to True when status indicates a real "SUBSCRIBED" ack; on CHANNEL_ERROR/CLOSED/TIMED_OUT set subscribed_success to False (and keep self._connected False) and after wait_for completes raise an exception if subscribed_success is not True so _connect_and_subscribe() fails and upstream reconnect/backoff logic runs as intended.
460-484:⚠️ Potential issue | 🟠 MajorFilter Postgres-change wakeups to pending rows for this cluster.
Line 471 still subscribes account-wide, and Line 467 wakes the worker without checking
cluster_idorstatus. In pgchanges mode, unrelated updates for the same account can wake every cluster worker and trigger unnecessary claim cycles. Filter in the subscription if supported, or guard inside_on_pg_change()before callingon_new_pending().Suggested guard
def _on_pg_change(payload: Dict[str, Any]) -> None: try: change = payload.get("data", {}) or {} + record = change.get("record") or change.get("new") or change + if ( + record.get("cluster_id") != self.dal.cluster + or record.get("status") != "pending" + ): + return logging.info( "RealtimeManager: Postgres change notification: %s", change.get("type"), )🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@holmes/core/conversations_worker/realtime_manager.py` around lines 460 - 484, The subscription currently wakes on any account-level INSERT/UPDATE; narrow it to only pending rows for this cluster by either adding cluster_id and status to the filter passed to on_postgres_changes (augment account_id_filter with cluster_id=eq.{self.cluster_id} and status=eq.pending) or, if the driver/filtering doesn't support that, modify the _on_pg_change callback to inspect the change payload (inspect change.get("data"/"new"/"old") for matching cluster_id and status == "pending") and only call self.on_new_pending() when both account_id and cluster_id match this worker and status indicates a pending row; reference symbols: _on_pg_change, on_postgres_changes, account_id_filter, on_new_pending, CONVERSATIONS_TABLE, self.dal.account_id, self.cluster_id.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@tests/core/conversations_worker/test_worker_edge_cases.py`:
- Around line 23-84: Add explicit type hints to the untyped helpers and tests:
annotate _bare_worker to return ConversationWorker (or ConversationWorker | Any
if import constraints), annotate _task to return ConversationTask, add parameter
and return types for _run_process (e.g., worker: ConversationWorker, task:
ConversationTask, events: Sequence[dict] -> MagicMock or AsyncMock return type
as appropriate), and add -> None to test functions. Update imports to include
the concrete types (ConversationWorker, ConversationTask, Sequence) or use
typing.Any/typing.Sequence where necessary so mypy passes; ensure signatures
referenced in the file ( _bare_worker, _task, _run_process and the test
functions between the following block) are updated consistently.
---
Duplicate comments:
In `@holmes/core/conversations_worker/realtime_manager.py`:
- Around line 487-521: The subscribe handler currently marks the local Event and
returns on any status including errors, letting _connect_and_subscribe() succeed
even when not truly subscribed; modify the _on_subscribe callback and
surrounding logic (the asyncio.Event "subscribed", the _on_subscribe function,
and the await self._channel.subscribe(...) / wait_for(...) block) to track a
success flag (e.g., subscribed_success) and only set it to True when status
indicates a real "SUBSCRIBED" ack; on CHANNEL_ERROR/CLOSED/TIMED_OUT set
subscribed_success to False (and keep self._connected False) and after wait_for
completes raise an exception if subscribed_success is not True so
_connect_and_subscribe() fails and upstream reconnect/backoff logic runs as
intended.
- Around line 460-484: The subscription currently wakes on any account-level
INSERT/UPDATE; narrow it to only pending rows for this cluster by either adding
cluster_id and status to the filter passed to on_postgres_changes (augment
account_id_filter with cluster_id=eq.{self.cluster_id} and status=eq.pending)
or, if the driver/filtering doesn't support that, modify the _on_pg_change
callback to inspect the change payload (inspect change.get("data"/"new"/"old")
for matching cluster_id and status == "pending") and only call
self.on_new_pending() when both account_id and cluster_id match this worker and
status indicates a pending row; reference symbols: _on_pg_change,
on_postgres_changes, account_id_filter, on_new_pending, CONVERSATIONS_TABLE,
self.dal.account_id, self.cluster_id.
In `@holmes/core/conversations_worker/worker.py`:
- Around line 396-398: The worker currently loads all conversation events via
get_conversation_events and then calls _hydrate_task_from_events, which lets a
later request_sequence override the intended turn; change the data retrieval so
get_conversation_events (or its RPC) is called/scoped with the exact
task.request_sequence (or returns per-event seq values and the call filters by
seq) so only events matching task.request_sequence are returned to
_hydrate_task_from_events; update the DAL/RPC signature to accept an exact
sequence parameter or include seq on each event and filter client-side before
hydration so _process_conversation and the _hydrate_task_from_events call
operate only on the intended sequence.
- Around line 212-217: The worker currently calls
claim_conversations(self.holmes_id) in _try_claim_and_dispatch which eagerly
claims all pending conversations; change this to compute the worker's remaining
capacity (e.g., max_workers or executor_slots minus current active/queued tasks)
and pass that limit into claim_conversations (e.g.,
claim_conversations(self.holmes_id, limit)) so only up to remaining-capacity are
claimed; update the DAL method signature (claim_conversations) and any callers
accordingly to support a bounded-prefetch/remaining_capacity parameter to avoid
monopolizing work across workers.
- Around line 253-265: The block that calls self.dal.update_conversation_status
(used to transition queued→running) currently drops the popped task when it
returns False; instead, if update_conversation_status returns False you must
requeue the task back into self._queued_tasks and stop the current dispatch pass
so it can be retried later. Concretely: when
self.dal.update_conversation_status(conversation_id=task.conversation_id,
request_sequence=task.request_sequence, assignee=self.holmes_id,
status="running") yields False, push the same task back onto self._queued_tasks
(or equivalent queue structure) and break/return from the dispatch loop rather
than continue, preserving the queued/assigned row in the DAL for a later retry.
- Around line 63-91: The __init__ method needs an explicit return type and
several bare collection types must be parameterized: add "-> None" to the def
__init__(...) signature, import typing aliases (from typing import Optional,
Set, Deque, Any, List) and change the attributes to typed collections—e.g.
_active_conversation_ids: Set[str] = set(), _queued_tasks: Deque[Any] = deque(),
and keep _executor: Optional[ThreadPoolExecutor] and _active_lock:
threading.Lock as-is; also locate the history variables referenced around lines
494-495 and replace their bare list/dict types with typed counterparts (e.g.
List[Any] or List[dict] or a concrete HistoryRecord type) so mypy can validate
them.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 32118767-f059-439c-9780-ea77243bdc5e
📒 Files selected for processing (8)
holmes/common/env_vars.pyholmes/core/conversations_worker/models.pyholmes/core/conversations_worker/realtime_manager.pyholmes/core/conversations_worker/worker.pytests/core/conversations_worker/integration/test_conversation_integration.pytests/core/conversations_worker/test_worker_edge_cases.pytests/core/conversations_worker/test_worker_hydration.pytests/core/conversations_worker/test_worker_polling.py
✅ Files skipped from review due to trivial changes (2)
- holmes/common/env_vars.py
- tests/core/conversations_worker/integration/test_conversation_integration.py
🚧 Files skipped from review as they are similar to previous changes (2)
- tests/core/conversations_worker/test_worker_polling.py
- tests/core/conversations_worker/test_worker_hydration.py
There was a problem hiding this comment.
Actionable comments posted: 1
🧹 Nitpick comments (1)
tests/core/conversations_worker/test_worker_polling.py (1)
8-44: Add type annotations to helpers and test functions.
_make_worker_with_rtand alltest_*functions lack type hints. As per coding guidelines, "Type hints are required throughout the codebase (mypy configuration in pyproject.toml)."🧩 Suggested annotations
-def _make_worker_with_rt(connected: bool): +def _make_worker_with_rt(connected: bool) -> ConversationWorker: worker = ConversationWorker.__new__(ConversationWorker) rt = MagicMock() rt.is_connected.return_value = connected worker._realtime_manager = rt return worker -def test_realtime_connected_returns_true_when_manager_connected(): +def test_realtime_connected_returns_true_when_manager_connected() -> None: ... -def test_realtime_connected_false_when_no_manager(): +def test_realtime_connected_false_when_no_manager() -> None: ... -def test_realtime_connected_false_when_manager_disconnected(): +def test_realtime_connected_false_when_manager_disconnected() -> None: ... -def test_realtime_connected_false_when_is_connected_raises(): +def test_realtime_connected_false_when_is_connected_raises() -> None: ... -def test_connected_poll_is_reasonable_safety_net(): +def test_connected_poll_is_reasonable_safety_net() -> None: ...🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/core/conversations_worker/test_worker_polling.py` around lines 8 - 44, Add missing type annotations for the helper and test functions: annotate _make_worker_with_rt(connected: bool) -> ConversationWorker, add return types (-> None) and parameter types where applicable for all test_* functions (e.g., test_realtime_connected_returns_true_when_manager_connected() -> None, etc.), and type the rt variable as MagicMock | Any if needed; ensure ConversationWorker and CONVERSATION_WORKER_POLL_INTERVAL_SECONDS_WITH_REALTIME references remain unchanged and imports/types are resolved so mypy passes.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@holmes/core/conversations_worker/worker.py`:
- Around line 471-492: The terminal_events tuple currently omits "error",
causing errored turns to be treated as non-terminal and allowing stale user
messages to be re-run; update terminal_events in worker.py to include "error"
(i.e., terminal_events = ("ai_answer_end", "approval_required", "error")) so
that already_answered correctly detects an errored completion and
task.user_message_data is not re-hydrated, and add a unit test exercising the
case where the last user_message is followed by an "error" event to assert that
already_answered is True and _process_conversation does not re-run the LLM path;
also add a short inline comment by terminal_events explaining why "error" is
considered terminal.
---
Nitpick comments:
In `@tests/core/conversations_worker/test_worker_polling.py`:
- Around line 8-44: Add missing type annotations for the helper and test
functions: annotate _make_worker_with_rt(connected: bool) -> ConversationWorker,
add return types (-> None) and parameter types where applicable for all test_*
functions (e.g., test_realtime_connected_returns_true_when_manager_connected()
-> None, etc.), and type the rt variable as MagicMock | Any if needed; ensure
ConversationWorker and CONVERSATION_WORKER_POLL_INTERVAL_SECONDS_WITH_REALTIME
references remain unchanged and imports/types are resolved so mypy passes.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 19b687e8-5362-41a3-b512-2b3a9bc6bbcf
📒 Files selected for processing (8)
holmes/common/env_vars.pyholmes/core/conversations_worker/models.pyholmes/core/conversations_worker/realtime_manager.pyholmes/core/conversations_worker/worker.pytests/core/conversations_worker/integration/test_conversation_integration.pytests/core/conversations_worker/test_worker_edge_cases.pytests/core/conversations_worker/test_worker_hydration.pytests/core/conversations_worker/test_worker_polling.py
✅ Files skipped from review due to trivial changes (2)
- holmes/common/env_vars.py
- tests/core/conversations_worker/integration/test_conversation_integration.py
🚧 Files skipped from review as they are similar to previous changes (2)
- tests/core/conversations_worker/test_worker_hydration.py
- holmes/core/conversations_worker/models.py
|
@claude review once |
There was a problem hiding this comment.
Given the scope (new ~2k-line async worker subsystem touching realtime WebSocket handling, thread-pool dispatch, JWT refresh, and event-batch compaction), this really wants a human pass — I've left inline comments on the issues I spotted (the only non-nit is the realtime reconnect backoff that's computed but never awaited in the main loop, so a Supabase outage busy-loops reconnect attempts). Not approving.
Extended reasoning...
Overview\n\nThis PR introduces an entirely new M2 Conversation Worker subsystem with ~2000 lines of production code across holmes/core/conversations_worker/ (worker.py, realtime_manager.py, event_publisher.py, models.py), plus extensions to supabase_dal.py, server.py lifecycle wiring, new env vars, and a sizable unit + integration test suite. The subsystem claims pending Conversations rows atomically, dispatches them to a bounded ThreadPoolExecutor, streams results through the existing chat pipeline, and batches/compacts ConversationEvents back to Supabase. It also manages a Supabase Realtime subscription (broadcast or pgchanges mode) in a background asyncio thread, with JWT auth refresh, WebSocket proxy support, and exponential-backoff reconnect.\n\n### Security risks\n\nMostly adjacent rather than direct. The realtime path handles JWT re-sign-in and set_auth for long-lived WebSockets; the broadcast channel is scoped by account_id:cluster_id (correct, but depends on the server-side topic ACL); the HTTP CONNECT proxy patch monkey-patches realtime._async.client.connect based on env vars; and error-event posting was specifically scrubbed so raw exceptions don't leak to client-visible events. None of these are new auth/crypto primitives, but the surface area is meaningful and the RLS/SECURITY DEFINER contract with the Supabase RPCs needs trust.\n\n### Level of scrutiny\n\nHigh. This is production-critical infrastructure: an always-on background worker that owns a significant chunk of the user-facing conversation lifecycle. Concurrency correctness matters (capacity gating, dispatch/shutdown races, presence-sequence TOCTOU, reconnect backoff under outage), and regressions would manifest as stuck conversations, duplicate work across replicas, or silent polling fallback. The change is too large and too central for bot approval — a human who owns this subsystem should sign off.\n\n### Other factors\n\nTest coverage is substantial (unit tests for publisher, hydration, lifecycle, polling, edge cases, DAL contract; integration tests against real Supabase). CodeRabbit has iterated extensively with the author and most of its findings are marked addressed. The one real bug I'd flag (already posted inline) is the reconnect-backoff gap in realtime_manager._run() — the else branch computes backoff and logs it but never awaits _async_stop.wait(timeout=backoff), so after _full_reconnect() sets _channel = None, _channel_needs_reconnect() immediately returns True on the next loop iteration and reconnect attempts fire as fast as the network will reject them. The remaining inline comments are nits.
… import - 🔴 Reconnect backoff actually applied: after _full_reconnect() fails, await asyncio.wait_for(_async_stop.wait(), timeout=backoff) now sleeps for the computed backoff before retrying. Previously the loop would immediately re-enter _channel_needs_reconnect() (channel is None after failed reconnect), busy-looping on sign_in + WS connect. - 🟡 Orphaned queued tasks: when _build_task_from_conversation_row returns None (parse failure), the conversation was stuck in queued state forever. Now posts an error event and marks it as failed. - 🟡 Stale comment: CONVERSATION_WORKER_USE_REALTIME_BROADCAST default is True (broadcast mode), not False. Comment corrected. - 🟡 supabase_dal.py: hoisted `import time` to top-level (stdlib, no circular risk). Kept ConversationReassignedError as inline import with a comment — top-level import creates a circular chain via conversations_worker.__init__ → worker → config → supabase_dal. https://claude.ai/code/session_013q1t5sReMaC3SboNv3LZ1m Signed-off-by: Claude <noreply@anthropic.com>
Conflicts resolved: - supabase_dal.py: keep both CONVERSATIONS_TABLE and OAUTH_TOKENS_TABLE - holmes_status.py: keep both realtime metadata fields and namespace field - pyproject.toml: keep both conversation_worker and manual markers Signed-off-by: Claude <noreply@anthropic.com>
…2477/git/HolmesGPT/holmesgpt into claude/implement-m2-holmes-bzoqT
The master merge included the Skills PR (#1953) which renamed Config.get_runbook_catalog() to get_skill_catalog() and the build_chat_messages parameter from runbooks= to skills=. https://claude.ai/code/session_013q1t5sReMaC3SboNv3LZ1m Signed-off-by: Claude <noreply@anthropic.com>
Summary
Introduces the M2 Conversation Worker, a new background service that asynchronously processes pending conversations from Supabase. The worker claims conversations, runs them through the existing chat pipeline, and publishes results as real-time events back to the database.
Key Changes
Core Worker Implementation
holmes/core/conversations_worker/worker.py): Main worker class that:Real-time Event Publishing
holmes/core/conversations_worker/event_publisher.py): Batches and publishes stream events:Real-time Subscriptions
holmes/core/conversations_worker/realtime_manager.py): Manages Supabase Realtime connections:Data Models
holmes/core/conversations_worker/models.py): Represents a claimed conversation with extracted metadata, user ask, history, and configurationDatabase Integration
SupabaseDalwith M2-specific RPC methods:claim_conversations(): Atomically claim pending conversationspost_conversation_events(): Batch post events with sequence trackingcomplete_conversation(): Mark conversation as completed/failedget_conversation_events(): Fetch conversation event historyConfiguration
holmes/common/env_vars.py:ENABLE_CONVERSATION_WORKER: Feature flagCONVERSATION_WORKER_MAX_CONCURRENT: Max concurrent conversations (default: 10)CONVERSATION_WORKER_POLL_INTERVAL_SECONDS: Polling interval (default: 5s)CONVERSATION_WORKER_EVENT_BATCH_INTERVAL_SECONDS: Event batch interval (default: 0.5s)CONVERSATION_WORKER_REALTIME_ENABLED: Enable real-time subscriptionsServer Integration
server.pyto start/stop ConversationWorker on application lifecycleNotable Implementation Details
Tests
https://claude.ai/code/session_013q1t5sReMaC3SboNv3LZ1m
Summary by CodeRabbit
New Features
Server
DAL
Models
Tests
Chores