diff --git a/src/openhound/core/clients/bloodhound_enterprise.py b/src/openhound/core/clients/bloodhound_enterprise.py index fb61c552..8d86eda3 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, @@ -47,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") @@ -64,6 +68,38 @@ 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: + """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=" + 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, + 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.", + 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 cc218f7e..8ecd5a5a 100644 --- a/src/openhound/core/clients/models/jobs.py +++ b/src/openhound/core/clients/models/jobs.py @@ -1,6 +1,9 @@ -from pydantic import BaseModel from datetime import datetime from enum import StrEnum +from typing import Any +from uuid import UUID + +from pydantic import BaseModel class Job(BaseModel): @@ -53,6 +56,41 @@ class JobsEnd(BaseModel): data: Job +class CollectorJob(BaseModel): + 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: UUID | 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: UUID | 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 + skip: int + limit: int + data: CollectorJobsAvailableData + + class ManagementOperationType(StrEnum): SUPPORT_BUNDLE = "support_bundle" diff --git a/src/openhound/scheduler/service.py b/src/openhound/scheduler/service.py index e64a4b6e..a8831f03 100644 --- a/src/openhound/scheduler/service.py +++ b/src/openhound/scheduler/service.py @@ -10,6 +10,7 @@ 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, @@ -170,6 +171,17 @@ def check_jobs(self) -> Job | 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.") diff --git a/tests/test_bhe_job_scheduling.py b/tests/test_bhe_job_scheduling.py index b9177947..ba0cd5ca 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 @@ -21,6 +22,8 @@ ManagedJobOutcome, ) from openhound.core.clients.models.jobs import ( + CollectorJob, + CollectorJobsAvailable, ManagementOperation, ManagementOperationStatus, ManagementOperationType, @@ -57,6 +60,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 +83,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 +242,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": @@ -242,8 +264,6 @@ def mock_request(method, url, **kwargs): token_id="test-id", collector_name="openhound-faker", ) - - def test_client_update_sends_metadata(mock_service, mock_bloodhound_api, monkeypatch): monkeypatch.setattr( bloodhound_enterprise.socket, "gethostname", lambda: "test-host" @@ -298,6 +318,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 +584,104 @@ def test_jobs_no_jobs_available(mock_service, mock_bloodhound_api): assert mock_service.check_jobs() is None +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" + ) + + job = mock_service.check_managed_collector_jobs() + + assert job is not None + 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(): + 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 == 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 + 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 == "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 + 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" + + +@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") + + payload["data"]["jobs"][1][field] = "not-a-uuid" + + with pytest.raises(ValueError): + CollectorJobsAvailable.model_validate(payload) + + +@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] + + with pytest.raises(ValueError): + CollectorJobsAvailable.model_validate(payload) + + +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_re_raised( + mock_service, mock_bloodhound_api, caplog +): + mock_bloodhound_api.app.state.collector_job_queue_error_status = 503 + + with caplog.at_level(logging.ERROR), pytest.raises(BloodHoundHTTPError): + mock_service.check_managed_collector_jobs() + + 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 + + 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 00000000..4b71256f --- /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 00000000..7e59a5c9 --- /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": "ready", + "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 00000000..7b68a33e --- /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": "ready", + "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": null, + "claimed_at": null, + "claim_expires_at": null, + "created_at": "2026-02-20T09:45:00Z", + "updated_at": "2026-02-20T10:50:00Z" + } + ] + } +}