diff --git a/xiaobai-datahub/datahub/admin_api.py b/xiaobai-datahub/datahub/admin_api.py index fb284f5..41f162b 100644 --- a/xiaobai-datahub/datahub/admin_api.py +++ b/xiaobai-datahub/datahub/admin_api.py @@ -145,10 +145,18 @@ class AdminAPI: result = self.pipeline.ingest_reference(day) elif dataset in OFFICIAL_DATASETS or dataset == STOCKS_DATASET: # Manual same-day republish must rebuild the full A/B boundary. - result = self.pipeline.force_republish_boundary(dataset, day) - failures = self.pipeline.eod_failures(result) - if failures: - raise ApiError("FAILED_PRECONDITION", "; ".join(failures)) + # Gate failures and mid-switch exceptions both surface as + # FAILED_PRECONDITION so the admin API never leaks raw + # transaction errors to the client. + try: + result = self.pipeline.force_republish_boundary(dataset, day) + failures = self.pipeline.eod_failures(result) + if failures: + raise ApiError("FAILED_PRECONDITION", "; ".join(failures)) + except ApiError: + raise + except Exception as exc: + raise ApiError("FAILED_PRECONDITION", str(exc)) from exc else: raise ApiError("INVALID_ARGUMENT", f"unsupported backfill dataset: {dataset}") self.pipeline.audit(actor, "backfill", f"{dataset}:{day}", json.dumps({"ok": True})) diff --git a/xiaobai-datahub/tests/test_atomic_release.py b/xiaobai-datahub/tests/test_atomic_release.py index 49a1ffc..8bccf7d 100644 --- a/xiaobai-datahub/tests/test_atomic_release.py +++ b/xiaobai-datahub/tests/test_atomic_release.py @@ -367,6 +367,29 @@ class ForceBoundaryEntryTests(unittest.TestCase): with self.assertRaises(ApiError): admin.backfill("daily", TRADE_DATE, "wrong", f"daily:{TRADE_DATE}", "tester") + def test_admin_backfill_switch_crash_is_failed_precondition(self) -> None: + from datahub.admin_api import AdminAPI + from datahub.auth import AuthService + from datahub.crypto import SecretVault + from datahub.scheduler import Scheduler + from datahub.serving import ApiError + + vault = SecretVault(self.pipe.settings.encryption_key) + auth = AuthService(self.db, vault, self.pipe.settings.api_token, "StartPass1") + admin = AdminAPI(self.db, self.pipe, Scheduler(self.db, self.pipe), auth) + before = publications_map(self.db, TRADE_DATE) + + def explode() -> None: + raise RuntimeError("killed mid-switch") + + self.pipe.before_commit = explode + with self.assertRaises(ApiError) as ctx: + admin.backfill("valuation", TRADE_DATE, "StartPass1", f"valuation:{TRADE_DATE}", "tester") + self.assertEqual(ctx.exception.code, "FAILED_PRECONDITION") + self.assertIn("killed mid-switch", ctx.exception.message) + # previous complete A/B versions keep serving + self.assertEqual(publications_map(self.db, TRADE_DATE), before) + if __name__ == "__main__": unittest.main()