fix: receive emergency controls while ordered chat acceptance settles

This commit is contained in:
Ouroboros 2026-09-26 06:12:30 +03:00
parent 192a5780a1
commit 80fdcee156
6 changed files with 345 additions and 40 deletions

View file

@ -142,7 +142,7 @@ Rationale: `server.py` owns process startup/lifespan/static mounting, while `gat
### WebSocket protocol
`/ws` is the live browser delivery channel, not a durable state owner: queue, task, Project, review, skill, settings, cost, and update modules persist their own truth, and REST/history endpoints reconstruct it after reload or disconnection. `gateway/contracts.py` describes the frozen envelope shapes and message-type index for Python/JavaScript parity; `gateway/ws.py` performs the actual transport checks — incoming text must decode to a JSON object, extension types must parse as an owned namespace, and built-in `chat` or `command` frames must carry a non-empty payload before entering the message bridge. A `chat` frame's acceptance (ingress lock, locked canonical row, enqueue, echo) runs off the event loop through `gateway._helpers.run_sync_to_completion`, so a held lock or a slow append never stalls other handlers, while a `command` frame stays an inline queue put. The authentication middleware above governs socket admission; the public socket never receives the Host Service token or the owned Claudexor daemon token.
`/ws` is the live browser delivery channel, not a durable state owner: queue, task, Project, review, skill, settings, cost, and update modules persist their own truth, and REST/history endpoints reconstruct it after reload or disconnection. `gateway/contracts.py` describes the frozen envelope shapes and message-type index for Python/JavaScript parity; `gateway/ws.py` performs the actual transport checks — incoming text must decode to a JSON object, extension types must parse as an owned namespace, and built-in `chat` or `command` frames must carry a non-empty payload before entering the message bridge. A `chat` frame's acceptance (ingress lock, locked canonical row, enqueue, echo) runs off the event loop through `gateway._helpers.run_sync_to_completion`, so a held lock or a slow append never stalls other handlers. Each socket chains its chat acceptances so they settle in receive order, but its receive loop never awaits them: a `command` frame (Panic and Restart share the SPA's one socket with chat) stays an inline queue put, admitted even while that socket's earlier chat is held. Disconnect or cancellation returns only after every received chat settled (`settle_to_completion`). The authentication middleware above governs socket admission; the public socket never receives the Host Service token or the owned Claudexor daemon token.
The browser constructs one socket for the whole SPA. Feature modules subscribe before connection, and the initial complete Project chat-id set is fetched before the first open so an early Project frame cannot be mistaken for Main traffic. `ws.on(type, listener)` stores listeners in insertion-ordered sets and returns a disposer; emission uses a listener snapshot, so a listener added during dispatch does not receive the current frame and disposing one listener cannot skip its neighbor. Every decoded frame first reaches the generic `message` event and then its type-specific event, which lets Widgets consume reviewed namespaced events without duplicating the socket.

View file

@ -63,6 +63,6 @@ Startup and throttled maintenance reconcile three residue classes, the ~600 s pa
Cooperative project checkpointing has two equivalent quiescence triggers: a host-minted genesis or cooperative tree is checked when its root settles with no live descendants, and again when the last child settles beneath an already-terminal root — a root-scope budget stop terminalizes the root before its children, so a root-only trigger would see a live tree once and never return. The bounded git chain runs on a daemon thread, revalidates quiescence under the queue lock immediately before mutation, and replays a trigger that arrives during an in-flight check. Only host-minted project roots are eligible: owner-attached folders are never auto-committed, credential-shaped files stay excluded and disclosed, and every material success, skip or error receives a durable receipt.
Bridge intake is batched: `LocalChatBridge.get_updates` blocks only for the first item, then drains a bounded ready snapshot. Every update keeps a monotonic id and its own reply transport (`activate_update_transport`); malformed items are logged and skipped individually. Web owner ingress durably writes the canonical `chat.jsonl` row with `log_chat(require_write=True, ensure_record_boundary=True)` (as the named ingress does, so a torn prior tail cannot swallow it) before queue handoff or `ingress_accepted` echo, passing that row as the dequeue witness. Each async ingress-lock caller (`gateway/ws.py` chat, Host Service delivery, late quiz answers) runs via `gateway._helpers.run_sync_to_completion` off the ASGI loop: unrelated HTTP and sockets stay responsive, later frames on a socket remain ordered, and cancellation settles row → queue → echo. A command on another socket remains an inline queue put. A failed handler requeues the unprocessed tail with its ids for the next read; accepted `/restart` does likewise, while `/panic` never delays the hard stop for handback. Requeue is memory-only and lost on process exit; accepted rows on disk survive without automatic replay.
Bridge intake is batched: `LocalChatBridge.get_updates` blocks only for the first item, then drains a bounded ready snapshot. Every update keeps a monotonic id and its own reply transport (`activate_update_transport`); malformed items are logged and skipped individually. Web owner ingress durably writes the canonical `chat.jsonl` row with `log_chat(require_write=True, ensure_record_boundary=True)` (as the named ingress does, so a torn prior tail cannot swallow it) before queue handoff or `ingress_accepted` echo, passing that row as the dequeue witness. Each async ingress-lock caller (`gateway/ws.py` chat, Host Service delivery, late quiz answers) runs via `gateway._helpers.run_sync_to_completion` off the ASGI loop: unrelated HTTP and sockets stay responsive, a socket's later chats remain ordered, and cancellation settles row → queue → echo. A command on any socket remains an inline queue put. A failed handler requeues the unprocessed tail with its ids for the next read; accepted `/restart` does likewise, while `/panic` never delays the hard stop for handback. Requeue is memory-only and lost on process exit; accepted rows on disk survive without automatic replay.
The bridge recognizes `/panic`, `/restart`, `/review`, `/evolve [on|off]`, `/bg [start|stop|status]` and `/status`; all other text enters ordinary agent routing. External transports may invoke these commands only with positive owner identity and a transport-specific owner-chat binding, and the commands reuse runtime-mode, queue, cancellation and typed-result authority rather than implementing parallel control paths. Runtime logs rotate on the same supervisor tick and archive readers preserve their retained timelines; only explicitly isolated devtool roots may use the narrow rotation sentinel from §1.

View file

@ -106,7 +106,7 @@ Run roots are append-only outside `repo/` and live `data/`; the focused contract
than waiting for a growing file to reach EOF. HTTP admission and
materialization run their whole blocking operation off the event loop
(`gateway._helpers.run_sync_to_completion`); every async ingress-lock
caller uses it for locked row → queue → echo. Cancellation settles before release; other sockets
caller uses it for locked row → queue → echo. Cancellation settles before release; receive loops
stay responsive. Cancelling an HTTP waiter never cancels the admitted task. Directory exports carry a complete relative member/size/SHA manifest
plus a streamed ZIP (outputs above 50 MiB included); a changed file or
missing member is an explicit capture failure, while genesis LISTING is

View file

@ -35,25 +35,33 @@ async def run_sync_to_completion(function, /, *args, **kwargs):
try:
return await asyncio.shield(worker)
except asyncio.CancelledError:
# Starlette streams use level cancellation; shield that scope while
# also tolerating repeated raw asyncio Task.cancel() calls.
with anyio.CancelScope(shield=True):
while not worker.done():
try:
await asyncio.shield(worker)
except asyncio.CancelledError:
pass
except Exception:
break
try:
worker.result()
except Exception:
logging.getLogger(__name__).debug(
"Request worker failed while cancellation settled", exc_info=True,
)
await settle_to_completion(worker)
raise
async def settle_to_completion(task: asyncio.Task) -> None:
"""Wait until an owned task is done, through the caller's cancellation.
Its failure is logged, not raised; a cancelled task still raises.
"""
# Starlette streams use level cancellation; shield that scope while
# also tolerating repeated raw asyncio Task.cancel() calls.
with anyio.CancelScope(shield=True):
while not task.done():
try:
await asyncio.shield(task)
except asyncio.CancelledError:
pass
except Exception:
break
try:
task.result()
except Exception:
logging.getLogger(__name__).debug(
"Request worker failed while cancellation settled", exc_info=True,
)
def read_rotated_jsonl_entries(
live: pathlib.Path,
archive_dir: pathlib.Path,
@ -180,5 +188,5 @@ def stage_initial_task_attachments(
__all__ = (
"coerce_bool", "coerce_int", "iter_jsonl_objects", "json_error", "json_exception",
"read_rotated_jsonl_entries", "request_json_or", "request_drive_root", "request_repo_dir",
"stage_initial_task_attachments", "run_sync_to_completion",
"stage_initial_task_attachments", "run_sync_to_completion", "settle_to_completion",
)

View file

@ -3,6 +3,7 @@
from __future__ import annotations
import asyncio
import functools
import inspect
import json
import logging
@ -13,7 +14,7 @@ from typing import Any
from starlette.websockets import WebSocket, WebSocketDisconnect
from ouroboros.config import DATA_DIR
from ouroboros.gateway._helpers import run_sync_to_completion
from ouroboros.gateway._helpers import run_sync_to_completion, settle_to_completion
from ouroboros.utils import utc_now_iso
log = logging.getLogger(__name__)
@ -315,12 +316,45 @@ async def _dispatch_extension_message(
return True
def _initialization_notice() -> str:
return json.dumps({
"type": "chat",
"role": "system", "system_type": "initialization_notice",
"content": "⚠️ System is still initializing. Please wait a moment and try again.",
"ts": utc_now_iso(),
})
async def _accept_chat_after(websocket: WebSocket, previous: asyncio.Task | None, accept) -> None:
"""Run one chat frame's acceptance once this socket's previous one settled.
The web acceptance takes the single host's ingress lock and a locked durable
append (log_chat(require_write=True)); either may wait behind a skill delivery
scanning retained chat or a slow disk, so it runs off the ASGI loop, and the
receive loop never awaits it: a command frame on the same socket (Panic,
Restart) is admitted meanwhile. ``run_sync_to_completion`` keeps custody of
row → queue → echo; a failure answers the sender with the initialization notice.
"""
if previous is not None:
await asyncio.wait((previous,)) # its failure was its own notice
try:
await run_sync_to_completion(accept)
except Exception:
try:
await websocket.send_text(_initialization_notice())
except Exception:
log.debug("WebSocket closed before its chat failure notice", exc_info=True)
async def ws_endpoint(websocket: WebSocket) -> None:
await websocket.accept()
with _ws_lock:
_ws_clients.append(websocket)
total = len(_ws_clients)
log.info("WebSocket client connected (total: %d)", total)
# Tail of this socket's chat acceptances: each waits for the one before it,
# so its chat frames settle in receive order while the loop keeps receiving.
accepting: asyncio.Task | None = None
try:
while True:
data = await websocket.receive_text()
@ -391,15 +425,7 @@ async def ws_endpoint(websocket: WebSocket) -> None:
# land in chat.jsonl at the canonical-row writer.
client_surface["received_at"] = utc_now_iso()
task_metadata["client_surface"] = client_surface
# The web acceptance takes the single host's ingress lock and
# a locked durable append (log_chat(require_write=True)); either
# may wait behind a skill delivery scanning retained chat or a
# slow disk, so it must not run on the ASGI loop. The settled
# worker wait keeps its custody: a cancelled socket task still
# lets the canonical row, the queue item and the echo complete
# in that order, and this socket's frames stay ordered because
# the next receive waits for it.
await run_sync_to_completion(
accept = functools.partial(
bridge.ui_send,
payload,
sender_session_id=str(msg.get("sender_session_id", "") or ""),
@ -411,17 +437,14 @@ async def ws_endpoint(websocket: WebSocket) -> None:
chat_id=thread_id,
project_id=str(msg.get("project_id", "") or ""),
)
accepting = asyncio.create_task(_accept_chat_after(websocket, accepting, accept))
else:
# A command is a queue put only (no lock, no durable write):
# it stays inline once received; another socket is not blocked by owner ingress.
# A command is a queue put only (no lock, no durable write). It is
# admitted on receipt, never behind this socket's pending chat
# acceptances: Panic and Restart ride it on the SPA's one socket.
bridge.ui_send(payload, broadcast=False)
except Exception:
await websocket.send_text(json.dumps({
"type": "chat",
"role": "system", "system_type": "initialization_notice",
"content": "⚠️ System is still initializing. Please wait a moment and try again.",
"ts": utc_now_iso(),
}))
await websocket.send_text(_initialization_notice())
except WebSocketDisconnect:
pass
except Exception as exc:
@ -434,6 +457,10 @@ async def ws_endpoint(websocket: WebSocket) -> None:
pass
total = len(_ws_clients)
log.info("WebSocket client disconnected (total: %d)", total)
if accepting is not None:
# Received chat frames keep custody through disconnect or cancellation:
# the socket task returns only after each settled row → queue → echo.
await settle_to_completion(accepting)
__all__ = [

View file

@ -8,8 +8,10 @@ event loop with it: every other HTTP and WebSocket handler waited, Stop and
Panic included. The acceptance now runs off the loop through
``gateway._helpers.run_sync_to_completion`` and keeps its custody: the row, the
queue item and the echo still complete in that order, even when the socket task
is cancelled while the acceptance stands on the lock. A command on another
socket stays an inline queue put; a later frame on this socket waits its turn.
is cancelled while the acceptance stands on the lock. A later chat frame on
this socket waits its turn, but the receive loop never waits: the SPA sends
Panic and Restart on the SAME socket as chat, so a command frame is admitted
inline while that socket's chat acceptance is still held (review F2).
A late quiz answer (``POST /api/decisions`` on an expired card) is forwarded
through the named ingress, which waits for that same lock — behind a socket
@ -24,10 +26,12 @@ import json
import threading
import time
import anyio
import pytest
from starlette.applications import Starlette
from starlette.routing import Route, WebSocketRoute
from starlette.testclient import TestClient
from starlette.websockets import WebSocketDisconnect
from ouroboros.gateway.state import api_health
from ouroboros.gateway.ws import ws_endpoint
@ -216,6 +220,272 @@ def test_an_acceptance_failure_still_answers_the_socket(bridge):
assert echoes == [] and bridge.get_updates(offset=0, timeout=0) == []
def _record_handoffs(bridge, release: threading.Event | None = None):
"""Instrument the bridge callback the socket hands every frame to: record each
chat acceptance entering it (held until ``release`` when given) and each command
admitted. Nothing consumes the queue, so an admitted /panic is never executed."""
handoffs: list[tuple[str, str]] = []
chat_entered = threading.Event()
original = bridge.ui_send
def ui_send(text, **kwargs):
if kwargs.get("broadcast") is False:
handoffs.append(("command", text))
else:
handoffs.append(("chat", kwargs["client_message_id"]))
chat_entered.set()
if release is not None:
assert release.wait(5), "the test never released the held acceptance"
return original(text, **kwargs)
bridge.ui_send = ui_send
return handoffs, chat_entered
def _wait_for(predicate, timeout: float = 3.0) -> bool:
deadline = time.monotonic() + timeout
while not predicate() and time.monotonic() < deadline:
time.sleep(0.01)
return predicate()
def _command_frame(cmd: str) -> str:
return json.dumps({"type": "command", "cmd": cmd})
def _queued(bridge) -> list[str]:
with bridge._inbox.mutex:
return [item["client_message_id"] or item["text"] for item in bridge._inbox.queue]
def test_a_same_socket_panic_is_admitted_while_its_chat_acceptance_is_held(bridge):
"""Consumer regression (TZ-1 PR1 review F2): ``confirmAndSendPanic`` uses the one
SPA socket that also carries chat. With that socket's chat acceptance standing on
a held ingress lock, its next frame — Panic — must still be received and admitted
on one event loop; the Emergency Stop never waits for ordinary acceptance work."""
held, release, timed_out = threading.Event(), threading.Event(), threading.Event()
frames, echoed = _witness_custody(bridge, chat_echoes=1)
handoffs, chat_entered = _record_handoffs(bridge)
def hold_ingress(): # a skill delivery mid-scan under the single host's ingress lock
with message_bus._INGRESS_LOCK:
held.set()
if not release.wait(timeout=6.0):
timed_out.set()
holder = threading.Thread(target=hold_ingress, name="ingress-holder", daemon=True)
holder.start()
assert held.wait(5)
app = Starlette(routes=[WebSocketRoute("/ws", ws_endpoint)])
try:
with TestClient(app) as client, client.websocket_connect("/ws") as owner:
owner.send_text(_chat_frame("held-1"))
assert chat_entered.wait(5), "the chat frame never reached the bridge"
owner.send_text(_command_frame("/panic"))
assert _wait_for(lambda: ("command", "/panic") in handoffs), (
"Panic waited behind this socket's held chat acceptance")
assert not timed_out.is_set() and message_bus._INGRESS_LOCK.locked()
assert frames == [] and _row_ids(bridge.drive) == [] and _queued(bridge) == ["/panic"]
release.set()
holder.join(5)
assert echoed.wait(5), frames
finally:
release.set()
assert handoffs == [("chat", "held-1"), ("command", "/panic")]
[(echo, queued, rows)] = frames
assert echo["client_message_id"] == "held-1" and echo["ingress_accepted"] is True
assert queued == ["/panic", "held-1"] and rows == ["held-1"] # row → queue → echo
def test_same_socket_chat_frames_settle_in_receive_order_while_commands_pass(bridge):
"""Later chat frames on the socket wait for the held one — they never race it for
the lock — while commands between them are admitted on receipt. After release the
rows, queue items and echoes follow receive order, each echo after its row and item."""
release = threading.Event()
frames, echoed = _witness_custody(bridge, chat_echoes=3)
handoffs, chat_entered = _record_handoffs(bridge, release)
app = Starlette(routes=[WebSocketRoute("/ws", ws_endpoint)])
try:
with TestClient(app) as client, client.websocket_connect("/ws") as owner:
owner.send_text(_chat_frame("c-1", "first"))
assert chat_entered.wait(5)
owner.send_text(_chat_frame("c-2", "second"))
owner.send_text(_command_frame("/restart"))
owner.send_text(_chat_frame("c-3", "third"))
owner.send_text(_command_frame("/status")) # its admission proves c-3 was received
assert _wait_for(lambda: len(handoffs) == 3), handoffs
assert handoffs == [("chat", "c-1"), ("command", "/restart"), ("command", "/status")]
assert frames == [] and _row_ids(bridge.drive) == []
release.set()
assert echoed.wait(5), frames
finally:
release.set()
assert handoffs[3:] == [("chat", "c-2"), ("chat", "c-3")]
assert _row_ids(bridge.drive) == ["c-1", "c-2", "c-3"]
assert [frame["client_message_id"] for frame, *_ in frames] == ["c-1", "c-2", "c-3"]
for frame, queued, rows in frames:
assert frame["client_message_id"] in queued and frame["client_message_id"] in rows, frames
updates = bridge.get_updates(offset=0, timeout=1)
assert [u["message"]["text"] for u in updates] == ["/restart", "/status", "first", "second", "third"]
def test_a_disconnect_settles_every_received_chat_before_the_socket_task_returns(bridge):
"""The owner closes the socket (TestClient then cancels the session's anyio scope)
while its first chat is held and a second waits behind it: the socket task returns
only after both settled row → queue → echo, in receive order."""
release = threading.Event()
frames, echoed = _witness_custody(bridge, chat_echoes=2)
handoffs, chat_entered = _record_handoffs(bridge, release)
app = Starlette(routes=[WebSocketRoute("/ws", ws_endpoint)])
with TestClient(app) as client:
owner = client.websocket_connect("/ws")
owner.__enter__()
closer = threading.Thread(target=owner.__exit__, args=(None, None, None), daemon=True)
try:
owner.send_text(_chat_frame("d-1", "first"))
assert chat_entered.wait(5)
owner.send_text(_chat_frame("d-2", "second"))
owner.send_text(_command_frame("/status")) # its admission proves d-2 was received
assert _wait_for(lambda: ("command", "/status") in handoffs), handoffs
closer.start() # disconnect, then the session's anyio scope is cancelled
closer.join(0.3)
assert closer.is_alive(), "the socket task returned with received chat frames unsettled"
assert frames == [] and _row_ids(bridge.drive) == []
release.set()
closer.join(5)
assert not closer.is_alive()
assert echoed.wait(5), frames
finally: # a failed assertion still closes the session instead of hanging the client
release.set()
if closer.ident is None:
closer.start()
closer.join(10)
assert _row_ids(bridge.drive) == ["d-1", "d-2"]
assert [frame["client_message_id"] for frame, *_ in frames] == ["d-1", "d-2"]
for frame, queued, rows in frames:
assert frame["client_message_id"] in queued and frame["client_message_id"] in rows, frames
class _ClosingSocket(_OpenSocket):
"""A socket whose client disconnects after its frames; later sends fail."""
def __init__(self, frames, *, disconnect: bool):
super().__init__(frames)
self.disconnect = disconnect
self.closed = False
async def receive_text(self):
if self.frames or not self.disconnect:
return await super().receive_text()
self.closed = True
raise WebSocketDisconnect(1001)
async def send_text(self, text):
if self.closed:
raise RuntimeError("websocket is closed")
await super().send_text(text)
@pytest.mark.parametrize("disconnect", [False, True])
def test_a_failed_acceptance_answers_in_order_and_the_next_chat_still_lands(bridge, disconnect):
"""A refused acceptance answers its sender with the initialization notice before
the next chat on the socket is accepted. If the socket is already gone, the
undeliverable notice is dropped quietly, the next chat still lands, and the
socket task ends cleanly — no task exception is left unretrieved."""
release = threading.Event()
timeline: list[str] = []
bridge._broadcast_fn = lambda payload: timeline.append(f"echo:{payload['client_message_id']}")
original = bridge.ui_send
def refuse_first(text, **kwargs):
if kwargs["client_message_id"] == "refused-1":
assert release.wait(5)
raise OSError("disk refused")
return original(text, **kwargs)
bridge.ui_send = refuse_first
socket = _ClosingSocket([_chat_frame("refused-1"), _chat_frame("ok-2", "next")], disconnect=disconnect)
sent = socket.send_text
async def send_text(text):
await sent(text)
timeline.append(f"sent:{json.loads(text)['system_type']}")
socket.send_text = send_text
unretrieved: list[dict] = []
async def main():
asyncio.get_running_loop().set_exception_handler(lambda _loop, context: unretrieved.append(context))
task = asyncio.create_task(ws_endpoint(socket))
await asyncio.sleep(0.05)
assert not task.done() and timeline == []
release.set()
deadline = time.monotonic() + 5.0
while "echo:ok-2" not in timeline and time.monotonic() < deadline:
await asyncio.sleep(0.01)
if disconnect:
await asyncio.wait_for(task, 5) # returns cleanly, nothing re-raised
else:
task.cancel()
with pytest.raises(asyncio.CancelledError):
await asyncio.wait_for(task, 5)
asyncio.run(main())
expected = ["echo:ok-2"] if disconnect else ["sent:initialization_notice", "echo:ok-2"]
assert timeline == expected
assert _row_ids(bridge.drive) == ["ok-2"] and not unretrieved
[update] = bridge.get_updates(offset=0, timeout=1)
assert update["message"]["text"] == "next"
@pytest.mark.parametrize("mode", ["asyncio", "anyio"])
def test_a_cancelled_socket_settles_every_received_chat_in_order(bridge, mode):
"""Cancel the socket task — raw ``Task.cancel()`` twice, or anyio's level
cancellation — while its first chat is held and a second waits behind it: the
task stays in custody until both settled in receive order."""
release = threading.Event()
frames, _echoed = _witness_custody(bridge, chat_echoes=2)
handoffs, chat_entered = _record_handoffs(bridge, release)
socket = _OpenSocket([_chat_frame("x-1", "first"), _chat_frame("x-2", "second")])
scopes: list[anyio.CancelScope] = []
async def endpoint():
with anyio.CancelScope() as scope:
scopes.append(scope)
await ws_endpoint(socket)
async def main():
task = asyncio.create_task(endpoint())
assert await asyncio.to_thread(chat_entered.wait, 5)
await asyncio.sleep(0.01) # the socket is now waiting for a third frame
if mode == "anyio":
scopes[0].cancel()
else:
task.cancel()
await asyncio.sleep(0)
task.cancel()
await asyncio.sleep(0.05)
assert not task.done(), "the cancelled socket task abandoned its received chats"
assert frames == [] and handoffs == [("chat", "x-1")]
release.set()
if mode == "anyio":
await asyncio.wait_for(task, 5) # the scope absorbs its own cancellation
else:
with pytest.raises(asyncio.CancelledError):
await asyncio.wait_for(task, 5)
try:
asyncio.run(main())
finally:
release.set()
assert handoffs == [("chat", "x-1"), ("chat", "x-2")]
assert _row_ids(bridge.drive) == ["x-1", "x-2"]
assert [frame["client_message_id"] for frame, *_ in frames] == ["x-1", "x-2"]
for frame, queued, rows in frames:
assert frame["client_message_id"] in queued and frame["client_message_id"] in rows, frames
_LATE_ID = "quiz_late_answer:task-1:q1"
_LATE_BODY = {"request_id": "r1", "decision_id": "quiz:task-1:q1", "option_index": 1}