feat(mma): Implement AsyncEventQueue in events.py
This commit is contained in:
30
events.py
30
events.py
@@ -1,7 +1,8 @@
|
|||||||
"""
|
"""
|
||||||
Decoupled event emission system for cross-module communication.
|
Decoupled event emission system for cross-module communication.
|
||||||
"""
|
"""
|
||||||
from typing import Callable, Any, Dict, List
|
import asyncio
|
||||||
|
from typing import Callable, Any, Dict, List, Tuple
|
||||||
|
|
||||||
class EventEmitter:
|
class EventEmitter:
|
||||||
"""
|
"""
|
||||||
@@ -35,3 +36,30 @@ class EventEmitter:
|
|||||||
if event_name in self._listeners:
|
if event_name in self._listeners:
|
||||||
for callback in self._listeners[event_name]:
|
for callback in self._listeners[event_name]:
|
||||||
callback(*args, **kwargs)
|
callback(*args, **kwargs)
|
||||||
|
|
||||||
|
class AsyncEventQueue:
|
||||||
|
"""
|
||||||
|
Asynchronous event queue for decoupled communication using asyncio.Queue.
|
||||||
|
"""
|
||||||
|
def __init__(self):
|
||||||
|
"""Initializes the AsyncEventQueue with an internal asyncio.Queue."""
|
||||||
|
self._queue: asyncio.Queue = asyncio.Queue()
|
||||||
|
|
||||||
|
async def put(self, event_name: str, payload: Any = None):
|
||||||
|
"""
|
||||||
|
Puts an event into the queue.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
event_name: The name of the event.
|
||||||
|
payload: Optional data associated with the event.
|
||||||
|
"""
|
||||||
|
await self._queue.put((event_name, payload))
|
||||||
|
|
||||||
|
async def get(self) -> Tuple[str, Any]:
|
||||||
|
"""
|
||||||
|
Gets an event from the queue.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
A tuple containing (event_name, payload).
|
||||||
|
"""
|
||||||
|
return await self._queue.get()
|
||||||
|
|||||||
@@ -18,7 +18,8 @@ files = [
|
|||||||
"tests/test_headless_service.py",
|
"tests/test_headless_service.py",
|
||||||
"tests/test_performance_monitor.py",
|
"tests/test_performance_monitor.py",
|
||||||
"tests/test_token_usage.py",
|
"tests/test_token_usage.py",
|
||||||
"tests/test_layout_reorganization.py"
|
"tests/test_layout_reorganization.py",
|
||||||
|
"tests/test_async_events.py"
|
||||||
]
|
]
|
||||||
|
|
||||||
[categories.conductor]
|
[categories.conductor]
|
||||||
|
|||||||
47
tests/test_async_events.py
Normal file
47
tests/test_async_events.py
Normal file
@@ -0,0 +1,47 @@
|
|||||||
|
import asyncio
|
||||||
|
import pytest
|
||||||
|
from events import AsyncEventQueue
|
||||||
|
|
||||||
|
def test_async_event_queue_put_get():
|
||||||
|
"""Verify that an event can be asynchronously put and retrieved from the queue."""
|
||||||
|
async def run_test():
|
||||||
|
queue = AsyncEventQueue()
|
||||||
|
event_name = "test_event"
|
||||||
|
payload = {"data": "hello"}
|
||||||
|
|
||||||
|
await queue.put(event_name, payload)
|
||||||
|
ret_name, ret_payload = await queue.get()
|
||||||
|
|
||||||
|
assert ret_name == event_name
|
||||||
|
assert ret_payload == payload
|
||||||
|
|
||||||
|
asyncio.run(run_test())
|
||||||
|
|
||||||
|
def test_async_event_queue_multiple():
|
||||||
|
"""Verify that multiple events can be asynchronously put and retrieved in order."""
|
||||||
|
async def run_test():
|
||||||
|
queue = AsyncEventQueue()
|
||||||
|
|
||||||
|
await queue.put("event1", 1)
|
||||||
|
await queue.put("event2", 2)
|
||||||
|
|
||||||
|
name1, val1 = await queue.get()
|
||||||
|
name2, val2 = await queue.get()
|
||||||
|
|
||||||
|
assert name1 == "event1"
|
||||||
|
assert val1 == 1
|
||||||
|
assert name2 == "event2"
|
||||||
|
assert val2 == 2
|
||||||
|
|
||||||
|
asyncio.run(run_test())
|
||||||
|
|
||||||
|
def test_async_event_queue_none_payload():
|
||||||
|
"""Verify that an event with None payload works correctly."""
|
||||||
|
async def run_test():
|
||||||
|
queue = AsyncEventQueue()
|
||||||
|
await queue.put("no_payload")
|
||||||
|
name, payload = await queue.get()
|
||||||
|
assert name == "no_payload"
|
||||||
|
assert payload is None
|
||||||
|
|
||||||
|
asyncio.run(run_test())
|
||||||
Reference in New Issue
Block a user