From 1e870c397b8c2f704652ece114ff43ddf2e4f950 Mon Sep 17 00:00:00 2001 From: Xore Date: Fri, 25 Sep 2026 12:28:27 +0200 Subject: [PATCH] feat(llm-worker): record terminal observation coverage --- docs/llm-worker/README.md | 20 ++++++++ llm-worker/tests/test_worker.py | 83 +++++++++++++++++++++++++++++++++ llm-worker/worker.py | 67 +++++++++++++++++++++++++- 3 files changed, 168 insertions(+), 2 deletions(-) diff --git a/docs/llm-worker/README.md b/docs/llm-worker/README.md index 70c12f364..7f4706f3e 100644 --- a/docs/llm-worker/README.md +++ b/docs/llm-worker/README.md @@ -109,6 +109,26 @@ itself instead of trusting a filename. - after bounded retries, an error annotation is written while raw ingestion remains unaffected. +## Session terminal observation + +Each `session_accumulator` document records a bounded `terminal_observation` +value: `close_observed` means a `cowrie.session.closed` event was captured, +`idle_finalized` means the worker finalized the session through its idle +readiness path, and `open_after_window` means the observation window ended +without a close event. `capture_coverage` is `observed` only for a captured +close and is `missing` otherwise, so a close with zero commands remains +distinguishable from no close. Existing `command_count`, `auth_success`, +`closed`, and `duration_seconds` fields retain their contract. + +The worker's session scan reports covered and excluded session counts in its +cycle status: close-only scanner sessions are intentionally excluded from the +accumulator population, so future aggregates must use that explicit denominator +and exclusions rather than treating the accumulator population as all sessions. +Terminal absence is ambiguous. It must remain unknown and must not be described +as attacker abandonment, deliberate disengagement, automation, or an +AI/automated attacker; a captured close does not supply a termination reason +that the sensor did not emit. + The cross-sensor decoder/correlation expansion remains tracked by [#154](https://github.com/Xore/APIARY/issues/154). Dashboard delivery is #150, and the customizable analyzer workbench is explicitly tracked by diff --git a/llm-worker/tests/test_worker.py b/llm-worker/tests/test_worker.py index 64fb4df09..e51389cf1 100644 --- a/llm-worker/tests/test_worker.py +++ b/llm-worker/tests/test_worker.py @@ -196,6 +196,65 @@ def test_close_and_login_state_are_recorded(self): self.assertTrue(accumulator.auth_success) self.assertTrue(accumulator.closed) self.assertEqual(accumulator.duration_seconds, 42.5) + document = accumulator.document() + self.assertEqual(document["terminal_observation"], worker.TERMINAL_OBSERVATION_CLOSE) + self.assertEqual(document["capture_coverage"], worker.CAPTURE_COVERAGE_OBSERVED) + + def test_captured_close_is_distinct_from_idle_finalization(self): + accumulator = worker.SessionAccumulator("session-fixture") + self.assertEqual(accumulator.document()["terminal_observation"], worker.TERMINAL_OBSERVATION_OPEN) + close = self.event(event_id="cowrie.session.closed") + close["_id"] = "close" + accumulator.add_event(close, 12000) + closed_document = accumulator.document() + self.assertEqual(closed_document["command_count"], 0) + self.assertEqual(closed_document["terminal_observation"], worker.TERMINAL_OBSERVATION_CLOSE) + self.assertEqual(closed_document["capture_coverage"], worker.CAPTURE_COVERAGE_OBSERVED) + + def test_idle_finalization_is_recorded_without_inventing_a_close(self): + fake_es = MagicMock() + fake_es.search.return_value = { + "hits": { + "hits": [ + { + "_id": "session-state-id", + "_source": { + "session_id": "session-fixture", + "command_count": 5, + "auth_success": False, + "closed": False, + "duration_seconds": 0.0, + "finalized": False, + }, + } + ] + } + } + llm_worker = worker.LLMWorker(config(), es=fake_es) + ready = llm_worker.ready_sessions() + self.assertEqual(len(ready), 1) + document = ready[0][1].document() + self.assertEqual(document["terminal_observation"], worker.TERMINAL_OBSERVATION_IDLE) + self.assertEqual(document["capture_coverage"], worker.CAPTURE_COVERAGE_MISSING) + self.assertFalse(document["closed"]) + self.assertEqual(document["duration_seconds"], 0.0) + + def test_still_open_after_observation_window_remains_unknown(self): + accumulator = worker.SessionAccumulator("session-fixture") + document = accumulator.document() + self.assertEqual(document["terminal_observation"], worker.TERMINAL_OBSERVATION_OPEN) + self.assertEqual(document["capture_coverage"], worker.CAPTURE_COVERAGE_MISSING) + self.assertFalse(document["closed"]) + self.assertEqual(document["command_count"], 0) + + def test_state_mapping_keeps_compatibility_and_bounded_observation_fields(self): + properties = worker.state_mapping()["mappings"]["properties"] + self.assertEqual(properties["command_count"], {"type": "integer"}) + self.assertEqual(properties["auth_success"], {"type": "boolean"}) + self.assertEqual(properties["closed"], {"type": "boolean"}) + self.assertEqual(properties["duration_seconds"], {"type": "float"}) + self.assertEqual(properties["terminal_observation"], {"type": "keyword"}) + self.assertEqual(properties["capture_coverage"], {"type": "keyword"}) def test_collection_uses_stable_cowrie_event_ids_not_container_sensor(self): fake_es = MagicMock() @@ -238,6 +297,15 @@ def test_close_only_scanner_session_does_not_create_accumulator_state(self): self.assertEqual(llm_worker.collect_session_events(), 1) fake_es.index.assert_not_called() save_checkpoint.assert_called_once_with("sessions", "2026-08-01T12:00:00Z") + self.assertEqual( + llm_worker.session_capture_coverage, + { + "selected_event_count": 1, + "usable_event_count": 1, + "covered_session_count": 0, + "excluded_session_count": 1, + }, + ) class ProductionCanaryTests(unittest.TestCase): @@ -341,6 +409,12 @@ def _worker(self): w.analyze_ready_sessions = lambda: 3 w.analyze_payloads = lambda: 2 w.payload_scan_truncated = False + w.session_capture_coverage = { + "selected_event_count": 7, + "usable_event_count": 7, + "covered_session_count": 2, + "excluded_session_count": 1, + } return w def test_a_failing_stage_keeps_the_others_results(self): @@ -362,6 +436,15 @@ def test_a_healthy_cycle_reports_no_stage_errors(self): result = w.run_once() self.assertNotIn("stage_errors", result) self.assertEqual(result["reports"], 1) + self.assertEqual( + result["session_capture_coverage"], + { + "selected_event_count": 7, + "usable_event_count": 7, + "covered_session_count": 2, + "excluded_session_count": 1, + }, + ) def test_only_the_exception_type_is_recorded_never_its_message(self): # These exceptions come from a model fed attacker-controlled text; a diff --git a/llm-worker/worker.py b/llm-worker/worker.py index 387939f9d..2a37c63ea 100644 --- a/llm-worker/worker.py +++ b/llm-worker/worker.py @@ -78,6 +78,24 @@ "cowrie.login.success", "cowrie.session.closed", ) +TERMINAL_OBSERVATION_CLOSE = "close_observed" +TERMINAL_OBSERVATION_IDLE = "idle_finalized" +TERMINAL_OBSERVATION_OPEN = "open_after_window" +TERMINAL_OBSERVATION_VALUES = frozenset( + { + TERMINAL_OBSERVATION_CLOSE, + TERMINAL_OBSERVATION_IDLE, + TERMINAL_OBSERVATION_OPEN, + } +) +CAPTURE_COVERAGE_OBSERVED = "observed" +CAPTURE_COVERAGE_MISSING = "missing" +CAPTURE_COVERAGE_VALUES = frozenset( + { + CAPTURE_COVERAGE_OBSERVED, + CAPTURE_COVERAGE_MISSING, + } +) LOG = logging.getLogger("llm-worker") AnnotationT = TypeVar("AnnotationT", bound=StrictAnnotation) @@ -512,12 +530,20 @@ class SessionAccumulator: auth_success: bool = False closed: bool = False duration_seconds: float = 0.0 + terminal_observation: str = TERMINAL_OBSERVATION_OPEN + capture_coverage: str = CAPTURE_COVERAGE_MISSING event_hashes: list[str] = field(default_factory=list) finalized: bool = False attempts: int = 0 @classmethod def from_document(cls, source: dict[str, Any], session_id: str) -> "SessionAccumulator": + terminal_observation = source.get("terminal_observation") + if terminal_observation not in TERMINAL_OBSERVATION_VALUES: + terminal_observation = TERMINAL_OBSERVATION_OPEN + capture_coverage = source.get("capture_coverage") + if capture_coverage not in CAPTURE_COVERAGE_VALUES: + capture_coverage = CAPTURE_COVERAGE_MISSING return cls( session_id=session_id, first_seen=str(source.get("first_seen") or ""), @@ -529,6 +555,8 @@ def from_document(cls, source: dict[str, Any], session_id: str) -> "SessionAccum auth_success=bool(source.get("auth_success")), closed=bool(source.get("closed")), duration_seconds=max(0.0, float(source.get("duration_seconds") or 0.0)), + terminal_observation=terminal_observation, + capture_coverage=capture_coverage, event_hashes=[str(value) for value in source.get("event_hashes", []) if isinstance(value, str)][-400:], finalized=bool(source.get("finalized")), attempts=max(0, int(source.get("attempts") or 0)), @@ -574,6 +602,8 @@ def add_event(self, hit: dict[str, Any], max_content_chars: int) -> None: self.commands = self.commands[:100] + self.commands[-99:] + [cleaned] if event_id == "cowrie.session.closed": self.closed = True + self.terminal_observation = TERMINAL_OBSERVATION_CLOSE + self.capture_coverage = CAPTURE_COVERAGE_OBSERVED try: self.duration_seconds = max(0.0, float(nested(source, "honeypot", "duration") or 0.0)) except (TypeError, ValueError): @@ -592,6 +622,8 @@ def document(self) -> dict[str, Any]: "auth_success": self.auth_success, "closed": self.closed, "duration_seconds": self.duration_seconds, + "terminal_observation": self.terminal_observation, + "capture_coverage": self.capture_coverage, "event_hashes": self.event_hashes, "finalized": self.finalized, "attempts": self.attempts, @@ -670,6 +702,8 @@ def state_mapping() -> dict[str, Any]: "auth_success": {"type": "boolean"}, "closed": {"type": "boolean"}, "duration_seconds": {"type": "float"}, + "terminal_observation": {"type": "keyword"}, + "capture_coverage": {"type": "keyword"}, "event_hashes": {"type": "keyword", "index": False, "doc_values": False}, "finalized": {"type": "boolean"}, "attempts": {"type": "integer"}, @@ -694,6 +728,12 @@ def __init__(self, config: Config, es: Elasticsearch | None = None, model: Ollam self.config = config self.es = es self.model = model + self.session_capture_coverage: dict[str, int] = { + "selected_event_count": 0, + "usable_event_count": 0, + "covered_session_count": 0, + "excluded_session_count": 0, + } if not config.dry_run: self.es = es or Elasticsearch(config.es_host, request_timeout=30) self.model = model or OllamaClient(config) @@ -735,6 +775,12 @@ def load_accumulator(self, session_id: str) -> SessionAccumulator: def collect_session_events(self) -> int: assert self.es is not None + self.session_capture_coverage = { + "selected_event_count": 0, + "usable_event_count": 0, + "covered_session_count": 0, + "excluded_session_count": 0, + } since = self.load_checkpoint("sessions") response = self.es.search( index="honeypot-v2-*", @@ -772,6 +818,8 @@ def collect_session_events(self) -> int: usable = 0 accumulated = 0 close_only_skipped = 0 + covered_session_ids: set[str] = set() + excluded_session_ids: set[str] = set() for hit in hits: source = hit.get("_source") if isinstance(hit.get("_source"), dict) else {} session_id = bounded_string(nested(source, "honeypot", "session"), 128) @@ -779,6 +827,7 @@ def collect_session_events(self) -> int: if not session_id or not timestamp: continue usable += 1 + covered_session_ids.add(session_id) accumulator = self.load_accumulator(session_id) event_id = bounded_string(nested(source, "honeypot", "eventid"), 120).lower() # Most Cowrie connections close without ever reaching a command. @@ -792,6 +841,8 @@ def collect_session_events(self) -> int: ): latest = max(latest, timestamp) close_only_skipped += 1 + covered_session_ids.discard(session_id) + excluded_session_ids.add(session_id) continue if not accumulator.finalized: accumulator.add_event(hit, self.config.max_content_chars) @@ -800,6 +851,12 @@ def collect_session_events(self) -> int: latest = max(latest, timestamp) if latest: self.save_checkpoint("sessions", latest) + self.session_capture_coverage = { + "selected_event_count": len(hits), + "usable_event_count": usable, + "covered_session_count": len(covered_session_ids), + "excluded_session_count": len(excluded_session_ids), + } LOG.info( "session scan selected=%d usable=%d accumulated=%d close_only_skipped=%d checkpoint_advanced=%s", len(hits), @@ -839,7 +896,11 @@ def ready_sessions(self) -> list[tuple[str, SessionAccumulator]]: source = hit.get("_source") if isinstance(hit.get("_source"), dict) else {} session_id = source.get("session_id") if isinstance(session_id, str) and session_id: - ready.append((str(hit.get("_id")), SessionAccumulator.from_document(source, session_id))) + accumulator = SessionAccumulator.from_document(source, session_id) + if not accumulator.closed: + accumulator.terminal_observation = TERMINAL_OBSERVATION_IDLE + accumulator.capture_coverage = CAPTURE_COVERAGE_MISSING + ready.append((str(hit.get("_id")), accumulator)) return ready def base_document( @@ -1190,7 +1251,7 @@ def analyze_daily_report(self) -> int: self.es.index(index=ANALYSIS_INDEX, id=report_id, document=document) return 1 - def run_once(self) -> dict[str, int | str | bool]: + def run_once(self) -> dict[str, Any]: if self.config.dry_run: run_selftest() return {"mode": "dry-run", "selftest": True, "sessions": 0, "payloads": 0, "reports": 0} @@ -1226,6 +1287,8 @@ def run_once(self) -> dict[str, int | str | bool]: "payloads": payloads, "reports": reports, } + if self.config.session_enabled and not stage_errors.get("session_events"): + result["session_capture_coverage"] = dict(self.session_capture_coverage) # #2228: disclose a truncated payload scan on the status document # itself, not just in .env.example -- a run that hit either cap must # never read as full coverage.