From b31669d4e286c4cbd92a17d76dc924d4dc0286e8 Mon Sep 17 00:00:00 2001 From: Adam Wright Date: Sat, 3 Oct 2026 22:50:55 +0000 Subject: [PATCH] GSA correctness fixes from review (area 2) - The first label given is group 1. It was the alphabetical first, so the baseline -- and every Up and Down -- depended on spelling ('WT, KO' compared against KO). The result now says which groups were compared, in which order. (Which way ReactomeGSA's sign points is not stated in its docs, so it is not claimed.) - An R write.table header, one cell short of the rows, kept its first sample; header[1:] dropped it and misaligned every label. - A CSV upload is converted to tab-separated before it is submitted; the service reads only tabs, and a CSV went through unconverted. - Files are decoded the same way at validation and at submit: UTF-8, UTF-16, or Windows-1252. A 1252 file passed validation, was labelled, then failed every retry at submit. - Validation is bounded (1,000 samples; blank lines count towards the read limit) and, with the matrix read at submit, runs off the event loop. - The wait retries a failed poll or result fetch up to four times; one 502 used to lose a finished run. A failure now gives the analysis id. Checked end to end in a browser against a recording stand-in for ReactomeGSA: the melanoma TSV and the same matrix as CSV both submit tab-separated with 16 samples and group1=MOCK, and finish. Co-Authored-By: Claude Opus 5.5 --- src/gsa/chainlit_flow.py | 14 ++++-- src/gsa/chat.py | 17 ++++++- src/gsa/job.py | 44 +++++++++++++++-- src/gsa/upload.py | 94 +++++++++++++++++++++++++++++------- tests/gsa/test_gsa_chat.py | 18 +++++++ tests/gsa/test_gsa_job.py | 44 ++++++++++++++++- tests/gsa/test_gsa_upload.py | 48 ++++++++++++++++++ 7 files changed, 251 insertions(+), 28 deletions(-) diff --git a/src/gsa/chainlit_flow.py b/src/gsa/chainlit_flow.py index b5aaff0..fab7247 100644 --- a/src/gsa/chainlit_flow.py +++ b/src/gsa/chainlit_flow.py @@ -13,6 +13,7 @@ summary at all. """ +import asyncio from pathlib import Path from typing import Any, Literal, Protocol, TypedDict @@ -150,7 +151,9 @@ async def _run( client: GsaClient | None, ) -> None: try: - matrix = validate(path) + # Off the event loop: reading and checking a 20 MB file stalled every + # session for seconds (review, area 2). + matrix = await asyncio.to_thread(validate, path) except UploadRejectedError as refusal: await send(str(refusal)) return @@ -223,7 +226,12 @@ async def _run( return except Exception: logger.exception("gsa analysis failed", extra={"analysis": analysis_id}) - await send("Something went wrong while waiting for the analysis.") + # With the id: the analysis may well have finished, and the upload is + # already deleted, so this is the only way back to it (review, 2). + await send( + "Something went wrong while waiting for the analysis. It may still " + f"finish on Reactome's side, as analysis `{analysis_id}`." + ) return # `finished.for_model` is deliberately not used here. @@ -243,6 +251,6 @@ async def _run( # way to tell "sent and not shown" from "never sent". This line is the # difference. logger.info("gsa result sending", extra={"analysis": finished.analysis_id}) - await send(chat.describe_result(finished)) + await send(chat.describe_result(finished, grouping)) await send_file(finished.table_path) logger.info("gsa result delivered", extra={"analysis": finished.analysis_id}) diff --git a/src/gsa/chat.py b/src/gsa/chat.py index dc7fbd2..893ae89 100644 --- a/src/gsa/chat.py +++ b/src/gsa/chat.py @@ -114,7 +114,11 @@ def parse_grouping(reply: str, sample_count: int) -> Grouping: ) canonical = [seen[label.casefold()] for label in labels] - group1, group2 = sorted(seen.values()) + # The first label to appear is the reference. It was the alphabetical + # first, so which group was the baseline -- and so every Up and Down -- + # depended on spelling: "WT, KO" compared against KO (review, area 2). + # People list the control first, as the measured MOCK/MCM run did. + group1, group2 = list(seen.values()) return Grouping(labels=canonical, group1=group1, group2=group2) @@ -132,7 +136,7 @@ def describe_progress(status: AnalysisStatus) -> str: return f"Running the analysis — {detail}" -def describe_result(finished: Finished) -> str: +def describe_result(finished: Finished, grouping: Grouping | None = None) -> str: """What the person reads. Never a prompt. The Pathway Browser link is included *here* and not in anything the @@ -148,6 +152,15 @@ def describe_result(finished: Finished) -> str: f"**{significant:,} of {total:,} pathways** are significant at FDR < 0.05.", "", ] + if top and grouping is not None: + # The table's Direction was never said to be relative to anything. + # Only what is certain is stated: which groups, in which order. + # ReactomeGSA's own docs do not say which way its sign points. + lines += [ + f"Compared: **{escape(grouping.group1)}** (group 1, the first " + f"label you gave) with **{escape(grouping.group2)}** (group 2).", + "", + ] if top: lines += [ "| Pathway | Direction | FDR |", diff --git a/src/gsa/job.py b/src/gsa/job.py index 60baef2..d08c4db 100644 --- a/src/gsa/job.py +++ b/src/gsa/job.py @@ -12,18 +12,23 @@ promises nothing, and `await_result` is where success or failure is decided. """ +import asyncio import contextlib import time from collections.abc import Awaitable, Callable from dataclasses import dataclass from pathlib import Path -from typing import Any +from typing import Any, TypeVar + +import httpx from gsa import results as gsa_results -from gsa.client import AnalysisStatus, GsaClient, GsaError +from gsa.client import AnalysisStatus, GsaClient, GsaError, GsaNotReadyError from gsa.upload import Matrix, discard from util.logging import logging +T = TypeVar("T") + logger = logging.getLogger(__name__) #: A PADOG run with 1,000 permutations took minutes in the measured run. @@ -170,7 +175,7 @@ async def submit_uploaded_matrix( method=method, dataset_name="uploaded", dataset_type=dataset_type, - matrix=matrix.text, + matrix=await asyncio.to_thread(lambda: matrix.text), samples=matrix.samples, analysis_group=analysis_group, group1=group1, @@ -254,6 +259,35 @@ def prune_results( return removed +#: Consecutive failures tolerated on one request before giving up, and the +#: pause between tries. A run takes up to 30 minutes and ~150 polls; with no +#: tolerance, one 502 or timeout anywhere lost a finished analysis the reader +#: could not get back (review, area 2: ~17% of 30-minute runs at a 0.1% +#: per-request failure rate). +TRANSIENT_ATTEMPTS = 4 +TRANSIENT_PAUSE_SECONDS = 5.0 + + +async def _patiently(call: Callable[[], Awaitable[T]]) -> T: + """Retry a request through brief outages; a real refusal still raises.""" + for attempt in range(1, TRANSIENT_ATTEMPTS + 1): + try: + return await call() + except (GsaNotReadyError, AnalysisFailedError): + raise + except (GsaError, httpx.HTTPError) as exc: + if attempt == TRANSIENT_ATTEMPTS: + raise + logger.warning( + "gsa request failed (%s); retrying, attempt %d of %d", + type(exc).__name__, + attempt, + TRANSIENT_ATTEMPTS, + ) + await _sleep(TRANSIENT_PAUSE_SECONDS * attempt) + raise AssertionError("unreachable") + + async def await_result( client: GsaClient, analysis_id: str, @@ -275,7 +309,7 @@ async def await_result( polls = 0 while True: - status = await client.analysis_status(analysis_id) + status = await _patiently(lambda: client.analysis_status(analysis_id)) polls += 1 if on_progress is not None: await on_progress(status) @@ -294,7 +328,7 @@ async def await_result( ) await _sleep(poll_interval) - parsed = gsa_results.parse(await client.result(analysis_id)) + parsed = gsa_results.parse(await _patiently(lambda: client.result(analysis_id))) if not parsed.pathways: # Complete, and yet nothing to report. Better to say so than to # hand back an empty file and an exact-sounding zero. diff --git a/src/gsa/upload.py b/src/gsa/upload.py index c5b87fe..87f3235 100644 --- a/src/gsa/upload.py +++ b/src/gsa/upload.py @@ -14,6 +14,8 @@ """ import contextlib +import csv +import io import os from dataclasses import dataclass from pathlib import Path @@ -46,11 +48,24 @@ class Matrix: #: the first thing anyone does with a gene count is show it to someone. #: "unknown" is a fact; "-1 genes" is a bug wearing a number. gene_count: int | None + #: How the upload separates cells; `text` always hands the service tabs. + delimiter: str = "\t" @property def text(self) -> str: - """The matrix itself. Never logged, never shown, never prompted.""" - return self.path.read_text() + """The matrix, tab-separated. Never logged, never shown, never prompted. + + Decoded the way `validate` decoded it -- a strict read here let a + Windows-1252 file pass validation and then fail at submit, every + retry (review, area 2) -- and tab-separated whatever was uploaded: + the service reads only tabs, and a CSV went through unconverted. + A blocking read: call it with `asyncio.to_thread`. + """ + text = decode(self.path.read_bytes()) + if self.delimiter == "\t": + return text + rows = csv.reader(io.StringIO(text)) + return "\n".join("\t".join(cell.strip() for cell in row) for row in rows if row) def max_upload_bytes() -> int: @@ -64,6 +79,29 @@ def max_upload_bytes() -> int: return value if value > 0 else DEFAULT_MAX_UPLOAD_BYTES +#: More samples than any expression study here; a bound on what one upload +#: can make the server hold (a 20 MB one-line header was 7 million names). +MAX_SAMPLES = 1000 +#: Lines read before deciding it is a matrix, blank ones included: blank +#: lines used to skip the early exit, and 21 million of them stalled every +#: session for two seconds. +MAX_LINES_READ = 5000 + + +def decode(raw: bytes) -> str: + """The file's text: UTF-8 (with or without BOM), UTF-16, or Windows-1252. + + Excel on Windows writes 1252 or UTF-16; reading those as UTF-8 with + replacement showed the reader garbled sample names to label. + """ + if raw.startswith((b"\xff\xfe", b"\xfe\xff")): + return raw.decode("utf-16") + try: + return raw.decode("utf-8-sig") + except UnicodeDecodeError: + return raw.decode("cp1252", errors="replace") + + def _split(line: str) -> list[str]: # Tab first: it is what the service wants and what every export # produces. Comma only if there is no tab at all, because a TSV cell can @@ -90,22 +128,39 @@ def validate(path: Path) -> Matrix: raise UploadRejectedError("That file is empty.") header: list[str] = [] + first_row: list[str] = [] rows = 0 counted = True - with path.open(encoding="utf-8", errors="replace") as handle: - for index, line in enumerate(handle): - if not line.strip(): - continue - if not header: - header = _split(line.rstrip("\n")) - continue - rows += 1 - if index > 5000 and rows > MIN_DATA_ROWS: - # Enough to know it is a matrix. Counting every gene of a - # 20,000-row file to answer "is this a matrix" is work - # nobody asked for. - counted = False - break + delimiter = "\t" + # Iterated, not split: 21 million blank lines as a list is the problem + # this loop's bound exists to avoid. + for index, raw_line in enumerate(io.StringIO(decode(path.read_bytes()))): + line = raw_line.rstrip("\r\n") + if index > MAX_LINES_READ: + # Enough to know whether it is a matrix. Counting every gene of + # a 20,000-row file is work nobody asked for -- and a file of + # blank lines must not make us read all of it. + if rows <= MIN_DATA_ROWS: + raise UploadRejectedError( + "That file is mostly empty lines. It needs a header row " + "naming the samples, then one row per gene." + ) + counted = False + break + if not line.strip(): + continue + if not header: + delimiter = "\t" if "\t" in line else "," + header = _split(line) + if len(header) > MAX_SAMPLES + 1: + raise UploadRejectedError( + f"That file has {len(header) - 1:,} columns. This handles up " + f"to {MAX_SAMPLES:,} samples." + ) + continue + if not first_row: + first_row = _split(line) + rows += 1 if len(header) < MIN_COLUMNS: raise UploadRejectedError( @@ -118,7 +173,11 @@ def validate(path: Path) -> Matrix: # The first header cell labels the gene column and is often blank -- # the measured example's header starts with a tab. - samples = [name.strip() for name in header[1:] if name.strip()] + # R's write.table writes no cell for the gene column, so the header is + # one short of the rows; taking header[1:] then dropped the first + # sample, and its label would have gone to the wrong column (review, 2). + names = header if len(first_row) == len(header) + 1 else header[1:] + samples = [name.strip() for name in names if name.strip()] if len(samples) < 2: raise UploadRejectedError( "There is only one sample in that file. A gene set analysis " @@ -130,6 +189,7 @@ def validate(path: Path) -> Matrix: size_bytes=size, samples=samples, gene_count=rows if counted else None, + delimiter=delimiter, ) diff --git a/tests/gsa/test_gsa_chat.py b/tests/gsa/test_gsa_chat.py index d31bb56..22b3e3c 100644 --- a/tests/gsa/test_gsa_chat.py +++ b/tests/gsa/test_gsa_chat.py @@ -289,3 +289,21 @@ def test_the_matrix_only_reply_offers_no_gene_list() -> None: assert "Attach" in chat.HOW_TO_RUN_GSA_WITH_A_MATRIX assert "list of genes" not in chat.HOW_TO_RUN_GSA_WITH_A_MATRIX assert "list of genes" in chat.HOW_TO_RUN_GSA + + +@pytest.mark.parametrize( + ("reply", "reference", "other"), + [ + ("MOCK, MOCK, MCM, MCM", "MOCK", "MCM"), + ("WT, KO, WT, KO", "WT", "KO"), + ("control, control, Treated, Treated", "control", "Treated"), + ("healthy, disease", "healthy", "disease"), + ], +) +def test_the_first_label_given_is_group_one( + reply: str, reference: str, other: str +) -> None: + # It was the alphabetical first, so the baseline depended on spelling + # (review, area 2). + grouping = chat.parse_grouping(reply, len(reply.split(","))) + assert (grouping.group1, grouping.group2) == (reference, other) diff --git a/tests/gsa/test_gsa_job.py b/tests/gsa/test_gsa_job.py index a8b0cbb..080d2fb 100644 --- a/tests/gsa/test_gsa_job.py +++ b/tests/gsa/test_gsa_job.py @@ -17,7 +17,7 @@ import pytest from gsa import job -from gsa.client import AnalysisStatus, DatasetSummary, LoadingStatus +from gsa.client import AnalysisStatus, DatasetSummary, GsaError, LoadingStatus from gsa.upload import Matrix FIXTURE = Path(__file__).parent / "result_fixture.json" @@ -472,3 +472,45 @@ async def test_writing_a_result_prunes_the_directory_first(tmp_path: Path) -> No assert not stale.exists() assert finished.table_path.exists() + + +# --- review, area 2 ----------------------------------------------------------- + + +class Flaky(StubClient): + """Fails the first `failures` status polls, then behaves.""" + + def __init__(self, failures: int, **kwargs: Any) -> None: + super().__init__(**kwargs) + self.failures = failures + + async def analysis_status(self, analysis_id: str) -> AnalysisStatus: + if self.failures: + self.failures -= 1 + raise GsaError("GET /status returned 502") + return await super().analysis_status(analysis_id) + + +@asyncio_test +async def test_a_brief_outage_while_waiting_does_not_lose_the_run( + tmp_path: Path, +) -> None: + # One 502 anywhere in a 30-minute wait used to end it, and the upload was + # already deleted: the finished result was out of reach. + finished = await job.await_result( + # A fixed count: one derived from TRANSIENT_ATTEMPTS followed a sabotage. + Flaky(2), # type: ignore[arg-type] + "an-1", + out_dir=tmp_path, + ) + assert finished.analysis_id == "an-1" + + +@asyncio_test +async def test_a_lasting_outage_still_ends_the_wait(tmp_path: Path) -> None: + with pytest.raises(GsaError): + await job.await_result( + Flaky(job.TRANSIENT_ATTEMPTS), # type: ignore[arg-type] + "an-1", + out_dir=tmp_path, + ) diff --git a/tests/gsa/test_gsa_upload.py b/tests/gsa/test_gsa_upload.py index bd701e3..d540817 100644 --- a/tests/gsa/test_gsa_upload.py +++ b/tests/gsa/test_gsa_upload.py @@ -120,3 +120,51 @@ def test_a_small_file_still_reports_a_real_count(tmp_path: Path) -> None: # The control: if every file reported None, the test above would pass # and the count would be useless. assert upload.validate(write(tmp_path, GOOD)).gene_count == 2 + + +# --- review, area 2 ----------------------------------------------------------- + + +def test_an_r_style_header_keeps_its_first_sample(tmp_path: Path) -> None: + # write.table writes no cell for the gene column; header[1:] dropped S1. + text = "S1\tS2\tS3\tS4\nTP53\t10\t12\t30\t33\nMDM2\t1\t2\t3\t4\nBAX\t5\t6\t7\t8\n" + assert upload.validate(write(tmp_path, text)).samples == ["S1", "S2", "S3", "S4"] + + +def test_a_csv_reaches_the_service_tab_separated(tmp_path: Path) -> None: + text = 'gene,ctrl_1,"treated, day 1"\nTP53,1,2\nMDM2,3,4\nBAX,5,6\n' + matrix = upload.validate(write(tmp_path, text, "m.csv")) + assert matrix.text.splitlines()[0] == "gene\tctrl_1\ttreated, day 1" + assert "," not in matrix.text.splitlines()[1] + + +@pytest.mark.parametrize("encoding", ["cp1252", "utf-16"]) +def test_windows_encodings_read_the_same_at_validation_and_at_submit( + tmp_path: Path, encoding: str +) -> None: + # Validation decoded with replacement, submission strictly: such a file + # passed, was labelled, then failed every retry. + text = "\tcontrôle_1\tcontrôle_2\tbehandelt\nTP53\t1\t2\t3\nMDM2\t4\t5\t6\nBAX\t7\t8\t9\n" + path = tmp_path / "m.tsv" + path.write_bytes(text.encode(encoding)) + matrix = upload.validate(path) + assert matrix.samples == ["contrôle_1", "contrôle_2", "behandelt"] + assert "contrôle_1" in matrix.text + + +def test_too_many_columns_is_refused(tmp_path: Path) -> None: + header = "\t" + "\t".join(f"s{i}" for i in range(upload.MAX_SAMPLES + 5)) + with pytest.raises(upload.UploadRejectedError, match="columns"): + upload.validate( + write(tmp_path, header + "\nTP53" + "\t1" * (upload.MAX_SAMPLES + 5) + "\n") + ) + + +def test_blank_lines_cannot_make_validation_read_the_whole_file(tmp_path: Path) -> None: + import time + + text = "\tS1\tS2\n" + "\n" * 3_000_000 + "TP53\t1\t2\n" + started = time.perf_counter() + with pytest.raises(upload.UploadRejectedError): + upload.validate(write(tmp_path, text)) + assert time.perf_counter() - started < 1.0