-
Notifications
You must be signed in to change notification settings - Fork 513
ROB-3670 - Add M2 Conversation Worker for async conversation processing #1903
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from all commits
Commits
Show all changes
83 commits
Select commit
Hold shift + click to select a range
27d53e2
feat(m2): add ConversationWorker for supabase-backed conversations
claude 201c789
feat(m2): route realtime WebSocket through HTTPS proxy when configured
claude 95ebc00
feat(m2): handle resume-only followups + add coverage for all event t…
claude 7b7f36f
feat(m2): wire per-conversation Presence + fix cluster Presence config
claude 370c913
refactor(m2): tune defaults, gate polling on realtime, simplify DAL
claude 97643a3
refactor(m2): align with updated plan (assignee param, global seq, RP…
claude 4afa6d2
fix(m2): pass _account_id (now first param) to all RPC calls
claude c6ba742
refactor(m2): five small correctness/robustness fixes
claude 71cf001
refactor(m2): seven small code-quality fixes
claude e3323ac
refactor(m2): periodic JWT refresh + document shutdown strategy
claude d45ad7e
feat(m2): add queued state, update_conversation_status, max_concurrent=5
claude 421a106
Merge remote-tracking branch 'origin/master' into claude/implement-m2…
claude 3a13d96
feat(m2): advertise supports_realtime_conversations in HolmesStatus
claude 60b4577
fix(m2): guard against None session in realtime connect, fix arg order
claude 029c6ae
feat(m2): post error events to ConversationEvents before failing
claude c207a83
test(m2): add conversation worker integration test suite
claude 707a59f
feat(m2): compact prior events on ai_answer_end and approval_required
claude 0d1be7c
fix(m2): sanitize error events, guard dispatch during shutdown
claude 30a1bc9
Merge remote-tracking branch 'origin/master' into claude/implement-m2…
claude 01bd215
fix(m2): fixture cleanup, replace sleep with polling, dynamic stress NUM
claude 67c0ad8
fix(m2): add dispatch lock to serialize dispatch with shutdown
claude 61896bb
fix(m2): don't drop events on flush failure, retry on None seq
claude 9a9dd14
refactor(m2): _terminal_to_status returns 'completed' for approval
claude 9794256
fix(m2): sticky compact flag, consume fails on unsaved events, retry …
claude 79685c0
fix(m2): skip final flush after ConversationReassignedError
claude ada3703
Merge branch 'master' into claude/implement-m2-holmes-bzoqT
naomi-robusta c483972
feat(m2): race-safe presence, topic helpers, config knobs, flush cleanup
claude b78538b
Merge remote-tracking branch 'origin/claude/implement-m2-holmes-bzoqT…
claude 0cea647
refactor(m2): include account_id in conversation presence topic
claude c6b2dfc
fix(m2): prune presence map on leave, gate conv subscribe on SUBSCRIBED
claude 64950d4
Merge remote-tracking branch 'origin/master' into claude/implement-m2…
claude 648041b
refactor(m2): single per-account presence channel with debounced tracks
claude 7718f6d
fix(m2): split presence and pg-changes channels + retry transient DAL…
claude 2e16992
refactor(m2): per-account pg-changes topic, drop cluster liveness
claude 4fbf31e
Merge remote-tracking branch 'origin/master' into claude/implement-m2…
claude 2f5b476
style(m2): move in-function imports to module top in test_realtime_ma…
claude 50e8271
fix(m2): fix proxy URL builder for missing port/password, strengthen …
claude 7c0aad3
refactor(m2): only advertise presence when conversations are active
claude b4a0aed
fix(m2): log suppressed on_new_pending exceptions at debug level
claude 66162e5
fix(m2): atomic check-and-pop in leave_conversation_presence
claude 2ff2398
refactor(m2): remove all presence code, keep only pg-changes
claude 9c65f0e
style: replace EN DASH with hyphen-minus in test docstring
claude 2f64ce1
Merge remote-tracking branch 'origin/master' into claude/implement-m2…
claude 86760d7
feat(m2): support both pgchanges and broadcast for conversation submi…
claude c235758
feat(m2): flip env var to CONVERSATION_WORKER_USE_REALTIME_BROADCAST,
claude 323e80d
fix(m2): reuse persistent WS for broadcast in integration tests
claude dbb1881
fix(m2): log teardown exceptions with tracebacks instead of silencing
claude ac71d38
fix(m2): wake async sleep on stop, requeue transient failures, move i…
claude 1a626f0
Merge branch 'master' into claude/implement-m2-holmes-bzoqT
naomi-robusta aa109ea
rename broadcast event to pending_conversations
claude 5e837a2
Merge branch 'claude/implement-m2-holmes-bzoqT' of http://127.0.0.1:5…
claude 59abcd7
comment cleanup
naomi-robusta 0e405ba
Merge branch 'claude/implement-m2-holmes-bzoqT' of github.com:HolmesG…
naomi-robusta 8290710
use broadcast by default
naomi-robusta 3d2e5f8
fix stale env var name in realtime_manager docstring
claude cffabbc
Merge branch 'claude/implement-m2-holmes-bzoqT' of http://127.0.0.1:3…
claude 57feb19
update docstring: broadcast is now the default subscription mode
claude fb6335f
drop keyword-only marker on RealtimeManager.__init__
claude 5dc75ed
close realtime client cleanly on stop() and connect failure
claude b0112c7
fix realtime URL double-suffix, broadcast setup error handling, shutd…
claude ca2d467
Merge remote-tracking branch 'origin/master' into claude/implement-m2…
claude 279409f
add test_max_concurrent_never_exceeded + review fixes
claude c9dae88
add broadcast health check diagnostic script
claude 872d9da
detect channel closure and auto-reconnect realtime subscription
claude cb56698
force sign_in on reconnect instead of relying on get_session auto-ref…
claude a9fa337
fix reconnect storm, subscribe waiter, client leak, task drop
claude 4c25b64
retry initial connect with backoff instead of exiting thread
claude a84bec0
rename misleading test: test_failed → test_successful_has_no_error
claude f607043
revert
naomi-robusta 6dabd0f
cleanup
naomi-robusta 4477a13
cleanup
naomi-robusta d5d50d9
promote poll/reconnect constants to env vars, remove unused env vars
claude 60b11e2
Merge branch 'claude/implement-m2-holmes-bzoqT' of http://127.0.0.1:3…
claude db0046a
support frontend_tools, response_format, behavior_controls in user_me…
claude 8f9fd63
simplify ConversationTask: store raw user_message_data instead of mir…
claude f791dfc
guard against processing already-answered user_message
claude 58d7fe6
Merge remote-tracking branch 'origin/master' into claude/implement-m2…
claude afd11cd
fix reconnect busy-loop, orphaned queued tasks, stale comment, inline…
claude 910fa57
cleanup
naomi-robusta aa508f1
merge origin/master: resolve conflicts keeping both sides
claude 52a1aca
Merge branch 'claude/implement-m2-holmes-bzoqT' of http://127.0.0.1:2…
claude 8755b3d
fix: get_runbook_catalog renamed to get_skill_catalog in Skills PR
claude 0ea081c
Merge branch 'master' into claude/implement-m2-holmes-bzoqT
naomi-robusta File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
Empty file.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,3 @@ | ||
| from holmes.core.conversations_worker.worker import ConversationWorker | ||
|
|
||
| __all__ = ["ConversationWorker"] |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,192 @@ | ||
| import logging | ||
| import threading | ||
| import time | ||
| from datetime import datetime, timezone | ||
| from typing import Any, Dict, Generator, List, Optional, TYPE_CHECKING | ||
|
|
||
| from holmes.core.conversations_worker.models import ConversationReassignedError | ||
| from holmes.utils.stream import StreamEvents, StreamMessage | ||
|
|
||
| if TYPE_CHECKING: | ||
| from holmes.core.supabase_dal import SupabaseDal | ||
|
|
||
|
|
||
| # Events that end a turn and must be flushed immediately. The publisher | ||
| # reports the last one it saw back to the worker so the conversation status | ||
| # can be set appropriately. | ||
| _TERMINAL_EVENTS = { | ||
| StreamEvents.ANSWER_END, | ||
| StreamEvents.APPROVAL_REQUIRED, | ||
| StreamEvents.ERROR, | ||
| } | ||
|
|
||
| # Events that should cause an immediate flush. Terminal events end a turn; | ||
| # CONVERSATION_HISTORY_COMPACTED isn't terminal but carries the same | ||
| # "history snapshot + prior events superseded" semantics so it's flushed + | ||
| # compacted with the same logic. | ||
| _FLUSH_IMMEDIATELY_EVENTS = _TERMINAL_EVENTS | { | ||
| StreamEvents.CONVERSATION_HISTORY_COMPACTED | ||
| } | ||
|
|
||
| # Events whose `messages` array carries the full conversation history | ||
| # snapshot — all prior events are superseded and should be marked compacted. | ||
| _COMPACT_ON_FLUSH_EVENTS = { | ||
|
naomi-robusta marked this conversation as resolved.
|
||
| StreamEvents.ANSWER_END, | ||
| StreamEvents.APPROVAL_REQUIRED, | ||
| StreamEvents.CONVERSATION_HISTORY_COMPACTED, | ||
| } | ||
|
|
||
|
|
||
| class ConversationEventPublisher: | ||
| """ | ||
| Consumes StreamMessage events from call_stream() and batches them | ||
| into ConversationEvents rows in Supabase. | ||
| """ | ||
|
|
||
| def __init__( | ||
| self, | ||
| dal: "SupabaseDal", | ||
| conversation_id: str, | ||
| assignee: str, | ||
| request_sequence: int, | ||
| batch_interval_seconds: float = 1.0, | ||
| ): | ||
| self.dal = dal | ||
| self.conversation_id = conversation_id | ||
| self.assignee = assignee | ||
| self.request_sequence = request_sequence | ||
| self.batch_interval_seconds = batch_interval_seconds | ||
|
|
||
| self._pending_events: List[Dict[str, Any]] = [] | ||
| self._last_flush_time: float = time.monotonic() | ||
| self._last_retry_time: float = 0.0 | ||
| self._lock = threading.Lock() | ||
|
|
||
| self._last_terminal_event: Optional[StreamEvents] = None | ||
|
|
||
| # Sticky compact flag: set when a compact flush is attempted but the | ||
| # DAL returns None. Ensures the compact intent is preserved across | ||
| # retries and the final drain. | ||
| self._pending_compact: bool = False | ||
|
|
||
| def consume( | ||
| self, | ||
| stream: Generator[StreamMessage, None, None], | ||
| ) -> Optional[StreamEvents]: | ||
| """ | ||
| Drain the stream generator, batching events and writing them to the DB. | ||
| Returns the terminal StreamEvents value observed, or None if the stream ended | ||
| without a terminal event. | ||
| Raises ConversationReassignedError if the conversation was reassigned mid-stream. | ||
| """ | ||
| reassigned = False | ||
| try: | ||
| for message in stream: | ||
| self._append_event(message) | ||
| if message.event in _TERMINAL_EVENTS: | ||
| self._last_terminal_event = message.event | ||
| # Flush on terminal events immediately, or when interval elapses | ||
| if message.event in _FLUSH_IMMEDIATELY_EVENTS: | ||
| # ai_answer_end / approval_required / compacted carry a | ||
| # full conversation history snapshot in their messages | ||
| # array, so all prior events are superseded → compact. | ||
| if message.event in _COMPACT_ON_FLUSH_EVENTS: | ||
| self._pending_compact = True | ||
| self._flush() | ||
|
naomi-robusta marked this conversation as resolved.
|
||
| elif ( | ||
|
naomi-robusta marked this conversation as resolved.
|
||
| time.monotonic() - self._last_flush_time | ||
| >= self.batch_interval_seconds | ||
| and time.monotonic() - self._last_retry_time | ||
| >= self.batch_interval_seconds | ||
| ): | ||
| self._flush() | ||
| except ConversationReassignedError: | ||
| reassigned = True | ||
| raise | ||
| finally: | ||
| # Final drain of any remaining events — skip if the conversation | ||
| # was reassigned, since our assignee/sequence are stale and writing | ||
| # would either fail or race with the new owner. | ||
| if not reassigned: | ||
| self._flush() | ||
|
|
||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
| # If events remain unsaved after the stream is fully consumed, the | ||
| # terminal batch was lost (repeated None returns). Surface this to the | ||
| # caller so the conversation is marked failed rather than completed. | ||
| with self._lock: | ||
| remaining = len(self._pending_events) | ||
| if remaining > 0: | ||
| logging.error( | ||
| "consume() finished with %d unsaved events for conversation %s", | ||
| remaining, | ||
| self.conversation_id, | ||
| ) | ||
| return None | ||
|
|
||
|
naomi-robusta marked this conversation as resolved.
|
||
| return self._last_terminal_event | ||
|
|
||
| def _append_event(self, message: StreamMessage) -> None: | ||
| with self._lock: | ||
| self._pending_events.append( | ||
| { | ||
| "event": message.event.value, | ||
| "data": message.data, | ||
| "ts": datetime.now(timezone.utc).isoformat(), | ||
| } | ||
| ) | ||
|
|
||
| def _flush(self) -> None: | ||
| with self._lock: | ||
| if not self._pending_events: | ||
| return | ||
| # Snapshot but don't clear yet — only clear after a successful post. | ||
| events_to_flush = list(self._pending_events) | ||
| compact = self._pending_compact | ||
|
|
||
| try: | ||
|
naomi-robusta marked this conversation as resolved.
|
||
| seq = self.dal.post_conversation_events( | ||
| conversation_id=self.conversation_id, | ||
| assignee=self.assignee, | ||
| request_sequence=self.request_sequence, | ||
| events=events_to_flush, | ||
| compact=compact, | ||
| ) | ||
| except ConversationReassignedError: | ||
| raise | ||
| except Exception as e: | ||
| # The RPCs prefix mismatch errors (status / assignee / request_sequence) | ||
| # with "MISMATCH " — promote those to ConversationReassignedError so | ||
| # the worker can exit the processing loop cleanly. | ||
| if "mismatch" in str(e).lower(): | ||
| raise ConversationReassignedError(str(e)) from e | ||
| raise | ||
|
|
||
| if seq is None: | ||
| # DAL returned None (disabled or unexpected empty response). | ||
| # Keep events and compact flag in memory so the next flush retries. | ||
| # Update _last_retry_time to throttle retries independently of | ||
| # normal flush timing. | ||
| self._last_retry_time = time.monotonic() | ||
| logging.warning( | ||
| "post_conversation_events returned None for conversation %s — " | ||
| "events retained for retry (%d events, compact=%s)", | ||
| self.conversation_id, | ||
| len(events_to_flush), | ||
| compact, | ||
| ) | ||
| return | ||
|
|
||
| # Success — remove the flushed events and clear the compact flag. | ||
| # New events may have been appended while the RPC was in flight, | ||
| # so we remove only the count we just posted. | ||
| with self._lock: | ||
| del self._pending_events[: len(events_to_flush)] | ||
| self._pending_compact = False | ||
| self._last_flush_time = time.monotonic() | ||
| logging.debug( | ||
| "Posted %d events to conversation %s (seq=%s, compact=%s)", | ||
| len(events_to_flush), | ||
| self.conversation_id, | ||
| seq, | ||
| compact, | ||
| ) | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,39 @@ | ||
| from enum import Enum | ||
| from typing import Any, Dict, List, Optional | ||
|
|
||
| from pydantic import BaseModel, Field | ||
|
|
||
|
|
||
| class ConversationStatus(str, Enum): | ||
| PENDING = "pending" | ||
| QUEUED = "queued" | ||
| RUNNING = "running" | ||
| COMPLETED = "completed" | ||
| FAILED = "failed" | ||
| STOPPED = "stopped" | ||
|
|
||
|
|
||
| class ConversationTask(BaseModel): | ||
| """A claimed conversation ready for processing.""" | ||
|
|
||
| conversation_id: str | ||
| account_id: str | ||
| cluster_id: str | ||
| origin: str | ||
| request_sequence: int | ||
| metadata: Dict[str, Any] = Field(default_factory=dict) | ||
| title: Optional[str] = None | ||
|
|
||
| # Raw data from the latest user_message event. Used to construct | ||
| # ChatRequest without duplicating every field. | ||
| user_message_data: Dict[str, Any] = Field(default_factory=dict) | ||
|
naomi-robusta marked this conversation as resolved.
|
||
|
|
||
| # Reconstructed from prior terminal events (ai_answer_end / approval_required). | ||
| conversation_history: Optional[List[Dict[str, Any]]] = None | ||
|
|
||
|
|
||
| class ConversationReassignedError(Exception): | ||
| """Raised when the conversation's assignee/request_sequence no longer matches ours.""" | ||
|
|
||
|
|
||
| EVENT_USER_MESSAGE = "user_message" | ||
|
naomi-robusta marked this conversation as resolved.
|
||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.