Spaces:
Running
Running
File size: 4,764 Bytes
e4ecac1 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 | """The de novo GO-term prediction can also be an unfinished InterProScan job.
``function_hints`` runs inside the ``uniprot`` annotation bundle, so its
background poller has to merge the finished prediction into the bundle of
composition + GO terms instead of replacing the whole step payload the way the
domains poller does. These tests pin the merge, the hold-open contract, and the
upstream-failure verdict.
"""
import asyncio
import pytest
from app.routers import pipeline_v2 as pv
@pytest.fixture(autouse=True)
def _clean_registry():
pv._INTERPRO_PENDING.clear()
pv._INTERPRO_DEFERRED.clear()
pv._INTERPRO_POLL_TASKS.clear()
yield
pv._INTERPRO_PENDING.clear()
pv._INTERPRO_DEFERRED.clear()
pv._INTERPRO_POLL_TASKS.clear()
def _install_job(job_id: str, data: dict):
pv._jobs[job_id] = {
"status": "running",
"steps": {"uniprot": {"status": "running", "progress": 50, "data": data}},
"context": {},
}
def _no_mirror(monkeypatch):
monkeypatch.setattr(pv, "_mirror", lambda job_id: None)
monkeypatch.setattr(pv, "_persist_v2_final", lambda *a, **k: None)
async def _noop_finalize(job_id, context):
pass
monkeypatch.setattr(pv, "_finalize_context", _noop_finalize)
def _mark(job_id: str):
return lambda step, status, **kw: pv._set_step_status(job_id, step, status, **kw)
def test_function_hints_poll_merges_into_annotation_bundle(monkeypatch):
"""The finished prediction lands next to composition, not in its place."""
job_id = "aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa"
_install_job(job_id, {"_de_novo": True, "composition": {"aa": 1}})
_no_mirror(monkeypatch)
prediction = {
"status": "inferred",
"go_terms": [{"go_id": "GO:0003674", "name": "molecular function", "namespace": "MF"}],
"source": "interpro2go",
}
async def _ok(_job_id):
return {"status": "complete", "source": "interproscan6", "domains": []}
monkeypatch.setattr("app.services.de_novo.await_interpro_result", _ok)
monkeypatch.setattr(
"app.tools.function_predict.prediction_from_result",
lambda result, seq, pdb_id="de_novo": prediction,
)
context: dict = {}
async def drive():
pv._schedule_function_hints_poll(
"ipr_job_1", _mark(job_id), job_id, context, "ACDEFGHIKLM", None
)
assert pv._interpro_pending(job_id) == 1
pv._complete_job(job_id, context, None)
assert pv._get_job(job_id)["status"] == "running", (
"parent must stay open while the GO prediction is finishing"
)
await asyncio.gather(*list(pv._INTERPRO_POLL_TASKS))
asyncio.run(drive())
bundle = pv._get_job(job_id)["steps"]["uniprot"]["data"]
assert bundle["composition"] == {"aa": 1}, "composition must survive the merge"
assert bundle["function_hints"] is prediction
assert pv._get_job(job_id)["steps"]["uniprot"]["status"] == "complete"
assert context["uniprot"]["function_hints"] is prediction
assert pv._get_job(job_id)["status"] == "complete", "parent released after the last poll"
assert job_id not in pv._INTERPRO_PENDING
def test_function_hints_upstream_failure_is_evidence_unavailable(monkeypatch):
"""A scan that fails upstream must read as unavailable, never 'no GO terms'."""
job_id = "bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb"
_install_job(job_id, {"_de_novo": True, "composition": {"aa": 1}})
_no_mirror(monkeypatch)
async def _boom(_job_id):
return {"status": "failed", "source": "interproscan6",
"error": "InterProScan job failed upstream: ERROR"}
monkeypatch.setattr("app.services.de_novo.await_interpro_result", _boom)
monkeypatch.setattr(
"app.tools.function_predict.prediction_from_result",
lambda result, seq, pdb_id="de_novo": {
"status": "evidence_unavailable",
"go_terms": [],
"source": "interpro2go",
"retrieval_failures": [result["error"]],
},
)
context: dict = {}
async def drive():
pv._schedule_function_hints_poll(
"ipr_job_2", _mark(job_id), job_id, context, "ACDEFGHIKLM", None
)
pv._complete_job(job_id, context, None)
await asyncio.gather(*list(pv._INTERPRO_POLL_TASKS))
asyncio.run(drive())
fh = pv._get_job(job_id)["steps"]["uniprot"]["data"]["function_hints"]
assert fh["status"] == "evidence_unavailable"
assert "failed upstream" in fh["retrieval_failures"][0]
assert pv._get_job(job_id)["steps"]["uniprot"]["status"] == "complete", (
"an upstream failure is a verdict, not a step crash: the bundle is done"
)
assert pv._get_job(job_id)["status"] == "complete"
|