Skip to content

Latest commit

 

History

3 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Shuffle Worker

The Shuffle Worker is the workflow orchestration and execution engine for Shuffle. It executes actions across multiple environments (Docker, Kubernetes, Docker Swarm, and Standalone Sandboxes) and is designed to operate in two distinct modes:

  1. Standalone Worker: Run as a container (Docker/K8s) spawned by Orborus, or as a standalone CLI executable (worker.go).
  2. Embedded Library (Backend-Injected): Imported directly by shuffle/backend/go-app as github.com/shuffle/shuffle-worker/pkg, making external worker containers optional.

Architecture & Package Structure

The repository is organized into focused, modular packages:

shuffle-worker/
├── worker.go              # Standalone CLI entrypoint (package main -> worker.RunMain())
├── Dockerfile             # Alpine/scratch container builder
├── go.mod / go.sum
└── pkg/                   # Modular worker engine library (package worker)
    ├── config.go          # Config struct & NewConfigFromEnv()
    ├── types.go           # Shared package types, image caches, and sync primitives
    ├── engine.go          # Worker engine, New(), StartWorkflowExecution(), HandleExecutionResult()
    ├── dag.go             # Graph topology, branch evaluation, subflow barriers, validation
    ├── docker.go          # Docker container deployment (DeployContainer, deployApp)
    ├── k8s.go             # Kubernetes pod deployment and lifecycle
    ├── swarm.go           # Docker Swarm service deployment and auto-scaler
    ├── subflow.go         # Subflow polling, backoff strategy, and stream wrappers
    ├── server.go          # HTTP server, queue handlers, stream endpoints
    └── standalone.go      # Sandboxed script engine (in-process / direct python virtualenvs)

Backend Integration Pattern

Why Dynamic Runtime Checks (No Static Global Flags)

In Shuffle, execution environments are heterogeneous:

  • Some environments target Orborus Docker (standalone worker container spawned via Docker socket).
  • Some environments target Orborus Kubernetes (standalone worker pod spawned in a cluster).
  • Some environments target remote on-premises workers.
  • Some environments target Backend Injection (in-process orchestration without worker containers).

Because a single Shuffle backend handles workflows destined for different runtime locations simultaneously, execution routing cannot rely on a static environment variable (e.g. SHUFFLE_BUILTIN_WORKER). Instead, the backend performs a dynamic runtime check per execution based on the environment configuration (env.Type, env.Executor, or target runtime location).


Step 1: Add Dependency in shuffle/backend/go-app/go.mod

Add the module (or local replace during development):

require (
    github.com/shuffle/shuffle-worker v0.0.1
)

// For local workspace development:
replace github.com/shuffle/shuffle-worker => ../../../shuffle-worker

Step 2: Initialize Embedded Engine in backend/go-app/main.go

Initialize the embedded worker singleton during backend startup with Standalone: false so that it handles orchestration in-process without invoking os.Exit() on completion:

package main

import (
    worker "github.com/shuffle/shuffle-worker/pkg"
)

var embeddedWorker *worker.Worker

func initEmbeddedWorker() {
    cfg := worker.NewConfigFromEnv()
    cfg.Standalone = false // Embedded mode: do not terminate host backend on execution completion
    embeddedWorker = worker.New(cfg)
}

Step 3: Dynamic Dispatch in backend/go-app/walkoff.go (handleExecution)

When a workflow execution is triggered, inspect the target environment and determine whether to queue it for an external Orborus worker or orchestrate it directly via the embedded worker:

// In handleExecution(id string, workflow shuffle.Workflow, request *http.Request, orgId string):

// 1. Resolve environment runtime capability
isBackendInjected := false
for _, env := range allEnvs {
    if env.Name == action.Environment {
        // Dynamic check: Check environment type or executor configuration
        if env.Type == "builtin" || env.Type == "backend" || env.Executor == "backend" {
            isBackendInjected = true
        }
        break
    }
}

