Skip to main content

hooks

Hook infrastructure for Bitfount.

Attributes: on_pod_init_error: Decorator to be used on the Pod.__init__ method. on_pod_startup_error: Decorator to be used on the Pod.start method.

Module​

Functions​

get_child_process_hook_factories​

def get_child_process_hook_factories() ‑> list[collections.abc.Callable[[], None]]:

Collect hook factories for re-registration in a spawned child process.

Iterates all registered hooks across every hook type. Any hook whose class defines a child_process_factory instance method (returning a no-arg callable) will contribute a factory to the returned list.

This is the mechanism by which the orchestrator (or any consumer) can ensure its hooks are present in the child worker process that runs the protocol. The returned callables must be picklable (plain module-level functions or functools.partial over such functions).

get_child_process_hook_factory_paths​

def get_child_process_hook_factory_paths() ‑> list[str]:

Collect hook factories as "module:qualname" import paths.

The dotted-path form of get_child_process_hook_factories, for a child this process cannot pickle to. A background DAG child is launched by Prefect from a served deployment and receives only JSON, so a factory reaches it as a string it imports for itself.

Only a factory that is addressable by import survives the trip: a plain module-level function is, a lambda and a closure are not. A functools.partial is addressable through the function it wraps — the child calls the factory with no arguments either way — but only while it binds nothing, because bound arguments do not travel as an import path. That matters because get_child_process_hook_factories names a partial over a module-level function as a supported factory form, for the pickling the v8 worker's child does instead of this.

A non-addressable factory is skipped with a warning rather than failing the collection, because the alternative is a child that refuses to run over a hook it could have done without.

Returns One "module:qualname" string per addressable factory.

get_hooks​

def get_hooks(    type: HookType,) ‑> list[bitfount.hooks.AlgorithmHookProtocol] | list[bitfount.hooks.PodHookProtocol] | list[bitfount.hooks.ProtocolHookProtocol] | list[bitfount.hooks.ModellerHookProtocol] | list[bitfount.hooks.DatasourceHookProtocol]:

Get all registered hooks of a particular type.

Arguments

  • type: The type of hook to get.

Returns A list of hooks of the provided type.

Raises

  • ValueError: If the provided type is not a valid hook type.

import_hook_factory​

def import_hook_factory(path: str) ‑> collections.abc.Callable[[], None]:

Import the hook factory named by a "module:qualname" path.

Arguments

  • path: A path produced by get_child_process_hook_factory_paths.

Returns The factory the path names.

Raises

  • ValueError: If path is not of the form "module:qualname", or names something that is not callable.
  • ImportError: If the module cannot be imported.
  • AttributeError: If the module has no such attribute.

Classes​

BaseAlgorithmHook​

class BaseAlgorithmHook():

Base algorithm hook class.

Initialise the hook.

Variables​

  • type : HookType - Return the hook type.

Methods​


on_init_end​

def on_init_end(self, algorithm: _BaseAlgorithm, *args: Any, **kwargs: Any) ‑> None:

Run the hook at the very end of algorithm initialisation.

on_init_start​

def on_init_start(self, algorithm: _BaseAlgorithm, *args: Any, **kwargs: Any) ‑> None:

Run the hook at the very start of algorithm initialisation.

on_progress​

