| 1 | from helpers.api import ApiHandler, Request, Response |
| 2 | from helpers import message_queue as mq |
| 3 | from agent import AgentContext |
| 4 | from helpers.state_monitor_integration import mark_dirty_for_context |
| 5 | |
| 6 | class MessageQueueSend(ApiHandler): |
| 7 | """Send queued message(s) immediately.""" |
| 8 | |
| 9 | async def process(self, input: dict, request: Request) -> dict | Response: |
| 10 | context = AgentContext.get(input.get("context", "")) |
| 11 | if not context: |
| 12 | return Response("Context not found", status=404) |
| 13 | |
| 14 | if not mq.has_queue(context): |
| 15 | return {"ok": True, "message": "Queue empty"} |
| 16 | |
| 17 | item_id = input.get("item_id") |
| 18 | send_all = input.get("send_all", False) |
| 19 | |
| 20 | if send_all: |
| 21 | count = mq.send_all_aggregated(context) |
| 22 | if count: |
| 23 | mark_dirty_for_context(context.id, reason="message_queue_send_all") |
| 24 | return {"ok": True, "sent_count": count} |
| 25 | |
| 26 | # Send single item |
| 27 | item = mq.pop_item(context, item_id) if item_id else mq.pop_first(context) |
| 28 | if not item: |
| 29 | return Response("Item not found", status=404) |
| 30 | |
| 31 | mq.send_message(context, item) |
| 32 | mark_dirty_for_context(context.id, reason="message_queue_send") |
| 33 | return {"ok": True, "sent_item_id": item["id"]} |