Bound assigned-issue reminder snapshots #730
|
|
@ -403,15 +403,39 @@ async def work_page(stream: str, page: int = 1, limit: int = 50) -> dict:
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
async def assigned_issue_snapshot(*, limit: int = 50, max_pages: int = 20) -> dict:
|
async def assigned_issue_snapshot(
|
||||||
|
*,
|
||||||
|
limit: int = 50,
|
||||||
|
max_pages: int = 20,
|
||||||
|
max_concurrency: int = 4,
|
||||||
|
deadline_seconds: float = 5.0,
|
||||||
|
) -> dict:
|
||||||
"""Load a complete, bounded assigned-issue snapshot for deadline dispatch."""
|
"""Load a complete, bounded assigned-issue snapshot for deadline dispatch."""
|
||||||
items = []
|
async with asyncio.timeout(deadline_seconds):
|
||||||
for page in range(1, max_pages + 1):
|
first = await work_page("issue", 1, limit)
|
||||||
result = await work_page("issue", page, limit)
|
page_count = max(1, (first["total"] + limit - 1) // limit)
|
||||||
items.extend(result["items"])
|
if page_count > max_pages:
|
||||||
if not result["has_more"]:
|
raise ValueError("Assigned issue snapshot exceeds the scan limit")
|
||||||
return {"items": items, "complete": True}
|
semaphore = asyncio.Semaphore(max(1, max_concurrency))
|
||||||
return {"items": [], "complete": False}
|
|
||||||
|
async def load(page: int) -> dict:
|
||||||
|
async with semaphore:
|
||||||
|
return await work_page("issue", page, limit)
|
||||||
|
|
||||||
|
remaining = await asyncio.gather(
|
||||||
|
*(load(page) for page in range(2, page_count + 1))
|
||||||
|
)
|
||||||
|
pages = [first, *remaining]
|
||||||
|
if any(page.get("total") != first["total"] for page in pages[1:]):
|
||||||
|
raise ValueError("Assigned issue pagination changed during the scan")
|
||||||
|
items = [item for page in pages for item in page["items"]]
|
||||||
|
issue_ids = [
|
||||||
|
item.get("id") for item in items
|
||||||
|
if isinstance(item, dict) and isinstance(item.get("id"), int) and item["id"] > 0
|
||||||
|
]
|
||||||
|
if len(items) != first["total"] or len(set(issue_ids)) != first["total"]:
|
||||||
|
raise ValueError("Assigned issue snapshot has an incomplete issue set")
|
||||||
|
return {"items": items, "complete": True}
|
||||||
|
|
||||||
|
|
||||||
def _normalize_global_search_item(item: Any, kind: str) -> dict | None:
|
def _normalize_global_search_item(item: Any, kind: str) -> dict | None:
|
||||||
|
|
|
||||||
|
|
@ -270,23 +270,92 @@ def test_existing_deadline_preferences_migrate_to_two_day_horizon(tmp_path):
|
||||||
@pytest.mark.anyio
|
@pytest.mark.anyio
|
||||||
async def test_assigned_deadline_snapshot_is_pagination_complete(monkeypatch):
|
async def test_assigned_deadline_snapshot_is_pagination_complete(monkeypatch):
|
||||||
pages = {
|
pages = {
|
||||||
1: {"items": [{"id": 1}], "has_more": True},
|
1: {"items": [{"id": 1}], "total": 2, "has_more": True},
|
||||||
2: {"items": [{"id": 2}], "has_more": False},
|
2: {"items": [{"id": 2}], "total": 2, "has_more": False},
|
||||||
}
|
}
|
||||||
|
|
||||||
async def work_page(stream, page, limit):
|
async def work_page(stream, page, limit):
|
||||||
assert stream == "issue"
|
assert stream == "issue"
|
||||||
assert limit == 50
|
assert limit == 1
|
||||||
return pages[page]
|
return pages[page]
|
||||||
|
|
||||||
monkeypatch.setattr(gitea_proxy, "work_page", work_page)
|
monkeypatch.setattr(gitea_proxy, "work_page", work_page)
|
||||||
|
|
||||||
assert await gitea_proxy.assigned_issue_snapshot() == {
|
assert await gitea_proxy.assigned_issue_snapshot(limit=1) == {
|
||||||
"items": [{"id": 1}, {"id": 2}],
|
"items": [{"id": 1}, {"id": 2}],
|
||||||
"complete": True,
|
"complete": True,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.anyio
|
||||||
|
async def test_assigned_deadline_snapshot_fetches_remaining_pages_concurrently(monkeypatch):
|
||||||
|
active = 0
|
||||||
|
peak = 0
|
||||||
|
remaining_started = asyncio.Event()
|
||||||
|
|
||||||
|
async def work_page(stream, page, limit):
|
||||||
|
nonlocal active, peak
|
||||||
|
assert stream == "issue"
|
||||||
|
assert limit == 1
|
||||||
|
if page == 1:
|
||||||
|
return {"items": [{"id": 1}], "total": 5, "has_more": True}
|
||||||
|
active += 1
|
||||||
|
peak = max(peak, active)
|
||||||
|
if active == 2:
|
||||||
|
remaining_started.set()
|
||||||
|
await asyncio.wait_for(remaining_started.wait(), timeout=0.2)
|
||||||
|
await asyncio.sleep(0)
|
||||||
|
active -= 1
|
||||||
|
return {"items": [{"id": page}], "total": 5, "has_more": page < 5}
|
||||||
|
|
||||||
|
monkeypatch.setattr(gitea_proxy, "work_page", work_page)
|
||||||
|
|
||||||
|
assert await gitea_proxy.assigned_issue_snapshot(limit=1, max_concurrency=2) == {
|
||||||
|
"items": [{"id": 1}, {"id": 2}, {"id": 3}, {"id": 4}, {"id": 5}],
|
||||||
|
"complete": True,
|
||||||
|
}
|
||||||
|
assert peak == 2
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.anyio
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"remaining",
|
||||||
|
[
|
||||||
|
{"items": [{"id": 2}], "total": 3, "has_more": True},
|
||||||
|
{"items": [{"id": 1}], "total": 2, "has_more": False},
|
||||||
|
],
|
||||||
|
)
|
||||||
|
async def test_assigned_deadline_snapshot_rejects_changed_or_duplicate_pages(
|
||||||
|
monkeypatch, remaining
|
||||||
|
):
|
||||||
|
async def work_page(_stream, page, _limit):
|
||||||
|
if page == 1:
|
||||||
|
return {"items": [{"id": 1}], "total": 2, "has_more": True}
|
||||||
|
return remaining
|
||||||
|
|
||||||
|
monkeypatch.setattr(gitea_proxy, "work_page", work_page)
|
||||||
|
|
||||||
|
with pytest.raises(ValueError, match="incomplete|changed"):
|
||||||
|
await gitea_proxy.assigned_issue_snapshot(limit=1)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.anyio
|
||||||
|
async def test_assigned_deadline_snapshot_enforces_aggregate_deadline(monkeypatch):
|
||||||
|
cancelled = asyncio.Event()
|
||||||
|
|
||||||
|
async def work_page(_stream, _page, _limit):
|
||||||
|
try:
|
||||||
|
await asyncio.sleep(60)
|
||||||
|
finally:
|
||||||
|
cancelled.set()
|
||||||
|
|
||||||
|
monkeypatch.setattr(gitea_proxy, "work_page", work_page)
|
||||||
|
|
||||||
|
with pytest.raises(TimeoutError):
|
||||||
|
await gitea_proxy.assigned_issue_snapshot(deadline_seconds=0.01)
|
||||||
|
assert cancelled.is_set()
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.anyio
|
@pytest.mark.anyio
|
||||||
async def test_competing_workers_send_one_deadline_digest(tmp_path):
|
async def test_competing_workers_send_one_deadline_digest(tmp_path):
|
||||||
path = tmp_path / "push.sqlite3"
|
path = tmp_path / "push.sqlite3"
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue
Block a user