From b6c7ec83e41ff78f42988b546433392e8d22b1a1 Mon Sep 17 00:00:00 2001 From: timmy Date: Fri, 7 Aug 2026 08:58:07 +0000 Subject: [PATCH] feat: make Find Work production-schema-safe (#181) --- frontend/pick-work.js | 3 + src/gitea_proxy.py | 124 ++++++++++++++++---------- src/main.py | 81 ++++++++++++++--- tests/test_gitea_transport.py | 35 ++++++++ tests/test_gitea_work_search.py | 151 +++++++++++++++++++++++++++++--- tests/test_issue_api.py | 2 +- tests/test_my_work.py | 6 +- 7 files changed, 331 insertions(+), 71 deletions(-) diff --git a/frontend/pick-work.js b/frontend/pick-work.js index f7be839..ce10ffe 100644 --- a/frontend/pick-work.js +++ b/frontend/pick-work.js @@ -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; }); diff --git a/src/gitea_proxy.py b/src/gitea_proxy.py index a54f450..414e00e 100644 --- a/src/gitea_proxy.py +++ b/src/gitea_proxy.py @@ -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") diff --git a/src/main.py b/src/main.py index 072d901..ec615f5 100644 --- a/src/main.py +++ b/src/main.py @@ -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) diff --git a/tests/test_gitea_transport.py b/tests/test_gitea_transport.py index f7fc3cb..16a1888 100644 --- a/tests/test_gitea_transport.py +++ b/tests/test_gitea_transport.py @@ -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)"] diff --git a/tests/test_gitea_work_search.py b/tests/test_gitea_work_search.py index 1624eaf..3af39b4 100644 --- a/tests/test_gitea_work_search.py +++ b/tests/test_gitea_work_search.py @@ -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 diff --git a/tests/test_issue_api.py b/tests/test_issue_api.py index 5435a18..60ba6e0 100644 --- a/tests/test_issue_api.py +++ b/tests/test_issue_api.py @@ -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": diff --git a/tests/test_my_work.py b/tests/test_my_work.py index c0596f4..d8bb841 100644 --- a/tests/test_my_work.py +++ b/tests/test_my_work.py @@ -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"]