Make Find Work production-schema-safe and truthfully paginated #182

Merged
rockachopa merged 1 commits from timmy/181-find-work-schema-pagination into main 2026-08-07 08:59:02 +00:00
7 changed files with 331 additions and 71 deletions

View File

@ -83,8 +83,11 @@ function createFindWork({ fetchJson, onItems, onPagination, onStatus }) {
available = available.filter(candidate =>
candidate.repository !== item.repository || candidate.number !== item.number
);
pagination.total = Math.max(available.length, pagination.total - 1);
pagination.has_more = available.length < pagination.total;
previewed.delete(itemKey(item));
onItems(available.slice());
onPagination({ ...pagination });
onStatus('Assigned ' + key + ' to you.');
return confirmed;
}).finally(() => { claimRequest = null; });

View File

@ -179,48 +179,70 @@ async def work_page(stream: str, page: int = 1, limit: int = 50) -> dict:
}
async def available_issue_page(page: int = 1, limit: int = 50) -> dict:
"""Return one bounded page of open, unassigned issues visible to the user."""
response = await _get_client().get(
"/api/v1/repos/issues/search",
headers=_auth(),
params={"state": "open", "type": "issues", "limit": limit, "page": page},
)
response.raise_for_status()
payload = response.json()
if not isinstance(payload, list):
raise ValueError("Gitea available issue search response was not a list")
def _normalize_available_issue(item: Any) -> dict | None:
if (
not isinstance(item, dict)
or item.get("state") != "open"
or item.get("pull_request") is not None
or item.get("assignees") not in (None, [])
):
return None
labels_value = item.get("labels")
labels = labels_value if isinstance(labels_value, list) else []
repository_value = item.get("repository")
repository = repository_value if isinstance(repository_value, dict) else {}
return {
"id": item.get("id"),
"number": item.get("number"),
"title": item.get("title", "") if isinstance(item.get("title"), str) else "",
"body": item.get("body", "") if isinstance(item.get("body"), str) else "",
"state": "open",
"repository": repository.get("full_name", "")
if isinstance(repository.get("full_name"), str) else "",
"labels": [
label["name"] for label in labels
if isinstance(label, dict) and isinstance(label.get("name"), str)
],
"assignees": [],
"updated_at": item.get("updated_at", "")
if isinstance(item.get("updated_at"), str) else "",
"url": _safe_web_url(item.get("html_url")),
}
items = []
for item in payload:
if (
not isinstance(item, dict)
or item.get("state") != "open"
or "pull_request" in item
or item.get("assignees") not in (None, [])
async def available_issue_snapshot(max_pages: int = 10, upstream_limit: int = 50) -> list[dict]:
"""Load, filter, and globally rank a bounded snapshot of available issues."""
items: list[dict] = []
seen_ids: set[Any] = set()
upstream_total: int | None = None
for upstream_page in range(1, max_pages + 1):
response = await _get_client().get(
"/api/v1/repos/issues/search",
headers=_auth(),
params={
"state": "open", "type": "issues",
"limit": upstream_limit, "page": upstream_page,
},
)
response.raise_for_status()
payload = response.json()
if not isinstance(payload, list):
raise ValueError("Gitea available issue search response was not a list")
if upstream_total is None:
try:
upstream_total = max(0, int(response.headers["X-Total-Count"]))
except (KeyError, TypeError, ValueError):
upstream_total = None
for raw_item in payload:
item = _normalize_available_issue(raw_item)
if item is not None and item["id"] not in seen_ids:
seen_ids.add(item["id"])
items.append(item)
loaded = upstream_page * upstream_limit
if not payload or len(payload) < upstream_limit or (
upstream_total is not None and loaded >= upstream_total
):
continue
labels_value = item.get("labels")
labels = labels_value if isinstance(labels_value, list) else []
repository_value = item.get("repository")
repository = repository_value if isinstance(repository_value, dict) else {}
items.append({
"id": item.get("id"),
"number": item.get("number"),
"title": item.get("title", "") if isinstance(item.get("title"), str) else "",
"body": item.get("body", "") if isinstance(item.get("body"), str) else "",
"state": "open",
"repository": repository.get("full_name", "")
if isinstance(repository.get("full_name"), str) else "",
"labels": [
label["name"] for label in labels
if isinstance(label, dict) and isinstance(label.get("name"), str)
],
"assignees": [],
"updated_at": item.get("updated_at", "")
if isinstance(item.get("updated_at"), str) else "",
"url": _safe_web_url(item.get("html_url")),
})
break
priority = {"p0", "priority-high", "critical"}
items.sort(key=lambda item: (item["repository"], item["number"] or 0))
@ -228,11 +250,21 @@ async def available_issue_page(page: int = 1, limit: int = 50) -> dict:
items.sort(key=lambda item: (
0 if any(str(label).lower() in priority for label in item["labels"]) else 1
))
try:
total = max(len(items), int(response.headers.get("X-Total-Count", len(items))))
except (TypeError, ValueError):
total = len(items)
return {"items": items, "page": page, "total": total, "has_more": page * limit < total}
return items
async def available_issue_page(page: int = 1, limit: int = 50) -> dict:
"""Return one logical page from a bounded, globally ranked available-work scan."""
items = await available_issue_snapshot()
total = len(items)
start = (page - 1) * limit
page_items = items[start:start + limit]
return {
"items": page_items,
"page": page,
"total": total,
"has_more": start + len(page_items) < total,
}
def _page_metadata(result: dict) -> dict:
@ -575,7 +607,7 @@ async def claim_available_issue(repository: str, number: int) -> dict:
if (
not isinstance(issue, dict)
or issue.get("state") != "open"
or "pull_request" in issue
or issue.get("pull_request") is not None
or assignees_value not in (None, [])
):
raise IssueNotAvailableError("Issue is no longer available")

View File

@ -30,23 +30,27 @@ from src.views import router as frontend_router
@asynccontextmanager
async def lifespan(_app: FastAPI):
global _live_snapshot_task
global _live_snapshot_task, _available_issue_snapshot_task
gitea_proxy.start_client()
try:
yield
finally:
task = _live_snapshot_task
if task is not None and not task.done():
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
live_task = _live_snapshot_task
available_task = _available_issue_snapshot_task
for task in (live_task, available_task):
if task is not None and not task.done():
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
try:
await gitea_proxy.stop_client()
finally:
if _live_snapshot_task is task:
if _live_snapshot_task is live_task:
_live_snapshot_task = None
if _available_issue_snapshot_task is available_task:
_available_issue_snapshot_task = None
app = FastAPI(title="Stackchain Dashboard", lifespan=lifespan)
@ -66,6 +70,7 @@ BULK_NOTIFICATION_DEADLINE_SECONDS = 6.0
LIVE_SNAPSHOT_FRESHNESS_SECONDS = 8.0
LIVE_SNAPSHOT_RETRY_BASE_SECONDS = 5.0
LIVE_SNAPSHOT_RETRY_MAX_SECONDS = 60.0
AVAILABLE_ISSUE_SNAPSHOT_FRESHNESS_SECONDS = 15.0
FRONTEND_DIR = Path(__file__).resolve().parent.parent / "frontend"
_live_snapshot_task: asyncio.Task | None = None
_live_snapshot_value: dict | None = None
@ -85,6 +90,9 @@ _read_notification_ids: set[int] = set()
_issue_creation_operations: dict[
str, tuple[tuple[Any, ...], asyncio.Task, float]
] = {}
_available_issue_snapshot_task: asyncio.Task | None = None
_available_issue_snapshot_value: list[dict] | None = None
_available_issue_snapshot_created_at: float | None = None
class ContextPayloadError(ValueError):
@ -405,11 +413,52 @@ async def paged_work(
})
async def _refresh_available_issue_snapshot() -> list[dict]:
global _available_issue_snapshot_value, _available_issue_snapshot_created_at
result = await gitea_proxy.available_issue_snapshot()
_available_issue_snapshot_value = result
_available_issue_snapshot_created_at = time.monotonic()
return result
async def _available_issue_snapshot() -> tuple[list[dict], bool]:
global _available_issue_snapshot_task
now = time.monotonic()
if (
_available_issue_snapshot_value is not None
and _available_issue_snapshot_created_at is not None
and now - _available_issue_snapshot_created_at
< AVAILABLE_ISSUE_SNAPSHOT_FRESHNESS_SECONDS
):
return _available_issue_snapshot_value, False
if _available_issue_snapshot_task is None or _available_issue_snapshot_task.done():
_available_issue_snapshot_task = asyncio.create_task(
_refresh_available_issue_snapshot()
)
try:
return await asyncio.shield(_available_issue_snapshot_task), False
except Exception:
if _available_issue_snapshot_value is not None:
return _available_issue_snapshot_value, True
raise
def _available_issue_page(items: list[dict], page: int, limit: int = 50) -> dict:
start = (page - 1) * limit
page_items = items[start:start + limit]
return {
"items": page_items,
"page": page,
"total": len(items),
"has_more": start + len(page_items) < len(items),
}
@app.get("/api/v1/available-issues")
async def available_issues(page: int = Query(default=1, ge=1, le=100)) -> JSONResponse:
try:
result = await asyncio.wait_for(
gitea_proxy.available_issue_page(page), timeout=WORK_PAGE_TIMEOUT_SECONDS
items, stale = await asyncio.wait_for(
_available_issue_snapshot(), timeout=WORK_PAGE_TIMEOUT_SECONDS
)
except Exception:
return JSONResponse(
@ -417,6 +466,9 @@ async def available_issues(page: int = Query(default=1, ge=1, le=100)) -> JSONRe
status_code=503,
headers={"Retry-After": str(math.ceil(WORK_PAGE_TIMEOUT_SECONDS))},
)
result = _available_issue_page(items, page)
if stale:
result["stale"] = True
return JSONResponse(result)
@ -897,6 +949,13 @@ async def claim_available_issue(
status_code=503,
headers={"Retry-After": "1"},
)
if _available_issue_snapshot_value is not None:
_available_issue_snapshot_value[:] = [
item for item in _available_issue_snapshot_value
if not (
item.get("repository") == repository and item.get("number") == number
)
]
return JSONResponse(result)

View File

@ -79,3 +79,38 @@ async def test_application_shutdown_finishes_snapshot_before_closing_transport(m
main._live_snapshot_task = None
assert calls == ["start", "snapshot cancelled", "stop (done)"]
@pytest.mark.anyio
async def test_application_shutdown_cancels_available_work_scan_before_transport(monkeypatch):
calls = []
started = asyncio.Event()
async def active_scan():
started.set()
try:
await asyncio.Event().wait()
finally:
calls.append("available scan cancelled")
monkeypatch.setattr(gitea_proxy, "start_client", lambda: calls.append("start"))
async def stop_client():
task = main._available_issue_snapshot_task
calls.append(f"stop ({'done' if task is not None and task.done() else 'active'})")
monkeypatch.setattr(gitea_proxy, "stop_client", stop_client)
try:
async with main.app.router.lifespan_context(main.app):
main._available_issue_snapshot_task = asyncio.create_task(active_scan())
await started.wait()
finally:
task = main._available_issue_snapshot_task
if task is not None and not task.done():
task.cancel()
with pytest.raises(asyncio.CancelledError):
await task
main._available_issue_snapshot_task = None
assert calls == ["start", "available scan cancelled", "stop (done)"]

View File

@ -8,6 +8,20 @@ from src import gitea_proxy
from src import main
@pytest.fixture(autouse=True)
def reset_available_issue_snapshot():
main._available_issue_snapshot_task = None
main._available_issue_snapshot_value = None
main._available_issue_snapshot_created_at = None
yield
task = main._available_issue_snapshot_task
if task is not None and not task.done():
task.cancel()
main._available_issue_snapshot_task = None
main._available_issue_snapshot_value = None
main._available_issue_snapshot_created_at = None
@pytest.mark.anyio
async def test_work_page_preserves_total_and_reason_without_loading_other_pages():
requests = []
@ -50,10 +64,11 @@ async def test_available_issue_page_filters_assigned_and_pull_items_then_ranks_p
requests.append(str(request.url))
return httpx.Response(
200,
headers={"X-Total-Count": "77"},
headers={"X-Total-Count": "4"},
json=[
{"id": 1, "number": 1, "title": "Ordinary", "state": "open",
"updated_at": "2026-08-07T12:00:00Z", "assignees": [],
"updated_at": "2026-08-07T12:00:00Z", "assignees": None,
"pull_request": None,
"labels": [], "repository": {"full_name": "stackchain/api"}},
{"id": 2, "number": 2, "title": "Claimed", "state": "open",
"assignees": [{"login": "alex"}], "repository": {"full_name": "stackchain/api"}},
@ -62,6 +77,7 @@ async def test_available_issue_page_filters_assigned_and_pull_items_then_ranks_p
"repository": {"full_name": "stackchain/api"}},
{"id": 4, "number": 4, "title": "Critical", "state": "open",
"updated_at": "2026-08-07T10:00:00Z", "assignees": [],
"pull_request": None,
"labels": [{"name": "critical"}],
"repository": {"full_name": "stackchain/web"}},
],
@ -69,16 +85,54 @@ async def test_available_issue_page_filters_assigned_and_pull_items_then_ranks_p
gitea_proxy.start_client(transport=httpx.MockTransport(upstream))
try:
result = await gitea_proxy.available_issue_page(page=2)
result = await gitea_proxy.available_issue_page(page=1)
finally:
await gitea_proxy.stop_client()
assert requests == [
"http://127.0.0.1:3000/api/v1/repos/issues/search?state=open&type=issues&limit=50&page=2"
"http://127.0.0.1:3000/api/v1/repos/issues/search?state=open&type=issues&limit=50&page=1"
]
assert [item["title"] for item in result["items"]] == ["Critical", "Ordinary"]
assert result == {
"items": result["items"], "page": 2, "total": 77, "has_more": False,
"items": result["items"], "page": 1, "total": 2, "has_more": False,
}
@pytest.mark.anyio
async def test_available_issue_page_ranks_all_upstream_pages_before_logical_pagination():
requests = []
def issue(number, *, assigned=True, critical=False, updated_at="2026-08-07T12:00:00Z"):
return {
"id": number, "number": number, "title": f"Issue {number}", "state": "open",
"updated_at": updated_at, "assignees": [{"login": "alex"}] if assigned else None,
"pull_request": None,
"labels": [{"name": "critical"}] if critical else [],
"repository": {"full_name": "stackchain/api"},
}
first_page = [issue(number) for number in range(1, 51)]
first_page[0] = issue(1, assigned=False, updated_at="2026-08-07T13:00:00Z")
def upstream(request):
requests.append(str(request.url))
page = int(request.url.params["page"])
payload = first_page if page == 1 else [
issue(51, assigned=False, critical=True, updated_at="2026-08-07T10:00:00Z"),
issue(52, assigned=False, updated_at="2026-08-07T11:00:00Z"),
]
return httpx.Response(200, headers={"X-Total-Count": "52"}, json=payload)
gitea_proxy.start_client(transport=httpx.MockTransport(upstream))
try:
result = await gitea_proxy.available_issue_page(page=1, limit=2)
finally:
await gitea_proxy.stop_client()
assert len(requests) == 2
assert [item["number"] for item in result["items"]] == [51, 1]
assert result == {
"items": result["items"], "page": 1, "total": 3, "has_more": True,
}
@ -86,11 +140,11 @@ async def test_available_issue_page_filters_assigned_and_pull_items_then_ranks_p
async def test_available_issue_endpoint_is_bounded_retryable_and_no_store(monkeypatch):
calls = []
async def available(page):
calls.append(page)
return {"items": [{"number": 7}], "page": page, "total": 51, "has_more": True}
async def available():
calls.append(True)
return [{"number": number} for number in range(1, 52)]
monkeypatch.setattr(main.gitea_proxy, "available_issue_page", available)
monkeypatch.setattr(main.gitea_proxy, "available_issue_snapshot", available)
transport = httpx.ASGITransport(app=main.app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
response = await client.get("/api/v1/available-issues?page=1")
@ -98,9 +152,84 @@ async def test_available_issue_endpoint_is_bounded_retryable_and_no_store(monkey
assert response.status_code == 200
assert response.headers["cache-control"] == "no-store"
assert response.json() == {
"items": [{"number": 7}], "page": 1, "total": 51, "has_more": True,
"items": [{"number": number} for number in range(1, 51)],
"page": 1, "total": 51, "has_more": True,
}
assert calls == [True]
@pytest.mark.anyio
async def test_available_issue_endpoint_coalesces_cold_scan_and_reuses_it_for_pages(monkeypatch):
calls = 0
started = asyncio.Event()
release = asyncio.Event()
async def available():
nonlocal calls
calls += 1
started.set()
await release.wait()
return [{"number": number} for number in range(1, 64)]
monkeypatch.setattr(main.gitea_proxy, "available_issue_snapshot", available)
transport = httpx.ASGITransport(app=main.app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
first = asyncio.create_task(client.get("/api/v1/available-issues?page=1"))
second = asyncio.create_task(client.get("/api/v1/available-issues?page=2"))
await asyncio.wait_for(started.wait(), timeout=1)
release.set()
first_response, second_response = await asyncio.gather(first, second)
repeated_response = await client.get("/api/v1/available-issues?page=1")
assert calls == 1
assert len(first_response.json()["items"]) == 50
assert len(second_response.json()["items"]) == 13
assert second_response.json()["total"] == 63
assert second_response.json()["has_more"] is False
assert repeated_response.status_code == 200
@pytest.mark.anyio
async def test_available_issue_endpoint_retains_last_snapshot_on_refresh_failure(monkeypatch):
async def unavailable():
raise httpx.ConnectError("offline")
main._available_issue_snapshot_value = [{"number": 7}]
main._available_issue_snapshot_created_at = 0.0
monkeypatch.setattr(main.gitea_proxy, "available_issue_snapshot", unavailable)
transport = httpx.ASGITransport(app=main.app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
response = await client.get("/api/v1/available-issues?page=1")
assert response.status_code == 200
assert response.json() == {
"items": [{"number": 7}], "page": 1, "total": 1,
"has_more": False, "stale": True,
}
@pytest.mark.anyio
async def test_confirmed_claim_is_removed_from_retained_available_snapshot(monkeypatch):
main._available_issue_snapshot_value = [
{"number": 7, "repository": "stackchain/api"},
{"number": 8, "repository": "stackchain/web"},
]
main._available_issue_snapshot_created_at = 10**12
async def claim(repository, number):
return {"repository": repository, "number": number, "assignees": ["timmy"]}
monkeypatch.setattr(main.gitea_proxy, "claim_available_issue", claim)
transport = httpx.ASGITransport(app=main.app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
claimed = await client.patch("/api/v1/repos/stackchain/api/issues/7/claim")
available = await client.get("/api/v1/available-issues?page=1")
assert claimed.status_code == 200
assert available.json() == {
"items": [{"number": 8, "repository": "stackchain/web"}],
"page": 1, "total": 1, "has_more": False,
}
assert calls == [1]
@pytest.mark.anyio

View File

@ -393,7 +393,7 @@ async def test_gitea_claim_available_issue_rechecks_then_confirms_authenticated_
if request.method == "GET" and request.url.path.endswith("/issues/17"):
return httpx.Response(200, json={
"id": 81, "number": 17, "title": "Available", "state": "open",
"assignees": [], "labels": [],
"assignees": None, "pull_request": None, "labels": [],
"html_url": "https://forge.example/stackchain/api/issues/17",
})
if request.method == "GET" and request.url.path == "/api/v1/user":

View File

@ -245,6 +245,7 @@ let calls = 0;
let release;
const states = [];
const statuses = [];
const pages = [];
const controller = createFindWork({{
fetchJson: (url, options) => {{
calls += 1;
@ -254,7 +255,7 @@ const controller = createFindWork({{
}}); }});
}},
onItems: items => states.push(items),
onPagination: () => {{}},
onPagination: page => pages.push(page),
onStatus: status => statuses.push(status),
}});
controller.reset({{items:[
@ -266,7 +267,7 @@ const first = controller.claim(item);
const duplicate = controller.claim(item);
release();
Promise.all([first,duplicate]).then(results => process.stdout.write(JSON.stringify({{
calls,states,statuses,results,remaining:controller.items()
calls,states,statuses,pages,results,remaining:controller.items()
}})));
"""
result = subprocess.run(
@ -278,6 +279,7 @@ Promise.all([first,duplicate]).then(results => process.stdout.write(JSON.stringi
assert [item["number"] for item in output["remaining"]] == [8]
assert [item["number"] for item in output["states"][-1]] == [8]
assert output["statuses"] == ["Assigning stackchain/api#7…", "Assigned stackchain/api#7 to you."]
assert output["pages"][-1] == {"page": 1, "total": 1, "has_more": False}
assert output["results"][0]["assignees"] == ["timmy"]
assert output["results"][1]["assignees"] == ["timmy"]