173 lines
5.5 KiB
Python
173 lines
5.5 KiB
Python
"""Server scheduler composition (T12): config mapping and store wiring.
|
|
|
|
Scheduling stays off unless the config section or the CLI flag enables
|
|
it. When enabled, the service runs over the server's own stores behind
|
|
one composition lock, and the transport lifespan owns start/stop.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from pathlib import Path
|
|
|
|
from wf_config import WorkflowConfigFile
|
|
from wf_scheduling.lifecycle import DrainReport, SchedulerServiceConfig
|
|
from wf_scheduling.ownership import SchedulerOwnership
|
|
from wf_server.context import build_local_static_workflow_server
|
|
from wf_server.scheduling import (
|
|
build_scheduler_service,
|
|
scheduler_lifespan,
|
|
server_scheduler_config,
|
|
)
|
|
from wf_transport_rpc_http import create_rpc_app
|
|
|
|
|
|
def test_server_scheduler_config_disabled_by_default() -> None:
|
|
assert server_scheduler_config(None, False) is None
|
|
|
|
bare = WorkflowConfigFile.model_validate({"version": 1})
|
|
assert server_scheduler_config(bare, False) is None
|
|
|
|
file_disabled = WorkflowConfigFile.model_validate(
|
|
{"version": 1, "server": {"scheduler": {"enabled": False}}}
|
|
)
|
|
assert server_scheduler_config(file_disabled, False) is None
|
|
|
|
|
|
def test_server_scheduler_config_flag_enables_defaults_without_config() -> None:
|
|
resolved = server_scheduler_config(None, True)
|
|
|
|
assert resolved is not None
|
|
assert resolved.poll_interval_s == 1.0
|
|
assert resolved.capacity == 4
|
|
assert resolved.drain_grace_s == 30.0
|
|
assert resolved.auto_tick is True
|
|
|
|
|
|
def test_server_scheduler_config_maps_file_values() -> None:
|
|
config = WorkflowConfigFile.model_validate(
|
|
{
|
|
"version": 1,
|
|
"server": {
|
|
"scheduler": {
|
|
"enabled": True,
|
|
"poll_interval_s": 2.5,
|
|
"max_concurrent_runs": 8,
|
|
"drain_grace_s": 60.0,
|
|
},
|
|
},
|
|
}
|
|
)
|
|
|
|
resolved = server_scheduler_config(config, False)
|
|
|
|
assert resolved is not None
|
|
assert resolved.poll_interval_s == 2.5
|
|
assert resolved.capacity == 8
|
|
assert resolved.drain_grace_s == 60.0
|
|
assert resolved.auto_tick is True
|
|
|
|
|
|
def test_server_scheduler_config_flag_overrides_disabled_section() -> None:
|
|
config = WorkflowConfigFile.model_validate(
|
|
{
|
|
"version": 1,
|
|
"server": {
|
|
"scheduler": {"enabled": False, "max_concurrent_runs": 2},
|
|
},
|
|
}
|
|
)
|
|
|
|
resolved = server_scheduler_config(config, True)
|
|
|
|
assert resolved is not None
|
|
assert resolved.capacity == 2
|
|
assert resolved.auto_tick is True
|
|
|
|
|
|
def test_build_scheduler_service_wires_server_stores(tmp_path: Path) -> None:
|
|
server = build_local_static_workflow_server(tmp_path)
|
|
service = build_scheduler_service(server, SchedulerServiceConfig(auto_tick=False))
|
|
|
|
assert service.schedule_store.root == server.config.store_root
|
|
assert (
|
|
service.schedule_store.schedules_dir == server.config.store_root / "schedules"
|
|
)
|
|
assert service.run_store is server.stores.run_store
|
|
assert service.artifact_store is server.stores.artifact_store
|
|
assert service.runtime is server.context.runtime
|
|
assert service.ownership.lock_path == server.config.store_root / "scheduler.lock"
|
|
|
|
|
|
def test_build_scheduler_service_shares_schedule_store_with_server_api(
|
|
tmp_path: Path,
|
|
) -> None:
|
|
"""Enabled scheduler composition exposes the same store to API and poller."""
|
|
server = build_local_static_workflow_server(tmp_path)
|
|
service = build_scheduler_service(server, SchedulerServiceConfig(auto_tick=False))
|
|
|
|
assert server.api.schedules._schedule_store() is service.schedule_store
|
|
|
|
|
|
async def test_scheduler_service_start_stop_on_server_stores(
|
|
tmp_path: Path,
|
|
) -> None:
|
|
server = build_local_static_workflow_server(tmp_path)
|
|
service = build_scheduler_service(server, SchedulerServiceConfig(auto_tick=False))
|
|
|
|
await service.start()
|
|
try:
|
|
assert service.running is True
|
|
assert service.ownership.covers(
|
|
service.schedule_store.root, service.run_store.root
|
|
)
|
|
finally:
|
|
report = await service.stop()
|
|
|
|
assert isinstance(report, DrainReport)
|
|
assert service.running is False
|
|
# The lock is released: a fresh owner can acquire the same composition.
|
|
probe = SchedulerOwnership(tmp_path, owner="probe")
|
|
probe.acquire()
|
|
try:
|
|
assert probe.held is True
|
|
finally:
|
|
probe.release()
|
|
|
|
|
|
async def test_scheduler_lifespan_releases_lock_on_exit(tmp_path: Path) -> None:
|
|
server = build_local_static_workflow_server(tmp_path)
|
|
resolved = server_scheduler_config(None, True)
|
|
assert resolved is not None
|
|
|
|
async with scheduler_lifespan(server, resolved) as service:
|
|
assert service.running is True
|
|
|
|
probe = SchedulerOwnership(tmp_path, owner="probe")
|
|
probe.acquire()
|
|
try:
|
|
assert probe.held is True
|
|
finally:
|
|
probe.release()
|
|
|
|
|
|
async def test_rpc_scheduler_lifespan_factory_starts_and_stops_service(
|
|
tmp_path: Path,
|
|
) -> None:
|
|
"""The RPC app receives a callable that lazily owns scheduler startup."""
|
|
server = build_local_static_workflow_server(tmp_path)
|
|
config = SchedulerServiceConfig(auto_tick=False)
|
|
app = create_rpc_app(
|
|
server,
|
|
lifespan=lambda _app: scheduler_lifespan(server, config),
|
|
)
|
|
|
|
async with app.router.lifespan_context(app):
|
|
assert server.api.schedules._schedule_store().root == tmp_path
|
|
|
|
probe = SchedulerOwnership(tmp_path, owner="probe")
|
|
probe.acquire()
|
|
try:
|
|
assert probe.held is True
|
|
finally:
|
|
probe.release()
|