aboutsummaryrefslogtreecommitdiff
path: root/jb/flow/events.py
diff options
context:
space:
mode:
Diffstat (limited to 'jb/flow/events.py')
-rw-r--r--jb/flow/events.py94
1 files changed, 19 insertions, 75 deletions
diff --git a/jb/flow/events.py b/jb/flow/events.py
index 2825cb1..bee63fd 100644
--- a/jb/flow/events.py
+++ b/jb/flow/events.py
@@ -1,18 +1,16 @@
-import logging
import time
from concurrent import futures
-from concurrent.futures import ThreadPoolExecutor, Executor
-from typing import Optional, cast, TypedDict
+from concurrent.futures import Executor, ThreadPoolExecutor
+from typing import cast
import redis
from jb.config import (
- JB_EVENTS_STREAM,
CONSUMER_GROUP,
CONSUMER_NAME,
- JB_EVENTS_FAILED_STREAM,
+ JB_EVENTS_STREAM,
)
-from jb.decorators import REDIS
+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
@@ -20,13 +18,6 @@ from jb.models.event import MTurkEvent
StreamMessages = list[tuple[str, list[tuple[bytes, dict[bytes, bytes]]]]]
-class PendingEntry(TypedDict):
- message_id: bytes
- consumer: bytes
- time_since_delivered: int
- times_delivered: int
-
-
def process_mturk_events_task():
executor = ThreadPoolExecutor(max_workers=5)
create_consumer_group()
@@ -34,21 +25,11 @@ def process_mturk_events_task():
try:
process_mturk_events(executor=executor)
except Exception as e:
- logging.exception(e)
+ LOG.exception(e)
finally:
time.sleep(1)
-def handle_pending_msgs_task():
- while True:
- try:
- handle_pending_msgs()
- except Exception as e:
- logging.exception(e)
- finally:
- time.sleep(60)
-
-
def process_mturk_events(executor: Executor):
while True:
n = process_mturk_events_chunk(executor=executor)
@@ -58,7 +39,10 @@ def process_mturk_events(executor: Executor):
def create_consumer_group():
try:
- REDIS.xgroup_create(JB_EVENTS_STREAM, CONSUMER_GROUP, id="0", mkstream=True)
+ 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
@@ -66,8 +50,9 @@ def create_consumer_group():
raise
-def process_mturk_events_chunk(executor: Executor) -> Optional[int]:
- msgs_raw = REDIS.xreadgroup(
+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: ">"},
@@ -89,66 +74,25 @@ def process_mturk_events_chunk(executor: Executor) -> Optional[int]:
executor.submit(process_assignment_submitted_event, event, str(msg_id))
)
else:
- logging.info(f"Discarding {event}")
- REDIS.xdel(JB_EVENTS_STREAM, msg_id)
+ 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 AMTM, AM, HM, BM
+ 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:
- logging.exception(f"{event.amt_assignment_id=}, {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.xackdel(JB_EVENTS_STREAM, CONSUMER_GROUP, msg_id)
-
-
-def handle_pending_msgs():
- # TODO!: This doesn't run at all.
-
- # Looks in the redis queue for msgs that
- # are pending (read by a consumer but not ACK). These prob failed.
- # Below is from chatgpt, idk if it works
- pending = cast(
- list[PendingEntry],
- REDIS.xpending_range(
- JB_EVENTS_STREAM, CONSUMER_GROUP, min="-", max="+", count=10
- ),
- )
-
- for entry in pending:
- msg_id = entry["message_id"]
- # Claim message if idle > 10 sec
- if entry["idle"] > 10_000: # milliseconds
- claimed = REDIS.xclaim(
- JB_EVENTS_STREAM,
- CONSUMER_GROUP,
- CONSUMER_NAME,
- min_idle_time=10_000,
- message_ids=[msg_id],
- )
- for cid, data in claimed:
- msg_json = data["data"]
- event = MTurkEvent.model_validate_json(msg_json)
- if event.event_type == "AssignmentSubmitted":
- # Try to process it again. If it fails, add
- # it to the failed stream, so maybe we can fix
- # and try again?
- try:
- process_assignment_submitted_event(event, cid)
- REDIS.xack(JB_EVENTS_STREAM, CONSUMER_GROUP, cid)
- except Exception as e:
- logging.exception(e)
- REDIS.xadd(JB_EVENTS_FAILED_STREAM, data)
- REDIS.xack(JB_EVENTS_STREAM, CONSUMER_GROUP, cid)
- else:
- logging.info(f"Discarding {event}")
- REDIS.xdel(JB_EVENTS_STREAM, msg_id)
+ redis_client.xackdel(JB_EVENTS_STREAM, CONSUMER_GROUP, msg_id)