rl-llm-wiki/rl-bucket-sync
0
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 