Drain authored mutations before shutdown #462
|
|
@ -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
|
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
|
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
|
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
|
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
|
`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
|
should preserve the `Idempotency-Key` header with the unchanged route and payload until a
|
||||||
|
|
|
||||||
26
src/main.py
26
src/main.py
|
|
@ -43,6 +43,28 @@ from src.later_store import LaterStore
|
||||||
from src.today_store import TodayPlanFull, TodayStore
|
from src.today_store import TodayPlanFull, TodayStore
|
||||||
from src.views import FRONTEND_BUILD, router as frontend_router
|
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
|
@asynccontextmanager
|
||||||
async def lifespan(_app: FastAPI):
|
async def lifespan(_app: FastAPI):
|
||||||
global _live_snapshot_task, _available_issue_snapshot_task
|
global _live_snapshot_task, _available_issue_snapshot_task
|
||||||
|
|
@ -60,6 +82,7 @@ async def lifespan(_app: FastAPI):
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
pass
|
pass
|
||||||
try:
|
try:
|
||||||
|
await _drain_authored_action_operations()
|
||||||
await gitea_proxy.stop_client()
|
await gitea_proxy.stop_client()
|
||||||
finally:
|
finally:
|
||||||
if _live_snapshot_task is live_task:
|
if _live_snapshot_task is live_task:
|
||||||
|
|
@ -83,6 +106,9 @@ AUTHORED_ACTION_IDEMPOTENCY_MAX_ENTRIES = 256
|
||||||
AUTHORED_ACTION_LEDGER_LOCK_TIMEOUT_SECONDS = float(
|
AUTHORED_ACTION_LEDGER_LOCK_TIMEOUT_SECONDS = float(
|
||||||
os.getenv("STACKCHAIN_IDEMPOTENCY_LOCK_TIMEOUT_SECONDS", "0.1")
|
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_MUTATION_TIMEOUT_SECONDS = 5.0
|
||||||
NOTIFICATION_PAGE_TIMEOUT_SECONDS = 5.0
|
NOTIFICATION_PAGE_TIMEOUT_SECONDS = 5.0
|
||||||
WORK_PAGE_TIMEOUT_SECONDS = 5.0
|
WORK_PAGE_TIMEOUT_SECONDS = 5.0
|
||||||
|
|
|
||||||
|
|
@ -306,3 +306,79 @@ async def test_application_shutdown_cancels_available_work_scan_before_transport
|
||||||
main._available_issue_snapshot_task = None
|
main._available_issue_snapshot_task = None
|
||||||
|
|
||||||
assert calls == ["start", "available scan cancelled", "stop (done)"]
|
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"]
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue
Block a user