aboutsummaryrefslogtreecommitdiff
path: root/jb/flow
diff options
context:
space:
mode:
Diffstat (limited to 'jb/flow')
-rw-r--r--jb/flow/assignment_tasks.py17
-rw-r--r--jb/flow/events.py9
-rw-r--r--jb/flow/tasks.py17
3 files changed, 17 insertions, 26 deletions
diff --git a/jb/flow/assignment_tasks.py b/jb/flow/assignment_tasks.py
index bdebb3d..b3c820a 100644
--- a/jb/flow/assignment_tasks.py
+++ b/jb/flow/assignment_tasks.py
@@ -1,5 +1,4 @@
-import logging
-
+from jb.decorators import LOG
from jb.flow.monitoring import emit_assignment_event, emit_error_event
from jb.managers.amt import (
REJECT_MESSAGE_UNKNOWN_ASSIGNMENT,
@@ -27,7 +26,7 @@ def process_assignment_submitted(
#
# Step 1: Attempt to get the Assignment out of the API
#
- logging.info(f"{event=}")
+ LOG.info(f"{event=}")
# This is the assignment model from AMT. In the DB, we should only have
# the AssignmentStub
@@ -42,7 +41,7 @@ def process_assignment_submitted(
# It is not found in amt, either it is invalid, not yet submitted, or
# already been approved/rejected, so we just do nothing ...
# todo: maybe we confirm its state matches what we have in the db
- logging.warning(f"No assignment found on AMT: {event.amt_assignment_id}")
+ LOG.warning(f"No assignment found on AMT: {event.amt_assignment_id}")
emit_error_event(
event_type="assignment_not_found_in_amt",
amt_hit_type_id=event.amt_hit_type_id,
@@ -70,9 +69,7 @@ def review_hit(amtm: AMTManager, hm: HitManager, assignment: Assignment) -> None
hit, _ = amtm.get_hit_if_exists(amt_hit_id=assignment.amt_hit_id)
if hit is None:
- logging.warning(
- f"Hit not found when trying to review hit: {assignment.amt_hit_id}"
- )
+ LOG.warning(f"Hit not found when trying to review hit: {assignment.amt_hit_id}")
return
# Update the db
@@ -101,7 +98,7 @@ def reject_assignment(
event_type="failed_to_reject_assignment",
amt_hit_type_id=amt_hit_type_id,
)
- logging.exception(f"Failed to reject assignment: {amt_assignment_id}")
+ LOG.exception(f"Failed to reject assignment: {amt_assignment_id}")
# We just rejected this assignment, get it from amazon again
assignment = amtm.get_assignment(amt_assignment_id=amt_assignment_id)
@@ -112,7 +109,7 @@ def reject_assignment(
# need to create as assignment first ...
stub = am.get_stub_if_exists(amt_assignment_id=assignment.amt_assignment_id)
if stub is None:
- logging.warning(
+ LOG.warning(
f"Rejected assignment doesn't exist in DB. Creating ... : {amt_assignment_id}"
)
# Even if the assignment doesn't exist, the hit must ...
@@ -123,5 +120,5 @@ def reject_assignment(
emit_assignment_event(
status=AssignmentStatus.Rejected, amt_hit_type_id=amt_hit_type_id, reason=msg
)
- logging.warning(f"Rejected assignment: {amt_assignment_id}")
+ LOG.warning(f"Rejected assignment: {amt_assignment_id}")
return assignment
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,
diff --git a/jb/flow/tasks.py b/jb/flow/tasks.py
index 24e96d4..c555021 100644
--- a/jb/flow/tasks.py
+++ b/jb/flow/tasks.py
@@ -1,19 +1,14 @@
-import logging
import time
from typing import TypedDict, cast
from generalresearch.config import is_debug
-from jb.decorators import AMTM, HM, HQM, HTM, pg_config
+from jb.decorators import AMTM, HM, HQM, HTM, LOG, pg_config
from jb.flow.maintenance import check_hit_status
from jb.flow.monitoring import emit_hit_event, write_hit_gauge
from jb.models.definitions import HitStatus
from jb.models.hit import Hit, HitQuestion, HitType
-logging.basicConfig()
-logger = logging.getLogger()
-logger.setLevel(logging.INFO)
-
class HitRow(TypedDict):
amt_hit_id: str
@@ -34,7 +29,7 @@ def check_stale_hits():
params={"status": HitStatus.Assignable.value},
)
for hit in cast(list[HitRow], res):
- logging.info(f"check_stale_hits: {hit["amt_hit_id"]}")
+ LOG.info(f"check_stale_hits: {hit["amt_hit_id"]}")
check_hit_status(
amtm=AMTM,
amt_hit_id=hit["amt_hit_id"],
@@ -56,7 +51,7 @@ def check_expired_hits():
params={"status": HitStatus.Assignable.value},
)
for hit in cast(list[HitRow], res):
- logging.info(f"check_expired_hits: {hit["amt_hit_id"]}")
+ LOG.info(f"check_expired_hits: {hit["amt_hit_id"]}")
check_hit_status(
amtm=AMTM,
amt_hit_id=hit["amt_hit_id"],
@@ -87,7 +82,7 @@ def refill_hits() -> None:
assert hit_type.amt_hit_type_id
active_count = HM.get_active_count(hit_type_id=hit_type.id)
- logging.info(
+ LOG.info(
f"HitType: {hit_type.amt_hit_type_id}, {hit_type.min_active=}, active_count={active_count}"
)
write_hit_gauge(
@@ -97,7 +92,7 @@ def refill_hits() -> None:
)
if active_count < hit_type.min_active:
cnt_todo = hit_type.min_active - active_count
- logging.info(f"Refilling {cnt_todo} hits")
+ LOG.info(f"Refilling {cnt_todo} hits")
for _ in range(cnt_todo):
create_hit_from_hittype(hit_type)
@@ -109,6 +104,6 @@ def refill_hits_task():
check_stale_hits()
refill_hits()
except Exception as e:
- logging.exception(e)
+ LOG.exception(e)
finally:
time.sleep(5 * 60)