Skip to content
Closed
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
23 changes: 19 additions & 4 deletions backend/app/app_jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,14 +26,29 @@ def runner_script() -> Path:
return Path(__file__).resolve().parent.parent / "scripts" / "app-job-runner.py"


def runner_command(app_id: int, job_path: Path) -> list[str]:
return [sys.executable, str(runner_script()), str(app_id), str(job_path)]
def runner_command(
app_id: int, job_path: Path, *, wait_for_ready: bool = False,
) -> list[str]:
"""Build the common supervisor command for one app job.

Bootstrap installs happen inside FastAPI's lifespan, before the server can
answer the capability calls the supervisor makes. Only that launch path
needs to wait for the already-defined readiness contract; ordinary cron and
manual jobs run against an already-serving backend.
"""
command = [sys.executable, str(runner_script())]
if wait_for_ready:
command.append("--wait-for-ready")
command.extend((str(app_id), str(job_path)))
return command

def launch_app_job(app_id: int, job_path: Path, source_dir: Path):

def launch_app_job(
app_id: int, job_path: Path, source_dir: Path, *, wait_for_ready: bool = False,
):
"""Launch the common wrapper detached from the API worker's pipes."""
return subprocess.Popen(
runner_command(app_id, job_path),
runner_command(app_id, job_path, wait_for_ready=wait_for_ready),
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
cwd=str(source_dir),
Expand Down
14 changes: 12 additions & 2 deletions backend/app/install.py
Original file line number Diff line number Diff line change
Expand Up @@ -3349,8 +3349,18 @@ async def install_from_manifest(
try:
from app.app_jobs import launch_app_job
source = Path(app.source_dir)
launch_app_job(app.id, source / job_name, source)
warnings.append("initialization started")
# Bootstrap runs inside FastAPI lifespan, before this backend can answer
# the supervisor's scoped capability calls. Keep that ordering detail in
# the generic runner: it waits for the existing readiness signal before
# starting. Interactive installs already happen against a live server.
wait_for_ready = source == "bootstrap"
launch_app_job(
app.id, source / job_name, source, wait_for_ready=wait_for_ready,
)
warnings.append(
"initialization waiting for startup readiness"
if wait_for_ready else "initialization started"
)
except Exception as exc:
log.exception("install: initialization job failed to start")
warnings.append(f"initialization failed to start — {exc!r}")
Expand Down
40 changes: 36 additions & 4 deletions backend/scripts/app-job-runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import subprocess
import sys
import tempfile
import time
import urllib.request
import uuid
from pathlib import Path
Expand All @@ -29,6 +30,7 @@
# cron never restarts the container for us.
SUPERVISOR_LOG = DATA_DIR / "cron-logs" / "app-jobs.log"
SUPERVISOR_LOG_CAP = 2 * 1024 * 1024
READY_WAIT_SECONDS = 90


def _log(app_id: object, message: str) -> None:
Expand All @@ -52,6 +54,29 @@ def _start_ticks(pid: int) -> int:
return int(tail[19])


def _wait_for_ready(timeout_seconds: int = READY_WAIT_SECONDS) -> bool:
"""Wait only for the platform startup dependency bootstrap jobs require.

A bootstrap install runs during FastAPI lifespan, while the app-job runner
needs the backend to mint a scoped token and return job context. `/api/ready`
is the platform's existing readiness contract; polling it here avoids a
startup ordering race without adding a second scheduler or retry system.
"""
deadline = time.monotonic() + max(0, timeout_seconds)
while True:
try:
request = urllib.request.Request(f"{API_BASE_URL}/api/ready")
with urllib.request.urlopen(request, timeout=2) as response:
if response.status == 200:
return True
except Exception:
pass
remaining = deadline - time.monotonic()
if remaining <= 0:
return False
time.sleep(min(1, remaining))


def _atomic_json(path: Path, value: dict) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
fd, tmp = tempfile.mkstemp(dir=str(path.parent), prefix=".lease-", suffix=".tmp")
Expand Down Expand Up @@ -243,11 +268,15 @@ def _sandboxed_command(


def run() -> int:
if len(sys.argv) != 3 or not re.fullmatch(r"[0-9]+", sys.argv[1]):
_log(sys.argv[1] if len(sys.argv) > 1 else "?", "rejected: bad argv")
argv = sys.argv[1:]
wait_for_ready = argv[:1] == ["--wait-for-ready"]
if wait_for_ready:
argv = argv[1:]
if len(argv) != 2 or not re.fullmatch(r"[0-9]+", argv[0]):
_log(argv[0] if argv else "?", "rejected: bad argv")
return 2
app_id = int(sys.argv[1])
job = Path(sys.argv[2])
app_id = int(argv[0])
job = Path(argv[1])
if job.is_symlink():
_log(app_id, f"rejected: symlinked job {job}")
return 2
Expand Down Expand Up @@ -283,6 +312,9 @@ def run() -> int:
"job": str(resolved),
})
try:
if wait_for_ready and not _wait_for_ready():
_log(app_id, "failed: timed out waiting for platform readiness")
return 4
app_token = _mint_app_token(app_id)
if not app_token:
_log(app_id, "failed: could not mint app token (backend down or bad service token)")
Expand Down
5 changes: 4 additions & 1 deletion backend/tests/test_app_capabilities.py
Original file line number Diff line number Diff line change
Expand Up @@ -237,5 +237,8 @@ def test_matching_digest_is_persisted_with_explicit_system_identity(
assert app.system_app is True
assert app.capability_contract == contract
assert response.json()["capability_contract"] == contract
launch.assert_called_once()
launch.assert_called_once_with(
app.id, Path(app.source_dir) / "memory-job.sh", Path(app.source_dir),
wait_for_ready=False,
)
assert "initialization started" in response.json()["warnings"]
64 changes: 64 additions & 0 deletions backend/tests/test_app_jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,15 @@ def test_cron_parser_resolves_supervised_command_to_real_job():
)


def test_only_bootstrap_commands_request_a_readiness_wait():
job = Path("/data/apps/memory/memory-job.sh")

assert app_jobs.runner_command(57, job, wait_for_ready=True)[-3:] == [
"--wait-for-ready", "57", str(job),
]
assert app_jobs.runner_command(57, job)[-2:] == ["57", str(job)]


def test_terminate_verifies_start_ticks_before_signalling(monkeypatch):
data_dir = Path(get_settings().data_dir)
leases = data_dir / "run" / "app-jobs" / "57"
Expand Down Expand Up @@ -109,6 +118,61 @@ def urlopen(request, timeout):
}


def test_bootstrap_waits_for_ready_before_minting_a_job_token(
tmp_path, monkeypatch,
):
runner = _load_runner()
data_dir = tmp_path / "data"
source = data_dir / "apps" / "memory"
source.mkdir(parents=True)
job = source / "memory-job.sh"
job.write_text("#!/bin/sh\nexit 0\n")
monkeypatch.setattr(runner, "DATA_DIR", data_dir)
monkeypatch.setattr(runner.os, "getsid", lambda _pid: os.getpid())
events = []
monkeypatch.setattr(
runner, "_wait_for_ready", lambda: events.append("ready") or True,
)
monkeypatch.setattr(
runner, "_mint_app_token", lambda _app_id: events.append("mint") or "token",
)
monkeypatch.setattr(runner, "_app_is_live", lambda *_args: True)
monkeypatch.setattr(
runner, "_job_context", lambda *_args: {"source_dir": str(source)},
)
monkeypatch.setattr(
runner.subprocess, "Popen", lambda *_args, **_kwargs: types.SimpleNamespace(wait=lambda: 0),
)
monkeypatch.setattr(runner.sys, "argv", [
"app-job-runner.py", "--wait-for-ready", "57", str(job),
])

assert runner.run() == 0
assert events == ["ready", "mint"]


def test_bootstrap_readiness_timeout_never_mints_a_job_token(
tmp_path, monkeypatch,
):
runner = _load_runner()
data_dir = tmp_path / "data"
source = data_dir / "apps" / "memory"
source.mkdir(parents=True)
job = source / "memory-job.sh"
job.write_text("#!/bin/sh\nexit 0\n")
monkeypatch.setattr(runner, "DATA_DIR", data_dir)
monkeypatch.setattr(runner.os, "getsid", lambda _pid: os.getpid())
monkeypatch.setattr(runner, "_wait_for_ready", lambda: False)
minted = []
monkeypatch.setattr(runner, "_mint_app_token", lambda _app_id: minted.append(True))
monkeypatch.setattr(runner.sys, "argv", [
"app-job-runner.py", "--wait-for-ready", "57", str(job),
])

assert runner.run() == 4
assert minted == []


def test_wrapper_publishes_lease_before_live_check_and_cleans_it(
tmp_path, monkeypatch,
):
Expand Down