fix(roi): continue syncing accessible repositories

This commit is contained in:
moe-berri 2026-09-30 12:37:30 -07:00
parent 3dc7848c0b
commit 85f558fbb8
6 changed files with 156 additions and 11 deletions

Binary file not shown.

After

Width:  |  Height:  |  Size: 57 KiB

View file

@ -81,6 +81,7 @@ def _summarize_person(
key: str,
spend: tuple[ROISpendRecord, ...],
pulls: tuple[tuple[ROIPullRecord, str, str], ...],
complete_scope: bool,
) -> ROIPersonSummary:
spend_rows: Final = tuple(
row for row in spend if _person_key(normalize_email(row["email"]), "gateway:" + row["user_id"]) == key
@ -111,11 +112,14 @@ def _summarize_person(
pending_prs=pending_count,
match_methods=methods,
eligible=eligible,
cost_per_hour=spend_total / hours if eligible and hours > 0 and spend_total is not None else None,
cost_per_hour=spend_total / hours
if complete_scope and eligible and hours > 0 and spend_total is not None
else None,
)
def summarize(report: ROIReport, mappings: Mapping[str, str]) -> ROISummary:
complete_scope: Final = not report.get("unavailable_repos", ())
observed: Final = frozenset(
normalized for normalized in (normalize_email(row["email"]) for row in report["spend"]) if normalized
)
@ -134,6 +138,7 @@ def summarize(report: ROIReport, mappings: Mapping[str, str]) -> ROISummary:
key,
report["spend"],
matched_pulls,
complete_scope,
)
for key in sorted(people_keys)
)
@ -181,8 +186,8 @@ def summarize(report: ROIReport, mappings: Mapping[str, str]) -> ROISummary:
total_spend=total_spend,
total_output_hours=total_output_hours,
excluded_spend=max(0.0, total_spend - matched_spend),
cost_per_hour=matched_spend / output_hours if output_hours else None,
hours_per_dollar=output_hours / matched_spend if matched_spend else None,
cost_per_hour=matched_spend / output_hours if complete_scope and output_hours else None,
hours_per_dollar=output_hours / matched_spend if complete_scope and matched_spend else None,
merged_prs=len(pull_summaries),
estimated_prs=sum(person["estimated_prs"] for person in people),
matched_prs=sum(pull["matched"] for pull in pull_summaries),

View file

