test_t3_e2e.py python
204 lines 7.4 KB
Raw
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