Merge pull request 'Keep Find Work instant while its catalog refreshes' (#214) from timmy/213-instant-find-work into main
All checks were successful
CI / lint (push) Successful in 17s
Release / release-candidate (push) Successful in 4s
CI / build-frontend (push) Successful in 4s

This commit is contained in:
timmy 2026-08-07 17:16:30 +00:00
commit c0cc6ec46e
6 changed files with 227 additions and 24 deletions

View File

@ -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.';
}

View File

@ -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;

View File

@ -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))

View File

@ -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)

View File

@ -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,
}

View File

@ -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