MCPcopy Create free account
hub / github.com/omkarcloud/botasaurus / MasterExecutor

Class MasterExecutor

botasaurus_server/botasaurus_server/master_executor.py:6–59  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

4from .k8s import K8s
5
6class MasterExecutor(TaskExecutor):
7
8 def load(self):
9 self.lock = Lock()
10 self.k8s = K8s()
11 # Initialize capacities
12 self.current_capacity = {"request": 0, "browser": 0, "task": 0}
13 self.max_capacity = self.calculate_max_capacity()
14
15 def get_max_running_count(self, scraper_type):
16 return self.max_capacity[scraper_type]
17
18 def get_current_running_count(self, scraper_type):
19 return self.current_capacity[scraper_type]
20
21 def run_task_and_update_state(self, scraper_type, task_json):
22 node = self.k8s.get_min_capacity_node(scraper_type)
23
24 # Send task to the selected node
25 self.k8s.run_task_on_node(task_json, node['node_name'])
26
27 # Update capacities
28 self.increment_master_capacity(node, scraper_type)
29
30 def increment_master_capacity(self, node , scraper_type):
31 with self.lock:
32 node['current_capacity'][scraper_type] += 1
33 self.current_capacity[scraper_type] += 1
34
35 def decrement_master_capacity(self, node , scraper_type):
36 with self.lock:
37 node['current_capacity'][scraper_type] -= 1
38 self.current_capacity[scraper_type] -= 1
39
40 def calculate_max_capacity(self):
41 limit = Server.get_rate_limit()
42 replicas = len(self.k8s.nodes)
43 return {"request": limit["request"] * replicas,"task": limit["task"] * replicas, "browser": limit["browser"] * replicas}
44
45 def on_success(self, task_id, task_type, task_result, node_name,scraper_name, data):
46 node = self.k8s.get_node(node_name)
47 self.decrement_master_capacity(node, task_type)
48
49 # Further processing of task_result can be done here
50 self.mark_task_as_success(task_id, task_result, Server.cache,scraper_name, data)
51 self.update_parent_task(task_id, task_result)
52
53 def on_failure(self, task_id, task_type, task_result, node_name):
54 node = self.k8s.get_node(node_name)
55 self.decrement_master_capacity(node, task_type)
56
57 # Further processing of task_result can be done here
58 self.mark_task_as_failure(task_id, task_result)
59 self.update_parent_task(task_id, [])

Callers 1

executor.pyFile · 0.70

Calls

no outgoing calls

Tested by

no test coverage detected