CoolFace
Apppublic

rl-llm-wiki/rl-bucket-sync

sourceHugging Faceupdated 3mo agoView on Hugging Face
0likes
test_wiki_queue.py103 linesDownload Raw Back to tests
1"""The discovery queue: AtomicQueue mutations, atomic claim under contention,2and the HTTP endpoints."""3import threading4 5 6def _wiki_env(make_env, **kw):7    return make_env(WIKI_ENABLED="true", WIKI_DATASET="test-org/kb", **kw)8 9 10def _ids_by_status(env):11    out = {}12    for r in env.read_model.records("queue"):13        out.setdefault(r.frontmatter.get("status"), []).append(r.frontmatter.get("id"))14    return out15 16 17def test_add_dedup_bumps_ref_count(make_env):18    env = _wiki_env(make_env)19    q = env.queue20    assert q.add([{"id": "arxiv:1", "title": "One"}, {"id": "arxiv:2"}]) == {"added": 2, "updated": 0}21    assert q.add([{"id": "arxiv:1", "discovered_from": ["arxiv:9"]}]) == {"added": 0, "updated": 1}22    fm = {r.frontmatter["id"]: r.frontmatter for r in env.read_model.records("queue")}23    assert fm["arxiv:1"]["ref_count"] == 124    assert fm["arxiv:1"]["discovered_from"] == ["arxiv:9"]25 26 27def test_claim_prefers_highest_ref_count_then_exhausts(make_env):28    env = _wiki_env(make_env)29    q = env.queue30    q.add([{"id": "arxiv:1"}, {"id": "arxiv:2"}])31    q.add([{"id": "arxiv:2"}])  # arxiv:2 now ref_count 1 -> higher priority32    assert q.claim("reader-a").id == "arxiv:2"33    assert q.claim("reader-b").id == "arxiv:1"34    assert q.claim("reader-c") is None  # frontier exhausted35 36 37def test_skip_and_mark_processed(make_env):38    env = _wiki_env(make_env)39    q = env.queue40    q.add([{"id": "arxiv:1"}, {"id": "arxiv:2"}])41    assert q.skip("reader-a", "arxiv:1", "out of scope") is True42    q.mark_processed("arxiv:2", "sources/arxiv-2.md")43    by_status = _ids_by_status(env)44    assert by_status["skipped"] == ["arxiv:1"]45    assert by_status["processed"] == ["arxiv:2"]46 47 48def test_expired_lease_is_reaped(make_env):49    env = _wiki_env(make_env, QUEUE_CLAIM_LEASE_S="0")  # lease expires immediately50    q = env.queue51    q.add([{"id": "arxiv:1"}])52    claimed = q.claim("reader-a")53    assert claimed.id == "arxiv:1"54    # lease already expired -> reaped back to the frontier, claimable again55    assert q.reap() == 156    assert q.claim("reader-b").id == "arxiv:1"57 58 59def test_claim_is_atomic_under_contention(make_env):60    env = _wiki_env(make_env)61    q = env.queue62    n = 863    q.add([{"id": f"arxiv:{i}"} for i in range(n)])64    got, lock, barrier = [], threading.Lock(), threading.Barrier(n)65 66    def worker(name):67        barrier.wait()68        it = q.claim(name)69        if it:70            with lock:71                got.append(it.id)72 73    threads = [threading.Thread(target=worker, args=(f"a{i}",)) for i in range(n)]74    for t in threads:75        t.start()76    for t in threads:77        t.join()78    # every source claimed exactly once — no double-claim79    assert len(got) == n and len(set(got)) == n80 81 82def test_queue_endpoints(make_env):83    env = _wiki_env(make_env)84    client = env.client85    assert client.post("/v1/queue:add", json={"items": [{"id": "arxiv:7", "title": "Seven"}]}).json() == {86        "added": 1, "updated": 087    }88    listing = client.get("/v1/queue").json()89    assert listing["count"] == 1 and listing["items"][0]["id"] == "arxiv:7"90 91    claim = client.post("/v1/queue:claim", json={"agent_id": "reader-a"}).json()92    assert claim["claimed"] is True and claim["id"] == "arxiv:7" and claim["lease_seconds"] > 093 94    skip = client.post("/v1/queue:skip", json={"agent_id": "reader-a", "id": "arxiv:7", "reason": "dup"})95    assert skip.status_code == 20096    assert client.get("/v1/queue").json()["by_status"].get("skipped") == 197 98 99def test_queue_404_when_wiki_disabled(env):100    # the default env fixture has wiki mode off101    assert env.client.post("/v1/queue:add", json={"items": []}).status_code == 404102    assert env.client.get("/v1/queue").status_code == 404103