aboutsummaryrefslogtreecommitdiff
path: root/jb/flow/events.py
diff options
context:
space:
mode:
authorMax Nanis2026-09-10 01:00:39 -0700
committerMax Nanis2026-09-10 01:00:39 -0700
commit4dca7296742b607e74f16e2f6484c51163a41ace (patch)
tree0839c16d0905deb587d3ee25ff93fd3ccf7eeeb6 /jb/flow/events.py
parent832aecaddce80e312095ecdb572d7756eb9df5e9 (diff)
downloadamt-jb-4dca7296742b607e74f16e2f6484c51163a41ace.tar.gz
amt-jb-4dca7296742b607e74f16e2f6484c51163a41ace.zip
using model_validator on GRLSettings. Allows null default values, then to asser them on load. Required so pydantic_settings can be loaded in tests without params
Diffstat (limited to 'jb/flow/events.py')
-rw-r--r--jb/flow/events.py9
1 files changed, 4 insertions, 5 deletions
diff --git a/jb/flow/events.py b/jb/flow/events.py
index 0eb91fd..f224c01 100644
--- a/jb/flow/events.py
+++ b/jb/flow/events.py
@@ -1,4 +1,3 @@
-import logging
import time
from concurrent import futures
from concurrent.futures import Executor, ThreadPoolExecutor
@@ -11,7 +10,7 @@ from jb.config import (
CONSUMER_NAME,
JB_EVENTS_STREAM,
)
-from jb.decorators import REDIS
+from jb.decorators import LOG, REDIS
from jb.flow.assignment_tasks import process_assignment_submitted
from jb.flow.monitoring import emit_error_event
from jb.models.event import MTurkEvent
@@ -26,7 +25,7 @@ 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)
@@ -71,7 +70,7 @@ def process_mturk_events_chunk(executor: Executor) -> int | None:
executor.submit(process_assignment_submitted_event, event, str(msg_id))
)
else:
- logging.info(f"Discarding {event}")
+ LOG.info(f"Discarding {event}")
REDIS.xdel(JB_EVENTS_STREAM, msg_id)
futures.wait(fs, timeout=60)
@@ -84,7 +83,7 @@ def process_assignment_submitted_event(event: MTurkEvent, msg_id: str):
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,