Skip to content
Merged
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
36 changes: 36 additions & 0 deletions src/openhound/core/clients/bloodhound_enterprise.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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")

Expand All @@ -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})
Expand Down
40 changes: 39 additions & 1 deletion src/openhound/core/clients/models/jobs.py
Original file line number Diff line number Diff line change
@@ -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):
Expand Down Expand Up @@ -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"

Expand Down
12 changes: 12 additions & 0 deletions src/openhound/scheduler/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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.")
Expand Down
143 changes: 140 additions & 3 deletions tests/test_bhe_job_scheduling.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -21,6 +22,8 @@
ManagedJobOutcome,
)
from openhound.core.clients.models.jobs import (
CollectorJob,
CollectorJobsAvailable,
ManagementOperation,
ManagementOperationStatus,
ManagementOperationType,
Expand Down Expand Up @@ -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 = []
Expand All @@ -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
Expand Down Expand Up @@ -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":
Expand All @@ -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"
Expand Down Expand Up @@ -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
):
Expand Down Expand Up @@ -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()
Expand Down
8 changes: 8 additions & 0 deletions tests/test_data/api/jobs/collector_jobs_available_empty.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
{
"count": 0,
"skip": 0,
"limit": 1,
"data": {
"jobs": []
}
}
36 changes: 36 additions & 0 deletions tests/test_data/api/jobs/collector_jobs_available_with_job.json
Original file line number Diff line number Diff line change
@@ -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"
}
]
}
}
Loading