def on_progress(    self,    algorithm: _BaseAlgorithm,    context: TaskContext | None,    step: str,    current_epoch: int | None = None,    max_epochs: int | None = None,    step_info: str | None = None,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook at an arbitrary progress checkpoint during algorithm run.

on_run_end​

def on_run_end(    self,    algorithm: _BaseAlgorithm,    context: TaskContext | None,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook at the very end of algorithm run.

on_run_start​

def on_run_start(    self,    algorithm: _BaseAlgorithm,    context: TaskContext | None,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook at the very start of algorithm run.

on_train_epoch_end​

def on_train_epoch_end(    self,    current_epoch: int,    min_epochs: int | None,    max_epochs: int | None,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook at the end of an epoch testing.

Only applicable for training algorithms.

on_train_epoch_start​

def on_train_epoch_start(    self,    current_epoch: int,    min_epochs: int | None,    max_epochs: int | None,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook at the start of an epoch testing.

Only applicable for training algorithms.

BasePodHook​

class BasePodHook():

Base pod hook class.

Initialise the hook.

Variables​

  • type : HookType - Return the hook type.

Methods​


on_batches_complete​

def on_batches_complete(    self,    task_id: str,    modeller_username: str,    total_batches: int,    total_files: int,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook when all batches are processed but before resilience starts.

on_configuring_task​

def on_configuring_task(self, task_id: str, *args: Any, **kwargs: Any) ‑> None:

Run the hook when the worker is configuring the task.

on_file_filter_progress​

def on_file_filter_progress(    self, total_files: int, total_skipped: int, *args: Any, **kwargs: Any,) ‑> None:

Run the hook when filtering files to track progress.

on_file_process_end​

def on_file_process_end(    self,    datasource: FileSystemIterableSource,    file_num: int,    total_num_files: int | None,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook when a file processing ends.

on_file_process_start​

def on_file_process_start(    self,    datasource: FileSystemIterableSource,    file_num: int,    total_num_files: int | None,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook when a file starts to be processed.

on_files_partition​

def on_files_partition(    self,    datasource: FileSystemIterableSource,    total_num_files: int | None,    batch_size: int,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook when we partition files to be processed.

on_model_download_end​

def on_model_download_end(    self, task_id: str | None, model_id: str, source: str, *args: Any, **kwargs: Any,) ‑> None:

Run when a model download completes.

on_model_download_start​

def on_model_download_start(    self, task_id: str | None, model_id: str, source: str, *args: Any, **kwargs: Any,) ‑> None:

Run when a model download starts.

on_model_upload_end​

def on_model_upload_end(    self, task_id: str | None, model_id: str, operation: str, *args: Any, **kwargs: Any,) ‑> None:

Run when a model upload completes.

on_model_upload_start​

def on_model_upload_start(    self, task_id: str | None, model_id: str, operation: str, *args: Any, **kwargs: Any,) ‑> None:

Run when a model upload starts.

on_pod_init_end​

def on_pod_init_end(self, pod: Pod, *args: Any, **kwargs: Any) ‑> None:

Run the hook at the end of pod initialisation.

on_pod_init_error​

def on_pod_init_error(    self, pod: Pod, exception: BaseException, *args: Any, **kwargs: Any,) ‑> None:

Run the hook if an uncaught exception is raised during pod initialisation.

Raises

  • NotImplementedError: If the hook is not implemented. This is to ensure that underlying exceptions are not swallowed if the hook is not implemented. This error is caught further up the chain and the underlying exception is raised instead.

on_pod_init_progress​

def on_pod_init_progress(    self,    pod: Pod,    message: str,    datasource_name: str | None = None,    base_datasource_names: list[str] | None = None,    pod_db_enabled: bool | None = None,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook at key points of pod initialisation.

on_pod_init_start​

def on_pod_init_start(    self, pod: Pod, pod_name: str, username: str | None = None, *args: Any, **kwargs: Any,) ‑> None:

Run the hook at the very start of pod initialisation.

on_pod_shutdown_end​

def on_pod_shutdown_end(self, pod: Pod, *args: Any, **kwargs: Any) ‑> None:

Run the hook at the very end of pod shutdown.

on_pod_shutdown_start​

def on_pod_shutdown_start(self, pod: Pod, *args: Any, **kwargs: Any) ‑> None:

Run the hook at the very start of pod shutdown.

on_pod_startup_end​

def on_pod_startup_end(self, pod: Pod, *args: Any, **kwargs: Any) ‑> None:

Run the hook at the end of pod startup.

on_pod_startup_error​

def on_pod_startup_error(    self, pod: Pod, exception: BaseException, *args: Any, **kwargs: Any,) ‑> None:

Run the hook if an uncaught exception is raised during pod startup.

Raises

  • NotImplementedError: If the hook is not implemented. This is to ensure that underlying exceptions are not swallowed if the hook is not implemented. This error is caught further up the chain and the underlying exception is raised instead.

on_pod_startup_start​

def on_pod_startup_start(self, pod: Pod, *args: Any, **kwargs: Any) ‑> None:

Run the hook at the very start of pod startup.

on_pod_task_data_check​

def on_pod_task_data_check(    self, task_id: str, message: str, *args: Any, **kwargs: Any,) ‑> None:

Run the hook at start of a job request to check that the pod has data.

on_preparing_data​

def on_preparing_data(self, task_id: str, *args: Any, **kwargs: Any) ‑> None:

Run the hook when the worker is about to prepare data for a task.

on_preparing_data_batches​

def on_preparing_data_batches(self, task_id: str, *args: Any, **kwargs: Any) ‑> None:

Run the hook when the worker is preparing data batches.

on_process_spawn_error​

def on_process_spawn_error(    self, process_name: str, attempts: int, error_message: str, *args: Any, **kwargs: Any,) ‑> None:

Run when a background process fails to spawn after all retries.

on_resilience_complete​

def on_resilience_complete(    self,    task_id: str,    modeller_username: str,    total_attempted: int,    total_succeeded: int,    total_failed: int,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook when individual file retry phase is complete.

on_resilience_progress​

def on_resilience_progress(    self,    task_id: str,    modeller_username: str,    current_file: int,    total_files: int,    file_name: str,    success: bool,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook for each individual file retry attempt.

on_resilience_start​

def on_resilience_start(    self,    task_id: str,    modeller_username: str,    total_failed_files: int,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook when individual file retry phase begins.

on_task_abort​

def on_task_abort(    self,    pod: Pod,    message: str,    task_id: str,    project_id: str | None,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook when there is an exception in a task.

on_task_end​

def on_task_end(self, pod: Pod, task_id: str, *args: Any, **kwargs: Any) ‑> None:

Run the hook when a new task is received at the end.

on_task_error​

def on_task_error(    self,    pod: Pod,    exception: BaseException,    task_id: str,    project_id: str | None,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook when there is an exception in a task.

on_task_progress​

def on_task_progress(self, task_id: str, message: str, *args: Any, **kwargs: Any) ‑> None:

Run the hook at key points of the task.

on_task_start​

def on_task_start(    self,    pod: Pod,    task_id: str,    project_id: str | None,    modeller_username: str,    protocol_name: str,    save_path: str | None = None,    primary_results_path: str | None = None,    dataset_name: str | None = None,    *args: Any,    **kwargs: Any,) ‑> None:

Run the hook when a new task is received at the start.

HookType​

class HookType(*args, **kwds):

Enum for hook types.

Ancestors​

Variables​

  • static ALGORITHM
  • static DATASOURCE
  • static MODELLER
  • static POD
  • static PROTOCOL