Coverage for server / services / teacher / question_import.py: 96%
78 statements
« prev ^ index » next coverage.py v7.13.4, created at 2026-10-04 09:33 +0000
« prev ^ index » next coverage.py v7.13.4, created at 2026-10-04 09:33 +0000
1"""Client for the Question Ingest service.
3The Teacher-Student frontend must never call the ingest service directly — it talks only
4to this API, which proxies. (The existing bulkUpload.users.* endpoints in the client call
5a bulk service straight from the browser; that is the pattern this deliberately avoids.)
7Retries are added here because no existing httpx call site in this repo has any: a single
8transient connection error would otherwise surface to a teacher as a failed import of a
9file that was perfectly good.
11Developer: Allan Ninal
12"""
14import asyncio
15import logging
16import mimetypes
17import os
18from typing import Any
20import httpx
22logger = logging.getLogger(__name__)
24# Never hardcode the host — this differs per environment and the service may move.
25INGEST_BASE_URL = os.getenv("QUESTION_INGEST_URL", "http://127.0.0.1:8100")
27# Upload can carry a large file; status polls should be quick.
28UPLOAD_TIMEOUT = float(os.getenv("QUESTION_INGEST_UPLOAD_TIMEOUT", "120"))
29STATUS_TIMEOUT = float(os.getenv("QUESTION_INGEST_STATUS_TIMEOUT", "15"))
31# Only idempotent, transport-level failures are retried. A 4xx is the teacher's file
32# being wrong and must reach them unchanged rather than being tried three times.
33# Commit had been borrowing STATUS_TIMEOUT, which is the budget for a 2-second poll, not
34# for sending a whole reviewed batch back. With the retry loop that meant a worst case of
35# 15 + 0.5 + 15 + 1.0 + 15 = 46.5s inside this call alone — against nginx's 60s
36# proxy_read_timeout, and before this API had written a single question. The browser then
37# got a 504 while the API carried on writing, so nobody could tell how many of 22
38# questions had been created.
39#
40# 30s is far beyond what the work costs (the ingest commit is pure CPU — it validates and
41# maps, it does no database I/O — and this API's own writes go to a database milliseconds
42# away), so it fires only when something is genuinely wrong rather than merely slow.
43COMMIT_TIMEOUT = float(os.getenv("QUESTION_INGEST_COMMIT_TIMEOUT", "30"))
45RETRY_ATTEMPTS = 3
46RETRY_BACKOFF_SECONDS = 0.5
49class IngestUnavailable(RuntimeError):
50 """The ingest service could not be reached. Distinct from it rejecting the file."""
53class IngestRejected(RuntimeError):
54 """The ingest service rejected the request. Carries its status and body through."""
56 def __init__(self, status_code: int, detail: Any) -> None:
57 super().__init__(f"ingest returned {status_code}")
58 self.status_code = status_code
59 self.detail = detail
62async def _request(
63 method: str, path: str, *, timeout: float, retry_on_timeout: bool = True, **kwargs
64) -> Any:
65 url = f"{INGEST_BASE_URL.rstrip('/')}{path}"
66 last_error: Exception | None = None
68 for attempt in range(1, RETRY_ATTEMPTS + 1):
69 try:
70 async with httpx.AsyncClient(timeout=timeout) as client:
71 response = await client.request(method, url, **kwargs)
72 except (httpx.ConnectError, httpx.ReadTimeout, httpx.WriteTimeout) as exc:
73 # Transport failure — worth retrying, the request may never have arrived.
74 #
75 # Safe is not the same as worth doing. A commit passes retry_on_timeout=False:
76 # a read timeout there means the service is already struggling, and spending
77 # the caller's remaining budget on two more long waits only turns a slow
78 # answer into no answer at all. A CONNECT error is different — it fails
79 # immediately and usually means a restart mid-request — so that is still
80 # retried. Committing the same job twice is harmless either way: the ingest
81 # service makes it idempotent and both banks skip questions already stored.
82 timed_out = isinstance(exc, (httpx.ReadTimeout, httpx.WriteTimeout))
83 last_error = exc
84 if timed_out and not retry_on_timeout:
85 break
86 if attempt < RETRY_ATTEMPTS:
87 await asyncio.sleep(RETRY_BACKOFF_SECONDS * attempt)
88 continue
89 break
90 except httpx.HTTPError as exc: # anything else is not retryable
91 raise IngestUnavailable(f"{type(exc).__name__}: {exc}") from exc
93 if response.status_code >= 500:
94 # The service is unwell; a retry may land on a healthy worker.
95 last_error = IngestUnavailable(f"ingest returned {response.status_code}")
96 if attempt < RETRY_ATTEMPTS:
97 await asyncio.sleep(RETRY_BACKOFF_SECONDS * attempt)
98 continue
99 break
101 if response.status_code >= 400:
102 # Pass the rejection through untouched: it explains which row failed and why,
103 # and re-trying it would produce the same answer three times.
104 raise IngestRejected(response.status_code, _safe_json(response))
106 return _safe_json(response)
108 raise IngestUnavailable(str(last_error) if last_error else "ingest unreachable")
111# The ingest service scopes a job to whoever created it, and reads that id from this
112# header. Sending it is what stops two teachers who import the same district-issued CSV
113# from sharing one job — where whichever committed first blocked the other.
114OWNER_HEADER = "X-Import-Owner"
117def _owner_headers(owner: str | None) -> dict[str, str]:
118 """The caller-id header, or nothing when we have no id to send.
120 Omitted rather than sent empty: the ingest service treats a missing header as an
121 unowned job and a malformed one as a 422, so a blank value would turn a request it
122 would happily accept into a rejection.
123 """
124 return {OWNER_HEADER: str(owner)} if owner else {}
127def _safe_json(response: httpx.Response) -> Any:
128 try:
129 return response.json()
130 except ValueError:
131 return {"detail": response.text[:500]}
134# What a browser can send about a file, beyond the file. A STAAR released-item PDF names
135# the TEKS, the points and the answer for every item, but never the grade or the
136# difficulty — and neither is guessable, so they are asked for at upload and applied to
137# the batch. Without them every imported question arrives unsaveable.
138_BATCH_DEFAULTS = ("grade_level", "difficulty")
141def _media_type(filename: str) -> str:
142 """The content type to send the ingest service for this file.
144 It used to be hardcoded to "text/csv" for every upload, including .docx and .pdf.
145 The ingest service reads the filename suffix rather than the content type, so it
146 happened to work — but it is a lie on the wire, and the kind that turns into a
147 silent failure the day anything between here and there starts believing it.
148 """
149 return mimetypes.guess_type(filename)[0] or "application/octet-stream"
152async def create_import(
153 filename: str,
154 content: bytes,
155 question_type: str,
156 owner: str | None = None,
157 *,
158 grade_level: int | None = None,
159 difficulty: str | None = None,
160) -> dict:
161 data = {
162 "question_type": question_type,
163 # Tells the ingest service WHICH PORTAL this came from. The two populations
164 # are not interchangeable: teacher access is gated on the district plan
165 # (which defaults to free), staff access has no feature gate at all. Without
166 # this, a scan rate measured almost entirely from one is read as a decision
167 # about the other.
168 "source": "teacher",
169 }
170 # Built conditionally, never as `{"grade_level": grade_level}` with a None: httpx
171 # serialises None in a form field as the literal string "None", which the ingest
172 # service then rejects as an invalid integer — turning "the teacher did not pick a
173 # grade" into a 422 on the upload itself.
174 for name, value in zip(_BATCH_DEFAULTS, (grade_level, difficulty), strict=True):
175 if value is not None and str(value).strip() != "":
176 data[name] = str(value)
178 return await _request(
179 "POST",
180 "/imports",
181 timeout=UPLOAD_TIMEOUT,
182 files={"file": (filename, content, _media_type(filename))},
183 data=data,
184 headers=_owner_headers(owner),
185 )
188async def create_import_from_text(
189 text: str, question_type: str, owner: str | None = None
190) -> dict:
191 """Phase 5: a single pasted question, no file.
193 Uses the upload timeout rather than the status one: the work is the same parse, and
194 a slow ingest service should not turn a valid paste into a spurious failure.
195 """
196 return await _request(
197 "POST",
198 "/imports/text",
199 timeout=UPLOAD_TIMEOUT,
200 json={"text": text, "question_type": question_type, "source": "teacher"},
201 headers=_owner_headers(owner),
202 )
205async def get_import(job_id: str, owner: str | None = None) -> dict:
206 return await _request(
207 "GET", f"/imports/{job_id}", timeout=STATUS_TIMEOUT, headers=_owner_headers(owner)
208 )
211async def get_import_questions(job_id: str, owner: str | None = None) -> dict:
212 return await _request(
213 "GET",
214 f"/imports/{job_id}/questions",
215 timeout=STATUS_TIMEOUT,
216 headers=_owner_headers(owner),
217 )
220async def commit_import(job_id: str, payload: dict, owner: str | None = None) -> dict:
221 """Send the reviewed batch back and get the create payloads.
223 Its own timeout, and no retry on a slow answer: see COMMIT_TIMEOUT. The browser is
224 waiting on this call, and everything after it in this request — the duplicate check
225 and one write per question — still has to happen inside nginx's 60s ceiling.
226 """
227 return await _request(
228 "POST",
229 f"/imports/{job_id}/commit",
230 timeout=COMMIT_TIMEOUT,
231 retry_on_timeout=False,
232 json=payload,
233 headers=_owner_headers(owner),
234 )
237async def ingest_ready() -> bool:
238 """True when the ingest service reports itself ready. Used by health surfaces."""
239 try:
240 await _request("GET", "/ready", timeout=STATUS_TIMEOUT)
241 return True
242 except (IngestUnavailable, IngestRejected):
243 return False