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
4 changes: 4 additions & 0 deletions .github/scripts/dispatch.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,8 @@ def sanitize(s: str) -> str:
COVERALLS_TOKEN = os.environ['COVERALLS_REPO_TOKEN']
PR_NUMBER = os.environ.get('PR_NUMBER', '')
JOB_TIMEOUT = int(os.environ.get('JOB_TIMEOUT_SECONDS', '7200'))
OPENALYX_PASSWORD = os.environ.get('OPENALYX_PASSWORD', '')
OPENALYX_USER = os.environ.get('OPENALYX_USER', 'intbrainlab')
# Base image for the Lightning job. python:3.12 verified to ship git/bash/pip.
# uv (installed at runtime) provides the actual test Python via PY_VERSION, so
# this base Python version is only used to bootstrap `pip install uv`.
Expand Down Expand Up @@ -83,6 +85,8 @@ def sanitize(s: str) -> str:
'PY_VERSION': PY_VERSION,
'REPO_URL': REPO_URL,
'INTEGRATION_DATA_DIR': INTEGRATION_DATA_DIR,
'OPENALYX_PASSWORD': OPENALYX_PASSWORD,
'OPENALYX_USER': OPENALYX_USER,
# --- Coveralls auth + parallel grouping ---
'COVERALLS_REPO_TOKEN': COVERALLS_TOKEN,
'COVERALLS_PARALLEL': 'true',
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/integration-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ jobs:
# --- consumed by dispatch.py / Job.run ---
LIGHTNING_TEAMSPACE: ${{ vars.LIGHTNING_TEAMSPACE }} # "owner/teamspace"
INTEGRATION_DATA_DIR: ${{ vars.INTEGRATION_DATA_DIR }}
OPENALYX_PASSWORD: ${{ secrets.ALYX_PWD }}
REPO_URL: https://github.com/${{ github.repository }}.git
COVERALLS_REPO_TOKEN: ${{ secrets.COVERALLS_REPO_TOKEN }}
PY_VERSION: ${{ matrix.python_version }}
Expand Down
2 changes: 2 additions & 0 deletions .github/workflows/unit-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,8 @@ jobs:
python-version: 3.13
env:
ONE_SAVE_ON_DELETE: false
OPENALYX_USER: intbrainlab
OPENALYX_PASSWORD: ${{ secrets.ALYX_PWD }}
steps:
- uses: actions/checkout@v7
- name: Set up Python ${{ matrix.python-version }}
Expand Down
15 changes: 15 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,9 +7,24 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [3.5.0] Unreleased

### Added
- `ibllib.pipes.spec.TaskSpec`: plain-data task specification for creating Alyx tasks without importing the task class
- `ibllib.pipes.routing`: route tasks to environments by their executable; `ENV_PATHS` maps environment labels to server venvs
- `ibllib.pipes.plan`: plan the tasks of other repositories (e.g. mpci) in the environment that runs them, via `python -m ibllib.pipes.plan`

### Changed
- SpikeSortingLoaders
- merging clusters with channels doesn't require reading spikes if metrics are available
- `local_server.task_queue` and `list_queued_envs` determine a task's environment from its executable and only import task classes of the requested environments
- `dynamic_pipeline.make_pipeline` no longer imports mpci; mesoscope tasks are planned by `mpci.alyx.pipeline:plan` in the mpci env (requires mpci with the `plan` function)
- `local_server.job_creator` keeps the raw_session.flag file if an external task planner fails, so the missing tasks are created on the next run
- `Pipeline.create_alyx_tasks` creates tasks in dependency order and computes task levels from the parents

### Fixed
- Task time outs are now stored on Alyx: tasks were posted with the key `time_out_sec` instead of the model field `time_out_secs`. `Pipeline.create_alyx_tasks` asserts that time outs don't exceed the Alyx field maximum (32767 s) before creating any tasks

### Removed
- `dynamic_pipeline.get_mesoscope_tasks`

## [4.0.1] 2026-05-22

