aboutsummaryrefslogtreecommitdiff
path: root/jb/views
diff options
context:
space:
mode:
authorMax Nanis2026-09-13 19:03:55 +0000
committerMax Nanis2026-09-13 19:03:55 +0000
commit8fbe8d439b418796932aaa88297e006c714e4d13 (patch)
tree6c800476edc1e771fc570b559d483df8d59d4324 /jb/views
parent12f6fee851b68e86af658dfa17e4a0daed457dd1 (diff)
parent2c94f248d2438071a918fa9a30bf114ef9aa29b4 (diff)
downloadamt-jb-8fbe8d439b418796932aaa88297e006c714e4d13.tar.gz
amt-jb-8fbe8d439b418796932aaa88297e006c714e4d13.zip
Merges pull request #3
Off of Amazon!!!
Diffstat (limited to 'jb/views')
-rw-r--r--jb/views/auth.py203
-rw-r--r--jb/views/common.py68
-rw-r--r--jb/views/tasks.py78
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