diff --git a/functions/onprem/orborus/go.mod b/functions/onprem/orborus/go.mod index 4f3177e8..1c6a7a5a 100644 --- a/functions/onprem/orborus/go.mod +++ b/functions/onprem/orborus/go.mod @@ -4,7 +4,7 @@ go 1.24.0 toolchain go1.24.4 -//replace github.com/shuffle/shuffle-shared => ../../../../shuffle-shared +replace github.com/shuffle/shuffle-shared => ../../../../shuffle-shared require ( github.com/docker/docker v28.3.3+incompatible diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index 08c1557b..457b36bc 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -118,7 +118,7 @@ var pipelineApikey = os.Getenv("SHUFFLE_PIPELINE_AUTH") var pipelineUrl = os.Getenv("SHUFFLE_PIPELINE_URL") var executionIds = []string{} -var pipelines = []shuffle.PipelineInfoMini{} +var pipelines = []shuffle.PipelineInfo{} var namespacemade = false // For K8s var skipPipelineMount = false var tenzirDisabled = false @@ -2495,7 +2495,7 @@ func main() { for _, incRequest := range executionRequests.Data { // Looking for specific jobs - if incRequest.Type == "PIPELINE_CREATE" || incRequest.Type == "PIPELINE_START" || incRequest.Type == "PIPELINE_STOP" || incRequest.Type == "PIPELINE_DELETE" { + if incRequest.Type == "PIPELINE_CREATE" || incRequest.Type == "PIPELINE_START" || incRequest.Type == "PIPELINE_STOP" || incRequest.Type == "PIPELINE_DELETE" || incRequest.Type == "PIPELINE_UPDATE" { log.Printf("[INFO] Handling pipeline request from backend: '%s' with argument '%s'", incRequest.Type, incRequest.ExecutionArgument) //os.Setenv("SHUFFLE_SKIP_PIPELINES", "false") @@ -2521,8 +2521,8 @@ func main() { } else if incRequest.Type == "CATEGORY_UPDATE" { os.Setenv("SHUFFLE_SKIP_PIPELINES", "false") - tenzirDisabled = false + tenzirDisabled = false err = handleFileCategoryChange() if err != nil { log.Printf("[ERROR] Failed to download the file category: %s", err) @@ -2576,7 +2576,7 @@ func main() { if strings.Contains(fmt.Sprintf("%s", err), "node available") { // Disabling until UI is updated //os.Setenv("SHUFFLE_SKIP_PIPELINES", "true") - tenzirDisabled = true + //tenzirDisabled = true log.Printf("[ERROR] Failed to start tenzir, reason: %s", err) err = shuffle.CreateOrgNotification( @@ -2824,7 +2824,7 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { // no need of execution arguments for STOP and DELETE if (incRequest.Type != "PIPELINE_STOP" && incRequest.Type != "PIPELINE_DELETE") && len(incRequest.ExecutionArgument) == 0 { - log.Printf("[ERROR] No execution argument found for pipeline create. Skipping") + log.Printf("[ERROR] No execution argument found for pipeline type %s. Skipping", incRequest.Type) return errors.New("no execution argument found for pipeline create. Skipping") } @@ -2835,6 +2835,7 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { } command := incRequest.ExecutionArgument + pipelines = []shuffle.PipelineInfo{} if incRequest.Type == "PIPELINE_CREATE" { log.Printf("[INFO] Should delete -> recreate new pipeline with id %#v", identifier) //err := deployPipeline(image, identifier, command) @@ -2880,20 +2881,24 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { if err != nil { if err.Error() == "no existing pipeline found with name" { log.Printf("[INFO] Starting a new pipeline with command '%s' and identifier '%s'", command, identifier) - _, CreateErr := createPipeline(command, identifier) - return CreateErr + var createErr error + pipelineId, createErr = createPipeline(command, identifier) + if createErr != nil { + return createErr + } + } else { + log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err) + return err } - - log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err) - return err } + log.Printf("[INFO] Starting existing pipeline with ID %s", pipelineId) _, err = updatePipelineState(command, pipelineId, "start") if err != nil { log.Printf("[ERROR] Failed to start Pipeline: %s reason:%s ", pipelineId, err) return err } else { - log.Printf("[INFO] Successfully started the Pipeline: %s", pipelineId) + log.Printf("[INFO] Successfully started pipeline: %s", pipelineId) } } else { @@ -3085,8 +3090,9 @@ func createAndStartTenzirNode(ctx context.Context, containerName, imageName stri { Type: "bind", Source: tenzirStorageFolder, - Target: "/var/lib/tenzir/", + Target: "/tmp", }, + /* { Type: "bind", Source: tenzirStorageFolder, @@ -3097,6 +3103,7 @@ func createAndStartTenzirNode(ctx context.Context, containerName, imageName stri Source: tenzirStorageFolder, Target: "/var/cache/tenzir/", }, + */ }, VolumeDriver: "local", RestartPolicy: container.RestartPolicy{ @@ -3104,6 +3111,12 @@ func createAndStartTenzirNode(ctx context.Context, containerName, imageName stri }, } + if os.Getenv("SHUFFLE_DISABLE_SYSLOG") == "true" { + hostConfig.PortBindings = nat.PortMap{ + "5160/tcp": []nat.PortBinding{{HostPort: "5160"}}, + } + } + if skipPipelineMount { hostConfig.Mounts = []mount.Mount{} } @@ -3289,21 +3302,18 @@ func createPipeline(command, identifier string) (string, error) { //command = "from file /var/lib/tenzir/sysmon_logs.ndjson read json | sigma /var/lib/tenzir/rule.yaml" //command = "from file /var/lib/tenzir/sysmon_logs.ndjson read json | import" + // Make sure to escape them + //if strings.Contains(command, "/") { + // command = strings.ReplaceAll("\\\"", "", command) + // command = strings.ReplaceAll(command, "\"", "") + //} + requestBody := map[string]interface{}{ "definition": command, "name": identifier, "hidden": false, "retry_delay": "500.0ms", - "autostart": map[string]bool{ - //"created": true, - "completed": false, - "failed": false, - }, - "autodelete": map[string]bool{ - "completed": false, - "failed": false, - "stopped": false, - }, + "unstoppable": true, } requestBodyJSON, err := json.Marshal(requestBody) @@ -3339,7 +3349,9 @@ func createPipeline(command, identifier string) (string, error) { } if strings.Contains(string(body), "error") { - log.Printf("[ERROR] Pipeline creation response (%d): %s", resp.StatusCode, string(body)) + log.Printf("[ERROR] Pipeline creation error resp (%d): %s", resp.StatusCode, string(body)) + } else { + log.Printf("[DEBUG] Pipeline creation debug (%d): %s", resp.StatusCode, string(body)) } defer resp.Body.Close() @@ -3365,37 +3377,39 @@ func createPipeline(command, identifier string) (string, error) { return "", errors.New("Pipeline ID not found or empty in the response. See error logs.") } - id := response.ID - return id, nil + return response.ID, nil } func updatePipelineState(command, pipelineId, action string) (string, error) { url := fmt.Sprintf("%s/api/v0/pipeline/update", pipelineUrl) forwardMethod := "POST" - requestBody := map[string]interface{}{ "id": pipelineId, - "definition": command, "action": action, + + /* "autostart": map[string]bool{ "created": true, - "completed": true, - "failed": true, + "completed": false, + "failed": false, }, "autodelete": map[string]bool{ "completed": false, "failed": false, "stopped": false, }, + */ } requestBodyJSON, err := json.Marshal(requestBody) if err != nil { return "", err } - forwardData := bytes.NewBuffer(requestBodyJSON) + log.Printf("[INFO] Updating pipeline %s with action %s to ensure it starts. Body: %s", pipelineId, action, string(requestBodyJSON)) + + forwardData := bytes.NewBuffer(requestBodyJSON) req, err := http.NewRequest( forwardMethod, url, @@ -3478,7 +3492,7 @@ func deletePipeline(pipelineId string) error { log.Printf("[INFO] Pipeline with ID: %s deleted successfully", pipelineId) - pipelines = []shuffle.PipelineInfoMini{} + pipelines = []shuffle.PipelineInfo{} return nil } @@ -3739,7 +3753,7 @@ func removePath(containerName, path string) error { func sendPipelineHealthStatus() (shuffle.LakeConfig, error) { pipelinePayload := shuffle.LakeConfig{ Enabled: false, - Pipelines: []shuffle.PipelineInfoMini{}, + Pipelines: []shuffle.PipelineInfo{}, } if tenzirDisabled { @@ -3751,18 +3765,9 @@ func sendPipelineHealthStatus() (shuffle.LakeConfig, error) { if len(pipelines) == 0 || randint == 0 { pipelineDef, err := listPipelines() - if err == nil { - for _, pipeline := range pipelineDef { - pipelinePayload.Pipelines = append(pipelinePayload.Pipelines, shuffle.PipelineInfoMini{ - ID: pipeline.ID, - Name: pipeline.Name, - Definition: pipeline.Definition, - TotalRuns: pipeline.TotalRuns, - CreatedAt: pipeline.CreatedAt, - }) - } - - pipelines = pipelinePayload.Pipelines + if err == nil || len(pipelines) > 0 { + pipelines = pipelineDef + pipelinePayload.Pipelines = pipelines } } else { pipelinePayload.Pipelines = pipelines @@ -3775,7 +3780,7 @@ func sendPipelineHealthStatus() (shuffle.LakeConfig, error) { log.Printf("[ERROR] Tenzir node connection problem: %s", err) } else { - tenzirDisabled = true + //tenzirDisabled = true log.Printf("[WARNING] Disabling pipelines: %s. You will need to restart the Orborus to fix this.", err) }