diff --git a/TAP/tap.py b/TAP/tap.py index 3442d23..61213c1 100755 --- a/TAP/tap.py +++ b/TAP/tap.py @@ -1280,24 +1280,42 @@ def __run__(self, **kwargs): # } end async bogus value case # - if self.debug: - logging.debug('') - logging.debug ('call writeStatusMsg') - logging.debug (f'statuspath= {self.statuspath:s}') + if (self.param['phase'] == 'RUN'): + # + # Publish the job, answer the client, and keep running + # the query in a detached child. + # - self.__writeStatusMsg__(self.statuspath, self.statdict, - self.param) + self.__respondAsyncAndDetach__() - if self.debug: - logging.debug('') - logging.debug ('call printAsyncResponse') + if self.debug: + logging.debug('') + logging.debug('async response sent; job running detached') - self.__printAsyncResponse__(self.statusurl) + else: - if self.debug: - logging.debug('') - logging.debug('returned printAsyncResponse') + # + # ABORT and unrecognized phases are terminal: the branch + # above has already decided the phase, so write it, + # answer the client and stop. Falling through to the + # query path below would run the job the client just + # aborted and overwrite ABORTED (or ERROR) with + # COMPLETED. + # + + if self.debug: + logging.debug('') + logging.debug ('terminal async phase= ' + f"{self.statdict['phase']:s}") + logging.debug (f'statuspath= {self.statuspath:s}') + + self.__writeStatusMsg__(self.statuspath, self.statdict, + self.param) + + self.__printAsyncResponse__(self.statusurl) + + sys.exit() # # } end async submit case @@ -2726,25 +2744,213 @@ def __printSyncResponse__(self, status, msg, resulturl, format, **kwargs): def __printAsyncResponse__(self, statusurl, **kwargs): # - # async: return statusurl and kill the parent process + # async: point the client at the job's status URL. + # + # The body is one line, but it still needs a Content-Length. This + # is an nph- script, so nothing downstream supplies one, and a + # reverse proxy in front of the CGI has no other way to tell + # where the response ends. Earlier versions omitted it and let + # the connection dying stand in for the end of the message. # - print("HTTP/1.1 303 See Other\r") - print("Location: %s\r\n\r" % statusurl) - print("Redirect Location: %s" % statusurl) + body = 'Redirect Location: %s\n' % statusurl + + sys.stdout.write('HTTP/1.1 303 See Other\r\n') + sys.stdout.write('Location: %s\r\n' % statusurl) + sys.stdout.write('Content-Type: text/plain\r\n') + sys.stdout.write('Content-Length: %d\r\n' + % len(body.encode('utf-8'))) + sys.stdout.write('Connection: close\r\n') + sys.stdout.write('\r\n') + sys.stdout.write(body) sys.stdout.flush() - time.sleep(2.0) + if self.debug: + logging.debug('') + logging.debug(f'async response sent: statusurl= {statusurl:s}') + + return + + def __respondAsyncAndDetach__(self, **kwargs): + + # + # { An async submit has to finish an HTTP response now and keep + # executing the query afterwards, and under CGI those two pull + # against each other: the web server completes the response + # when the script closes stdout, and mod_cgi terminates the + # script once the request is cleaned up. The process that runs + # the query can therefore be neither the one holding stdout nor + # part of the request. + # + # So fork. The parent names the child as the job's runId, + # writes the status document, sends the 303 and exits, which + # ends the response the way any other CGI would. The child + # leaves the request's process group, points its standard + # streams away from the server pipe, waits for the parent to + # confirm the job is published, and returns to run the query. # - # Shut down parent program + # Earlier versions sent the response and then SIGKILLed + # os.getppid(). Under CGI that parent is the web server child + # serving the request: killing it truncates the response, and + # behind a reverse proxy it leaves the proxy holding a dead + # upstream connection, which the proxy reports as 503 on the + # next request or two routed over it. # - os.kill(os.getppid(), signal.SIGKILL) + try: + readfd, writefd = os.pipe() + pid = os.fork() + + except OSError as e: + + # + # No fork available: answer the request from this process and + # run the job here. The response is still well formed; the + # job now lives and dies with the request. + # + + logging.error(f'Could not fork async worker: {str(e)}') + + self.__writeStatusMsg__(self.statuspath, self.statdict, + self.param) + + self.__printAsyncResponse__(self.statusurl) + + return + + if (pid > 0): + # + # { parent: publish the job, answer the client, exit + # + os.close(readfd) + + self.statdict['process_id'] = pid + + published = False + + try: + self.__writeStatusMsg__(self.statuspath, self.statdict, + self.param) + + published = True + + self.__printAsyncResponse__(self.statusurl) + + except Exception as e: + + # + # Most likely the client or the proxy hung up while the + # response was going out. Nothing can be sent about it + # now; the job's own fate is decided in the finally + # clause below. + # + + logging.error(f'Async response failed: {str(e)}') + + finally: + + # + # The child blocks on this byte, so a fast query cannot + # overwrite the status document written above. Once the + # job is published the child has to be released even if + # the response itself failed: a client that hung up does + # not make the job go away, and a child that exits here + # would leave the job stuck in EXECUTING forever. If the + # status write is what failed, the byte is withheld and + # the child exits, because there is no job document for + # it to update. + # + + if published: + try: + os.write(writefd, b'1') + + except OSError as e: + logging.error( + f'Could not release async worker: {str(e)}') + + os.close(writefd) + + sys.exit(0) + # + # } end parent + # + + # + # { child: detach from the request, then run the query + # + os.close(writefd) + + self.pid = os.getpid() + self.statdict['process_id'] = self.pid + + self.__detachFromServer__() + + published = b'' + + try: + published = os.read(readfd, 1) + + except OSError as e: + logging.error(f'Async worker handshake failed: {str(e)}') + + os.close(readfd) + + if (len(published) == 0): + + # + # The parent exited before it published the job, so there is + # no status document for this run to update. + # + + logging.error('Async parent exited before publishing the job') + + os._exit(1) if self.debug: logging.debug('') - logging.debug('parent process killed') + logging.debug(f'async worker detached: pid= {self.pid:d}') + + return + # + # } end child + # + + + def __detachFromServer__(self, **kwargs): + + # + # Give up the web server's request context: leave the request's + # process group so the server cannot reap this process along with + # the request, and replace the inherited stdio with /dev/null so + # the response the parent just sent is neither held open by nor + # corrupted from here. + # + + try: + os.setsid() + + except OSError as e: + logging.error(f'setsid failed in async worker: {str(e)}') + + try: + sys.stdout.flush() + sys.stderr.flush() + + except (OSError, ValueError) as e: + logging.error(f'Could not flush async worker stdio: {str(e)}') + + devnull = os.open(os.devnull, os.O_RDWR) + + try: + os.dup2(devnull, 0) + os.dup2(devnull, 1) + os.dup2(devnull, 2) + + finally: + if (devnull > 2): + os.close(devnull) return diff --git a/TAP/vositables.py b/TAP/vositables.py index d414081..2b15001 100644 --- a/TAP/vositables.py +++ b/TAP/vositables.py @@ -211,7 +211,7 @@ def __init__(self, **kwargs): # # { Connect to DBMS # - if('connectInfo' in kwargs): + if('connectInfo' in kwargs): self.connectInfo = kwargs['connectInfo'] else: self.msg = 'Required connectInfo dict is missing.' diff --git a/requirements-test.txt b/requirements-test.txt index 343732f..66043b2 100644 --- a/requirements-test.txt +++ b/requirements-test.txt @@ -10,6 +10,10 @@ configobj sqlparse astropy beautifulsoup4 +# tap.py asks BeautifulSoup for the 'lxml' tree builder by name when it +# reads a job status document (TAP/tap.py:904, 2162), so bs4 alone is not +# enough for the async flow. +lxml xmltodict # pyneid is not currently used — tests/test_pyneid_compat.py is skipped diff --git a/tests/conftest.py b/tests/conftest.py index 6aa0cf6..e0ab97c 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -66,7 +66,8 @@ def tap_conf(fixture_root: Path) -> Path: template = (FIXTURES / "TAP.conf.template").read_text() conf = template.format( TEST_WORKDIR=str(fixture_root / "workdir"), - TEST_HTTP_URL=f"http://127.0.0.1:{port}", + TEST_HTTP_HOST="127.0.0.1", + TEST_HTTP_PORT=str(port), TEST_DB_PATH=str(fixture_root / "test_data.db"), TEST_TAP_SCHEMA=str(fixture_root / "tap_schema.db"), ) @@ -147,6 +148,12 @@ def _run_cgi(self): ) stdout, stderr = proc.communicate(input=body) + # A CGI that dies has only stderr to say so, and swallowing it + # makes a 500 from the script indistinguishable from a 500 the + # script meant to send. Surface it on the test runner's stderr. + if stderr: + sys.stderr.write(stderr.decode("utf-8", "replace")) + # Forward raw subprocess stdout to the socket. The script is # responsible for emitting a valid HTTP status line + headers. try: diff --git a/tests/fixtures/TAP.conf.template b/tests/fixtures/TAP.conf.template index e0280d2..6c2c4a9 100644 --- a/tests/fixtures/TAP.conf.template +++ b/tests/fixtures/TAP.conf.template @@ -5,7 +5,8 @@ # # Placeholders: # {TEST_WORKDIR} — working dir for TAP request artifacts -# {TEST_HTTP_URL} — http://localhost:NNNN (fixture-allocated) +# {TEST_HTTP_HOST} — 127.0.0.1 +# {TEST_HTTP_PORT} — fixture-allocated port # {TEST_DB_PATH} — absolute path to test_data.db # {TEST_TAP_SCHEMA} — absolute path to tap_schema.db @@ -13,8 +14,12 @@ TAP_WORKDIR={TEST_WORKDIR} TAP_WORKURL=/workspace - HTTP_URL={TEST_HTTP_URL} - HTTP_PORT=0 + # configparam appends HTTP_PORT to HTTP_URL for any port other than + # 80/443, and that composed value is what the async status and result + # URLs are built from. Keep the two apart here or the job URLs come + # back with the port twice. + HTTP_URL=http://{TEST_HTTP_HOST} + HTTP_PORT={TEST_HTTP_PORT} CGI_PGM=TAP diff --git a/tests/test_async_flow.py b/tests/test_async_flow.py new file mode 100644 index 0000000..b8e2201 --- /dev/null +++ b/tests/test_async_flow.py @@ -0,0 +1,212 @@ +"""End-to-end coverage of the UWS async flow. + +An async TAP job takes three round trips, the same sequence pyvo and +pyNEID drive: + +1. POST to ``/async`` with the query. The server creates the job in + PENDING and answers 303 with ``Location: ``. +2. POST ``PHASE=RUN`` to ``/phase``. The server starts the + job and answers 303 to the same status URL. +3. GET ```` until the phase reaches COMPLETED, then read the + result named in the status document. + +None of this could be covered before the detach fix. ``tap.py`` ended +every async response by SIGKILLing ``os.getppid()`` — under this fixture +that is the process running the test suite, and in production it is the +web server child serving the request. Step 2 therefore either hung (the +CGI never closed the stdout pipe the server was reading) or took the +test runner down with it. +""" +from __future__ import annotations + +import time +import xml.etree.ElementTree as ET +from pathlib import Path + +import requests +from _helpers import tap_async_url + +UWS = {"uws": "http://www.ivoa.net/xml/UWS/v1.0", + "xlink": "http://www.w3.org/1999/xlink"} + +QUERY = "select * from data_l0" + + +def _submit(server: str) -> str: + """Round trip 1. Returns the job's status URL.""" + resp = requests.post( + tap_async_url(server), + data={"query": QUERY, "format": "ipac"}, + allow_redirects=False, + timeout=30, + ) + assert resp.status_code == 303, resp.text[:500] + statusurl = resp.headers["Location"] + assert statusurl, "303 carried no Location header" + return statusurl + + +def _status(statusurl: str) -> ET.Element: + resp = requests.get(statusurl, timeout=30) + assert resp.status_code == 200, resp.text[:500] + return ET.fromstring(resp.text) + + +def _phase(statusurl: str) -> str: + return _status(statusurl).find("uws:phase", UWS).text + + +def test_async_submit_creates_pending_job(tap_server): + statusurl = _submit(tap_server) + assert _phase(statusurl) == "PENDING" + + +def test_async_run_starts_the_job(tap_server): + """Round trip 2 answers promptly and does not take the server with it.""" + statusurl = _submit(tap_server) + + resp = requests.post( + f"{statusurl}/phase", + data={"PHASE": "RUN"}, + allow_redirects=False, + timeout=30, + ) + + assert resp.status_code == 303, resp.text[:500] + assert resp.headers["Location"] == statusurl + + +def test_async_run_response_is_self_delimiting(tap_server): + """The 303 must be framed by its headers, not by the connection dying. + + A reverse proxy in front of the CGI needs a Content-Length (or a + chunked body) to relay the response. The old code emitted a body + with neither and relied on the connection being destroyed to mark + the end of the message. + """ + statusurl = _submit(tap_server) + + resp = requests.post( + f"{statusurl}/phase", + data={"PHASE": "RUN"}, + allow_redirects=False, + timeout=30, + ) + + assert "Content-Length" in resp.headers + assert int(resp.headers["Content-Length"]) == len(resp.content) + + +def test_async_job_runs_to_completion(tap_server, fixture_root: Path): + statusurl = _submit(tap_server) + requests.post( + f"{statusurl}/phase", + data={"PHASE": "RUN"}, + allow_redirects=False, + timeout=30, + ) + + deadline = time.time() + 30 + phase = "" + while time.time() < deadline: + phase = _phase(statusurl) + if phase in ("COMPLETED", "ERROR", "ABORTED"): + break + time.sleep(0.2) + + assert phase == "COMPLETED", f"job ended in phase {phase}" + + job = _status(statusurl) + href = job.find("uws:results/uws:result", UWS).get( + "{http://www.w3.org/1999/xlink}href") + assert href + + # The fixture server does not serve TAP_WORKURL (Apache does that in + # production), so read the result out of the workspace on disk. + # tap.py puts each job under /TAP/. + workspace = job.find("uws:jobId", UWS).text + results = list((fixture_root / "workdir" / "TAP" / workspace).glob("*")) + assert results, f"no result artifacts in workspace {workspace}" + written = [p for p in results if p.name != "status.xml" and p.stat().st_size] + assert written, f"job completed but wrote no result file: {results}" + + +def _post_phase(statusurl: str, phase: str): + return requests.post( + f"{statusurl}/phase", + data={"PHASE": phase}, + allow_redirects=False, + timeout=30, + ) + + +def _results(fixture_root: Path, workspace: str) -> list: + workdir = fixture_root / "workdir" / "TAP" / workspace + return [p for p in workdir.glob("*") if p.name != "status.xml"] + + +def test_async_abort_is_terminal(tap_server, fixture_root: Path): + """ABORT must end the job, not run it. + + Every phase shares one response path at the end of the async block, + and that path used to fall through into the query. So a POST of + PHASE=ABORT wrote ABORTED, answered the client, and then ran the + job anyway, overwriting the phase the client had just been told. + """ + statusurl = _submit(tap_server) + assert _phase(statusurl) == "PENDING" + + resp = _post_phase(statusurl, "ABORT") + assert resp.status_code == 303, resp.text[:500] + + job = _status(statusurl) + assert job.find("uws:phase", UWS).text == "ABORTED" + + workspace = job.find("uws:jobId", UWS).text + time.sleep(1.0) + assert _phase(statusurl) == "ABORTED", "the aborted job kept running" + assert not _results(fixture_root, workspace), \ + "the aborted job wrote a result" + + +def test_async_unknown_phase_is_terminal(tap_server, fixture_root: Path): + """An unrecognized phase reports ERROR and runs nothing.""" + statusurl = _submit(tap_server) + + resp = _post_phase(statusurl, "SPIN") + assert resp.status_code == 303, resp.text[:500] + + job = _status(statusurl) + assert job.find("uws:phase", UWS).text == "ERROR" + + workspace = job.find("uws:jobId", UWS).text + time.sleep(1.0) + assert _phase(statusurl) == "ERROR", "the rejected job kept running" + assert not _results(fixture_root, workspace), \ + "the rejected job wrote a result" + + +def test_async_detach_does_not_signal_the_web_server(): + """Regression guard against the SIGKILL-the-parent detach. + + Asserted on the source rather than on behavior on purpose: the + behavioral assertion is "the process running this suite is still + alive", which a suite that has been SIGKILLed cannot make. + """ + import tokenize + + from TAP import tap + + # Comments and strings are stripped before the check so the code can + # still explain in prose why it no longer signals its parent. + with open(tap.__file__, "rb") as fp: + src = "".join( + tok.string + for tok in tokenize.tokenize(fp.readline) + if tok.type not in (tokenize.COMMENT, tokenize.STRING)) + + assert "getppid" not in src, ( + "tap.py signals its parent process. Under CGI the parent is the " + "web server child serving the request; killing it drops the " + "response and, behind a proxy, poisons the upstream connection." + ) diff --git a/tests/test_pyneid_compat.py b/tests/test_pyneid_compat.py index e927e0a..883a3b7 100644 --- a/tests/test_pyneid_compat.py +++ b/tests/test_pyneid_compat.py @@ -1,54 +1,40 @@ """pyNEID compatibility regression tests. -Currently all skipped — the test fixture's HTTP handler doesn't yet -support TAP's full async flow end-to-end. +Still skipped, but for a different reason than before. -Why these matter: pyNEID's `query_adql` always uses the async TAP -endpoint, which works in three round-trips: +pyNEID's ``query_adql`` always uses the async TAP endpoint, which works +in three round trips: -1. POST to /TAP/async with the ADQL query → server returns 303 See Other +1. POST to /TAP/async with the ADQL query -> server returns 303 See Other with a Location: header. -2. POST to with PHASE=RUN → server starts the job, returns - another 303 to the same status URL. +2. POST to /phase with PHASE=RUN -> server starts the job and + returns another 303 to the same status URL. 3. GET the status URL repeatedly until phase=COMPLETED, then GET the - results URL embedded in the status XML. - -Steps 2–3 depend on the server persisting per-request workspace state -(status.xml, results) under TAP_WORKDIR keyed by a workspace ID -embedded in the URL path. The current `_NphCGIHandler` in -tests/conftest.py routes by URL prefix but doesn't model the workspace -lifecycle — each request is a fresh CGI subprocess that doesn't know -about prior request state. Production runs under Apache, which manages -the same workdir mount across requests so the state naturally -persists, but reproducing that fidelity in CI needs: - -- Sticky working directory across requests (already mostly there since - we set `os.chdir(fixture_root)` once at fixture setup) -- Correct PATH_INFO routing for the second and subsequent calls in the - async flow (the Location URL pyNEID gets back may not match the URL - pattern the handler routes today) -- Possibly more debug — when this was tried locally, pyNEID hung at - step 2 with no observable error in the TAP debug log - -Equivalent direct-HTTP coverage of the same code paths exists in -tests/test_security_e2e.py and tests/test_http_endpoint.py — those -exercise the sync endpoint, which uses a single request and doesn't -need the workspace dance. The pyNEID layer adds value as a regression -check on top of those (it's the exact wrapper production traffic hits), -but is not blocking. - -To wire this up later: trace the async POST/poll loop, fix the handler -routing, then drop the skip decorator. + result URL embedded in the status XML. + +This module used to say the fixture handler could not carry that +sequence, and recorded that "pyNEID hung at step 2 with no observable +error in the TAP debug log". The hang was not the fixture. Step 2 ended +by SIGKILLing os.getppid(), which under the fixture is the test runner +and in production is the web server child serving the request, so the +client was left waiting on a pipe nobody would close. That is fixed, and +the sequence above now has direct coverage in tests/test_async_flow.py. + +What is still missing here is pyNEID itself: it is not a test dependency +and is not installed in CI, so there is nothing to import. Adding it +means weighing a heavier test dependency (and its own network defaults) +against coverage that tests/test_async_flow.py already provides at the +HTTP level. Until that call is made, this module stays skipped. + +To wire it up later: add pyNEID to requirements-test.txt, write the +round-trip test against the fixture server, and drop the skip below. """ from __future__ import annotations import pytest -# Once the async fixture support lands, drop this module-level skip and -# re-add the actual test bodies (a snapshot of the intended tests is in -# the git history of this file or in the design plan). pytest.skip( - "Async TAP flow not yet supported by the test fixture handler. " - "See module docstring for what's needed.", + "pyNEID is not a test dependency; the async flow it drives is covered " + "directly in tests/test_async_flow.py. See module docstring.", allow_module_level=True, )