Skip to main content

worker_process

Contains classes and methods for spawning the Worker in a dedicated process.

The parent owns the real worker mailbox and message-service connection. The child receives a proxy mailbox that forwards transport-sensitive operations back to the parent over IPC.

Module​

Functions​

_run_protocol_in_child​

def _run_protocol_in_child(    config: _WorkerProcessConfig,    result_queue: multiprocessing.Queue[_ChildResultMessage],    request_queue: Any,    response_queue: multiprocessing.Queue[_ParentResponse] | None = None,    modeller_ready_event: multiprocessing.synchronize.Event | None = None,    heartbeat: Heartbeat | None = None,):

Target function for the child worker multiprocessing.Process.

Must be a plain module-level function (not a method) so it can be pickled by the spawn start method.

All exceptions are caught and placed on result_queue; this function must never raise.

Arguments

  • config: Serialisable configuration for reconstructing the worker.
  • result_queue: Queue for sending results back to the parent process.
  • request_queue: Queue for child requests that the parent broker fulfills, or the legacy modeller_ready_event when called via older test wrappers.
  • response_queue: Queue for parent responses to child IPC requests. Optional for backwards-compatible test wrappers.
  • modeller_ready_event: Cross-process event set by the parent when the modeller's TASK_START message is received. A daemon thread in this child watches it and sets mailbox.modeller_ready. Optional for backwards-compatible test wrappers.
  • heartbeat: Shared runtime state reserved for child liveness tracking in later milestones.

_run_protocol_in_child_target​

def _run_protocol_in_child_target(    config: _WorkerProcessConfig,    result_queue: multiprocessing.Queue[_ChildResultMessage],    request_queue: Any,    response_queue: multiprocessing.Queue[_ParentResponse] | None = None,    modeller_ready_event: multiprocessing.synchronize.Event | None = None,    heartbeat: Heartbeat | None = None,) ‑> Any:

Target function for the child worker multiprocessing.Process.

Must be a plain module-level function (not a method) so it can be pickled by the spawn start method.

All exceptions are caught and placed on result_queue; this function must never raise.

Arguments

  • config: Serialisable configuration for reconstructing the worker.
  • result_queue: Queue for sending results back to the parent process.
  • request_queue: Queue for child requests that the parent broker fulfills, or the legacy modeller_ready_event when called via older test wrappers.
  • response_queue: Queue for parent responses to child IPC requests. Optional for backwards-compatible test wrappers.
  • modeller_ready_event: Cross-process event set by the parent when the modeller's TASK_START message is received. A daemon thread in this child watches it and sets mailbox.modeller_ready. Optional for backwards-compatible test wrappers.
  • heartbeat: Shared runtime state reserved for child liveness tracking in later milestones.

Classes​

ProtocolTaskBatchRun​

class ProtocolTaskBatchRun():

Hook to report dataset statistics before and after protocol run.

Initialise the hook.

Methods​


on_run_end​

def on_run_end(    self,    protocol: bitfount.federated.protocols.base._BaseProtocol,    context: TaskContext,    *args: Any,    **kwargs: Any,) ‑> None:

Runs after protocol run to report dataset statistics.

SaveFailedFilesToDatabase​

class SaveFailedFilesToDatabase():

Hook to save failed files to database.

Initialise the hook.

Methods​


on_resilience_end​

def on_resilience_end(    self,    protocol: bitfount.federated.protocols.base._BaseProtocol,    context: TaskContext,    *args: Any,    **kwargs: Any,) ‑> None:

Runs after protocol run to save failed files to database.

SaveResultsToDatabase​

class SaveResultsToDatabase():

Hook to save protocol results to database.

Initialise the hook.

Methods​


on_run_end​

def on_run_end(    self,    protocol: bitfount.federated.protocols.base._BaseProtocol,    context: TaskContext,    *args: Any,    **kwargs: Any,) ‑> None:

Runs after protocol run to save results to database.

_ChildResultMessage​

class _ChildResultMessage(    status: bitfount.federated.worker_process.WorkerProcessStatus,    datasource_state: dict[str, typing.Any] | None = None,):

Base class for child-to-parent terminal messages on the result queue.

Variables​

  • static datasource_state : dict[str, typing.Any] | None
  • static status : bitfount.federated.worker_process.WorkerProcessStatus

_WorkerProcessConfig​

class _WorkerProcessConfig(    *,    ms_config: MessageServiceConfig,    modeller_mailbox_id: str,    modeller_name: str,    pod_mailbox_ids: dict[str, str],    task_id: str,    aes_encryption_key: bytes,    datasource: BaseSource,    datasource_name: str,    serialized_protocol: SerializedProtocol,    serialized_protocol_bytes: bytes,    parent_pod_identifier: str,    mailbox_pod_identifier: str,    schema: BitfountSchema,    batched_execution: bool,    test_run: bool,    project_id: str | None,    run_on_new_data_only: bool,    force_rerun_failed_files: bool,    inference_limits: dict[str, InferenceLimits],    model_urls: dict[str, ModelURLs],    pod_vitals: bitfount.federated.pod_vitals._PodVitals | None = None,    pod_dp: DPPodConfig | None = None,    project_db_connector: ProjectDbConnector | None = None,    data_identifier: str | None = None,    secrets: APIKeys | RefreshableJWT | dict[typing.Literal[''codeBlockAnchor[bitfount](/api/bitfount/index)'', 'ehr'], APIKeys | RefreshableJWT] | None = None,    ehr_config: bitfount.federated.types.NextGenEHRConfig | SMARTStandaloneEHRConfig | SMARTBackendEHRConfig | None = None,    supports_mps: bool | None = None,    data_splitter: DatasetSplitter | None = None,    task_hash: str | None = None,    username: str | None = None,    spawn_started_monotonic: float | None = None,    pod_public_keys_pem: dict[str, bytes] | None = None,    private_key_pem: bytes | None = None,    hook_factories: list[collections.abc.Callable[[], None]] = [],):

A serialisable Dataclass that holds config to reinstantiate a Worker process.

All fields are picklable. The only _Worker attribute intentionally absent is mailbox - the subprocess builds a fresh _WorkerMailbox from the fields below.

Variables​

  • static aes_encryption_key : bytes
  • static batched_execution : bool
  • static data_identifier : str | None
  • static datasource_name : str
  • static force_rerun_failed_files : bool
  • static mailbox_pod_identifier : str
  • static modeller_mailbox_id : str
  • static modeller_name : str
  • static parent_pod_identifier : str
  • static pod_mailbox_ids : dict[str, str]
  • static pod_public_keys_pem : dict[str, bytes] | None
  • static pod_vitals : bitfount.federated.pod_vitals._PodVitals | None
  • static private_key_pem : bytes | None
  • static project_id : str | None
  • static run_on_new_data_only : bool
  • static serialized_protocol_bytes : bytes
  • static spawn_started_monotonic : float | None
  • static supports_mps : bool | None
  • static task_hash : str | None
  • static task_id : str
  • static test_run : bool
  • static username : str | None