From 6cb58db676a7cd5c145d7a5cf0892bd9015aef47 Mon Sep 17 00:00:00 2001 From: Stran Dutton Date: Thu, 17 Sep 2026 14:45:28 -0500 Subject: [PATCH 1/6] Check for available managed collector jobs - BED-9512  Conflicts:  src/openhound/core/clients/bloodhound_enterprise.py --- .../core/clients/bloodhound_enterprise.py | 22 +++++++++++ src/openhound/core/clients/models/jobs.py | 39 ++++++++++++++++++- src/openhound/scheduler/service.py | 30 ++++++++++---- 3 files changed, 83 insertions(+), 8 deletions(-) diff --git a/src/openhound/core/clients/bloodhound_enterprise.py b/src/openhound/core/clients/bloodhound_enterprise.py index fb61c55..56e17cf 100644 --- a/src/openhound/core/clients/bloodhound_enterprise.py +++ b/src/openhound/core/clients/bloodhound_enterprise.py @@ -10,12 +10,14 @@ from enum import Enum from pathlib import Path from typing import Any, TypeVar +from urllib.parse import quote import requests from openhound.core.clients.bloodhound import BloodHound, BloodHoundHTTPError from openhound.core.clients.models.jobs import ( ArtifactUploadSession, + CollectorJobsAvailable, JobsAvailable, JobsCurrent, JobsEnd, @@ -64,6 +66,26 @@ def jobs_current(self) -> JobsCurrent: response = self.request(method="GET", path=path) return JobsCurrent.model_validate(response.json()) + def available_collector_jobs(self, job_key: str) -> CollectorJobsAvailable: + encoded_job_key = quote(f"eq:{job_key}", safe="") + path = ( + "/api/v2/collector-job-queue/available?limit=1&job_key=" + f"{encoded_job_key}" + ) + logger.debug( + "Polling managed collector job queue.", + extra={"endpoint": path, "job_key": job_key}, + ) + try: + response = self.request(method="GET", path=path) + except (BloodHoundHTTPError, requests.RequestException): + logger.exception( + "Managed collector job queue request failed.", + extra={"endpoint": path, "job_key": job_key}, + ) + raise + return CollectorJobsAvailable.model_validate(response.json()) + def start_job(self, job_id: int) -> JobStart: path = "/api/v2/jobs/start" body = json.dumps({"id": job_id}) diff --git a/src/openhound/core/clients/models/jobs.py b/src/openhound/core/clients/models/jobs.py index cc218f7..db256a1 100644 --- a/src/openhound/core/clients/models/jobs.py +++ b/src/openhound/core/clients/models/jobs.py @@ -1,6 +1,8 @@ -from pydantic import BaseModel from datetime import datetime from enum import StrEnum +from typing import Any + +from pydantic import BaseModel class Job(BaseModel): @@ -53,6 +55,41 @@ class JobsEnd(BaseModel): data: Job +class CollectorJob(BaseModel): + id: str + job_schedule_id: int | None + job_profile_id: int | None + job_type_id: int + job_key: str + params_version: str + params: dict[str, Any] + scope_client_id: str | None + secret_key_id: str | None + priority: int + status: str + run_at: datetime + unclaimed_deadline_at: datetime + attempts: int + max_attempts: int + last_failure: str | None + claimed_by: str | None + claimed_at: datetime | None + claim_expires_at: datetime | None + created_at: datetime + updated_at: datetime + + +class CollectorJobsAvailableData(BaseModel): + jobs: list[CollectorJob] + + +class CollectorJobsAvailable(BaseModel): + count: int | None = None + skip: int | None = None + limit: int | None = None + data: CollectorJobsAvailableData + + class ManagementOperationType(StrEnum): SUPPORT_BUNDLE = "support_bundle" diff --git a/src/openhound/scheduler/service.py b/src/openhound/scheduler/service.py index e64a4b6..ccd380f 100644 --- a/src/openhound/scheduler/service.py +++ b/src/openhound/scheduler/service.py @@ -7,9 +7,11 @@ from pathlib import Path import openhound +import openhound.config as config import openhound.core.logging as openhound_logging from openhound.core.clients.bloodhound_enterprise import BloodHoundEnterprise, JobStatus from openhound.core.clients.models.jobs import ( + CollectorJob, Job, ManagementOperation, ManagementOperationStatus, @@ -90,8 +92,8 @@ def _subprocess_collect(collector_name: str, job_id: int) -> Result: class Service: """Base scheduler service that checks for available jobs in BloodHound Enterprise. - Runs on a simple loop every X seconds and checks for available jobs. If a job is available, - a subprocess is started for the configured collector to run the DLT/OpenHound pipeline. + Runs on a simple loop every X seconds and checks for available jobs. In unmanaged mode (default), + an available job starts a subprocess for the configured collector to run the DLT/OpenHound pipeline. """ def __init__( @@ -143,13 +145,23 @@ def _reset_executor(self) -> None: logger.exception("Error shutting down broken executor.") self.executor = ProcessPoolExecutor(max_workers=1, max_tasks_per_child=1) - def check_jobs(self) -> Job | None: + def check_jobs(self) -> Job | CollectorJob | None: """Checks BloodHound enterprise for available jobs. These can either be new jobs or jobs currently started and not finished/stopped. Returns: - Job | None: Returns a Job object if there is a new or existing job available, otherwise returns None. + Job | CollectorJob | None: Returns the first available job, otherwise returns None. """ logger.info("Checking for new jobs in BloodHound Enterprise.") + + if config.is_managed(): + available_jobs = self.client.available_collector_jobs(self.collector_name) + if available_jobs.data.jobs: + logger.info( + "Managed queue job available: %s", available_jobs.data.jobs[0].id + ) + return available_jobs.data.jobs[0] + return None + new_jobs = self.client.jobs_available if new_jobs.data: logger.info(f"New job available: {new_jobs.data[0].id}") @@ -309,9 +321,13 @@ def _poll(self) -> None: return try: - available_job = self.check_jobs() - if available_job: - self._start_job(available_job) + if config.is_managed(): + # Queue jobs are returned for the later claim/validation handoff. + self.check_jobs() + else: + available_job = self.check_jobs() + if available_job: + self._start_job(available_job) except Exception: logger.exception("Error checking for or starting jobs.") else: From c82bd710a1eaf1f47126f39ac1a0183239492e14 Mon Sep 17 00:00:00 2001 From: Stran Dutton Date: Mon, 21 Sep 2026 16:41:40 -0500 Subject: [PATCH 2/6] Add timeout, cleanup, and tests --- .../core/clients/bloodhound_enterprise.py | 11 +- tests/test_bhe_job_scheduling.py | 199 +++++++++++++++++- .../jobs/collector_jobs_available_empty.json | 8 + .../collector_jobs_available_with_job.json | 36 ++++ .../collector_jobs_available_with_jobs.json | 59 ++++++ 5 files changed, 311 insertions(+), 2 deletions(-) create mode 100644 tests/test_data/api/jobs/collector_jobs_available_empty.json create mode 100644 tests/test_data/api/jobs/collector_jobs_available_with_job.json create mode 100644 tests/test_data/api/jobs/collector_jobs_available_with_jobs.json diff --git a/src/openhound/core/clients/bloodhound_enterprise.py b/src/openhound/core/clients/bloodhound_enterprise.py index 56e17cf..4ff3dbd 100644 --- a/src/openhound/core/clients/bloodhound_enterprise.py +++ b/src/openhound/core/clients/bloodhound_enterprise.py @@ -49,6 +49,8 @@ class ManagedJobOutcome(str, Enum): SUPPORT_BUNDLE_RETRY_DELAY_SECONDS = 2 SUPPORT_BUNDLE_CONNECT_TIMEOUT_SECONDS = 10 SUPPORT_BUNDLE_READ_TIMEOUT_SECONDS = 120 +MANAGED_COLLECTOR_JOB_AVAILABLE_CONNECT_TIMEOUT_SECONDS = 10 +MANAGED_COLLECTOR_JOB_AVAILABLE_READ_TIMEOUT_SECONDS = 20 T = TypeVar("T") @@ -77,7 +79,14 @@ def available_collector_jobs(self, job_key: str) -> CollectorJobsAvailable: extra={"endpoint": path, "job_key": job_key}, ) try: - response = self.request(method="GET", path=path) + response = self.request( + method="GET", + path=path, + timeout=( + MANAGED_COLLECTOR_JOB_AVAILABLE_CONNECT_TIMEOUT_SECONDS, + MANAGED_COLLECTOR_JOB_AVAILABLE_READ_TIMEOUT_SECONDS, + ), + ) except (BloodHoundHTTPError, requests.RequestException): logger.exception( "Managed collector job queue request failed.", diff --git a/tests/test_bhe_job_scheduling.py b/tests/test_bhe_job_scheduling.py index b917794..72ee51b 100644 --- a/tests/test_bhe_job_scheduling.py +++ b/tests/test_bhe_job_scheduling.py @@ -21,6 +21,8 @@ ManagedJobOutcome, ) from openhound.core.clients.models.jobs import ( + CollectorJob, + CollectorJobsAvailable, ManagementOperation, ManagementOperationStatus, ManagementOperationType, @@ -57,6 +59,13 @@ def mock_bloodhound_api(): app.state.end_payload = None app.state.collector_job_end_requests = [] app.state.start_payload = None + app.state.jobs_available_requests = 0 + app.state.jobs_current_requests = 0 + app.state.collector_job_queue_requests = [] + app.state.collector_job_queue_response = load_json( + "collector_jobs_available_empty.json" + ) + app.state.collector_job_queue_error_status = None app.state.client_update_payload = None app.state.ingested_edges = 0 app.state.management_operations = [] @@ -73,14 +82,23 @@ def mock_bloodhound_api(): @app.get("/api/v2/jobs/available") async def jobs_available(): + app.state.jobs_available_requests += 1 if not app.state.job_started: return load_json("jobs_available_with_job.json") return load_json("jobs_available_empty.json") @app.get("/api/v2/jobs/current") async def jobs_current(): + app.state.jobs_current_requests += 1 return Response(status_code=404) + @app.get("/api/v2/collector-job-queue/available") + async def collector_job_queue(request: Request): + app.state.collector_job_queue_requests.append(dict(request.query_params)) + if app.state.collector_job_queue_error_status is not None: + return Response(status_code=app.state.collector_job_queue_error_status) + return app.state.collector_job_queue_response + @app.post("/api/v2/jobs/start") async def start_job(body: dict): app.state.job_started = True @@ -223,7 +241,10 @@ def shutdown(self, *args, **kwargs): return None def mock_request(method, url, **kwargs): - path = urlsplit(url).path + parsed_url = urlsplit(url) + path = parsed_url.path + if parsed_url.query: + path = f"{path}?{parsed_url.query}" if method.upper() == "GET": return mock_bloodhound_api.get(path) if method.upper() == "POST": @@ -235,6 +256,7 @@ def mock_request(method, url, **kwargs): monkeypatch.setattr("requests.request", mock_request) monkeypatch.setattr(scheduler_service, "ProcessPoolExecutor", DummyExecutor) + monkeypatch.setattr(scheduler_service.config, "is_managed", lambda: False) return Service( bhe_uri="http://localhost:8000", @@ -244,6 +266,11 @@ def mock_request(method, url, **kwargs): ) +@pytest.fixture +def managed_mode(mock_service, monkeypatch): + monkeypatch.setattr(scheduler_service.config, "is_managed", lambda: True) + + def test_client_update_sends_metadata(mock_service, mock_bloodhound_api, monkeypatch): monkeypatch.setattr( bloodhound_enterprise.socket, "gethostname", lambda: "test-host" @@ -298,6 +325,25 @@ def capture_request(*args, **kwargs): ) +def test_available_collector_jobs_uses_bounded_timeouts(mock_service, monkeypatch): + captured = {} + + def capture_request(*args, **kwargs): + captured.update(kwargs) + return SimpleNamespace( + json=lambda: load_json("collector_jobs_available_empty.json") + ) + + monkeypatch.setattr(mock_service.client, "request", capture_request) + + mock_service.client.available_collector_jobs("openhound-faker") + + assert captured["timeout"] == ( + bloodhound_enterprise.MANAGED_COLLECTOR_JOB_AVAILABLE_CONNECT_TIMEOUT_SECONDS, + bloodhound_enterprise.MANAGED_COLLECTOR_JOB_AVAILABLE_READ_TIMEOUT_SECONDS, + ) + + def test_client_update_uses_unknown_when_hostname_lookup_fails( mock_service, mock_bloodhound_api, monkeypatch ): @@ -545,6 +591,157 @@ def test_jobs_no_jobs_available(mock_service, mock_bloodhound_api): assert mock_service.check_jobs() is None +def test_managed_mode_polls_collector_job_queue( + mock_service, mock_bloodhound_api, managed_mode +): + mock_bloodhound_api.app.state.collector_job_queue_response = load_json( + "collector_jobs_available_with_job.json" + ) + + job = mock_service.check_jobs() + + assert job is not None + assert job.id == "11111111-1111-1111-1111-111111111111" + assert mock_bloodhound_api.app.state.collector_job_queue_requests + assert mock_bloodhound_api.app.state.jobs_available_requests == 0 + + +def test_unmanaged_mode_never_polls_collector_job_queue( + mock_service, mock_bloodhound_api, monkeypatch +): + monkeypatch.setattr(scheduler_service.config, "is_managed", lambda: False) + + job = mock_service.check_jobs() + + assert job is not None + assert job.id == 123 + assert mock_bloodhound_api.app.state.collector_job_queue_requests == [] + assert mock_bloodhound_api.app.state.jobs_available_requests == 1 + + +def test_managed_queue_request_uses_collector_key_and_first_page( + mock_service, mock_bloodhound_api, managed_mode +): + mock_service.check_jobs() + + assert mock_bloodhound_api.app.state.collector_job_queue_requests == [ + {"limit": "1", "job_key": "eq:openhound-faker"} + ] + + +def test_collector_job_queue_response_parses_contract_fields(): + payload = load_json("collector_jobs_available_with_jobs.json") + response = CollectorJobsAvailable.model_validate(payload) + job = response.data.jobs[0] + + assert response.count == 2 + assert response.skip == 0 + assert response.limit == 2 + assert job == CollectorJob.model_validate(payload["data"]["jobs"][0]) + assert job.id == "11111111-1111-1111-1111-111111111111" + assert job.job_schedule_id is None + assert job.job_profile_id is None + assert job.job_type_id == 7 + assert job.job_key == "openhound-faker" + assert job.params_version == "v2" + assert job.params == { + "domains": ["example.com", "example.org"], + "include_deleted": False, + "nested": {"depth": 2}, + } + assert job.scope_client_id is None + assert job.secret_key_id is None + assert job.priority == 10 + assert job.status == "running" + assert job.run_at.isoformat() == "2026-02-20T10:00:00+00:00" + assert job.unclaimed_deadline_at.isoformat() == "2026-02-20T10:05:00+00:00" + assert job.attempts == 0 + assert job.max_attempts == 3 + assert job.last_failure is None + assert job.claimed_by is None + assert job.claimed_at is None + assert job.claim_expires_at is None + assert job.created_at.isoformat() == "2026-02-20T09:00:00+00:00" + assert job.updated_at.isoformat() == "2026-02-20T09:30:00+00:00" + + +def test_managed_mode_selects_first_queue_job_unchanged( + mock_service, managed_mode, monkeypatch +): + payload = load_json("collector_jobs_available_with_jobs.json") + available_jobs = CollectorJobsAvailable.model_validate(payload) + monkeypatch.setattr( + mock_service.client, + "available_collector_jobs", + lambda job_key: available_jobs, + ) + + selected = mock_service.check_jobs() + + assert selected == CollectorJob.model_validate(payload["data"]["jobs"][0]) + assert selected != CollectorJob.model_validate(payload["data"]["jobs"][1]) + + +def test_managed_mode_returns_no_job_for_empty_queue( + mock_service, mock_bloodhound_api, managed_mode +): + assert mock_service.check_jobs() is None + assert len(mock_bloodhound_api.app.state.collector_job_queue_requests) == 1 + + +def test_running_job_prevents_managed_queue_poll( + mock_service, mock_bloodhound_api, managed_mode +): + mock_service.job_running = 123 + + mock_service._poll() + + assert mock_bloodhound_api.app.state.collector_job_queue_requests == [] + assert mock_bloodhound_api.app.state.jobs_current_requests == 1 + + +def test_queue_http_error_is_logged_and_next_poll_continues( + mock_service, mock_bloodhound_api, managed_mode, caplog +): + mock_bloodhound_api.app.state.collector_job_queue_error_status = 503 + + with caplog.at_level(logging.ERROR): + mock_service._poll() + + queue_errors = [ + record + for record in caplog.records + if record.name == "openhound.core.clients.bloodhound_enterprise" + and record.levelno == logging.ERROR + and record.getMessage() == "Managed collector job queue request failed." + ] + assert len(queue_errors) == 1 + + mock_bloodhound_api.app.state.collector_job_queue_error_status = None + mock_service._poll() + + assert len(mock_bloodhound_api.app.state.collector_job_queue_requests) == 2 + + +def test_each_queue_request_is_logged( + mock_service, mock_bloodhound_api, managed_mode, caplog +): + with caplog.at_level( + logging.DEBUG, logger="openhound.core.clients.bloodhound_enterprise" + ): + mock_service.client.available_collector_jobs("openhound-faker") + mock_service.client.available_collector_jobs("openhound-faker") + + queue_polls = [ + record + for record in caplog.records + if record.name == "openhound.core.clients.bloodhound_enterprise" + and record.getMessage() == "Polling managed collector job queue." + ] + assert len(queue_polls) == 2 + assert all(record.levelno == logging.DEBUG for record in queue_polls) + + def test_poll_starts_new_job(mock_service, mock_bloodhound_api, monkeypatch): """Similar to test_jobs_starts_new_job but using the poll method""" submitted = Future() diff --git a/tests/test_data/api/jobs/collector_jobs_available_empty.json b/tests/test_data/api/jobs/collector_jobs_available_empty.json new file mode 100644 index 0000000..4b71256 --- /dev/null +++ b/tests/test_data/api/jobs/collector_jobs_available_empty.json @@ -0,0 +1,8 @@ +{ + "count": 0, + "skip": 0, + "limit": 1, + "data": { + "jobs": [] + } +} diff --git a/tests/test_data/api/jobs/collector_jobs_available_with_job.json b/tests/test_data/api/jobs/collector_jobs_available_with_job.json new file mode 100644 index 0000000..d146dd1 --- /dev/null +++ b/tests/test_data/api/jobs/collector_jobs_available_with_job.json @@ -0,0 +1,36 @@ +{ + "count": 1, + "skip": 0, + "limit": 1, + "data": { + "jobs": [ + { + "id": "11111111-1111-1111-1111-111111111111", + "job_schedule_id": null, + "job_profile_id": null, + "job_type_id": 7, + "job_key": "openhound-faker", + "params_version": "v2", + "params": { + "domains": ["example.com", "example.org"], + "include_deleted": false, + "nested": {"depth": 2} + }, + "scope_client_id": null, + "secret_key_id": null, + "priority": 10, + "status": "running", + "run_at": "2026-02-20T10:00:00Z", + "unclaimed_deadline_at": "2026-02-20T10:05:00Z", + "attempts": 0, + "max_attempts": 3, + "last_failure": null, + "claimed_by": null, + "claimed_at": null, + "claim_expires_at": null, + "created_at": "2026-02-20T09:00:00Z", + "updated_at": "2026-02-20T09:30:00Z" + } + ] + } +} diff --git a/tests/test_data/api/jobs/collector_jobs_available_with_jobs.json b/tests/test_data/api/jobs/collector_jobs_available_with_jobs.json new file mode 100644 index 0000000..a33c294 --- /dev/null +++ b/tests/test_data/api/jobs/collector_jobs_available_with_jobs.json @@ -0,0 +1,59 @@ +{ + "count": 2, + "skip": 0, + "limit": 2, + "data": { + "jobs": [ + { + "id": "11111111-1111-1111-1111-111111111111", + "job_schedule_id": null, + "job_profile_id": null, + "job_type_id": 7, + "job_key": "openhound-faker", + "params_version": "v2", + "params": { + "domains": ["example.com", "example.org"], + "include_deleted": false, + "nested": {"depth": 2} + }, + "scope_client_id": null, + "secret_key_id": null, + "priority": 10, + "status": "running", + "run_at": "2026-02-20T10:00:00Z", + "unclaimed_deadline_at": "2026-02-20T10:05:00Z", + "attempts": 0, + "max_attempts": 3, + "last_failure": null, + "claimed_by": null, + "claimed_at": null, + "claim_expires_at": null, + "created_at": "2026-02-20T09:00:00Z", + "updated_at": "2026-02-20T09:30:00Z" + }, + { + "id": "22222222-2222-2222-2222-222222222222", + "job_schedule_id": 42, + "job_profile_id": 84, + "job_type_id": 8, + "job_key": "openhound-faker", + "params_version": "v3", + "params": {"mode": "full"}, + "scope_client_id": "33333333-3333-3333-3333-333333333333", + "secret_key_id": "secret-1", + "priority": 5, + "status": "ready", + "run_at": "2026-02-20T11:00:00Z", + "unclaimed_deadline_at": "2026-02-20T11:05:00Z", + "attempts": 1, + "max_attempts": 4, + "last_failure": "previous attempt failed", + "claimed_by": "44444444-4444-4444-4444-444444444444", + "claimed_at": "2026-02-20T10:55:00Z", + "claim_expires_at": "2026-02-20T10:59:00Z", + "created_at": "2026-02-20T09:45:00Z", + "updated_at": "2026-02-20T10:50:00Z" + } + ] + } +} From 13ceec89b80857651f8e667cadfb2010336a358b Mon Sep 17 00:00:00 2001 From: Stran Dutton Date: Mon, 21 Sep 2026 16:52:40 -0500 Subject: [PATCH 3/6] added comments --- src/openhound/core/clients/bloodhound_enterprise.py | 5 +++++ src/openhound/scheduler/service.py | 3 ++- 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/src/openhound/core/clients/bloodhound_enterprise.py b/src/openhound/core/clients/bloodhound_enterprise.py index 4ff3dbd..8d86eda 100644 --- a/src/openhound/core/clients/bloodhound_enterprise.py +++ b/src/openhound/core/clients/bloodhound_enterprise.py @@ -69,6 +69,11 @@ def jobs_current(self) -> JobsCurrent: return JobsCurrent.model_validate(response.json()) def available_collector_jobs(self, job_key: str) -> CollectorJobsAvailable: + """Return the next job from BHE's managed collector queue. + + This endpoint is used only by managed OpenHound deployments and is not + part of the standard open-source/self-hosted collector workflow. + """ encoded_job_key = quote(f"eq:{job_key}", safe="") path = ( "/api/v2/collector-job-queue/available?limit=1&job_key=" diff --git a/src/openhound/scheduler/service.py b/src/openhound/scheduler/service.py index ccd380f..ab84c8c 100644 --- a/src/openhound/scheduler/service.py +++ b/src/openhound/scheduler/service.py @@ -153,11 +153,12 @@ def check_jobs(self) -> Job | CollectorJob | None: """ logger.info("Checking for new jobs in BloodHound Enterprise.") + # This 'managed' flow only applies to OpenHound deployments that are hosted by SpecterOps. if config.is_managed(): available_jobs = self.client.available_collector_jobs(self.collector_name) if available_jobs.data.jobs: logger.info( - "Managed queue job available: %s", available_jobs.data.jobs[0].id + "New managed queue job available: %s", available_jobs.data.jobs[0].id ) return available_jobs.data.jobs[0] return None From 494de645efb993f800693680a44d59d0686a5bae Mon Sep 17 00:00:00 2001 From: Stran Dutton Date: Tue, 22 Sep 2026 12:00:55 -0500 Subject: [PATCH 4/6] separating into independent func, removing calls --- src/openhound/scheduler/service.py | 40 ++++++++++++++---------------- tests/test_bhe_job_scheduling.py | 18 +++++++------- 2 files changed, 27 insertions(+), 31 deletions(-) diff --git a/src/openhound/scheduler/service.py b/src/openhound/scheduler/service.py index ab84c8c..b415907 100644 --- a/src/openhound/scheduler/service.py +++ b/src/openhound/scheduler/service.py @@ -92,8 +92,8 @@ def _subprocess_collect(collector_name: str, job_id: int) -> Result: class Service: """Base scheduler service that checks for available jobs in BloodHound Enterprise. - Runs on a simple loop every X seconds and checks for available jobs. In unmanaged mode (default), - an available job starts a subprocess for the configured collector to run the DLT/OpenHound pipeline. + Runs on a simple loop every X seconds and checks for available jobs. If a job is available, + a subprocess is started for the configured collector to run the DLT/OpenHound pipeline. """ def __init__( @@ -145,24 +145,13 @@ def _reset_executor(self) -> None: logger.exception("Error shutting down broken executor.") self.executor = ProcessPoolExecutor(max_workers=1, max_tasks_per_child=1) - def check_jobs(self) -> Job | CollectorJob | None: + def check_jobs(self) -> Job | None: """Checks BloodHound enterprise for available jobs. These can either be new jobs or jobs currently started and not finished/stopped. Returns: - Job | CollectorJob | None: Returns the first available job, otherwise returns None. + Job | None: Returns a Job object if there is a new or existing job available, otherwise returns None. """ logger.info("Checking for new jobs in BloodHound Enterprise.") - - # This 'managed' flow only applies to OpenHound deployments that are hosted by SpecterOps. - if config.is_managed(): - available_jobs = self.client.available_collector_jobs(self.collector_name) - if available_jobs.data.jobs: - logger.info( - "New managed queue job available: %s", available_jobs.data.jobs[0].id - ) - return available_jobs.data.jobs[0] - return None - new_jobs = self.client.jobs_available if new_jobs.data: logger.info(f"New job available: {new_jobs.data[0].id}") @@ -183,6 +172,17 @@ def check_jobs(self) -> Job | CollectorJob | None: return None + def check_managed_collector_jobs(self) -> CollectorJob | None: + """Return the first available job from the managed collector queue.""" + logger.info("Checking for new managed collector jobs in BloodHound Enterprise.") + available_jobs = self.client.available_collector_jobs(self.collector_name) + if available_jobs.data.jobs: + logger.info( + "New managed queue job available: %s", available_jobs.data.jobs[0].id + ) + return available_jobs.data.jobs[0] + return None + def check_management(self) -> ManagementOperation | None: """Return the first pending support-bundle operation, if any.""" logger.info("Checking for management operations in BloodHound Enterprise.") @@ -322,13 +322,9 @@ def _poll(self) -> None: return try: - if config.is_managed(): - # Queue jobs are returned for the later claim/validation handoff. - self.check_jobs() - else: - available_job = self.check_jobs() - if available_job: - self._start_job(available_job) + available_job = self.check_jobs() + if available_job: + self._start_job(available_job) except Exception: logger.exception("Error checking for or starting jobs.") else: diff --git a/tests/test_bhe_job_scheduling.py b/tests/test_bhe_job_scheduling.py index 72ee51b..25f2c30 100644 --- a/tests/test_bhe_job_scheduling.py +++ b/tests/test_bhe_job_scheduling.py @@ -591,14 +591,14 @@ def test_jobs_no_jobs_available(mock_service, mock_bloodhound_api): assert mock_service.check_jobs() is None -def test_managed_mode_polls_collector_job_queue( +def test_managed_collector_job_check_polls_queue( mock_service, mock_bloodhound_api, managed_mode ): mock_bloodhound_api.app.state.collector_job_queue_response = load_json( "collector_jobs_available_with_job.json" ) - job = mock_service.check_jobs() + job = mock_service.check_managed_collector_jobs() assert job is not None assert job.id == "11111111-1111-1111-1111-111111111111" @@ -622,7 +622,7 @@ def test_unmanaged_mode_never_polls_collector_job_queue( def test_managed_queue_request_uses_collector_key_and_first_page( mock_service, mock_bloodhound_api, managed_mode ): - mock_service.check_jobs() + mock_service.check_managed_collector_jobs() assert mock_bloodhound_api.app.state.collector_job_queue_requests == [ {"limit": "1", "job_key": "eq:openhound-faker"} @@ -676,7 +676,7 @@ def test_managed_mode_selects_first_queue_job_unchanged( lambda job_key: available_jobs, ) - selected = mock_service.check_jobs() + selected = mock_service.check_managed_collector_jobs() assert selected == CollectorJob.model_validate(payload["data"]["jobs"][0]) assert selected != CollectorJob.model_validate(payload["data"]["jobs"][1]) @@ -685,7 +685,7 @@ def test_managed_mode_selects_first_queue_job_unchanged( def test_managed_mode_returns_no_job_for_empty_queue( mock_service, mock_bloodhound_api, managed_mode ): - assert mock_service.check_jobs() is None + assert mock_service.check_managed_collector_jobs() is None assert len(mock_bloodhound_api.app.state.collector_job_queue_requests) == 1 @@ -700,13 +700,13 @@ def test_running_job_prevents_managed_queue_poll( assert mock_bloodhound_api.app.state.jobs_current_requests == 1 -def test_queue_http_error_is_logged_and_next_poll_continues( +def test_queue_http_error_is_logged_and_next_check_continues( mock_service, mock_bloodhound_api, managed_mode, caplog ): mock_bloodhound_api.app.state.collector_job_queue_error_status = 503 - with caplog.at_level(logging.ERROR): - mock_service._poll() + with caplog.at_level(logging.ERROR), pytest.raises(BloodHoundHTTPError): + mock_service.check_managed_collector_jobs() queue_errors = [ record @@ -718,7 +718,7 @@ def test_queue_http_error_is_logged_and_next_poll_continues( assert len(queue_errors) == 1 mock_bloodhound_api.app.state.collector_job_queue_error_status = None - mock_service._poll() + mock_service.check_managed_collector_jobs() assert len(mock_bloodhound_api.app.state.collector_job_queue_requests) == 2 From 7713cb329f86b9ad72d22b384ec0fd3608b1ec68 Mon Sep 17 00:00:00 2001 From: Stran Dutton Date: Tue, 22 Sep 2026 15:21:24 -0500 Subject: [PATCH 5/6] remove unused import --- src/openhound/scheduler/service.py | 1 - 1 file changed, 1 deletion(-) diff --git a/src/openhound/scheduler/service.py b/src/openhound/scheduler/service.py index b415907..a8831f0 100644 --- a/src/openhound/scheduler/service.py +++ b/src/openhound/scheduler/service.py @@ -7,7 +7,6 @@ from pathlib import Path import openhound -import openhound.config as config import openhound.core.logging as openhound_logging from openhound.core.clients.bloodhound_enterprise import BloodHoundEnterprise, JobStatus from openhound.core.clients.models.jobs import ( From 03c4d5e1696f25771463e7f345777bfe4def590f Mon Sep 17 00:00:00 2001 From: Stran Dutton Date: Wed, 23 Sep 2026 16:47:28 -0500 Subject: [PATCH 6/6] align test data and models with contract --- src/openhound/core/clients/models/jobs.py | 13 ++- tests/test_bhe_job_scheduling.py | 110 ++++-------------- .../collector_jobs_available_with_job.json | 2 +- .../collector_jobs_available_with_jobs.json | 8 +- 4 files changed, 37 insertions(+), 96 deletions(-) diff --git a/src/openhound/core/clients/models/jobs.py b/src/openhound/core/clients/models/jobs.py index db256a1..8ecd5a5 100644 --- a/src/openhound/core/clients/models/jobs.py +++ b/src/openhound/core/clients/models/jobs.py @@ -1,6 +1,7 @@ from datetime import datetime from enum import StrEnum from typing import Any +from uuid import UUID from pydantic import BaseModel @@ -56,14 +57,14 @@ class JobsEnd(BaseModel): class CollectorJob(BaseModel): - id: str + id: UUID job_schedule_id: int | None job_profile_id: int | None job_type_id: int job_key: str params_version: str params: dict[str, Any] - scope_client_id: str | None + scope_client_id: UUID | None secret_key_id: str | None priority: int status: str @@ -72,7 +73,7 @@ class CollectorJob(BaseModel): attempts: int max_attempts: int last_failure: str | None - claimed_by: str | None + claimed_by: UUID | None claimed_at: datetime | None claim_expires_at: datetime | None created_at: datetime @@ -84,9 +85,9 @@ class CollectorJobsAvailableData(BaseModel): class CollectorJobsAvailable(BaseModel): - count: int | None = None - skip: int | None = None - limit: int | None = None + count: int + skip: int + limit: int data: CollectorJobsAvailableData diff --git a/tests/test_bhe_job_scheduling.py b/tests/test_bhe_job_scheduling.py index 25f2c30..ba0cd5c 100644 --- a/tests/test_bhe_job_scheduling.py +++ b/tests/test_bhe_job_scheduling.py @@ -8,6 +8,7 @@ from pathlib import Path from types import SimpleNamespace from urllib.parse import urlsplit +from uuid import UUID import pytest import requests @@ -256,7 +257,6 @@ def mock_request(method, url, **kwargs): monkeypatch.setattr("requests.request", mock_request) monkeypatch.setattr(scheduler_service, "ProcessPoolExecutor", DummyExecutor) - monkeypatch.setattr(scheduler_service.config, "is_managed", lambda: False) return Service( bhe_uri="http://localhost:8000", @@ -264,13 +264,6 @@ def mock_request(method, url, **kwargs): token_id="test-id", collector_name="openhound-faker", ) - - -@pytest.fixture -def managed_mode(mock_service, monkeypatch): - monkeypatch.setattr(scheduler_service.config, "is_managed", lambda: True) - - def test_client_update_sends_metadata(mock_service, mock_bloodhound_api, monkeypatch): monkeypatch.setattr( bloodhound_enterprise.socket, "gethostname", lambda: "test-host" @@ -591,8 +584,8 @@ def test_jobs_no_jobs_available(mock_service, mock_bloodhound_api): assert mock_service.check_jobs() is None -def test_managed_collector_job_check_polls_queue( - mock_service, mock_bloodhound_api, managed_mode +def test_managed_collector_job_check_requests_first_queue_job( + mock_service, mock_bloodhound_api ): mock_bloodhound_api.app.state.collector_job_queue_response = load_json( "collector_jobs_available_with_job.json" @@ -601,32 +594,11 @@ def test_managed_collector_job_check_polls_queue( job = mock_service.check_managed_collector_jobs() assert job is not None - assert job.id == "11111111-1111-1111-1111-111111111111" - assert mock_bloodhound_api.app.state.collector_job_queue_requests - assert mock_bloodhound_api.app.state.jobs_available_requests == 0 - - -def test_unmanaged_mode_never_polls_collector_job_queue( - mock_service, mock_bloodhound_api, monkeypatch -): - monkeypatch.setattr(scheduler_service.config, "is_managed", lambda: False) - - job = mock_service.check_jobs() - - assert job is not None - assert job.id == 123 - assert mock_bloodhound_api.app.state.collector_job_queue_requests == [] - assert mock_bloodhound_api.app.state.jobs_available_requests == 1 - - -def test_managed_queue_request_uses_collector_key_and_first_page( - mock_service, mock_bloodhound_api, managed_mode -): - mock_service.check_managed_collector_jobs() - + assert job.id == UUID("11111111-1111-1111-1111-111111111111") assert mock_bloodhound_api.app.state.collector_job_queue_requests == [ {"limit": "1", "job_key": "eq:openhound-faker"} ] + assert mock_bloodhound_api.app.state.jobs_available_requests == 0 def test_collector_job_queue_response_parses_contract_fields(): @@ -638,7 +610,7 @@ def test_collector_job_queue_response_parses_contract_fields(): assert response.skip == 0 assert response.limit == 2 assert job == CollectorJob.model_validate(payload["data"]["jobs"][0]) - assert job.id == "11111111-1111-1111-1111-111111111111" + assert job.id == UUID("11111111-1111-1111-1111-111111111111") assert job.job_schedule_id is None assert job.job_profile_id is None assert job.job_type_id == 7 @@ -652,7 +624,7 @@ def test_collector_job_queue_response_parses_contract_fields(): assert job.scope_client_id is None assert job.secret_key_id is None assert job.priority == 10 - assert job.status == "running" + assert job.status == "ready" assert job.run_at.isoformat() == "2026-02-20T10:00:00+00:00" assert job.unclaimed_deadline_at.isoformat() == "2026-02-20T10:05:00+00:00" assert job.attempts == 0 @@ -665,43 +637,35 @@ def test_collector_job_queue_response_parses_contract_fields(): assert job.updated_at.isoformat() == "2026-02-20T09:30:00+00:00" -def test_managed_mode_selects_first_queue_job_unchanged( - mock_service, managed_mode, monkeypatch -): +@pytest.mark.parametrize("field", ["id", "scope_client_id", "claimed_by"]) +def test_collector_job_queue_response_rejects_invalid_uuids(field): payload = load_json("collector_jobs_available_with_jobs.json") - available_jobs = CollectorJobsAvailable.model_validate(payload) - monkeypatch.setattr( - mock_service.client, - "available_collector_jobs", - lambda job_key: available_jobs, - ) - selected = mock_service.check_managed_collector_jobs() + payload["data"]["jobs"][1][field] = "not-a-uuid" - assert selected == CollectorJob.model_validate(payload["data"]["jobs"][0]) - assert selected != CollectorJob.model_validate(payload["data"]["jobs"][1]) + with pytest.raises(ValueError): + CollectorJobsAvailable.model_validate(payload) -def test_managed_mode_returns_no_job_for_empty_queue( - mock_service, mock_bloodhound_api, managed_mode -): - assert mock_service.check_managed_collector_jobs() is None - assert len(mock_bloodhound_api.app.state.collector_job_queue_requests) == 1 +@pytest.mark.parametrize("field", ["scope_client_id", "claimed_by"]) +def test_collector_job_queue_response_requires_nullable_uuid_fields(field): + payload = load_json("collector_jobs_available_with_jobs.json") + del payload["data"]["jobs"][0][field] -def test_running_job_prevents_managed_queue_poll( - mock_service, mock_bloodhound_api, managed_mode -): - mock_service.job_running = 123 + with pytest.raises(ValueError): + CollectorJobsAvailable.model_validate(payload) - mock_service._poll() - assert mock_bloodhound_api.app.state.collector_job_queue_requests == [] - assert mock_bloodhound_api.app.state.jobs_current_requests == 1 +def test_managed_collector_job_check_returns_no_job_for_empty_queue( + mock_service, mock_bloodhound_api +): + assert mock_service.check_managed_collector_jobs() is None + assert len(mock_bloodhound_api.app.state.collector_job_queue_requests) == 1 -def test_queue_http_error_is_logged_and_next_check_continues( - mock_service, mock_bloodhound_api, managed_mode, caplog +def test_queue_http_error_is_logged_and_re_raised( + mock_service, mock_bloodhound_api, caplog ): mock_bloodhound_api.app.state.collector_job_queue_error_status = 503 @@ -717,30 +681,6 @@ def test_queue_http_error_is_logged_and_next_check_continues( ] assert len(queue_errors) == 1 - mock_bloodhound_api.app.state.collector_job_queue_error_status = None - mock_service.check_managed_collector_jobs() - - assert len(mock_bloodhound_api.app.state.collector_job_queue_requests) == 2 - - -def test_each_queue_request_is_logged( - mock_service, mock_bloodhound_api, managed_mode, caplog -): - with caplog.at_level( - logging.DEBUG, logger="openhound.core.clients.bloodhound_enterprise" - ): - mock_service.client.available_collector_jobs("openhound-faker") - mock_service.client.available_collector_jobs("openhound-faker") - - queue_polls = [ - record - for record in caplog.records - if record.name == "openhound.core.clients.bloodhound_enterprise" - and record.getMessage() == "Polling managed collector job queue." - ] - assert len(queue_polls) == 2 - assert all(record.levelno == logging.DEBUG for record in queue_polls) - def test_poll_starts_new_job(mock_service, mock_bloodhound_api, monkeypatch): """Similar to test_jobs_starts_new_job but using the poll method""" diff --git a/tests/test_data/api/jobs/collector_jobs_available_with_job.json b/tests/test_data/api/jobs/collector_jobs_available_with_job.json index d146dd1..7e59a5c 100644 --- a/tests/test_data/api/jobs/collector_jobs_available_with_job.json +++ b/tests/test_data/api/jobs/collector_jobs_available_with_job.json @@ -19,7 +19,7 @@ "scope_client_id": null, "secret_key_id": null, "priority": 10, - "status": "running", + "status": "ready", "run_at": "2026-02-20T10:00:00Z", "unclaimed_deadline_at": "2026-02-20T10:05:00Z", "attempts": 0, diff --git a/tests/test_data/api/jobs/collector_jobs_available_with_jobs.json b/tests/test_data/api/jobs/collector_jobs_available_with_jobs.json index a33c294..7b68a33 100644 --- a/tests/test_data/api/jobs/collector_jobs_available_with_jobs.json +++ b/tests/test_data/api/jobs/collector_jobs_available_with_jobs.json @@ -19,7 +19,7 @@ "scope_client_id": null, "secret_key_id": null, "priority": 10, - "status": "running", + "status": "ready", "run_at": "2026-02-20T10:00:00Z", "unclaimed_deadline_at": "2026-02-20T10:05:00Z", "attempts": 0, @@ -48,9 +48,9 @@ "attempts": 1, "max_attempts": 4, "last_failure": "previous attempt failed", - "claimed_by": "44444444-4444-4444-4444-444444444444", - "claimed_at": "2026-02-20T10:55:00Z", - "claim_expires_at": "2026-02-20T10:59:00Z", + "claimed_by": null, + "claimed_at": null, + "claim_expires_at": null, "created_at": "2026-02-20T09:45:00Z", "updated_at": "2026-02-20T10:50:00Z" }