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

1"""Client for the Question Ingest service. 

2 

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.) 

6 

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. 

10 

11Developer: Allan Ninal 

12""" 

13 

14import asyncio 

15import logging 

16import mimetypes 

17import os 

18from typing import Any 

19 

20import httpx 

21 

22logger = logging.getLogger(__name__) 

23 

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") 

26 

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")) 

30 

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")) 

44 

45RETRY_ATTEMPTS = 3 

46RETRY_BACKOFF_SECONDS = 0.5 

47 

48 

49class IngestUnavailable(RuntimeError): 

50 """The ingest service could not be reached. Distinct from it rejecting the file.""" 

51 

52 

53class IngestRejected(RuntimeError): 

54 """The ingest service rejected the request. Carries its status and body through.""" 

55 

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 

60 

61 

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 

67 

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 

92 

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 

100 

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)) 

105 

106 return _safe_json(response) 

107 

108 raise IngestUnavailable(str(last_error) if last_error else "ingest unreachable") 

109 

110 

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" 

115 

116 

117def _owner_headers(owner: str | None) -> dict[str, str]: 

118 """The caller-id header, or nothing when we have no id to send. 

119 

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 {} 

125 

126 

127def _safe_json(response: httpx.Response) -> Any: 

128 try: 

129 return response.json() 

130 except ValueError: 

131 return {"detail": response.text[:500]} 

132 

133 

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") 

139 

140 

141def _media_type(filename: str) -> str: 

142 """The content type to send the ingest service for this file. 

143 

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" 

150 

151 

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) 

177 

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 ) 

186 

187 

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. 

192 

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 ) 

203 

204 

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 ) 

209 

210 

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 ) 

218 

219 

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. 

222 

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 ) 

235 

236 

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