@ -245,6 +245,54 @@ class _ProcessedPull(NamedTuple):
metadata_unavailable: bool = False
class _RepositoryPulls(NamedTuple):
repo: str
pulls: tuple[GitHubPullListItem, ...]
unavailable: bool = False
class _RepositoryBatch(NamedTuple):
queue: tuple[tuple[str, GitHubPullListItem], ...]
unavailable_repos: tuple[str, ...]
warnings: tuple[str, ...]
stage: str
async def _read_repository(github: GitHub, repo: str, start: date, end: date) -> _RepositoryPulls:
try:
return _RepositoryPulls(repo, await github.pulls(repo, start, end))
except SourceError:
return _RepositoryPulls(repo, (), unavailable=True)
async def _read_repositories(github: GitHub, repos: tuple[str, ...], start: date, end: date) -> _RepositoryBatch:
groups: Final = await asyncio.gather(*(_read_repository(github, repo, start, end) for repo in repos))
unavailable: Final = tuple(group.repo for group in groups if group.unavailable)
if len(unavailable) == len(repos):
raise SourceError(
"GitHub could not read any selected repository. No new report was published; "
"check repository access or try analysis again later."
)
queue: Final = tuple(chain.from_iterable(((group.repo, pull) for pull in group.pulls) for group in groups))
warnings: Final = (
(
(
f"Incomplete report: could not read {', '.join(unavailable)}. "
"Results include only accessible repositories. Spend-per-hour figures are unavailable until "
"all selected repositories can be read. Check repository access or run analysis again to retry."
),
)
if unavailable
else ()
)
return _RepositoryBatch(
queue,
unavailable,
warnings,
"Analysis complete with unavailable repositories" if unavailable else "Analysis complete",
)
def _processed_records(processed: tuple[_ProcessedPull, ...]) -> Mapping[int, ROIPullRecord]:
if processed and all(item.metadata_unavailable for item in processed):
raise SourceError(
@ -384,12 +432,8 @@ class SyncManager:
start: Final = end - timedelta(days=settings.backfill_days - 1)
spend: Final = await spend_reader(start, end)
self._update_status(phase="repositories", stage="Reading configured repositories")
pull_groups: Final = await asyncio.gather(*(github.pulls(repo, start, end) for repo in settings.repos))
queue: Final = tuple(
chain.from_iterable(
((repo, pull) for pull in pulls) for repo, pulls in zip(settings.repos, pull_groups, strict=True)
)
)
repositories: Final = await _read_repositories(github, settings.repos, start, end)
queue: Final = repositories.queue
context: Final = cache_context(settings, estimator_models)
previous: Final = await self._previous_report(repository)
previous_pulls: Final[Mapping[str, ROIPullRecord]] = MappingProxyType(
@ -500,7 +544,8 @@ class SyncManager:
spend=spend,
pulls=tuple(processed_by_index[index] for index in range(len(queue))),
settings_fingerprint=settings_fingerprint(settings),
warnings=(),
warnings=repositories.warnings,
unavailable_repos=repositories.unavailable_repos,
)
await github.close()
report_json: Final[Mapping[str, object]] = _JSON_OBJECT_ADAPTER.validate_python(
@ -514,7 +559,7 @@ class SyncManager:
{
"running": False,
"phase": "complete",
"stage": "Analysis complete",
"stage": repositories.stage,
"finished_at": self._clock().isoformat(),
}
)

View file

@ -205,6 +205,7 @@ class ROIReport(TypedDict):
pulls: ReadOnly[tuple[ROIPullRecord, ...]]
settings_fingerprint: ReadOnly[str]
warnings: NotRequired[ReadOnly[tuple[str, ...]]]
unavailable_repos: NotRequired[ReadOnly[tuple[str, ...]]]
id: NotRequired[ReadOnly[str]]

View file

@ -9,6 +9,7 @@ import httpx
import pytest
from pydantic import TypeAdapter
from litellm.proxy.roi_calculator.analytics import summarize
from litellm.proxy.roi_calculator.estimator import CompletionCaller
from litellm.proxy.roi_calculator.github import GitHubPullListItem
from litellm.proxy.roi_calculator.sync import SpendReader, SyncManager, read_spend
@ -457,3 +458,72 @@ async def test_one_unreadable_pr_preserves_other_estimates_in_report() -> None:
assert manager.status.phase == "complete"
assert manager.status.estimated == 1
assert manager.status.needs_attention == 1
def _repository_outage_transport(status: int, *, all_unavailable: bool = False) -> httpx.MockTransport:
baseline: Final = _transport()
def respond(request: httpx.Request) -> httpx.Response:
if request.url.path == "/repos/org/unavailable/pulls":
return httpx.Response(status, json=[] if status == 200 else {"message": "Repository unavailable"})
if all_unavailable and request.url.path.endswith("/pulls"):
return httpx.Response(status)
return baseline.handle_request(request)
return httpx.MockTransport(respond)
@pytest.mark.asyncio
@pytest.mark.parametrize("status", (403, 404, 429))
async def test_unavailable_repository_publishes_flagged_partial_report_and_recovers(status: int) -> None:
repository: Final = _ReportRepository()
manager: Final = SyncManager(clock=_fixed_now)
settings: Final = _settings().model_copy(update=MappingProxyType({"repos": ("org/repo", "org/unavailable")}))
assert await manager.start(
settings, repository, _spend_reader(), _completion(), _repository_outage_transport(status)
)
await _wait_until_finished(manager)
report: Final = TypeAdapter(ROIReport).validate_python(repository.values["roi_calculator_report"])
summary: Final = summarize(report, MappingProxyType({}))
assert manager.status.phase == "complete"
assert manager.status.estimated == 1
assert report["unavailable_repos"] == ("org/unavailable",)
assert "Incomplete report" in report["warnings"][0] and "org/unavailable" in report["warnings"][0]
assert report["pulls"][0]["estimate"]["status"] == "estimated"
assert summary["metrics"]["total_output_hours"] == 4
assert summary["metrics"]["cost_per_hour"] is None
assert summary["metrics"]["hours_per_dollar"] is None
assert all(person["cost_per_hour"] is None for person in summary["people"])
async def unexpected_completion(request: ROICompletionRequest) -> object:
raise AssertionError("The healthy repository's estimate must be reused after recovery")
assert await manager.start(
settings, repository, _spend_reader(), unexpected_completion, _repository_outage_transport(200)
)
await _wait_until_finished(manager)
recovered: Final = TypeAdapter(ROIReport).validate_python(repository.values["roi_calculator_report"])
assert recovered["unavailable_repos"] == ()
assert recovered["warnings"] == ()
assert manager.status.reused == 1
assert summarize(recovered, MappingProxyType({}))["metrics"]["cost_per_hour"] == 3
@pytest.mark.asyncio
async def test_all_repository_outage_preserves_previous_report() -> None:
repository: Final = _ReportRepository()
manager: Final = SyncManager(clock=_fixed_now)
settings: Final = _settings().model_copy(update=MappingProxyType({"repos": ("org/repo", "org/unavailable")}))
assert await manager.start(settings, repository, _spend_reader(), _completion(), _repository_outage_transport(200))
await _wait_until_finished(manager)
previous: Final = repository.values["roi_calculator_report"]
assert await manager.start(
settings, repository, _spend_reader(), _completion(), _repository_outage_transport(403, all_unavailable=True)
)
await _wait_until_finished(manager)
assert manager.status.phase == "error"
assert manager.status.error is not None and "No new report was published" in manager.status.error
assert repository.values["roi_calculator_report"] == previous

View file

@ -157,6 +157,30 @@ describe("ROICalculatorView", () => {
);
});
it("shows incomplete repository results without a spend-per-hour figure", async () => {
const warning = "Incomplete report: could not read org/unavailable. Spend-per-hour figures are unavailable.";
vi.mocked(apiClient.get).mockImplementation((path: string) => {
if (path === "/roi-calculator/settings") return Promise.resolve(settings);
if (path === "/roi-calculator/report") {
return Promise.resolve({
report: {
...summary,
warnings: [warning],
metrics: { ...summary.metrics, cost_per_hour: null, hours_per_dollar: null },
people: summary.people.map((person) => ({ ...person, cost_per_hour: null })),
},
});
}
return Promise.resolve(idleStatus);
});
render(<ROICalculatorView accessToken="token" />);
expect(await screen.findByRole("alert")).toHaveTextContent(warning);
expect(screen.getByRole("button", { name: "Open estimate for org/repo pull request 42" })).toBeInTheDocument();
expect(screen.queryByText("$3.00")).not.toBeInTheDocument();
});
it("lets a view-only admin read the report without write controls", async () => {
const runningStatus = {
...idleStatus,