Queue
Queue HTTP endpoints for the queue-mediated agent protocol.
Classes:
-
HeartbeatRequest–Request body for POST /queue/heartbeat.
-
NextTaskRequest–Request body for POST /queue/next.
Functions:
-
get_drain_event–Retrieve the drain asyncio.Event from the app state.
-
get_task_queue–Retrieve the TaskQueue from app state.
-
queue_complete–Store a completed task decision and signal the drain loop.
-
queue_heartbeat–Extend the claim window for an in-progress task.
-
queue_next–Claim and return the next pending task for the requesting agent.
HeartbeatRequest
pydantic-model
¶
Bases: BaseModel
Request body for POST /queue/heartbeat.
Attributes:
Parameters:
-
task_id(str) –
Show JSON schema:
Fields:
-
task_id(str)
NextTaskRequest
pydantic-model
¶
Bases: BaseModel
Request body for POST /queue/next.
Attributes:
Parameters:
-
agent_url(str) –
Show JSON schema:
Fields:
-
agent_url(str)
get_drain_event
¶
Retrieve the drain asyncio.Event from the app state.
Parameters:
-
request(Request) –The incoming FastAPI request.
Returns:
-
Event | None–The drain asyncio.Event, or None if not yet initialized.
get_task_queue
¶
get_task_queue(request: Request) -> TaskQueue
Retrieve the TaskQueue from app state.
Parameters:
-
request(Request) –The incoming FastAPI request.
Returns:
-
TaskQueue–The TaskQueue instance that is attached to the app state.
queue_complete
async
¶
queue_complete(decision: DecisionMessage, task_queue: TaskQueue = Depends(get_task_queue), drain_event: Event | None = Depends(get_drain_event)) -> Response
Store a completed task decision and signal the drain loop.
Parameters:
-
decision(DecisionMessage) –The agent's DecisionMessage.
-
task_queue(TaskQueue, default:Depends(get_task_queue)) –Task queue from the app state (injected).
-
drain_event(Event | None, default:Depends(get_drain_event)) –Drain loop event from the app state (injected).
Returns:
-
Response–202 Accepted.
queue_heartbeat
async
¶
queue_heartbeat(body: HeartbeatRequest, task_queue: TaskQueue = Depends(get_task_queue)) -> Response
Extend the claim window for an in-progress task.
Parameters:
-
body(HeartbeatRequest) –Request body containing the task ID.
-
task_queue(TaskQueue, default:Depends(get_task_queue)) –Task queue from app state (injected).
Returns:
-
Response–202 Accepted.
queue_next
async
¶
queue_next(body: NextTaskRequest, task_queue: TaskQueue = Depends(get_task_queue)) -> Response
Claim and return the next pending task for the requesting agent.
Parameters:
-
body(NextTaskRequest) –Request body containing the agent URL.
-
task_queue(TaskQueue, default:Depends(get_task_queue)) –Task queue from the app state (injected).
Returns:
-
Response–200 with TaskMessage JSON when a task is available, 204 when the queue is empty.