diff options
Diffstat (limited to 'jb/views')
| -rw-r--r-- | jb/views/auth.py | 203 | ||||
| -rw-r--r-- | jb/views/common.py | 68 | ||||
| -rw-r--r-- | jb/views/tasks.py | 78 |
3 files changed, 239 insertions, 110 deletions
diff --git a/jb/views/auth.py b/jb/views/auth.py new file mode 100644 index 0000000..ce9c1e0 --- /dev/null +++ b/jb/views/auth.py @@ -0,0 +1,203 @@ +from typing import Annotated +from urllib.parse import urlencode + +from fastapi import APIRouter, Depends, HTTPException, Response, status +from fastapi.responses import HTMLResponse, RedirectResponse + +from jb.api.auth import ( + SESSION_COOKIE_NAME, + create_session, + get_authenticated_user, +) +from jb.api.magic_token import ( + consume_amt_account_link_token, + consume_magic_token, + create_amt_account_link_token, + create_magic_token, +) +from jb.config import settings +from jb.decorators import LOG +from jb.dependencies import get_gr_api_manager +from jb.managers.email_manager import ( + get_or_create_contact, + send_amt_link_email, + send_login_email, +) +from jb.managers.gr_api import GRApiManager +from jb.models.auth import ( + AccountLogin, + AmtAccountLink, + MagicLinkExchangeRequest, + User, +) +from jb.settings import BASE_HTML + +auth_router = APIRouter(prefix="/auth", tags=["Auth"]) + + +@auth_router.post("/magic-link/request") +def request_magic_link(body: AccountLogin) -> dict[str, str]: + """Create a magic link.""" + email = str(body.email) + token = create_magic_token(user_email=email) + + if settings.debug: + query = urlencode({"token": token}) + return {"magic_link": f"{settings.base_url}auth/magic-link/?{query}"} + + send_login_email(email=email, magic_token=token) + return {"detail": "Link sent. Check your inbox and follow the link to log in."} + + +@auth_router.get("/magic-link/", response_class=HTMLResponse, include_in_schema=False) +def magic_link_landing_page( + gr_api: Annotated[GRApiManager, Depends(get_gr_api_manager)], + token: str | None = None, +) -> Response: + """Serve the SPA without redeeming the token; email prefetches are harmless.""" + if settings.debug: + if token is None: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail="token is required", + ) + response = RedirectResponse(url="/", status_code=status.HTTP_303_SEE_OTHER) + _exchange_magic_link(token, response, gr_api) + return response + return HTMLResponse( + BASE_HTML, + headers={ + "Cache-Control": "no-store", + "Referrer-Policy": "no-referrer", + "X-Robots-Tag": "noindex, nofollow", + }, + ) + + +@auth_router.post("/magic-link/exchange", status_code=status.HTTP_204_NO_CONTENT) +def exchange_magic_link( + body: MagicLinkExchangeRequest, + response: Response, + gr_api: Annotated[GRApiManager, Depends(get_gr_api_manager)], +) -> None: + """Exchange a magic link only after its landing page makes an explicit POST.""" + _exchange_magic_link(body.token, response, gr_api) + + +def _exchange_magic_link(token: str, response: Response, gr_api: GRApiManager) -> None: + user_email = consume_magic_token(token) + user = gr_api.ensure_user_exists(User.model_validate({"email": user_email})) + response.set_cookie( + key=SESSION_COOKIE_NAME, + value=create_session(user.product_user_id), + max_age=settings.session_token_ttl_seconds, + httponly=True, + secure=not settings.debug, + samesite="lax", + path="/", + ) + + +@auth_router.post("/link-amt/request") +def link_amt_account(body: AmtAccountLink) -> dict[str, str]: + """Link an AMT account and login.""" + email = str(body.email) + amt_worker_id = body.amt_worker_id + + # TODO! Prevent a user that's already transitioned their account, from being + # TODO! able to continuously create this special link token. + # TODO! Max Notes: this seems to be handled within the gr-api, and that + # TODO! can raise, the following line would / should fail if needed. + + token = create_amt_account_link_token(email=email, amt_worker_id=amt_worker_id) + + if settings.debug: + query = urlencode({"token": token}) + return {"magic_link": f"{settings.base_url}auth/link-amt/?{query}"} + + send_amt_link_email(email=email, magic_token=token) + return {} + + +@auth_router.get("/debug/", response_class=HTMLResponse, include_in_schema=False) +def link_amt_account_landing_page( + gr_api: Annotated[GRApiManager, Depends(get_gr_api_manager)], + token: str | None = None, +) -> HTMLResponse: + """Serve the account-link SPA without consuming the one-time token.""" + + # TODO! Try catch any of this, and if it fails, show the user a + # TODO! failed HTML page. As of now, it shows them a failed JSON response. + + if settings.debug: + if token is None: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail="token is required", + ) + + _response = RedirectResponse(url="/", status_code=status.HTTP_303_SEE_OTHER) + _exchange_amt_account_link(token=token, response=_response, gr_api=gr_api) + + return HTMLResponse( + BASE_HTML, + headers={ + "Cache-Control": "no-store", + "Referrer-Policy": "no-referrer", + "X-Robots-Tag": "noindex, nofollow", + }, + ) + + +@auth_router.post("/link-amt/exchange", status_code=status.HTTP_204_NO_CONTENT) +def exchange_amt_account_link( + body: MagicLinkExchangeRequest, + response: Response, + gr_api: Annotated[GRApiManager, Depends(get_gr_api_manager)], +) -> None: + """Validate the email link, then transition the bound AMT account.""" + try: + _exchange_amt_account_link(body.token, response, gr_api) + except ValueError as e: + LOG.error(f"Failed to exchange AMT account link: {e}") + raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(e)) + + +def _exchange_amt_account_link(token: str, response: Response, gr_api: GRApiManager): + token_data = consume_amt_account_link_token(token) + email = token_data.email + amt_worker_id = token_data.amt_worker_id + + user = User(email=email) + user = gr_api.transition_user_from_amt(user=user, amt_worker_id=amt_worker_id) + + # In Mautic, associate the email with the worker ID (AFTER the user has transitioned) + get_or_create_contact(email=email, amt_worker_id=amt_worker_id) + + session_token = create_session(user.product_user_id) + response.set_cookie( + key=SESSION_COOKIE_NAME, + value=session_token, + max_age=settings.session_token_ttl_seconds, + httponly=True, + secure=not settings.debug, + samesite="lax", + path="/", + ) + + +@auth_router.get("/session", response_model=User) +def get_session( + user: Annotated[User, Depends(get_authenticated_user)], +) -> User: + return user + + +@auth_router.delete("/session", status_code=status.HTTP_204_NO_CONTENT) +def delete_session( + response: Response, +) -> None: + # Logout is idempotent so clients can always discard their local token. + # JWT sessions are stateless, so logout discards the browser cookie. A token + # copied elsewhere remains valid until its short, configured expiration. + response.delete_cookie(key=SESSION_COOKIE_NAME, path="/") diff --git a/jb/views/common.py b/jb/views/common.py index 7011557..9eee453 100644 --- a/jb/views/common.py +++ b/jb/views/common.py @@ -1,19 +1,19 @@ 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 generalresearch.redis_helper import RedisConfig 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 get_redis_config +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,37 +29,24 @@ 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, - ) + return HTMLResponse(BASE_HTML) - 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) +RedisConfigDep = Annotated[RedisConfig, Depends(get_redis_config)] @common_router.post(path=f"/{settings.sns_path}/", include_in_schema=False) -async def mturk_notifications(request: Request): +async def mturk_notifications(request: Request, redis_config: RedisConfigDep): """ Our SNS topic will POST to this endpoint whenever we get a new message """ @@ -77,7 +64,7 @@ async def mturk_notifications(request: Request): case "Notification": msg = json.loads(message["Message"]) print("Received MTurk event:", msg) - enqueue_mturk_notifications(msg) + enqueue_mturk_notifications(msg=msg, redis_config=redis_config) case _: raise HTTPException(status_code=500, detail="Invalid JSON") @@ -85,10 +72,27 @@ 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], redis_config: RedisConfig) -> None: + redis_client = redis_config.create_redis_client() + 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()}) + + print("enqueue_mturk_notifications", event.model_dump_json()) + redis_client.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 |
