From c3a973bd0972968862072e692f0320b8a5bff0be Mon Sep 17 00:00:00 2001 From: Christian Date: Thu, 30 Jul 2026 12:21:03 +0200 Subject: [PATCH 1/3] Publish analysis results to platform for BYOK contribution Add contribution publisher and flow_orchestrator hook so backend analyses can write analysis.result events for platform ingest. Publisher lives in report_analyst_jobs/contribution.py for now; placement may move after review. --- report_analyst_jobs/contribution.py | 55 ++++++++++++++ .../flow_orchestrator.py | 42 ++++++++++- tests/test_contribution.py | 74 +++++++++++++++++++ 3 files changed, 167 insertions(+), 4 deletions(-) create mode 100644 report_analyst_jobs/contribution.py create mode 100644 tests/test_contribution.py diff --git a/report_analyst_jobs/contribution.py b/report_analyst_jobs/contribution.py new file mode 100644 index 00000000..ac667b35 --- /dev/null +++ b/report_analyst_jobs/contribution.py @@ -0,0 +1,55 @@ +"""Publish analysis results to platform (BYOK contribution mode).""" + +from __future__ import annotations + +import json +import logging +import os +import uuid +from typing import Any, Dict, Optional + +import nats + +logger = logging.getLogger(__name__) + + +async def publish_analysis_result( + resource_id: str, + results: Dict[str, Any], + provenance: Optional[Dict[str, Any]] = None, + analysis_config: Optional[Dict[str, Any]] = None, + owner_user_id: Optional[str] = None, + duration_ms: Optional[int] = None, +) -> str: + """Publish analysis.result.{id} for platform to persist.""" + request_id = str(uuid.uuid4()) + nats_url = os.getenv("NATS_URL", "nats://localhost:4222") + nats_token = os.getenv("NATS_TOKEN") + if nats_token and "@" not in nats_url: + protocol, rest = nats_url.split("://", 1) + nats_url = f"{protocol}://{nats_token}@{rest}" + + payload = { + "request_id": request_id, + "resource_id": resource_id, + "results": results, + "results_summary": results, + "provenance": provenance or {}, + "analysis_config": analysis_config or {}, + "owner_user_id": owner_user_id, + "duration_ms": duration_ms, + "source": os.getenv("NATS_USER", "report-analyst"), + } + + nc = await nats.connect(nats_url, connect_timeout=15) + try: + js = nc.jetstream() + subject = f"analysis.result.{request_id}" + try: + await js.publish(subject, json.dumps(payload).encode()) + except Exception: + await nc.publish(subject, json.dumps(payload).encode()) + logger.info("Published %s for resource %s", subject, resource_id) + return request_id + finally: + await nc.close() diff --git a/report_analyst_search_backend/flow_orchestrator.py b/report_analyst_search_backend/flow_orchestrator.py index 389301fd..064df446 100644 --- a/report_analyst_search_backend/flow_orchestrator.py +++ b/report_analyst_search_backend/flow_orchestrator.py @@ -258,11 +258,45 @@ def _analyze_local_with_features(self, chunks: List[Dict[str, Any]], questions: return self._analyze_local(chunks, questions) def _analyze_enhanced(self, chunks: List[Dict[str, Any]], questions: List[str]) -> AnalysisResult: - """Enhanced analysis with centralized LLM and data lake""" - # This would use NATS LLM and store in data lake - # For now, fallback to local analysis + """Enhanced analysis with platform contribution publish when backend is enabled.""" + import asyncio + import os + st.info("Enhanced analysis not fully implemented - using local analysis") - return self._analyze_local(chunks, questions) + result = self._analyze_local(chunks, questions) + if not (result.success and self.config.use_backend and result.results): + return result + + try: + from report_analyst_jobs.contribution import publish_analysis_result + + resource_id = chunks[0].get("resource_id") if chunks else None + if resource_id: + questions_list = result.results.get("questions", questions) + answers_list = result.results.get("answers", []) + asyncio.run( + publish_analysis_result( + resource_id=str(resource_id), + results={ + "answers": answers_list, + "questions": questions_list, + }, + provenance={ + "model": os.getenv("OPENAI_API_MODEL", "local"), + "provider": "report_analyst", + "mode": "byok_contribution" + if os.getenv("USE_BYOK_CONTRIBUTION", "").lower() in ("1", "true", "yes") + else "centralized", + }, + ) + ) + except ImportError: + pass + except Exception as exc: + logger.warning("Could not publish analysis result to platform: %s", exc) + + result.stored_in_backend = self.config.use_backend + return result def _configure_question_set(self, default_question_set: str) -> str: """Configure question set for backend analysis""" diff --git a/tests/test_contribution.py b/tests/test_contribution.py new file mode 100644 index 00000000..d1915459 --- /dev/null +++ b/tests/test_contribution.py @@ -0,0 +1,74 @@ +"""Unit tests for BYOK contribution publish (analysis.result NATS events).""" + +from __future__ import annotations + +import json +from unittest.mock import AsyncMock, MagicMock + +import pytest + +from report_analyst_jobs.contribution import publish_analysis_result + + +@pytest.mark.asyncio +async def test_publish_analysis_result_uses_analysis_result_subject_and_payload(monkeypatch): + mock_nc = MagicMock() + mock_js = AsyncMock() + mock_nc.jetstream.return_value = mock_js + mock_nc.close = AsyncMock() + + async def fake_connect(url, **kwargs): + return mock_nc + + monkeypatch.setenv("NATS_URL", "nats://localhost:4222") + monkeypatch.setenv("NATS_USER", "report-analyst-test") + monkeypatch.setattr("report_analyst_jobs.contribution.nats.connect", fake_connect) + + request_id = await publish_analysis_result( + resource_id="res-123", + results={"answers": ["answer one"], "questions": ["question one"]}, + provenance={"mode": "byok_contribution", "provider": "report_analyst"}, + owner_user_id="user-42", + duration_ms=1500, + ) + + mock_js.publish.assert_awaited_once() + subject, payload_bytes = mock_js.publish.call_args[0] + assert subject == f"analysis.result.{request_id}" + payload = json.loads(payload_bytes.decode()) + assert payload["request_id"] == request_id + assert payload["resource_id"] == "res-123" + assert payload["results"] == {"answers": ["answer one"], "questions": ["question one"]} + assert payload["results_summary"] == payload["results"] + assert payload["provenance"]["mode"] == "byok_contribution" + assert payload["owner_user_id"] == "user-42" + assert payload["duration_ms"] == 1500 + assert payload["source"] == "report-analyst-test" + mock_nc.close.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_publish_analysis_result_falls_back_to_core_publish(monkeypatch): + mock_nc = MagicMock() + mock_js = AsyncMock() + mock_js.publish.side_effect = RuntimeError("jetstream unavailable") + mock_nc.jetstream.return_value = mock_js + mock_nc.close = AsyncMock() + mock_nc.publish = AsyncMock() + + async def fake_connect(url, **kwargs): + return mock_nc + + monkeypatch.setenv("NATS_URL", "nats://localhost:4222") + monkeypatch.setattr("report_analyst_jobs.contribution.nats.connect", fake_connect) + + request_id = await publish_analysis_result( + resource_id="res-456", + results={"answers": ["a"], "questions": ["q"]}, + ) + + mock_nc.publish.assert_awaited_once() + subject, payload_bytes = mock_nc.publish.call_args[0] + assert subject == f"analysis.result.{request_id}" + payload = json.loads(payload_bytes.decode()) + assert payload["resource_id"] == "res-456" From 2aa32c025b0ab314df7b0d695f9adb6a889d23f8 Mon Sep 17 00:00:00 2001 From: Christian Date: Thu, 30 Jul 2026 15:16:04 +0200 Subject: [PATCH 2/3] Move BYOK contribution publisher into enterprise package Publisher lives in report_analyst_enterprise; search-backend keeps a thin optional import so open-core runs without the enterprise module. --- .../contribution.py | 5 ++++- report_analyst_search_backend/flow_orchestrator.py | 3 ++- tests/test_contribution.py | 8 ++++---- 3 files changed, 10 insertions(+), 6 deletions(-) rename {report_analyst_jobs => report_analyst_enterprise}/contribution.py (91%) diff --git a/report_analyst_jobs/contribution.py b/report_analyst_enterprise/contribution.py similarity index 91% rename from report_analyst_jobs/contribution.py rename to report_analyst_enterprise/contribution.py index ac667b35..69d94656 100644 --- a/report_analyst_jobs/contribution.py +++ b/report_analyst_enterprise/contribution.py @@ -1,4 +1,7 @@ -"""Publish analysis results to platform (BYOK contribution mode).""" +"""Publish analysis results to platform (BYOK contribution / enterprise). + +Enterprise-only: open-core callers import this optionally and no-op if absent. +""" from __future__ import annotations diff --git a/report_analyst_search_backend/flow_orchestrator.py b/report_analyst_search_backend/flow_orchestrator.py index 064df446..f348e15c 100644 --- a/report_analyst_search_backend/flow_orchestrator.py +++ b/report_analyst_search_backend/flow_orchestrator.py @@ -268,7 +268,7 @@ def _analyze_enhanced(self, chunks: List[Dict[str, Any]], questions: List[str]) return result try: - from report_analyst_jobs.contribution import publish_analysis_result + from report_analyst_enterprise.contribution import publish_analysis_result resource_id = chunks[0].get("resource_id") if chunks else None if resource_id: @@ -291,6 +291,7 @@ def _analyze_enhanced(self, chunks: List[Dict[str, Any]], questions: List[str]) ) ) except ImportError: + # Enterprise package not installed — contribution publish is optional. pass except Exception as exc: logger.warning("Could not publish analysis result to platform: %s", exc) diff --git a/tests/test_contribution.py b/tests/test_contribution.py index d1915459..93c59a6f 100644 --- a/tests/test_contribution.py +++ b/tests/test_contribution.py @@ -1,4 +1,4 @@ -"""Unit tests for BYOK contribution publish (analysis.result NATS events).""" +"""Unit tests for enterprise BYOK contribution publish (analysis.result NATS events).""" from __future__ import annotations @@ -7,7 +7,7 @@ import pytest -from report_analyst_jobs.contribution import publish_analysis_result +from report_analyst_enterprise.contribution import publish_analysis_result @pytest.mark.asyncio @@ -22,7 +22,7 @@ async def fake_connect(url, **kwargs): monkeypatch.setenv("NATS_URL", "nats://localhost:4222") monkeypatch.setenv("NATS_USER", "report-analyst-test") - monkeypatch.setattr("report_analyst_jobs.contribution.nats.connect", fake_connect) + monkeypatch.setattr("report_analyst_enterprise.contribution.nats.connect", fake_connect) request_id = await publish_analysis_result( resource_id="res-123", @@ -60,7 +60,7 @@ async def fake_connect(url, **kwargs): return mock_nc monkeypatch.setenv("NATS_URL", "nats://localhost:4222") - monkeypatch.setattr("report_analyst_jobs.contribution.nats.connect", fake_connect) + monkeypatch.setattr("report_analyst_enterprise.contribution.nats.connect", fake_connect) request_id = await publish_analysis_result( resource_id="res-456", From a1efee2a9e267e78e553e6b935888fb1682abe7d Mon Sep 17 00:00:00 2001 From: Christian Date: Thu, 20 Aug 2026 12:55:05 +0200 Subject: [PATCH 3/3] Fix black formatting and ruff BLE001 on BYOK publish path. - Parenthesize mode ternary so black --check and test_linting pass - Silence intentional broad excepts; use iterable unpack for options list --- report_analyst_enterprise/contribution.py | 2 +- .../flow_orchestrator.py | 18 ++++++++++-------- 2 files changed, 11 insertions(+), 9 deletions(-) diff --git a/report_analyst_enterprise/contribution.py b/report_analyst_enterprise/contribution.py index 69d94656..c9a6cccc 100644 --- a/report_analyst_enterprise/contribution.py +++ b/report_analyst_enterprise/contribution.py @@ -50,7 +50,7 @@ async def publish_analysis_result( subject = f"analysis.result.{request_id}" try: await js.publish(subject, json.dumps(payload).encode()) - except Exception: + except Exception: # noqa: BLE001 — JetStream may be absent; fall back to core NATS await nc.publish(subject, json.dumps(payload).encode()) logger.info("Published %s for resource %s", subject, resource_id) return request_id diff --git a/report_analyst_search_backend/flow_orchestrator.py b/report_analyst_search_backend/flow_orchestrator.py index f348e15c..f7acf81a 100644 --- a/report_analyst_search_backend/flow_orchestrator.py +++ b/report_analyst_search_backend/flow_orchestrator.py @@ -88,7 +88,7 @@ def process_document(self, uploaded_file) -> ProcessingResult: return self._process_complete_backend(uploaded_file) else: return ProcessingResult(success=False, error=f"Unknown flow type: {flow_type}") - except Exception as e: + except Exception as e: # noqa: BLE001 logger.error(f"Document processing failed: {e}") return ProcessingResult(success=False, error=str(e)) @@ -114,7 +114,7 @@ def analyze_document(self, chunks: List[Dict[str, Any]], questions: List[str]) - return self._analyze_enhanced(chunks, questions) else: return AnalysisResult(success=False, error=f"Analysis not supported for flow: {flow_type}") - except Exception as e: + except Exception as e: # noqa: BLE001 logger.error(f"Document analysis failed: {e}") return AnalysisResult(success=False, error=str(e)) @@ -284,16 +284,18 @@ def _analyze_enhanced(self, chunks: List[Dict[str, Any]], questions: List[str]) provenance={ "model": os.getenv("OPENAI_API_MODEL", "local"), "provider": "report_analyst", - "mode": "byok_contribution" - if os.getenv("USE_BYOK_CONTRIBUTION", "").lower() in ("1", "true", "yes") - else "centralized", + "mode": ( + "byok_contribution" + if os.getenv("USE_BYOK_CONTRIBUTION", "").lower() in ("1", "true", "yes") + else "centralized" + ), }, ) ) except ImportError: # Enterprise package not installed — contribution publish is optional. pass - except Exception as exc: + except Exception as exc: # noqa: BLE001 — contribution publish must not fail analysis logger.warning("Could not publish analysis result to platform: %s", exc) result.stored_in_backend = self.config.use_backend @@ -305,7 +307,7 @@ def _configure_question_set(self, default_question_set: str) -> str: # Get dynamic question set options if QUESTION_LOADER_AVAILABLE: - question_set_options = question_loader.get_question_set_options() + ["custom"] + question_set_options = [*question_loader.get_question_set_options(), "custom"] # Calculate index for default question set try: index = question_set_options.index(default_question_set) if default_question_set in question_set_options else 0 @@ -503,7 +505,7 @@ async def external_service_analysis( analysis_job_id=request_id, ) - except Exception as e: + except Exception as e: # noqa: BLE001 logger.error(f"External service analysis failed: {e}") return AnalysisResult(success=False, error=str(e))