Skip to content
Merged
Show file tree
Hide file tree
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
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
from naeural_core.business.default.web_app.fast_api_web_app import FastApiWebAppPlugin as BasePlugin



_CONFIG = {
**BasePlugin.CONFIG,

Expand All @@ -9,7 +10,7 @@
"PROCESS_DELAY": 0,

# 'PORT': 5001,
'ASSETS': 'extensions/business/fastapi/_ai4everyone',
'ASSETS': 'extensions/business/ai4e',
'JINJA_ARGS': {
'html_files': [
{
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
from ratio1 import Session, Pipeline, Instance
# from ratio1 import Session, Pipeline, Instance
from datetime import datetime
import numpy as np

Expand Down Expand Up @@ -63,9 +63,9 @@ def get_job_config(


class Job:
def __init__(self, owner, session: Session, job_id: str, node_id: str, pipeline: str, signature: str, instance: str):
def __init__(self, owner, job_id: str, node_id: str, pipeline: str, signature: str, instance: str):
self.owner = owner
self.session = session
# self.session = session
self.job_id = job_id
# TODO: instead of single values for instance_id and the rest, use a dict to
# have a set of properties for each subtask(each signature)
Expand Down Expand Up @@ -249,7 +249,6 @@ def maybe_update_data(self, data: dict, pipeline: str, signature: str):
self.train_status['total_grid_iterations'] = train_data.get('NR_ALL_GRID_ITER', self.train_status['total_grid_iterations'])
if 'TRAIN_FINAL' in data.keys():
self.train_final = data.get('TRAIN_FINAL', self.train_final)
self.deploy_job({})
# endif training finished

# endif status for training
Expand Down Expand Up @@ -333,24 +332,24 @@ def get_train_status(self):
# 'full_last_payload': self.train_status_full_payload
}

def __get_pipeline_and_instance(self, node=None, pipeline=None, signature=None, instance_id=None):
node = node or self.node_id
active_pipelines = self.session.get_active_pipelines(node=node)
if active_pipelines is None:
return None, None, f"Node_ID {node} not found"
pipeline = pipeline or self.pipeline
curr_pipeline = self.session.attach_to_pipeline(node=node, name=pipeline)
if curr_pipeline is None:
return None, None, f"Pipeline {pipeline} not found on {node}"
signature = signature or self.signature
instance_id = instance_id or self.instance_id
curr_instance = curr_pipeline.attach_to_plugin_instance(
signature=signature,
instance_id=instance_id
)
if curr_instance is None:
return curr_pipeline, None, f"Instance {instance_id} of {signature} not found on {pipeline}"
return curr_pipeline, curr_instance, ""
# def __get_pipeline_and_instance(self, node=None, pipeline=None, signature=None, instance_id=None):
# node = node or self.node_id
# active_pipelines = self.session.get_active_pipelines(node=node)
# if active_pipelines is None:
# return None, None, f"Node_ID {node} not found"
# pipeline = pipeline or self.pipeline
# curr_pipeline = self.session.attach_to_pipeline(node=node, name=pipeline)
# if curr_pipeline is None:
# return None, None, f"Pipeline {pipeline} not found on {node}"
# signature = signature or self.signature
# instance_id = instance_id or self.instance_id
# curr_instance = curr_pipeline.attach_to_plugin_instance(
# signature=signature,
# instance_id=instance_id
# )
# if curr_instance is None:
# return curr_pipeline, None, f"Instance {instance_id} of {signature} not found on {pipeline}"
# return curr_pipeline, curr_instance, ""

def send_instance_command(self, node=None, pipeline=None, signature=None, instance_id=None, **kwargs):
node = node or self.node_id
Expand Down Expand Up @@ -443,7 +442,7 @@ def publish_labels(self):
return self.send_instance_command(publish=True)

def start_train(self, body: dict):
self.session.P(f'Starting training for {self.job_id} on {self.node_id}...')
self.owner.P(f'Starting training for {self.job_id} on {self.node_id}...')
success, err_msg = self.send_instance_update(
config={
'START_TRAINING': True,
Expand Down Expand Up @@ -495,38 +494,58 @@ def deploy_configs(self, lst_allowed, deploy_ngrok_edge_label):
'MODEL_INSTANCE_ID': self.train_meta['MODEL_INSTANCE_ID']
},
'DESCRIPTION': self.description,
'OBJECTIVE_NAME': self.objective_name
'OBJECTIVE_NAME': self.objective_name,
"INSTANCE_ID": self.job_id,
}
pipeline = self.session.create_or_attach_to_pipeline(
node=chosen_node,
name=f'deploy_{self.job_id}',
data_source='ON_DEMAND_INPUT',
# plugins=[det_plugin_config]
)
instance = pipeline.create_or_attach_to_plugin_instance(
signature='ai4e_custom_inference_agent',
instance_id=self.job_id,
config=instance_config
inference_pipeline_config = {
"TYPE": "ON_DEMAND_INPUT",
"NAME": f"deploy_{self.job_id}",
"PLUGINS": [
{
"INSTANCES": [
instance_config
],
"SIGNATURE": "ai4e_custom_inference_agent",
}
]
}
self.owner.cmadpi_start_pipeline(
config=inference_pipeline_config,
node_address=chosen_node,
)
pipeline.deploy()
# END DETECTION PIPELINE
# START FASTAPI PIPELINE
fastapi_instance_id = "AI4EveryoneDeploys"
fastapi_instance_config = {
"RESPONSE_FORMAT": self.owner.cfg_response_format,
"NGROK_EDGE_LABEL": deploy_ngrok_edge_label,
"INSTANCE_ID": fastapi_instance_id,
}
pipeline = self.session.create_or_attach_to_pipeline(
node=chosen_node,
name=f'AI4EveryoneDeploys',
data_source='VOID',
)
instance = pipeline.create_or_attach_to_plugin_instance(
signature="AI4E_DEPLOY",
instance_id='AI4EveryoneDeploys',
config=fastapi_instance_config
fastapi_pipeline_name = f"AI4EveryoneDeploys"
fastapi_signature = "AI4E_DEPLOY"
fastapi_pipeline_config = {
"NAME": fastapi_pipeline_name,
"TYPE": "VOID", # Change to custom NetworkListener
"PLUGINS": [
{
"INSTANCES": [
{
fastapi_instance_config
}
],
"SIGNATURE": fastapi_signature,
}
]
}
self.owner.cmadpi_start_pipeline(
config=fastapi_pipeline_config,
node_address=chosen_node,
)
pipeline.deploy()
instance.send_instance_command(
self.send_instance_command(
node=chosen_node,
pipeline=fastapi_pipeline_name,
signature=fastapi_signature,
instance_id=fastapi_instance_id,
command='REGISTER',
command_params={
'CONFIG': instance_config
Expand All @@ -539,8 +558,9 @@ def deploy_job(self, body):
if self.deployed:
return True, "Job already deployed"
self.started_deploying = True
lst_allowed = self.session.get_allowed_nodes()
self.session.P(f"Allowed nodes: {lst_allowed}")
# lst_allowed = self.session.get_allowed_nodes()
lst_allowed = self.owner.netmon.accessible_nodes()
self.owner.P(f"Allowed nodes: {lst_allowed}")
if len(lst_allowed) == 0:
self.started_deploying = False
return False, "No node available at the moment."
Expand Down
Loading