diff options
| author | Max Nanis | 2026-09-10 01:00:39 -0700 |
|---|---|---|
| committer | Max Nanis | 2026-09-10 01:00:39 -0700 |
| commit | 4dca7296742b607e74f16e2f6484c51163a41ace (patch) | |
| tree | 0839c16d0905deb587d3ee25ff93fd3ccf7eeeb6 /jb/flow/events.py | |
| parent | 832aecaddce80e312095ecdb572d7756eb9df5e9 (diff) | |
| download | amt-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.py | 9 |
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, |
