mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Update build runner documentation
This commit is contained in:
@@ -26,15 +26,15 @@ overridden with your logic.
|
||||
|
||||
These are the following methods which need to be implemented:
|
||||
|
||||
1. \_\_init\_\_(app, nworkers, \*\*kwargs)
|
||||
1. ``__init__(app, nworkers, **kwargs)``
|
||||
|
||||
2. queue\_job(job\_wrapper)
|
||||
2. ``queue_job(job_wrapper)``
|
||||
|
||||
3. check\_watched\_item(job\_state)
|
||||
3. ``check_watched_item(job_state)``
|
||||
|
||||
4. stop\_job(job)
|
||||
4. ``stop_job(job)``
|
||||
|
||||
5. recover(job, job\_wrapper)
|
||||
5. ``recover(job, job_wrapper)``
|
||||
|
||||
The big picture
|
||||
---------------
|
||||
@@ -45,8 +45,8 @@ framework and the external executor framework. To know, when and how
|
||||
these methods are invoked, we will see about the implementation of
|
||||
parent class and process lifecycle of the runner.
|
||||
|
||||
Implementation of parent class (galaxy.jobs.runners.\_\_init\_\_.py)
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
Implementation of parent class (``galaxy.jobs.runners.__init__.py``)
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
- .. rubric:: Class Inheritance structure
|
||||
:name: class-inheritance-structure
|
||||
@@ -64,28 +64,28 @@ purpose.
|
||||
Runner Methods in detail
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
1. \_\_init\_\_ method - STAGE 1
|
||||
1. ``__init__`` method - STAGE 1
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
Input params:
|
||||
|
||||
1. app
|
||||
1. ``app``
|
||||
|
||||
2. nworkers (Number of threads specified in job\_conf)
|
||||
2. ``nworkers`` (Number of threads specified in ``job_conf``)
|
||||
|
||||
3. \*\*kwargs (Variable length argument)
|
||||
3. ``**kwargs`` (Variable length argument)
|
||||
|
||||
Output params: NA
|
||||
|
||||
The input params are read from job\_conf.xml and passed to the runner by
|
||||
The input params are read from ``job_conf.xml`` and passed to the runner by
|
||||
the Galaxy framework. Configuration of where to run jobs and external
|
||||
runner configuration is performed in the job\_conf.xml file. More
|
||||
information about job\_conf.xml is available
|
||||
runner configuration is performed in the ``job_conf.xml`` file. More
|
||||
information about ``job_conf.xml`` is available
|
||||
`here <https://galaxyproject.org/admin/config/jobs/>`__.
|
||||
|
||||
Have a look at the sample job\_conf.xml:
|
||||
Have a look at the sample ``job_conf.xml``:
|
||||
|
||||
::
|
||||
.. code-block:: xml
|
||||
|
||||
<job_conf>
|
||||
<plugins>
|
||||
@@ -107,30 +107,30 @@ Have a look at the sample job\_conf.xml:
|
||||
</destinations>
|
||||
</job_conf>
|
||||
|
||||
The following steps are followed to manipulate the data in job\_conf.xml
|
||||
The following steps are followed to manipulate the data in ``job_conf.xml``
|
||||
|
||||
A: Define structure of data under plugin tag (plugin tag in
|
||||
job\_conf.xml) as a dictionary.
|
||||
``job_conf.xml``) as a dictionary.
|
||||
|
||||
::
|
||||
.. code-block:: python
|
||||
|
||||
runner_param_specs = dict(user=dict(map=str), key=dict(map=str))
|
||||
|
||||
B: Update the dictionary structure in kwargs.
|
||||
|
||||
::
|
||||
.. code-block:: python
|
||||
|
||||
kwargs.update({'runner_param_specs': runner_param_specs})
|
||||
|
||||
C: Now call the parent constructor to assign the values.
|
||||
|
||||
::
|
||||
.. code-block:: python
|
||||
|
||||
super(GodockerJobRunner, self).__init__(app, nworkers, **kwargs)
|
||||
|
||||
D: The assigned values can be accessed in runner in the following way.
|
||||
|
||||
::
|
||||
.. code-block:: python
|
||||
|
||||
print self.runner_params["user"]
|
||||
print self.runner_params["key"]
|
||||
@@ -148,50 +148,50 @@ for initialization.
|
||||
Finally the worker threads and monitor threads are invoked for galaxy to
|
||||
listen for incoming tool submissions.
|
||||
|
||||
::
|
||||
.. code-block:: python
|
||||
|
||||
self._init_monitor_thread()
|
||||
self._init_worker_threads()
|
||||
|
||||
2. queue\_job method - STAGE 2
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
2. ``queue_job`` method - STAGE 2
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
Input params: job\_wrapper (Object of
|
||||
Input params: ``job_wrapper`` (Object of
|
||||
`galaxy.jobs.JobWrapper <https://github.com/galaxyproject/galaxy/blob/dev/lib/galaxy/jobs/__init__.py#L743>`__)
|
||||
|
||||
Output params: None
|
||||
|
||||
galaxy.jobs.JobWrapper is a Wrapper around 'model.Job' with convenience
|
||||
``galaxy.jobs.JobWrapper`` is a Wrapper around 'model.Job' with convenience
|
||||
methods for running processes and state management.
|
||||
|
||||
- Functioning of queue\_job method.
|
||||
- Functioning of ``queue_job`` method.
|
||||
|
||||
A. prepare\_job() method is invoked to do some sanity checks that all runners' queue\_job() methods are
|
||||
A. ``prepare_job()`` method is invoked to do some sanity checks that all runners' ``queue_job()`` methods are
|
||||
likely to want to do and also to build runner command line for that
|
||||
job. Initial state and configuration of the job are set and every
|
||||
data is associated with **job\_wrapper**.
|
||||
|
||||
B. Submit job to the external runner and return the jobid. Accessing
|
||||
jobs data (tool submitted in Galaxy webframework) is purely from
|
||||
job\_wrapper. eg: job\_wrapper.get\_state() -> gives state of a job
|
||||
``job_wrapper``. eg: ``job_wrapper.get_state()`` -> gives state of a job
|
||||
(queued/running/failed/success/...)
|
||||
|
||||
Let us look at a means of accessing external runner's configuration
|
||||
present under destination tag of job\_conf.xml in the above example.
|
||||
present under destination tag of ``job_conf.xml`` in the above example.
|
||||
|
||||
::
|
||||
.. code-block:: python
|
||||
|
||||
job_destination = job_wrapper.job_destination
|
||||
docker_cpu = int(job_destination.params["docker_cpu"])
|
||||
docker_ram = int(job_destination.params["docker_memory"])
|
||||
|
||||
A special case: User Story: A docker based external runner is present. A
|
||||
default docker image for execution is set in job\_conf.xml. A tool can
|
||||
default docker image for execution is set in ``job_conf.xml``. A tool can
|
||||
also specify the docker image for its execution. Specification in tool
|
||||
is given more priority than the default specification. To achieve such a
|
||||
functionality. We can use the following statement:
|
||||
|
||||
::
|
||||
.. code-block:: python
|
||||
|
||||
docker_image = self._find_container(job_wrapper).container_id
|
||||
|
||||
@@ -200,20 +200,20 @@ image/container/os..
|
||||
|
||||
C. After successful submission of job in the external runner, submit the
|
||||
job to Galaxy framework. To do that,make an object of
|
||||
AsynchronousJobState and put it in monitor\_queue.
|
||||
AsynchronousJobState and put it in ``monitor_queue``.
|
||||
|
||||
::
|
||||
.. code-block:: python
|
||||
|
||||
ajs = AsynchronousJobState(files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper, job_id=job_id, job_destination=job_destination)
|
||||
self.monitor_queue.put(ajs)
|
||||
|
||||
3. check\_watched\_item method - STAGE 3
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
3. ``check_watched_item`` method - STAGE 3
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
Input params: job\_state (Object of
|
||||
Input params: ``job_state`` (Object of
|
||||
`galaxy.jobs.runners.AsynchronousJobState <https://github.com/galaxyproject/galaxy/blob/dev/lib/galaxy/jobs/runners/__init__.py#L400>`__)
|
||||
|
||||
Output params: AsynchronousJobState object
|
||||
Output params: ``AsynchronousJobState`` object
|
||||
|
||||
Without going into much detail, assume there is a queue to track the status of every job. eg:
|
||||
|
||||
@@ -221,12 +221,12 @@ Without going into much detail, assume there is a queue to track the status of e
|
||||
:align: center
|
||||
|
||||
The galaxy framework updates the status of a job by iterating through the
|
||||
queue. During the iteration, it calls check\_watched\_item method with the job.
|
||||
queue. During the iteration, it calls ``check_watched_item`` method with the job.
|
||||
Your responsibility will be to get the status of execution of the job from the
|
||||
external runner and return the updated status of the job, and also to
|
||||
copy the output files for the completed jobs.
|
||||
|
||||
Updated result after an iteration (after invocation of check\_watched\_item 6 times):
|
||||
Updated result after an iteration (after invocation of ``check_watched_item`` 6 times):
|
||||
|
||||
.. image:: queue_b.png
|
||||
:align: center
|
||||
@@ -236,17 +236,17 @@ Note: Iterating through the queue is already taken care by the framework.
|
||||
|
||||
To inform galaxy about the status of the job:
|
||||
|
||||
- Get the job status from external runner using the job\_id.
|
||||
- Get the job status from external runner using the ``job_id``.
|
||||
|
||||
- Check if the job is queued/running/completed.. etc. A general structure is provided below.
|
||||
|
||||
- Call self.mark\_as\_finished(job\_state), if the job has been successfully executed.
|
||||
- Call ``self.mark_as_finished(job_state)``, if the job has been successfully executed.
|
||||
|
||||
- Call self.mark\_as\_failed(job\_state), if the job has failed during execution.
|
||||
- Call ``self.mark_as_failed(job_state)``, if the job has failed during execution.
|
||||
|
||||
- To change state of a job, change job\_state.running and job\_state.job\_wrapper.change\_state()
|
||||
- To change state of a job, change ``job_state.running`` and ``job_state.job_wrapper.change_state()``
|
||||
|
||||
::
|
||||
.. code-block:: python
|
||||
|
||||
def check_watched_item(self, job_state):
|
||||
!job_status = get_task_from_external_runner(job_state.job_id)
|
||||
@@ -276,23 +276,23 @@ Note:
|
||||
|
||||
- Methods prefixed with ! are user-defined methods.
|
||||
|
||||
- Return value is job\_state for running,pending jobs and None for rest of the states of jobs.
|
||||
- Return value is ``job_state`` for running, pending jobs and None for rest of the states of jobs.
|
||||
|
||||
create\_log\_files() are nothing but copying the files (error\_file,
|
||||
output\_file, exit\_code\_file) from external runner's directory to
|
||||
``create_log_files()`` are nothing but copying the files (``error_file``,
|
||||
``output_file``, ``exit_code_file``) from external runner's directory to
|
||||
working directory of Galaxy.
|
||||
|
||||
Source of the files are from the output directory of your external
|
||||
runner. Destination of the files will be:
|
||||
|
||||
- output file -> job\_state.output\_file.
|
||||
- output file -> ``job_state.output_file``.
|
||||
|
||||
- error file -> job\_state.error\_file.
|
||||
- error file -> ``job_state.error_file``.
|
||||
|
||||
- exit code file -> job\_state.exit\_code\_file.
|
||||
- exit code file -> ``job_state.exit_code_file``.
|
||||
|
||||
4. stop\_job method - STAGE 4
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
4. ``stop_job`` method - STAGE 4
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
Input params: job (Object of
|
||||
`galaxy.model.Job <https://github.com/galaxyproject/galaxy/blob/dev/lib/galaxy/model/__init__.py#L344>`__)
|
||||
@@ -305,20 +305,20 @@ runner.
|
||||
When an user requests to stop the execution of job in Galaxy framework,
|
||||
a call is made to the external runner to stop the job execution.
|
||||
|
||||
The job\_id of the job to be deleted is accessed by
|
||||
The ``job_id`` of the job to be deleted is accessed by
|
||||
|
||||
::
|
||||
.. code-block:: python
|
||||
|
||||
job.id
|
||||
|
||||
5. recover method - STAGE 5
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
5. ``recover`` method - STAGE 5
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
Input params:
|
||||
|
||||
- job (Object of `galaxy.model.Job <https://github.com/galaxyproject/galaxy/blob/dev/lib/galaxy/model/__init__.py#L344>`__).
|
||||
- ``job`` (Object of `galaxy.model.Job <https://github.com/galaxyproject/galaxy/blob/dev/lib/galaxy/model/__init__.py#L344>`__).
|
||||
|
||||
- job\_wrapper (Object of `galaxy.jobs.JobWrapper <https://github.com/galaxyproject/galaxy/blob/dev/lib/galaxy/jobs/__init__.py#L743>`__).
|
||||
- ``job_wrapper`` (Object of `galaxy.jobs.JobWrapper <https://github.com/galaxyproject/galaxy/blob/dev/lib/galaxy/jobs/__init__.py#L743>`__).
|
||||
|
||||
|
||||
Output params: None
|
||||
@@ -327,12 +327,12 @@ Functionality: Recovers jobs stuck in the queued/running state when
|
||||
Galaxy started.
|
||||
|
||||
This method is invoked by Galaxy at the time of startup. Jobs in Running
|
||||
& Queued status in Galaxy are put in the monitor\_queue by creating an
|
||||
AsynchronousJobState object.
|
||||
& Queued status in Galaxy are put in the ``monitor_queue`` by creating an
|
||||
``AsynchronousJobState`` object.
|
||||
|
||||
The following is a generic code snippet for recover method.
|
||||
The following is a generic code snippet for ``recover`` method.
|
||||
|
||||
::
|
||||
.. code-block:: python
|
||||
|
||||
ajs = AsynchronousJobState(files_dir=job_wrapper.working_directory, job_wrapper=job_wrapper)
|
||||
ajs.job_id = str(job_wrapper.job_id)
|
||||
|
||||
Reference in New Issue
Block a user