Skip to content
Open
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
4 changes: 4 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -95,3 +95,7 @@ containers/ # Docker configuration
## Deployment

This service is deployed via the [deploy repository](https://github.com/All-Hands-AI/deploy). Docker images are automatically built and pushed to `ghcr.io/openhands/automation` on every push to main and on tags.

### Multi-worker / multi-replica

Uvicorn starts the cron scheduler, run dispatcher, and staleness watchdog once per worker process. For multi-worker deployments, run those loops in a single dedicated process (`AUTOMATION_ENABLE_BACKGROUND_WORKERS=true`, the default) and set `AUTOMATION_ENABLE_BACKGROUND_WORKERS=false` on request-serving replicas so they only handle HTTP.
118 changes: 72 additions & 46 deletions openhands/automation/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,75 @@
logger = logging.getLogger("automation.app")


def start_background_worker_tasks(
app: FastAPI,
settings,
shutdown_event: asyncio.Event,
) -> list[tuple[str, asyncio.Task]]:
"""Start scheduler/dispatcher/watchdog when this process owns them.

Returns the started ``(name, task)`` pairs so lifespan can await them on
shutdown. When ``settings.enable_background_workers`` is False, no tasks
are created (request-serving replicas in a multi-worker deployment).
"""
app.state.scheduler_task = None
app.state.dispatcher_task = None
app.state.watchdog_task = None

if not settings.enable_background_workers:
logger.info(
"Background workers disabled "
"(AUTOMATION_ENABLE_BACKGROUND_WORKERS=false); "
"scheduler/dispatcher/watchdog will not start in this process"
)
return []

# Scheduler: polls automations and creates PENDING runs
scheduler_task = asyncio.create_task(
scheduler_loop(
app.state.session_factory,
interval_seconds=settings.scheduler_interval_seconds,
shutdown_event=shutdown_event,
)
)
app.state.scheduler_task = scheduler_task
logger.info("Background scheduler started")

# Dispatcher: picks up PENDING runs and dispatches them
if not settings.base_url:
logger.warning(
"AUTOMATION_BASE_URL not set — using localhost. "
"Sandboxes in the cloud won't be able to reach this URL."
)
dispatcher_task = asyncio.create_task(
dispatcher_loop(
app.state.session_factory,
settings=settings,
interval_seconds=settings.dispatcher_interval_seconds,
shutdown_event=shutdown_event,
)
)
app.state.dispatcher_task = dispatcher_task
logger.info("Background dispatcher started")

# Watchdog: marks stale RUNNING runs as FAILED
watchdog_task = asyncio.create_task(
watchdog_loop(
app.state.session_factory,
settings=settings,
shutdown_event=shutdown_event,
)
)
app.state.watchdog_task = watchdog_task
logger.info("Background watchdog started")

return [
("scheduler", scheduler_task),
("dispatcher", dispatcher_task),
("watchdog", watchdog_task),
]


@asynccontextmanager
async def lifespan(app: FastAPI):
"""Application startup/shutdown lifecycle."""
Expand Down Expand Up @@ -117,61 +186,18 @@ async def lifespan(app: FastAPI):
msg = f"SQLite migration failed. Database may be inconsistent: {e}"
raise RuntimeError(msg) from e

# Start the background scheduler and dispatcher
shutdown_event = asyncio.Event()
app.state.shutdown_event = shutdown_event

# Scheduler: polls automations and creates PENDING runs
scheduler_task = asyncio.create_task(
scheduler_loop(
app.state.session_factory,
interval_seconds=settings.scheduler_interval_seconds,
shutdown_event=shutdown_event,
)
)
app.state.scheduler_task = scheduler_task
logger.info("Background scheduler started")

# Dispatcher: picks up PENDING runs and dispatches them
if not settings.base_url:
logger.warning(
"AUTOMATION_BASE_URL not set — using localhost. "
"Sandboxes in the cloud won't be able to reach this URL."
)
dispatcher_task = asyncio.create_task(
dispatcher_loop(
app.state.session_factory,
settings=settings,
interval_seconds=settings.dispatcher_interval_seconds,
shutdown_event=shutdown_event,
)
)
app.state.dispatcher_task = dispatcher_task
logger.info("Background dispatcher started")

# Watchdog: marks stale RUNNING runs as FAILED
watchdog_task = asyncio.create_task(
watchdog_loop(
app.state.session_factory,
settings=settings,
shutdown_event=shutdown_event,
)
)
app.state.watchdog_task = watchdog_task
logger.info("Background watchdog started")
background_tasks = start_background_worker_tasks(app, settings, shutdown_event)

yield

# Shutdown
logger.info("Shutting down background tasks...")
shutdown_event.set()

# Wait for all tasks to exit gracefully
for task_name, task in [
("scheduler", scheduler_task),
("dispatcher", dispatcher_task),
("watchdog", watchdog_task),
]:
# Wait for started tasks to exit gracefully (no-op when workers disabled)
for task_name, task in background_tasks:
try:
await asyncio.wait_for(task, timeout=5.0)
except TimeoutError:
Expand Down
9 changes: 9 additions & 0 deletions openhands/automation/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -327,6 +327,9 @@ class ServiceSettings(BaseSettings):
AUTOMATION_WORKSPACE_BASE: Base workspace directory (local mode default)

# Background workers
AUTOMATION_ENABLE_BACKGROUND_WORKERS: Start scheduler/dispatcher/watchdog
in this process (default: True). Set False on request-serving replicas
when a dedicated background process owns the loops (see #286).
AUTOMATION_SCHEDULER_INTERVAL_SECONDS: Scheduler poll interval (default: 60)
AUTOMATION_SCHEDULER_BATCH_SIZE: Scheduler batch size (default: 50)
AUTOMATION_DISPATCHER_INTERVAL_SECONDS: Dispatcher poll interval (default: 10)
Expand Down Expand Up @@ -410,6 +413,12 @@ class ServiceSettings(BaseSettings):
openhands_api_base_url: str = "https://app.all-hands.dev"

# Background workers
# When True (default), this process runs scheduler/dispatcher/watchdog.
# Uvicorn runs lifespan once per worker process, so multi-worker / multi-replica
# deployments should enable this on exactly one dedicated process and disable
# it on request-serving instances to avoid duplicate polling and sandbox-API
# fan-out (see OpenHands/automation#286).
enable_background_workers: bool = True
scheduler_interval_seconds: int = 60
scheduler_batch_size: int = 50
dispatcher_interval_seconds: int = 10
Expand Down
182 changes: 182 additions & 0 deletions tests/test_background_workers.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,182 @@
"""Regression tests for single-owner background workers (#286).

Uvicorn runs lifespan once per worker process. Without a gate, N workers start
N schedulers / dispatchers / watchdogs against the same DB. These tests pin the
``enable_background_workers`` ownership model from #286.
"""

from __future__ import annotations

import asyncio
from unittest.mock import MagicMock, patch

import pytest
from fastapi import FastAPI

from openhands.automation.app import start_background_worker_tasks
from openhands.automation.config import ServiceSettings


async def _idle_loop(*_args, shutdown_event: asyncio.Event, **_kwargs) -> None:
"""Stand-in for scheduler/dispatcher/watchdog that exits on shutdown."""
await shutdown_event.wait()


def _settings(*, enable: bool) -> ServiceSettings:
return ServiceSettings(
enable_background_workers=enable,
base_url="https://example.test",
scheduler_interval_seconds=60,
dispatcher_interval_seconds=10,
watchdog_interval_seconds=60,
)


def _app_with_session_factory() -> FastAPI:
app = FastAPI()
app.state.session_factory = MagicMock(name="session_factory")
return app


@pytest.mark.asyncio
async def test_disabled_starts_no_background_workers() -> None:
"""Request-serving workers must not start scheduler/dispatcher/watchdog."""
app = _app_with_session_factory()
shutdown_event = asyncio.Event()

with (
patch("openhands.automation.app.scheduler_loop", new=_idle_loop),
patch("openhands.automation.app.dispatcher_loop", new=_idle_loop),
patch("openhands.automation.app.watchdog_loop", new=_idle_loop),
):
tasks = start_background_worker_tasks(
app, _settings(enable=False), shutdown_event
)

assert tasks == []
assert app.state.scheduler_task is None
assert app.state.dispatcher_task is None
assert app.state.watchdog_task is None


@pytest.mark.asyncio
async def test_enabled_starts_exactly_one_of_each_worker() -> None:
"""The dedicated background process owns one of each loop."""
app = _app_with_session_factory()
shutdown_event = asyncio.Event()

with (
patch("openhands.automation.app.scheduler_loop", new=_idle_loop),
patch("openhands.automation.app.dispatcher_loop", new=_idle_loop),
patch("openhands.automation.app.watchdog_loop", new=_idle_loop),
):
tasks = start_background_worker_tasks(
app, _settings(enable=True), shutdown_event
)

try:
names = [name for name, _task in tasks]
assert names == ["scheduler", "dispatcher", "watchdog"]
assert app.state.scheduler_task is tasks[0][1]
assert app.state.dispatcher_task is tasks[1][1]
assert app.state.watchdog_task is tasks[2][1]
assert all(not task.done() for _name, task in tasks)
finally:
shutdown_event.set()
await asyncio.gather(*(task for _name, task in tasks))


@pytest.mark.asyncio
async def test_multi_worker_simulation_single_owner() -> None:
"""Simulate replicas×workers: only the process with the flag owns loops.

Three "workers" start; two are request-serving (flag off) and one is the
dedicated background owner (flag on). Across the deployment there must be
exactly one scheduler, one dispatcher, and one watchdog.
"""
worker_flags = (False, False, True)
started: list[tuple[str, asyncio.Task]] = []
shutdown_events: list[asyncio.Event] = []

with (
patch("openhands.automation.app.scheduler_loop", new=_idle_loop),
patch("openhands.automation.app.dispatcher_loop", new=_idle_loop),
patch("openhands.automation.app.watchdog_loop", new=_idle_loop),
):
for enable in worker_flags:
app = _app_with_session_factory()
shutdown_event = asyncio.Event()
shutdown_events.append(shutdown_event)
started.extend(
start_background_worker_tasks(
app, _settings(enable=enable), shutdown_event
)
)

try:
by_name: dict[str, list[asyncio.Task]] = {
"scheduler": [],
"dispatcher": [],
"watchdog": [],
}
for name, task in started:
by_name[name].append(task)

assert len(by_name["scheduler"]) == 1
assert len(by_name["dispatcher"]) == 1
assert len(by_name["watchdog"]) == 1
assert len(started) == 3
finally:
for event in shutdown_events:
event.set()
if started:
await asyncio.gather(*(task for _name, task in started))


@pytest.mark.asyncio
async def test_multi_worker_all_enabled_still_duplicates() -> None:
"""Document the operator invariant: enabling on every worker still fans out.

The gate does not elect a leader. If every uvicorn worker leaves the flag
at its default (True), each still starts a full set of loops — same as
before #286 for single-process deploys, but wrong for multi-worker.
"""
started: list[tuple[str, asyncio.Task]] = []
shutdown_events: list[asyncio.Event] = []

with (
patch("openhands.automation.app.scheduler_loop", new=_idle_loop),
patch("openhands.automation.app.dispatcher_loop", new=_idle_loop),
patch("openhands.automation.app.watchdog_loop", new=_idle_loop),
):
for _ in range(3):
app = _app_with_session_factory()
shutdown_event = asyncio.Event()
shutdown_events.append(shutdown_event)
started.extend(
start_background_worker_tasks(
app, _settings(enable=True), shutdown_event
)
)

try:
assert len(started) == 9 # 3 workers × 3 loops
assert sum(1 for name, _ in started if name == "scheduler") == 3
assert sum(1 for name, _ in started if name == "dispatcher") == 3
assert sum(1 for name, _ in started if name == "watchdog") == 3
finally:
for event in shutdown_events:
event.set()
await asyncio.gather(*(task for _name, task in started))


def test_enable_background_workers_defaults_true() -> None:
"""Local / single-process deploys keep prior behavior without new env."""
assert ServiceSettings().enable_background_workers is True


def test_enable_background_workers_env_false(monkeypatch: pytest.MonkeyPatch) -> None:
monkeypatch.setenv("AUTOMATION_ENABLE_BACKGROUND_WORKERS", "false")
# ServiceSettings reads env at construction via pydantic-settings
settings = ServiceSettings()
assert settings.enable_background_workers is False
Loading