// 2. Dispatch based on runtime location
if isBackendInjected && embeddedWorker != nil {
    log.Printf("[INFO][%s] Orchestrating execution directly via embedded backend worker", workflowExecution.ExecutionId)
    workflowExecution.Status = "EXECUTING"
    
    // Save execution state in database
    _ = shuffle.SetWorkflowExecution(ctx, workflowExecution, true)
    
    // Start workflow DAG execution in-process
    go embeddedWorker.StartWorkflowExecution(ctx, workflowExecution)
} else if execInfo.OnpremExecution {
    // Standard Orborus path: add to database queue for external Docker/K8s worker
    for _, environment := range execInfo.Environments {
        log.Printf("[INFO][%s] Queuing execution for external Orborus environment: %s", workflowExecution.ExecutionId, environment)
        executionRequest := shuffle.ExecutionRequest{
            ExecutionId:   workflowExecution.ExecutionId,
            WorkflowId:    workflowExecution.Workflow.ID,
            Authorization: workflowExecution.Authorization,
            Environments:  execInfo.Environments,
            Priority:      workflowExecution.Priority,
        }
        _ = shuffle.SetWorkflowQueue(ctx, executionRequest, environment)
    }
}

Step 4: Handle Action Results & DAG Progression in walkoff.go (runWorkflowExecutionTransaction)

When an app container/process finishes, it posts its result back to the backend at /api/v1/workflows/{id}/executions/{id}.

If the workflow is orchestrated by the embedded worker, notify HandleExecutionResult to resolve dependencies, check branch conditions, and trigger downstream actions:

// In runWorkflowExecutionTransaction(...) around line 764:

if setExecution || workflowExecution.Status == "FINISHED" || workflowExecution.Status == "ABORTED" || workflowExecution.Status == "FAILURE" {
    err = shuffle.SetWorkflowExecution(ctx, *workflowExecution, dbSave)
    if err != nil {
        resp.WriteHeader(401)
        resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err)))
        return
    }

    // Dynamic continuation: If this execution is handled by the backend worker,
    // evaluate DAG conditions and dispatch next actions:
    if embeddedWorker != nil && isBackendOrchestrated(workflowExecution.ExecutionId) {
        if workflowExecution.Status == "EXECUTING" {
            go embeddedWorker.HandleExecutionResult(ctx, *workflowExecution)
        }
    }
}

Action Execution Engines

Within both standalone and embedded modes, action execution is handled across different backend engines via pkg.Execute(ctx, engine, action, workflowExecution):

Engine Type Supported Aliases Execution Mechanism Ideal For
standalone standalone, local, process Local Python virtualenv in an isolated directory sandbox Low latency (~30ms), serverless / lightweight runners, systems without Docker
docker docker, container Ephemeral Docker container spun up per action (DeployContainer) Standard containerized apps, apps requiring root/system packages
swarm swarm, docker-swarm Long-running Docker Swarm services with HTTP dispatch (/api/v1/run) High-throughput distributed nodes
kubernetes kubernetes, k8s Kubernetes Pods/Jobs with HTTP dispatch (/api/v1/run) Production Kubernetes clusters

Standalone Engine Sandboxing & Security

When engine is set to "standalone", the worker executes the app's Python script without Docker while enforcing containment:

  1. Python Isolated Mode (-I):
    • Strips parent process PYTHON* environment variables.
    • Disables user site-packages (~/.local/...), strictly confining imports to the app's isolated virtualenv (<sandboxDir>/<app>/venv).
  2. Hermetic Environment (cmd.Env):
    • Host secrets, backend authorization tokens, and cloud metadata environment variables are not inherited by the Python process.
    • HOME, TMPDIR, and PATH are locked to the app's private sandbox folder.
  3. Directory Confinement (cmd.Dir):
    • The process working directory is locked to the app's assigned folder.
  4. OS-Level Containment Hooks:
    • macOS: Wrapped with /usr/bin/sandbox-exec (Seatbelt) denying filesystem writes outside the sandbox directory.
    • Linux: Wrapped with bwrap (Bubblewrap) if installed on the host.

Configuration Reference

Environment Variable Description Default
SHUFFLE_APP_EXECUTION_TYPE Default engine to use when resolving defaults (standalone, docker, swarm, kubernetes) docker
SHUFFLE_SANDBOX_DIR Base directory for standalone app scripts, virtualenvs, and temp files /tmp/shuffle-sandboxes
SHUFFLE_APP_SDK_TIMEOUT Default execution timeout for app actions (seconds) 30
KUBECONFIG Custom path to Kubernetes config file (when running outside cluster) ~/.kube/config
SHUFFLE_SWARM_NETWORK_NAME Docker network name for Swarm / container bridge communication shuffle_swarm_network

About

The Shuffle Worker for remote orchestration and automation, originally controlled in shuffle/shuffle/functions/onprem/worker

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages