Cai
2026-08-09 9f6cddff222fe1e1f77a3b66c7b673fec2695999
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
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()