from __future__ import annotations
|
|
from pathlib import Path
|
from concurrent.futures import ThreadPoolExecutor
|
import tempfile
|
import unittest
|
|
from hibor_fast_collection.models import ContractError, ErrorCode
|
from hibor_fast_collection.quota import AppQuotaObservation, QuotaLedger
|
|
|
class QuotaTests(unittest.TestCase):
|
def test_initialize_conservative_floor_and_lifecycle(self):
|
with tempfile.TemporaryDirectory() as tmp:
|
ledger = QuotaLedger(Path(tmp) / "quota.csv", quota_date="2026-07-29")
|
snap = ledger.initialize(task_id="T", requester_role="R", handoff_id="H", known_floor=3)
|
self.assertEqual((snap.external_floor, snap.safe_available, snap.row_count), (3, 22, 1))
|
reservation = ledger.reserve(task_id="T", requester_role="R", handoff_id="H", run_id="RUN",
|
slot_id="S1", report_identity="RID")
|
self.assertTrue(reservation.allowed)
|
self.assertEqual(reservation.snapshot.active, 1)
|
terminal = ledger.terminal(reservation, event_type="CONSUME_CONFIRMED", task_id="T",
|
requester_role="R", handoff_id="H", run_id="RUN", slot_id="S1",
|
report_identity="RID", note="confirmed")
|
self.assertEqual(terminal["event_family"], "QUOTA_TERMINAL")
|
artifact = ledger.artifact(task_id="T", requester_role="R", handoff_id="H", run_id="RUN",
|
slot_id="S1", report_identity="RID",
|
quota_terminal_event_id=terminal["event_id"], success=True, note="ok")
|
self.assertEqual(artifact["event_family"], "ARTIFACT_TERMINAL")
|
after = ledger.snapshot()
|
self.assertEqual((after.confirmed, after.active, after.effective_consumed), (1, 0, 4))
|
|
def test_reservation_replay_and_terminal_conflict(self):
|
with tempfile.TemporaryDirectory() as tmp:
|
ledger = QuotaLedger(Path(tmp) / "quota.csv", quota_date="2026-07-29")
|
ledger.initialize(task_id="T", requester_role="R", handoff_id="H")
|
args = dict(task_id="T", requester_role="R", handoff_id="H", run_id="RUN",
|
slot_id="S", report_identity="RID")
|
one = ledger.reserve(**args)
|
two = ledger.reserve(**args)
|
self.assertTrue(two.replayed)
|
self.assertFalse(two.allowed)
|
self.assertEqual(two.stop_code, ErrorCode.QUOTA_REPLAY_CONFLICT)
|
ledger.terminal(one, event_type="RELEASE", note="none", **args)
|
three = ledger.reserve(**args)
|
self.assertEqual((three.allowed, three.replayed, three.terminal_event_type),
|
(False, True, "RELEASE"))
|
with self.assertRaises(ContractError):
|
ledger.terminal(one, event_type="CONSUME_CONFIRMED", note="conflict", **args)
|
|
def test_app_reconcile_is_durable_and_reservation_bound(self):
|
with tempfile.TemporaryDirectory() as tmp:
|
ledger = QuotaLedger(Path(tmp) / "quota.csv", quota_date="2026-07-29")
|
ledger.initialize(task_id="T", requester_role="R", handoff_id="H", known_floor=3)
|
observation = AppQuotaObservation.build(
|
device_serial="emulator-5554", quota_date="2026-07-29",
|
captured_at_utc="2026-07-29T01:02:03Z", visible_remaining=7,
|
ui_snapshot_fingerprint="f" * 64,
|
)
|
event = ledger.observe_app_remaining(task_id="T", requester_role="R",
|
handoff_id="H", observation=observation)
|
replay = ledger.observe_app_remaining(task_id="T", requester_role="R",
|
handoff_id="H", observation=observation)
|
self.assertEqual(event, replay)
|
self.assertEqual(ledger.snapshot().effective_consumed, 23)
|
reservation = ledger.reserve(task_id="T", requester_role="R", handoff_id="H",
|
run_id="RUN", slot_id="S", report_identity="RID",
|
observation_event_id=event["event_id"])
|
self.assertTrue(reservation.allowed)
|
with self.assertRaises(ContractError):
|
ledger.reserve(task_id="T2", requester_role="R", handoff_id="H2",
|
run_id="RUN2", slot_id="S2", report_identity="RID2",
|
observation_event_id="0" * 64)
|
|
def test_safe_zero_blocks_new_reservation(self):
|
with tempfile.TemporaryDirectory() as tmp:
|
ledger = QuotaLedger(Path(tmp) / "quota.csv", quota_date="2026-07-29")
|
ledger.initialize(task_id="T", requester_role="R", handoff_id="H", known_floor=25)
|
result = ledger.reserve(task_id="T", requester_role="R", handoff_id="H", run_id="R",
|
slot_id="S", report_identity="RID")
|
self.assertFalse(result.allowed)
|
self.assertEqual(result.stop_code, ErrorCode.QUOTA_EXHAUSTED)
|
self.assertEqual(ledger.snapshot().row_count, 1)
|
|
def test_concurrent_last_safe_slot_has_one_winner(self):
|
with tempfile.TemporaryDirectory() as tmp:
|
path = Path(tmp) / "quota.csv"
|
ledger = QuotaLedger(path, quota_date="2026-07-29")
|
ledger.initialize(task_id="T", requester_role="R", handoff_id="H", known_floor=24)
|
def reserve(index: int):
|
return QuotaLedger(path, quota_date="2026-07-29").reserve(
|
task_id=f"T{index}", requester_role="R", handoff_id=f"H{index}",
|
run_id=f"RUN{index}", slot_id=f"S{index}", report_identity=f"RID{index}")
|
with ThreadPoolExecutor(max_workers=2) as pool:
|
results = list(pool.map(reserve, (1, 2)))
|
self.assertEqual(sum(result.allowed for result in results), 1)
|
self.assertEqual(ledger.snapshot().active, 1)
|
|
|
if __name__ == "__main__": unittest.main()
|