mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge branch 'dev' of https://github.com/galaxyproject/galaxy into realtimetools
This commit is contained in:
@@ -1,4 +1,4 @@
|
||||
#!/bin/sh
|
||||
|
||||
cd "$(dirname "$0")"/../..
|
||||
python ./scripts/cleanup_datasets/cleanup_datasets.py -d 10 -3 -r "$@" >> ./scripts/cleanup_datasets/purge_datasets.log
|
||||
python ./scripts/cleanup_datasets/cleanup_datasets.py -d 60 -3 -r "$@" >> ./scripts/cleanup_datasets/purge_datasets.log
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
#!/bin/sh
|
||||
|
||||
cd "$(dirname "$0")"/../..
|
||||
python ./scripts/cleanup_datasets/cleanup_datasets.py -d 10 -5 -r "$@" >> ./scripts/cleanup_datasets/purge_folders.log
|
||||
python ./scripts/cleanup_datasets/cleanup_datasets.py -d 60 -5 -r "$@" >> ./scripts/cleanup_datasets/purge_folders.log
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
#!/bin/sh
|
||||
|
||||
cd "$(dirname "$0")"/../..
|
||||
python ./scripts/cleanup_datasets/cleanup_datasets.py -d 10 -2 -r "$@" >> ./scripts/cleanup_datasets/purge_histories.log
|
||||
python ./scripts/cleanup_datasets/cleanup_datasets.py -d 60 -2 -r "$@" >> ./scripts/cleanup_datasets/purge_histories.log
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
#!/bin/sh
|
||||
|
||||
cd "$(dirname "$0")"/../..
|
||||
python ./scripts/cleanup_datasets/cleanup_datasets.py -d 10 -4 -r "$@" >> ./scripts/cleanup_datasets/purge_libraries.log
|
||||
python ./scripts/cleanup_datasets/cleanup_datasets.py -d 60 -4 -r "$@" >> ./scripts/cleanup_datasets/purge_libraries.log
|
||||
|
||||
@@ -1,56 +0,0 @@
|
||||
#!/usr/bin/env python
|
||||
"""
|
||||
Updates dataset.size column.
|
||||
Remember to backup your database before running.
|
||||
|
||||
Deprecated - this doesn't work with modern Galaxy configurations options.
|
||||
"""
|
||||
from __future__ import print_function
|
||||
|
||||
import os
|
||||
import sys
|
||||
|
||||
from six.moves import configparser
|
||||
|
||||
sys.path.insert(1, os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir, 'lib')))
|
||||
|
||||
import galaxy.app
|
||||
|
||||
assert sys.version_info[:2] >= (2, 6)
|
||||
|
||||
|
||||
def usage(prog):
|
||||
print("usage: %s galaxy.ini" % prog)
|
||||
print("""
|
||||
Updates the dataset.size column. Users are advised to backup the database before
|
||||
running.
|
||||
""")
|
||||
|
||||
|
||||
def main():
|
||||
if len(sys.argv) != 2 or sys.argv[1] == "-h" or sys.argv[1] == "--help":
|
||||
usage(sys.argv[0])
|
||||
sys.exit()
|
||||
ini_file = sys.argv.pop(1)
|
||||
conf_parser = configparser.ConfigParser({'here': os.getcwd()})
|
||||
conf_parser.read(ini_file)
|
||||
configuration = {}
|
||||
for key, value in conf_parser.items("app:main"):
|
||||
configuration[key] = value
|
||||
app = galaxy.app.UniverseApplication(global_conf=ini_file, **configuration)
|
||||
|
||||
# Step through Datasets, determining size on disk for each.
|
||||
print("Determining the size of each dataset...")
|
||||
for row in app.model.Dataset.table.select().execute():
|
||||
purged = app.model.Dataset.get(row.id).purged
|
||||
file_size = app.model.Dataset.get(row.id).file_size
|
||||
if file_size is None and not purged:
|
||||
size_on_disk = app.model.Dataset.get(row.id).get_size()
|
||||
print("Updating Dataset.%d with file_size: %d" % (row.id, size_on_disk))
|
||||
app.model.Dataset.table.update(app.model.Dataset.table.c.id == row.id).execute(file_size=size_on_disk)
|
||||
app.shutdown()
|
||||
sys.exit(0)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -1,66 +0,0 @@
|
||||
#!/usr/bin/env python
|
||||
# Dan Blankenberg
|
||||
"""
|
||||
Updates metadata in the database to match rev 1891.
|
||||
|
||||
Remember to backup your database before running.
|
||||
|
||||
Deprecated - this doesn't work with modern Galaxy configurations options.
|
||||
"""
|
||||
from __future__ import print_function
|
||||
|
||||
import os
|
||||
import sys
|
||||
|
||||
from six.moves import configparser
|
||||
|
||||
sys.path.insert(1, os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir, 'lib')))
|
||||
|
||||
import galaxy.app
|
||||
import galaxy.datatypes.tabular
|
||||
|
||||
assert sys.version_info[:2] >= (2, 6)
|
||||
|
||||
|
||||
def usage(prog):
|
||||
print("usage: %s galaxy.ini" % prog)
|
||||
print("""
|
||||
Updates the metadata in the database to match rev 1981.
|
||||
|
||||
Remember to backup your database before running.
|
||||
""")
|
||||
|
||||
|
||||
def main():
|
||||
if len(sys.argv) != 2 or sys.argv[1] == "-h" or sys.argv[1] == "--help":
|
||||
usage(sys.argv[0])
|
||||
sys.exit()
|
||||
ini_file = sys.argv.pop(1)
|
||||
conf_parser = configparser.ConfigParser({'here': os.getcwd()})
|
||||
conf_parser.read(ini_file)
|
||||
configuration = {}
|
||||
for key, value in conf_parser.items("app:main"):
|
||||
configuration[key] = value
|
||||
app = galaxy.app.UniverseApplication(global_conf=ini_file, **configuration)
|
||||
|
||||
# Search out tabular datatypes (and subclasses) and initialize metadata
|
||||
print("Seeking out tabular based files and initializing metadata")
|
||||
for row in app.model.Dataset.table.select().execute():
|
||||
data = app.model.Dataset.get(row.id)
|
||||
if issubclass(type(data.datatype), type(app.datatypes_registry.get_datatype_by_extension('tabular'))):
|
||||
print(row.id, data.extension)
|
||||
# Call meta_data for all tabular files
|
||||
# special case interval type where we do not want to overwrite chr, start, end, etc assignments
|
||||
if issubclass(type(data.datatype), type(app.datatypes_registry.get_datatype_by_extension('interval'))):
|
||||
galaxy.datatypes.tabular.Tabular().set_meta(data)
|
||||
else:
|
||||
data.set_meta()
|
||||
app.model.context.add(data)
|
||||
app.model.context.flush()
|
||||
|
||||
app.shutdown()
|
||||
sys.exit(0)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -164,17 +164,11 @@ fi
|
||||
# activate virtualenv or conda env, sets $GALAXY_VIRTUAL_ENV and $GALAXY_CONDA_ENV
|
||||
setup_python
|
||||
|
||||
if [ $SET_VENV -eq 1 ] && [ -z "$VIRTUAL_ENV" ] && [ -z "$CONDA_DEFAULT_ENV" ]; then
|
||||
if [ $SET_VENV -eq 1 ] && [ -z "$VIRTUAL_ENV" ]; then
|
||||
echo "ERROR: A virtualenv cannot be found. Please create a virtualenv in $GALAXY_VIRTUAL_ENV, or activate one."
|
||||
exit 1
|
||||
fi
|
||||
|
||||
# this shouldn't happen, but check just in case
|
||||
if [ -z "$VIRTUAL_ENV" ] && [ "$CONDA_DEFAULT_ENV" = "base" ] || [ "$CONDA_DEFAULT_ENV" = "root" ]; then
|
||||
echo "ERROR: Conda is in 'base' environment, refusing to continue"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
: ${GALAXY_WHEELS_INDEX_URL:="https://wheels.galaxyproject.org/simple"}
|
||||
: ${PYPI_INDEX_URL:="https://pypi.python.org/simple"}
|
||||
: ${GALAXY_DEV_REQUIREMENTS:="./lib/galaxy/dependencies/dev-requirements.txt"}
|
||||
@@ -207,7 +201,7 @@ if [ $SKIP_CLIENT_BUILD -eq 0 ]; then
|
||||
# If git is not used and static/client_build_hash.txt is present, next
|
||||
# client rebuilds must be done manually by the admin
|
||||
if [ "$GIT_BRANCH" = "0" ]; then
|
||||
echo "Skipping Galaxy client build because git is not in use and the client build state cannot be compared against local changes. If you have made local modifications, then manual client builds will be required."
|
||||
echo "Skipping Galaxy client build because git is not in use and the client build state cannot be compared against local changes. If you have made local modifications, then manual client builds will be required. See ./client/README.md for more information."
|
||||
SKIP_CLIENT_BUILD=1
|
||||
else
|
||||
# Check if anything has changed in client/ since the last build
|
||||
@@ -234,7 +228,7 @@ if [ -n "$VIRTUAL_ENV" ]; then
|
||||
elif [ -n "$CONDA_DEFAULT_ENV" ] && [ -n "$CONDA_EXE" ]; then
|
||||
if ! in_conda_env "$(command -v node)"; then
|
||||
echo "Installing node into '$CONDA_DEFAULT_ENV' Conda environment with conda."
|
||||
$CONDA_EXE install --yes --override-channels --channel conda-forge --channel defaults --name "$CONDA_DEFAULT_ENV" node="$NODE_VERSION"
|
||||
$CONDA_EXE install --yes --override-channels --channel conda-forge --channel defaults --name "$CONDA_DEFAULT_ENV" nodejs="$NODE_VERSION"
|
||||
fi
|
||||
fi
|
||||
|
||||
|
||||
@@ -113,6 +113,10 @@ setup_python() {
|
||||
if [ "$CONDA_DEFAULT_ENV" != "$GALAXY_CONDA_ENV" ]; then
|
||||
conda_activate
|
||||
fi
|
||||
if [ "$CONDA_DEFAULT_ENV" = "base" ] || [ "$CONDA_DEFAULT_ENV" = "root" ]; then
|
||||
echo "ERROR: Conda is in 'base' environment, refusing to continue"
|
||||
exit 1
|
||||
fi
|
||||
fi
|
||||
fi
|
||||
|
||||
|
||||
@@ -15,6 +15,7 @@ database will be constructed.
|
||||
.. seealso: galaxy.ini, specifically the settings: database_connection and
|
||||
database file
|
||||
"""
|
||||
import logging
|
||||
import os.path
|
||||
import sys
|
||||
|
||||
@@ -25,6 +26,9 @@ from galaxy.model.orm.scripts import get_config
|
||||
from galaxy.model.tool_shed_install.migrate.check import create_or_verify_database as create_install_db
|
||||
from galaxy.webapps.tool_shed.model.migrate.check import create_or_verify_database as create_tool_shed_db
|
||||
|
||||
logging.basicConfig(level=logging.DEBUG)
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def invoke_create():
|
||||
config = get_config(sys.argv)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
""" This script parses Galaxy or Tool Shed config file for database connection
|
||||
and then delegates to sqlalchemy_migrate shell main function in
|
||||
migrate.versioning.shell. """
|
||||
import logging
|
||||
import os.path
|
||||
import sys
|
||||
|
||||
@@ -10,6 +11,9 @@ sys.path.insert(1, os.path.abspath(os.path.join(os.path.dirname(__file__), os.pa
|
||||
|
||||
from galaxy.model.orm.scripts import get_config
|
||||
|
||||
logging.basicConfig(level=logging.DEBUG)
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def invoke_migrate_main():
|
||||
# Migrate has its own args, so cannot use argparse
|
||||
|
||||
@@ -8,7 +8,6 @@ import sys
|
||||
sys.path.insert(1, os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, 'lib')))
|
||||
|
||||
import galaxy.config
|
||||
from galaxy.model.util import pgcalc
|
||||
from galaxy.objectstore import build_object_store_from_config
|
||||
from galaxy.util import nice_size
|
||||
from galaxy.util.script import app_properties_from_args, populate_config_args
|
||||
@@ -41,16 +40,15 @@ def quotacheck(sa_session, users, engine):
|
||||
sa_session.refresh(user)
|
||||
current = user.get_disk_usage()
|
||||
print(user.username, '<' + user.email + '>:', end=' ')
|
||||
if engine not in ('postgres', 'postgresql'):
|
||||
new = user.calculate_disk_usage()
|
||||
sa_session.refresh(user)
|
||||
# usage changed while calculating, do it again
|
||||
if user.get_disk_usage() != current:
|
||||
print('usage changed while calculating, trying again...')
|
||||
return quotacheck(sa_session, user, engine)
|
||||
|
||||
if not args.dryrun:
|
||||
# Apply new disk usage
|
||||
user.calculate_and_set_disk_usage()
|
||||
# And fetch
|
||||
new = user.get_disk_usage()
|
||||
else:
|
||||
new = pgcalc(sa_session, user.id, dryrun=args.dryrun)
|
||||
# yes, still a small race condition between here and the flush
|
||||
new = user.calculate_disk_usage()
|
||||
|
||||
print('old usage:', nice_size(current), 'change:', end=' ')
|
||||
if new in (current, None):
|
||||
print('none')
|
||||
@@ -59,10 +57,6 @@ def quotacheck(sa_session, users, engine):
|
||||
print('+%s' % (nice_size(new - current)))
|
||||
else:
|
||||
print('-%s' % (nice_size(current - new)))
|
||||
if not args.dryrun and engine not in ('postgres', 'postgresql'):
|
||||
user.set_disk_usage(new)
|
||||
sa_session.add(user)
|
||||
sa_session.flush()
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
|
||||
@@ -1,65 +0,0 @@
|
||||
"""Deprecated - only works with older ini files.
|
||||
"""
|
||||
from sys import argv
|
||||
|
||||
from six.moves.configparser import ConfigParser
|
||||
|
||||
REPLACE_PROPERTIES = ["file_path", "database_connection", "new_file_path"]
|
||||
MAIN_SECTION = "app:main"
|
||||
|
||||
|
||||
def sync():
|
||||
# Add or replace the relevant properites from galaxy.ini
|
||||
# into reports.ini
|
||||
reports_config_file = "config/reports.ini"
|
||||
if len(argv) > 1:
|
||||
reports_config_file = argv[1]
|
||||
|
||||
universe_config_file = "config/galaxy.ini"
|
||||
if len(argv) > 2:
|
||||
universe_config_file = argv[2]
|
||||
|
||||
parser = ConfigParser()
|
||||
parser.read(universe_config_file)
|
||||
|
||||
with open(reports_config_file, "r") as f:
|
||||
reports_config_lines = f.readlines()
|
||||
|
||||
replaced_properties = set([])
|
||||
with open(reports_config_file, "w") as f:
|
||||
# Write all properties from reports config replacing as
|
||||
# needed.
|
||||
for reports_config_line in reports_config_lines:
|
||||
(line, replaced_property) = get_synced_line(reports_config_line, parser)
|
||||
if replaced_property:
|
||||
replaced_properties.add(replaced_property)
|
||||
f.write(line)
|
||||
|
||||
# If any properties appear in universe config and not in
|
||||
# reports write these as well.
|
||||
for replacement_property in REPLACE_PROPERTIES:
|
||||
if parser.has_option(MAIN_SECTION, replacement_property) and \
|
||||
not (replacement_property in replaced_properties):
|
||||
f.write(get_universe_line(replacement_property, parser))
|
||||
|
||||
|
||||
def get_synced_line(reports_line, universe_config):
|
||||
# Cycle through properties to replace and perform replacement on
|
||||
# this line if needed.
|
||||
synced_line = reports_line
|
||||
replaced_property = None
|
||||
for replacement_property in REPLACE_PROPERTIES:
|
||||
if reports_line.startswith(replacement_property) and \
|
||||
universe_config.has_option(MAIN_SECTION, replacement_property):
|
||||
synced_line = get_universe_line(replacement_property, universe_config)
|
||||
replaced_property = replacement_property
|
||||
break
|
||||
return (synced_line, replaced_property)
|
||||
|
||||
|
||||
def get_universe_line(property_name, universe_config):
|
||||
return "%s=%s\n" % (property_name, universe_config.get(MAIN_SECTION, property_name))
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
sync()
|
||||
@@ -1,306 +0,0 @@
|
||||
#!/usr/bin/env python
|
||||
"""
|
||||
Downloads files to temp locations. This script is invoked by the Transfer
|
||||
Manager (galaxy.jobs.transfer_manager) and should not normally be invoked by
|
||||
hand.
|
||||
|
||||
This is deprecated - it only works with older ini configurations of Galaxy.
|
||||
"""
|
||||
import json
|
||||
import logging
|
||||
import optparse
|
||||
import os
|
||||
import random
|
||||
import sys
|
||||
import tempfile
|
||||
import threading
|
||||
import time
|
||||
|
||||
try:
|
||||
import pexpect
|
||||
except ImportError:
|
||||
pexpect = None
|
||||
|
||||
from daemon import DaemonContext
|
||||
from six.moves import (
|
||||
configparser,
|
||||
socketserver
|
||||
)
|
||||
from six.moves.urllib.error import URLError
|
||||
from six.moves.urllib.request import urlopen
|
||||
from sqlalchemy import create_engine, MetaData, Table
|
||||
from sqlalchemy.orm import scoped_session, sessionmaker
|
||||
|
||||
galaxy_root = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir))
|
||||
sys.path.insert(1, os.path.join(galaxy_root, 'lib'))
|
||||
|
||||
import galaxy.model
|
||||
from galaxy.util import bunch
|
||||
from galaxy.util.json import jsonrpc_response, validate_jsonrpc_request
|
||||
|
||||
PEXPECT_IMPORT_MESSAGE = ('The Python pexpect package is required to use this '
|
||||
'feature, please install it')
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
log.setLevel(logging.DEBUG)
|
||||
handler = logging.StreamHandler(sys.stdout)
|
||||
log.addHandler(handler)
|
||||
|
||||
debug = False
|
||||
slow = False
|
||||
|
||||
|
||||
class ArgHandler(object):
|
||||
"""
|
||||
Collect command line flags.
|
||||
"""
|
||||
def __init__(self):
|
||||
self.parser = optparse.OptionParser()
|
||||
self.parser.add_option('-c', '--config', dest='config', help='Path to Galaxy config file (config/galaxy.ini)',
|
||||
default=os.path.abspath(os.path.join(galaxy_root, 'config/galaxy.ini')))
|
||||
self.parser.add_option('-d', '--debug', action='store_true', dest='debug', help="Debug (don't detach)")
|
||||
self.parser.add_option('-s', '--slow', action='store_true', dest='slow', help="Transfer slowly (for debugging)")
|
||||
self.opts = None
|
||||
|
||||
def parse(self):
|
||||
self.opts, args = self.parser.parse_args()
|
||||
if len(args) != 1:
|
||||
log.error('usage: transfer.py <transfer job id>')
|
||||
sys.exit(1)
|
||||
try:
|
||||
self.transfer_job_id = int(args[0])
|
||||
except TypeError:
|
||||
log.error('The provided transfer job ID is not an integer: %s' % args[0])
|
||||
sys.exit(1)
|
||||
if self.opts.debug:
|
||||
global debug
|
||||
debug = True
|
||||
log.setLevel(logging.DEBUG)
|
||||
if self.opts.slow:
|
||||
global slow
|
||||
slow = True
|
||||
|
||||
|
||||
class GalaxyApp(object):
|
||||
"""
|
||||
A shell Galaxy App to provide access to the Galaxy configuration and
|
||||
model/database.
|
||||
"""
|
||||
def __init__(self, config_file):
|
||||
self.config = configparser.ConfigParser(dict(database_file='database/universe.sqlite',
|
||||
file_path='database/files',
|
||||
transfer_worker_port_range='12275-12675',
|
||||
transfer_worker_log=None))
|
||||
self.config.read(config_file)
|
||||
self.model = bunch.Bunch()
|
||||
self.connect_database()
|
||||
|
||||
def connect_database(self):
|
||||
# Avoid loading the entire model since doing so is exceptionally slow
|
||||
default_dburl = 'sqlite:///%s?isolation_level=IMMEDIATE' % self.config.get('app:main', 'database_file')
|
||||
try:
|
||||
dburl = self.config.get('app:main', 'database_connection')
|
||||
except configparser.NoOptionError:
|
||||
dburl = default_dburl
|
||||
engine = create_engine(dburl)
|
||||
metadata = MetaData(engine)
|
||||
self.sa_session = scoped_session(sessionmaker(bind=engine, autoflush=False, autocommit=True))
|
||||
self.model.TransferJob = galaxy.model.TransferJob
|
||||
self.model.TransferJob.table = Table("transfer_job", metadata, autoload=True)
|
||||
|
||||
def get_transfer_job(self, id):
|
||||
return self.sa_session.query(self.model.TransferJob).get(int(id))
|
||||
|
||||
|
||||
class ListenerServer(socketserver.ThreadingTCPServer):
|
||||
"""
|
||||
The listener will accept state requests and new transfers for as long as
|
||||
the manager is running.
|
||||
"""
|
||||
def __init__(self, port_range, RequestHandlerClass, app, transfer_job, state_result):
|
||||
self.state_result = state_result
|
||||
# Try random ports until a free one is found
|
||||
while True:
|
||||
random_port = random.choice(port_range)
|
||||
try:
|
||||
super(ListenerServer, self).__init__(('localhost', random_port), RequestHandlerClass)
|
||||
log.info('Listening on port %s' % random_port)
|
||||
break
|
||||
except Exception as e:
|
||||
log.warning('Tried binding port %s: %s' % (random_port, str(e)))
|
||||
transfer_job.socket = random_port
|
||||
app.sa_session.add(transfer_job)
|
||||
app.sa_session.flush()
|
||||
|
||||
|
||||
class ListenerRequestHandler(socketserver.BaseRequestHandler):
|
||||
"""
|
||||
Handle state or transfer requests received on the socket.
|
||||
"""
|
||||
def handle(self):
|
||||
request = self.request.recv(8192)
|
||||
response = {}
|
||||
valid, request, response = validate_jsonrpc_request(request, ('get_state', ), ())
|
||||
if valid:
|
||||
self.request.send(json.dumps(jsonrpc_response(request=request, result=self.server.state_result.result)))
|
||||
else:
|
||||
error_msg = 'Unable to serve request: %s' % response['error']['message']
|
||||
if 'data' in response['error']:
|
||||
error_msg += ': %s' % response['error']['data']
|
||||
log.error(error_msg)
|
||||
log.debug('Original request was: %s' % request)
|
||||
|
||||
|
||||
class StateResult(object):
|
||||
"""
|
||||
A mutable container for the 'result' portion of JSON-RPC responses to state requests.
|
||||
"""
|
||||
def __init__(self, result=None):
|
||||
self.result = result
|
||||
|
||||
|
||||
def transfer(app, transfer_job_id):
|
||||
transfer_job = app.get_transfer_job(transfer_job_id)
|
||||
if transfer_job is None:
|
||||
log.error('Invalid transfer job ID: %s' % transfer_job_id)
|
||||
return False
|
||||
port_range = app.config.get('app:main', 'transfer_worker_port_range')
|
||||
try:
|
||||
port_range = [int(p) for p in port_range.split('-')]
|
||||
except Exception as e:
|
||||
log.error('Invalid port range set in transfer_worker_port_range: %s: %s' % (port_range, str(e)))
|
||||
return False
|
||||
protocol = transfer_job.params['protocol']
|
||||
if protocol not in ('http', 'https', 'scp'):
|
||||
log.error('Unsupported protocol: %s' % protocol)
|
||||
return False
|
||||
state_result = StateResult(result=dict(state=transfer_job.states.RUNNING, info='Transfer process starting up.'))
|
||||
listener_server = ListenerServer(range(port_range[0], port_range[1] + 1), ListenerRequestHandler, app, transfer_job, state_result)
|
||||
# daemonize here (if desired)
|
||||
if not debug:
|
||||
daemon_context = DaemonContext(files_preserve=[listener_server.fileno()], working_directory=os.getcwd())
|
||||
daemon_context.open()
|
||||
# If this fails, it'll never be detected. Hopefully it won't fail since it succeeded once.
|
||||
app.connect_database() # daemon closed the database fd
|
||||
transfer_job = app.get_transfer_job(transfer_job_id)
|
||||
listener_thread = threading.Thread(target=listener_server.serve_forever)
|
||||
listener_thread.setDaemon(True)
|
||||
listener_thread.start()
|
||||
# Store this process' pid so unhandled deaths can be handled by the restarter
|
||||
transfer_job.pid = os.getpid()
|
||||
app.sa_session.add(transfer_job)
|
||||
app.sa_session.flush()
|
||||
terminal_state = None
|
||||
if protocol in ['http', 'https']:
|
||||
for transfer_result_dict in http_transfer(transfer_job):
|
||||
state_result.result = transfer_result_dict
|
||||
if transfer_result_dict['state'] in transfer_job.terminal_states:
|
||||
terminal_state = transfer_result_dict
|
||||
elif protocol in ['scp']:
|
||||
# Transfer the file using scp
|
||||
transfer_result_dict = scp_transfer(transfer_job)
|
||||
# Handle the state of the transfer
|
||||
state = transfer_result_dict['state']
|
||||
state_result.result = transfer_result_dict
|
||||
if state in transfer_job.terminal_states:
|
||||
terminal_state = transfer_result_dict
|
||||
if terminal_state is not None:
|
||||
transfer_job.state = terminal_state['state']
|
||||
for name in ['info', 'path']:
|
||||
if name in terminal_state:
|
||||
transfer_job.__setattr__(name, terminal_state[name])
|
||||
else:
|
||||
transfer_job.state = transfer_job.states.ERROR
|
||||
transfer_job.info = 'Unknown error encountered by transfer worker.'
|
||||
app.sa_session.add(transfer_job)
|
||||
app.sa_session.flush()
|
||||
return True
|
||||
|
||||
|
||||
def http_transfer(transfer_job):
|
||||
"""Plugin" for handling http(s) transfers."""
|
||||
url = transfer_job.params['url']
|
||||
assert url.startswith('http://') or url.startswith('https://')
|
||||
try:
|
||||
f = urlopen(url)
|
||||
except URLError as e:
|
||||
yield dict(state=transfer_job.states.ERROR, info='Unable to open URL: %s' % str(e))
|
||||
return
|
||||
size = f.info().getheader('Content-Length')
|
||||
if size is not None:
|
||||
size = int(size)
|
||||
chunksize = 1024 * 1024
|
||||
if slow:
|
||||
chunksize = 1024
|
||||
read = 0
|
||||
last = 0
|
||||
try:
|
||||
fh, fn = tempfile.mkstemp()
|
||||
except Exception as e:
|
||||
yield dict(state=transfer_job.states.ERROR, info='Unable to create temporary file for transfer: %s' % str(e))
|
||||
return
|
||||
log.debug('Writing %s to %s, size is %s' % (url, fn, size or 'unknown'))
|
||||
try:
|
||||
while True:
|
||||
chunk = f.read(chunksize)
|
||||
if not chunk:
|
||||
break
|
||||
os.write(fh, chunk)
|
||||
read += chunksize
|
||||
if size is not None and read < size:
|
||||
percent = int(float(read) / size * 100)
|
||||
if percent != last:
|
||||
yield dict(state=transfer_job.states.PROGRESS, read=read, percent='%s' % percent)
|
||||
last = percent
|
||||
elif size is None:
|
||||
yield dict(state=transfer_job.states.PROGRESS, read=read)
|
||||
if slow:
|
||||
time.sleep(1)
|
||||
os.close(fh)
|
||||
yield dict(state=transfer_job.states.DONE, path=fn)
|
||||
except Exception as e:
|
||||
yield dict(state=transfer_job.states.ERROR, info='Error during file transfer: %s' % str(e))
|
||||
return
|
||||
return
|
||||
|
||||
|
||||
def scp_transfer(transfer_job):
|
||||
"""Plugin" for handling scp transfers using pexpect"""
|
||||
def print_ticks(d):
|
||||
pass
|
||||
host = transfer_job.params['host']
|
||||
user_name = transfer_job.params['user_name']
|
||||
password = transfer_job.params['password']
|
||||
file_path = transfer_job.params['file_path']
|
||||
if pexpect is None:
|
||||
return dict(state=transfer_job.states.ERROR, info=PEXPECT_IMPORT_MESSAGE)
|
||||
try:
|
||||
fh, fn = tempfile.mkstemp()
|
||||
except Exception as e:
|
||||
return dict(state=transfer_job.states.ERROR, info='Unable to create temporary file for transfer: %s' % str(e))
|
||||
try:
|
||||
# TODO: add the ability to determine progress of the copy here like we do in the http_transfer above.
|
||||
cmd = "scp %s@%s:'%s' '%s'" % (user_name,
|
||||
host,
|
||||
file_path.replace(' ', r'\ '),
|
||||
fn)
|
||||
pexpect.run(cmd, events={'.ssword:*': password + '\r\n',
|
||||
pexpect.TIMEOUT: print_ticks},
|
||||
timeout=10)
|
||||
return dict(state=transfer_job.states.DONE, path=fn)
|
||||
except Exception as e:
|
||||
return dict(state=transfer_job.states.ERROR, info='Error during file transfer: %s' % str(e))
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
arg_handler = ArgHandler()
|
||||
arg_handler.parse()
|
||||
app = GalaxyApp(arg_handler.opts.config)
|
||||
|
||||
log.debug('Initiating transfer...')
|
||||
if transfer(app, arg_handler.transfer_job_id):
|
||||
log.debug('Finished')
|
||||
else:
|
||||
log.error('Error in transfer process...')
|
||||
sys.exit(1)
|
||||
sys.exit(0)
|
||||
Reference in New Issue
Block a user