aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--jb/config.py1
-rw-r--r--jb/flow/assignment_tasks.py451
-rw-r--r--jb/flow/events.py64
-rw-r--r--jb/main.py1
-rw-r--r--jb/managers/gr_api.py1
-rw-r--r--jb/managers/thl.py49
-rw-r--r--jb/models/auth.py2
-rw-r--r--jb/settings.py3
-rw-r--r--jb/views/common.py56
-rw-r--r--jb/views/tasks.py78
10 files changed, 55 insertions, 651 deletions
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