SDK & ORborus update
This commit is contained in:
@@ -685,6 +685,7 @@ class AppBase:
|
|||||||
if not "User-Agent" in headers:
|
if not "User-Agent" in headers:
|
||||||
headers["User-Agent"] = "Shuffle App"
|
headers["User-Agent"] = "Shuffle App"
|
||||||
|
|
||||||
|
self.logger.info(f"[DEBUG][{self.current_execution_id}] Starting to send result to {url}")
|
||||||
try:
|
try:
|
||||||
finished = False
|
finished = False
|
||||||
ret = {}
|
ret = {}
|
||||||
@@ -702,8 +703,8 @@ class AppBase:
|
|||||||
proxies=self.proxy_config,
|
proxies=self.proxy_config,
|
||||||
)
|
)
|
||||||
|
|
||||||
#self.logger.info(f"""[DEBUG] Successful result request: Status= {ret.status_code} (break on 200/201) & Action status: {action_result["status"]}. Response= {ret.text}""")
|
|
||||||
if ret.status_code == 200 or ret.status_code == 201:
|
if ret.status_code == 200 or ret.status_code == 201:
|
||||||
|
self.logger.info(f"[DEBUG][{self.current_execution_id}] Successful send_result request: Status= {ret.status_code} & Action status: {action_result['status']}. Response= {ret.text}")
|
||||||
finished = True
|
finished = True
|
||||||
break
|
break
|
||||||
else:
|
else:
|
||||||
@@ -1258,6 +1259,7 @@ class AppBase:
|
|||||||
# self.send_result(action_result, headers, stream_path)
|
# self.send_result(action_result, headers, stream_path)
|
||||||
# return
|
# return
|
||||||
|
|
||||||
|
self.logger.info(f"[DEBUG][{self.current_execution_id}] Pre function with {len(param_multiplier)} multipliers")
|
||||||
for subparams in param_multiplier:
|
for subparams in param_multiplier:
|
||||||
#self.logger.info(f"SUBPARAMS IN MULTI: {subparams}")
|
#self.logger.info(f"SUBPARAMS IN MULTI: {subparams}")
|
||||||
tmp = ""
|
tmp = ""
|
||||||
@@ -1355,6 +1357,8 @@ class AppBase:
|
|||||||
#ret = ret[0]
|
#ret = ret[0]
|
||||||
self.logger.info("[DEBUG] DONT make list of 1 into 0!!")
|
self.logger.info("[DEBUG] DONT make list of 1 into 0!!")
|
||||||
|
|
||||||
|
self.logger.info(f"[DEBUG][%s] Done with execution recursion %d times" % (self.current_execution_id, len(param_multiplier)))
|
||||||
|
|
||||||
#self.logger.info("Return from execution: %s" % ret)
|
#self.logger.info("Return from execution: %s" % ret)
|
||||||
if ret == None:
|
if ret == None:
|
||||||
results.append("")
|
results.append("")
|
||||||
@@ -1923,14 +1927,14 @@ class AppBase:
|
|||||||
"""
|
"""
|
||||||
Generate strings contained in nested (), indexing i = level
|
Generate strings contained in nested (), indexing i = level
|
||||||
"""
|
"""
|
||||||
if len(re.findall("\(", string)) == len(re.findall("\)", string)):
|
if len(re.findall('(', string)) == len(re.findall(')', string)):
|
||||||
LeftRightIndex = [x for x in zip(
|
LeftRightIndex = [x for x in zip(
|
||||||
[Left.start()+1 for Left in re.finditer('\(', string)],
|
[Left.start()+1 for Left in re.finditer('(', string)],
|
||||||
reversed([Right.start() for Right in re.finditer('\)', string)]))]
|
reversed([Right.start() for Right in re.finditer(')', string)]))]
|
||||||
|
|
||||||
elif len(re.findall("\(", string)) > len(re.findall("\)", string)):
|
elif len(re.findall('(', string)) > len(re.findall(')', string)):
|
||||||
return parse_nested_param(string + ')', level)
|
return parse_nested_param(string + ')', level)
|
||||||
elif len(re.findall("\(", string)) < len(re.findall("\)", string)):
|
elif len(re.findall('(', string)) < len(re.findall('"', string)):
|
||||||
return parse_nested_param('(' + string, level)
|
return parse_nested_param('(' + string, level)
|
||||||
else:
|
else:
|
||||||
return 'Failed to parse params'
|
return 'Failed to parse params'
|
||||||
@@ -3388,7 +3392,6 @@ class AppBase:
|
|||||||
elif callable(func):
|
elif callable(func):
|
||||||
try:
|
try:
|
||||||
if len(action["parameters"]) < 1:
|
if len(action["parameters"]) < 1:
|
||||||
#result = await func()
|
|
||||||
result = func()
|
result = func()
|
||||||
else:
|
else:
|
||||||
# Potentially parse JSON here
|
# Potentially parse JSON here
|
||||||
|
|||||||
@@ -101,9 +101,10 @@ var swarmConfig = os.Getenv("SHUFFLE_SWARM_CONFIG")
|
|||||||
var swarmNetworkName = os.Getenv("SHUFFLE_SWARM_NETWORK_NAME")
|
var swarmNetworkName = os.Getenv("SHUFFLE_SWARM_NETWORK_NAME")
|
||||||
var orborusLabel = os.Getenv("SHUFFLE_ORBORUS_LABEL")
|
var orborusLabel = os.Getenv("SHUFFLE_ORBORUS_LABEL")
|
||||||
var memcached = os.Getenv("SHUFFLE_MEMCACHED")
|
var memcached = os.Getenv("SHUFFLE_MEMCACHED")
|
||||||
var tenzirUrl = os.Getenv("SHUFFLE_TENZIR_URL")
|
|
||||||
var apiKey = os.Getenv("AUTH_FOR_ORBORUS")
|
var apiKey = os.Getenv("AUTH_FOR_ORBORUS")
|
||||||
|
|
||||||
|
var pipelineUrl = os.Getenv("SHUFFLE_PIPELINE_URL")
|
||||||
|
|
||||||
var executionIds = []string{}
|
var executionIds = []string{}
|
||||||
var namespacemade = false // For K8s
|
var namespacemade = false // For K8s
|
||||||
|
|
||||||
@@ -1807,9 +1808,22 @@ func main() {
|
|||||||
log.Printf("[WARNING] Defaulting to environment name %s. Set environment variable ENVIRONMENT_NAME to change. This should be the same as in the frontend action.", environment)
|
log.Printf("[WARNING] Defaulting to environment name %s. Set environment variable ENVIRONMENT_NAME to change. This should be the same as in the frontend action.", environment)
|
||||||
}
|
}
|
||||||
|
|
||||||
if tenzirUrl == "" {
|
if pipelineUrl == "" {
|
||||||
tenzirUrl = "http://localhost:5160"
|
pipelineUrl = "http://localhost:5160"
|
||||||
log.Printf("[WARNING] SHUFFLE_TENZIR_URL not set, falling back to default URL: %s",tenzirUrl)
|
|
||||||
|
// Find the IP in baseUrl. Base format is http://<ip>:<port>
|
||||||
|
if baseUrl != "" {
|
||||||
|
urlSplit := strings.Split(baseUrl, "://")
|
||||||
|
if len(urlSplit) > 1 {
|
||||||
|
// Find the IP
|
||||||
|
ipSplit := strings.Split(urlSplit[1], ":")
|
||||||
|
if len(ipSplit) > 0 {
|
||||||
|
pipelineUrl = fmt.Sprintf("http://%s:5160", ipSplit[0])
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
log.Printf("[WARNING] SHUFFLE_PIPELINE_URL not set, falling back to default URL: %s. If BASE_URL is set, we use the external IP for that",pipelineUrl)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -2035,7 +2049,7 @@ func main() {
|
|||||||
|
|
||||||
err := deployTenzirNode()
|
err := deployTenzirNode()
|
||||||
if err != nil{
|
if err != nil{
|
||||||
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
|
log.Printf("[ERROR] Failed to deploy CATEGORY UPDATE, reason: %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
err = handleFileCategoryChange()
|
err = handleFileCategoryChange()
|
||||||
@@ -2049,7 +2063,7 @@ func main() {
|
|||||||
fileName := incRequest.ExecutionArgument
|
fileName := incRequest.ExecutionArgument
|
||||||
err := deployTenzirNode()
|
err := deployTenzirNode()
|
||||||
if err != nil{
|
if err != nil{
|
||||||
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
|
log.Printf("[ERROR] Failed to deploy DISABLE SIGMA FILE, reason: %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
err = disableRule(fileName)
|
err = disableRule(fileName)
|
||||||
@@ -2063,7 +2077,7 @@ func main() {
|
|||||||
fileName := incRequest.ExecutionArgument
|
fileName := incRequest.ExecutionArgument
|
||||||
err := deployTenzirNode()
|
err := deployTenzirNode()
|
||||||
if err != nil{
|
if err != nil{
|
||||||
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
|
log.Printf("[ERROR] Failed to deploy ENABLE SIGMA FILE, reason: %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
err = enableRule(fileName)
|
err = enableRule(fileName)
|
||||||
@@ -2074,10 +2088,9 @@ func main() {
|
|||||||
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
|
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
|
||||||
|
|
||||||
} else if incRequest.Type == "DISABLE_SIGMA_FOLDER" {
|
} else if incRequest.Type == "DISABLE_SIGMA_FOLDER" {
|
||||||
|
|
||||||
err := deployTenzirNode()
|
err := deployTenzirNode()
|
||||||
if err != nil{
|
if err != nil{
|
||||||
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
|
log.Printf("[ERROR] Failed to deploy DISABLE SIGMA FOLDER, reason: %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
err = removeAllFiles()
|
err = removeAllFiles()
|
||||||
@@ -2090,7 +2103,7 @@ func main() {
|
|||||||
log.Printf("[INFO] Got job to start tenzir")
|
log.Printf("[INFO] Got job to start tenzir")
|
||||||
err := deployTenzirNode()
|
err := deployTenzirNode()
|
||||||
if err != nil{
|
if err != nil{
|
||||||
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
|
log.Printf("[ERROR] Failed to deploy the pipeline, reason: %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
|
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
|
||||||
@@ -2446,7 +2459,7 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error {
|
|||||||
|
|
||||||
err := deployTenzirNode()
|
err := deployTenzirNode()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
|
log.Printf("[ERROR] Failed to deploy the pipeline, reason: %s", err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2557,7 +2570,7 @@ func deployTenzirNode() error {
|
|||||||
// Check if image exists
|
// Check if image exists
|
||||||
_, _, err := dockercli.ImageInspectWithRaw(ctx, imageName)
|
_, _, err := dockercli.ImageInspectWithRaw(ctx, imageName)
|
||||||
if dockerclient.IsErrNotFound(err) {
|
if dockerclient.IsErrNotFound(err) {
|
||||||
log.Printf("[DEBUG] pulling image %s", imageName)
|
log.Printf("[DEBUG] Pulling image %s", imageName)
|
||||||
pullOptions := image.PullOptions{}
|
pullOptions := image.PullOptions{}
|
||||||
out, err := dockercli.ImagePull(ctx, imageName, pullOptions)
|
out, err := dockercli.ImagePull(ctx, imageName, pullOptions)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -2713,7 +2726,7 @@ func createNetworkIfNotExists(ctx context.Context, networkName, subnet, gateway
|
|||||||
func checkTenzirNode() error {
|
func checkTenzirNode() error {
|
||||||
retries := 5
|
retries := 5
|
||||||
retryInterval := 3 * time.Second
|
retryInterval := 3 * time.Second
|
||||||
url := fmt.Sprintf("%s/api/v0/ping",tenzirUrl)
|
url := fmt.Sprintf("%s/api/v0/ping",pipelineUrl)
|
||||||
forwardMethod := "POST"
|
forwardMethod := "POST"
|
||||||
|
|
||||||
client := http.Client{}
|
client := http.Client{}
|
||||||
@@ -2739,7 +2752,7 @@ func createPipeline(command, identifier string) (string, error) {
|
|||||||
toBeDeleted := false
|
toBeDeleted := false
|
||||||
pipelineId, err := searchPipeline(identifier)
|
pipelineId, err := searchPipeline(identifier)
|
||||||
|
|
||||||
url := fmt.Sprintf("%s/api/v0/pipeline/create", tenzirUrl)
|
url := fmt.Sprintf("%s/api/v0/pipeline/create", pipelineUrl)
|
||||||
forwardMethod := "POST"
|
forwardMethod := "POST"
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -2849,7 +2862,7 @@ func createPipeline(command, identifier string) (string, error) {
|
|||||||
|
|
||||||
func updatePipelineState(command, pipelineId, action string) (string, error) {
|
func updatePipelineState(command, pipelineId, action string) (string, error) {
|
||||||
|
|
||||||
url := fmt.Sprintf("%s/api/v0/pipeline/update", tenzirUrl)
|
url := fmt.Sprintf("%s/api/v0/pipeline/update", pipelineUrl)
|
||||||
forwardMethod := "POST"
|
forwardMethod := "POST"
|
||||||
|
|
||||||
requestBody := map[string]interface{}{
|
requestBody := map[string]interface{}{
|
||||||
@@ -2919,7 +2932,7 @@ func deletePipeline(pipelineId string) error {
|
|||||||
"id": pipelineId,
|
"id": pipelineId,
|
||||||
}
|
}
|
||||||
|
|
||||||
url := fmt.Sprintf("%s/api/v0/pipeline/delete", tenzirUrl)
|
url := fmt.Sprintf("%s/api/v0/pipeline/delete", pipelineUrl)
|
||||||
forwardMethod := "POST"
|
forwardMethod := "POST"
|
||||||
|
|
||||||
requestBodyJSON, err := json.Marshal(requestBody)
|
requestBodyJSON, err := json.Marshal(requestBody)
|
||||||
@@ -2967,7 +2980,7 @@ func searchPipeline(identifier string) (string, error) {
|
|||||||
|
|
||||||
var reqBody []byte
|
var reqBody []byte
|
||||||
|
|
||||||
url := fmt.Sprintf("%s/api/v0/pipeline/list", tenzirUrl)
|
url := fmt.Sprintf("%s/api/v0/pipeline/list", pipelineUrl)
|
||||||
|
|
||||||
resp, err := http.Post(url, "application/json", bytes.NewBuffer(reqBody))
|
resp, err := http.Post(url, "application/json", bytes.NewBuffer(reqBody))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
Reference in New Issue
Block a user