From f57fb1350cfddb1ba2103bbbb24fd2179a64fb9b Mon Sep 17 00:00:00 2001 From: Evan Jarrett Date: Fri, 3 Apr 2026 16:29:10 -0500 Subject: [PATCH] matrix expansion and grpc for workflows --- .tangled/workflows/release-helm.yaml | 23 + .tangled/workflows/release.yaml | 3 +- .tangled/workflows/workflow-amd64.yaml | 25 - .tangled/workflows/workflow-arm64.yaml | 21 - Dockerfile | 4 +- Makefile | 33 +- api/v1alpha1/spindleset_types.go | 69 ++ api/v1alpha1/zz_generated.deepcopy.go | 47 ++ buf.gen.yaml | 10 + buf.yaml | 9 + cmd/controller/main.go | 93 ++- cmd/runner/main.go | 299 ++++++-- config/crd/bases/loom.j5t.io_spindlesets.yaml | 106 ++- config/gateway/httproute.yaml | 2 +- config/manager/grpc_service.yaml | 20 + config/manager/kustomization.yaml | 3 +- config/manager/loom-config.yaml | 2 +- config/manager/manager.yaml | 8 +- config/manager/service.yaml | 2 +- go.mod | 13 +- go.sum | 18 +- helm/loom/Chart.yaml | 4 +- helm/loom/crds/loom.j5t.io_spindlesets.yaml | 109 ++- helm/loom/templates/grpc_service.yaml | 17 + helm/loom/values.yaml | 6 +- internal/controller/spindleset_controller.go | 267 +++++-- internal/engine/kubernetes_engine.go | 693 +++++++++++------- internal/grpc/artifacts.go | 187 +++++ internal/grpc/hub.go | 142 ++++ internal/grpc/server.go | 136 ++++ internal/jobbuilder/job_template.go | 18 + internal/pb/loom/v1/loom.pb.go | 635 ++++++++++++++++ internal/pb/loom/v1/loom_grpc.pb.go | 129 ++++ proto/loom/v1/loom.proto | 77 ++ 34 files changed, 2722 insertions(+), 508 deletions(-) create mode 100644 .tangled/workflows/release-helm.yaml delete mode 100644 .tangled/workflows/workflow-amd64.yaml delete mode 100644 .tangled/workflows/workflow-arm64.yaml create mode 100644 buf.gen.yaml create mode 100644 buf.yaml create mode 100644 config/manager/grpc_service.yaml create mode 100644 helm/loom/templates/grpc_service.yaml create mode 100644 internal/grpc/artifacts.go create mode 100644 internal/grpc/hub.go create mode 100644 internal/grpc/server.go create mode 100644 internal/pb/loom/v1/loom.pb.go create mode 100644 internal/pb/loom/v1/loom_grpc.pb.go create mode 100644 proto/loom/v1/loom.proto diff --git a/.tangled/workflows/release-helm.yaml b/.tangled/workflows/release-helm.yaml new file mode 100644 index 0000000..2bd1b66 --- /dev/null +++ b/.tangled/workflows/release-helm.yaml @@ -0,0 +1,23 @@ +when: + - event: ["push"] + tag: ["v*"] + +engine: kubernetes +image: alpine/helm:latest +architecture: amd64 + +environment: + IMAGE_REGISTRY: buoy.cr + +steps: + - name: Login to registry + command: | + echo "${APP_PASSWORD}" | helm registry login \ + -u "${TANGLED_REPO_DID}" \ + --password-stdin \ + ${IMAGE_REGISTRY} + + - name: Package and push Helm chart + command: | + helm package helm/loom --version ${TANGLED_REF_NAME#v} --app-version ${TANGLED_REF_NAME#v} + helm push loom-${TANGLED_REF_NAME#v}.tgz oci://${IMAGE_REGISTRY}/${TANGLED_REPO_DID}/charts diff --git a/.tangled/workflows/release.yaml b/.tangled/workflows/release.yaml index 9a101ef..ed44e66 100644 --- a/.tangled/workflows/release.yaml +++ b/.tangled/workflows/release.yaml @@ -10,7 +10,7 @@ image: quay.io/buildah/stable:latest architecture: amd64 environment: - IMAGE_REGISTRY: atcr.io + IMAGE_REGISTRY: buoy.cr steps: - name: Login to registry @@ -30,3 +30,4 @@ steps: buildah push ${IMAGE_REGISTRY}/${TANGLED_REPO_DID}/${TANGLED_REPO_NAME}:latest buildah push ${IMAGE_REGISTRY}/${TANGLED_REPO_DID}/${TANGLED_REPO_NAME}:${TANGLED_REF_NAME} + diff --git a/.tangled/workflows/workflow-amd64.yaml b/.tangled/workflows/workflow-amd64.yaml deleted file mode 100644 index 21e128d..0000000 --- a/.tangled/workflows/workflow-amd64.yaml +++ /dev/null @@ -1,25 +0,0 @@ -when: - - event: ["push"] - tag: ["v*"] - branch : ["*"] - -engine: kubernetes -image: golang:1.25-trixie -architecture: amd64 - -environment: - IMAGE_REGISTRY: atcr.io - -steps: - - name: test environment vars - command: | - printenv - - - name: Login to registry - command: | - echo "${APP_PASSWORD}" | buildah login \ - -u "${TANGLED_REPO_DID}" \ - --password-stdin \ - ${IMAGE_REGISTRY} - - diff --git a/.tangled/workflows/workflow-arm64.yaml b/.tangled/workflows/workflow-arm64.yaml deleted file mode 100644 index f7fb7b9..0000000 --- a/.tangled/workflows/workflow-arm64.yaml +++ /dev/null @@ -1,21 +0,0 @@ -when: - - event: ["push"] - tag: ["v*"] - branch : ["*"] - -engine: kubernetes -image: golang:1.25-trixie -architecture: arm64 - -steps: - - name: build manager binary - command: | - make build - - - name: verify build artifacts - command: | - ls -lh bin/ - - - name: hello - command: | - echo "hello" \ No newline at end of file diff --git a/Dockerfile b/Dockerfile index 480a746..da557bc 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,6 +1,6 @@ # Build both binaries # Use BUILDPLATFORM so Go runs natively, cross-compile for target arch -FROM --platform=$BUILDPLATFORM golang:1.25 AS builder +FROM --platform=$BUILDPLATFORM golang:1.25-trixie AS builder ARG TARGETOS ARG TARGETARCH @@ -40,7 +40,7 @@ RUN CC=$(if [ "$TARGETARCH" = "arm64" ] && [ "$BUILDARCH" = "amd64" ]; then echo go build -a -ldflags='-s -w' -o manager ./cmd/controller # Unified image with both binaries -FROM gcr.io/distroless/base-debian12:nonroot +FROM gcr.io/distroless/base-debian13:nonroot COPY --from=builder /workspace/loom/manager /manager COPY --from=builder /workspace/loom/loom-runner /loom-runner diff --git a/Makefile b/Makefile index fdeac28..3b881d2 100644 --- a/Makefile +++ b/Makefile @@ -3,7 +3,7 @@ # To re-generate a bundle for another specific version without changing the standard setup, you can: # - use the VERSION as arg of the bundle target (e.g make bundle VERSION=0.0.2) # - use environment variables to overwrite this value (e.g export VERSION=0.0.2) -VERSION ?= 0.0.1 +VERSION ?= 0.1.5 # CHANNELS define the bundle channels used in the bundle. # Add a new line here if you would like to change its default config. (E.g CHANNELS = "candidate,fast,stable") @@ -50,7 +50,7 @@ endif # This is useful for CI or a project to utilize a specific version of the operator-sdk toolkit. OPERATOR_SDK_VERSION ?= v1.41.1 # Image URL to use all building/pushing image targets -IMG ?= atcr.io/evan.jarrett.net/loom:latest +IMG ?= buoy.cr/evan.jarrett.net/loom:latest # Get the currently used golang install path (in GOPATH/bin, unless GOBIN is set) ifeq (,$(shell go env GOBIN)) @@ -100,6 +100,10 @@ manifests: controller-gen ## Generate WebhookConfiguration, ClusterRole and Cust generate: controller-gen ## Generate code containing DeepCopy, DeepCopyInto, and DeepCopyObject method implementations. $(CONTROLLER_GEN) object:headerFile="hack/boilerplate.go.txt" paths="./..." +.PHONY: proto +proto: ## Generate protobuf and gRPC code. + buf generate + .PHONY: fmt fmt: ## Run go fmt against code. go fmt ./... @@ -211,10 +215,10 @@ setup-buildx: ## Set up buildx builder with credential access for multi-arch bui .PHONY: test-registry-auth test-registry-auth: ## Test registry authentication before building - @echo "Testing registry authentication for atcr.io..." + @echo "Testing registry authentication for buoy.cr..." @if [ -f /usr/local/sbin/docker-credential-atcr ]; then \ echo "Testing credential helper..."; \ - echo "atcr.io" | docker-credential-atcr get && echo "✓ Credential helper working!" || echo "✗ Credential helper failed"; \ + echo "buoy.cr" | docker-credential-atcr get && echo "✓ Credential helper working!" || echo "✗ Credential helper failed"; \ else \ echo "⚠ Credential helper not found at /usr/local/sbin/docker-credential-atcr"; \ fi @@ -222,30 +226,13 @@ test-registry-auth: ## Test registry authentication before building @echo "Testing Docker config..." @if [ -f $(HOME)/.docker/config.json ]; then \ echo "✓ Docker config exists at $(HOME)/.docker/config.json"; \ - cat $(HOME)/.docker/config.json | grep -q "atcr.io" && echo "✓ atcr.io found in config" || echo "⚠ atcr.io not found in config"; \ + cat $(HOME)/.docker/config.json | grep -q "buoy.cr" && echo "✓ buoy.cr found in config" || echo "⚠ buoy.cr not found in config"; \ else \ echo "✗ Docker config not found"; \ fi @echo "" @echo "Testing registry access with docker pull (this will fail if auth is broken)..." - @$(CONTAINER_TOOL) pull atcr.io/evan.jarrett.net/loom-runner:latest 2>/dev/null && echo "✓ Can pull from registry!" || echo "⚠ Cannot pull from registry (may not exist yet)" - -# PLATFORMS defines the target platforms for the manager image be built to provide support to multiple -# architectures. (i.e. make docker-buildx IMG=myregistry/mypoperator:0.0.1). To use this option you need to: -# - be able to use docker buildx. More info: https://docs.docker.com/build/buildx/ -# - have enabled BuildKit. More info: https://docs.docker.com/develop/develop-images/build_enhancements/ -# - be able to push the image to your registry (i.e. if you do not set a valid value via IMG=> then the export will fail) -# To adequately provide solutions that are compatible with multiple platforms, you should consider using this option. -PLATFORMS ?= linux/arm64,linux/amd64,linux/s390x,linux/ppc64le -.PHONY: docker-buildx -docker-buildx: ## Build and push docker image for the manager for cross-platform support - # copy existing Dockerfile and insert --platform=${BUILDPLATFORM} into Dockerfile.cross, and preserve the original Dockerfile - sed -e '1 s/\(^FROM\)/FROM --platform=\$$\{BUILDPLATFORM\}/; t' -e ' 1,// s//FROM --platform=\$$\{BUILDPLATFORM\}/' Dockerfile > Dockerfile.cross - - $(CONTAINER_TOOL) buildx create --name loom-builder - $(CONTAINER_TOOL) buildx use loom-builder - - cd .. && $(CONTAINER_TOOL) buildx build --push --platform=$(PLATFORMS) --tag ${IMG} -f loom/Dockerfile.cross . - - $(CONTAINER_TOOL) buildx rm loom-builder - rm Dockerfile.cross + @$(CONTAINER_TOOL) pull buoy.cr/evan.jarrett.net/loom-runner:latest 2>/dev/null && echo "✓ Can pull from registry!" || echo "⚠ Cannot pull from registry (may not exist yet)" .PHONY: build-installer build-installer: manifests generate kustomize ## Generate a consolidated YAML with CRDs and deployment. diff --git a/api/v1alpha1/spindleset_types.go b/api/v1alpha1/spindleset_types.go index 95a33c5..8833653 100644 --- a/api/v1alpha1/spindleset_types.go +++ b/api/v1alpha1/spindleset_types.go @@ -61,8 +61,14 @@ type PipelineRunSpec struct { Secrets []SecretData `json:"secrets,omitempty"` // Workflows is the list of workflows to execute in this pipeline. + // For multi-arch workflows, this contains one entry per matrix leg plus an optional final entry. // +kubebuilder:validation:MinItems=1 Workflows []WorkflowSpec `json:"workflows"` + + // MultiArch indicates this pipeline run contains multi-arch workflows. + // When true, the controller creates per-architecture Jobs and gates the final Job. + // +optional + MultiArch bool `json:"multiArch,omitempty"` } // SecretData represents a single secret key-value pair for injection into Jobs. @@ -81,6 +87,8 @@ type SecretData struct { // WorkflowSpec defines a workflow to execute as part of a pipeline. // This is the canonical workflow definition that matches the .tangled/workflows/*.yaml format. +// For multi-arch workflows, the engine expands the matrix and creates one WorkflowSpec per leg. +// Each leg has a single Image and Architecture; the matrix metadata lives in PipelineRunSpec. type WorkflowSpec struct { // Name is the workflow filename (e.g., "workflow-amd64.yaml"). // +kubebuilder:validation:Required @@ -110,6 +118,36 @@ type WorkflowSpec struct { // Dependencies specifies external dependencies for the workflow. // +optional Dependencies *WorkflowDependencies `json:"dependencies,omitempty"` + + // Final defines steps that run once after all matrix legs complete. + // Only valid on multi-arch workflows. The engine sets this on the dedicated final WorkflowSpec. + // +optional + Final *FinalSpec `json:"final,omitempty"` + + // IsMatrixLeg indicates this WorkflowSpec was generated from a matrix expansion. + // +optional + IsMatrixLeg bool `json:"isMatrixLeg,omitempty"` + + // IsFinal indicates this WorkflowSpec represents the final step of a multi-arch workflow. + // +optional + IsFinal bool `json:"isFinal,omitempty"` +} + +// FinalSpec defines steps that run once after all matrix legs complete successfully. +type FinalSpec struct { + // Architecture is the target architecture for the final steps. + // +kubebuilder:validation:Required + // +kubebuilder:validation:Enum=amd64;arm64 + Architecture string `json:"architecture"` + + // Image is the container image for the final steps. + // If empty, uses the first image from the matrix. + // +optional + Image string `json:"image,omitempty"` + + // Steps is the ordered list of steps to execute after all matrix legs complete. + // +kubebuilder:validation:MinItems=1 + Steps []WorkflowStep `json:"steps"` } // WorkflowStep defines a single step in a workflow. @@ -224,6 +262,7 @@ type WorkflowStatus struct { Name string `json:"name"` // JobName is the name of the Kubernetes Job created for this workflow. + // For multi-arch workflows, this is empty; use MatrixLegStatuses instead. // +optional JobName string `json:"jobName,omitempty"` @@ -238,6 +277,36 @@ type WorkflowStatus struct { // CompletionTime is when the workflow finished. // +optional CompletionTime *metav1.Time `json:"completionTime,omitempty"` + + // MatrixLegStatuses tracks per-architecture Job statuses for multi-arch workflows. + // +optional + MatrixLegStatuses []MatrixLegStatus `json:"matrixLegStatuses,omitempty"` + + // FinalJobName is the name of the final Job (for multi-arch workflows). + // +optional + FinalJobName string `json:"finalJobName,omitempty"` + + // FinalPhase is the phase of the final Job. + // +optional + FinalPhase string `json:"finalPhase,omitempty"` +} + +// MatrixLegStatus tracks the status of a single matrix leg Job. +type MatrixLegStatus struct { + // Architecture is the target architecture for this leg. + Architecture string `json:"architecture"` + + // Image is the container image used for this leg. + // +optional + Image string `json:"image,omitempty"` + + // JobName is the name of the Kubernetes Job for this leg. + // +optional + JobName string `json:"jobName,omitempty"` + + // Phase is the current phase (Pending, Running, Succeeded, Failed). + // +optional + Phase string `json:"phase,omitempty"` } // +kubebuilder:object:root=true diff --git a/api/v1alpha1/zz_generated.deepcopy.go b/api/v1alpha1/zz_generated.deepcopy.go index aaf5f9c..4219111 100644 --- a/api/v1alpha1/zz_generated.deepcopy.go +++ b/api/v1alpha1/zz_generated.deepcopy.go @@ -26,6 +26,43 @@ import ( runtime "k8s.io/apimachinery/pkg/runtime" ) +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *FinalSpec) DeepCopyInto(out *FinalSpec) { + *out = *in + if in.Steps != nil { + in, out := &in.Steps, &out.Steps + *out = make([]WorkflowStep, len(*in)) + for i := range *in { + (*in)[i].DeepCopyInto(&(*out)[i]) + } + } +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new FinalSpec. +func (in *FinalSpec) DeepCopy() *FinalSpec { + if in == nil { + return nil + } + out := new(FinalSpec) + in.DeepCopyInto(out) + return out +} + +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *MatrixLegStatus) DeepCopyInto(out *MatrixLegStatus) { + *out = *in +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new MatrixLegStatus. +func (in *MatrixLegStatus) DeepCopy() *MatrixLegStatus { + if in == nil { + return nil + } + out := new(MatrixLegStatus) + in.DeepCopyInto(out) + return out +} + // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *PipelineRunSpec) DeepCopyInto(out *PipelineRunSpec) { *out = *in @@ -288,6 +325,11 @@ func (in *WorkflowSpec) DeepCopyInto(out *WorkflowSpec) { *out = new(WorkflowDependencies) (*in).DeepCopyInto(*out) } + if in.Final != nil { + in, out := &in.Final, &out.Final + *out = new(FinalSpec) + (*in).DeepCopyInto(*out) + } } // DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new WorkflowSpec. @@ -311,6 +353,11 @@ func (in *WorkflowStatus) DeepCopyInto(out *WorkflowStatus) { in, out := &in.CompletionTime, &out.CompletionTime *out = (*in).DeepCopy() } + if in.MatrixLegStatuses != nil { + in, out := &in.MatrixLegStatuses, &out.MatrixLegStatuses + *out = make([]MatrixLegStatus, len(*in)) + copy(*out, *in) + } } // DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new WorkflowStatus. diff --git a/buf.gen.yaml b/buf.gen.yaml new file mode 100644 index 0000000..6e6d01a --- /dev/null +++ b/buf.gen.yaml @@ -0,0 +1,10 @@ +version: v2 +plugins: + - local: protoc-gen-go + out: internal/pb + opt: + - paths=source_relative + - local: protoc-gen-go-grpc + out: internal/pb + opt: + - paths=source_relative diff --git a/buf.yaml b/buf.yaml new file mode 100644 index 0000000..c7e30e3 --- /dev/null +++ b/buf.yaml @@ -0,0 +1,9 @@ +version: v2 +modules: + - path: proto +lint: + use: + - STANDARD +breaking: + use: + - FILE diff --git a/cmd/controller/main.go b/cmd/controller/main.go index 78b73aa..dc2baa5 100644 --- a/cmd/controller/main.go +++ b/cmd/controller/main.go @@ -22,6 +22,7 @@ import ( _ "embed" "flag" "fmt" + "net" "os" "path/filepath" @@ -37,6 +38,7 @@ import ( clientgoscheme "k8s.io/client-go/kubernetes/scheme" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/certwatcher" + "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/healthz" "sigs.k8s.io/controller-runtime/pkg/log/zap" "sigs.k8s.io/controller-runtime/pkg/metrics/filters" @@ -50,6 +52,7 @@ import ( loomv1alpha1 "tangled.org/evan.jarrett.net/loom/api/v1alpha1" "tangled.org/evan.jarrett.net/loom/internal/controller" "tangled.org/evan.jarrett.net/loom/internal/engine" + loomgrpc "tangled.org/evan.jarrett.net/loom/internal/grpc" // +kubebuilder:scaffold:imports ) @@ -177,7 +180,7 @@ func convertToResourceProfiles(profiles []ResourceProfileConfig) ([]loomv1alpha1 // initializeSpindle creates a spindle server with KubernetesEngine func initializeSpindle( - ctx context.Context, cfg *config.Config, mgr ctrl.Manager, loomCfg *LoomConfig, + ctx context.Context, cfg *config.Config, mgr ctrl.Manager, loomCfg *LoomConfig, hub *loomgrpc.Hub, artifacts *loomgrpc.ArtifactStore, ) (*spindle.Spindle, error) { // Initialize Kubernetes engine // Get namespace from environment (injected via Downward API) @@ -203,8 +206,8 @@ func initializeSpindle( return nil, fmt.Errorf("failed to create spindle: %w", err) } - // Now create kubernetes engine with access to vault - kubeEngine := engine.NewKubernetesEngine(mgr.GetClient(), mgr.GetConfig(), namespace, template, s.Vault()) + // Now create kubernetes engine with access to vault and gRPC hub + kubeEngine := engine.NewKubernetesEngine(mgr.GetClient(), mgr.GetConfig(), namespace, template, s.Vault(), hub, artifacts) // Register the engine with spindle by adding to the engines map s.Engines()["kubernetes"] = kubeEngine @@ -385,8 +388,22 @@ func main() { os.Exit(1) } + // Create gRPC hub for runner communication + hub := loomgrpc.NewHub() + + // Create artifact store for pipeline artifacts (scratch directory) + artifactDir := "/scratch/artifacts" + if dir := os.Getenv("LOOM_ARTIFACT_DIR"); dir != "" { + artifactDir = dir + } + artifactStore, err := loomgrpc.NewArtifactStore(artifactDir) + if err != nil { + setupLog.Error(err, "failed to create artifact store") + os.Exit(1) + } + // Initialize spindle server with KubernetesEngine - s, err := initializeSpindle(ctx, spindleCfg, mgr, loomCfg) + s, err := initializeSpindle(ctx, spindleCfg, mgr, loomCfg, hub, artifactStore) if err != nil { setupLog.Error(err, "failed to initialize spindle") os.Exit(1) @@ -396,6 +413,26 @@ func main() { setupLog.Info("spindle server initialized successfully") + // Start gRPC server for runner communication + grpcAddr := ":9090" + if addr := os.Getenv("LOOM_GRPC_ADDR"); addr != "" { + grpcAddr = addr + } + + grpcServer := loomgrpc.NewServer(hub, artifactStore) + go func() { + lis, err := net.Listen("tcp", grpcAddr) + if err != nil { + setupLog.Error(err, "failed to listen for gRPC", "address", grpcAddr) + os.Exit(1) + } + setupLog.Info("starting gRPC server", "address", grpcAddr) + if err := grpcServer.Serve(lis); err != nil { + setupLog.Error(err, "gRPC server error") + } + }() + defer grpcServer.GracefulStop() + // Start spindle HTTP server in background go func() { setupLog.Info("starting spindle HTTP server", "address", spindleCfg.Server.ListenAddr) @@ -407,16 +444,52 @@ func main() { // Get loom image from environment (used for runner init container) loomImage := os.Getenv("LOOM_IMAGE") if loomImage == "" { - loomImage = "atcr.io/evan.jarrett.net/loom:latest" // default fallback + loomImage = "buoy.cr/evan.jarrett.net/loom:latest" // default fallback + } + + // Discover the gRPC service address that runner pods will use to reach the operator. + podNamespace := os.Getenv("POD_NAMESPACE") + if podNamespace == "" { + podNamespace = "default" + } + operatorAddr := os.Getenv("LOOM_OPERATOR_ADDR") + if operatorAddr == "" { + // Find the gRPC service by label in our namespace + var services corev1.ServiceList + if err := mgr.GetAPIReader().List(context.Background(), &services, + client.InNamespace(podNamespace), + client.MatchingLabels{ + "app.kubernetes.io/name": "loom", + "app.kubernetes.io/component": "grpc", + }, + ); err != nil { + setupLog.Error(err, "failed to discover gRPC service") + os.Exit(1) + } + if len(services.Items) == 0 { + setupLog.Error(nil, "no gRPC service found with label app.kubernetes.io/component=grpc") + os.Exit(1) + } + svc := services.Items[0] + grpcPort := int32(9090) + for _, p := range svc.Spec.Ports { + if p.Name == "grpc" { + grpcPort = p.Port + break + } + } + operatorAddr = fmt.Sprintf("%s.%s.svc.cluster.local:%d", svc.Name, podNamespace, grpcPort) + setupLog.Info("discovered gRPC service", "address", operatorAddr) } // Setup controller with spindle components if err := (&controller.SpindleSetReconciler{ - Client: mgr.GetClient(), - Scheme: mgr.GetScheme(), - Config: mgr.GetConfig(), - Spindle: s, - LoomImage: loomImage, + Client: mgr.GetClient(), + Scheme: mgr.GetScheme(), + Config: mgr.GetConfig(), + Spindle: s, + LoomImage: loomImage, + OperatorAddr: operatorAddr, }).SetupWithManager(mgr); err != nil { setupLog.Error(err, "unable to create controller", "controller", "SpindleSet") os.Exit(1) diff --git a/cmd/runner/main.go b/cmd/runner/main.go index 81620ab..fcdb250 100644 --- a/cmd/runner/main.go +++ b/cmd/runner/main.go @@ -8,9 +8,14 @@ import ( "io" "os" "os/exec" + "path/filepath" + + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" "tangled.org/core/spindle/models" loomv1alpha1 "tangled.org/evan.jarrett.net/loom/api/v1alpha1" + pb "tangled.org/evan.jarrett.net/loom/internal/pb/loom/v1" ) // simpleStep implements the models.Step interface @@ -19,22 +24,43 @@ type simpleStep struct { command string } -// extendedLogLine extends models.LogLine with exit code for error reporting -type extendedLogLine struct { - models.LogLine - ExitCode int `json:"exit_code,omitempty"` +func (s *simpleStep) Name() string { return s.name } +func (s *simpleStep) Command() string { return s.command } +func (s *simpleStep) Kind() models.StepKind { + return models.StepKindUser } -func (s *simpleStep) Name() string { - return s.name +// grpcEmitter sends events to the operator over gRPC and also writes to stdout. +type grpcEmitter struct { + stream grpc.BidiStreamingClient[pb.ConnectRequest, pb.ConnectResponse] } -func (s *simpleStep) Command() string { - return s.command +func (e *grpcEmitter) sendStepControl(stepID int, status string, exitCode int) { + if e.stream != nil { + _ = e.stream.Send(&pb.ConnectRequest{ + Event: &pb.ConnectRequest_StepControl{ + StepControl: &pb.StepControl{ + StepId: int32(stepID), + Status: status, + ExitCode: int32(exitCode), + }, + }, + }) + } } -func (s *simpleStep) Kind() models.StepKind { - return models.StepKindUser +func (e *grpcEmitter) sendLogLine(stepID int, streamName, content string) { + if e.stream != nil { + _ = e.stream.Send(&pb.ConnectRequest{ + Event: &pb.ConnectRequest_LogLine{ + LogLine: &pb.LogLine{ + StepId: int32(stepID), + Stream: streamName, + Content: content, + }, + }, + }) + } } func main() { @@ -76,7 +102,6 @@ func installSelf(dst string) error { return fmt.Errorf("failed to copy: %w", err) } - // Make executable if err := os.Chmod(dst, 0755); err != nil { return fmt.Errorf("failed to chmod: %w", err) } @@ -96,6 +121,14 @@ func run() error { return fmt.Errorf("failed to parse workflow spec: %w", err) } + // Connect to operator via gRPC + emitter, cleanup, err := connectToOperator(workflow) + if err != nil { + // gRPC connection failure is fatal — the operator won't see our events + return fmt.Errorf("failed to connect to operator: %w", err) + } + defer cleanup() + // Set up environment variables if workflow.Environment != nil { for k, v := range workflow.Environment { @@ -105,26 +138,90 @@ func run() error { } } + // Create artifacts directory so user steps can write to it (non-root container + // can't create /artifacts at root itself). Safe to create even when unused. + artifactsDir := os.Getenv("LOOM_ARTIFACTS") + if artifactsDir != "" { + if err := os.MkdirAll(artifactsDir, 0755); err != nil { + return fmt.Errorf("failed to create artifacts directory: %w", err) + } + } + + // For final jobs, download artifacts from matrix legs before executing steps + if os.Getenv("LOOM_FINAL") == "true" && artifactsDir != "" { + fmt.Fprintf(os.Stderr, "downloading artifacts from matrix legs...\n") + if err := downloadArtifacts(emitter, artifactsDir); err != nil { + return fmt.Errorf("failed to download artifacts: %w", err) + } + fmt.Fprintf(os.Stderr, "artifacts downloaded to %s\n", artifactsDir) + } + // Execute each step ctx := context.Background() for i, step := range workflow.Steps { - if err := executeStep(ctx, i, step); err != nil { + if err := executeStep(ctx, i, step, emitter); err != nil { return fmt.Errorf("step %d (%s) failed: %w", i, step.Name, err) } } + // Upload artifacts if this is a matrix leg + if os.Getenv("LOOM_MATRIX_LEG") == "true" && artifactsDir != "" { + if err := uploadArtifacts(emitter, artifactsDir); err != nil { + fmt.Fprintf(os.Stderr, "WARNING: failed to upload artifacts: %v\n", err) + } + } + return nil } -func executeStep(ctx context.Context, stepID int, step loomv1alpha1.WorkflowStep) error { - // Create a simple step for logging - simpleStep := &simpleStep{ - name: step.Name, - command: step.Command, +// connectToOperator establishes the gRPC connection and sends the identity message. +func connectToOperator(workflow loomv1alpha1.WorkflowSpec) (*grpcEmitter, func(), error) { + addr := os.Getenv("LOOM_OPERATOR_ADDR") + if addr == "" { + return nil, nil, fmt.Errorf("LOOM_OPERATOR_ADDR environment variable not set") } - // Emit step start event - emitControlEvent(stepID, simpleStep, models.StepStatusStart) + conn, err := grpc.NewClient(addr, grpc.WithTransportCredentials(insecure.NewCredentials())) + if err != nil { + return nil, nil, fmt.Errorf("failed to create gRPC client: %w", err) + } + + client := pb.NewLoomRunnerServiceClient(conn) + stream, err := client.Connect(context.Background()) + if err != nil { + conn.Close() + return nil, nil, fmt.Errorf("failed to open gRPC stream: %w", err) + } + + // Send identity message + pipelineID := os.Getenv("TANGLED_PIPELINE_ID") + if pipelineID == "" { + pipelineID = workflow.Environment["TANGLED_PIPELINE_ID"] + } + + if err := stream.Send(&pb.ConnectRequest{ + PipelineId: pipelineID, + WorkflowName: workflow.Name, + Architecture: workflow.Architecture, + }); err != nil { + conn.Close() + return nil, nil, fmt.Errorf("failed to send identity: %w", err) + } + + emitter := &grpcEmitter{stream: stream} + cleanup := func() { + _ = stream.CloseSend() + conn.Close() + } + + fmt.Fprintf(os.Stderr, "connected to operator at %s\n", addr) + return emitter, cleanup, nil +} + +func executeStep(ctx context.Context, stepID int, step loomv1alpha1.WorkflowStep, emitter *grpcEmitter) error { + // Emit step start — gRPC + stdout + emitter.sendStepControl(stepID, "start", 0) + emitStdoutControl(stepID, &simpleStep{name: step.Name, command: step.Command}, models.StepStatusStart) // Set step-specific environment variables if step.Environment != nil { @@ -136,13 +233,11 @@ func executeStep(ctx context.Context, stepID int, step loomv1alpha1.WorkflowStep } // Create command that auto-sources LOOM_ENV if it exists - // Users can write "VAR=value" to this file to share env vars between steps wrappedCommand := `if [ -f "$LOOM_ENV" ]; then set -a; source "$LOOM_ENV"; set +a; fi; ` + step.Command cmd := exec.CommandContext(ctx, "bash", "-c", wrappedCommand) cmd.Dir = "/tangled/workspace" cmd.Env = append(os.Environ(), "LOOM_ENV=/tangled/workspace/.loom-env") - // Capture stdout and stderr stdout, err := cmd.StdoutPipe() if err != nil { return fmt.Errorf("failed to create stdout pipe: %w", err) @@ -153,21 +248,18 @@ func executeStep(ctx context.Context, stepID int, step loomv1alpha1.WorkflowStep return fmt.Errorf("failed to create stderr pipe: %w", err) } - // Start the command if err := cmd.Start(); err != nil { - emitControlEvent(stepID, simpleStep, models.StepStatusEnd) + emitter.sendStepControl(stepID, "end", 1) return fmt.Errorf("failed to start command: %w", err) } // Stream stdout and stderr concurrently done := make(chan error, 2) - go streamOutput(stdout, stepID, "stdout", done) - go streamOutput(stderr, stepID, "stderr", done) + go streamOutput(stdout, stepID, "stdout", emitter, done) + go streamOutput(stderr, stepID, "stderr", emitter, done) - // Wait for both streams to complete for i := 0; i < 2; i++ { if err := <-done; err != nil { - // Log error but don't fail - we still want to wait for the command fmt.Fprintf(os.Stderr, "WARNING: error streaming output: %v\n", err) } } @@ -183,8 +275,9 @@ func executeStep(ctx context.Context, stepID int, step loomv1alpha1.WorkflowStep } } - // Emit step end event with exit code for error reporting - emitControlEventWithCode(stepID, simpleStep, models.StepStatusEnd, exitCode) + // Emit step end — gRPC + stdout + emitter.sendStepControl(stepID, "end", exitCode) + emitStdoutControlWithCode(stepID, &simpleStep{name: step.Name, command: step.Command}, models.StepStatusEnd, exitCode) if exitCode != 0 { return fmt.Errorf("command exited with code %d", exitCode) @@ -193,50 +286,166 @@ func executeStep(ctx context.Context, stepID int, step loomv1alpha1.WorkflowStep return nil } -func streamOutput(reader io.Reader, stepID int, stream string, done chan<- error) { +func streamOutput(reader io.Reader, stepID int, streamName string, emitter *grpcEmitter, done chan<- error) { scanner := bufio.NewScanner(reader) - // Increase buffer size for long lines buf := make([]byte, 0, 64*1024) scanner.Buffer(buf, 1024*1024) for scanner.Scan() { line := scanner.Text() - emitDataEvent(stepID, stream, line) + // Send over gRPC (primary channel) + emitter.sendLogLine(stepID, streamName, line) + // Also emit to stdout for kubectl logs + emitStdoutData(stepID, streamName, line) } done <- scanner.Err() } -func emitControlEvent(stepID int, step models.Step, status models.StepStatus) { +// Stdout emitters — for kubectl logs debugging. Not consumed by the operator. + +func emitStdoutControl(stepID int, step models.Step, status models.StepStatus) { logLine := models.NewControlLogLine(stepID, step, status) - emitJSON(logLine) + emitStdoutJSON(logLine) } -// emitControlEventWithCode emits a control event with an exit code for error reporting -func emitControlEventWithCode(stepID int, step models.Step, status models.StepStatus, exitCode int) { +func emitStdoutControlWithCode(stepID int, step models.Step, status models.StepStatus, exitCode int) { logLine := models.NewControlLogLine(stepID, step, status) - extended := extendedLogLine{ - LogLine: logLine, - ExitCode: exitCode, + type extended struct { + models.LogLine + ExitCode int `json:"exit_code,omitempty"` } - data, err := json.Marshal(extended) + data, err := json.Marshal(extended{LogLine: logLine, ExitCode: exitCode}) if err != nil { - fmt.Fprintf(os.Stderr, "ERROR: failed to marshal JSON: %v\n", err) return } fmt.Println(string(data)) } -func emitDataEvent(stepID int, stream, content string) { +func emitStdoutData(stepID int, stream, content string) { logLine := models.NewDataLogLine(stepID, content, stream) - emitJSON(logLine) + emitStdoutJSON(logLine) } -func emitJSON(logLine models.LogLine) { +func emitStdoutJSON(logLine models.LogLine) { data, err := json.Marshal(logLine) if err != nil { - fmt.Fprintf(os.Stderr, "ERROR: failed to marshal JSON: %v\n", err) return } fmt.Println(string(data)) } + +// uploadArtifacts walks the artifacts directory and streams all files to the operator. +func uploadArtifacts(emitter *grpcEmitter, dir string) error { + info, err := os.Stat(dir) + if os.IsNotExist(err) { + return nil // No artifacts to upload + } + if err != nil { + return err + } + if !info.IsDir() { + return nil + } + + const chunkSize = 32 * 1024 // 32KB + + return filepath.Walk(dir, func(path string, info os.FileInfo, err error) error { + if err != nil || info.IsDir() { + return err + } + + relPath, err := filepath.Rel(dir, path) + if err != nil { + return err + } + + f, err := os.Open(path) + if err != nil { + return fmt.Errorf("failed to open artifact %s: %w", relPath, err) + } + defer f.Close() + + buf := make([]byte, chunkSize) + for { + n, readErr := f.Read(buf) + isEOF := readErr == io.EOF + + if emitter.stream != nil { + _ = emitter.stream.Send(&pb.ConnectRequest{ + Event: &pb.ConnectRequest_ArtifactChunk{ + ArtifactChunk: &pb.ArtifactChunk{ + Path: relPath, + Data: buf[:n], + Eof: isEOF, + }, + }, + }) + } + + if isEOF { + break + } + if readErr != nil { + return fmt.Errorf("failed to read artifact %s: %w", relPath, readErr) + } + } + + fmt.Fprintf(os.Stderr, "uploaded artifact: %s\n", relPath) + return nil + }) +} + +// downloadArtifacts receives artifact files from the operator into the local artifacts directory. +// Used by final jobs to receive artifacts from matrix legs. +func downloadArtifacts(emitter *grpcEmitter, dir string) error { + if emitter.stream == nil { + return nil + } + + for { + resp, err := emitter.stream.Recv() + if err == io.EOF { + return nil + } + if err != nil { + return fmt.Errorf("failed to receive artifact: %w", err) + } + + ad, ok := resp.Event.(*pb.ConnectResponse_ArtifactData) + if !ok { + continue + } + + data := ad.ArtifactData + + // Sentinel: empty path with Eof=true signals "all artifacts sent". + if data.Path == "" && data.Eof { + return nil + } + + targetDir := filepath.Join(dir, data.SourceArchitecture) + targetPath := filepath.Join(targetDir, data.Path) + + if err := os.MkdirAll(filepath.Dir(targetPath), 0755); err != nil { + return fmt.Errorf("failed to create artifact directory: %w", err) + } + + f, err := os.OpenFile(targetPath, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644) + if err != nil { + return fmt.Errorf("failed to open artifact file: %w", err) + } + + if len(data.Data) > 0 { + if _, err := f.Write(data.Data); err != nil { + f.Close() + return fmt.Errorf("failed to write artifact: %w", err) + } + } + f.Close() + + if data.Eof { + fmt.Fprintf(os.Stderr, "downloaded artifact: %s/%s\n", data.SourceArchitecture, data.Path) + } + } +} diff --git a/config/crd/bases/loom.j5t.io_spindlesets.yaml b/config/crd/bases/loom.j5t.io_spindlesets.yaml index 8ba1c14..1d22154 100644 --- a/config/crd/bases/loom.j5t.io_spindlesets.yaml +++ b/config/crd/bases/loom.j5t.io_spindlesets.yaml @@ -77,6 +77,11 @@ spec: items: type: string type: array + multiArch: + description: |- + MultiArch indicates this pipeline run contains multi-arch workflows. + When true, the controller creates per-architecture Jobs and gates the final Job. + type: boolean pipelineID: description: PipelineID is the unique identifier for this pipeline run from the knot. @@ -110,12 +115,15 @@ spec: container entirely. type: boolean workflows: - description: Workflows is the list of workflows to execute in - this pipeline. + description: |- + Workflows is the list of workflows to execute in this pipeline. + For multi-arch workflows, this contains one entry per matrix leg plus an optional final entry. items: description: |- WorkflowSpec defines a workflow to execute as part of a pipeline. This is the canonical workflow definition that matches the .tangled/workflows/*.yaml format. + For multi-arch workflows, the engine expands the matrix and creates one WorkflowSpec per leg. + Each leg has a single Image and Architecture; the matrix metadata lives in PipelineRunSpec. properties: architecture: description: Architecture is the target architecture for @@ -141,10 +149,65 @@ spec: description: Environment contains workflow-level environment variables. type: object + final: + description: |- + Final defines steps that run once after all matrix legs complete. + Only valid on multi-arch workflows. The engine sets this on the dedicated final WorkflowSpec. + properties: + architecture: + description: Architecture is the target architecture + for the final steps. + enum: + - amd64 + - arm64 + type: string + image: + description: |- + Image is the container image for the final steps. + If empty, uses the first image from the matrix. + type: string + steps: + description: Steps is the ordered list of steps to execute + after all matrix legs complete. + items: + description: WorkflowStep defines a single step in + a workflow. + properties: + command: + description: Command is the shell command to execute. + type: string + environment: + additionalProperties: + type: string + description: Environment contains step-specific + environment variables. + type: object + name: + description: Name is the human-readable name of + the step. + type: string + required: + - command + - name + type: object + minItems: 1 + type: array + required: + - architecture + - steps + type: object image: description: Image is the container image to use for executing the workflow steps. type: string + isFinal: + description: IsFinal indicates this WorkflowSpec represents + the final step of a multi-arch workflow. + type: boolean + isMatrixLeg: + description: IsMatrixLeg indicates this WorkflowSpec was + generated from a matrix expansion. + type: boolean name: description: Name is the workflow filename (e.g., "workflow-amd64.yaml"). type: string @@ -1360,10 +1423,45 @@ spec: description: CompletionTime is when the workflow finished. format: date-time type: string + finalJobName: + description: FinalJobName is the name of the final Job (for + multi-arch workflows). + type: string + finalPhase: + description: FinalPhase is the phase of the final Job. + type: string jobName: - description: JobName is the name of the Kubernetes Job created - for this workflow. + description: |- + JobName is the name of the Kubernetes Job created for this workflow. + For multi-arch workflows, this is empty; use MatrixLegStatuses instead. type: string + matrixLegStatuses: + description: MatrixLegStatuses tracks per-architecture Job statuses + for multi-arch workflows. + items: + description: MatrixLegStatus tracks the status of a single + matrix leg Job. + properties: + architecture: + description: Architecture is the target architecture for + this leg. + type: string + image: + description: Image is the container image used for this + leg. + type: string + jobName: + description: JobName is the name of the Kubernetes Job + for this leg. + type: string + phase: + description: Phase is the current phase (Pending, Running, + Succeeded, Failed). + type: string + required: + - architecture + type: object + type: array name: description: Name is the workflow name. type: string diff --git a/config/gateway/httproute.yaml b/config/gateway/httproute.yaml index f7f17dd..a97be29 100644 --- a/config/gateway/httproute.yaml +++ b/config/gateway/httproute.yaml @@ -11,5 +11,5 @@ spec: - loom.jarrett.net rules: - backendRefs: - - name: loom-loom-spindle-service + - name: loom-spindle-service port: 6555 diff --git a/config/manager/grpc_service.yaml b/config/manager/grpc_service.yaml new file mode 100644 index 0000000..07bc363 --- /dev/null +++ b/config/manager/grpc_service.yaml @@ -0,0 +1,20 @@ +--- +apiVersion: v1 +kind: Service +metadata: + name: controller-manager-grpc + namespace: system + labels: + app.kubernetes.io/name: loom + app.kubernetes.io/component: grpc + app.kubernetes.io/managed-by: kustomize +spec: + selector: + control-plane: controller-manager + app.kubernetes.io/name: loom + ports: + - name: grpc + port: 9090 + protocol: TCP + targetPort: 9090 + type: ClusterIP diff --git a/config/manager/kustomization.yaml b/config/manager/kustomization.yaml index 150f228..9116162 100644 --- a/config/manager/kustomization.yaml +++ b/config/manager/kustomization.yaml @@ -1,11 +1,12 @@ resources: - manager.yaml - service.yaml +- grpc_service.yaml - pvc.yaml - loom-config.yaml apiVersion: kustomize.config.k8s.io/v1beta1 kind: Kustomization images: - name: controller - newName: atcr.io/evan.jarrett.net/loom + newName: buoy.cr/evan.jarrett.net/loom newTag: latest diff --git a/config/manager/loom-config.yaml b/config/manager/loom-config.yaml index a820ba5..cb1bd26 100644 --- a/config/manager/loom-config.yaml +++ b/config/manager/loom-config.yaml @@ -1,7 +1,7 @@ apiVersion: v1 kind: ConfigMap metadata: - name: loom-config + name: config namespace: system data: config.yaml: | diff --git a/config/manager/manager.yaml b/config/manager/manager.yaml index d0799a1..6287394 100644 --- a/config/manager/manager.yaml +++ b/config/manager/manager.yaml @@ -66,7 +66,7 @@ spec: - /manager args: - --health-probe-bind-address=:8081 - image: atcr.io/evan.jarrett.net/loom:latest + image: buoy.cr/evan.jarrett.net/loom:latest imagePullPolicy: Always name: manager env: @@ -75,7 +75,7 @@ spec: fieldRef: fieldPath: metadata.namespace - name: LOOM_IMAGE - value: "atcr.io/evan.jarrett.net/loom:latest" + value: "buoy.cr/evan.jarrett.net/loom:latest" - name: SPINDLE_SERVER_HOSTNAME value: "loom.jarrett.net" - name: SPINDLE_SERVER_OWNER @@ -119,6 +119,8 @@ spec: - name: loom-config mountPath: /etc/loom readOnly: true + - name: scratch + mountPath: /scratch volumes: - name: spindle-logs persistentVolumeClaim: @@ -129,5 +131,7 @@ spec: - name: loom-config configMap: name: loom-config + - name: scratch + emptyDir: {} serviceAccountName: controller-manager terminationGracePeriodSeconds: 10 diff --git a/config/manager/service.yaml b/config/manager/service.yaml index f5c41bd..34af82f 100644 --- a/config/manager/service.yaml +++ b/config/manager/service.yaml @@ -2,7 +2,7 @@ apiVersion: v1 kind: Service metadata: - name: loom-spindle-service + name: spindle-service namespace: system labels: app.kubernetes.io/name: loom diff --git a/go.mod b/go.mod index 9e79536..38e93ee 100644 --- a/go.mod +++ b/go.mod @@ -5,8 +5,11 @@ go 1.25.5 require ( github.com/cenkalti/backoff/v4 v4.3.0 github.com/cyphar/filepath-securejoin v0.6.1 + github.com/go-logr/logr v1.4.3 github.com/onsi/ginkgo/v2 v2.28.1 - github.com/onsi/gomega v1.39.0 + github.com/onsi/gomega v1.39.1 + google.golang.org/grpc v1.80.0 + google.golang.org/protobuf v1.36.11 gopkg.in/yaml.v3 v3.0.1 k8s.io/api v0.35.3 k8s.io/apimachinery v0.35.3 @@ -17,7 +20,7 @@ require ( require ( cel.dev/expr v0.25.1 // indirect - github.com/Blank-Xu/sql-adapter v1.2.1 // indirect + github.com/Blank-Xu/sql-adapter v1.1.1 // indirect github.com/Masterminds/semver/v3 v3.4.0 // indirect github.com/Microsoft/go-winio v0.6.2 // indirect github.com/antlr4-go/antlr/v4 v4.13.1 // indirect @@ -48,8 +51,7 @@ require ( github.com/bluesky-social/jetstream v0.0.0-20260226214936-e0274250f654 // indirect github.com/bmatcuk/doublestar/v4 v4.10.0 // indirect github.com/carlmjohnson/versioninfo v0.22.5 // indirect - github.com/casbin/casbin/v2 v2.135.0 // indirect - github.com/casbin/casbin/v3 v3.10.0 // indirect + github.com/casbin/casbin/v2 v2.103.0 // indirect github.com/casbin/govaluate v1.10.0 // indirect github.com/cenkalti/backoff/v5 v5.0.3 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect @@ -79,7 +81,6 @@ require ( github.com/go-git/go-git/v5 v5.17.2 // indirect github.com/go-jose/go-jose/v4 v4.1.4 // indirect github.com/go-logfmt/logfmt v0.6.1 // indirect - github.com/go-logr/logr v1.4.3 // indirect github.com/go-logr/stdr v1.2.2 // indirect github.com/go-logr/zapr v1.3.0 // indirect github.com/go-openapi/jsonpointer v0.22.5 // indirect @@ -211,8 +212,6 @@ require ( gomodules.xyz/jsonpatch/v2 v2.5.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260401024825-9d38bb4040a9 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260401024825-9d38bb4040a9 // indirect - google.golang.org/grpc v1.80.0 // indirect - google.golang.org/protobuf v1.36.11 // indirect gopkg.in/evanphx/json-patch.v4 v4.13.0 // indirect gopkg.in/fsnotify.v1 v1.4.7 // indirect gopkg.in/inf.v0 v0.9.1 // indirect diff --git a/go.sum b/go.sum index 6056d12..b03cef4 100644 --- a/go.sum +++ b/go.sum @@ -41,8 +41,8 @@ cloud.google.com/go/storage v1.10.0/go.mod h1:FLPqc6j+Ki4BU591ie1oL6qBQGu2Bl/tZ9 dmitri.shuralyov.com/gpu/mtl v0.0.0-20190408044501-666a987793e9/go.mod h1:H6x//7gZCb22OMCxBHrMx7a5I7Hp++hsVxbQ4BYO7hU= github.com/Azure/go-ansiterm v0.0.0-20230124172434-306776ec8161 h1:L/gRVlceqvL25UVaW/CKtUDjefjrs0SPonmDGUVOYP0= github.com/Azure/go-ansiterm v0.0.0-20230124172434-306776ec8161/go.mod h1:xomTg63KZ2rFqZQzSB4Vz2SUXa1BpHTVz9L5PTmPC4E= -github.com/Blank-Xu/sql-adapter v1.2.1 h1:Gl9CZI3PCDLg2EKvmYFbieOe95IRJMuruj7AL9JXsLk= -github.com/Blank-Xu/sql-adapter v1.2.1/go.mod h1:Duskd1ORzVkmxOxk6i6HSAdmASjqVhg9fcAefibnrns= +github.com/Blank-Xu/sql-adapter v1.1.1 h1:+g7QXU9sl/qT6Po97teMpf3GjAO0X9aFaqgSePXvYko= +github.com/Blank-Xu/sql-adapter v1.1.1/go.mod h1:o2g8EZhZ3TudnYEGDkoU+3jCTCgDgx1o/Ig5ajKkaLY= github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03qcyfWMU= github.com/BurntSushi/xgb v0.0.0-20160522181843-27f122750802/go.mod h1:IVnqGOEym/WlBOVXweHU+Q+/VP0lqqI8lqeDx9IjBqo= github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1Xbatp0= @@ -112,7 +112,7 @@ github.com/bluesky-social/indigo v0.0.0-20260318212431-cbaa83aee9dd/go.mod h1:VG github.com/bluesky-social/jetstream v0.0.0-20260226214936-e0274250f654 h1:OK76FcHhZp8ohjRB0OMWgti0oYAWFlt3KDQcIkH1pfI= github.com/bluesky-social/jetstream v0.0.0-20260226214936-e0274250f654/go.mod h1:vt8kVRKtvrBspt9G38wDD8+BotjIMO8u8IYoVnyE4zY= github.com/bmatcuk/doublestar/v4 v4.6.1/go.mod h1:xBQ8jztBU6kakFMg+8WGxn0c6z1fTSPVIjEY1Wr7jzc= -github.com/bmatcuk/doublestar/v4 v4.9.1/go.mod h1:xBQ8jztBU6kakFMg+8WGxn0c6z1fTSPVIjEY1Wr7jzc= +github.com/bmatcuk/doublestar/v4 v4.7.1/go.mod h1:xBQ8jztBU6kakFMg+8WGxn0c6z1fTSPVIjEY1Wr7jzc= github.com/bmatcuk/doublestar/v4 v4.10.0 h1:zU9WiOla1YA122oLM6i4EXvGW62DvKZVxIe6TYWexEs= github.com/bmatcuk/doublestar/v4 v4.10.0/go.mod h1:xBQ8jztBU6kakFMg+8WGxn0c6z1fTSPVIjEY1Wr7jzc= github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= @@ -121,10 +121,10 @@ github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= github.com/carlmjohnson/versioninfo v0.22.5 h1:O00sjOLUAFxYQjlN/bzYTuZiS0y6fWDQjMRvwtKgwwc= github.com/carlmjohnson/versioninfo v0.22.5/go.mod h1:QT9mph3wcVfISUKd0i9sZfVrPviHuSF+cUtLjm2WSf8= -github.com/casbin/casbin/v2 v2.135.0 h1:6BLkMQiGotYyS5yYeWgW19vxqugUlvHFkFiLnLR/bxk= -github.com/casbin/casbin/v2 v2.135.0/go.mod h1:FmcfntdXLTcYXv/hxgNntcRPqAbwOG9xsism0yXT+18= -github.com/casbin/casbin/v3 v3.10.0 h1:039ORla55vCeIZWd0LfzWFt1yiEA5X4W41xBW2bQuHs= -github.com/casbin/casbin/v3 v3.10.0/go.mod h1:5rJbQr2e6AuuDDNxnPc5lQlC9nIgg6nS1zYwKXhpHC8= +github.com/casbin/casbin/v2 v2.100.0/go.mod h1:LO7YPez4dX3LgoTCqSQAleQDo0S0BeZBDxYnPUl95Ng= +github.com/casbin/casbin/v2 v2.103.0 h1:dHElatNXNrr8XcseUov0ZSiWjauwmZZE6YMV3eU1yic= +github.com/casbin/casbin/v2 v2.103.0/go.mod h1:Ee33aqGrmES+GNL17L0h9X28wXuo829wnNUnS0edAco= +github.com/casbin/govaluate v1.2.0/go.mod h1:G/UnbIjZk/0uMNaLwZZmFQrR72tYRZWQkO70si/iR7A= github.com/casbin/govaluate v1.3.0/go.mod h1:G/UnbIjZk/0uMNaLwZZmFQrR72tYRZWQkO70si/iR7A= github.com/casbin/govaluate v1.10.0 h1:ffGw51/hYH3w3rZcxO/KcaUIDOLP84w7nsidMVgaDG0= github.com/casbin/govaluate v1.10.0/go.mod h1:G/UnbIjZk/0uMNaLwZZmFQrR72tYRZWQkO70si/iR7A= @@ -584,8 +584,8 @@ github.com/onsi/gomega v1.22.1/go.mod h1:x6n7VNe4hw0vkyYUM4mjIXx3JbLiPaBPNgB7PRQ github.com/onsi/gomega v1.24.0/go.mod h1:Z/NWtiqwBrwUt4/2loMmHL63EDLnYHmVbuBpDr2vQAg= github.com/onsi/gomega v1.24.1/go.mod h1:3AOiACssS3/MajrniINInwbfOOtfZvplPzuRSmvt1jM= github.com/onsi/gomega v1.25.0/go.mod h1:r+zV744Re+DiYCIPRlYOTxn0YkOLcAnW8k1xXdMPGhM= -github.com/onsi/gomega v1.39.0 h1:y2ROC3hKFmQZJNFeGAMeHZKkjBL65mIZcvrLQBF9k6Q= -github.com/onsi/gomega v1.39.0/go.mod h1:ZCU1pkQcXDO5Sl9/VVEGlDyp+zm0m1cmeG5TOzLgdh4= +github.com/onsi/gomega v1.39.1 h1:1IJLAad4zjPn2PsnhH70V4DKRFlrCzGBNrNaru+Vf28= +github.com/onsi/gomega v1.39.1/go.mod h1:hL6yVALoTOxeWudERyfppUcZXjMwIMLnuSfruD2lcfg= github.com/openbao/openbao/api/v2 v2.5.1 h1:Br79D6L20SbAa5P7xqENxmvv8LyI4HoKosPy7klhn4o= github.com/openbao/openbao/api/v2 v2.5.1/go.mod h1:Dh5un77tqGgMbmlVEqjqN+8/dMyUohnkaQVg/wXW0Ig= github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U= diff --git a/helm/loom/Chart.yaml b/helm/loom/Chart.yaml index 326fa7b..7882876 100644 --- a/helm/loom/Chart.yaml +++ b/helm/loom/Chart.yaml @@ -2,8 +2,8 @@ apiVersion: v2 name: loom description: A Kubernetes operator that runs CI/CD pipelines from tangled.org type: application -version: 0.0.1 -appVersion: "0.0.1" +version: 0.1.5 +appVersion: "0.1.5" home: https://github.com/tangled-sh/loom sources: - https://github.com/tangled-sh/loom diff --git a/helm/loom/crds/loom.j5t.io_spindlesets.yaml b/helm/loom/crds/loom.j5t.io_spindlesets.yaml index ccd9d61..1d22154 100644 --- a/helm/loom/crds/loom.j5t.io_spindlesets.yaml +++ b/helm/loom/crds/loom.j5t.io_spindlesets.yaml @@ -77,6 +77,11 @@ spec: items: type: string type: array + multiArch: + description: |- + MultiArch indicates this pipeline run contains multi-arch workflows. + When true, the controller creates per-architecture Jobs and gates the final Job. + type: boolean pipelineID: description: PipelineID is the unique identifier for this pipeline run from the knot. @@ -110,12 +115,15 @@ spec: container entirely. type: boolean workflows: - description: Workflows is the list of workflows to execute in - this pipeline. + description: |- + Workflows is the list of workflows to execute in this pipeline. + For multi-arch workflows, this contains one entry per matrix leg plus an optional final entry. items: description: |- WorkflowSpec defines a workflow to execute as part of a pipeline. This is the canonical workflow definition that matches the .tangled/workflows/*.yaml format. + For multi-arch workflows, the engine expands the matrix and creates one WorkflowSpec per leg. + Each leg has a single Image and Architecture; the matrix metadata lives in PipelineRunSpec. properties: architecture: description: Architecture is the target architecture for @@ -141,10 +149,65 @@ spec: description: Environment contains workflow-level environment variables. type: object + final: + description: |- + Final defines steps that run once after all matrix legs complete. + Only valid on multi-arch workflows. The engine sets this on the dedicated final WorkflowSpec. + properties: + architecture: + description: Architecture is the target architecture + for the final steps. + enum: + - amd64 + - arm64 + type: string + image: + description: |- + Image is the container image for the final steps. + If empty, uses the first image from the matrix. + type: string + steps: + description: Steps is the ordered list of steps to execute + after all matrix legs complete. + items: + description: WorkflowStep defines a single step in + a workflow. + properties: + command: + description: Command is the shell command to execute. + type: string + environment: + additionalProperties: + type: string + description: Environment contains step-specific + environment variables. + type: object + name: + description: Name is the human-readable name of + the step. + type: string + required: + - command + - name + type: object + minItems: 1 + type: array + required: + - architecture + - steps + type: object image: description: Image is the container image to use for executing the workflow steps. type: string + isFinal: + description: IsFinal indicates this WorkflowSpec represents + the final step of a multi-arch workflow. + type: boolean + isMatrixLeg: + description: IsMatrixLeg indicates this WorkflowSpec was + generated from a matrix expansion. + type: boolean name: description: Name is the workflow filename (e.g., "workflow-amd64.yaml"). type: string @@ -1240,9 +1303,10 @@ spec: operator: description: |- Operator represents a key's relationship to the value. - Valid operators are Exists and Equal. Defaults to Equal. + Valid operators are Exists, Equal, Lt, and Gt. Defaults to Equal. Exists is equivalent to wildcard for value, so that a pod can tolerate all taints of a particular category. + Lt and Gt perform numeric comparisons (requires feature gate TaintTolerationComparisonOperators). type: string tolerationSeconds: description: |- @@ -1359,10 +1423,45 @@ spec: description: CompletionTime is when the workflow finished. format: date-time type: string + finalJobName: + description: FinalJobName is the name of the final Job (for + multi-arch workflows). + type: string + finalPhase: + description: FinalPhase is the phase of the final Job. + type: string jobName: - description: JobName is the name of the Kubernetes Job created - for this workflow. + description: |- + JobName is the name of the Kubernetes Job created for this workflow. + For multi-arch workflows, this is empty; use MatrixLegStatuses instead. type: string + matrixLegStatuses: + description: MatrixLegStatuses tracks per-architecture Job statuses + for multi-arch workflows. + items: + description: MatrixLegStatus tracks the status of a single + matrix leg Job. + properties: + architecture: + description: Architecture is the target architecture for + this leg. + type: string + image: + description: Image is the container image used for this + leg. + type: string + jobName: + description: JobName is the name of the Kubernetes Job + for this leg. + type: string + phase: + description: Phase is the current phase (Pending, Running, + Succeeded, Failed). + type: string + required: + - architecture + type: object + type: array name: description: Name is the workflow name. type: string diff --git a/helm/loom/templates/grpc_service.yaml b/helm/loom/templates/grpc_service.yaml new file mode 100644 index 0000000..542cfd9 --- /dev/null +++ b/helm/loom/templates/grpc_service.yaml @@ -0,0 +1,17 @@ +apiVersion: v1 +kind: Service +metadata: + name: {{ include "loom.fullname" . }}-grpc + namespace: {{ .Release.Namespace }} + labels: + {{- include "loom.labels" . | nindent 4 }} + app.kubernetes.io/component: grpc +spec: + type: ClusterIP + selector: + {{- include "loom.controllerLabels" . | nindent 4 }} + ports: + - name: grpc + port: {{ .Values.grpc.port | default 9090 }} + protocol: TCP + targetPort: {{ .Values.grpc.port | default 9090 }} diff --git a/helm/loom/values.yaml b/helm/loom/values.yaml index 2ec5440..cd7e7e7 100644 --- a/helm/loom/values.yaml +++ b/helm/loom/values.yaml @@ -2,7 +2,7 @@ # Image configuration image: - repository: atcr.io/evan.jarrett.net/loom + repository: buoy.cr/evan.jarrett.net/loom pullPolicy: Always # Overrides the image tag whose default is the chart appVersion. tag: "" @@ -105,6 +105,10 @@ service: type: ClusterIP port: 6555 +# gRPC service for runner communication +grpc: + port: 9090 + # RBAC configuration rbac: create: true diff --git a/internal/controller/spindleset_controller.go b/internal/controller/spindleset_controller.go index f20515f..a71c810 100644 --- a/internal/controller/spindleset_controller.go +++ b/internal/controller/spindleset_controller.go @@ -24,6 +24,7 @@ import ( "time" "github.com/cenkalti/backoff/v4" + "github.com/go-logr/logr" "tangled.org/core/spindle" "tangled.org/core/spindle/models" @@ -33,6 +34,7 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/client-go/rest" + "k8s.io/client-go/util/retry" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" @@ -58,6 +60,9 @@ type SpindleSetReconciler struct { // Set from LOOM_IMAGE environment variable LoomImage string + // OperatorAddr is the gRPC address of the operator for runner communication + OperatorAddr string + // Track watched Jobs for status reporting watchedJobs sync.Map // map[string]models.WorkflowId } @@ -175,6 +180,11 @@ func (r *SpindleSetReconciler) Reconcile(ctx context.Context, req ctrl.Request) logger.Error(err, "Failed to monitor job statuses") } + // For multi-arch pipelines, create final Jobs once all matrix legs succeed + if err := r.ensureFinalJobs(ctx, spindleSet); err != nil { + logger.Error(err, "Failed to ensure final Jobs") + } + if jobsErr != nil { return ctrl.Result{}, jobsErr } @@ -264,20 +274,24 @@ func (r *SpindleSetReconciler) retryCreate(ctx context.Context, obj client.Objec return backoff.Retry(operation, backoff.WithContext(bo, ctx)) } -// updateStatus updates the SpindleSet status based on current Jobs +// updateStatus updates the SpindleSet status based on current Jobs. +// Status updates retry on conflict to handle races with concurrent reconciles +// (e.g., triggered by Job creation/updates via Owns). func (r *SpindleSetReconciler) updateStatus(ctx context.Context, spindleSet *loomv1alpha1.SpindleSet) error { logger := log.FromContext(ctx) - // Re-fetch the SpindleSet to get the latest version before updating status - // This avoids optimistic concurrency conflicts when the object was modified - // by another reconciliation loop (e.g., triggered by Job creation/updates) - latestSpindleSet := &loomv1alpha1.SpindleSet{} - if err := r.Get(ctx, client.ObjectKeyFromObject(spindleSet), latestSpindleSet); err != nil { - return fmt.Errorf("failed to fetch latest SpindleSet: %w", err) - } - // Use the latest version for all subsequent operations - spindleSet = latestSpindleSet + return retry.RetryOnConflict(retry.DefaultRetry, func() error { + // Re-fetch the SpindleSet on each attempt to get the latest version + latestSpindleSet := &loomv1alpha1.SpindleSet{} + if err := r.Get(ctx, client.ObjectKeyFromObject(spindleSet), latestSpindleSet); err != nil { + return fmt.Errorf("failed to fetch latest SpindleSet: %w", err) + } + return r.computeAndApplyStatus(ctx, logger, latestSpindleSet) + }) +} +// computeAndApplyStatus recomputes status from Jobs and applies it to spindleSet. +func (r *SpindleSetReconciler) computeAndApplyStatus(ctx context.Context, logger logr.Logger, spindleSet *loomv1alpha1.SpindleSet) error { // List all Jobs owned by this SpindleSet jobList := &batchv1.JobList{} if err := r.List(ctx, jobList, client.InNamespace(spindleSet.Namespace), client.MatchingLabels{ @@ -474,83 +488,190 @@ func (r *SpindleSetReconciler) ensurePipelineJobs(ctx context.Context, spindleSe return fmt.Errorf("failed to list nodes: %w", err) } - // Convert workflow steps to jobbuilder format and create Jobs for each workflow + // Convert workflow steps to jobbuilder format and create Jobs for each workflow. + // For multi-arch pipelines, final Jobs are gated on all matrix leg Jobs succeeding. for _, workflowSpec := range pipelineRun.Workflows { - // Check if Job already exists - jobName := fmt.Sprintf("spindle-%s-%s", pipelineRun.PipelineID, workflowSpec.Name) - if len(jobName) > 63 { - jobName = jobName[:63] + // For multi-arch: skip final workflows initially — they are created + // by ensureFinalJobs once all matrix leg Jobs have succeeded. + if workflowSpec.IsFinal { + continue } - existingJob := &batchv1.Job{} - err := r.Get(ctx, client.ObjectKey{ - Name: jobName, - Namespace: spindleSet.Namespace, - }, existingJob) - - if err == nil { - // Job already exists - logger.V(1).Info("Job already exists for workflow", "workflow", workflowSpec.Name, "job", jobName) - continue + if err := r.createWorkflowJob(ctx, spindleSet, pipelineRun, workflowSpec, secretName, secretKeys, &nodeList); err != nil { + return err } + } - if !apierrors.IsNotFound(err) { - return fmt.Errorf("failed to check for existing job: %w", err) - } - - // Convert workflow steps to jobbuilder format - jobSteps := make([]jobbuilder.WorkflowStep, 0, len(workflowSpec.Steps)) - for _, step := range workflowSpec.Steps { - jobSteps = append(jobSteps, jobbuilder.WorkflowStep{ - Name: step.Name, - Command: step.Command, - Env: step.Environment, - }) - } - - // Build Job configuration - jobConfig := jobbuilder.WorkflowConfig{ - WorkflowName: workflowSpec.Name, - PipelineID: pipelineRun.PipelineID, - SpindleSetName: spindleSet.Name, - Image: workflowSpec.Image, - LoomImage: r.LoomImage, - Architecture: workflowSpec.Architecture, - Steps: jobSteps, - WorkflowSpec: workflowSpec, // Pass full workflow spec to runner - CloneCommands: pipelineRun.CloneCommands, - SkipClone: pipelineRun.SkipClone, - SecretName: secretName, // Name of K8s Secret to inject (empty if no secrets) - SecretKeys: secretKeys, // Secret env var names for log masking - Template: spindleSet.Spec.Template, - Namespace: spindleSet.Namespace, - } - - // Create the Job - job, err := jobbuilder.BuildJob(jobConfig, &nodeList) - if err != nil { - return fmt.Errorf("failed to build job for workflow %s: %w", workflowSpec.Name, err) + return nil +} + +// createWorkflowJob creates a single Kubernetes Job for a workflow spec. +func (r *SpindleSetReconciler) createWorkflowJob(ctx context.Context, spindleSet *loomv1alpha1.SpindleSet, pipelineRun *loomv1alpha1.PipelineRunSpec, workflowSpec loomv1alpha1.WorkflowSpec, secretName string, secretKeys []string, nodeList *corev1.NodeList) error { + logger := log.FromContext(ctx) + + // Check if Job already exists + jobName := fmt.Sprintf("spindle-%s-%s", pipelineRun.PipelineID, workflowSpec.Name) + if len(jobName) > 63 { + jobName = jobName[:63] + } + + existingJob := &batchv1.Job{} + err := r.Get(ctx, client.ObjectKey{ + Name: jobName, + Namespace: spindleSet.Namespace, + }, existingJob) + + if err == nil { + logger.V(1).Info("Job already exists for workflow", "workflow", workflowSpec.Name, "job", jobName) + return nil + } + + if !apierrors.IsNotFound(err) { + return fmt.Errorf("failed to check for existing job: %w", err) + } + + // Convert workflow steps to jobbuilder format + jobSteps := make([]jobbuilder.WorkflowStep, 0, len(workflowSpec.Steps)) + for _, step := range workflowSpec.Steps { + jobSteps = append(jobSteps, jobbuilder.WorkflowStep{ + Name: step.Name, + Command: step.Command, + Env: step.Environment, + }) + } + + // Build Job configuration + jobConfig := jobbuilder.WorkflowConfig{ + WorkflowName: workflowSpec.Name, + PipelineID: pipelineRun.PipelineID, + SpindleSetName: spindleSet.Name, + Image: workflowSpec.Image, + LoomImage: r.LoomImage, + Architecture: workflowSpec.Architecture, + Steps: jobSteps, + WorkflowSpec: workflowSpec, + CloneCommands: pipelineRun.CloneCommands, + SkipClone: pipelineRun.SkipClone, + SecretName: secretName, + SecretKeys: secretKeys, + Template: spindleSet.Spec.Template, + Namespace: spindleSet.Namespace, + OperatorAddr: r.OperatorAddr, + } + + // Create the Job + job, err := jobbuilder.BuildJob(jobConfig, nodeList) + if err != nil { + return fmt.Errorf("failed to build job for workflow %s: %w", workflowSpec.Name, err) + } + + // Set SpindleSet as owner of the Job + if err := controllerutil.SetControllerReference(spindleSet, job, r.Scheme); err != nil { + return fmt.Errorf("failed to set controller reference: %w", err) + } + + logger.Info("Creating Job for workflow", "workflow", workflowSpec.Name, "job", job.Name) + if err := r.retryCreate(ctx, job); err != nil { + if apierrors.IsAlreadyExists(err) { + logger.Info("Job already exists, skipping creation", "workflow", workflowSpec.Name, "job", job.Name) + return nil } + return fmt.Errorf("failed to create job for workflow %s: %w", workflowSpec.Name, err) + } - // Set SpindleSet as owner of the Job - if err := controllerutil.SetControllerReference(spindleSet, job, r.Scheme); err != nil { - return fmt.Errorf("failed to set controller reference: %w", err) + logger.Info("Job created successfully", "workflow", workflowSpec.Name, "job", job.Name) + return nil +} + +// ensureFinalJobs creates final Jobs for multi-arch pipelines once all matrix legs have succeeded. +// This is called during reconciliation after job status monitoring. +func (r *SpindleSetReconciler) ensureFinalJobs(ctx context.Context, spindleSet *loomv1alpha1.SpindleSet) error { + logger := log.FromContext(ctx) + pipelineRun := spindleSet.Spec.PipelineRun + + if !pipelineRun.MultiArch { + return nil + } + + // Find the final workflow spec + var finalSpec *loomv1alpha1.WorkflowSpec + for i := range pipelineRun.Workflows { + if pipelineRun.Workflows[i].IsFinal { + finalSpec = &pipelineRun.Workflows[i] + break } + } + if finalSpec == nil { + return nil // No final step defined + } - logger.Info("Creating Job for workflow", "workflow", workflowSpec.Name, "job", job.Name) - if err := r.retryCreate(ctx, job); err != nil { - if apierrors.IsAlreadyExists(err) { - // Job already exists (possibly from previous deployment), skip - logger.Info("Job already exists, skipping creation", "workflow", workflowSpec.Name, "job", job.Name) - continue + // Check if final Job already exists + finalJobName := fmt.Sprintf("spindle-%s-%s", pipelineRun.PipelineID, finalSpec.Name) + if len(finalJobName) > 63 { + finalJobName = finalJobName[:63] + } + existingJob := &batchv1.Job{} + if err := r.Get(ctx, client.ObjectKey{Name: finalJobName, Namespace: spindleSet.Namespace}, existingJob); err == nil { + return nil // Already created + } + + // Check if all matrix leg Jobs have succeeded + jobList := &batchv1.JobList{} + if err := r.List(ctx, jobList, + client.InNamespace(spindleSet.Namespace), + client.MatchingLabels{"loom.j5t.io/spindleset": spindleSet.Name}, + ); err != nil { + return fmt.Errorf("failed to list jobs: %w", err) + } + + legCount := 0 + succeededCount := 0 + for _, job := range jobList.Items { + // Count only matrix leg jobs (not the final job) + wfName := job.Labels["loom.j5t.io/workflow"] + for _, wf := range pipelineRun.Workflows { + if wf.Name == wfName && wf.IsMatrixLeg { + legCount++ + if job.Status.Succeeded > 0 { + succeededCount++ + } + break } - return fmt.Errorf("failed to create job for workflow %s: %w", workflowSpec.Name, err) } + } - logger.Info("Job created successfully", "workflow", workflowSpec.Name, "job", job.Name) + // Count expected legs + expectedLegs := 0 + for _, wf := range pipelineRun.Workflows { + if wf.IsMatrixLeg { + expectedLegs++ + } } - return nil + if legCount < expectedLegs || succeededCount < expectedLegs { + logger.V(1).Info("Waiting for matrix legs to complete", + "expected", expectedLegs, "found", legCount, "succeeded", succeededCount) + return nil // Not all legs have completed yet + } + + logger.Info("All matrix legs succeeded, creating final Job", "legs", expectedLegs) + + // Build secret info + secretName := "" + var secretKeys []string + if len(pipelineRun.Secrets) > 0 { + secretName = fmt.Sprintf("%s-secrets", spindleSet.Name) + for _, s := range pipelineRun.Secrets { + secretKeys = append(secretKeys, s.Key) + } + } + + var nodeList corev1.NodeList + if err := r.List(ctx, &nodeList); err != nil { + return fmt.Errorf("failed to list nodes: %w", err) + } + + return r.createWorkflowJob(ctx, spindleSet, pipelineRun, *finalSpec, secretName, secretKeys, &nodeList) } // cleanupOrphanedJobs cleans up Jobs without a matching SpindleSet diff --git a/internal/engine/kubernetes_engine.go b/internal/engine/kubernetes_engine.go index 29fe292..8f72f13 100644 --- a/internal/engine/kubernetes_engine.go +++ b/internal/engine/kubernetes_engine.go @@ -1,22 +1,15 @@ package engine import ( - "bufio" "context" - "encoding/json" "fmt" - "io" "maps" "strings" - "sync" "time" securejoin "github.com/cyphar/filepath-securejoin" "gopkg.in/yaml.v3" - batchv1 "k8s.io/api/batch/v1" - corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/log" @@ -26,20 +19,34 @@ import ( "tangled.org/core/spindle/secrets" loomv1alpha1 "tangled.org/evan.jarrett.net/loom/api/v1alpha1" + loomgrpc "tangled.org/evan.jarrett.net/loom/internal/grpc" ) -// workflowLogStream holds the state for streaming logs from a workflow's pod -type workflowLogStream struct { - scanner *bufio.Scanner - stream io.ReadCloser - pod *corev1.Pod - podPhase corev1.PodPhase // Track pod phase at stream creation time +// syntheticStep is a minimal implementation of models.Step used to emit +// ControlWriter entries for matrix-leg user steps (which are invisible to +// the upstream spindle framework because loom wraps them in synthetic +// "Matrix build" / "Final" framework steps). +type syntheticStep struct { + name string + command string + kind models.StepKind } -// extendedLogLine extends models.LogLine with exit code for error reporting -type extendedLogLine struct { - models.LogLine - ExitCode int `json:"exit_code,omitempty"` +func (s syntheticStep) Name() string { return s.name } +func (s syntheticStep) Command() string { return s.command } +func (s syntheticStep) Kind() models.StepKind { return s.kind } + +// matrixLegLogStepID returns a collision-free log step id for a matrix leg's +// user step. The offset pushes us past the two framework steps (0, 1) that +// the upstream engine already emits ("Matrix build", "Final"). +func matrixLegLogStepID(legIdx, userStepIdx int) int { + return 1000 + legIdx*100 + userStepIdx +} + +// finalLogStepID returns a collision-free log step id for a final-phase step, +// placed after the matrix-leg id range. +func finalLogStepID(numLegs, stepIdx int) int { + return 1000 + (numLegs+1)*100 + stepIdx } // KubernetesEngine implements the spindle Engine interface for Kubernetes Jobs. @@ -49,33 +56,82 @@ type KubernetesEngine struct { namespace string template loomv1alpha1.SpindleTemplate vault secrets.Manager + hub *loomgrpc.Hub + artifacts *loomgrpc.ArtifactStore // Track created SpindleSets for cleanup spindleSets map[string]*loomv1alpha1.SpindleSet - - // Active log streams per workflow - persist across RunStep calls - logStreams map[string]*workflowLogStream - streamMutex sync.RWMutex } // NewKubernetesEngine creates a new Kubernetes-based spindle engine. -func NewKubernetesEngine(k8sClient client.Client, config *rest.Config, namespace string, template loomv1alpha1.SpindleTemplate, vault secrets.Manager) *KubernetesEngine { +func NewKubernetesEngine(k8sClient client.Client, config *rest.Config, namespace string, template loomv1alpha1.SpindleTemplate, vault secrets.Manager, hub *loomgrpc.Hub, artifacts *loomgrpc.ArtifactStore) *KubernetesEngine { return &KubernetesEngine{ client: k8sClient, config: config, namespace: namespace, template: template, vault: vault, + hub: hub, + artifacts: artifacts, spindleSets: make(map[string]*loomv1alpha1.SpindleSet), - logStreams: make(map[string]*workflowLogStream), } } +// StringOrSlice is a YAML type that accepts either a single string or an array of strings. +type StringOrSlice []string + +func (s *StringOrSlice) UnmarshalYAML(node *yaml.Node) error { + if node.Kind == yaml.ScalarNode { + *s = []string{node.Value} + return nil + } + var slice []string + if err := node.Decode(&slice); err != nil { + return err + } + *s = slice + return nil +} + +// rawWorkflowSpec is used for initial YAML parsing before matrix expansion. +// It handles the polymorphic image/architecture fields. +type rawWorkflowSpec struct { + Image StringOrSlice `yaml:"image"` + Architecture StringOrSlice `yaml:"architecture"` + Steps []loomv1alpha1.WorkflowStep `yaml:"steps"` + When []loomv1alpha1.WorkflowWhen `yaml:"when"` + Environment map[string]string `yaml:"environment"` + Dependencies *loomv1alpha1.WorkflowDependencies `yaml:"dependencies"` + Final *loomv1alpha1.FinalSpec `yaml:"final"` +} + +// MatrixLeg represents a single combination of image and architecture. +type MatrixLeg struct { + Image string + Architecture string +} + +// expandMatrix computes the cartesian product of images and architectures. +func expandMatrix(images, architectures StringOrSlice) []MatrixLeg { + var legs []MatrixLeg + for _, img := range images { + for _, arch := range architectures { + legs = append(legs, MatrixLeg{Image: img, Architecture: arch}) + } + } + return legs +} + // kubernetesWorkflowData holds pre-computed data for workflow execution. // Built in InitWorkflow, consumed in SetupWorkflow. type kubernetesWorkflowData struct { Spec loomv1alpha1.WorkflowSpec CloneStep models.CloneStep // empty if clone should be skipped + + // Multi-arch fields + IsMultiArch bool + MatrixLegs []MatrixLeg + FinalSpec *loomv1alpha1.FinalSpec } // SimpleStep implements the models.Step interface. @@ -100,27 +156,50 @@ func (s SimpleStep) Kind() models.StepKind { // InitWorkflow parses the workflow YAML and initializes a Workflow model. // Pipeline environment variables (TANGLED_*) are injected into workflow.Environment // by the framework after this method returns. +// +// For multi-arch workflows (image or architecture is an array), the matrix is expanded +// and stored in kubernetesWorkflowData. The framework sees synthetic steps: +// - Step 0: "Matrix build" (engine fans out to N parallel Jobs, waits for all) +// - Step 1: "Final" (engine creates final Job, waits for completion) — only if final block exists func (e *KubernetesEngine) InitWorkflow(twf tangled.Pipeline_Workflow, tpl tangled.Pipeline) (*models.Workflow, error) { - // Parse the Raw YAML into the unified WorkflowSpec type - var spec loomv1alpha1.WorkflowSpec - if err := yaml.Unmarshal([]byte(twf.Raw), &spec); err != nil { + // Parse YAML with polymorphic image/architecture handling + var raw rawWorkflowSpec + if err := yaml.Unmarshal([]byte(twf.Raw), &raw); err != nil { return nil, fmt.Errorf("failed to parse workflow YAML: %w", err) } - // Set the workflow name from the tangled workflow - spec.Name = twf.Name - - // Validate required fields - if spec.Image == "" { + if len(raw.Image) == 0 { return nil, fmt.Errorf("workflow must specify an 'image' field") } + if len(raw.Architecture) == 0 { + raw.Architecture = StringOrSlice{"amd64"} + } + + // Build clone step + var cloneStep models.CloneStep + if twf.Clone == nil || !twf.Clone.Skip { + cloneStep = models.BuildCloneStep(twf, *tpl.TriggerMetadata, false) + } + + // Determine if this is a multi-arch workflow + legs := expandMatrix(raw.Image, raw.Architecture) + isMultiArch := len(legs) > 1 || raw.Final != nil + + if isMultiArch { + return e.initMultiArchWorkflow(twf.Name, raw, legs, cloneStep) + } - // Default architecture to amd64 if not specified - if spec.Architecture == "" { - spec.Architecture = "amd64" + // Single-arch workflow: existing behavior + spec := loomv1alpha1.WorkflowSpec{ + Name: twf.Name, + Image: raw.Image[0], + Architecture: raw.Architecture[0], + Steps: raw.Steps, + When: raw.When, + Environment: raw.Environment, + Dependencies: raw.Dependencies, } - // Convert steps to models.Step interface modelSteps := make([]models.Step, 0, len(spec.Steps)) for _, stepSpec := range spec.Steps { modelSteps = append(modelSteps, SimpleStep{ @@ -130,36 +209,81 @@ func (e *KubernetesEngine) InitWorkflow(twf tangled.Pipeline_Workflow, tpl tangl }) } - // Build clone step (uses upstream models.BuildCloneStep which is self-contained) - var cloneStep models.CloneStep - devMode := false // TODO: Make this configurable - - if twf.Clone == nil || !twf.Clone.Skip { - cloneStep = models.BuildCloneStep(twf, *tpl.TriggerMetadata, devMode) - } - - // Store pre-computed workflow data workflowData := &kubernetesWorkflowData{ Spec: spec, CloneStep: cloneStep, } - // Set engine-specific environment variables on the workflow - // These will be merged with pipeline env vars by the framework workflowEnv := map[string]string{ "TANGLED_ARCHITECTURE": spec.Architecture, - // HOME must be writable; we run as user 10000 so default /root won't work - "HOME": "/tmp", + "HOME": "/tmp", } - workflow := &models.Workflow{ + return &models.Workflow{ Steps: modelSteps, Name: twf.Name, Data: workflowData, Environment: workflowEnv, + }, nil +} + +// initMultiArchWorkflow creates a Workflow with synthetic steps for matrix execution. +func (e *KubernetesEngine) initMultiArchWorkflow(name string, raw rawWorkflowSpec, legs []MatrixLeg, cloneStep models.CloneStep) (*models.Workflow, error) { + // Use the first leg's architecture for the spec that gets stored + // (the actual per-leg specs are built in SetupWorkflow) + spec := loomv1alpha1.WorkflowSpec{ + Name: name, + Image: raw.Image[0], + Architecture: raw.Architecture[0], + Steps: raw.Steps, + When: raw.When, + Environment: raw.Environment, + Dependencies: raw.Dependencies, + Final: raw.Final, + } + + // Build synthetic steps that the framework will iterate over + var modelSteps []models.Step + + // Step 0: matrix build phase (fans out to N parallel Jobs) + archList := make([]string, len(legs)) + for i, leg := range legs { + archList[i] = leg.Architecture + } + modelSteps = append(modelSteps, SimpleStep{ + StepName: fmt.Sprintf("Matrix build (%s)", strings.Join(archList, ", ")), + StepCommand: "# internal: matrix fan-out", + StepKind: models.StepKindUser, + }) + + // Step 1: final phase (if final block exists) + if raw.Final != nil { + modelSteps = append(modelSteps, SimpleStep{ + StepName: "Final", + StepCommand: "# internal: final step", + StepKind: models.StepKindUser, + }) + } + + workflowData := &kubernetesWorkflowData{ + Spec: spec, + CloneStep: cloneStep, + IsMultiArch: true, + MatrixLegs: legs, + FinalSpec: raw.Final, + } + + workflowEnv := map[string]string{ + "TANGLED_ARCHITECTURE": raw.Architecture[0], + "HOME": "/tmp", } - return workflow, nil + return &models.Workflow{ + Steps: modelSteps, + Name: name, + Data: workflowData, + Environment: workflowEnv, + }, nil } // SetupWorkflow creates a SpindleSet CR for the workflow. @@ -223,13 +347,72 @@ func (e *KubernetesEngine) SetupWorkflow(ctx context.Context, wid models.Workflo } // Build PipelineRunSpec from pre-computed data - // Knot is extracted from the pipeline ID provided by the framework skipClone := len(data.CloneStep.Commands()) == 0 + var workflows []loomv1alpha1.WorkflowSpec + + if data.IsMultiArch { + // Expand matrix into per-leg WorkflowSpecs + for _, leg := range data.MatrixLegs { + legSpec := loomv1alpha1.WorkflowSpec{ + Name: fmt.Sprintf("%s-%s", data.Spec.Name, leg.Architecture), + Image: leg.Image, + Architecture: leg.Architecture, + Steps: data.Spec.Steps, + When: data.Spec.When, + Environment: maps.Clone(data.Spec.Environment), + Dependencies: data.Spec.Dependencies, + IsMatrixLeg: true, + } + // Override per-leg env vars + if legSpec.Environment == nil { + legSpec.Environment = make(map[string]string) + } + legSpec.Environment["TANGLED_ARCHITECTURE"] = leg.Architecture + legSpec.Environment["TANGLED_IMAGE"] = leg.Image + legSpec.Environment["LOOM_MATRIX_LEG"] = "true" + legSpec.Environment["LOOM_ARTIFACTS"] = "/artifacts" + workflows = append(workflows, legSpec) + } + + // Add final WorkflowSpec if specified + if data.FinalSpec != nil { + finalImage := data.FinalSpec.Image + if finalImage == "" { + finalImage = data.MatrixLegs[0].Image + } + finalSpec := loomv1alpha1.WorkflowSpec{ + Name: fmt.Sprintf("%s-final", data.Spec.Name), + Image: finalImage, + Architecture: data.FinalSpec.Architecture, + Steps: data.FinalSpec.Steps, + Environment: maps.Clone(data.Spec.Environment), + IsFinal: true, + } + if finalSpec.Environment == nil { + finalSpec.Environment = make(map[string]string) + } + finalSpec.Environment["TANGLED_ARCHITECTURE"] = data.FinalSpec.Architecture + finalSpec.Environment["LOOM_FINAL"] = "true" + finalSpec.Environment["LOOM_ARTIFACTS"] = "/artifacts" + workflows = append(workflows, finalSpec) + } + + logger.Info("Expanded multi-arch workflow", "legs", len(data.MatrixLegs), "hasFinal", data.FinalSpec != nil) + } else { + singleSpec := data.Spec + if singleSpec.Environment == nil { + singleSpec.Environment = make(map[string]string) + } + singleSpec.Environment["LOOM_ARTIFACTS"] = "/artifacts" + workflows = []loomv1alpha1.WorkflowSpec{singleSpec} + } + pipelineRun := &loomv1alpha1.PipelineRunSpec{ PipelineID: wid.Rkey, SkipClone: skipClone, Secrets: repoSecrets, - Workflows: []loomv1alpha1.WorkflowSpec{data.Spec}, + Workflows: workflows, + MultiArch: data.IsMultiArch, } // Add clone commands if not skipping @@ -342,281 +525,243 @@ func (e *KubernetesEngine) DestroyWorkflow(ctx context.Context, wid models.Workf // Remove from tracking map delete(e.spindleSets, wid.String()) - // Close any open log streams for this workflow - e.closeLogStream(wid) + // Clean up artifacts for this pipeline + if e.artifacts != nil { + if err := e.artifacts.Cleanup(wid.PipelineId.AtUri().String()); err != nil { + logger.Error(err, "Failed to clean up artifacts") + } + } logger.Info("SpindleSet cleaned up successfully") return nil } -// RunStep streams logs for the specific step and waits for that step to complete. -// For Kubernetes engine, all steps run in a single Job, but we stream logs incrementally -// as each step executes. Each RunStep call blocks until that step's "end" control event is received. +// RunStep waits for step completion events from the runner via the gRPC hub. +// For single-arch workflows, blocks until the step's "end" control event. +// For multi-arch workflows: +// - idx 0: waits for ALL matrix leg runners to complete all their steps +// - idx 1: waits for the final runner to complete all its steps func (e *KubernetesEngine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.Workflow, idx int, wfSecrets []secrets.UnlockedSecret, wfLogger models.WorkflowLogger) error { logger := log.FromContext(ctx).WithValues("workflow", wid.Name, "pipeline", wid.Rkey, "step", idx) - // Query for the Job created by SpindleSetReconciler (only on first step) - var job *batchv1.Job - if idx == 0 { - spindleSet, err := e.getSpindleSet(ctx, wid) - if err != nil { - return err - } - if spindleSet == nil { - return fmt.Errorf("no SpindleSet found for workflow %s", wid.String()) - } - - // Wait for Job to be created by controller - deadline := time.Now().Add(5 * time.Minute) - for { - if time.Now().After(deadline) { - return fmt.Errorf("timeout waiting for Job to be created by controller") - } - - jobList := &batchv1.JobList{} - err := e.client.List(ctx, jobList, - client.InNamespace(e.namespace), - client.MatchingLabels{ - "loom.j5t.io/spindleset": spindleSet.Name, - "loom.j5t.io/workflow": w.Name, - "loom.j5t.io/pipeline-id": wid.Rkey, - }) - if err != nil { - return fmt.Errorf("failed to list jobs: %w", err) - } - - if len(jobList.Items) > 0 { - job = &jobList.Items[0] - break - } - - time.Sleep(2 * time.Second) - } - - logger.Info("Found Job for workflow", "jobName", job.Name) + data, ok := w.Data.(*kubernetesWorkflowData) + if !ok { + return fmt.Errorf("invalid workflow data type") } - // Get or create log stream (creates on first step, reuses on subsequent steps) - stream, err := e.getOrCreateLogStream(ctx, wid, job) - if err != nil { - return fmt.Errorf("failed to get log stream: %w", err) + if data.IsMultiArch { + return e.runMultiArchStep(ctx, wid, data, idx, wfLogger) } - // Read from stream until this step's end event - if wfLogger != nil { - logger.Info("Reading logs for step", "stepID", idx) - if err := e.readUntilStepEnd(ctx, stream, idx, w, wfLogger); err != nil { - logger.Error(err, "Failed to read step logs") - // Clean up stream on error - e.closeLogStream(wid) - return fmt.Errorf("failed to read logs for step %d: %w", idx, err) - } - logger.Info("Step completed", "stepID", idx) + // Single-arch: wait for one runner. + // PipelineID must match what the runner sends — the framework injects + // TANGLED_PIPELINE_ID as the full AT URI (see spindle/models.PipelineEnvVars), + // and the runner uses that value when registering with the hub. + key := loomgrpc.RunnerKey{ + PipelineID: wid.PipelineId.AtUri().String(), + WorkflowName: wid.Name, + Architecture: data.Spec.Architecture, } - // Clean up stream after last step - if idx == len(w.Steps)-1 { - logger.Info("Last step completed, closing log stream") - e.closeLogStream(wid) + if idx == 0 { + logger.Info("waiting for runner to connect", "key", key.String()) + select { + case <-e.hub.WaitForRunner(key): + logger.Info("runner connected", "key", key.String()) + case <-ctx.Done(): + return fmt.Errorf("context canceled while waiting for runner: %w", ctx.Err()) + case <-time.After(10 * time.Minute): + return fmt.Errorf("timeout waiting for runner to connect") + } } - return nil + return e.waitForRunnerStep(ctx, key, idx, idx, wfLogger) } -// getOrCreateLogStream gets an existing log stream or creates a new one for step 0 -func (e *KubernetesEngine) getOrCreateLogStream(ctx context.Context, wid models.WorkflowId, job *batchv1.Job) (*workflowLogStream, error) { - widKey := wid.String() +// runMultiArchStep handles RunStep for multi-arch workflows. +func (e *KubernetesEngine) runMultiArchStep(ctx context.Context, wid models.WorkflowId, data *kubernetesWorkflowData, idx int, wfLogger models.WorkflowLogger) error { + logger := log.FromContext(ctx).WithValues("workflow", wid.Name, "pipeline", wid.Rkey, "step", idx) - // Check if stream already exists - e.streamMutex.RLock() - stream, exists := e.logStreams[widKey] - e.streamMutex.RUnlock() + if idx == 0 { + // Matrix build phase: wait for all leg runners to complete all their steps + logger.Info("waiting for matrix leg runners", "legs", len(data.MatrixLegs)) - if exists { - return stream, nil - } + type legResult struct { + leg MatrixLeg + err error + } + results := make(chan legResult, len(data.MatrixLegs)) + + for legIdx, leg := range data.MatrixLegs { + go func() { + legName := fmt.Sprintf("%s-%s", data.Spec.Name, leg.Architecture) + key := loomgrpc.RunnerKey{ + PipelineID: wid.PipelineId.AtUri().String(), + WorkflowName: legName, + Architecture: leg.Architecture, + } - // Create new stream - logger := log.FromContext(ctx).WithValues("workflow", wid.Name, "pipeline", wid.Rkey) + // Wait for this leg's runner to connect + select { + case <-e.hub.WaitForRunner(key): + logger.Info("matrix leg runner connected", "key", key.String()) + case <-ctx.Done(): + results <- legResult{leg: leg, err: ctx.Err()} + return + case <-time.After(10 * time.Minute): + results <- legResult{leg: leg, err: fmt.Errorf("timeout waiting for runner %s", key.String())} + return + } - // Create kubernetes clientset for log streaming - clientset, err := kubernetes.NewForConfig(e.config) - if err != nil { - return nil, fmt.Errorf("failed to create kubernetes clientset: %w", err) - } + // Wait for all steps in this leg to complete, emitting per-leg + // control log lines so each architecture's output gets its own + // step section in the rendered log. + for stepIdx, userStep := range data.Spec.Steps { + logStepID := matrixLegLogStepID(legIdx, stepIdx) + sStep := syntheticStep{ + name: fmt.Sprintf("%s (%s)", userStep.Name, leg.Architecture), + command: userStep.Command, + kind: models.StepKindUser, + } + if wfLogger != nil { + _, _ = wfLogger.ControlWriter(logStepID, sStep, models.StepStatusStart).Write([]byte{0}) + } + err := e.waitForRunnerStep(ctx, key, stepIdx, logStepID, wfLogger) + if wfLogger != nil { + _, _ = wfLogger.ControlWriter(logStepID, sStep, models.StepStatusEnd).Write([]byte{0}) + } + if err != nil { + results <- legResult{leg: leg, err: fmt.Errorf("leg %s step %q: %w", leg.Architecture, userStep.Name, err)} + return + } + } - // Wait for pod to be created - var pod *corev1.Pod - deadline := time.Now().Add(2 * time.Minute) - for { - if time.Now().After(deadline) { - return nil, fmt.Errorf("timeout waiting for pod to be created") + results <- legResult{leg: leg, err: nil} + }() } - pods := &corev1.PodList{} - err := e.client.List(ctx, pods, client.InNamespace(job.Namespace), client.MatchingLabels(job.Spec.Template.Labels)) - if err != nil { - return nil, fmt.Errorf("failed to list pods: %w", err) + // Collect results from all legs + var errs []error + for range data.MatrixLegs { + result := <-results + if result.err != nil { + logger.Error(result.err, "matrix leg failed", "arch", result.leg.Architecture) + errs = append(errs, result.err) + } else { + logger.Info("matrix leg completed", "arch", result.leg.Architecture) + } } - if len(pods.Items) > 0 { - pod = &pods.Items[0] - break + if len(errs) > 0 { + return fmt.Errorf("matrix build failed: %v", errs[0]) } - time.Sleep(1 * time.Second) + logger.Info("all matrix legs completed successfully") + return nil } - logger.Info("Found pod for job", "podName", pod.Name) - - // Wait for pod to be running (or completed) - deadline = time.Now().Add(5 * time.Minute) - for { - if time.Now().After(deadline) { - return nil, fmt.Errorf("timeout waiting for pod to start") + if idx == 1 && data.FinalSpec != nil { + // Final phase: wait for the final runner + finalName := fmt.Sprintf("%s-final", data.Spec.Name) + key := loomgrpc.RunnerKey{ + PipelineID: wid.PipelineId.AtUri().String(), + WorkflowName: finalName, + Architecture: data.FinalSpec.Architecture, } - currentPod := &corev1.Pod{} - err := e.client.Get(ctx, client.ObjectKey{ - Namespace: pod.Namespace, - Name: pod.Name, - }, currentPod) - if err != nil { - return nil, fmt.Errorf("failed to get pod: %w", err) + logger.Info("waiting for final runner to connect", "key", key.String()) + select { + case <-e.hub.WaitForRunner(key): + logger.Info("final runner connected", "key", key.String()) + case <-ctx.Done(): + return fmt.Errorf("context canceled while waiting for final runner: %w", ctx.Err()) + case <-time.After(10 * time.Minute): + return fmt.Errorf("timeout waiting for final runner to connect") } - if currentPod.Status.Phase == corev1.PodRunning || currentPod.Status.Phase == corev1.PodSucceeded || currentPod.Status.Phase == corev1.PodFailed { - pod = currentPod - break + // Stream artifacts from matrix legs to the final runner + if e.artifacts != nil { + rs := e.hub.Get(key) + if rs != nil { + logger.Info("streaming artifacts to final runner", "pipeline", wid.Rkey) + if err := e.artifacts.StreamToRunner(wid.PipelineId.AtUri().String(), rs.SendToRunner); err != nil { + logger.Error(err, "failed to stream artifacts to final runner") + } + } } - time.Sleep(1 * time.Second) - } - - // Only use Follow mode for running pods. For completed pods, we need to read - // existing logs (Follow:true only streams NEW logs after connection). - shouldFollow := pod.Status.Phase == corev1.PodRunning - if !shouldFollow { - logger.Info("Pod already completed, reading existing logs", "podName", pod.Name, "phase", pod.Status.Phase) - } else { - logger.Info("Pod is running, streaming logs", "podName", pod.Name, "phase", pod.Status.Phase) - } - - // Stream logs from the main container (not init containers) - req := clientset.CoreV1().Pods(pod.Namespace).GetLogs(pod.Name, &corev1.PodLogOptions{ - Container: "runner", - Follow: shouldFollow, - }) - - logStream, err := req.Stream(ctx) - if err != nil { - return nil, fmt.Errorf("failed to open log stream: %w", err) - } - - // Create scanner - scanner := bufio.NewScanner(logStream) - buf := make([]byte, 0, 64*1024) - scanner.Buffer(buf, 1024*1024) + // Wait for all final steps, emitting per-step control log lines so the + // final phase's user steps are visible in the rendered log (the + // framework only emits a single "Final" control entry above us). + for stepIdx, finalStep := range data.FinalSpec.Steps { + logStepID := finalLogStepID(len(data.MatrixLegs), stepIdx) + sStep := syntheticStep{ + name: fmt.Sprintf("%s (final)", finalStep.Name), + command: finalStep.Command, + kind: models.StepKindUser, + } + if wfLogger != nil { + _, _ = wfLogger.ControlWriter(logStepID, sStep, models.StepStatusStart).Write([]byte{0}) + } + err := e.waitForRunnerStep(ctx, key, stepIdx, logStepID, wfLogger) + if wfLogger != nil { + _, _ = wfLogger.ControlWriter(logStepID, sStep, models.StepStatusEnd).Write([]byte{0}) + } + if err != nil { + return fmt.Errorf("final step %q: %w", finalStep.Name, err) + } + } - // Create and store stream - stream = &workflowLogStream{ - scanner: scanner, - stream: logStream, - pod: pod, - podPhase: pod.Status.Phase, + logger.Info("final steps completed") + return nil } - e.streamMutex.Lock() - e.logStreams[widKey] = stream - e.streamMutex.Unlock() - - return stream, nil + return fmt.Errorf("unexpected step index %d for multi-arch workflow", idx) } -// closeLogStream closes and removes a log stream -func (e *KubernetesEngine) closeLogStream(wid models.WorkflowId) { - widKey := wid.String() - - e.streamMutex.Lock() - defer e.streamMutex.Unlock() - - if stream, exists := e.logStreams[widKey]; exists { - stream.stream.Close() - delete(e.logStreams, widKey) +// waitForRunnerStep reads events from a runner's gRPC stream until a specific +// step completes. runnerStepID is matched against StepID fields on runner +// events to filter out other steps' events; logStepID is the step id passed to +// wfLogger.DataWriter when forwarding log content. For single-arch workflows +// these are the same; for matrix legs they differ so each leg's logs land in +// a distinct UI step. +func (e *KubernetesEngine) waitForRunnerStep(ctx context.Context, key loomgrpc.RunnerKey, runnerStepID, logStepID int, wfLogger models.WorkflowLogger) error { + rs := e.hub.Get(key) + if rs == nil { + return fmt.Errorf("runner not connected for %s", key.String()) } -} - -// readUntilStepEnd reads from the log stream until the end event for the specified step -func (e *KubernetesEngine) readUntilStepEnd(ctx context.Context, stream *workflowLogStream, stepID int, workflow *models.Workflow, wfLogger models.WorkflowLogger) error { - scanner := stream.scanner - - for scanner.Scan() { - line := scanner.Text() - - // Try to parse as extendedLogLine from the runner binary (includes exit_code) - var logLine extendedLogLine - if err := json.Unmarshal([]byte(line), &logLine); err != nil { - // Not JSON or parse error - skip - continue - } - - // Validate step index - if logLine.StepId < 0 || logLine.StepId >= len(workflow.Steps) { - continue - } - - // Only process events for the current step - if logLine.StepId != stepID { - // Got event for a different step - this shouldn't happen in sequential execution - // but log it and continue - continue - } - switch logLine.Kind { - case models.LogKindControl: - // Use control events from runner for flow control only - // Don't write them - the core spindle engine writes control events - if logLine.StepStatus == models.StepStatusEnd { - // Check exit code before returning success - if logLine.ExitCode != 0 { - return fmt.Errorf("step %d failed with exit code %d", stepID, logLine.ExitCode) + for { + select { + case evt := <-rs.Steps: + if evt.StepID != runnerStepID { + continue + } + if evt.Status == "end" { + if evt.ExitCode != 0 { + return fmt.Errorf("step %d failed with exit code %d", runnerStepID, evt.ExitCode) } return nil } - // For "start" events, just continue reading - case models.LogKindData: - // Log output from step - if logLine.Stream == "" { - logLine.Stream = "stdout" // Default to stdout + case logEvt := <-rs.Logs: + if logEvt.StepID != runnerStepID || wfLogger == nil { + continue } - dataWriter := wfLogger.DataWriter(logLine.StepId, logLine.Stream) - _, _ = dataWriter.Write([]byte(logLine.Content + "\n")) - } - } + stream := logEvt.Stream + if stream == "" { + stream = "stdout" + } + dataWriter := wfLogger.DataWriter(logStepID, stream) + _, _ = dataWriter.Write([]byte(logEvt.Content + "\n")) - if err := scanner.Err(); err != nil { - // EOF or context canceled is expected when pod terminates - if err != io.EOF && !strings.Contains(err.Error(), "context canceled") { - return fmt.Errorf("error reading logs: %w", err) - } - } + case <-rs.Done: + return fmt.Errorf("runner disconnected before step %d completed", runnerStepID) - // Scanner ended without seeing step end event. - // Re-check current pod status - it may have completed since we started streaming. - currentPod := &corev1.Pod{} - if err := e.client.Get(ctx, client.ObjectKey{Namespace: stream.pod.Namespace, Name: stream.pod.Name}, currentPod); err == nil { - if currentPod.Status.Phase == corev1.PodSucceeded { - // Pod succeeded - treat as success even without control event - return nil - } - if currentPod.Status.Phase == corev1.PodFailed { - return fmt.Errorf("pod failed before step %d completed", stepID) + case <-ctx.Done(): + return fmt.Errorf("context canceled during step %d: %w", runnerStepID, ctx.Err()) } } - - // Pod status unknown or still running but stream ended unexpectedly - return fmt.Errorf("log stream ended before step %d completed", stepID) } // Ensure KubernetesEngine implements the Engine interface diff --git a/internal/grpc/artifacts.go b/internal/grpc/artifacts.go new file mode 100644 index 0000000..d2fb7e2 --- /dev/null +++ b/internal/grpc/artifacts.go @@ -0,0 +1,187 @@ +package grpc + +import ( + "fmt" + "io" + "os" + "path/filepath" + "strings" + "sync" + + "sigs.k8s.io/controller-runtime/pkg/log" + + pb "tangled.org/evan.jarrett.net/loom/internal/pb/loom/v1" +) + +const artifactChunkSize = 32 * 1024 // 32KB chunks + +// ArtifactStore manages artifact files on disk, organized by pipeline/architecture. +type ArtifactStore struct { + baseDir string + mu sync.RWMutex + + // Track open file writers for streaming artifact chunks + writers map[string]*os.File + writersMu sync.Mutex +} + +// NewArtifactStore creates a new artifact store at the given base directory. +func NewArtifactStore(baseDir string) (*ArtifactStore, error) { + if err := os.MkdirAll(baseDir, 0755); err != nil { + return nil, fmt.Errorf("failed to create artifact store directory: %w", err) + } + return &ArtifactStore{ + baseDir: baseDir, + writers: make(map[string]*os.File), + }, nil +} + +// artifactDir returns the directory for a pipeline/architecture combination. +func (s *ArtifactStore) artifactDir(pipelineID, architecture string) string { + return filepath.Join(s.baseDir, pipelineID, architecture) +} + +// writerKey creates a unique key for tracking open file writers. +func writerKey(pipelineID, architecture, path string) string { + return fmt.Sprintf("%s/%s/%s", pipelineID, architecture, path) +} + +// WriteChunk writes an artifact chunk to disk. Creates the file on first chunk, +// appends on subsequent chunks, and closes on EOF. +func (s *ArtifactStore) WriteChunk(pipelineID, architecture string, chunk ArtifactEvent) error { + logger := log.Log.WithName("artifacts") + + dir := s.artifactDir(pipelineID, architecture) + fullPath := filepath.Join(dir, chunk.Path) + + // Security: ensure the path doesn't escape the artifact directory + cleanPath, err := filepath.Rel(dir, fullPath) + if err != nil || strings.HasPrefix(cleanPath, "..") { + return fmt.Errorf("invalid artifact path: %s", chunk.Path) + } + + key := writerKey(pipelineID, architecture, chunk.Path) + + s.writersMu.Lock() + defer s.writersMu.Unlock() + + f, exists := s.writers[key] + if !exists { + // Create directory structure and open file + if err := os.MkdirAll(filepath.Dir(fullPath), 0755); err != nil { + return fmt.Errorf("failed to create artifact directory: %w", err) + } + f, err = os.Create(fullPath) + if err != nil { + return fmt.Errorf("failed to create artifact file: %w", err) + } + s.writers[key] = f + logger.Info("receiving artifact", "pipeline", pipelineID, "arch", architecture, "path", chunk.Path) + } + + if len(chunk.Data) > 0 { + if _, err := f.Write(chunk.Data); err != nil { + return fmt.Errorf("failed to write artifact chunk: %w", err) + } + } + + if chunk.EOF { + f.Close() + delete(s.writers, key) + logger.Info("artifact received", "pipeline", pipelineID, "arch", architecture, "path", chunk.Path) + } + + return nil +} + +// StreamToRunner sends all artifacts for a pipeline to a runner's SendToRunner channel. +// Used by final jobs to receive artifacts from all matrix legs. +func (s *ArtifactStore) StreamToRunner(pipelineID string, sendCh chan<- *pb.ConnectResponse) error { + logger := log.Log.WithName("artifacts") + + pipelineDir := filepath.Join(s.baseDir, pipelineID) + if _, err := os.Stat(pipelineDir); os.IsNotExist(err) { + logger.Info("no artifacts found for pipeline", "pipeline", pipelineID) + return nil + } + + // Walk each architecture directory + archDirs, err := os.ReadDir(pipelineDir) + if err != nil { + return fmt.Errorf("failed to read pipeline artifacts: %w", err) + } + + for _, archDir := range archDirs { + if !archDir.IsDir() { + continue + } + arch := archDir.Name() + archPath := filepath.Join(pipelineDir, arch) + + err := filepath.Walk(archPath, func(path string, info os.FileInfo, err error) error { + if err != nil || info.IsDir() { + return err + } + + relPath, err := filepath.Rel(archPath, path) + if err != nil { + return err + } + + f, err := os.Open(path) + if err != nil { + return fmt.Errorf("failed to open artifact %s: %w", path, err) + } + defer f.Close() + + buf := make([]byte, artifactChunkSize) + for { + n, readErr := f.Read(buf) + isEOF := readErr == io.EOF + + sendCh <- &pb.ConnectResponse{ + Event: &pb.ConnectResponse_ArtifactData{ + ArtifactData: &pb.ArtifactData{ + SourceArchitecture: arch, + Path: relPath, + Data: buf[:n], + Eof: isEOF, + }, + }, + } + + if isEOF { + break + } + if readErr != nil { + return fmt.Errorf("failed to read artifact %s: %w", path, readErr) + } + } + + logger.Info("streamed artifact to runner", "pipeline", pipelineID, "arch", arch, "path", relPath) + return nil + }) + if err != nil { + return err + } + } + + // Send a sentinel message so the runner knows artifact streaming is complete. + // An empty Path with Eof=true signals "all artifacts sent". + sendCh <- &pb.ConnectResponse{ + Event: &pb.ConnectResponse_ArtifactData{ + ArtifactData: &pb.ArtifactData{ + Path: "", + Eof: true, + }, + }, + } + + return nil +} + +// Cleanup removes all artifacts for a pipeline. +func (s *ArtifactStore) Cleanup(pipelineID string) error { + dir := filepath.Join(s.baseDir, pipelineID) + return os.RemoveAll(dir) +} diff --git a/internal/grpc/hub.go b/internal/grpc/hub.go new file mode 100644 index 0000000..a4cf49f --- /dev/null +++ b/internal/grpc/hub.go @@ -0,0 +1,142 @@ +package grpc + +import ( + "fmt" + "sync" + + pb "tangled.org/evan.jarrett.net/loom/internal/pb/loom/v1" +) + +// RunnerKey uniquely identifies a runner connection. +type RunnerKey struct { + PipelineID string + WorkflowName string + Architecture string +} + +func (k RunnerKey) String() string { + return fmt.Sprintf("%s/%s/%s", k.PipelineID, k.WorkflowName, k.Architecture) +} + +// StepEvent represents a step lifecycle event received from a runner. +type StepEvent struct { + StepID int + Status string // "start" or "end" + ExitCode int +} + +// LogEvent represents a log line received from a runner. +type LogEvent struct { + StepID int + Stream string // "stdout" or "stderr" + Content string +} + +// ArtifactEvent represents an artifact chunk received from a runner. +type ArtifactEvent struct { + Path string + Data []byte + EOF bool +} + +// RunnerStream holds the channels for a single runner connection. +// The gRPC server writes to these channels; the engine reads from them. +type RunnerStream struct { + Steps chan StepEvent + Logs chan LogEvent + + // SendToRunner allows the engine to send messages back to the runner. + // The gRPC server reads from this channel and sends to the runner. + SendToRunner chan *pb.ConnectResponse + + // Done is closed when the runner disconnects. + Done chan struct{} +} + +func newRunnerStream() *RunnerStream { + return &RunnerStream{ + Steps: make(chan StepEvent, 64), + Logs: make(chan LogEvent, 256), + SendToRunner: make(chan *pb.ConnectResponse, 64), + Done: make(chan struct{}), + } +} + +// Hub manages active runner connections and provides channels for the engine +// to consume events from runners. +type Hub struct { + mu sync.RWMutex + streams map[string]*RunnerStream + + // waiters are channels that get notified when a runner with a given key connects. + waitersMu sync.Mutex + waiters map[string][]chan struct{} +} + +// NewHub creates a new Hub. +func NewHub() *Hub { + return &Hub{ + streams: make(map[string]*RunnerStream), + waiters: make(map[string][]chan struct{}), + } +} + +// Register creates a new RunnerStream for the given key. +// Called by the gRPC server when a runner connects. +func (h *Hub) Register(key RunnerKey) *RunnerStream { + h.mu.Lock() + stream := newRunnerStream() + h.streams[key.String()] = stream + h.mu.Unlock() + + // Notify any waiters + h.waitersMu.Lock() + if waiters, ok := h.waiters[key.String()]; ok { + for _, ch := range waiters { + close(ch) + } + delete(h.waiters, key.String()) + } + h.waitersMu.Unlock() + + return stream +} + +// Unregister removes a runner stream and closes its Done channel. +// Called by the gRPC server when a runner disconnects. +func (h *Hub) Unregister(key RunnerKey) { + h.mu.Lock() + defer h.mu.Unlock() + + if stream, ok := h.streams[key.String()]; ok { + close(stream.Done) + delete(h.streams, key.String()) + } +} + +// Get returns the RunnerStream for the given key, or nil if not connected. +func (h *Hub) Get(key RunnerKey) *RunnerStream { + h.mu.RLock() + defer h.mu.RUnlock() + return h.streams[key.String()] +} + +// WaitForRunner returns a channel that is closed when a runner with the given key connects. +// If the runner is already connected, returns a closed channel immediately. +func (h *Hub) WaitForRunner(key RunnerKey) <-chan struct{} { + h.mu.RLock() + if _, ok := h.streams[key.String()]; ok { + h.mu.RUnlock() + ch := make(chan struct{}) + close(ch) + return ch + } + h.mu.RUnlock() + + h.waitersMu.Lock() + defer h.waitersMu.Unlock() + + ch := make(chan struct{}) + h.waiters[key.String()] = append(h.waiters[key.String()], ch) + return ch +} diff --git a/internal/grpc/server.go b/internal/grpc/server.go new file mode 100644 index 0000000..0b24008 --- /dev/null +++ b/internal/grpc/server.go @@ -0,0 +1,136 @@ +package grpc + +import ( + "fmt" + "io" + "net" + + "google.golang.org/grpc" + "sigs.k8s.io/controller-runtime/pkg/log" + + pb "tangled.org/evan.jarrett.net/loom/internal/pb/loom/v1" +) + +// Server is the gRPC server that accepts runner connections. +type Server struct { + pb.UnimplementedLoomRunnerServiceServer + + hub *Hub + artifacts *ArtifactStore + grpcServer *grpc.Server +} + +// NewServer creates a new gRPC server with the given hub and artifact store. +func NewServer(hub *Hub, artifacts *ArtifactStore) *Server { + s := &Server{ + hub: hub, + artifacts: artifacts, + grpcServer: grpc.NewServer(), + } + pb.RegisterLoomRunnerServiceServer(s.grpcServer, s) + return s +} + +// Serve starts the gRPC server on the given listener. +func (s *Server) Serve(lis net.Listener) error { + return s.grpcServer.Serve(lis) +} + +// GracefulStop gracefully stops the gRPC server. +func (s *Server) GracefulStop() { + s.grpcServer.GracefulStop() +} + +// Connect handles a bidirectional stream from a runner. +func (s *Server) Connect(stream grpc.BidiStreamingServer[pb.ConnectRequest, pb.ConnectResponse]) error { + logger := log.Log.WithName("grpc") + + // First message must contain identity fields + msg, err := stream.Recv() + if err != nil { + return fmt.Errorf("failed to receive initial message: %w", err) + } + + key := RunnerKey{ + PipelineID: msg.PipelineId, + WorkflowName: msg.WorkflowName, + Architecture: msg.Architecture, + } + + if key.PipelineID == "" || key.WorkflowName == "" || key.Architecture == "" { + return fmt.Errorf("first message must include pipeline_id, workflow_name, and architecture") + } + + logger.Info("runner connected", "key", key.String()) + + // Register this runner + rs := s.hub.Register(key) + defer func() { + s.hub.Unregister(key) + logger.Info("runner disconnected", "key", key.String()) + }() + + // Process the first message's event (if any) + s.processEvent(key, rs, msg) + + // Start goroutine to send responses back to the runner + go func() { + for { + select { + case resp, ok := <-rs.SendToRunner: + if !ok { + return + } + if err := stream.Send(resp); err != nil { + logger.Error(err, "failed to send to runner", "key", key.String()) + return + } + case <-rs.Done: + return + } + } + }() + + // Read events from the runner + for { + msg, err := stream.Recv() + if err == io.EOF { + return nil + } + if err != nil { + return fmt.Errorf("recv error from %s: %w", key.String(), err) + } + s.processEvent(key, rs, msg) + } +} + +// processEvent routes a runner event to the appropriate channel. +func (s *Server) processEvent(key RunnerKey, rs *RunnerStream, msg *pb.ConnectRequest) { + logger := log.Log.WithName("grpc") + + switch evt := msg.Event.(type) { + case *pb.ConnectRequest_StepControl: + rs.Steps <- StepEvent{ + StepID: int(evt.StepControl.StepId), + Status: evt.StepControl.Status, + ExitCode: int(evt.StepControl.ExitCode), + } + case *pb.ConnectRequest_LogLine: + rs.Logs <- LogEvent{ + StepID: int(evt.LogLine.StepId), + Stream: evt.LogLine.Stream, + Content: evt.LogLine.Content, + } + case *pb.ConnectRequest_ArtifactChunk: + // Persist artifact to disk; final jobs stream it back via StreamToRunner. + if s.artifacts != nil { + if err := s.artifacts.WriteChunk(key.PipelineID, key.Architecture, ArtifactEvent{ + Path: evt.ArtifactChunk.Path, + Data: evt.ArtifactChunk.Data, + EOF: evt.ArtifactChunk.Eof, + }); err != nil { + logger.Error(err, "failed to write artifact chunk", "key", key.String()) + } + } + } +} diff --git a/internal/jobbuilder/job_template.go b/internal/jobbuilder/job_template.go index 43c8bf1..11fe5dc 100644 --- a/internal/jobbuilder/job_template.go +++ b/internal/jobbuilder/job_template.go @@ -67,6 +67,10 @@ type WorkflowConfig struct { // Namespace is the Kubernetes namespace for the Job Namespace string + + // OperatorAddr is the gRPC address of the Loom operator for runner communication. + // Injected as LOOM_OPERATOR_ADDR env var into the runner container. + OperatorAddr string } // nodeMatchesSelector returns true if at least one node has all the labels in selector. @@ -256,6 +260,10 @@ func BuildJob(config WorkflowConfig, nodes *corev1.NodeList) (*batchv1.Job, erro Name: "LOOM_SECRET_KEYS", Value: strings.Join(config.SecretKeys, ","), }, + corev1.EnvVar{ + Name: "LOOM_OPERATOR_ADDR", + Value: config.OperatorAddr, + }, ), // Inject repository secrets via envFrom if available @@ -510,6 +518,10 @@ func buildRunnerVolumeMounts(config WorkflowConfig) []corev1.VolumeMount { MountPath: "/home/runner", SubPath: "runner", }, + { + Name: "artifacts", + MountPath: "/artifacts", + }, } // Mount registry credentials if specified @@ -563,6 +575,12 @@ func buildVolumes(config WorkflowConfig) []corev1.Volume { EmptyDir: &corev1.EmptyDirVolumeSource{}, }, }, + { + Name: "artifacts", + VolumeSource: corev1.VolumeSource{ + EmptyDir: &corev1.EmptyDirVolumeSource{}, + }, + }, } // Add registry credentials volume if specified diff --git a/internal/pb/loom/v1/loom.pb.go b/internal/pb/loom/v1/loom.pb.go new file mode 100644 index 0000000..3c731d7 --- /dev/null +++ b/internal/pb/loom/v1/loom.pb.go @@ -0,0 +1,635 @@ +// Code generated by protoc-gen-go. DO NOT EDIT. +// versions: +// protoc-gen-go v1.36.8 +// protoc (unknown) +// source: loom/v1/loom.proto + +package loomv1 + +import ( + protoreflect "google.golang.org/protobuf/reflect/protoreflect" + protoimpl "google.golang.org/protobuf/runtime/protoimpl" + reflect "reflect" + sync "sync" + unsafe "unsafe" +) + +const ( + // Verify that this generated code is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(20 - protoimpl.MinVersion) + // Verify that runtime/protoimpl is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) +) + +// ConnectRequest is sent from the runner to the operator. +type ConnectRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Identity fields — must be set on the first message, optional on subsequent. + PipelineId string `protobuf:"bytes,1,opt,name=pipeline_id,json=pipelineId,proto3" json:"pipeline_id,omitempty"` + WorkflowName string `protobuf:"bytes,2,opt,name=workflow_name,json=workflowName,proto3" json:"workflow_name,omitempty"` + Architecture string `protobuf:"bytes,3,opt,name=architecture,proto3" json:"architecture,omitempty"` + // Types that are valid to be assigned to Event: + // + // *ConnectRequest_StepControl + // *ConnectRequest_LogLine + // *ConnectRequest_ArtifactChunk + Event isConnectRequest_Event `protobuf_oneof:"event"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ConnectRequest) Reset() { + *x = ConnectRequest{} + mi := &file_loom_v1_loom_proto_msgTypes[0] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ConnectRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ConnectRequest) ProtoMessage() {} + +func (x *ConnectRequest) ProtoReflect() protoreflect.Message { + mi := &file_loom_v1_loom_proto_msgTypes[0] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ConnectRequest.ProtoReflect.Descriptor instead. +func (*ConnectRequest) Descriptor() ([]byte, []int) { + return file_loom_v1_loom_proto_rawDescGZIP(), []int{0} +} + +func (x *ConnectRequest) GetPipelineId() string { + if x != nil { + return x.PipelineId + } + return "" +} + +func (x *ConnectRequest) GetWorkflowName() string { + if x != nil { + return x.WorkflowName + } + return "" +} + +func (x *ConnectRequest) GetArchitecture() string { + if x != nil { + return x.Architecture + } + return "" +} + +func (x *ConnectRequest) GetEvent() isConnectRequest_Event { + if x != nil { + return x.Event + } + return nil +} + +func (x *ConnectRequest) GetStepControl() *StepControl { + if x != nil { + if x, ok := x.Event.(*ConnectRequest_StepControl); ok { + return x.StepControl + } + } + return nil +} + +func (x *ConnectRequest) GetLogLine() *LogLine { + if x != nil { + if x, ok := x.Event.(*ConnectRequest_LogLine); ok { + return x.LogLine + } + } + return nil +} + +func (x *ConnectRequest) GetArtifactChunk() *ArtifactChunk { + if x != nil { + if x, ok := x.Event.(*ConnectRequest_ArtifactChunk); ok { + return x.ArtifactChunk + } + } + return nil +} + +type isConnectRequest_Event interface { + isConnectRequest_Event() +} + +type ConnectRequest_StepControl struct { + StepControl *StepControl `protobuf:"bytes,4,opt,name=step_control,json=stepControl,proto3,oneof"` +} + +type ConnectRequest_LogLine struct { + LogLine *LogLine `protobuf:"bytes,5,opt,name=log_line,json=logLine,proto3,oneof"` +} + +type ConnectRequest_ArtifactChunk struct { + ArtifactChunk *ArtifactChunk `protobuf:"bytes,6,opt,name=artifact_chunk,json=artifactChunk,proto3,oneof"` +} + +func (*ConnectRequest_StepControl) isConnectRequest_Event() {} + +func (*ConnectRequest_LogLine) isConnectRequest_Event() {} + +func (*ConnectRequest_ArtifactChunk) isConnectRequest_Event() {} + +// ConnectResponse is sent from the operator to the runner. +type ConnectResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Types that are valid to be assigned to Event: + // + // *ConnectResponse_Ack + // *ConnectResponse_ArtifactData + Event isConnectResponse_Event `protobuf_oneof:"event"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ConnectResponse) Reset() { + *x = ConnectResponse{} + mi := &file_loom_v1_loom_proto_msgTypes[1] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ConnectResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ConnectResponse) ProtoMessage() {} + +func (x *ConnectResponse) ProtoReflect() protoreflect.Message { + mi := &file_loom_v1_loom_proto_msgTypes[1] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ConnectResponse.ProtoReflect.Descriptor instead. +func (*ConnectResponse) Descriptor() ([]byte, []int) { + return file_loom_v1_loom_proto_rawDescGZIP(), []int{1} +} + +func (x *ConnectResponse) GetEvent() isConnectResponse_Event { + if x != nil { + return x.Event + } + return nil +} + +func (x *ConnectResponse) GetAck() *Ack { + if x != nil { + if x, ok := x.Event.(*ConnectResponse_Ack); ok { + return x.Ack + } + } + return nil +} + +func (x *ConnectResponse) GetArtifactData() *ArtifactData { + if x != nil { + if x, ok := x.Event.(*ConnectResponse_ArtifactData); ok { + return x.ArtifactData + } + } + return nil +} + +type isConnectResponse_Event interface { + isConnectResponse_Event() +} + +type ConnectResponse_Ack struct { + Ack *Ack `protobuf:"bytes,1,opt,name=ack,proto3,oneof"` +} + +type ConnectResponse_ArtifactData struct { + ArtifactData *ArtifactData `protobuf:"bytes,2,opt,name=artifact_data,json=artifactData,proto3,oneof"` +} + +func (*ConnectResponse_Ack) isConnectResponse_Event() {} + +func (*ConnectResponse_ArtifactData) isConnectResponse_Event() {} + +// StepControl signals step lifecycle transitions. +type StepControl struct { + state protoimpl.MessageState `protogen:"open.v1"` + StepId int32 `protobuf:"varint,1,opt,name=step_id,json=stepId,proto3" json:"step_id,omitempty"` + // "start" or "end" + Status string `protobuf:"bytes,2,opt,name=status,proto3" json:"status,omitempty"` + // Exit code of the step command (only meaningful when status = "end"). + ExitCode int32 `protobuf:"varint,3,opt,name=exit_code,json=exitCode,proto3" json:"exit_code,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *StepControl) Reset() { + *x = StepControl{} + mi := &file_loom_v1_loom_proto_msgTypes[2] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *StepControl) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*StepControl) ProtoMessage() {} + +func (x *StepControl) ProtoReflect() protoreflect.Message { + mi := &file_loom_v1_loom_proto_msgTypes[2] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use StepControl.ProtoReflect.Descriptor instead. +func (*StepControl) Descriptor() ([]byte, []int) { + return file_loom_v1_loom_proto_rawDescGZIP(), []int{2} +} + +func (x *StepControl) GetStepId() int32 { + if x != nil { + return x.StepId + } + return 0 +} + +func (x *StepControl) GetStatus() string { + if x != nil { + return x.Status + } + return "" +} + +func (x *StepControl) GetExitCode() int32 { + if x != nil { + return x.ExitCode + } + return 0 +} + +// LogLine carries a single line of step output. +type LogLine struct { + state protoimpl.MessageState `protogen:"open.v1"` + StepId int32 `protobuf:"varint,1,opt,name=step_id,json=stepId,proto3" json:"step_id,omitempty"` + // "stdout" or "stderr" + Stream string `protobuf:"bytes,2,opt,name=stream,proto3" json:"stream,omitempty"` + Content string `protobuf:"bytes,3,opt,name=content,proto3" json:"content,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *LogLine) Reset() { + *x = LogLine{} + mi := &file_loom_v1_loom_proto_msgTypes[3] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *LogLine) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*LogLine) ProtoMessage() {} + +func (x *LogLine) ProtoReflect() protoreflect.Message { + mi := &file_loom_v1_loom_proto_msgTypes[3] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use LogLine.ProtoReflect.Descriptor instead. +func (*LogLine) Descriptor() ([]byte, []int) { + return file_loom_v1_loom_proto_rawDescGZIP(), []int{3} +} + +func (x *LogLine) GetStepId() int32 { + if x != nil { + return x.StepId + } + return 0 +} + +func (x *LogLine) GetStream() string { + if x != nil { + return x.Stream + } + return "" +} + +func (x *LogLine) GetContent() string { + if x != nil { + return x.Content + } + return "" +} + +// ArtifactChunk streams a file from runner to operator in chunks. +// The runner sends one or more chunks per file, with eof=true on the last chunk. +type ArtifactChunk struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Relative path within the artifacts directory (e.g., "bin/myapp"). + Path string `protobuf:"bytes,1,opt,name=path,proto3" json:"path,omitempty"` + Data []byte `protobuf:"bytes,2,opt,name=data,proto3" json:"data,omitempty"` + Eof bool `protobuf:"varint,3,opt,name=eof,proto3" json:"eof,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ArtifactChunk) Reset() { + *x = ArtifactChunk{} + mi := &file_loom_v1_loom_proto_msgTypes[4] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ArtifactChunk) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ArtifactChunk) ProtoMessage() {} + +func (x *ArtifactChunk) ProtoReflect() protoreflect.Message { + mi := &file_loom_v1_loom_proto_msgTypes[4] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ArtifactChunk.ProtoReflect.Descriptor instead. +func (*ArtifactChunk) Descriptor() ([]byte, []int) { + return file_loom_v1_loom_proto_rawDescGZIP(), []int{4} +} + +func (x *ArtifactChunk) GetPath() string { + if x != nil { + return x.Path + } + return "" +} + +func (x *ArtifactChunk) GetData() []byte { + if x != nil { + return x.Data + } + return nil +} + +func (x *ArtifactChunk) GetEof() bool { + if x != nil { + return x.Eof + } + return false +} + +// ArtifactData streams artifact files from operator to runner (for final jobs). +// The operator sends collected artifacts from matrix legs to the final runner. +type ArtifactData struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Architecture of the source matrix leg (e.g., "amd64"). + SourceArchitecture string `protobuf:"bytes,1,opt,name=source_architecture,json=sourceArchitecture,proto3" json:"source_architecture,omitempty"` + // Relative path within the artifacts directory. + Path string `protobuf:"bytes,2,opt,name=path,proto3" json:"path,omitempty"` + Data []byte `protobuf:"bytes,3,opt,name=data,proto3" json:"data,omitempty"` + Eof bool `protobuf:"varint,4,opt,name=eof,proto3" json:"eof,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ArtifactData) Reset() { + *x = ArtifactData{} + mi := &file_loom_v1_loom_proto_msgTypes[5] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ArtifactData) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ArtifactData) ProtoMessage() {} + +func (x *ArtifactData) ProtoReflect() protoreflect.Message { + mi := &file_loom_v1_loom_proto_msgTypes[5] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ArtifactData.ProtoReflect.Descriptor instead. +func (*ArtifactData) Descriptor() ([]byte, []int) { + return file_loom_v1_loom_proto_rawDescGZIP(), []int{5} +} + +func (x *ArtifactData) GetSourceArchitecture() string { + if x != nil { + return x.SourceArchitecture + } + return "" +} + +func (x *ArtifactData) GetPath() string { + if x != nil { + return x.Path + } + return "" +} + +func (x *ArtifactData) GetData() []byte { + if x != nil { + return x.Data + } + return nil +} + +func (x *ArtifactData) GetEof() bool { + if x != nil { + return x.Eof + } + return false +} + +// Ack acknowledges receipt of a runner event. +type Ack struct { + state protoimpl.MessageState `protogen:"open.v1"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *Ack) Reset() { + *x = Ack{} + mi := &file_loom_v1_loom_proto_msgTypes[6] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *Ack) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*Ack) ProtoMessage() {} + +func (x *Ack) ProtoReflect() protoreflect.Message { + mi := &file_loom_v1_loom_proto_msgTypes[6] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use Ack.ProtoReflect.Descriptor instead. +func (*Ack) Descriptor() ([]byte, []int) { + return file_loom_v1_loom_proto_rawDescGZIP(), []int{6} +} + +var File_loom_v1_loom_proto protoreflect.FileDescriptor + +const file_loom_v1_loom_proto_rawDesc = "" + + "\n" + + "\x12loom/v1/loom.proto\x12\aloom.v1\"\xae\x02\n" + + "\x0eConnectRequest\x12\x1f\n" + + "\vpipeline_id\x18\x01 \x01(\tR\n" + + "pipelineId\x12#\n" + + "\rworkflow_name\x18\x02 \x01(\tR\fworkflowName\x12\"\n" + + "\farchitecture\x18\x03 \x01(\tR\farchitecture\x129\n" + + "\fstep_control\x18\x04 \x01(\v2\x14.loom.v1.StepControlH\x00R\vstepControl\x12-\n" + + "\blog_line\x18\x05 \x01(\v2\x10.loom.v1.LogLineH\x00R\alogLine\x12?\n" + + "\x0eartifact_chunk\x18\x06 \x01(\v2\x16.loom.v1.ArtifactChunkH\x00R\rartifactChunkB\a\n" + + "\x05event\"z\n" + + "\x0fConnectResponse\x12 \n" + + "\x03ack\x18\x01 \x01(\v2\f.loom.v1.AckH\x00R\x03ack\x12<\n" + + "\rartifact_data\x18\x02 \x01(\v2\x15.loom.v1.ArtifactDataH\x00R\fartifactDataB\a\n" + + "\x05event\"[\n" + + "\vStepControl\x12\x17\n" + + "\astep_id\x18\x01 \x01(\x05R\x06stepId\x12\x16\n" + + "\x06status\x18\x02 \x01(\tR\x06status\x12\x1b\n" + + "\texit_code\x18\x03 \x01(\x05R\bexitCode\"T\n" + + "\aLogLine\x12\x17\n" + + "\astep_id\x18\x01 \x01(\x05R\x06stepId\x12\x16\n" + + "\x06stream\x18\x02 \x01(\tR\x06stream\x12\x18\n" + + "\acontent\x18\x03 \x01(\tR\acontent\"I\n" + + "\rArtifactChunk\x12\x12\n" + + "\x04path\x18\x01 \x01(\tR\x04path\x12\x12\n" + + "\x04data\x18\x02 \x01(\fR\x04data\x12\x10\n" + + "\x03eof\x18\x03 \x01(\bR\x03eof\"y\n" + + "\fArtifactData\x12/\n" + + "\x13source_architecture\x18\x01 \x01(\tR\x12sourceArchitecture\x12\x12\n" + + "\x04path\x18\x02 \x01(\tR\x04path\x12\x12\n" + + "\x04data\x18\x03 \x01(\fR\x04data\x12\x10\n" + + "\x03eof\x18\x04 \x01(\bR\x03eof\"\x05\n" + + "\x03Ack2U\n" + + "\x11LoomRunnerService\x12@\n" + + "\aConnect\x12\x17.loom.v1.ConnectRequest\x1a\x18.loom.v1.ConnectResponse(\x010\x01B6Z4tangled.org/evan.jarrett.net/loom/internal/pb/loomv1b\x06proto3" + +var ( + file_loom_v1_loom_proto_rawDescOnce sync.Once + file_loom_v1_loom_proto_rawDescData []byte +) + +func file_loom_v1_loom_proto_rawDescGZIP() []byte { + file_loom_v1_loom_proto_rawDescOnce.Do(func() { + file_loom_v1_loom_proto_rawDescData = protoimpl.X.CompressGZIP(unsafe.Slice(unsafe.StringData(file_loom_v1_loom_proto_rawDesc), len(file_loom_v1_loom_proto_rawDesc))) + }) + return file_loom_v1_loom_proto_rawDescData +} + +var file_loom_v1_loom_proto_msgTypes = make([]protoimpl.MessageInfo, 7) +var file_loom_v1_loom_proto_goTypes = []any{ + (*ConnectRequest)(nil), // 0: loom.v1.ConnectRequest + (*ConnectResponse)(nil), // 1: loom.v1.ConnectResponse + (*StepControl)(nil), // 2: loom.v1.StepControl + (*LogLine)(nil), // 3: loom.v1.LogLine + (*ArtifactChunk)(nil), // 4: loom.v1.ArtifactChunk + (*ArtifactData)(nil), // 5: loom.v1.ArtifactData + (*Ack)(nil), // 6: loom.v1.Ack +} +var file_loom_v1_loom_proto_depIdxs = []int32{ + 2, // 0: loom.v1.ConnectRequest.step_control:type_name -> loom.v1.StepControl + 3, // 1: loom.v1.ConnectRequest.log_line:type_name -> loom.v1.LogLine + 4, // 2: loom.v1.ConnectRequest.artifact_chunk:type_name -> loom.v1.ArtifactChunk + 6, // 3: loom.v1.ConnectResponse.ack:type_name -> loom.v1.Ack + 5, // 4: loom.v1.ConnectResponse.artifact_data:type_name -> loom.v1.ArtifactData + 0, // 5: loom.v1.LoomRunnerService.Connect:input_type -> loom.v1.ConnectRequest + 1, // 6: loom.v1.LoomRunnerService.Connect:output_type -> loom.v1.ConnectResponse + 6, // [6:7] is the sub-list for method output_type + 5, // [5:6] is the sub-list for method input_type + 5, // [5:5] is the sub-list for extension type_name + 5, // [5:5] is the sub-list for extension extendee + 0, // [0:5] is the sub-list for field type_name +} + +func init() { file_loom_v1_loom_proto_init() } +func file_loom_v1_loom_proto_init() { + if File_loom_v1_loom_proto != nil { + return + } + file_loom_v1_loom_proto_msgTypes[0].OneofWrappers = []any{ + (*ConnectRequest_StepControl)(nil), + (*ConnectRequest_LogLine)(nil), + (*ConnectRequest_ArtifactChunk)(nil), + } + file_loom_v1_loom_proto_msgTypes[1].OneofWrappers = []any{ + (*ConnectResponse_Ack)(nil), + (*ConnectResponse_ArtifactData)(nil), + } + type x struct{} + out := protoimpl.TypeBuilder{ + File: protoimpl.DescBuilder{ + GoPackagePath: reflect.TypeOf(x{}).PkgPath(), + RawDescriptor: unsafe.Slice(unsafe.StringData(file_loom_v1_loom_proto_rawDesc), len(file_loom_v1_loom_proto_rawDesc)), + NumEnums: 0, + NumMessages: 7, + NumExtensions: 0, + NumServices: 1, + }, + GoTypes: file_loom_v1_loom_proto_goTypes, + DependencyIndexes: file_loom_v1_loom_proto_depIdxs, + MessageInfos: file_loom_v1_loom_proto_msgTypes, + }.Build() + File_loom_v1_loom_proto = out.File + file_loom_v1_loom_proto_goTypes = nil + file_loom_v1_loom_proto_depIdxs = nil +} diff --git a/internal/pb/loom/v1/loom_grpc.pb.go b/internal/pb/loom/v1/loom_grpc.pb.go new file mode 100644 index 0000000..04a7bbf --- /dev/null +++ b/internal/pb/loom/v1/loom_grpc.pb.go @@ -0,0 +1,129 @@ +// Code generated by protoc-gen-go-grpc. DO NOT EDIT. +// versions: +// - protoc-gen-go-grpc v1.5.1 +// - protoc (unknown) +// source: loom/v1/loom.proto + +package loomv1 + +import ( + context "context" + grpc "google.golang.org/grpc" + codes "google.golang.org/grpc/codes" + status "google.golang.org/grpc/status" +) + +// This is a compile-time assertion to ensure that this generated file +// is compatible with the grpc package it is being compiled against. +// Requires gRPC-Go v1.64.0 or later. +const _ = grpc.SupportPackageIsVersion9 + +const ( + LoomRunnerService_Connect_FullMethodName = "/loom.v1.LoomRunnerService/Connect" +) + +// LoomRunnerServiceClient is the client API for LoomRunnerService service. +// +// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream. +// +// LoomRunnerService is the bidirectional communication channel between runner pods +// and the Loom operator. Runners connect on startup and stream all events +// (step control, log output, artifacts) over this single connection. +type LoomRunnerServiceClient interface { + // Connect establishes a bidirectional stream between a runner and the operator. + // The runner identifies itself in the first message and then streams events. + // The operator sends acknowledgements and commands back. + Connect(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[ConnectRequest, ConnectResponse], error) +} + +type loomRunnerServiceClient struct { + cc grpc.ClientConnInterface +} + +func NewLoomRunnerServiceClient(cc grpc.ClientConnInterface) LoomRunnerServiceClient { + return &loomRunnerServiceClient{cc} +} + +func (c *loomRunnerServiceClient) Connect(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[ConnectRequest, ConnectResponse], error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + stream, err := c.cc.NewStream(ctx, &LoomRunnerService_ServiceDesc.Streams[0], LoomRunnerService_Connect_FullMethodName, cOpts...) + if err != nil { + return nil, err + } + x := &grpc.GenericClientStream[ConnectRequest, ConnectResponse]{ClientStream: stream} + return x, nil +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type LoomRunnerService_ConnectClient = grpc.BidiStreamingClient[ConnectRequest, ConnectResponse] + +// LoomRunnerServiceServer is the server API for LoomRunnerService service. +// All implementations must embed UnimplementedLoomRunnerServiceServer +// for forward compatibility. +// +// LoomRunnerService is the bidirectional communication channel between runner pods +// and the Loom operator. Runners connect on startup and stream all events +// (step control, log output, artifacts) over this single connection. +type LoomRunnerServiceServer interface { + // Connect establishes a bidirectional stream between a runner and the operator. + // The runner identifies itself in the first message and then streams events. + // The operator sends acknowledgements and commands back. + Connect(grpc.BidiStreamingServer[ConnectRequest, ConnectResponse]) error + mustEmbedUnimplementedLoomRunnerServiceServer() +} + +// UnimplementedLoomRunnerServiceServer must be embedded to have +// forward compatible implementations. +// +// NOTE: this should be embedded by value instead of pointer to avoid a nil +// pointer dereference when methods are called. +type UnimplementedLoomRunnerServiceServer struct{} + +func (UnimplementedLoomRunnerServiceServer) Connect(grpc.BidiStreamingServer[ConnectRequest, ConnectResponse]) error { + return status.Errorf(codes.Unimplemented, "method Connect not implemented") +} +func (UnimplementedLoomRunnerServiceServer) mustEmbedUnimplementedLoomRunnerServiceServer() {} +func (UnimplementedLoomRunnerServiceServer) testEmbeddedByValue() {} + +// UnsafeLoomRunnerServiceServer may be embedded to opt out of forward compatibility for this service. +// Use of this interface is not recommended, as added methods to LoomRunnerServiceServer will +// result in compilation errors. +type UnsafeLoomRunnerServiceServer interface { + mustEmbedUnimplementedLoomRunnerServiceServer() +} + +func RegisterLoomRunnerServiceServer(s grpc.ServiceRegistrar, srv LoomRunnerServiceServer) { + // If the following call pancis, it indicates UnimplementedLoomRunnerServiceServer was + // embedded by pointer and is nil. This will cause panics if an + // unimplemented method is ever invoked, so we test this at initialization + // time to prevent it from happening at runtime later due to I/O. + if t, ok := srv.(interface{ testEmbeddedByValue() }); ok { + t.testEmbeddedByValue() + } + s.RegisterService(&LoomRunnerService_ServiceDesc, srv) +} + +func _LoomRunnerService_Connect_Handler(srv interface{}, stream grpc.ServerStream) error { + return srv.(LoomRunnerServiceServer).Connect(&grpc.GenericServerStream[ConnectRequest, ConnectResponse]{ServerStream: stream}) +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type LoomRunnerService_ConnectServer = grpc.BidiStreamingServer[ConnectRequest, ConnectResponse] + +// LoomRunnerService_ServiceDesc is the grpc.ServiceDesc for LoomRunnerService service. +// It's only intended for direct use with grpc.RegisterService, +// and not to be introspected or modified (even as a copy) +var LoomRunnerService_ServiceDesc = grpc.ServiceDesc{ + ServiceName: "loom.v1.LoomRunnerService", + HandlerType: (*LoomRunnerServiceServer)(nil), + Methods: []grpc.MethodDesc{}, + Streams: []grpc.StreamDesc{ + { + StreamName: "Connect", + Handler: _LoomRunnerService_Connect_Handler, + ServerStreams: true, + ClientStreams: true, + }, + }, + Metadata: "loom/v1/loom.proto", +} diff --git a/proto/loom/v1/loom.proto b/proto/loom/v1/loom.proto new file mode 100644 index 0000000..d473e05 --- /dev/null +++ b/proto/loom/v1/loom.proto @@ -0,0 +1,77 @@ +syntax = "proto3"; + +package loom.v1; + +option go_package = "tangled.org/evan.jarrett.net/loom/internal/pb/loomv1"; + +// LoomRunnerService is the bidirectional communication channel between runner pods +// and the Loom operator. Runners connect on startup and stream all events +// (step control, log output, artifacts) over this single connection. +service LoomRunnerService { + // Connect establishes a bidirectional stream between a runner and the operator. + // The runner identifies itself in the first message and then streams events. + // The operator sends acknowledgements and commands back. + rpc Connect(stream ConnectRequest) returns (stream ConnectResponse); +} + +// ConnectRequest is sent from the runner to the operator. +message ConnectRequest { + // Identity fields — must be set on the first message, optional on subsequent. + string pipeline_id = 1; + string workflow_name = 2; + string architecture = 3; + + oneof event { + StepControl step_control = 4; + LogLine log_line = 5; + ArtifactChunk artifact_chunk = 6; + } +} + +// ConnectResponse is sent from the operator to the runner. +message ConnectResponse { + oneof event { + Ack ack = 1; + ArtifactData artifact_data = 2; + } +} + +// StepControl signals step lifecycle transitions. +message StepControl { + int32 step_id = 1; + // "start" or "end" + string status = 2; + // Exit code of the step command (only meaningful when status = "end"). + int32 exit_code = 3; +} + +// LogLine carries a single line of step output. +message LogLine { + int32 step_id = 1; + // "stdout" or "stderr" + string stream = 2; + string content = 3; +} + +// ArtifactChunk streams a file from runner to operator in chunks. +// The runner sends one or more chunks per file, with eof=true on the last chunk. +message ArtifactChunk { + // Relative path within the artifacts directory (e.g., "bin/myapp"). + string path = 1; + bytes data = 2; + bool eof = 3; +} + +// ArtifactData streams artifact files from operator to runner (for final jobs). +// The operator sends collected artifacts from matrix legs to the final runner. +message ArtifactData { + // Architecture of the source matrix leg (e.g., "amd64"). + string source_architecture = 1; + // Relative path within the artifacts directory. + string path = 2; + bytes data = 3; + bool eof = 4; +} + +// Ack acknowledges receipt of a runner event. +message Ack {} -- 2.51.2