From 0dd3932d14162312ed23ae56ac2a9bb81acbbeef Mon Sep 17 00:00:00 2001 From: timmy Date: Mon, 10 Aug 2026 07:12:51 +0000 Subject: [PATCH] fix: drain authored actions on shutdown (#461) --- README.md | 7 +++- src/main.py | 26 ++++++++++++ tests/test_gitea_transport.py | 76 +++++++++++++++++++++++++++++++++++ 3 files changed, 108 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 07d4c35..f987359 100644 --- a/README.md +++ b/README.md @@ -93,7 +93,12 @@ persistent, writable service directory (or set `STACKCHAIN_IDEMPOTENCY_DB` to an SQLite path); the local default is `.stackchain-state/idempotency.sqlite3`. Ledger reads and writes run outside the request event loop, and lock admission is bounded to 100 ms by default. Tune it with `STACKCHAIN_IDEMPOTENCY_LOCK_TIMEOUT_SECONDS`; keep the value below -route deadlines. Reservation contention returns retryable HTTP 503 with `Retry-After: 1`. +route deadlines. During a controlled shutdown, the server gives in-flight authored mutations +five seconds to finish and persist their ledger result before it closes the Gitea transport. +Set `STACKCHAIN_AUTHORED_ACTION_SHUTDOWN_GRACE_SECONDS` to match the service manager's +shutdown budget; when that deadline expires, remaining operations are cancelled and their +pending reservations continue to fail closed rather than being retried automatically. +Reservation contention returns retryable HTTP 503 with `Retry-After: 1`. If contention occurs after the upstream mutation, the dashboard fails closed with `Retry-After: 5` and asks the caller to verify the result before retrying. Direct API callers should preserve the `Idempotency-Key` header with the unchanged route and payload until a diff --git a/src/main.py b/src/main.py index 946d666..8f03a82 100644 --- a/src/main.py +++ b/src/main.py @@ -43,6 +43,28 @@ from src.later_store import LaterStore from src.today_store import TodayPlanFull, TodayStore from src.views import FRONTEND_BUILD, router as frontend_router + +async def _drain_authored_action_operations() -> None: + tasks = { + operation[1] + for operation in _authored_action_operations.values() + if not operation[1].done() + } + if not tasks: + return + try: + await asyncio.wait_for( + asyncio.gather(*tasks, return_exceptions=True), + timeout=AUTHORED_ACTION_SHUTDOWN_GRACE_SECONDS, + ) + except TimeoutError: + pass + finally: + for key, operation in list(_authored_action_operations.items()): + if operation[1] in tasks and operation[1].done(): + _authored_action_operations.pop(key, None) + + @asynccontextmanager async def lifespan(_app: FastAPI): global _live_snapshot_task, _available_issue_snapshot_task @@ -60,6 +82,7 @@ async def lifespan(_app: FastAPI): except asyncio.CancelledError: pass try: + await _drain_authored_action_operations() await gitea_proxy.stop_client() finally: if _live_snapshot_task is live_task: @@ -83,6 +106,9 @@ AUTHORED_ACTION_IDEMPOTENCY_MAX_ENTRIES = 256 AUTHORED_ACTION_LEDGER_LOCK_TIMEOUT_SECONDS = float( os.getenv("STACKCHAIN_IDEMPOTENCY_LOCK_TIMEOUT_SECONDS", "0.1") ) +AUTHORED_ACTION_SHUTDOWN_GRACE_SECONDS = float( + os.getenv("STACKCHAIN_AUTHORED_ACTION_SHUTDOWN_GRACE_SECONDS", "5.0") +) NOTIFICATION_MUTATION_TIMEOUT_SECONDS = 5.0 NOTIFICATION_PAGE_TIMEOUT_SECONDS = 5.0 WORK_PAGE_TIMEOUT_SECONDS = 5.0 diff --git a/tests/test_gitea_transport.py b/tests/test_gitea_transport.py index 85e8353..660d5be 100644 --- a/tests/test_gitea_transport.py +++ b/tests/test_gitea_transport.py @@ -306,3 +306,79 @@ async def test_application_shutdown_cancels_available_work_scan_before_transport main._available_issue_snapshot_task = None assert calls == ["start", "available scan cancelled", "stop (done)"] + + +@pytest.mark.anyio +async def test_application_shutdown_drains_authored_mutation_before_transport(monkeypatch): + calls = [] + started = asyncio.Event() + release = asyncio.Event() + + async def authored_mutation(): + started.set() + await release.wait() + calls.append("mutation persisted") + return {"id": 461} + + monkeypatch.setattr(gitea_proxy, "start_client", lambda: calls.append("start")) + + async def stop_client(): + calls.append("stop") + + monkeypatch.setattr(gitea_proxy, "stop_client", stop_client) + context = main.app.router.lifespan_context(main.app) + await context.__aenter__() + task = asyncio.create_task(authored_mutation()) + main._authored_action_operations["shutdown-drain-461"] = ( + ("comment", "stackchain/dashboard", 461), + task, + 0.0, + ) + await started.wait() + + shutdown = asyncio.create_task(context.__aexit__(None, None, None)) + await asyncio.sleep(0) + + assert not shutdown.done() + assert calls == ["start"] + + release.set() + await shutdown + + assert calls == ["start", "mutation persisted", "stop"] + + +@pytest.mark.anyio +async def test_application_shutdown_bounds_and_settles_stalled_authored_mutation(monkeypatch): + calls = [] + started = asyncio.Event() + + async def stalled_mutation(): + started.set() + try: + await asyncio.Event().wait() + finally: + calls.append("mutation cancelled") + + monkeypatch.setattr(main, "AUTHORED_ACTION_SHUTDOWN_GRACE_SECONDS", 0.01) + monkeypatch.setattr(gitea_proxy, "start_client", lambda: calls.append("start")) + + async def stop_client(): + calls.append("stop") + + monkeypatch.setattr(gitea_proxy, "stop_client", stop_client) + context = main.app.router.lifespan_context(main.app) + await context.__aenter__() + task = asyncio.create_task(stalled_mutation()) + main._authored_action_operations["stalled-shutdown-461"] = ( + ("review", "stackchain/dashboard", 461), + task, + 0.0, + ) + await started.wait() + + await asyncio.wait_for(context.__aexit__(None, None, None), timeout=0.2) + + assert task.cancelled() + assert main._authored_action_operations == {} + assert calls == ["start", "mutation cancelled", "stop"]