MCPcopy Create free account
hub / github.com/closeio/tasktiger / Worker

Class Worker

tasktiger/worker.py:57–1038  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

55
56
57class Worker:
58 def __init__(
59 self,
60 tiger: "TaskTiger",
61 queues: Optional[List[str]] = None,
62 exclude_queues: Optional[List[str]] = None,
63 single_worker_queues: Optional[List[str]] = None,
64 max_workers_per_queue: Optional[int] = None,
65 store_tracebacks: Optional[bool] = None,
66 executor_class: Optional[Type[Executor]] = None,
67 ) -> None:
68 """
69 Internal method to initialize a worker.
70 """
71
72 self.tiger = tiger
73 bound_logger = tiger.log.bind(pid=os.getpid())
74 assert isinstance(bound_logger, BoundLogger)
75 self.log = bound_logger
76
77 self.connection = tiger.connection
78 self.scripts = tiger.scripts
79 self.config = tiger.config
80 self._key = tiger._key
81 self._did_work = True
82 self._last_task_check = 0.0
83 self.stats_thread: Optional[StatsThread] = None
84 self.id = str(uuid.uuid4())
85
86 if executor_class is None:
87 executor_class = ForkExecutor
88 self.executor = executor_class(self)
89
90 if queues:
91 self.only_queues = set(queues)
92 elif self.config["ONLY_QUEUES"]:
93 self.only_queues = set(self.config["ONLY_QUEUES"])
94 else:
95 self.only_queues = set()
96
97 if exclude_queues:
98 self.exclude_queues = set(exclude_queues)
99 elif self.config["EXCLUDE_QUEUES"]:
100 self.exclude_queues = set(self.config["EXCLUDE_QUEUES"])
101 else:
102 self.exclude_queues = set()
103
104 if single_worker_queues:
105 self.single_worker_queues = set(single_worker_queues)
106 elif self.config["SINGLE_WORKER_QUEUES"]:
107 self.single_worker_queues = set(self.config["SINGLE_WORKER_QUEUES"])
108 else:
109 self.single_worker_queues = set()
110
111 if max_workers_per_queue:
112 self.max_workers_per_queue: Optional[int] = max_workers_per_queue
113 else:
114 self.max_workers_per_queue = None

Calls

no outgoing calls