test_t3_e2e.py
python
sha256:fc4c9ad652d1fff3dc508cb6ea02ee710ee6dfc4cb3761291d9900b5e029ea8a
feat(slice-7): T3 dataset review lifecycle, job queue, prov…
Human
minor
⚠ breaking
41 days ago
| 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"], "pending_review") |
| 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 | {"datasetId": "e2e-rejected-v1"}, |
| 97 | ) |
| 98 | self._json(f"{self._base}/datasets/e2e-rejected-v1/submit", "POST") |
| 99 | rejected = self._json( |
| 100 | f"{self._base}/datasets/e2e-rejected-v1/review", |
| 101 | "POST", |
| 102 | {"action": "reject", "reasonCode": "DUPLICATE_SUBMISSION"}, |
| 103 | ) |
| 104 | self.assertEqual(rejected["status"], "rejected") |
| 105 | self.assertEqual(rejected["rejectionReasonCode"], "DUPLICATE_SUBMISSION") |
| 106 | # Confirm no free text in the response body. |
| 107 | self.assertNotIn("caller message", json.dumps(rejected)) |
| 108 | |
| 109 | def test_e2e_t3_queue_state_endpoint_returns_counts(self) -> None: |
| 110 | """GET /training/queue returns a JSON object with queue metrics.""" |
| 111 | |
| 112 | state = self._json(f"{self._base}/training/queue", "GET") |
| 113 | self.assertIn("queuedCount", state) |
| 114 | self.assertIn("runningCount", state) |
| 115 | self.assertIn("activeCount", state) |
| 116 | self.assertIn("maxConcurrentRunning", state) |
| 117 | |
| 118 | def test_e2e_t3_expiry_tombstone_provenance_readable_over_http(self) -> None: |
| 119 | """After TTL expiry the provenance endpoint still returns 200.""" |
| 120 | |
| 121 | policy = {"policyClass": "ephemeral", "ttlSeconds": 60} |
| 122 | created = self._json( |
| 123 | f"{self._base}/training/jobs", |
| 124 | "POST", |
| 125 | valid_payload("e2e-expiry-prov", policy), |
| 126 | ) |
| 127 | job_id = str(created["id"]) |
| 128 | prov_before = self._json( |
| 129 | f"{self._base}/training/jobs/{job_id}/provenance", "GET" |
| 130 | ) |
| 131 | |
| 132 | # Trigger expiry via the service (direct call, not via HTTP). |
| 133 | self._service.sweep_expired_artifacts(datetime.now(UTC) + timedelta(seconds=120)) |
| 134 | |
| 135 | tombstone = self._json(f"{self._base}/training/jobs/{job_id}", "GET") |
| 136 | self.assertEqual(tombstone["status"], "deleted") |
| 137 | |
| 138 | prov_after = self._json( |
| 139 | f"{self._base}/training/jobs/{job_id}/provenance", "GET" |
| 140 | ) |
| 141 | self.assertEqual(prov_before["jobId"], prov_after["jobId"]) |
| 142 | self.assertEqual( |
| 143 | prov_before["artifactHash"], prov_after["artifactHash"] |
| 144 | ) |
| 145 | |
| 146 | def test_e2e_t3_explicit_delete_wipes_provenance_over_http(self) -> None: |
| 147 | """After explicit DELETE the provenance endpoint returns 404.""" |
| 148 | |
| 149 | created = self._json( |
| 150 | f"{self._base}/training/jobs", |
| 151 | "POST", |
| 152 | valid_payload("e2e-explicit-delete"), |
| 153 | ) |
| 154 | job_id = str(created["id"]) |
| 155 | arts = self._json( |
| 156 | f"{self._base}/training/jobs/{job_id}/artifacts", "GET" |
| 157 | ) |
| 158 | artifact_id = str(arts["artifacts"][0]["id"]) |
| 159 | |
| 160 | self._json( |
| 161 | f"{self._base}/training/jobs/{job_id}/artifacts/{artifact_id}", |
| 162 | "DELETE", |
| 163 | ) |
| 164 | with self.assertRaises(HTTPError) as raised: |
| 165 | self._json( |
| 166 | f"{self._base}/training/jobs/{job_id}/provenance", "GET" |
| 167 | ) |
| 168 | self.assertEqual(raised.exception.code, 404) |
| 169 | raised.exception.close() |
| 170 | |
| 171 | def test_e2e_t3_unapproved_dataset_returns_403_over_http(self) -> None: |
| 172 | """Job submission against an unapproved dataset returns HTTP 403.""" |
| 173 | |
| 174 | ds_store = DatasetStore() |
| 175 | ds_store.register("unapproved-e2e-ds") |
| 176 | service = TrainingApiService( |
| 177 | TrainingJobStore(), dataset_store=ds_store |
| 178 | ) |
| 179 | server = ThreadingHTTPServer(("127.0.0.1", 0), make_handler(service)) |
| 180 | thread = threading.Thread(target=server.serve_forever, daemon=True) |
| 181 | thread.start() |
| 182 | base = f"http://127.0.0.1:{server.server_port}" |
| 183 | try: |
| 184 | with self.assertRaises(HTTPError) as raised: |
| 185 | self._json( |
| 186 | f"{base}/training/jobs", |
| 187 | "POST", |
| 188 | { |
| 189 | "idempotencyKey": "e2e-403-test", |
| 190 | "datasetId": "unapproved-e2e-ds", |
| 191 | "modelId": "fixture-tiny-llm", |
| 192 | "requestedBy": "e2e-test", |
| 193 | }, |
| 194 | ) |
| 195 | self.assertEqual(raised.exception.code, 403) |
| 196 | raised.exception.close() |
| 197 | finally: |
| 198 | server.shutdown() |
| 199 | server.server_close() |
| 200 | thread.join(timeout=2) |
| 201 | |
| 202 | |
| 203 | if __name__ == "__main__": |
| 204 | unittest.main() |
File History
1 commit
sha256:fc4c9ad652d1fff3dc508cb6ea02ee710ee6dfc4cb3761291d9900b5e029ea8a
feat(slice-7): T3 dataset review lifecycle, job queue, prov…
Human
minor
⚠
41 days ago