mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-22 03:31:02 +02:00
Flow stop
This commit is contained in:
parent
950c624013
commit
19a66ec296
1 changed files with 46 additions and 10 deletions
|
|
@ -64,7 +64,6 @@ class FlowConfig:
|
||||||
async def handle_start_flow(self, msg):
|
async def handle_start_flow(self, msg):
|
||||||
|
|
||||||
def repl_template(tmp):
|
def repl_template(tmp):
|
||||||
print("REPL")
|
|
||||||
return tmp.replace(
|
return tmp.replace(
|
||||||
"{class}", msg.class_name
|
"{class}", msg.class_name
|
||||||
).replace(
|
).replace(
|
||||||
|
|
@ -83,17 +82,11 @@ class FlowConfig:
|
||||||
|
|
||||||
variant = repl_template(variant)
|
variant = repl_template(variant)
|
||||||
|
|
||||||
print(">>", processor, variant)
|
|
||||||
|
|
||||||
print(">>>", v)
|
|
||||||
|
|
||||||
v = {
|
v = {
|
||||||
repl_template(k2): repl_template(v2)
|
repl_template(k2): repl_template(v2)
|
||||||
for k2, v2 in v.items()
|
for k2, v2 in v.items()
|
||||||
}
|
}
|
||||||
|
|
||||||
print("<<<", v)
|
|
||||||
|
|
||||||
if processor in self.config["flows-active"]:
|
if processor in self.config["flows-active"]:
|
||||||
target = json.loads(self.config["flows-active"][processor])
|
target = json.loads(self.config["flows-active"][processor])
|
||||||
else:
|
else:
|
||||||
|
|
@ -104,10 +97,8 @@ class FlowConfig:
|
||||||
|
|
||||||
self.config["flows-active"][processor] = json.dumps(target)
|
self.config["flows-active"][processor] = json.dumps(target)
|
||||||
|
|
||||||
print(cls)
|
|
||||||
|
|
||||||
self.config["flows"][msg.flow_id] = {
|
self.config["flows"][msg.flow_id] = {
|
||||||
"description": cls["description"],
|
"description": msg.description,
|
||||||
"class-name": msg.class_name,
|
"class-name": msg.class_name,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -119,6 +110,51 @@ class FlowConfig:
|
||||||
|
|
||||||
async def handle_stop_flow(self, msg):
|
async def handle_stop_flow(self, msg):
|
||||||
|
|
||||||
|
def repl_template(tmp):
|
||||||
|
return tmp.replace(
|
||||||
|
"{class}", msg.class_name
|
||||||
|
).replace(
|
||||||
|
"{id}", msg.flow_id
|
||||||
|
)
|
||||||
|
|
||||||
|
cls = json.loads(self.config["flow-classes"][msg.class_name])
|
||||||
|
|
||||||
|
plumb = {}
|
||||||
|
|
||||||
|
for kind in ("flow"):
|
||||||
|
|
||||||
|
for k, v in cls[kind].items():
|
||||||
|
|
||||||
|
processor, variant = k.split(":", 1)
|
||||||
|
|
||||||
|
variant = repl_template(variant)
|
||||||
|
|
||||||
|
if processor in self.config["flows-active"]:
|
||||||
|
target = json.loads(self.config["flows-active"][processor])
|
||||||
|
else:
|
||||||
|
target = {}
|
||||||
|
|
||||||
|
if variant in target:
|
||||||
|
del target[variant]
|
||||||
|
|
||||||
|
self.config["flows-active"][processor] = json.dumps(target)
|
||||||
|
|
||||||
|
if msg.flow_id in self.config["flows"]:
|
||||||
|
del self.config["flows"][msg.flow_id]
|
||||||
|
|
||||||
|
await self.config.push()
|
||||||
|
|
||||||
|
return FlowResponse(
|
||||||
|
error = None,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
flow = self.config["flows"][msg.flow_id]
|
flow = self.config["flows"][msg.flow_id]
|
||||||
|
|
||||||
return FlowResponse(
|
return FlowResponse(
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue