Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 20 additions & 0 deletions docs/llm-worker/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
83 changes: 83 additions & 0 deletions llm-worker/tests/test_worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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):
Expand Down Expand Up @@ -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):
Expand All @@ -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
Expand Down
67 changes: 65 additions & 2 deletions llm-worker/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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 ""),
Expand All @@ -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)),
Expand Down Expand Up @@ -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):
Expand All @@ -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,
Expand Down Expand Up @@ -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"},
Expand All @@ -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)
Expand Down Expand Up @@ -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-*",
Expand Down Expand Up @@ -772,13 +818,16 @@ 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)
timestamp = bounded_string(source.get("@timestamp"), 64)
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.
Expand All @@ -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)
Expand All @@ -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),
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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}
Expand Down Expand Up @@ -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.
Expand Down
Loading