Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
248 changes: 227 additions & 21 deletions TAP/tap.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down
2 changes: 1 addition & 1 deletion TAP/vositables.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.'
Expand Down
4 changes: 4 additions & 0 deletions requirements-test.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
9 changes: 8 additions & 1 deletion tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"),
)
Expand Down Expand Up @@ -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:
Expand Down
11 changes: 8 additions & 3 deletions tests/fixtures/TAP.conf.template
Original file line number Diff line number Diff line change
Expand Up @@ -5,16 +5,21 @@
#
# 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

[WEB]
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


Expand Down
Loading
Loading