aboutsummaryrefslogtreecommitdiff
path: root/jb/flow/events.py
diff options
context:
space:
mode:
authorstuppie2026-09-02 15:08:19 -0600
committerstuppie2026-09-02 15:08:19 -0600
commit81cf8de5dbcb530718eb8687b5bfab44517a6757 (patch)
treeabd2266fbffd2d50e01b0ef8d29a6fb9c9242fe8 /jb/flow/events.py
parent81261e52931d055df5830e29b9bf5ef81ba9134e (diff)
downloadamt-jb-81cf8de5dbcb530718eb8687b5bfab44517a6757.tar.gz
amt-jb-81cf8de5dbcb530718eb8687b5bfab44517a6757.zip
stripping out some amt stuff. reject all submitted assignments. add direct work view
Diffstat (limited to 'jb/flow/events.py')
-rw-r--r--jb/flow/events.py64
1 files changed, 5 insertions, 59 deletions
diff --git a/jb/flow/events.py b/jb/flow/events.py
index 2825cb1..04c4bb6 100644
--- a/jb/flow/events.py
+++ b/jb/flow/events.py
@@ -1,16 +1,15 @@
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 TypedDict, 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.flow.assignment_tasks import process_assignment_submitted
@@ -39,16 +38,6 @@ def process_mturk_events_task():
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)
@@ -66,7 +55,7 @@ def create_consumer_group():
raise
-def process_mturk_events_chunk(executor: Executor) -> Optional[int]:
+def process_mturk_events_chunk(executor: Executor) -> int | None:
msgs_raw = REDIS.xreadgroup(
groupname=CONSUMER_GROUP,
consumername=CONSUMER_NAME,
@@ -97,7 +86,7 @@ def process_mturk_events_chunk(executor: Executor) -> Optional[int]:
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
try:
process_assignment_submitted(amtm=AMTM, am=AM, hm=HM, bm=BM, event=event)
@@ -109,46 +98,3 @@ def process_assignment_submitted_event(event: MTurkEvent, msg_id: str):
)
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)