File size: 4,742 Bytes
cc7ed3f
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
5bb077e
 
cc7ed3f
 
5bb077e
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
cc7ed3f
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
135
"""
Shared EBI MSA REST client, parameterized by tool base URL.

Verified live (2026-08-03) for clustalo / muscle / kalign / mafft / tcoffee:
all five expose the identical contract — POST ``{base}/run`` with
``email``/``stype``/``sequence``, poll ``{base}/status/{job}`` until FINISHED,
then fetch ``result/{job}/fa`` (FASTA alignment) and ``result/{job}/phylotree``
(Newick). The base URL is passed in full (per tool), never as a query param.
"""

from __future__ import annotations

import asyncio
import logging

import httpx

logger = logging.getLogger(__name__)

EBI_TOOLS: dict[str, str] = {
    "clustalo": "https://www.ebi.ac.uk/Tools/services/rest/clustalo",
    "muscle": "https://www.ebi.ac.uk/Tools/services/rest/muscle",
    "kalign": "https://www.ebi.ac.uk/Tools/services/rest/kalign",
    "mafft": "https://www.ebi.ac.uk/Tools/services/rest/mafft",
    "tcoffee": "https://www.ebi.ac.uk/Tools/services/rest/tcoffee",
}

POLL_INTERVAL = 1
MAX_POLLS = 60
TREE_TYPES = ["phylotree"]

# Tools tried in order when the primary EBI job fails (all share the same REST contract).
EBI_FALLBACK_ORDER = ["clustalo", "muscle"]


async def run_ebi_msa_best_effort(
    sequence: str,
    stype: str = "protein",
    email: str = "bioflow@example.com",
    tools: list[str] | None = None,
) -> dict:
    """Run EBI MSA trying each tool in ``tools`` until one succeeds.

    Raises ValueError only when every tool fails, with a combined message.
    Returns the same shape as :func:`run_ebi_msa` (``method`` names the tool
    that actually produced the alignment).
    """
    tools = tools or EBI_FALLBACK_ORDER
    errors: list[str] = []
    for tool in tools:
        base_url = EBI_TOOLS.get(tool)
        if not base_url:
            errors.append(f"unknown tool {tool!r}")
            continue
        try:
            result = await run_ebi_msa(
                base_url=base_url,
                sequence=sequence,
                stype=stype,
                email=email,
            )
            result["method"] = tool
            return result
        except Exception as e:
            errors.append(f"{tool}: {e}")
            logger.warning("EBI MSA tool %s failed (%s); trying next…", tool, e)
    raise ValueError("EBI MSA failed on all tools: " + " | ".join(errors))


async def run_ebi_msa(
    base_url: str,
    sequence: str,
    stype: str = "protein",
    email: str = "bioflow@example.com",
) -> dict:
    """Submit ``sequence`` (FASTA) to an EBI MSA tool and wait for the alignment.

    Raises ValueError with a human-readable message on any failure.
    Returns ``{"job_id", "aln_fasta", "phylotree", "method"}``.
    """
    method = base_url.rstrip("/").rsplit("/", 1)[-1]

    async with httpx.AsyncClient(timeout=30) as client:
        submit_resp = await client.post(
            f"{base_url}/run",
            data={"email": email, "stype": stype, "sequence": sequence},
            headers={"Accept": "text/plain"},
        )
        if submit_resp.status_code != 200:
            detail = submit_resp.text[:200] if submit_resp.text else "no response body"
            raise ValueError(f"EBI submission failed (HTTP {submit_resp.status_code}): {detail}")
        job_id = submit_resp.text.strip()
        logger.info("EBI MSA job submitted (%s): %s", method, job_id)

        for _ in range(MAX_POLLS):
            await asyncio.sleep(POLL_INTERVAL)
            try:
                status_resp = await client.get(f"{base_url}/status/{job_id}")
                status = status_resp.text.strip()
            except Exception as e:
                logger.warning("EBI status poll failed: %s", e)
                continue
            logger.info("EBI MSA status (%s/%s): %s", method, job_id, status)
            if status == "FINISHED":
                break
            if status == "ERROR":
                raise ValueError(f"EBI {method} job failed")
        else:
            raise ValueError(f"EBI {method} alignment timed out")

        await asyncio.sleep(1)

        fa_resp = await client.get(
            f"{base_url}/result/{job_id}/fa", headers={"Accept": "text/plain"}
        )
        if fa_resp.status_code != 200 or not fa_resp.text.strip():
            raise ValueError("Failed to fetch alignment result from EBI")
        aln_fasta = fa_resp.text

        phylotree = ""
        for t in TREE_TYPES:
            tr = await client.get(
                f"{base_url}/result/{job_id}/{t}", headers={"Accept": "text/plain"}
            )
            if tr.status_code == 200 and tr.text.strip():
                phylotree = tr.text
                break

    return {
        "job_id": job_id,
        "aln_fasta": aln_fasta,
        "phylotree": phylotree,
        "method": method,
    }