From 81cf8de5dbcb530718eb8687b5bfab44517a6757 Mon Sep 17 00:00:00 2001 From: stuppie Date: Wed, 2 Sep 2026 15:08:19 -0600 Subject: stripping out some amt stuff. reject all submitted assignments. add direct work view --- jb/config.py | 1 - jb/flow/assignment_tasks.py | 451 ++------------------------------------------ jb/flow/events.py | 64 +------ jb/main.py | 1 - jb/managers/gr_api.py | 1 + jb/managers/thl.py | 49 +---- jb/models/auth.py | 2 + jb/settings.py | 3 - jb/views/common.py | 56 +++--- jb/views/tasks.py | 78 -------- 10 files changed, 55 insertions(+), 651 deletions(-) delete mode 100644 jb/views/tasks.py diff --git a/jb/config.py b/jb/config.py index a91494d..359f108 100644 --- a/jb/config.py +++ b/jb/config.py @@ -30,7 +30,6 @@ SUBSCRIPTION = { } JB_EVENTS_STREAM = "amt_jb_events" -JB_EVENTS_FAILED_STREAM = "amt_jb_events_failed" CONSUMER_GROUP = "amt-jb-0" # We'll only have 1 consumer atm, change this if we don't CONSUMER_NAME = "amt-jb-0" diff --git a/jb/flow/assignment_tasks.py b/jb/flow/assignment_tasks.py index 2db20f8..bdebb3d 100644 --- a/jb/flow/assignment_tasks.py +++ b/jb/flow/assignment_tasks.py @@ -1,37 +1,16 @@ import logging -import math -from datetime import timedelta -from typing import Optional -from generalresearch.models.thl.definitions import PayoutStatus, StatusCode1 -from generalresearch.models.thl.wallet.cashout_method import CashoutRequestInfo -from generalresearch.currency import USDCent - -from jb.flow.monitoring import emit_error_event, emit_assignment_event, emit_bonus_event +from jb.flow.monitoring import emit_assignment_event, emit_error_event from jb.managers.amt import ( - AMTManager, REJECT_MESSAGE_UNKNOWN_ASSIGNMENT, - REJECT_MESSAGE_NO_WORK, - NO_WORK_APPROVAL_MESSAGE, - REJECT_MESSAGE_BADDIE, - APPROVAL_MESSAGE, - BONUS_MESSAGE, -) -from jb.managers.thl import ( - get_user_blocked, - get_task_status, - user_cashout_request, - manage_pending_cashout, - get_user_blocked_or_not_exists, - get_wallet_balance_if_non_negative, + AMTManager, ) +from jb.managers.assignment import AssignmentManager +from jb.managers.bonus import BonusManager +from jb.managers.hit import HitManager from jb.models.assignment import Assignment from jb.models.definitions import AssignmentStatus from jb.models.event import MTurkEvent -from jb.managers.assignment import AssignmentManager -from jb.managers.hit import HitManager -from jb.managers.bonus import BonusManager -from jb.config import settings def process_assignment_submitted( @@ -42,10 +21,7 @@ def process_assignment_submitted( event: MTurkEvent, ) -> None: """ - Called either directly or from the SNS Notification that a - HIT was submitted - - :return: None + Reject any submitted assignments """ # @@ -71,90 +47,21 @@ def process_assignment_submitted( event_type="assignment_not_found_in_amt", amt_hit_type_id=event.amt_hit_type_id, ) - return None + return # Even if the assignment doesn't exist, the hit must ... hit = hm.get_from_amt_id(amt_hit_id=assignment.amt_hit_id) - # - # Step 2: Attempt to get the Assignment out of the DB - # - # Now, we need to confirm it is something that we have in the db. If not, - # that means either something broke, or some funny business is happening - # (maybe a baddie is submitting an assignment without doing any work). - stub = am.get_stub_if_exists(amt_assignment_id=assignment.amt_assignment_id) - if stub is None: - # When they visited the "work" page, it should have created an - # AssignmentStub in the db. If that doesn't exist, something bad - # happened. - logging.warning(f"No assignment found in DB: {event.amt_assignment_id}") - emit_error_event( - event_type="assignment_stub_not_found_in_db", - amt_hit_type_id=event.amt_hit_type_id, - ) - reject_assignment( - amtm=amtm, - am=am, - hm=hm, - amt_assignment_id=assignment.amt_assignment_id, - msg=REJECT_MESSAGE_UNKNOWN_ASSIGNMENT, - amt_hit_type_id=hit.amt_hit_type_id, - ) - review_hit(amtm=amtm, hm=hm, assignment=assignment) - return None - - assert assignment.amt_assignment_id == event.amt_assignment_id - assert assignment.amt_hit_id == event.amt_hit_id - assert assignment.amt_hit_id == stub.amt_hit_id - assert assignment.amt_worker_id == stub.amt_worker_id - amt_assignment_id = assignment.amt_assignment_id - amt_worker_id = assignment.amt_worker_id - - # We don't have a TSID associated with the assignment until we the - # assignment is submitted. - am.update_answer(assignment=assignment) - - # check if the user is blocked by thl - if get_user_blocked_or_not_exists(amt_worker_id=amt_worker_id): - logging.warning( - f"User {amt_worker_id} blocked or not exists. Rejecting: {amt_assignment_id}" - ) - emit_error_event( - event_type="assignment_submitted_user_blocked_or_not_exists", - amt_hit_type_id=event.amt_hit_type_id, - ) - reject_assignment( - amtm=amtm, - am=am, - hm=hm, - amt_assignment_id=amt_assignment_id, - msg=REJECT_MESSAGE_BADDIE, - amt_hit_type_id=hit.amt_hit_type_id, - ) - review_hit(amtm=amtm, hm=hm, assignment=assignment) - return None - - if assignment.tsid is None: - assignment = handle_assignment_w_no_work( - amtm=amtm, am=am, hm=hm, assignment=assignment - ) - else: - # We need to validate the work exists on thl, and if so, approve - assignment = handle_assignment_w_work( - amtm=amtm, am=am, hm=hm, assignment=assignment - ) - - # - # Step 4: Tell Amazon we've reviewed the HIT, and update the DB - # + reject_assignment( + amtm=amtm, + am=am, + hm=hm, + amt_assignment_id=assignment.amt_assignment_id, + msg=REJECT_MESSAGE_UNKNOWN_ASSIGNMENT, + amt_hit_type_id=hit.amt_hit_type_id, + ) review_hit(amtm=amtm, hm=hm, assignment=assignment) - - if ( - assignment.tsid - and assignment.status == AssignmentStatus.Approved - and assignment.requester_feedback != NO_WORK_APPROVAL_MESSAGE - ): - return issue_worker_payment(amtm=amtm, hm=hm, bm=bm, assignment=assignment) + return def review_hit(amtm: AMTManager, hm: HitManager, assignment: Assignment) -> None: @@ -166,63 +73,12 @@ def review_hit(amtm: AMTManager, hm: HitManager, assignment: Assignment) -> None logging.warning( f"Hit not found when trying to review hit: {assignment.amt_hit_id}" ) - return None + return # Update the db hm.update_hit(hit) - return None - - -def handle_assignment_w_no_work( - amtm: AMTManager, am: AssignmentManager, hm: HitManager, assignment: Assignment -) -> Assignment: - """ - Called when an assignment is submitted without a wall event. - Not entirely clear why this happens. I think they accept a HIT, get no work - available for whatever reason, then report, and submit it. - - :return: The Assignment - """ - logging.warning( - f"Assignment submitted with no tsid: {assignment.amt_assignment_id}" - ) - amt_worker_id = assignment.amt_worker_id - amt_assignment_id = assignment.amt_assignment_id - hit = hm.get_from_amt_id(amt_hit_id=assignment.amt_hit_id) - emit_error_event( - event_type="assignment_submitted_no_work", - amt_hit_type_id=hit.amt_hit_type_id, - ) - - # They get 0 chances due to abuse - if True: - # if (am.missing_tsid_count(amt_worker_id=amt_worker_id) >= 3) or ( - # am.rejected_count(amt_worker_id=amt_worker_id) >= 3 - # or get_user_blocked(amt_worker_id=amt_worker_id) - # ): - assignment = reject_assignment( - amtm=amtm, - am=am, - hm=hm, - amt_assignment_id=amt_assignment_id, - msg=REJECT_MESSAGE_NO_WORK, - amt_hit_type_id=hit.amt_hit_type_id, - ) - # todo: we don't have a way to "block" a user (i.e. tattle to thl) - # make_block_worker_decision(user) - return assignment - - # Approve with a message explaining they shouldn't do it. - assignment = approve_assignment( - amtm=amtm, - am=am, - amt_assignment_id=amt_assignment_id, - msg=NO_WORK_APPROVAL_MESSAGE, - amt_hit_type_id=hit.amt_hit_type_id, - ) - - return assignment + return def reject_assignment( @@ -269,274 +125,3 @@ def reject_assignment( ) logging.warning(f"Rejected assignment: {amt_assignment_id}") return assignment - - -def approve_assignment( - amtm: AMTManager, - am: AssignmentManager, - amt_assignment_id: str, - msg: str, - amt_hit_type_id: str, - override_rejection: bool = False, -) -> Assignment: - # Approve in AMT, update db - - res = amtm.approve_assignment_if_possible( - amt_assignment_id=amt_assignment_id, - msg=msg, - override_rejection=override_rejection, - ) - if res is None: - # We failed to approve this assignment. Cannot distinguish between - # failed b/c assignment is already approved, or it is not possible. - emit_error_event( - event_type="failed_to_approve_assignment", - amt_hit_type_id=amt_hit_type_id, - ) - # The assignment might already be approved, the error msg is useless, so - # keep going. - # raise Exception(f"Failed to approve assignment: {amt_assignment_id}") - - # We just approved this assignment, get it from amazon again - assignment = amtm.get_assignment(amt_assignment_id=amt_assignment_id) - assert assignment.status == AssignmentStatus.Approved - # And update the db - am.approve(assignment=assignment) - emit_assignment_event( - status=AssignmentStatus.Approved, amt_hit_type_id=amt_hit_type_id, reason=msg - ) - logging.warning(f"Approved assignment: {amt_assignment_id}") - return assignment - - -def handle_assignment_w_work( - amtm: AMTManager, am: AssignmentManager, hm: HitManager, assignment: Assignment -) -> Assignment: - """ - Called when an assignment is submitted with a tsid. - - Check the tsid (thl status endpoint). Make sure it is finished, and - stuff matches (doesn't matter if not a complete) - - Try to submit a cashout request for the HIT payout (e.g. 5c) - """ - - amt_worker_id = assignment.amt_worker_id - amt_assignment_id = assignment.amt_assignment_id - tsid = assignment.tsid - assert ( - tsid is not None - ), "Assignment must have a tsid to be handled in handle_assignment_w_work" - - hit = hm.get_from_amt_id(amt_hit_id=assignment.amt_hit_id) - - tsr = get_task_status(tsid=tsid) - if ( - tsr is None - or tsr.status is None - or tsr.status_code_1 - in { - StatusCode1.SESSION_START_FAIL, - StatusCode1.SESSION_START_QUALITY_FAIL, - StatusCode1.SESSION_CONTINUE_QUALITY_FAIL, - } - ): - # TSID doesn't exist or work is not finished: - # Reject the assignment instead - if tsr is not None and tsr.status_code_1 in { - StatusCode1.SESSION_START_QUALITY_FAIL, - StatusCode1.SESSION_CONTINUE_QUALITY_FAIL, - }: - event_type = "assignment_submitted_quality_fail" - else: - event_type = "assignment_submitted_work_not_complete" - - emit_error_event( - event_type=event_type, - amt_hit_type_id=hit.amt_hit_type_id, - ) - assignment = reject_assignment( - amtm=amtm, - am=am, - hm=hm, - amt_assignment_id=amt_assignment_id, - msg=REJECT_MESSAGE_BADDIE, - amt_hit_type_id=hit.amt_hit_type_id, - ) - return assignment - - assert tsr.product_user_id == amt_worker_id - assert tsr.finished is not None - assert (tsr.finished - assignment.created_at) <= timedelta(minutes=90) - - # Request an AMT_ASSIGNMENT cashout for 1c - req = submit_and_approve_amt_assignment_request( - amt_worker_id=amt_worker_id, amount=hit.reward - ) - if req is None: - # Reject the assignment instead - logging.warning( - f"submit_and_approve_amt_assignment_request failed: {amt_assignment_id}" - ) - emit_error_event( - event_type="assignment_cashout_request_failed", - amt_hit_type_id=hit.amt_hit_type_id, - ) - assignment = reject_assignment( - amtm=amtm, - am=am, - hm=hm, - amt_assignment_id=amt_assignment_id, - msg=REJECT_MESSAGE_BADDIE, - amt_hit_type_id=hit.amt_hit_type_id, - ) - return assignment - - assert req.id - - # We've approved the HIT payment, now update the db to reflect this, and approve the assignment - assignment = approve_assignment( - amtm=amtm, - am=am, - amt_assignment_id=amt_assignment_id, - msg=APPROVAL_MESSAGE, - amt_hit_type_id=hit.amt_hit_type_id, - ) - # We complete after the assignment is approved - complete_res = manage_pending_cashout( - cashout_id=req.id, payout_status=PayoutStatus.COMPLETE - ) - if complete_res.status != PayoutStatus.COMPLETE: - # unclear wny this would happen - raise ValueError(f"Failed to complete cashout: {req.id}") - return assignment - - -def submit_and_approve_amt_assignment_request( - amt_worker_id: str, amount: USDCent -) -> Optional[CashoutRequestInfo]: - # If successful, returns the cashout id, otherwise, returns None - req = user_cashout_request( - amt_worker_id=amt_worker_id, - amount=amount, - cashout_method_id=settings.amt_assignment_cashout_method, - ) - assert req.id - - if req.status != PayoutStatus.PENDING: - return None - - approve_res = manage_pending_cashout(req.id, PayoutStatus.APPROVED) - if approve_res.status != PayoutStatus.APPROVED: - return None - - return req - - -def submit_and_approve_amt_bonus_request( - amt_worker_id: str, amount: USDCent -) -> Optional[CashoutRequestInfo]: - # If successful, returns the cashout id, otherwise, returns None - req = user_cashout_request( - amt_worker_id=amt_worker_id, - amount=amount, - cashout_method_id=settings.amt_bonus_cashout_method, - ) - assert req.id - - if req.status != PayoutStatus.PENDING: - return None - - approve_res = manage_pending_cashout(req.id, PayoutStatus.APPROVED) - if approve_res.status != PayoutStatus.APPROVED: - return None - - return req - - -def issue_worker_payment( - amtm: AMTManager, hm: HitManager, bm: BonusManager, assignment: Assignment -) -> None: - # For now, since we have no "I want my bonus" request/button. A user's - # balance will be sent out anytime they get an approved assignment. We - # don't need the task status, the tsid, nor the amount / user_payout - # that was paid, or anything. - # We just get the wallet balance and submit a cashout request if >0 - # then approve it, send the amt bonus, then complete it - amt_assignment_id = assignment.amt_assignment_id - hit = hm.get_from_amt_id(amt_hit_id=assignment.amt_hit_id) - wallet_balance = get_wallet_balance_if_non_negative( - amt_worker_id=assignment.amt_worker_id - ) - if not wallet_balance: - return None - amount = round_payment(amount=wallet_balance) - if not amount: - return None - - # Don't send more than $4.97 at a time. If they have a higher wallet balance, - # they just need to get another hit approved. - amount = min(amount, USDCent(4_97)) - - pe = submit_and_approve_amt_bonus_request( - amt_worker_id=assignment.amt_worker_id, amount=amount - ) - if pe is None: - logging.warning( - f"submit_and_approve_amt_bonus_request failed: {amt_assignment_id}" - ) - emit_error_event( - event_type="bonus_cashout_request_failed", - amt_hit_type_id=hit.amt_hit_type_id, - ) - return None - assert pe.id - - amtm.send_bonus( - amt_worker_id=assignment.amt_worker_id, - amt_assignment_id=assignment.amt_assignment_id, - amount=amount, - reason=BONUS_MESSAGE, - unique_request_token=pe.id, - ) - - # Confirm it was sent through amt - bonus = amtm.get_bonus( - amt_assignment_id=assignment.amt_assignment_id, payout_event_id=pe.id - ) - - if bonus is None: - logging.warning( - f"Failed to find bonus after sending it: {amt_assignment_id} {pe.id}" - ) - emit_error_event( - event_type="bonus_not_found_after_sending", - amt_hit_type_id=hit.amt_hit_type_id, - ) - return None - - # Create in DB - bm.create(bonus=bonus) - emit_bonus_event(amount=amount, amt_hit_type_id=hit.amt_hit_type_id) - - # Complete cashout - res = manage_pending_cashout(pe.id, PayoutStatus.COMPLETE) - if res.status != PayoutStatus.COMPLETE: - raise ValueError( - f"{assignment.amt_assignment_id} {pe.id=} manage_pending_cashout COMPLETE failed: {res=}" - ) - - -def round_payment(amount: USDCent) -> USDCent: - """ - Don't pay bonuses less than 7 cents, just add it to their wallet. - Round down bonuses (>=7 cents) to the nearest multiple of 5 - starting at 2 - """ - if amount < 7: - return USDCent(0) - - amt = (5 * math.floor((int(amount) - 2) / 5)) + 2 - - payout = USDCent(amt) - assert 0 <= payout <= 40_00, "Payout must be between $0.00 and $40.00" - - return payout 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) diff --git a/jb/main.py b/jb/main.py index 2b85eb7..70c98e9 100644 --- a/jb/main.py +++ b/jb/main.py @@ -54,7 +54,6 @@ def schedule_tasks(): from jb.flow.tasks import refill_hits_task Process(target=process_mturk_events_task).start() - # Process(target=handle_pending_msgs_task).start() Process(target=refill_hits_task).start() diff --git a/jb/managers/gr_api.py b/jb/managers/gr_api.py index a31f7da..20dc051 100644 --- a/jb/managers/gr_api.py +++ b/jb/managers/gr_api.py @@ -66,6 +66,7 @@ class GRApiManager: "product_user_id": res["product_user_id"], "email": res["metadata"].get("email_address"), "display_name": res["metadata"].get("display_name"), + 'blocked': res['blocked'], } ) diff --git a/jb/managers/thl.py b/jb/managers/thl.py index f0534db..e50fe76 100644 --- a/jb/managers/thl.py +++ b/jb/managers/thl.py @@ -16,46 +16,7 @@ from generalresearch.models.thl.definitions import PayoutStatus from typing import Optional import requests -# TODO: Organize this more with other endpoints (offerwall, cashout -# requests/approvals, etc). - - -def get_user_profile(amt_worker_id: str) -> UserProfile: - url = f"{settings.fsb_host}{settings.product_id}/user/{amt_worker_id}/profile/" - res = requests.get(url).json() - - if res.get("detail") == "user not found": - raise ValueError("user not found") - - user_profile = res["user_profile"] - # todo: these are computed fields, need a UserProfile parser - user_profile.pop("email_md5", None) - user_profile.pop("email_sha1", None) - user_profile.pop("email_sha256", None) - # todo: this contains computed fields inside each streak object - user_profile.pop("streaks", None) - # todo: this shouldn't be in here anyways ---v - user_profile["user"].pop("id", None) - - return UserProfile.model_validate(user_profile) - - -def get_user_blocked(amt_worker_id: str) -> bool: - # Not blocked if None - res = get_user_profile(amt_worker_id=amt_worker_id) - return res.user.blocked if res.user.blocked is not None else False - - -def get_user_blocked_or_not_exists(amt_worker_id: str) -> Optional[bool]: - try: - res = get_user_profile(amt_worker_id=amt_worker_id) - return res.user.blocked if res.user.blocked is not None else False - - except ValueError as e: - if e.args[0] == "user not found": - return True - - return None +from jb.models.auth import User def get_task_status(tsid: str) -> Optional[TaskStatusResponse]: @@ -68,18 +29,14 @@ def get_task_status(tsid: str) -> Optional[TaskStatusResponse]: def user_cashout_request( - amt_worker_id: str, amount: USDCent, cashout_method_id: str + user: User, amount: USDCent, cashout_method_id: str ) -> CashoutRequestInfo: - assert cashout_method_id in { - settings.amt_assignment_cashout_method, - settings.amt_bonus_cashout_method, - } assert isinstance(amount, USDCent) assert USDCent(0) < amount < USDCent(10_00) url = f"{settings.fsb_host}{settings.product_id}/cashout/" body: dict[str, str | int] = { - "bpuid": amt_worker_id, + "bpuid": user.product_user_id, "amount": int(amount), "cashout_method_id": cashout_method_id, } diff --git a/jb/models/auth.py b/jb/models/auth.py index d0e8386..fef5070 100644 --- a/jb/models/auth.py +++ b/jb/models/auth.py @@ -44,6 +44,8 @@ class User(BaseModel): display_name: str | None = Field(default=None, max_length=255) + blocked: bool = Field(default=False) + @field_validator("email", mode="before") @classmethod def normalize_email(cls, value: Any) -> Any: diff --git a/jb/settings.py b/jb/settings.py index 78b4b02..7fe2a5a 100644 --- a/jb/settings.py +++ b/jb/settings.py @@ -30,9 +30,6 @@ class AmtJbBaseSettings(BaseSettings): aws_owner_id: str = Field() aws_subscription_arn: str = Field() - amt_bonus_cashout_method: str = Field() - amt_assignment_cashout_method: str = Field() - class Settings(AmtJbBaseSettings): model_config = SettingsConfigDict( diff --git a/jb/views/common.py b/jb/views/common.py index 7011557..d3b3e93 100644 --- a/jb/views/common.py +++ b/jb/views/common.py @@ -1,19 +1,18 @@ import json -from typing import Dict, Any +from typing import Annotated, Any import requests -from fastapi import Request, APIRouter, HTTPException +from fastapi import APIRouter, Depends, HTTPException, Request from fastapi.responses import HTMLResponse from starlette.responses import RedirectResponse -from jb.config import settings, JB_EVENTS_STREAM -from jb.decorators import REDIS, HM -from jb.flow.monitoring import emit_assignment_event, emit_mturk_notification_event -from jb.models.definitions import AssignmentStatus +from jb.api.auth import get_authenticated_user +from jb.config import JB_EVENTS_STREAM, settings +from jb.decorators import REDIS +from jb.flow.monitoring import emit_mturk_notification_event +from jb.models.auth import User from jb.models.event import MTurkEvent from jb.settings import BASE_HTML -from jb.config import settings -from jb.views.tasks import process_request common_router = APIRouter(prefix="", tags=["API"], include_in_schema=True) @@ -29,32 +28,16 @@ async def work(request: Request): amt_hit_id = request.query_params.get("hitId", None) print(f"work: {amt_assignment_id=} {worker_id=} {amt_hit_id=}") - if not worker_id: + if ( + not worker_id + or not amt_assignment_id + or amt_assignment_id == "ASSIGNMENT_ID_NOT_AVAILABLE" + ): return RedirectResponse( url=f"/preview/?{request.url.query}" if request.url.query else "/preview/", status_code=302, ) - if amt_assignment_id is None or amt_assignment_id == "ASSIGNMENT_ID_NOT_AVAILABLE": - # Worker is previewing the HIT - amt_hit_type_id = "unknown" - if amt_hit_id: - hit = HM.get_from_amt_id(amt_hit_id=amt_hit_id) - amt_hit_type_id = hit.amt_hit_type_id - emit_assignment_event( - status=AssignmentStatus.PreviewState, amt_hit_type_id=amt_hit_type_id - ) - return RedirectResponse( - url=f"/preview/?{request.url.query}" if request.url.query else "/preview/", - status_code=302, - ) - - try: - # The Worker has accepted the HIT - process_request(request) - except Exception: - raise HTTPException(status_code=500, detail="Error processing request") - return HTMLResponse(BASE_HTML) @@ -85,10 +68,23 @@ async def mturk_notifications(request: Request): return {"status": "ok"} -def enqueue_mturk_notifications(msg: Dict[str, Any]) -> None: +def enqueue_mturk_notifications(msg: dict[str, Any]) -> None: for evt in msg["Events"]: event = MTurkEvent.from_sns(evt) emit_mturk_notification_event( event_type=event.event_type, amt_hit_type_id=event.amt_hit_type_id ) REDIS.xadd(JB_EVENTS_STREAM, {"data": event.model_dump_json()}) + + +@common_router.get(path="/work/direct/", response_class=HTMLResponse) +async def work_direct( + request: Request, + user: Annotated[User, Depends(get_authenticated_user)], +): + """ + View for direct work (not on AMT). + Makes sure user is authenticated. + """ + # todo: emit event + return HTMLResponse(BASE_HTML) diff --git a/jb/views/tasks.py b/jb/views/tasks.py deleted file mode 100644 index 15857c3..0000000 --- a/jb/views/tasks.py +++ /dev/null @@ -1,78 +0,0 @@ -from datetime import datetime, timezone, timedelta - -from fastapi import Request - -from jb.decorators import AMTM, AM, HM -from jb.flow.maintenance import check_hit_status -from jb.flow.monitoring import emit_assignment_event -from jb.models.assignment import AssignmentStub -from jb.models.definitions import AssignmentStatus - - -def process_request(request: Request) -> None: - """ - A worker has loaded the HIT (work) page and (probably) accepted the HIT. - AMT creates an assignment, tied to this hit and this worker. - Create it in the DB. - """ - amt_assignment_id = request.query_params.get("assignmentId", None) - if amt_assignment_id == "ASSIGNMENT_ID_NOT_AVAILABLE": - raise ValueError("shouldn't happen") - - amt_hit_id = request.query_params.get("hitId", None) - amt_worker_id = request.query_params.get("workerId", None) - print(f"process_request: {amt_assignment_id=} {amt_worker_id=} {amt_hit_id=}") - assert amt_worker_id and amt_hit_id and amt_assignment_id - - # Check that the HIT is still valid - hit = HM.get_from_amt_id_if_exists(amt_hit_id=amt_hit_id) - if not hit: - raise ValueError(f"Hit {amt_hit_id} not found in DB") - - _ = check_hit_status( - amtm=AMTM, amt_hit_id=amt_hit_id, amt_hit_type_id=hit.amt_hit_type_id - ) - - emit_assignment_event( - status=AssignmentStatus.Accepted, - amt_hit_type_id=hit.amt_hit_type_id, - ) - - # I think it won't be assignable anymore? idk - # assert hit_status == HitStatus.Assignable, f"hit {amt_hit_id} {hit_status=}. Expected Assignable" - - # I would like to verify in the AMT API that this assignment is valid, but there - # is no way to do that (until the assignment is submitted) - - # # Make an offerwall to create a user account... - # # todo: GSS: Do we really need to do this??? - # client_ip = get_client_ip(request) - # url = f"{settings.fsb_host}{settings.product_id}/offerwall/45b7228a7/" - # _ = requests.get( - # url, - # {"bpuid": amt_worker_id, "ip": client_ip, "n_bins": 1, "format": "json"}, - # ).json() - - # This assignment shouldn't already exist. If it does, just make sure it - # is all the same. - assignment_stub = AM.get_stub_if_exists(amt_assignment_id=amt_assignment_id) - if assignment_stub: - print(f"{assignment_stub=}") - assert assignment_stub.amt_worker_id == amt_worker_id - assert assignment_stub.amt_assignment_id == amt_assignment_id - assert assignment_stub.created_at > ( - datetime.now(tz=timezone.utc) - timedelta(minutes=90) - ) - return None - - assignment_stub = AssignmentStub( - amt_hit_id=amt_hit_id, - amt_worker_id=amt_worker_id, - amt_assignment_id=amt_assignment_id, - status=AssignmentStatus.Accepted, - hit_id=hit.id, - ) - - AM.create_stub(stub=assignment_stub) - - return None -- cgit v1.2.3