from .arvcontainer import ArvadosContainer, RunnerContainer
from .arvjob import ArvadosJob, RunnerJob, RunnerTemplate
from .arvtool import ArvadosCommandTool
-from .arvworkflow import ArvadosWorkflow
+from .arvworkflow import ArvadosWorkflow, upload_workflow
from .fsaccess import CollectionFsAccess
-from .arvworkflow import upload_workflow
from .perf import Perf
from cwltool.pack import pack
from arvados.api import OrderedJsonModel
logger = logging.getLogger('arvados.cwl-runner')
+metrics = logging.getLogger('arvados.cwl-runner.metrics')
logger.setLevel(logging.INFO)
+
class ArvCwlRunner(object):
"""Execute a CWL tool or workflow, submit work (using either jobs or
containers API), wait for them to complete, and report output.
self.work_api = work_api
self.stop_polling = threading.Event()
self.poll_api = None
+ self.pipeline = None
if self.work_api is None:
# todo: autodetect API to use.
self.cond.acquire()
j = self.processes[uuid]
logger.info("Job %s (%s) is %s", j.name, uuid, event["properties"]["new_attributes"]["state"])
- with Perf(logger, "done %s" % j.name):
+ with Perf(metrics, "done %s" % j.name):
j.done(event["properties"]["new_attributes"])
self.cond.notify()
finally:
# except when in cond.wait(), at which point on_message can update
# job state and process output callbacks.
+ loopperf = Perf(metrics, "jobiter")
+ loopperf.__enter__()
for runnable in jobiter:
+ loopperf.__exit__()
if runnable:
- with Perf(logger, "run"):
+ with Perf(metrics, "run"):
runnable.run(**kwargs)
else:
if self.processes:
else:
logger.error("Workflow is deadlocked, no runnable jobs and not waiting on any pending jobs.")
break
+ loopperf.__enter__()
+ loopperf.__exit__()
while self.processes:
self.cond.wait(1)
exgroup.add_argument("--quiet", action="store_true", help="Only print warnings and errors.")
exgroup.add_argument("--debug", action="store_true", help="Print even more logging")
+ parser.add_argument("--metrics", action="store_true", help="Print timing metrics")
+
parser.add_argument("--tool-help", action="store_true", help="Print command line help for tool")
exgroup = parser.add_mutually_exclusive_group()
def add_arv_hints():
cache = {}
res = pkg_resources.resource_stream(__name__, 'arv-cwl-schema.yml')
- cache["https://w3id.org/cwl/arv-cwl-schema.yml"] = res.read()
+ cache["http://arvados.org/cwl"] = res.read()
res.close()
_, cwlnames, _, _ = cwltool.process.get_schema("v1.0")
- _, extnames, _, _ = schema_salad.schema.load_schema("https://w3id.org/cwl/arv-cwl-schema.yml", cache=cache)
+ _, extnames, _, _ = schema_salad.schema.load_schema("http://arvados.org/cwl", cache=cache)
for n in extnames.names:
- cwlnames.add_name("http://arvados.org/cwl#"+n, "", extnames.get_name(n, ""))
+ if not cwlnames.has_name("http://arvados.org/cwl#"+n, ""):
+ cwlnames.add_name("http://arvados.org/cwl#"+n, "", extnames.get_name(n, ""))
def main(args, stdout, stderr, api_client=None):
parser = arg_parser()
logger.setLevel(logging.WARN)
logging.getLogger('arvados.arv-run').setLevel(logging.WARN)
+ if arvargs.metrics:
+ metrics.setLevel(logging.DEBUG)
+ logging.getLogger("cwltool.metrics").setLevel(logging.DEBUG)
+
arvargs.conformance_test = None
arvargs.use_container = True