diff --git a/cmd/kubernetes/kubernetes.go b/cmd/kubernetes/kubernetes.go new file mode 100644 index 000000000..9f7e82811 --- /dev/null +++ b/cmd/kubernetes/kubernetes.go @@ -0,0 +1,114 @@ +// Package kubernetes contains CLI commands for managing Funnel's Kubernetes resources. +package kubernetes + +import ( + "context" + "fmt" + "strings" + "syscall" + "time" + + cmdutil "github.com/ohsu-comp-bio/funnel/cmd/util" + k8sbackend "github.com/ohsu-comp-bio/funnel/compute/kubernetes" + "github.com/ohsu-comp-bio/funnel/config" + "github.com/ohsu-comp-bio/funnel/database/badger" + "github.com/ohsu-comp-bio/funnel/database/boltdb" + "github.com/ohsu-comp-bio/funnel/database/datastore" + "github.com/ohsu-comp-bio/funnel/database/dynamodb" + "github.com/ohsu-comp-bio/funnel/database/elastic" + "github.com/ohsu-comp-bio/funnel/database/mongodb" + "github.com/ohsu-comp-bio/funnel/database/postgres" + "github.com/ohsu-comp-bio/funnel/events" + "github.com/ohsu-comp-bio/funnel/logger" + "github.com/ohsu-comp-bio/funnel/tes" + "github.com/ohsu-comp-bio/funnel/util" + "github.com/spf13/cobra" +) + +// Cmd is the root "funnel kubernetes" command. +var Cmd = &cobra.Command{ + Use: "kubernetes", + Short: "Funnel Kubernetes management commands.", +} + +func init() { + Cmd.AddCommand(cleanupCmd()) +} + +func cleanupCmd() *cobra.Command { + var ( + configFile string + flagConf = config.EmptyConfig() + ) + + cmd := &cobra.Command{ + Use: "cleanup", + Short: "Delete orphaned Funnel-managed Kubernetes resources with no matching task.", + Long: `Scans Funnel-labeled Kubernetes resources (PVs, PVCs, ConfigMaps, ServiceAccounts, +Roles, RoleBindings) and deletes any whose task ID is no longer present or active +in the Funnel database. Intended to be run as a Kubernetes CronJob so that cleanup +is decoupled from the server lifecycle and multiple replicas do not race.`, + Args: cobra.NoArgs, + RunE: func(cmd *cobra.Command, args []string) error { + conf, err := cmdutil.MergeConfigFileWithFlags(configFile, flagConf) + if err != nil { + return fmt.Errorf("error processing config: %v", err) + } + + log := logger.NewLogger("kubernetes-cleanup", conf.Logger) + + ctx, cancel := context.WithCancel(context.Background()) + ctx = util.SignalContext(ctx, time.Millisecond*500, syscall.SIGINT, syscall.SIGTERM) + defer cancel() + + // Open only the database — no HTTP/gRPC server needed. + reader, err := openReader(ctx, conf) + if err != nil { + return fmt.Errorf("opening database: %v", err) + } + + // Build the K8s backend (connects to the cluster via in-cluster config). + // We pass a no-op event writer since this command only deletes resources + // and never needs to emit task state events. + backend, err := k8sbackend.NewBackend(ctx, conf, reader, &events.Logger{Log: log}, log) + if err != nil { + return fmt.Errorf("initializing kubernetes backend: %v", err) + } + + log.Info("Starting orphaned resource cleanup", + "namespace", conf.Kubernetes.JobsNamespace) + backend.CleanOrphanedResources(ctx) + log.Info("Orphaned resource cleanup complete") + return nil + }, + } + + f := cmd.Flags() + f.StringVarP(&configFile, "config", "c", "", "Path to Funnel config file") + cmd.SetGlobalNormalizationFunc(cmdutil.NormalizeFlags) + f.AddFlagSet(cmdutil.ServerFlags(flagConf, &configFile)) + + return cmd +} + +// openReader opens a read-only connection to the configured Funnel database. +func openReader(ctx context.Context, conf *config.Config) (tes.ReadOnlyServer, error) { + switch strings.ToLower(conf.Database) { + case "boltdb": + return boltdb.NewBoltDB(conf.BoltDB) + case "badger": + return badger.NewBadger(conf.Badger) + case "datastore": + return datastore.NewDatastore(conf.Datastore) + case "dynamodb": + return dynamodb.NewDynamoDB(conf.DynamoDB) + case "elastic": + return elastic.NewElastic(conf.Elastic) + case "mongodb": + return mongodb.NewMongoDB(conf.MongoDB) + case "postgres", "psql": + return postgres.NewPostgres(conf.Postgres) + default: + return nil, fmt.Errorf("unknown database: '%s'", conf.Database) + } +} diff --git a/cmd/root.go b/cmd/root.go index 84841311a..3fb90248e 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -5,6 +5,7 @@ import ( "github.com/ohsu-comp-bio/funnel/cmd/aws" "github.com/ohsu-comp-bio/funnel/cmd/examples" "github.com/ohsu-comp-bio/funnel/cmd/gce" + "github.com/ohsu-comp-bio/funnel/cmd/kubernetes" "github.com/ohsu-comp-bio/funnel/cmd/node" "github.com/ohsu-comp-bio/funnel/cmd/run" "github.com/ohsu-comp-bio/funnel/cmd/server" @@ -27,6 +28,7 @@ func init() { RootCmd.AddCommand(aws.Cmd) RootCmd.AddCommand(examples.Cmd) RootCmd.AddCommand(gce.Cmd) + RootCmd.AddCommand(kubernetes.Cmd) RootCmd.AddCommand(completionCmd) RootCmd.AddCommand(genMarkdownCmd) RootCmd.AddCommand(node.NewCommand()) diff --git a/compute/kubernetes/backend.go b/compute/kubernetes/backend.go index 5f77ced47..05f199d3a 100644 --- a/compute/kubernetes/backend.go +++ b/compute/kubernetes/backend.go @@ -171,55 +171,63 @@ func (b *Backend) createResources(ctx context.Context, task *tes.Task, config *c defer cancel() } - // If the task has inputs, outputs, or declared volumes, create a PVC so - // executor pods can share data via PVC subPath mounts. - if len(task.Inputs) > 0 || len(task.Outputs) > 0 || len(task.Volumes) > 0 { - b.log.Debug("creating Worker PV", "taskID", task.Id) - - // Create PV - err := resources.CreatePV(timeoutCtx, task.Id, - config, - b.client, b.log) - if err != nil { - _ = b.Cancel(context.Background(), task.Id) - return fmt.Errorf("creating Worker PV: %w", err) - } - - // Create PVC - b.log.Debug("creating Worker PVC", "taskID", task.Id) - err = resources.CreatePVC(timeoutCtx, task.Id, config, b.client, b.log) - if err != nil { - _ = b.Cancel(context.Background(), task.Id) - return fmt.Errorf("creating Worker PVC: %w", err) - } + // Create Worker Job first so its UID can be used as an owner reference on + // all subordinate namespaced resources, enabling automatic K8s GC cleanup. + b.log.Debug("creating Worker Job", "taskID", task.Id) + job, err := resources.CreateJob(timeoutCtx, task, config, b.client, b.log) + if err != nil { + _ = b.Cancel(context.Background(), task.Id) + return fmt.Errorf("creating Worker Job: %w", err) } - // err is declared here for use across all conditional resource-creation - // blocks below (ConfigMap, ServiceAccount, Role, RoleBinding, Job). - var err error + blockOwnerDeletion := true + isController := true + ownerRef := &metav1.OwnerReference{ + APIVersion: "batch/v1", + Kind: "Job", + Name: job.Name, + UID: job.UID, + BlockOwnerDeletion: &blockOwnerDeletion, + Controller: &isController, + } + // Create ConfigMap (only when a template is configured; deployments using a + // static shared ConfigMap via the WorkerTemplate volume spec skip this). if config.Kubernetes.ConfigMapTemplate != "" { b.log.Debug("creating Worker ConfigMap", "taskID", task.Id) - err = resources.CreateConfigMap(timeoutCtx, task.Id, config, b.client, b.log) + err = resources.CreateConfigMap(timeoutCtx, task.Id, config, b.client, b.log, ownerRef) if err != nil { _ = b.Cancel(context.Background(), task.Id) return fmt.Errorf("creating Worker ConfigMap: %w", err) } } + // Create ServiceAccount, Role, and RoleBinding only when templates are + // configured. Deployments that supply a pre-existing shared SA (e.g. via + // _WORKER_SA tag or a static Helm-managed SA) skip these steps entirely. + // External (user-managed) SAs are not owned by the Job — they outlive tasks. if config.Kubernetes.ServiceAccountTemplate != "" { saName := fmt.Sprintf("funnel-worker-sa-%s-%s", config.Kubernetes.JobsNamespace, task.Id) - if _, exists := task.Tags["_WORKER_SA"]; exists { - saName = task.Tags["_WORKER_SA"] + externalSA := false + if sa, exists := task.Tags["_WORKER_SA"]; exists && sa != "" { + saName = sa + externalSA = true } // TODO: Add error handler to handle case where Get fails for reasons other than `NotFound` // e.g. network issues, permission issues, etc. _, err = b.client.CoreV1().ServiceAccounts(config.Kubernetes.JobsNamespace).Get(timeoutCtx, saName, metav1.GetOptions{}) + + // ServiceAccount does not exist, create it if err != nil { b.log.Debug("Error getting ServiceAccount:", "ServiceAccount", saName, "taskID", task.Id, "error", err) b.log.Debug("Creating Worker ServiceAccount", "taskID", task.Id) - err = resources.CreateServiceAccount(timeoutCtx, task, config, b.client, b.log) + // Only set the owner reference for task-level SAs; external SAs are shared and must not be GC'd with the job. + saOwnerRef := ownerRef + if externalSA { + saOwnerRef = nil + } + err = resources.CreateServiceAccount(timeoutCtx, task, config, b.client, b.log, saOwnerRef) if err != nil { _ = b.Cancel(context.Background(), task.Id) return fmt.Errorf("creating Worker ServiceAccount: %w", err) @@ -231,7 +239,7 @@ func (b *Backend) createResources(ctx context.Context, task *tes.Task, config *c if config.Kubernetes.RoleTemplate != "" { b.log.Debug("creating Worker Role", "taskID", task.Id) - err = resources.CreateRole(timeoutCtx, task, config, b.client, b.log) + err = resources.CreateRole(timeoutCtx, task, config, b.client, b.log, ownerRef) if err != nil { _ = b.Cancel(context.Background(), task.Id) return fmt.Errorf("creating Worker Role: %w", err) @@ -240,19 +248,38 @@ func (b *Backend) createResources(ctx context.Context, task *tes.Task, config *c if config.Kubernetes.RoleBindingTemplate != "" { b.log.Debug("creating Worker RoleBinding", "taskID", task.Id) - err = resources.CreateRoleBinding(timeoutCtx, task, config, b.client, b.log) + err = resources.CreateRoleBinding(timeoutCtx, task, config, b.client, b.log, ownerRef) if err != nil { _ = b.Cancel(context.Background(), task.Id) return fmt.Errorf("creating Worker RoleBinding: %w", err) } } - // Create Worker Job - b.log.Debug("creating Worker Job", "taskID", task.Id) - err = resources.CreateJob(timeoutCtx, task, config, b.client, b.log) - if err != nil { - _ = b.Cancel(context.Background(), task.Id) - return fmt.Errorf("creating Worker Job: %w", err) + // If the task has inputs, outputs, or declared volumes, create a PVC so + // executor pods can share data via PVC subPath mounts. + if len(task.Inputs) > 0 || len(task.Outputs) > 0 || len(task.Volumes) > 0 { + b.log.Debug("creating Worker PV", "taskID", task.Id) + + // Check to make sure required configs are present + if config.GenericS3 == nil || len(config.GenericS3) == 0 || + config.GenericS3[0].Bucket == "" || config.GenericS3[0].Region == "" { + return fmt.Errorf("Bucket or Region not found in GenericS3 config when attempting to create resources for task: %#v", task) + } + + // Create PV (cluster-scoped — cannot be owned by a namespaced Job) + err = resources.CreatePV(timeoutCtx, task.Id, config, b.client, b.log) + if err != nil { + _ = b.Cancel(context.Background(), task.Id) + return fmt.Errorf("creating Worker PV: %w", err) + } + + // Create PVC + b.log.Debug("creating Worker PVC", "taskID", task.Id) + err = resources.CreatePVC(timeoutCtx, task.Id, config, b.client, b.log, ownerRef) + if err != nil { + _ = b.Cancel(context.Background(), task.Id) + return fmt.Errorf("creating Worker PVC: %w", err) + } } return nil @@ -262,6 +289,18 @@ func (b *Backend) createResources(ctx context.Context, task *tes.Task, config *c func (b *Backend) cleanResources(ctx context.Context, taskId string) error { var errs error + // Check whether this task used an externally-managed ServiceAccount (e.g. + // Gen3Workflow per-user SA supplied via _WORKER_SA tag). If so, skip SA + // deletion — the SA is shared across tasks and must not be torn down here. + externalSA := false + if b.database != nil { + if task, err := b.database.GetTask(ctx, &tes.GetTaskRequest{Id: taskId, View: tes.View_FULL.String()}); err == nil { + if saName, exists := task.Tags["_WORKER_SA"]; exists && saName != "" { + externalSA = true + } + } + } + // Delete Job b.log.Debug("deleting Job", "taskID", taskId) err := resources.DeleteJob(ctx, b.conf, taskId, b.client, b.log) @@ -294,7 +333,7 @@ func (b *Backend) cleanResources(ctx context.Context, taskId string) error { } // Delete ServiceAccount - err = resources.DeleteServiceAccount(ctx, taskId, b.conf.Kubernetes.JobsNamespace, b.client, b.log) + err = resources.DeleteServiceAccount(ctx, taskId, b.conf.Kubernetes.JobsNamespace, b.client, b.log, externalSA) if err != nil { errs = multierror.Append(errs, err) b.log.Error("deleting Worker ServiceAccount", "error", err) @@ -557,10 +596,6 @@ func (b *Backend) reconcile(ctx context.Context, rate time.Duration, disableClea } delete(failedJobEvents, taskID) } - - // Clean up all orphaned Funnel-managed resources whose task no longer exists in - // the DB. These can be left behind by server crashes or partial cleanup failures. - b.cleanOrphanedResources(ctx) } } } @@ -581,12 +616,13 @@ func (b *Backend) isResourceCleanupNeeded(ctx context.Context, taskID string) (b } } -// cleanOrphanedResources deletes any Funnel-managed Kubernetes resources that are not associated with an active task -// in the database. +// CleanOrphanedResources deletes any Funnel-managed Kubernetes resources that are not associated +// with an active task in the database. // -// This is a safety measure to prevent resource leaks from orphaned jobs whose tasks have been -// deleted or completed, but whose resources were not cleaned up due to transient errors or server crashes. -func (b *Backend) cleanOrphanedResources(ctx context.Context) { +// This is intended to be called as a one-shot operation (e.g. from a Kubernetes CronJob) rather +// than as a long-running goroutine, so that cleanup is decoupled from the Funnel server lifecycle +// and multiple server replicas do not race to clean the same resources simultaneously. +func (b *Backend) CleanOrphanedResources(ctx context.Context) { namespace := b.conf.Kubernetes.JobsNamespace taskIDs := make(map[string]struct{}) @@ -678,5 +714,9 @@ func (b *Backend) cleanOrphanedResources(ctx context.Context) { if err := b.cleanResources(ctx, taskID); err != nil { b.log.Error("backlog cleanup: failed to clean resources", "taskID", taskID, "error", err) } + + // Sleep briefly between deletions to avoid overwhelming the API server if there are many orphaned resources. + // Test this + see if better to sleep ~1 second for every 10 tasks... + time.Sleep(500 * time.Millisecond) } } diff --git a/compute/kubernetes/resources/configmap.go b/compute/kubernetes/resources/configmap.go index 23d7e9029..18df70f0a 100644 --- a/compute/kubernetes/resources/configmap.go +++ b/compute/kubernetes/resources/configmap.go @@ -21,7 +21,7 @@ import ( // This function is only called when ConfigMapTemplate is non-empty; deployments // that mount a static shared ConfigMap (e.g. "funnel-config") via the // WorkerTemplate volume spec should leave ConfigMapTemplate empty. -func CreateConfigMap(ctx context.Context, taskId string, conf *config.Config, client kubernetes.Interface, log *logger.Logger) error { +func CreateConfigMap(ctx context.Context, taskId string, conf *config.Config, client kubernetes.Interface, log *logger.Logger, ownerRef *metav1.OwnerReference) error { t, err := template.New(taskId).Parse(conf.Kubernetes.ConfigMapTemplate) if err != nil { return fmt.Errorf("parsing ConfigMapTemplate: %v", err) @@ -53,6 +53,10 @@ func CreateConfigMap(ctx context.Context, taskId string, conf *config.Config, cl return fmt.Errorf("ConfigMapTemplate did not produce a ConfigMap object") } + if ownerRef != nil { + cm.OwnerReferences = []metav1.OwnerReference{*ownerRef} + } + _, err = client.CoreV1().ConfigMaps(conf.Kubernetes.JobsNamespace).Create(ctx, cm, metav1.CreateOptions{}) if err != nil { return fmt.Errorf("%v", err) diff --git a/compute/kubernetes/resources/job.go b/compute/kubernetes/resources/job.go index c48c1495b..bd40069b8 100644 --- a/compute/kubernetes/resources/job.go +++ b/compute/kubernetes/resources/job.go @@ -35,20 +35,20 @@ func SanitizeLabelValue(s string) string { // Create the Funnel Worker job from kubernetes-template.yaml // Executor job is created in worker/kubernetes.go#Run -func CreateJob(ctx context.Context, task *tes.Task, conf *config.Config, client kubernetes.Interface, log *logger.Logger) error { +func CreateJob(ctx context.Context, task *tes.Task, conf *config.Config, client kubernetes.Interface, log *logger.Logger) (*v1.Job, error) { // Parse Worker Template log.Debug("Creating job from template", "template", conf.Kubernetes.WorkerTemplate) t, err := template.New(task.Id).Parse(conf.Kubernetes.WorkerTemplate) if err != nil { - return fmt.Errorf("%v", err) + return nil, fmt.Errorf("%v", err) } pods, err := client.CoreV1().Pods(conf.Kubernetes.Namespace).List(ctx, metav1.ListOptions{ LabelSelector: "app=funnel", }) if err != nil { - return fmt.Errorf("failed to list pods: %v", err) + return nil, fmt.Errorf("failed to list pods: %v", err) } var image string @@ -96,19 +96,19 @@ func CreateJob(ctx context.Context, task *tes.Task, conf *config.Config, client var buf bytes.Buffer err = t.Execute(&buf, templateData) if err != nil { - return fmt.Errorf("%v", err) + return nil, fmt.Errorf("%v", err) } log.Debug("Job template", "template", buf.String()) decode := scheme.Codecs.UniversalDeserializer().Decode obj, _, err := decode(buf.Bytes(), nil, nil) if err != nil { - return err + return nil, err } job, ok := obj.(*v1.Job) if !ok { - return fmt.Errorf("failed to decode job spec") + return nil, fmt.Errorf("failed to decode job spec") } // Ensure completed jobs are garbage-collected by the Kubernetes TTL Controller @@ -120,12 +120,12 @@ func CreateJob(ctx context.Context, task *tes.Task, conf *config.Config, client } log.Debug("Creating job", "Job", job.Name, "JobsNamespace", conf.Kubernetes.JobsNamespace) - _, err = client.BatchV1().Jobs(conf.Kubernetes.JobsNamespace).Create(ctx, job, metav1.CreateOptions{}) + created, err := client.BatchV1().Jobs(conf.Kubernetes.JobsNamespace).Create(ctx, job, metav1.CreateOptions{}) if err != nil { - return fmt.Errorf("%v", err) + return nil, fmt.Errorf("%v", err) } - return nil + return created, nil } // DeleteJob removes the worker job for a task and all associated executor jobs. diff --git a/compute/kubernetes/resources/pv.go b/compute/kubernetes/resources/pv.go index e38ccad93..12afe9f8f 100644 --- a/compute/kubernetes/resources/pv.go +++ b/compute/kubernetes/resources/pv.go @@ -5,11 +5,13 @@ import ( "context" "fmt" "text/template" + "time" "github.com/ohsu-comp-bio/funnel/config" "github.com/ohsu-comp-bio/funnel/logger" corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" "k8s.io/client-go/kubernetes/scheme" @@ -63,28 +65,46 @@ func CreatePV(ctx context.Context, taskId string, conf *config.Config, client ku return nil } -// Add this helper function for PV cleanup +// DeletePV removes the PV for a task, retrying on conflict errors that occur +// when another process (e.g. reconciler + cancel running concurrently) modifies +// the PV between our Get and Update calls. func DeletePV(ctx context.Context, taskID string, client kubernetes.Interface, log *logger.Logger) error { name := fmt.Sprintf("funnel-worker-pv-%s", taskID) - // The PV may not have been made. Some jobs with no I/O don't need a PV or it may have already been deleted. - pv, err := client.CoreV1().PersistentVolumes().Get(ctx, name, metav1.GetOptions{}) - if err != nil { - return nil - } - // Remove the pv-protection finalizer so Kubernetes allows deletion - if len(pv.Finalizers) > 0 { - pv.Finalizers = nil - _, err = client.CoreV1().PersistentVolumes().Update(ctx, pv, metav1.UpdateOptions{}) + const maxRetries = 5 + delay := 100 * time.Millisecond + for i := range maxRetries { + // The PV may not exist (no I/O task, or already deleted). + pv, err := client.CoreV1().PersistentVolumes().Get(ctx, name, metav1.GetOptions{}) if err != nil { - return fmt.Errorf("removing finalizers from PV %s: %v", name, err) + if errors.IsNotFound(err) { + return nil + } + return fmt.Errorf("getting PV %s: %v", name, err) } - } - log.Debug("deleting Worker PV", "taskID", taskID) - err = client.CoreV1().PersistentVolumes().Delete(ctx, name, metav1.DeleteOptions{}) - if err != nil { - return fmt.Errorf("%v", err) + // Remove the pv-protection finalizer so Kubernetes allows deletion. + if len(pv.Finalizers) > 0 { + pv.Finalizers = nil + _, err = client.CoreV1().PersistentVolumes().Update(ctx, pv, metav1.UpdateOptions{}) + if err != nil { + if errors.IsConflict(err) && i < maxRetries-1 { + log.Debug("conflict removing PV finalizers, retrying", "pv", name, "attempt", i+1) + time.Sleep(delay) + delay *= 2 + continue + } + return fmt.Errorf("removing finalizers from PV %s: %v", name, err) + } + } + + log.Debug("deleting Worker PV", "taskID", taskID) + err = client.CoreV1().PersistentVolumes().Delete(ctx, name, metav1.DeleteOptions{}) + if err != nil && !errors.IsNotFound(err) { + return fmt.Errorf("deleting PV %s: %v", name, err) + } + return nil } - return nil + + return fmt.Errorf("removing finalizers from PV %s: exceeded max retries", name) } diff --git a/compute/kubernetes/resources/pvc.go b/compute/kubernetes/resources/pvc.go index af1c07bbf..b78ab813e 100644 --- a/compute/kubernetes/resources/pvc.go +++ b/compute/kubernetes/resources/pvc.go @@ -5,10 +5,12 @@ import ( "context" "fmt" "text/template" + "time" "github.com/ohsu-comp-bio/funnel/config" "github.com/ohsu-comp-bio/funnel/logger" corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" "k8s.io/client-go/kubernetes/scheme" @@ -16,7 +18,7 @@ import ( // Create the Worker/Executor PVC from config/kubernetes-pvc.yaml // TODO: Move this config file to Helm Charts so users can see/customize it -func CreatePVC(ctx context.Context, taskId string, conf *config.Config, client kubernetes.Interface, log *logger.Logger) error { +func CreatePVC(ctx context.Context, taskId string, conf *config.Config, client kubernetes.Interface, log *logger.Logger, ownerRef *metav1.OwnerReference) error { jobNamespace := conf.Kubernetes.JobsNamespace @@ -55,6 +57,10 @@ func CreatePVC(ctx context.Context, taskId string, conf *config.Config, client k return fmt.Errorf("failed to decode PVC spec") } + if ownerRef != nil { + pvc.OwnerReferences = []metav1.OwnerReference{*ownerRef} + } + _, err = client.CoreV1().PersistentVolumeClaims(jobNamespace).Create(ctx, pvc, metav1.CreateOptions{}) if err != nil { return fmt.Errorf("%v", err) @@ -63,28 +69,46 @@ func CreatePVC(ctx context.Context, taskId string, conf *config.Config, client k return nil } -// Add this helper function for PVC cleanup +// DeletePVC removes the PVC for a task, retrying on conflict errors that occur +// when another process (e.g. reconciler + cancel running concurrently) modifies +// the PVC between our Get and Update calls. func DeletePVC(ctx context.Context, taskID string, namespace string, client kubernetes.Interface, log *logger.Logger) error { name := fmt.Sprintf("funnel-worker-pvc-%s", taskID) - pvc, err := client.CoreV1().PersistentVolumeClaims(namespace).Get(ctx, name, metav1.GetOptions{}) - if err != nil { - return nil - } - // Remove the pvc-protection finalizer so Kubernetes allows deletion - if len(pvc.Finalizers) > 0 { - pvc.Finalizers = nil - _, err = client.CoreV1().PersistentVolumeClaims(namespace).Update(ctx, pvc, metav1.UpdateOptions{}) + const maxRetries = 5 + delay := 100 * time.Millisecond + for i := range maxRetries { + // The PVC may not exist (no I/O task, or already deleted). + pvc, err := client.CoreV1().PersistentVolumeClaims(namespace).Get(ctx, name, metav1.GetOptions{}) if err != nil { - return fmt.Errorf("removing finalizers from PVC %s: %v", name, err) + if errors.IsNotFound(err) { + return nil + } + return fmt.Errorf("getting PVC %s: %v", name, err) } - } - log.Debug("deleting Worker PVC", "taskID", taskID) - err = client.CoreV1().PersistentVolumeClaims(namespace).Delete(ctx, name, metav1.DeleteOptions{}) - if err != nil { - return fmt.Errorf("deleting shared PVC: %v", err) + // Remove the pvc-protection finalizer so Kubernetes allows deletion. + if len(pvc.Finalizers) > 0 { + pvc.Finalizers = nil + _, err = client.CoreV1().PersistentVolumeClaims(namespace).Update(ctx, pvc, metav1.UpdateOptions{}) + if err != nil { + if errors.IsConflict(err) && i < maxRetries-1 { + log.Debug("conflict removing PVC finalizers, retrying", "pvc", name, "attempt", i+1) + time.Sleep(delay) + delay *= 2 + continue + } + return fmt.Errorf("removing finalizers from PVC %s: %v", name, err) + } + } + + log.Debug("deleting Worker PVC", "taskID", taskID) + err = client.CoreV1().PersistentVolumeClaims(namespace).Delete(ctx, name, metav1.DeleteOptions{}) + if err != nil && !errors.IsNotFound(err) { + return fmt.Errorf("deleting PVC %s: %v", name, err) + } + return nil } - return nil + return fmt.Errorf("removing finalizers from PVC %s: exceeded max retries", name) } diff --git a/compute/kubernetes/resources/resources_test.go b/compute/kubernetes/resources/resources_test.go index 791f869e2..d5b4639c3 100644 --- a/compute/kubernetes/resources/resources_test.go +++ b/compute/kubernetes/resources/resources_test.go @@ -120,7 +120,7 @@ data: funnel-worker.yaml: | placeholder` - err := CreateConfigMap(ctx, testTaskID, conf, fake.NewSimpleClientset(), l) + err := CreateConfigMap(ctx, testTaskID, conf, fake.NewSimpleClientset(), l, nil) if err != nil { t.Errorf("CreateConfigMap failed: %v", err) } @@ -197,7 +197,7 @@ func TestCreateJob(t *testing.T) { conf := config.DefaultConfig() conf.Kubernetes.JobsNamespace = jobsNamespace conf.Kubernetes.WorkerTemplate = minimalWorkerTemplate - err := CreateJob(ctx, task, conf, fake.NewSimpleClientset(), l) + _, err := CreateJob(ctx, task, conf, fake.NewSimpleClientset(), l) if err != nil { t.Errorf("CreateJob failed: %v", err) } @@ -273,7 +273,7 @@ func TestCreatePVC(t *testing.T) { conf := config.DefaultConfig() conf.Kubernetes.JobsNamespace = jobsNamespace conf.Kubernetes.PVCTemplate = minimalPVCTemplate - err := CreatePVC(ctx, testTaskID, conf, fake.NewSimpleClientset(), l) + err := CreatePVC(ctx, testTaskID, conf, fake.NewSimpleClientset(), l, nil) if err != nil { t.Errorf("CreatePVC failed: %v", err) } @@ -315,7 +315,7 @@ func TestCreateJobWithNoResources(t *testing.T) { conf := config.DefaultConfig() conf.Kubernetes.JobsNamespace = jobsNamespace conf.Kubernetes.WorkerTemplate = minimalWorkerTemplate - err := CreateJob(ctx, task, conf, fake.NewSimpleClientset(), l) + _, err := CreateJob(ctx, task, conf, fake.NewSimpleClientset(), l) if err != nil { t.Errorf("CreateJob failed with nil resources: %v", err) } @@ -357,7 +357,7 @@ func TestCreateServiceAccount(t *testing.T) { conf := config.DefaultConfig() conf.Kubernetes.JobsNamespace = jobsNamespace conf.Kubernetes.ServiceAccountTemplate = minimalServiceAccountTemplate - err := CreateServiceAccount(ctx, task, conf, fake.NewSimpleClientset(), l) + err := CreateServiceAccount(ctx, task, conf, fake.NewSimpleClientset(), l, nil) if err != nil { t.Errorf("CreateServiceAccount failed: %v", err) } @@ -370,6 +370,10 @@ func TestDeleteServiceAccount(t *testing.T) { ObjectMeta: metav1.ObjectMeta{ Name: "funnel-worker-sa-" + testTaskID, Namespace: namespace, + Labels: map[string]string{ + "app": "funnel", + "taskId": testTaskID, + }, }, } _, err := fakeClient.CoreV1().ServiceAccounts(namespace).Create(context.Background(), sa, metav1.CreateOptions{}) @@ -377,12 +381,55 @@ func TestDeleteServiceAccount(t *testing.T) { t.Fatalf("Failed to create test ServiceAccount: %v", err) } - err = DeleteServiceAccount(context.Background(), testTaskID, namespace, fakeClient, l) + err = DeleteServiceAccount(context.Background(), testTaskID, namespace, fakeClient, l, false) if err != nil { t.Errorf("DeleteServiceAccount failed: %v", err) } } +func TestDeleteServiceAccountInUse(t *testing.T) { + fakeClient := fake.NewSimpleClientset() + + sa := &corev1.ServiceAccount{ + ObjectMeta: metav1.ObjectMeta{ + Name: "funnel-worker-sa-" + testTaskID, + Namespace: namespace, + Labels: map[string]string{ + "app": "funnel", + "taskId": testTaskID, + }, + }, + } + _, err := fakeClient.CoreV1().ServiceAccounts(namespace).Create(context.Background(), sa, metav1.CreateOptions{}) + if err != nil { + t.Fatalf("Failed to create test ServiceAccount: %v", err) + } + + pod := &corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-pod", + Namespace: namespace, + }, + Spec: corev1.PodSpec{ + ServiceAccountName: sa.Name, + }, + } + _, err = fakeClient.CoreV1().Pods(namespace).Create(context.Background(), pod, metav1.CreateOptions{}) + if err != nil { + t.Fatalf("Failed to create test Pod: %v", err) + } + + err = DeleteServiceAccount(context.Background(), testTaskID, namespace, fakeClient, l, false) + if err == nil { + t.Fatal("expected DeleteServiceAccount to fail when ServiceAccount is in use") + } + + _, err = fakeClient.CoreV1().ServiceAccounts(namespace).Get(context.Background(), sa.Name, metav1.GetOptions{}) + if err != nil { + t.Fatalf("expected ServiceAccount to remain after failed delete: %v", err) + } +} + func TestCreateRole(t *testing.T) { task := &tes.Task{ Id: testTaskID, @@ -391,7 +438,7 @@ func TestCreateRole(t *testing.T) { conf := config.DefaultConfig() conf.Kubernetes.JobsNamespace = jobsNamespace conf.Kubernetes.RoleTemplate = minimalRoleTemplate - err := CreateRole(ctx, task, conf, fake.NewSimpleClientset(), l) + err := CreateRole(ctx, task, conf, fake.NewSimpleClientset(), l, nil) if err != nil { t.Errorf("CreateRole failed: %v", err) } @@ -425,7 +472,7 @@ func TestCreateRoleBinding(t *testing.T) { conf := config.DefaultConfig() conf.Kubernetes.JobsNamespace = jobsNamespace conf.Kubernetes.RoleBindingTemplate = minimalRoleBindingTemplate - err := CreateRoleBinding(ctx, task, conf, fake.NewSimpleClientset(), l) + err := CreateRoleBinding(ctx, task, conf, fake.NewSimpleClientset(), l, nil) if err != nil { t.Errorf("CreateRoleBinding failed: %v", err) } diff --git a/compute/kubernetes/resources/role.go b/compute/kubernetes/resources/role.go index 4dc416d20..0086652fa 100644 --- a/compute/kubernetes/resources/role.go +++ b/compute/kubernetes/resources/role.go @@ -17,7 +17,7 @@ import ( ) // Create the Worker/Executor Role from config/kubernetes-role.yaml -func CreateRole(ctx context.Context, task *tes.Task, conf *config.Config, client kubernetes.Interface, log *logger.Logger) error { +func CreateRole(ctx context.Context, task *tes.Task, conf *config.Config, client kubernetes.Interface, log *logger.Logger, ownerRef *metav1.OwnerReference) error { // Load templates t, err := template.New(task.Id).Parse(conf.Kubernetes.RoleTemplate) @@ -47,6 +47,10 @@ func CreateRole(ctx context.Context, task *tes.Task, conf *config.Config, client return fmt.Errorf("failed to verify Role spec") } + if ownerRef != nil { + role.OwnerReferences = []metav1.OwnerReference{*ownerRef} + } + _, err = client.RbacV1().Roles(conf.Kubernetes.JobsNamespace).Create(ctx, role, metav1.CreateOptions{}) if err != nil { return fmt.Errorf("failed to create Role: %v", err) diff --git a/compute/kubernetes/resources/roleBinding.go b/compute/kubernetes/resources/roleBinding.go index 4455d97d5..92f07c2c9 100644 --- a/compute/kubernetes/resources/roleBinding.go +++ b/compute/kubernetes/resources/roleBinding.go @@ -17,7 +17,7 @@ import ( ) // Create the Worker/Executor RoleBinding from config/kubernetes-rolebinding.yaml -func CreateRoleBinding(ctx context.Context, task *tes.Task, conf *config.Config, client kubernetes.Interface, log *logger.Logger) error { +func CreateRoleBinding(ctx context.Context, task *tes.Task, conf *config.Config, client kubernetes.Interface, log *logger.Logger, ownerRef *metav1.OwnerReference) error { // Load templates t, err := template.New(task.Id).Parse(conf.Kubernetes.RoleBindingTemplate) @@ -50,11 +50,15 @@ func CreateRoleBinding(ctx context.Context, task *tes.Task, conf *config.Config, return fmt.Errorf("failed to decode Role spec: %v", err) } - roleBinding, ok := obj.(*rbacv1.RoleBinding) // Change from corev1.Role to rbacv1.Role + roleBinding, ok := obj.(*rbacv1.RoleBinding) if !ok { return fmt.Errorf("failed to decode RoleBinding spec") } + if ownerRef != nil { + roleBinding.OwnerReferences = []metav1.OwnerReference{*ownerRef} + } + _, err = client.RbacV1().RoleBindings(conf.Kubernetes.JobsNamespace).Create(ctx, roleBinding, metav1.CreateOptions{}) if err != nil { return fmt.Errorf("failed to create RoleBinding: %v", err) diff --git a/compute/kubernetes/resources/serviceaccount.go b/compute/kubernetes/resources/serviceaccount.go index a22cac4b8..b17715a72 100644 --- a/compute/kubernetes/resources/serviceaccount.go +++ b/compute/kubernetes/resources/serviceaccount.go @@ -17,7 +17,7 @@ import ( ) // Create the Worker/Executor ServiceAccount from config/kubernetes-serviceaccount.yaml -func CreateServiceAccount(ctx context.Context, task *tes.Task, conf *config.Config, client kubernetes.Interface, log *logger.Logger) error { +func CreateServiceAccount(ctx context.Context, task *tes.Task, conf *config.Config, client kubernetes.Interface, log *logger.Logger, ownerRef *metav1.OwnerReference) error { // Load templates t, err := template.New(task.Id).Parse(conf.Kubernetes.ServiceAccountTemplate) @@ -60,6 +60,10 @@ func CreateServiceAccount(ctx context.Context, task *tes.Task, conf *config.Conf return fmt.Errorf("failed to cast to ServiceAccount spec") } + if ownerRef != nil { + sa.OwnerReferences = []metav1.OwnerReference{*ownerRef} + } + _, err = client.CoreV1().ServiceAccounts(conf.Kubernetes.JobsNamespace).Create(ctx, sa, metav1.CreateOptions{}) if err != nil { return fmt.Errorf("failed to create ServiceAccount: %v", err) @@ -68,8 +72,25 @@ func CreateServiceAccount(ctx context.Context, task *tes.Task, conf *config.Conf return nil } +func isServiceAccountAttachedToPods(ctx context.Context, saName, namespace string, client kubernetes.Interface) (bool, error) { + pods, err := client.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{ + FieldSelector: fmt.Sprintf("spec.serviceAccountName=%s", saName), + }) + if err != nil { + return false, fmt.Errorf("listing pods using ServiceAccount %s: %v", saName, err) + } + return len(pods.Items) > 0, nil +} + // DeleteServiceAccount deletes the ServiceAccount created for a task. -func DeleteServiceAccount(ctx context.Context, taskID string, namespace string, client kubernetes.Interface, log *logger.Logger) error { +// If externalSA is true the ServiceAccount is externally managed (e.g. a +// Gen3Workflow per-user SA supplied via the _WORKER_SA task tag) and must not +// be deleted by Funnel. +func DeleteServiceAccount(ctx context.Context, taskID string, namespace string, client kubernetes.Interface, log *logger.Logger, externalSA bool) error { + if externalSA { + log.Debug("skipping deletion of externally-managed ServiceAccount", "taskID", taskID) + return nil + } // ServiceAccount names are not available here without config, so we list by label. sas, err := client.CoreV1().ServiceAccounts(namespace).List(ctx, metav1.ListOptions{ LabelSelector: fmt.Sprintf("app=funnel,taskId=%s", taskID), @@ -78,6 +99,14 @@ func DeleteServiceAccount(ctx context.Context, taskID string, namespace string, return fmt.Errorf("listing ServiceAccounts for task %s: %v", taskID, err) } for _, sa := range sas.Items { + inUse, err := isServiceAccountAttachedToPods(ctx, sa.Name, namespace, client) + if err != nil { + return err + } + if inUse { + return fmt.Errorf("serviceAccount %s is still in use by active pod(s)", sa.Name) + } + log.Debug("deleting Worker ServiceAccount", "name", sa.Name, "taskID", taskID) if err := client.CoreV1().ServiceAccounts(namespace).Delete(ctx, sa.Name, metav1.DeleteOptions{}); err != nil { return fmt.Errorf("deleting ServiceAccount %s: %v", sa.Name, err) diff --git a/storage/generic_s3.go b/storage/generic_s3.go index 49c117b58..3875d3e68 100644 --- a/storage/generic_s3.go +++ b/storage/generic_s3.go @@ -2,11 +2,14 @@ package storage import ( "context" + "errors" "fmt" "io" "os" "path/filepath" "strings" + "syscall" + "time" "github.com/minio/minio-go/v7" "github.com/minio/minio-go/v7/pkg/credentials" @@ -268,11 +271,30 @@ func (s3 *GenericS3) Put(ctx context.Context, url, path string) (*Object, error) opts.ServerSideEncryption = SSEKMS } - // Check if the path is a directory - fileInfo, err := os.Stat(path) - if err != nil { - return nil, err + // Wait for the local path to become readable. When the local path is on a + // Mountpoint-for-S3 backed filesystem and the file was written by another + // mount instance (the executor pod), os.Stat returns EPERM until Mountpoint + // finishes flushing the write to S3. Poll until the file is readable or the + // timeout expires. + const mountpointFlushTimeout = 30 * time.Second + const mountpointFlushInterval = 2 * time.Second + deadline := time.Now().Add(mountpointFlushTimeout) + var fileInfo os.FileInfo + for { + var err error + fileInfo, err = os.Stat(path) + if err == nil { + break + } + if time.Now().Before(deadline) && (errors.Is(err, syscall.EPERM) || errors.Is(err, os.ErrNotExist)) { + logger.Debug("genericS3: waiting for output file to become readable", "path", path, "error", err) + time.Sleep(mountpointFlushInterval) + continue + } + return nil, fmt.Errorf("genericS3: putting object %s: %v", url, err) } + + // Check if the path is a directory if fileInfo.IsDir() { // Walk the directory and upload all files and subdirectories err = filepath.Walk(path, func(filePath string, info os.FileInfo, err error) error { @@ -288,7 +310,7 @@ func (s3 *GenericS3) Put(ctx context.Context, url, path string) (*Object, error) uploadPath := filepath.Join(u.path, relativePath) _, err = s3.client.FPutObject(ctx, u.bucket, uploadPath, filePath, opts) if err != nil { - return fmt.Errorf("genericS3: putting object %s: %v", url, err) + return fmt.Errorf("genericS3: putting nested object %s: %v", url, err) } } return nil @@ -297,7 +319,6 @@ func (s3 *GenericS3) Put(ctx context.Context, url, path string) (*Object, error) return nil, err } } else { - // Upload the file directly _, err = s3.client.FPutObject(ctx, u.bucket, u.path, path, opts) if err != nil { return nil, fmt.Errorf("genericS3: putting object %s: %v", url, err) diff --git a/tests/fixtures/quotes/single-element-single-quote.json b/tests/fixtures/quotes/single-element-single-quote.json new file mode 100644 index 000000000..7aa6ec9c8 --- /dev/null +++ b/tests/fixtures/quotes/single-element-single-quote.json @@ -0,0 +1,9 @@ +{ + "name": "Single-element exec — single quote in argument", + "executors": [ + { + "image": "alpine", + "command": ["echo \"Hello O'hare!\""] + } + ] +} diff --git a/tests/kubernetes/kubernetes_test.go b/tests/kubernetes/kubernetes_test.go index 816bd8bb3..91e13350b 100644 --- a/tests/kubernetes/kubernetes_test.go +++ b/tests/kubernetes/kubernetes_test.go @@ -358,7 +358,7 @@ func TestCreateJob_DefaultTTLIsSet(t *testing.T) { Resources: &tes.Resources{CpuCores: 1, RamGb: 1.0}, } - if err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { + if _, err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { t.Fatalf("CreateJob: %v", err) } @@ -383,7 +383,7 @@ func TestCreateJob_ExistingTTLIsPreserved(t *testing.T) { Resources: &tes.Resources{CpuCores: 1, RamGb: 1.0}, } - if err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateWithTTL), client, unitTestLog); err != nil { + if _, err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateWithTTL), client, unitTestLog); err != nil { t.Fatalf("CreateJob: %v", err) } @@ -408,7 +408,7 @@ func TestCreateJob_DefaultBackoffLimit(t *testing.T) { Resources: &tes.Resources{CpuCores: 1}, } - if err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { + if _, err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { t.Fatalf("CreateJob: %v", err) } @@ -436,7 +436,7 @@ func TestCreateJob_BackoffLimitFromBackendParameters(t *testing.T) { }, } - if err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { + if _, err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { t.Fatalf("CreateJob: %v", err) } @@ -463,7 +463,7 @@ func TestCreateJob_BackoffLimitInvalidValueFallsBackToDefault(t *testing.T) { }, } - if err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { + if _, err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { t.Fatalf("CreateJob: %v", err) } @@ -490,7 +490,7 @@ func TestCreateJob_NegativeBackoffLimitFallsBackToDefault(t *testing.T) { }, } - if err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { + if _, err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { t.Fatalf("CreateJob: %v", err) } @@ -516,7 +516,7 @@ func TestCreateJob_TaskNameLabelIsSanitized(t *testing.T) { Resources: &tes.Resources{CpuCores: 1}, } - if err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { + if _, err := resources.CreateJob(context.Background(), task, baseJobConfig(jobTemplateNoTTL), client, unitTestLog); err != nil { t.Fatalf("CreateJob: %v", err) } @@ -555,7 +555,7 @@ func TestCreateConfigMap_WithValidTemplate(t *testing.T) { conf.Kubernetes.ConfigMapTemplate = validConfigMapTemplate client := fake.NewSimpleClientset() - if err := resources.CreateConfigMap(context.Background(), unitTestTaskID, conf, client, unitTestLog); err != nil { + if err := resources.CreateConfigMap(context.Background(), unitTestTaskID, conf, client, unitTestLog, nil); err != nil { t.Fatalf("CreateConfigMap: %v", err) } @@ -574,7 +574,7 @@ func TestCreateConfigMap_InvalidTemplateSyntax(t *testing.T) { conf.Kubernetes.JobsNamespace = unitTestJobsNS conf.Kubernetes.ConfigMapTemplate = `{{.Unclosed` - err := resources.CreateConfigMap(context.Background(), unitTestTaskID, conf, fake.NewSimpleClientset(), unitTestLog) + err := resources.CreateConfigMap(context.Background(), unitTestTaskID, conf, fake.NewSimpleClientset(), unitTestLog, nil) if err == nil { t.Error("expected error for malformed template syntax; got nil") } @@ -594,7 +594,7 @@ spec: - name: c image: alpine ` - err := resources.CreateConfigMap(context.Background(), unitTestTaskID, conf, fake.NewSimpleClientset(), unitTestLog) + err := resources.CreateConfigMap(context.Background(), unitTestTaskID, conf, fake.NewSimpleClientset(), unitTestLog, nil) if err == nil { t.Error("expected error when template produces a non-ConfigMap object; got nil") } @@ -689,7 +689,7 @@ func TestCreatePVC_WithoutGenericS3DoesNotPanic(t *testing.T) { conf.Kubernetes.PVCTemplate = pvcTemplateHostPath client := fake.NewSimpleClientset() - if err := resources.CreatePVC(context.Background(), unitTestTaskID, conf, client, unitTestLog); err != nil { + if err := resources.CreatePVC(context.Background(), unitTestTaskID, conf, client, unitTestLog, nil); err != nil { t.Fatalf("CreatePVC without GenericS3: %v", err) } diff --git a/worker/kubernetes.go b/worker/kubernetes.go index c42b0acd6..0d682575e 100644 --- a/worker/kubernetes.go +++ b/worker/kubernetes.go @@ -5,6 +5,7 @@ import ( "context" "fmt" "io" + "strings" "text/template" "time" @@ -25,6 +26,8 @@ type KubernetesCommand struct { TaskId string JobId int StdinFile string + StdoutFile string + StderrFile string TaskTemplate string Namespace string // Funnel Server Namespace JobsNamespace string // Funnel Worker + Executor Namespace (default: Namespace) @@ -84,15 +87,34 @@ func (kcmd KubernetesCommand) Run(ctx context.Context) error { var cmd = kcmd.ShellCommand - if kcmd.StdinFile != "" { - cmd = append(cmd, "<", kcmd.StdinFile) + // When stdio redirects are present, collapse the command into a single + // shell string so the executor template's shell wrapper (which takes only + // index .Command 0) receives the full command including redirects. + hasRedirects := kcmd.StdinFile != "" || kcmd.StdoutFile != "" || kcmd.StderrFile != "" + if hasRedirects { + // Quote each argument to preserve spaces/special characters, then + // append the redirect operators (which must not be quoted). + parts := make([]string, len(cmd)) + for i, arg := range cmd { + parts[i] = strings.ReplaceAll(arg, "'", "'\\''") + parts[i] = "'" + parts[i] + "'" + } + shellCmd := strings.Join(parts, " ") + if kcmd.StdinFile != "" { + shellCmd += " < " + kcmd.StdinFile + } + if kcmd.StdoutFile != "" { + shellCmd += " > " + kcmd.StdoutFile + } + if kcmd.StderrFile != "" { + shellCmd += " 2> " + kcmd.StderrFile + } + cmd = []string{shellCmd} } - // Use a shell wrapper only when the command is a single element (i.e. a - // shell script string). When the caller provides multiple elements the - // array is passed directly as the container command+args so that spaces, - // quotes, and other special characters are preserved without any escaping. - useShell := len(cmd) == 1 + // Use a shell wrapper when the command is a single element (a shell script + // string) or when stdio redirects are present. + useShell := len(cmd) == 1 || hasRedirects templateData := map[string]interface{}{ "TaskId": taskId, diff --git a/worker/worker.go b/worker/worker.go index 9df7c820f..3089d788c 100644 --- a/worker/worker.go +++ b/worker/worker.go @@ -230,6 +230,8 @@ func (r *DefaultWorker) Run(pctx context.Context) (runerr error) { TaskId: task.Id, JobId: i, StdinFile: d.Stdin, + StdoutFile: d.Stdout, + StderrFile: d.Stderr, TaskTemplate: r.Executor.Template, Namespace: r.Executor.Namespace, JobsNamespace: r.Executor.JobsNamespace, @@ -279,7 +281,11 @@ func (r *DefaultWorker) Run(pctx context.Context) (runerr error) { } // Opens stdin/out/err files and updates those fields on "cmd". - if run.ok() || ignoreError { + // Skip for Kubernetes: the executor runs in a separate pod and writes + // stdout/stderr directly via its PVC mount. Creating host files here + // would poison the Mountpoint inode, making the file unreadable by + // the worker's mount instance (EPERM). + if (run.ok() || ignoreError) && r.Executor.Backend != "kubernetes" { run.syserr = r.openStepLogs(mapper, s, d) } @@ -311,24 +317,26 @@ func (r *DefaultWorker) Run(pctx context.Context) (runerr error) { } // Try to fix symlinks broken by docker filesystems. - if run.ok() { + if run.syserr == nil { for _, output := range mapper.Outputs { fixLinks(mapper, output.Path) } } - if run.ok() { + if run.syserr == nil { // Resolve wildcards in the output paths resolveWildcards(mapper) } - if run.ok() && r.Conf.ScratchPath != "" { + if run.syserr == nil && r.Conf.ScratchPath != "" { mapper.CopyOutputsToWorkDir(r.Conf.ScratchPath) } - // Upload outputs + // Upload outputs regardless of executor error — the user needs the output + // files (logs, stderr, etc.) to diagnose failures. Only skip on system errors + // where the worker itself is in a bad state. var outputLog []*tes.OutputFileLog - if run.ok() { + if run.syserr == nil { outputLog, run.syserr = UploadOutputs(ctx, mapper.Outputs, r.Store, event, int(r.Conf.MaxParallelTransfers)) }