Expand Down
2 changes: 1 addition & 1 deletion ibllib/ephys/sync_probes.py
Original file line number Diff line number Diff line change
Expand Up @@ -84,7 +84,7 @@ def get_sync_fronts(auxiliary_name):
# exits if sync label not found for current probe
if auxiliary_name not in sync_map:
return
isync = np.in1d(sync['channels'], np.array([sync_map[auxiliary_name]]))
isync = np.isin(sync['channels'], np.array([sync_map[auxiliary_name]]))
# only returns syncs if we get fronts for all probes
if np.all(~isync):
return
Expand Down
4 changes: 4 additions & 0 deletions ibllib/pipes/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,10 @@
inferring the acquisition hardware from the task protocol. The new session's pipeline tasks are
then registered for another process (or server) to query.

Tasks from repositories installed in other environments (e.g. mpci) are planned within their own
environment as plain-data :class:`spec.TaskSpec` objects (see :mod:`plan`), so the job creator
doesn't need to import them.

Another process calls :func:`local_server.task_queue` to get a list of queued tasks from Alyx, then
:func:`local_server.tasks_runner` to loop through tasks. Each task is run by calling
:func:`tasks.run_alyx_task` with a dictionary of task information, including the Task class and its
Expand Down
38 changes: 19 additions & 19 deletions ibllib/pipes/dynamic_pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
import ibllib.io.raw_data_loaders as rawio
import ibllib.io.session_params as sess_params
import ibllib.pipes.tasks as mtasks
from ibllib.pipes.plan import get_external_tasks
import ibllib.pipes.base_tasks as bstasks
import ibllib.pipes.widefield_tasks as wtasks
import ibllib.pipes.sync_tasks as stasks
Expand Down Expand Up @@ -516,8 +517,10 @@ def get_audio_tasks(acquisition_description, **kwargs):
devices = acquisition_description.get('devices', {})
audio_tasks = OrderedDict()
if 'microphone' in devices:
((microphone, micro_kwargs),) = devices['microphone'].items()
micro_kwargs['device_collection'] = micro_kwargs.pop('collection')
((microphone, micro_info),) = devices['microphone'].items()
# Rename the collection key without modifying the acquisition description
micro_kwargs = {k: v for k, v in micro_info.items() if k != 'collection'}
micro_kwargs['device_collection'] = micro_info['collection']
if sync_kwargs['sync'] == 'bpod':
audio_tasks['AudioRegisterRaw'] = type('AudioRegisterRaw', (atasks.AudioSync,), {})(
**kwargs, **sync_kwargs, **micro_kwargs, collection=micro_kwargs['device_collection']
Expand All @@ -534,8 +537,10 @@ def get_wfield_tasks(acquisition_description, sync_tasks, **kwargs):
wfield_tasks = OrderedDict()

if 'widefield' in devices:
((_, wfield_kwargs),) = devices['widefield'].items()
wfield_kwargs['device_collection'] = wfield_kwargs.pop('collection')
((_, wfield_info),) = devices['widefield'].items()
# Rename the collection key without modifying the acquisition description
wfield_kwargs = {k: v for k, v in wfield_info.items() if k != 'collection'}
wfield_kwargs['device_collection'] = wfield_info['collection']
wfield_tasks['WideFieldRegisterRaw'] = type('WidefieldRegisterRaw', (wtasks.WidefieldRegisterRaw,), {})(
**kwargs, **wfield_kwargs
)
Expand All @@ -558,16 +563,6 @@ def get_wfield_tasks(acquisition_description, sync_tasks, **kwargs):
return wfield_tasks


def get_mesoscope_tasks(acquisition_description, **kwargs):
if 'mesoscope' not in acquisition_description.get('devices', {}):
return OrderedDict()

import mpci.alyx.pipeline

pipe = mpci.alyx.pipeline.make_pipeline(acquisition_description, **kwargs)
return pipe.tasks


def get_photometry_tasks(acquisition_description, **kwargs):
devices = acquisition_description.get('devices', {})
photometry_tasks = OrderedDict()
Expand Down Expand Up @@ -650,7 +645,10 @@ def make_pipeline(session_path, **pkwargs):
Returns
-------
ibllib.pipes.tasks.Pipeline
A task pipeline object.
A task pipeline object. Tasks from other repositories (see
:data:`ibllib.pipes.plan.PLANNERS`) are :class:`ibllib.pipes.spec.TaskSpec` objects, planned
in the environment that runs them. If any of these planners fail, the pipeline
`planner_errors` attribute maps the device to the error message.
"""
# NB: this pattern is a pattern for dynamic class creation
# tasks['SyncPulses'] = type('SyncPulses', (epp.EphysPulses,), {})(session_path=session_path)
Expand Down Expand Up @@ -694,17 +692,19 @@ def make_pipeline(session_path, **pkwargs):
wfield_tasks = get_wfield_tasks(acquisition_description, sync_parent_tasks, **kwargs)
tasks.update(wfield_tasks)

# Mesoscope tasks
mesoscope_tasks = get_mesoscope_tasks(acquisition_description, **kwargs)
tasks.update(mesoscope_tasks)

# photometry tasks
# photometry_tasks = get_photometry_tasks(acquisition_description, **kwargs)
# tasks.update(photometry_tasks)

# Tasks from other repositories (e.g. mesoscope), planned as specs in the env that runs them
context = {'tasks': [t.to_spec().to_dict() for t in tasks.values()]}
external_tasks, planner_errors = get_external_tasks(acquisition_description, session_path, context=context)
tasks.update(external_tasks)

# combine: make pipeline and add tasks
p = mtasks.Pipeline(session_path=session_path, **pkwargs)
p.tasks = tasks
p.planner_errors = planner_errors
return p


Expand Down
125 changes: 94 additions & 31 deletions ibllib/pipes/local_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@

from ibllib import __version__ as ibllib_version
from ibllib.pipes import tasks
from ibllib.pipes.routing import task_env
from ibllib.time import date2isostr
from ibllib.oneibl.registration import IBLRegistrationClient
from ibllib.oneibl.data_handlers import get_local_data_repository
Expand Down Expand Up @@ -133,7 +134,15 @@ def job_creator(root_path, one=None, dry=False, rerun=False):
else:
rerun__status__in = ['Waiting']
pipe.create_alyx_tasks(rerun__status__in=rerun__status__in)
flag_file.unlink()
if pipe.planner_errors:
# Keep the flag file so that the missing tasks are created on the next run
_logger.error(
'Failed to plan %s tasks for session %s; keeping flag file to retry',
', '.join(pipe.planner_errors),
session_path.relative_to(root_path),
)
else:
flag_file.unlink()
if pipe is not None:
pipes.append(pipe)
except Exception:
Expand All @@ -147,6 +156,9 @@ def list_available_envs(root=Path.home() / 'Documents/PYTHON/envs'):
"""
List all the envs within `root` dir.

NB: Environment labels don't necessarily match the venv directory names; use
:func:`ibllib.pipes.routing.installed_envs` to list the installed environment labels.

Parameters
----------
root : str, pathlib.Path
Expand All @@ -164,27 +176,92 @@ def list_available_envs(root=Path.home() / 'Documents/PYTHON/envs'):
return [None]


def list_queued_envs(one=None):
def _waiting_tasks(alyx, lab=None):
"""
Query the waiting tasks of a lab for the local data repository.

Parameters
----------
alyx : one.webclient.AlyxClient
An Alyx instance.
lab : str
Lab name as per Alyx, otherwise try to infer from local Globus install.

Returns
-------
list of dict, None
A list of Alyx tasks with a 'Waiting' status, or None if the lab could not be determined.
"""
if lab is None:
_logger.debug('Trying to infer lab from globus installation')
lab = get_lab_from_endpoint_id(alyx=alyx)
if lab is None:
_logger.error('No lab provided or found')
return # if the lab is none, this will return empty tasks each time
data_repo = get_local_data_repository(alyx)
return alyx.rest(
'tasks', 'list', status='Waiting', django=f'session__lab__name__in,{lab},data_repository__name,{data_repo}', no_cache=True
)


def list_queued_envs(one=None, lab=None):
"""
The set of all envs in the list of waiting tasks.

The environment of each task is determined from its executable, without importing the task class.

Parameters
----------
one : one.api.OneAlyx
An instance of ONE.
lab : str
Lab name as per Alyx, otherwise try to infer from local Globus install.

Returns
-------
set
All environments required to process waiting tasks.
"""
one = one or ONE(mode='remote', cache_rest=None)
waiting_tasks = task_queue(mode='large', alyx=one.alyx, env=list_available_envs())
envs_in_queue = set()
for task_exe in map(lambda x: x['executable'], waiting_tasks):
envs_in_queue.add(tasks.str2class(task_exe).env)
return envs_in_queue
return {task_env(t['executable']) for t in _waiting_tasks(one.alyx, lab=lab) or []}


def is_job_size(task, mode):
"""
Check whether a task is of a given job size.

NB: This imports the task class.

Parameters
----------
task : dict
An Alyx task dictionary.
mode : {'all', 'small', 'large'}
The job size to check.

Returns
-------
bool
True if the task class job size matches the mode (or mode is 'all'). False if the task
class could not be imported.
"""
if mode == 'all':
return True
try:
return tasks.str2class(task['executable']).job_size == mode
except (ImportError, AttributeError):
_logger.error('Task %s not found in this env', task['executable'])
return False


def task_queue(mode='all', lab=None, alyx=None, env=(None,)):
"""
Query waiting jobs from the specified Lab

The environment of each task is determined from its executable (see
:func:`ibllib.pipes.routing.task_env`), so only the task classes of the given environments are
imported (to determine the job size).

Parameters
----------
mode : {'all', 'small', 'large'}
Expand All @@ -193,37 +270,23 @@ def task_queue(mode='all', lab=None, alyx=None, env=(None,)):
Lab name as per Alyx, otherwise try to infer from local Globus install.
alyx : one.webclient.AlyxClient
An Alyx instance.
env : list
One or more environments to filter by. See :prop:`ibllib.pipes.tasks.Task.env`.
env : str, list
One or more environment labels to filter by, where None is the base environment. See
:data:`ibllib.pipes.routing.ROUTES`.

Returns
-------
list of dict
A list of Alyx tasks associated with `lab` that have a 'Waiting' status.
"""

def predicate(task):
try:
classe = tasks.str2class(task['executable'])
return (mode == 'all' or classe.job_size == mode) and classe.env in env
except ModuleNotFoundError:
_logger.error('Task %s not found in this env', task['executable'])
return False

env = (env,) if env is None or isinstance(env, str) else tuple(env)
alyx = alyx or AlyxClient(cache_rest=None)
if lab is None:
_logger.debug('Trying to infer lab from globus installation')
lab = get_lab_from_endpoint_id(alyx=alyx)
if lab is None:
_logger.error('No lab provided or found')
return # if the lab is none, this will return empty tasks each time
data_repo = get_local_data_repository(alyx)
# Filter for tasks
waiting_tasks = alyx.rest(
'tasks', 'list', status='Waiting', django=f'session__lab__name__in,{lab},data_repository__name,{data_repo}', no_cache=True
)
# Filter tasks by size
filtered_tasks = filter(predicate, waiting_tasks)
waiting_tasks = _waiting_tasks(alyx, lab=lab)
if waiting_tasks is None:
return
# Filter tasks by environment, then by size
filtered_tasks = (t for t in waiting_tasks if task_env(t['executable']) in env)
filtered_tasks = filter(lambda t: is_job_size(t, mode), filtered_tasks)
# Order tasks by priority
sorted_tasks = sorted(filtered_tasks, key=lambda d: d['priority'], reverse=True)

Expand Down
Loading
Loading