From aaff109edeea18c34a2f1de6a61a2abbf9291f26 Mon Sep 17 00:00:00 2001 From: timmy Date: Fri, 7 Aug 2026 17:15:31 +0000 Subject: [PATCH] feat: keep Find Work instant during refresh (#213) --- frontend/index.html | 17 +++-- frontend/pick-work.js | 5 ++ src/gitea_proxy.py | 45 +++++++++---- src/main.py | 36 +++++++++-- tests/test_gitea_work_search.py | 111 +++++++++++++++++++++++++++++++- tests/test_my_work.py | 37 +++++++++++ 6 files changed, 227 insertions(+), 24 deletions(-) diff --git a/frontend/index.html b/frontend/index.html index b4a47fb..c7c2683 100644 --- a/frontend/index.html +++ b/frontend/index.html @@ -1660,15 +1660,20 @@ textarea { resize: vertical; min-height: 120px; } async function openFindWorkSheet() { findingWork = true; qs('#find-work-sheet').classList.add('open'); - qs('#find-work-list').textContent = ''; - qs('#find-work-status').textContent = 'Loading available issues…'; + const retainedItems = findWorkController.items(); + if (retainedItems.length) renderAvailableIssues(retainedItems); + else qs('#find-work-list').textContent = ''; + qs('#find-work-status').textContent = retainedItems.length ? + 'Refreshing available issues…' : 'Loading available issues…'; qs('#load-more-available').hidden = true; qs('#close-find-work').focus(); try { - await findWorkController.load(); - qs('#find-work-status').textContent = findWorkController.items().length ? - findWorkController.items().length + ' of ' + availablePagination.total + ' available issues loaded.' : - 'No unassigned issues are available.'; + const result = await findWorkController.load(); + if (!result?.stale) { + qs('#find-work-status').textContent = findWorkController.items().length ? + findWorkController.items().length + ' of ' + availablePagination.total + ' available issues loaded.' : + 'No unassigned issues are available.'; + } } catch (error) { qs('#find-work-status').textContent = error.message + ' Close and retry.'; } diff --git a/frontend/pick-work.js b/frontend/pick-work.js index ce10ffe..5b3b05c 100644 --- a/frontend/pick-work.js +++ b/frontend/pick-work.js @@ -39,6 +39,11 @@ function createFindWork({ fetchJson, onItems, onPagination, onStatus }) { headers: { Accept: 'application/json' }, }).then(result => { apply(result, append); + if (result?.refresh_failed === true) { + onStatus('Showing saved available work. Catalog refresh failed; retrying shortly.'); + } else if (result?.revalidating === true) { + onStatus('Showing saved available work while the catalog refreshes…'); + } return result; }).finally(() => { loadRequest = null; }); return loadRequest; diff --git a/src/gitea_proxy.py b/src/gitea_proxy.py index 97ab58d..e343a95 100644 --- a/src/gitea_proxy.py +++ b/src/gitea_proxy.py @@ -11,6 +11,7 @@ GITEA_URL = os.getenv("GITEA_URL", "http://127.0.0.1:3000").rstrip("/") GITEA_TOKEN = os.getenv("GITEA_TOKEN", "") REVIEW_DIFF_MAX_BYTES = 64 * 1024 REVIEW_DIFF_MAX_LINES = 400 +AVAILABLE_ISSUE_PAGE_CONCURRENCY = 3 _client: httpx.AsyncClient | None = None @@ -328,8 +329,8 @@ async def available_issue_snapshot(max_pages: int = 10, upstream_limit: int = 50 """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): + + async def load_page(upstream_page: int) -> tuple[list[Any], int | None]: response = await _get_client().get( "/api/v1/repos/issues/search", headers=_auth(), @@ -342,21 +343,41 @@ async def available_issue_snapshot(max_pages: int = 10, upstream_limit: int = 50 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 + try: + total = max(0, int(response.headers["X-Total-Count"])) + except (KeyError, TypeError, ValueError): + total = None + return payload, total + + first_payload, upstream_total = await load_page(1) + pages: list[list[Any]] = [first_payload] + if upstream_total is not None: + page_count = min(max_pages, max(1, (upstream_total + upstream_limit - 1) // upstream_limit)) + semaphore = asyncio.Semaphore(AVAILABLE_ISSUE_PAGE_CONCURRENCY) + + async def load_bounded(page: int) -> list[Any]: + async with semaphore: + payload, _ = await load_page(page) + return payload + + if page_count > 1: + pages.extend(await asyncio.gather(*( + load_bounded(page) for page in range(2, page_count + 1) + ))) + else: + previous = first_payload + for page in range(2, max_pages + 1): + if not previous or len(previous) < upstream_limit: + break + previous, _ = await load_page(page) + pages.append(previous) + + for payload in pages: 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 - ): - break priority = {"p0", "priority-high", "critical"} items.sort(key=lambda item: (item["repository"], item["number"] or 0)) diff --git a/src/main.py b/src/main.py index 0f9b051..3937098 100644 --- a/src/main.py +++ b/src/main.py @@ -76,6 +76,7 @@ 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 +AVAILABLE_ISSUE_SNAPSHOT_RETRY_SECONDS = 5.0 FRONTEND_DIR = Path(__file__).resolve().parent.parent / "frontend" _live_snapshot_task: asyncio.Task | None = None _live_snapshot_value: dict | None = None @@ -107,6 +108,7 @@ _idempotency_ledger = IdempotencyLedger( _available_issue_snapshot_task: asyncio.Task | None = None _available_issue_snapshot_value: list[dict] | None = None _available_issue_snapshot_created_at: float | None = None +_available_issue_snapshot_retry_at: float | None = None class ContextPayloadError(ValueError): @@ -600,13 +602,24 @@ async def paged_work( async def _refresh_available_issue_snapshot() -> list[dict]: global _available_issue_snapshot_value, _available_issue_snapshot_created_at + global _available_issue_snapshot_retry_at result = await gitea_proxy.available_issue_snapshot() _available_issue_snapshot_value = result _available_issue_snapshot_created_at = time.monotonic() + _available_issue_snapshot_retry_at = None return result -async def _available_issue_snapshot() -> tuple[list[dict], bool]: +def _observe_available_issue_refresh(task: asyncio.Task) -> None: + global _available_issue_snapshot_retry_at + if not task.cancelled(): + if task.exception() is not None: + _available_issue_snapshot_retry_at = ( + time.monotonic() + AVAILABLE_ISSUE_SNAPSHOT_RETRY_SECONDS + ) + + +async def _available_issue_snapshot() -> tuple[list[dict], bool, bool, bool]: global _available_issue_snapshot_task now = time.monotonic() if ( @@ -615,16 +628,25 @@ async def _available_issue_snapshot() -> tuple[list[dict], bool]: and now - _available_issue_snapshot_created_at < AVAILABLE_ISSUE_SNAPSHOT_FRESHNESS_SECONDS ): - return _available_issue_snapshot_value, False + return _available_issue_snapshot_value, False, False, False + if ( + _available_issue_snapshot_value is not None + and _available_issue_snapshot_retry_at is not None + and now < _available_issue_snapshot_retry_at + ): + return _available_issue_snapshot_value, True, False, True if _available_issue_snapshot_task is None or _available_issue_snapshot_task.done(): _available_issue_snapshot_task = asyncio.create_task( _refresh_available_issue_snapshot() ) + _available_issue_snapshot_task.add_done_callback(_observe_available_issue_refresh) + if _available_issue_snapshot_value is not None: + return _available_issue_snapshot_value, True, True, False try: - return await asyncio.shield(_available_issue_snapshot_task), False + return await asyncio.shield(_available_issue_snapshot_task), False, False, False except Exception: if _available_issue_snapshot_value is not None: - return _available_issue_snapshot_value, True + return _available_issue_snapshot_value, True, False, True raise @@ -642,7 +664,7 @@ def _available_issue_page(items: list[dict], page: int, limit: int = 50) -> dict @app.get("/api/v1/available-issues") async def available_issues(page: int = Query(default=1, ge=1, le=100)) -> JSONResponse: try: - items, stale = await asyncio.wait_for( + items, stale, revalidating, refresh_failed = await asyncio.wait_for( _available_issue_snapshot(), timeout=WORK_PAGE_TIMEOUT_SECONDS ) except Exception: @@ -654,6 +676,10 @@ async def available_issues(page: int = Query(default=1, ge=1, le=100)) -> JSONRe result = _available_issue_page(items, page) if stale: result["stale"] = True + if revalidating: + result["revalidating"] = True + if refresh_failed: + result["refresh_failed"] = True return JSONResponse(result) diff --git a/tests/test_gitea_work_search.py b/tests/test_gitea_work_search.py index 4e42b29..1642781 100644 --- a/tests/test_gitea_work_search.py +++ b/tests/test_gitea_work_search.py @@ -13,6 +13,7 @@ def reset_available_issue_snapshot(): main._available_issue_snapshot_task = None main._available_issue_snapshot_value = None main._available_issue_snapshot_created_at = None + main._available_issue_snapshot_retry_at = None yield task = main._available_issue_snapshot_task if task is not None and not task.done(): @@ -20,6 +21,7 @@ def reset_available_issue_snapshot(): main._available_issue_snapshot_task = None main._available_issue_snapshot_value = None main._available_issue_snapshot_created_at = None + main._available_issue_snapshot_retry_at = None @pytest.mark.anyio @@ -136,6 +138,53 @@ async def test_available_issue_page_ranks_all_upstream_pages_before_logical_pagi } +@pytest.mark.anyio +async def test_available_issue_snapshot_loads_known_remaining_pages_concurrently(): + active = 0 + peak_active = 0 + remaining_started = asyncio.Event() + release = asyncio.Event() + + async def upstream(request): + nonlocal active, peak_active + page = int(request.url.params["page"]) + if page == 1: + return httpx.Response( + 200, headers={"X-Total-Count": "200"}, + json=[{ + "id": number, "number": number, "title": f"Issue {number}", + "state": "open", "assignees": [], "pull_request": None, + "labels": [], "repository": {"full_name": "stackchain/api"}, + } for number in range(1, 51)], + ) + active += 1 + peak_active = max(peak_active, active) + if peak_active == 3: + remaining_started.set() + await release.wait() + active -= 1 + start = (page - 1) * 50 + 1 + return httpx.Response(200, headers={"X-Total-Count": "200"}, json=[{ + "id": number, "number": number, "title": f"Issue {number}", + "state": "open", "assignees": [], "pull_request": None, + "labels": [], "repository": {"full_name": "stackchain/api"}, + } for number in range(start, start + 50)]) + + gitea_proxy.start_client(transport=httpx.MockTransport(upstream)) + task = asyncio.create_task(gitea_proxy.available_issue_snapshot()) + try: + await asyncio.wait_for(remaining_started.wait(), timeout=1) + assert peak_active == 3 + release.set() + result = await asyncio.wait_for(task, timeout=1) + finally: + release.set() + await gitea_proxy.stop_client() + + assert len(result) == 200 + assert peak_active == 3 + + @pytest.mark.anyio async def test_available_issue_endpoint_is_bounded_retryable_and_no_store(monkeypatch): calls = [] @@ -204,7 +253,67 @@ async def test_available_issue_endpoint_retains_last_snapshot_on_refresh_failure assert response.status_code == 200 assert response.json() == { "items": [{"number": 7}], "page": 1, "total": 1, - "has_more": False, "stale": True, + "has_more": False, "stale": True, "revalidating": True, + } + + +@pytest.mark.anyio +async def test_available_issue_endpoint_serves_expired_snapshot_while_one_refresh_runs(monkeypatch): + calls = 0 + started = asyncio.Event() + release = asyncio.Event() + + async def refresh(): + nonlocal calls + calls += 1 + started.set() + await release.wait() + return [{"number": 8}] + + main._available_issue_snapshot_value = [{"number": 7}] + main._available_issue_snapshot_created_at = 0.0 + monkeypatch.setattr(main.gitea_proxy, "available_issue_snapshot", refresh) + transport = httpx.ASGITransport(app=main.app) + async with httpx.AsyncClient(transport=transport, base_url="http://test") as client: + first, second = await asyncio.gather( + client.get("/api/v1/available-issues?page=1"), + client.get("/api/v1/available-issues?page=1"), + ) + await asyncio.wait_for(started.wait(), timeout=1) + + assert calls == 1 + assert first.json() == second.json() == { + "items": [{"number": 7}], "page": 1, "total": 1, + "has_more": False, "stale": True, "revalidating": True, + } + release.set() + await asyncio.wait_for(main._available_issue_snapshot_task, timeout=1) + assert main._available_issue_snapshot_value == [{"number": 8}] + + +@pytest.mark.anyio +async def test_available_issue_endpoint_backs_off_after_background_refresh_failure(monkeypatch): + calls = 0 + + async def unavailable(): + nonlocal calls + calls += 1 + 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: + first = await client.get("/api/v1/available-issues?page=1") + await asyncio.gather(main._available_issue_snapshot_task, return_exceptions=True) + second = await client.get("/api/v1/available-issues?page=1") + + assert calls == 1 + assert first.json()["revalidating"] is True + assert second.json() == { + "items": [{"number": 7}], "page": 1, "total": 1, + "has_more": False, "stale": True, "refresh_failed": True, } diff --git a/tests/test_my_work.py b/tests/test_my_work.py index 20a2595..7168cb1 100644 --- a/tests/test_my_work.py +++ b/tests/test_my_work.py @@ -920,6 +920,40 @@ controller.load().then(() => controller.loadMore()).then(() => assert output["pages"][-1] == {"page": 2, "total": 2, "has_more": False} +def test_find_work_keeps_retained_cards_and_announces_refresh_freshness(): + script = f""" +const createFindWork = require({json.dumps(str(PICK_WORK))}); +const statuses = []; +const states = []; +const responses = [ + {{items:[{{id:1,number:1,repository:'stackchain/api'}}],page:1,total:1,has_more:false, + stale:true,revalidating:true}}, + {{items:[{{id:1,number:1,repository:'stackchain/api'}}],page:1,total:1,has_more:false, + stale:true,refresh_failed:true}}, +]; +const controller = createFindWork({{ + fetchJson: () => Promise.resolve(responses.shift()), + onItems: items => states.push(items), + onPagination: () => {{}}, + onStatus: status => statuses.push(status), +}}); +controller.load().then(() => controller.load()).then(() => + process.stdout.write(JSON.stringify({{statuses,states,items:controller.items()}})) +); +""" + result = subprocess.run( + ["node", "-e", script], check=True, capture_output=True, text=True + ) + output = json.loads(result.stdout) + + assert output["statuses"] == [ + "Showing saved available work while the catalog refreshes…", + "Showing saved available work. Catalog refresh failed; retrying shortly.", + ] + assert [item["number"] for item in output["items"]] == [1] + assert len(output["states"]) == 2 + + def test_find_work_preview_stays_with_issue_across_pagination_and_clears_when_claimed(): script = f""" const createFindWork = require({json.dumps(str(PICK_WORK))}); @@ -1000,6 +1034,9 @@ async def test_mobile_find_work_sheet_is_accessible_touch_sized_and_subpath_safe assert '.find-work-action { min-height:44px;' in html assert 'padding-bottom:calc(18px + env(safe-area-inset-bottom))' in html assert '@media(max-width:320px)' in html + assert 'const retainedItems = findWorkController.items();' in html + assert 'if (retainedItems.length) renderAvailableIssues(retainedItems);' in html + assert 'Refreshing available issues…' in html @pytest.mark.anyio -- 2.43.0