import time from concurrent import futures from concurrent.futures import Executor, ThreadPoolExecutor from typing import cast import redis from jb.config import ( CONSUMER_GROUP, CONSUMER_NAME, JB_EVENTS_STREAM, ) from jb.decorators import LOG, get_redis from jb.flow.assignment_tasks import process_assignment_submitted from jb.flow.monitoring import emit_error_event from jb.models.event import MTurkEvent StreamMessages = list[tuple[str, list[tuple[bytes, dict[bytes, bytes]]]]] def process_mturk_events_task(): executor = ThreadPoolExecutor(max_workers=5) create_consumer_group() while True: try: process_mturk_events(executor=executor) except Exception as e: LOG.exception(e) finally: time.sleep(1) def process_mturk_events(executor: Executor): while True: n = process_mturk_events_chunk(executor=executor) if n is None or n < 10: break def create_consumer_group(): try: redis_client = get_redis() redis_client.xgroup_create( JB_EVENTS_STREAM, CONSUMER_GROUP, id="0", mkstream=True ) except redis.exceptions.ResponseError as e: if "BUSYGROUP Consumer Group name already exists" in str(e): pass # group already exists else: raise def process_mturk_events_chunk(executor: Executor) -> int | None: redis_client = get_redis() msgs_raw = redis_client.xreadgroup( groupname=CONSUMER_GROUP, consumername=CONSUMER_NAME, streams={JB_EVENTS_STREAM: ">"}, count=10, ) if not msgs_raw: return None msgs = cast(StreamMessages, msgs_raw)[0][1] # The queue, we only have 1 fs = [] for msg in msgs: msg_id, data = msg msg_json: str = data["data"] event = MTurkEvent.model_validate_json(json_data=msg_json) if event.event_type == "AssignmentSubmitted": fs.append( executor.submit(process_assignment_submitted_event, event, str(msg_id)) ) else: LOG.info(f"Discarding {event}") redis_client.xdel(JB_EVENTS_STREAM, msg_id) futures.wait(fs, timeout=60) return len(msgs) def process_assignment_submitted_event(event: MTurkEvent, msg_id: str): from jb.decorators import AM, AMTM, BM, HM redis_client = get_redis() try: process_assignment_submitted(amtm=AMTM, am=AM, hm=HM, bm=BM, event=event) except Exception as e: LOG.exception(f"{event.amt_assignment_id=}, {e=}") emit_error_event( event_type="failed_process_assignment_submitted", amt_hit_type_id=event.amt_hit_type_id, ) redis_client.xackdel(JB_EVENTS_STREAM, CONSUMER_GROUP, msg_id)