Skip to content

Event Bus Topology

HydraFlow's EventBus fans out every published event to all fan-out subscribers: an argless subscribe() / subscription() (src/events.py) returns a queue that drains the whole bus. Since ADR-0114 (#10660) a subscriber MAY instead pass an optional per-type filter — subscribe(types={EventType.X, ...}) — so its queue only receives those types (filtered at publish time, so an ignored high-volume type never fills its bounded slots). The Subscribers column shows ★ all (fan-out) for every live event, plus any typed consumer that filters on it. An event is flagged ⚠️ (likely dead) only when it is a declared EventType with no publisher (the live-only EPHEMERAL_EVENT_TYPES allowlist is exempt).

Fan-out consumers (each receives every published event):

  • src.dashboard_routes._routes:create_router.websocket_endpoint
  • src.dashboard_routes._ws_stream:_serve_merged_ws
Event Publishers Subscribers
ADR_CONFORMANCE_UPDATE src.adr_conformance_loop:AdrConformanceLoop._emit_event ★ all (fan-out)
ADR_DRAFT_OPENED src.base_runner:BaseRunner._process_transcript_for_adr_draft ★ all (fan-out)
ADVERSARIAL_STAGE_CONVERGED src.adversarial_retry_loop:AdversarialRetryLoop._emit_stage_converged ★ all (fan-out)
ADVERSARIAL_STAGE_EXHAUSTED src.adversarial_retry_loop:AdversarialRetryLoop._emit_stage_exhausted ★ all (fan-out)
ADVERSARIAL_STAGE_STARTED src.adversarial_retry_loop:AdversarialRetryLoop._emit_stage_started ★ all (fan-out)
AGENT_ACTIVITY src.runner_utils:_stream_and_collect ★ all (fan-out)
BACKGROUND_WORKER_STATUS src.base_background_loop:BaseBackgroundLoop._execute_cycle
src.base_background_loop:BaseBackgroundLoop._report_cycle_failure
src.orchestrator_bg_workers:OrchestratorBGWorkersMixin._seed_background_worker_statuses
★ all (fan-out)
BASELINE_UPDATE src.baseline_policy:BaselinePolicy.check_approval
src.baseline_policy:BaselinePolicy.rollback
★ all (fan-out)
CI_CHECK src.pr_manager_ci:PRManagerCIMixin.wait_for_ci
src.reviewer._fixes:ReviewFixMixin.fix_ci
★ all (fan-out)
CONCERN_ADDRESSED ⚠️ — —
CONCERN_FORWARDED src.adversarial_retry_loop:AdversarialRetryLoop._emit_concerns_forwarded ★ all (fan-out)
DIAGNOSTIC_UPDATE src.diagnostic_loop:DiagnosticLoop._publish_update ★ all (fan-out)
EPIC_PROGRESS src.epic._detail:EpicDetailMixin.refresh_cache ★ all (fan-out)
EPIC_READY src.epic._detail:EpicDetailMixin.refresh_cache
src.epic._merge_order:EpicMergeOrderMixin._publish_ready_event
★ all (fan-out)
EPIC_RELEASED src.epic._release:EpicReleaseMixin._execute_release ★ all (fan-out)
EPIC_RELEASING src.epic._release:EpicReleaseMixin._execute_release ★ all (fan-out)
EPIC_UPDATE src.epic._manager:EpicManager._publish_update ★ all (fan-out)
ERROR src.base_background_loop:BaseBackgroundLoop._report_cycle_failure
src.orchestrator_loops:OrchestratorLoopsMixin._polling_loop
src.orchestrator_restart:OrchestratorRestartMixin._restart_loop
★ all (fan-out)
HITL_ESCALATION src.dashboard_routes._routes:create_router.request_changes
src.review_phase._insights:ReviewInsightsMixin._escalate_to_hitl
★ all (fan-out)
HITL_UPDATE src.dashboard_routes._hitl_routes:register._resolve_hitl_item
src.dashboard_routes._hitl_routes:register.hitl_correct
src.hitl_phase:HITLPhase._process_one_hitl
src.hitl_runner:HITLRunner.run
src.pr_unsticker._unsticker:PRUnsticker.unstick
★ all (fan-out)
ISSUE_CREATED src.pr_manager_issues:PRManagerIssuesMixin.create_issue ★ all (fan-out)
ISSUE_REFINEMENT_UPDATE src.issue_refinement_loop:IssueRefinementLoop._publish_refinement_event ★ all (fan-out)
LOOP_FITNESS_UPDATE src.fitness_scorecard_loop:FitnessScorecardLoop._do_work ★ all (fan-out)
MEMORY_SYNC ⚠️ — —
MERGE_UPDATE src.pr_manager:PRManager.merge_pr
src.pr_manager_promotion:PRManagerPromotionMixin.merge_promotion_pr
★ all (fan-out)
METRICS_UPDATE src.metrics_manager:MetricsManager.sync ★ all (fan-out)
ORCHESTRATOR_STATUS src.dashboard_routes._control_routes:register.start_orchestrator
src.factory_autostart:maybe_autostart_host
src.orchestrator_stats:OrchestratorStatsMixin._publish_status
★ all (fan-out)
PHASE_CHANGE src.server:_boot_factory ★ all (fan-out)
PIPELINE_SNAPSHOT src.issue_store:IssueStore._flush_pipeline_snapshot ★ all (fan-out)
PIPELINE_STATS src.orchestrator_stats:OrchestratorStatsMixin.emit_pipeline_stats ★ all (fan-out)
PLANNER_UPDATE src.planner:PlannerRunner._emit_status ★ all (fan-out)
PR_CREATED src.pr_manager:PRManager.create_pr
src.pr_manager_promotion:PRManagerPromotionMixin.create_promotion_pr
★ all (fan-out)
QUEUE_UPDATE src.issue_store:IssueStore._publish_queue_update_nowait
src.issue_store:IssueStore.refresh
src.mockworld.fakes.fake_issue_store:FakeIssueStore.refresh
★ all (fan-out)
RATCHET_TIGHTENED src.auto_tighten_loop:AutoTightenLoop._emit_tightened
src.auto_tighten_loop:AutoTightenLoop._emit_unattributed
★ all (fan-out)
REPORT_UPDATE src.report_issue_loop:ReportIssueLoop._emit_report_event ★ all (fan-out)
RETROSPECTIVE_UPDATE src.retrospective_loop:RetrospectiveLoop._publish_update ★ all (fan-out)
REVIEW_UPDATE src.merge_conflict_resolver:MergeConflictResolver._publish_review_status
src.phase_utils:publish_review_status
src.reviewer._fixes:ReviewFixMixin.fix_review_findings
src.reviewer._runner:ReviewRunner.review
★ all (fan-out)
SESSION_END src.orchestrator_lifecycle:OrchestratorLifecycleMixin._end_session ★ all (fan-out)
SESSION_START src.orchestrator_lifecycle:OrchestratorLifecycleMixin._start_session ★ all (fan-out)
SHIPPED_WITH_KNOWN_GAP src.post_merge_handler:PostMergeHandler._maybe_emit_shipped_with_known_gap ★ all (fan-out)
SUPERVISOR_OBSERVATION src.goal_supervisor_loop:GoalSupervisorLoop._emit ★ all (fan-out)
SYSTEM_ALERT src.close_verification:reconcile_false_close
src.cost_budget_alerts:check_daily_budget
src.cost_budget_alerts:check_issue_cost
src.dependabot_merge_loop:DependabotMergeLoop._do_work
src.epic._staleness:EpicStalenessMixin.check_stale_epics
src.health_monitor_loop._errors:HealthMonitorPersistentErrorMixin._check_persistent_worker_errors
src.health_monitor_loop._freshness:HealthMonitorFreshnessMixin._check_stale_code
src.health_monitor_loop._stall:HealthMonitorStallMixin._check_worker_staleness
src.health_monitor_loop._vitals:HealthMonitorFleetVitalsMixin._run_fleet_vitals
src.merge_state_watcher_loop:MergeStateWatcherLoop._alert_on_capture_gap
src.orchestrator_credits:OrchestratorCreditsMixin._maybe_engage_failover
src.orchestrator_credits:OrchestratorCreditsMixin._pause_for_credits
src.orchestrator_credits:OrchestratorCreditsMixin._probe_claude_for_switchback
src.orchestrator_credits:OrchestratorCreditsMixin._resume_loops_after_credit_pause
src.orchestrator_lifecycle:OrchestratorLifecycleMixin._deferred_pipeline_start
src.orchestrator_loops:OrchestratorLoopsMixin._polling_loop
src.orchestrator_restart:OrchestratorRestartMixin._handle_auth_error
src.post_merge_handler:PostMergeHandler._safe_hook
src.post_merge_handler:PostMergeHandler.handle_approved
src.prompt_gate_alerts:alert_prompt_gate_block
src.review_advisor:PostVerifyAdvisor._alarm_fail_open
src.runs_gc_loop:RunsGCLoop._tend_audit_chains
src.server:_check_and_publish_boot_gap
src.staging_promotion_loop:StagingPromotionLoop._compile_evidence_pack
src.staging_promotion_loop:StagingPromotionLoop._handle_open_promotion
src.triage:TriageRunner._emit_injection_alert
src.unpushed_branch_alert:check_and_alert_unpushed_branches
★ all (fan-out)
SYSTEM_REROUTE src.review_phase._adr:AdrReviewMixin._review_single_adr
src.review_phase._adr:AdrReviewMixin._run_post_verify_advisor_for_adr
src.triage_phase:TriagePhase._flow_route
★ all (fan-out)
TRANSCRIPT_LINE src.runner_utils:_stream_and_collect
src.triage:TriageRunner._emit_transcript
src.triage_phase:TriagePhase._hold_for_blockers
★ all (fan-out)
TRANSCRIPT_SUMMARY src.transcript_summarizer:TranscriptSummarizer._summarize_and_comment_inner ★ all (fan-out)
TRIAGE_UPDATE src.triage:TriageRunner._emit_status ★ all (fan-out)
TRIBAL_PROMOTION ⚠️ — —
VERIFICATION_JUDGE src.verification_judge:VerificationJudge.judge ★ all (fan-out)
VISUAL_GATE src.post_merge_handler:PostMergeHandler._run_visual_gate
src.review_phase._visual_gate:VisualGateMixin._emit_visual_gate_telemetry
src.review_phase._visual_gate:VisualGateMixin.check_visual_gate
★ all (fan-out)
WIKI_SUPERSEDES ⚠️ — —
WORKER_UPDATE src.agent._runner:AgentRunner._emit_status ★ all (fan-out)