Skip to content
Open
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
18 changes: 11 additions & 7 deletions zebrok/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,9 @@ class TaskQueueWorker:
Listens and receives tasks and uses a task runner to execute them
"""

def __init__(self, connection, runner) -> None:
def __init__(
self, connection: BaseSocketConnection, runner: BaseTaskRunner
) -> None:
self.slaves: List[Any] = []
self.connection = connection
self.socket = self.connection.socket
Expand Down Expand Up @@ -79,11 +81,11 @@ def stop(self) -> None:
self.current_slave = 0
self.connection.close()

def add_slave(self, worker) -> None:
def add_slave(self, slave_socket: Any) -> None:
"""
Adds a slave worker to a master worker
Adds a slave worker socket to a master worker
"""
self.slaves.append(worker)
self.slaves.append(slave_socket)


class WorkerInitializer:
Expand Down Expand Up @@ -159,7 +161,9 @@ def _initialize_workers(self) -> None:
master_worker,
)

def _create_master_worker(self, *settings: Any):
def _create_master_worker(
self, *settings: Any
) -> Tuple[BaseSocketConnection, TaskQueueWorker]:
"""
Creates the main worker
"""
Expand All @@ -172,8 +176,8 @@ def _initialize_slave_workers(
max_workers: int,
host: str,
port: int,
master_socket: Any,
master_worker: Any,
master_socket: BaseSocketConnection,
master_worker: TaskQueueWorker,
) -> None:
"""
Creates worker threads as slaves to be associated with the main worker
Expand Down
Loading