evgeniy778 commited on
Commit
12afd49
·
verified ·
1 Parent(s): 367243f

Add Dispatcher Version 1.0

Browse files
Files changed (1) hide show
  1. dispatcher.py +152 -0
dispatcher.py ADDED
@@ -0,0 +1,152 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ # =====================================================
2
+ # Apckeyl Framework
3
+ # Version 1.0
4
+ # dispatcher.py
5
+ # =====================================================
6
+
7
+ """
8
+ Apckeyl Dispatcher.
9
+
10
+ Version 1.0
11
+
12
+ Dispatcher отвечает за передачу задачи
13
+ из Control Plane в Compute Plane.
14
+
15
+ На этом этапе Dispatcher НЕ выполняет
16
+ реальный HTTP/API-запрос.
17
+
18
+ Он только:
19
+
20
+ - принимает задачу;
21
+ - проверяет выбранный Compute Module;
22
+ - формирует dispatch request;
23
+ - переводит задачу в статус dispatched.
24
+
25
+ Реальное подключение внешнего Space
26
+ будет добавлено через Compute Adapter.
27
+ """
28
+
29
+ from task_manager import TaskManager
30
+
31
+
32
+ # =====================================================
33
+ # Dispatcher
34
+ # =====================================================
35
+
36
+ class Dispatcher:
37
+
38
+ def __init__(
39
+ self,
40
+ task_manager=None,
41
+ ):
42
+
43
+ if task_manager is None:
44
+
45
+ task_manager = TaskManager()
46
+
47
+ self.task_manager = task_manager
48
+
49
+
50
+ # =================================================
51
+ # Prepare Dispatch
52
+ # =================================================
53
+
54
+ def prepare_dispatch(
55
+ self,
56
+ task_id,
57
+ ):
58
+
59
+ task = self.task_manager.get_task(
60
+ task_id
61
+ )
62
+
63
+ if task is None:
64
+
65
+ raise KeyError(
66
+ f"Unknown task: {task_id}"
67
+ )
68
+
69
+
70
+ # -------------------------------------------------
71
+ # Task must be created or queued
72
+ # -------------------------------------------------
73
+
74
+ allowed_statuses = {
75
+
76
+ "created",
77
+
78
+ "queued",
79
+ }
80
+
81
+ if task["status"] not in allowed_statuses:
82
+
83
+ raise RuntimeError(
84
+
85
+ f"Task {task_id} cannot be "
86
+ f"dispatched from status "
87
+ f"'{task['status']}'"
88
+ )
89
+
90
+
91
+ # -------------------------------------------------
92
+ # Build dispatch request
93
+ # -------------------------------------------------
94
+
95
+ dispatch_request = {
96
+
97
+ "task_id": task[
98
+ "task_id"
99
+ ],
100
+
101
+ "task_type": task[
102
+ "task_type"
103
+ ],
104
+
105
+ "module_id": task[
106
+ "module_id"
107
+ ],
108
+
109
+ "module_name": task[
110
+ "module_name"
111
+ ],
112
+
113
+ "payload": task[
114
+ "payload"
115
+ ],
116
+ }
117
+
118
+
119
+ # -------------------------------------------------
120
+ # Update task status
121
+ # -------------------------------------------------
122
+
123
+ self.task_manager.update_status(
124
+
125
+ task_id,
126
+
127
+ "dispatched",
128
+ )
129
+
130
+
131
+ return dispatch_request
132
+
133
+
134
+ # =================================================
135
+ # Get Task
136
+ # =================================================
137
+
138
+ def get_task(
139
+ self,
140
+ task_id,
141
+ ):
142
+
143
+ return self.task_manager.get_task(
144
+ task_id
145
+ )
146
+
147
+
148
+ # =====================================================
149
+ # Default Dispatcher
150
+ # =====================================================
151
+
152
+ dispatcher = Dispatcher()