test_t3_e2e.py file-level

at sha256:f · View file ↗ · Intel ↗

History
1 files
1 commits
0 hotspots
0 🧊 dead
0 💥 blast risk
sha256:c Add T4 job cancellation retry and validation · · Jun 11, 2026
1 """End-to-end tier tests — T3 full lifecycle including rejection path and expiry tombstone."""
2
3 from __future__ import annotations
4
5 import json
6 import threading
7 import unittest
8 from datetime import UTC, datetime, timedelta
9 from http.server import ThreadingHTTPServer
10 from urllib.error import HTTPError
11 from urllib.request import Request, urlopen
12
13 from scooling_lab_helpers import valid_payload
14
15 from scooling_lab.api import make_handler
16 from scooling_lab.dataset_review import DatasetStore, RejectionReasonCode
17 from scooling_lab.service import TrainingApiService
18 from scooling_lab.store import TrainingJobStore
19
20
21 class T3EndToEndTests(unittest.TestCase):
22 """E2E HTTP tests for the T3 dataset review and retention lifecycle."""
23
24 def setUp(self) -> None:
25 """Start a fresh server for each test."""
26
27 self._service = TrainingApiService(TrainingJobStore())
28 self._server = ThreadingHTTPServer(
29 ("127.0.0.1", 0), make_handler(self._service)
30 )
31 thread = threading.Thread(target=self._server.serve_forever, daemon=True)
32 thread.start()
33 self._thread = thread
34 self._base = f"http://127.0.0.1:{self._server.server_port}"
35
36 def tearDown(self) -> None:
37 """Shut down the server."""
38
39 self._server.shutdown()
40 self._server.server_close()
41 self._thread.join(timeout=2)
42
43 # ----------------------------------------------------------------- helpers
44
45 def _json(
46 self, url: str, method: str, payload: dict[str, object] | None = None
47 ) -> dict[str, object]:
48 body = None
49 headers = {"Content-Type": "application/json"}
50 if payload is not None:
51 body = json.dumps(payload).encode("utf-8")
52 req = Request(url, data=body, headers=headers, method=method)
53 with urlopen(req, timeout=5) as resp:
54 decoded = json.loads(resp.read().decode("utf-8"))
55 if not isinstance(decoded, dict):
56 raise AssertionError("expected JSON object")
57 return decoded
58
59 # ------------------------------------------------------------------ tests
60
61 def test_e2e_t3_register_approve_submit_job_over_http(self) -> None:
62 """Full dataset review flow works end-to-end via HTTP."""
63
64 # Register a new dataset.
65 reg = self._json(
66 f"{self._base}/datasets",
67 "POST",
68 {"datasetId": "e2e-dataset-v1"},
69 )
70 self.assertEqual(reg["status"], "registered")
71
72 # Submit for review.
73 submitted = self._json(
74 f"{self._base}/datasets/e2e-dataset-v1/submit", "POST"
75 )
76 self.assertEqual(submitted["status"], "approved")
77
78 # Approve it.
79 approved = self._json(
80 f"{self._base}/datasets/e2e-dataset-v1/review",
81 "POST",
82 {"action": "approve"},
83 )
84 self.assertEqual(approved["status"], "approved")
85
86 # Read back.
87 fetched = self._json(f"{self._base}/datasets/e2e-dataset-v1", "GET")
88 self.assertEqual(fetched["status"], "approved")
89
90 def test_e2e_t3_rejection_path_over_http(self) -> None:
91 """Dataset rejection carries enum reason code; no free text reflected."""
92
93 self._json(
94 f"{self._base}/datasets",
95 "POST",
96 {
97 "datasetId": "e2e-rejected-v1",
98 "rowCount": 0,
99 "declaredSchema": {
100 "exampleId": "string",
101 "inputTokenCount": "integer",
102 "outputTokenCount": "integer",
103 "split": "string",
104 },
105 },
106 )
107 self._json(f"{self._base}/datasets/e2e-rejected-v1/submit", "POST")
108 rejected = self._json(
109 f"{self._base}/datasets/e2e-rejected-v1/review",
110 "POST",
111 {"action": "reject", "reasonCode": "DUPLICATE_SUBMISSION"},
112 )
113 self.assertEqual(rejected["status"], "rejected")
114 self.assertEqual(rejected["rejectionReasonCode"], "SYNTHETIC_LIMIT")
115 # Confirm no free text in the response body.
116 self.assertNotIn("caller message", json.dumps(rejected))
117
118 def test_e2e_t3_queue_state_endpoint_returns_counts(self) -> None:
119 """GET /training/queue returns a JSON object with queue metrics."""
120
121 state = self._json(f"{self._base}/training/queue", "GET")
122 self.assertIn("queuedCount", state)
123 self.assertIn("runningCount", state)
124 self.assertIn("activeCount", state)
125 self.assertIn("maxConcurrentRunning", state)
126
127 def test_e2e_t3_expiry_tombstone_provenance_readable_over_http(self) -> None:
128 """After TTL expiry the provenance endpoint still returns 200."""
129
130 policy = {"policyClass": "ephemeral", "ttlSeconds": 60}
131 created = self._json(
132 f"{self._base}/training/jobs",
133 "POST",
134 valid_payload("e2e-expiry-prov", policy),
135 )
136 job_id = str(created["id"])
137 prov_before = self._json(
138 f"{self._base}/training/jobs/{job_id}/provenance", "GET"
139 )
140
141 # Trigger expiry via the service (direct call, not via HTTP).
142 self._service.sweep_expired_artifacts(datetime.now(UTC) + timedelta(seconds=120))
143
144 tombstone = self._json(f"{self._base}/training/jobs/{job_id}", "GET")
145 self.assertEqual(tombstone["status"], "deleted")
146
147 prov_after = self._json(
148 f"{self._base}/training/jobs/{job_id}/provenance", "GET"
149 )
150 self.assertEqual(prov_before["jobId"], prov_after["jobId"])
151 self.assertEqual(
152 prov_before["artifactHash"], prov_after["artifactHash"]
153 )
154
155 def test_e2e_t3_explicit_delete_wipes_provenance_over_http(self) -> None:
156 """After explicit DELETE the provenance endpoint returns 404."""
157
158 created = self._json(
159 f"{self._base}/training/jobs",
160 "POST",
161 valid_payload("e2e-explicit-delete"),
162 )
163 job_id = str(created["id"])
164 arts = self._json(
165 f"{self._base}/training/jobs/{job_id}/artifacts", "GET"
166 )
167 artifact_id = str(arts["artifacts"][0]["id"])
168
169 self._json(
170 f"{self._base}/training/jobs/{job_id}/artifacts/{artifact_id}",
171 "DELETE",
172 )
173 with self.assertRaises(HTTPError) as raised:
174 self._json(
175 f"{self._base}/training/jobs/{job_id}/provenance", "GET"
176 )
177 self.assertEqual(raised.exception.code, 404)
178 raised.exception.close()
179
180 def test_e2e_t3_unapproved_dataset_returns_403_over_http(self) -> None:
181 """Job submission against an unapproved dataset returns HTTP 403."""
182
183 ds_store = DatasetStore()
184 ds_store.register("unapproved-e2e-ds")
185 service = TrainingApiService(
186 TrainingJobStore(), dataset_store=ds_store
187 )
188 server = ThreadingHTTPServer(("127.0.0.1", 0), make_handler(service))
189 thread = threading.Thread(target=server.serve_forever, daemon=True)
190 thread.start()
191 base = f"http://127.0.0.1:{server.server_port}"
192 try:
193 with self.assertRaises(HTTPError) as raised:
194 self._json(
195 f"{base}/training/jobs",
196 "POST",
197 {
198 "idempotencyKey": "e2e-403-test",
199 "datasetId": "unapproved-e2e-ds",
200 "modelId": "fixture-tiny-llm",
201 "requestedBy": "e2e-test",
202 },
203 )
204 self.assertEqual(raised.exception.code, 403)
205 raised.exception.close()
206 finally:
207 server.shutdown()
208 server.server_close()
209 thread.join(timeout=2)
210
211
212 if __name__ == "__main__":
213 unittest.main()