From 104b56d7bc76ce7b4fc0144d73c3f0f70e562fc1 Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Wed, 15 Apr 2026 11:09:11 -0700 Subject: [PATCH 01/16] hotfix: Revert ConfigMap creation update Signed-off-by: Liam Beckman --- compute/kubernetes/resources/configmap.go | 52 ++++++----------------- 1 file changed, 13 insertions(+), 39 deletions(-) diff --git a/compute/kubernetes/resources/configmap.go b/compute/kubernetes/resources/configmap.go index c1d00f23b..b102404d6 100644 --- a/compute/kubernetes/resources/configmap.go +++ b/compute/kubernetes/resources/configmap.go @@ -1,18 +1,14 @@ package resources import ( - "bytes" "context" "fmt" - "strings" - "text/template" "github.com/ohsu-comp-bio/funnel/config" "github.com/ohsu-comp-bio/funnel/logger" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" - "k8s.io/client-go/kubernetes/scheme" ) func CreateConfigMap(ctx context.Context, taskId string, conf *config.Config, client kubernetes.Interface, log *logger.Logger) error { @@ -21,41 +17,19 @@ func CreateConfigMap(ctx context.Context, taskId string, conf *config.Config, cl return fmt.Errorf("marshaling config to ConfigMap: %v", err) } - indentFn := func(spaces int, s string) string { - pad := strings.Repeat(" ", spaces) - lines := strings.Split(s, "\n") - for i, line := range lines { - if line != "" { - lines[i] = pad + line - } - } - return strings.Join(lines, "\n") - } - - t, err := template.New(taskId).Funcs(template.FuncMap{"indent": indentFn}).Parse(conf.Kubernetes.ConfigMapTemplate) - if err != nil { - return fmt.Errorf("parsing template: %v", err) - } - - var buf bytes.Buffer - err = t.Execute(&buf, map[string]interface{}{ - "TaskId": taskId, - "Namespace": conf.Kubernetes.JobsNamespace, - "Data": string(configBytes), - }) - if err != nil { - return fmt.Errorf("%v", err) - } - - decode := scheme.Codecs.UniversalDeserializer().Decode - obj, _, err := decode(buf.Bytes(), nil, nil) - if err != nil { - return fmt.Errorf("decoding ConfigMap spec: %v", err) - } - - cm, ok := obj.(*corev1.ConfigMap) - if !ok { - return fmt.Errorf("failed to decode ConfigMap spec") + // Create the ConfigMap that will contain the Funnel Worker Config (`funnel-worker.yaml`) + cm := &corev1.ConfigMap{ + ObjectMeta: metav1.ObjectMeta{ + Name: fmt.Sprintf("funnel-worker-config-%s", taskId), + Namespace: conf.Kubernetes.JobsNamespace, + Labels: map[string]string{ + "app": "funnel", + "taskId": taskId, + }, + }, + Data: map[string]string{ + "funnel-worker.yaml": string(configBytes), + }, } _, err = client.CoreV1().ConfigMaps(conf.Kubernetes.JobsNamespace).Create(ctx, cm, metav1.CreateOptions{}) From 92edb9e8bb80f09ca6fa07d7fec45bffd54a6302 Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Thu, 16 Apr 2026 20:02:24 -0700 Subject: [PATCH 02/16] fix: Skip deletion of ServiceAccounts still in use Signed-off-by: Liam Beckman --- compute/kubernetes/backend.go | 12 +++++++++++- compute/kubernetes/resources/resources_test.go | 2 +- compute/kubernetes/resources/serviceaccount.go | 9 ++++++++- 3 files changed, 20 insertions(+), 3 deletions(-) diff --git a/compute/kubernetes/backend.go b/compute/kubernetes/backend.go index 072ff6cad..a1c2a93cd 100644 --- a/compute/kubernetes/backend.go +++ b/compute/kubernetes/backend.go @@ -266,6 +266,16 @@ 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 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) @@ -303,7 +313,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) diff --git a/compute/kubernetes/resources/resources_test.go b/compute/kubernetes/resources/resources_test.go index 893abe8cd..3d0e2a0d7 100644 --- a/compute/kubernetes/resources/resources_test.go +++ b/compute/kubernetes/resources/resources_test.go @@ -283,7 +283,7 @@ 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) } diff --git a/compute/kubernetes/resources/serviceaccount.go b/compute/kubernetes/resources/serviceaccount.go index 2a0322083..228f0e2c6 100644 --- a/compute/kubernetes/resources/serviceaccount.go +++ b/compute/kubernetes/resources/serviceaccount.go @@ -69,7 +69,14 @@ func CreateServiceAccount(ctx context.Context, task *tes.Task, conf *config.Conf } // 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), From dc02e7bbfc35f3a582bc3a22c7cf7c3c39db592a Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Tue, 21 Apr 2026 11:53:55 -0700 Subject: [PATCH 03/16] feat: Move orphaned resource cleanup to backend startup Signed-off-by: Liam Beckman Co-authored-by: Sai Shanmukha Narumanchi @nss10 --- compute/kubernetes/backend.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/compute/kubernetes/backend.go b/compute/kubernetes/backend.go index a1c2a93cd..81685a5b7 100644 --- a/compute/kubernetes/backend.go +++ b/compute/kubernetes/backend.go @@ -67,6 +67,10 @@ func NewBackend(ctx context.Context, conf *config.Config, reader tes.ReadOnlySer } if !conf.Kubernetes.DisableReconciler { + // 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) + rate := conf.Kubernetes.ReconcileRate.AsDuration() go b.reconcile(ctx, rate, conf.Kubernetes.DisableJobCleanup) } @@ -576,10 +580,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) } } } From 9665964260d5b9014939307041e5678fefbecd1a Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Wed, 22 Apr 2026 13:08:11 -0700 Subject: [PATCH 04/16] fix: Skip ServiceAccount deletion if in use by other pods Signed-off-by: Liam Beckman --- .../kubernetes/resources/resources_test.go | 47 +++++++++++++++++++ .../kubernetes/resources/serviceaccount.go | 23 +++++++++ 2 files changed, 70 insertions(+) diff --git a/compute/kubernetes/resources/resources_test.go b/compute/kubernetes/resources/resources_test.go index 893abe8cd..f59de6da4 100644 --- a/compute/kubernetes/resources/resources_test.go +++ b/compute/kubernetes/resources/resources_test.go @@ -276,6 +276,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{}) @@ -289,6 +293,49 @@ func TestDeleteServiceAccount(t *testing.T) { } } +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) + 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, diff --git a/compute/kubernetes/resources/serviceaccount.go b/compute/kubernetes/resources/serviceaccount.go index 2a0322083..677752c4e 100644 --- a/compute/kubernetes/resources/serviceaccount.go +++ b/compute/kubernetes/resources/serviceaccount.go @@ -68,6 +68,21 @@ func CreateServiceAccount(ctx context.Context, task *tes.Task, conf *config.Conf return nil } +func podsUsingServiceAccount(ctx context.Context, saName, namespace string, client kubernetes.Interface) ([]string, error) { + pods, err := client.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{}) + if err != nil { + return nil, fmt.Errorf("listing pods using ServiceAccount %s: %v", saName, err) + } + + var podNames []string + for _, pod := range pods.Items { + if pod.Spec.ServiceAccountName == saName { + podNames = append(podNames, pod.Name) + } + } + return podNames, 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 { // ServiceAccount names are not available here without config, so we list by label. @@ -78,6 +93,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 { + podNames, err := podsUsingServiceAccount(ctx, sa.Name, namespace, client) + if err != nil { + return err + } + if len(podNames) > 0 { + return fmt.Errorf("serviceAccount %s is still in use by pod(s): %v", sa.Name, podNames) + } + 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) From f96715023383f64a3808eaa2d011b311a5ba64bd Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Wed, 22 Apr 2026 13:08:11 -0700 Subject: [PATCH 05/16] fix: Skip ServiceAccount deletion if in use by other pods Signed-off-by: Liam Beckman --- .../kubernetes/resources/resources_test.go | 47 +++++++++++++++++++ .../kubernetes/resources/serviceaccount.go | 23 +++++++++ 2 files changed, 70 insertions(+) diff --git a/compute/kubernetes/resources/resources_test.go b/compute/kubernetes/resources/resources_test.go index 3d0e2a0d7..1c1d5873c 100644 --- a/compute/kubernetes/resources/resources_test.go +++ b/compute/kubernetes/resources/resources_test.go @@ -276,6 +276,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{}) @@ -289,6 +293,49 @@ func TestDeleteServiceAccount(t *testing.T) { } } +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) + 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, diff --git a/compute/kubernetes/resources/serviceaccount.go b/compute/kubernetes/resources/serviceaccount.go index 228f0e2c6..e93ec670d 100644 --- a/compute/kubernetes/resources/serviceaccount.go +++ b/compute/kubernetes/resources/serviceaccount.go @@ -68,6 +68,21 @@ func CreateServiceAccount(ctx context.Context, task *tes.Task, conf *config.Conf return nil } +func podsUsingServiceAccount(ctx context.Context, saName, namespace string, client kubernetes.Interface) ([]string, error) { + pods, err := client.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{}) + if err != nil { + return nil, fmt.Errorf("listing pods using ServiceAccount %s: %v", saName, err) + } + + var podNames []string + for _, pod := range pods.Items { + if pod.Spec.ServiceAccountName == saName { + podNames = append(podNames, pod.Name) + } + } + return podNames, nil +} + // DeleteServiceAccount deletes the ServiceAccount created for a task. // 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 @@ -85,6 +100,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 { + podNames, err := podsUsingServiceAccount(ctx, sa.Name, namespace, client) + if err != nil { + return err + } + if len(podNames) > 0 { + return fmt.Errorf("serviceAccount %s is still in use by pod(s): %v", sa.Name, podNames) + } + 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) From 2021d2e801e0c9f649e5eba4b7747a3d3ef6b378 Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Wed, 22 Apr 2026 13:08:11 -0700 Subject: [PATCH 06/16] feat: update cleanOrphanedResources to non-blocking goroutine Signed-off-by: Liam Beckman --- compute/kubernetes/backend.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/compute/kubernetes/backend.go b/compute/kubernetes/backend.go index 81685a5b7..4770b3cbe 100644 --- a/compute/kubernetes/backend.go +++ b/compute/kubernetes/backend.go @@ -69,7 +69,7 @@ func NewBackend(ctx context.Context, conf *config.Config, reader tes.ReadOnlySer if !conf.Kubernetes.DisableReconciler { // 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) + go b.cleanOrphanedResources(ctx) rate := conf.Kubernetes.ReconcileRate.AsDuration() go b.reconcile(ctx, rate, conf.Kubernetes.DisableJobCleanup) From 2bc183b41759f1ebc6a0fed3bc751ce2436ba171 Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Mon, 27 Apr 2026 10:41:08 -0700 Subject: [PATCH 07/16] fix: Address PR review comments on ServiceAccount pod-in-use check - Rename podsUsingServiceAccount -> isServiceAccountAttachedToPods, returning bool with early exit - Use FieldSelector to filter pods at the API level instead of fetching all pods and filtering in-process - Fix TestDeleteServiceAccountInUse missing externalSA argument Co-Authored-By: Claude Sonnet 4.6 Assisted-by: Claude Code:claude [Claude Code] --- .../kubernetes/resources/resources_test.go | 2 +- .../kubernetes/resources/serviceaccount.go | 23 ++++++++----------- 2 files changed, 10 insertions(+), 15 deletions(-) diff --git a/compute/kubernetes/resources/resources_test.go b/compute/kubernetes/resources/resources_test.go index 1c1d5873c..3bb8814cb 100644 --- a/compute/kubernetes/resources/resources_test.go +++ b/compute/kubernetes/resources/resources_test.go @@ -325,7 +325,7 @@ func TestDeleteServiceAccountInUse(t *testing.T) { t.Fatalf("Failed to create test Pod: %v", err) } - err = DeleteServiceAccount(context.Background(), testTaskID, namespace, fakeClient, l) + err = DeleteServiceAccount(context.Background(), testTaskID, namespace, fakeClient, l, false) if err == nil { t.Fatal("expected DeleteServiceAccount to fail when ServiceAccount is in use") } diff --git a/compute/kubernetes/resources/serviceaccount.go b/compute/kubernetes/resources/serviceaccount.go index e93ec670d..a606d9466 100644 --- a/compute/kubernetes/resources/serviceaccount.go +++ b/compute/kubernetes/resources/serviceaccount.go @@ -68,19 +68,14 @@ func CreateServiceAccount(ctx context.Context, task *tes.Task, conf *config.Conf return nil } -func podsUsingServiceAccount(ctx context.Context, saName, namespace string, client kubernetes.Interface) ([]string, error) { - pods, err := client.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{}) +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 nil, fmt.Errorf("listing pods using ServiceAccount %s: %v", saName, err) - } - - var podNames []string - for _, pod := range pods.Items { - if pod.Spec.ServiceAccountName == saName { - podNames = append(podNames, pod.Name) - } + return false, fmt.Errorf("listing pods using ServiceAccount %s: %v", saName, err) } - return podNames, nil + return len(pods.Items) > 0, nil } // DeleteServiceAccount deletes the ServiceAccount created for a task. @@ -100,12 +95,12 @@ 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 { - podNames, err := podsUsingServiceAccount(ctx, sa.Name, namespace, client) + inUse, err := isServiceAccountAttachedToPods(ctx, sa.Name, namespace, client) if err != nil { return err } - if len(podNames) > 0 { - return fmt.Errorf("serviceAccount %s is still in use by pod(s): %v", sa.Name, podNames) + 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) From 6f0733986f4869dcea4f3db08b39ecdf2ad16c28 Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Tue, 28 Apr 2026 16:40:07 -0700 Subject: [PATCH 08/16] feat: Add owner references to task resources for automatic K8s GC cleanup Set the worker Job as the owner of all namespaced task resources (ConfigMap, ServiceAccount, Role, RoleBinding, PVC) so Kubernetes garbage-collects them automatically when the Job is deleted. - CreateJob now returns the created Job so callers can read its UID - All Create* functions accept an optional *metav1.OwnerReference - createResources in backend.go creates the Job first, builds the owner ref, and threads it through subsequent resource creation - External (user-managed) SAs via _WORKER_SA tag receive no owner ref since they outlive individual tasks - PVs remain explicitly managed (cluster-scoped resources cannot be owned by a namespaced Job) Co-Authored-By: Claude Sonnet 4.6 Assisted-by: Claude Code:claude [Claude Code] --- compute/kubernetes/backend.go | 101 +++++++++++------- compute/kubernetes/resources/configmap.go | 6 +- compute/kubernetes/resources/job.go | 18 ++-- compute/kubernetes/resources/pvc.go | 6 +- .../kubernetes/resources/resources_test.go | 14 +-- compute/kubernetes/resources/role.go | 6 +- compute/kubernetes/resources/roleBinding.go | 8 +- .../kubernetes/resources/serviceaccount.go | 6 +- 8 files changed, 103 insertions(+), 62 deletions(-) diff --git a/compute/kubernetes/backend.go b/compute/kubernetes/backend.go index 4770b3cbe..51eba7ec9 100644 --- a/compute/kubernetes/backend.go +++ b/compute/kubernetes/backend.go @@ -175,38 +175,29 @@ 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) - - // 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 - 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 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) + } - // 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) - } + blockOwnerDeletion := true + isController := true + ownerRef := &metav1.OwnerReference{ + APIVersion: "batch/v1", + Kind: "Job", + Name: job.Name, + UID: job.UID, + BlockOwnerDeletion: &blockOwnerDeletion, + Controller: &isController, } // Create ConfigMap 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) b.log.Debug("creating Worker ConfigMap", "error", err) @@ -216,10 +207,12 @@ func (b *Backend) createResources(ctx context.Context, task *tes.Task, config *c // Create ServiceAccount: // - This should only be created if no such ServiceAccount with the same name exists // - ServiceAccount will still always need to be added to Worker Job and Executor - saName := "funnel-worker-sa-%s-%s" - saName = fmt.Sprintf(saName, config.Kubernetes.JobsNamespace, task.Id) - if _, exists := task.Tags["_WORKER_SA"]; exists { - saName = task.Tags["_WORKER_SA"] + // - External (user-managed) SAs are not owned by the Job — they outlive individual tasks + saName := fmt.Sprintf("funnel-worker-sa-%s-%s", config.Kubernetes.JobsNamespace, task.Id) + 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` @@ -230,7 +223,12 @@ func (b *Backend) createResources(ctx context.Context, task *tes.Task, config *c 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) @@ -241,7 +239,7 @@ func (b *Backend) createResources(ctx context.Context, task *tes.Task, config *c // Create Role 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) @@ -249,18 +247,37 @@ func (b *Backend) createResources(ctx context.Context, task *tes.Task, config *c // Create RoleBinding 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 @@ -697,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 b102404d6..40cf94de4 100644 --- a/compute/kubernetes/resources/configmap.go +++ b/compute/kubernetes/resources/configmap.go @@ -11,7 +11,7 @@ import ( "k8s.io/client-go/kubernetes" ) -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 { configBytes, err := config.ToYaml(conf) if err != nil { return fmt.Errorf("marshaling config to ConfigMap: %v", err) @@ -32,6 +32,10 @@ func CreateConfigMap(ctx context.Context, taskId string, conf *config.Config, cl }, } + 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 480b9dc24..a9173623f 100644 --- a/compute/kubernetes/resources/job.go +++ b/compute/kubernetes/resources/job.go @@ -18,20 +18,20 @@ import ( // 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) } res := task.GetResources() @@ -61,28 +61,28 @@ 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") } 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/pvc.go b/compute/kubernetes/resources/pvc.go index 81bf12d84..afcbc3ffa 100644 --- a/compute/kubernetes/resources/pvc.go +++ b/compute/kubernetes/resources/pvc.go @@ -16,7 +16,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 @@ -49,6 +49,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) diff --git a/compute/kubernetes/resources/resources_test.go b/compute/kubernetes/resources/resources_test.go index 3bb8814cb..848049ae8 100644 --- a/compute/kubernetes/resources/resources_test.go +++ b/compute/kubernetes/resources/resources_test.go @@ -36,7 +36,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) } @@ -111,7 +111,7 @@ func TestCreateJob(t *testing.T) { } conf := &config.Config{} - 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) } @@ -181,7 +181,7 @@ func TestDeletePV(t *testing.T) { func TestCreatePVC(t *testing.T) { conf := &config.Config{} - 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) } @@ -221,7 +221,7 @@ func TestCreateJobWithNoResources(t *testing.T) { } conf := &config.Config{} - 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) } @@ -263,7 +263,7 @@ func TestCreateServiceAccount(t *testing.T) { } conf := config.DefaultConfig() - 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) } @@ -342,7 +342,7 @@ func TestCreateRole(t *testing.T) { } conf := config.DefaultConfig() - 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) } @@ -374,7 +374,7 @@ func TestCreateRoleBinding(t *testing.T) { } conf := config.DefaultConfig() - 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 9ce180594..07cb9a180 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 7a4be2ed6..c2871539c 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 a606d9466..464b40d45 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 decode 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) From 4a96b7acb4d96fa9466de022f1f21b44dd13a318 Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Tue, 28 Apr 2026 19:22:22 -0700 Subject: [PATCH 09/16] fix: Add support for strong read-write consistency (S3 CSI Driver) Signed-off-by: Liam Beckman --- storage/generic_s3.go | 33 +++++++++++++++++++++++++++------ worker/kubernetes.go | 36 +++++++++++++++++++++++++++++------- worker/worker.go | 8 +++++++- 3 files changed, 63 insertions(+), 14 deletions(-) diff --git a/storage/generic_s3.go b/storage/generic_s3.go index 15724c88d..7762528e2 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" @@ -259,11 +262,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 { @@ -279,7 +301,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 @@ -288,7 +310,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/worker/kubernetes.go b/worker/kubernetes.go index e3ba6703e..8c46552a4 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) @@ -85,15 +88,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 1797426d2..68770830b 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, @@ -278,7 +280,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) } From 10bde5077bb2c0366c5fcc098140bad147b1e487 Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Wed, 29 Apr 2026 12:19:46 -0700 Subject: [PATCH 10/16] feat: Add support for exponential backoff calls to K8s Signed-off-by: Liam Beckman --- compute/kubernetes/resources/pv.go | 54 ++++++++++++++++++++--------- compute/kubernetes/resources/pvc.go | 52 ++++++++++++++++++--------- worker/worker.go | 12 ++++--- 3 files changed, 80 insertions(+), 38 deletions(-) diff --git a/compute/kubernetes/resources/pv.go b/compute/kubernetes/resources/pv.go index 293918dbe..393c1e184 100644 --- a/compute/kubernetes/resources/pv.go +++ b/compute/kubernetes/resources/pv.go @@ -5,11 +5,13 @@ import ( "context" "fmt" "html/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" @@ -56,28 +58,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 afcbc3ffa..b3b94ebbb 100644 --- a/compute/kubernetes/resources/pvc.go +++ b/compute/kubernetes/resources/pvc.go @@ -5,10 +5,12 @@ import ( "context" "fmt" "html/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" @@ -61,28 +63,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/worker/worker.go b/worker/worker.go index 68770830b..d98f44f19 100644 --- a/worker/worker.go +++ b/worker/worker.go @@ -316,24 +316,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)) } From 2bc8c538aa42d99b1d8e829e5281e59774ec73ec Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Wed, 29 Apr 2026 14:02:09 -0700 Subject: [PATCH 11/16] fix: Add missing test fixture for quote handling Signed-off-by: Liam Beckman --- tests/fixtures/quotes/single-element-single-quote.json | 9 +++++++++ 1 file changed, 9 insertions(+) create mode 100644 tests/fixtures/quotes/single-element-single-quote.json 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..08fd1a6a3 --- /dev/null +++ b/tests/fixtures/quotes/single-element-single-quote.json @@ -0,0 +1,9 @@ +{ + "name": "Multi-element exec — single quote in argument", + "executors": [ + { + "image": "alpine", + "command": ["echo Hello O'hare!"] + } + ] +} From 6946c9cf737628916ee6923a2849cde76643bb69 Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Fri, 1 May 2026 10:49:39 -0700 Subject: [PATCH 12/16] feat: Move cleanOrphanedResources to Helm Cron Job (background task) Signed-off-by: Liam Beckman --- cmd/root.go | 2 ++ compute/kubernetes/backend.go | 15 ++++++--------- 2 files changed, 8 insertions(+), 9 deletions(-) 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 51eba7ec9..d9ec63ed2 100644 --- a/compute/kubernetes/backend.go +++ b/compute/kubernetes/backend.go @@ -67,10 +67,6 @@ func NewBackend(ctx context.Context, conf *config.Config, reader tes.ReadOnlySer } if !conf.Kubernetes.DisableReconciler { - // 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. - go b.cleanOrphanedResources(ctx) - rate := conf.Kubernetes.ReconcileRate.AsDuration() go b.reconcile(ctx, rate, conf.Kubernetes.DisableJobCleanup) } @@ -617,12 +613,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{}) From 5a567a253002895756ce34dc3258642fa51a5ccd Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Fri, 1 May 2026 16:58:22 -0700 Subject: [PATCH 13/16] fix: Update K8s tests Signed-off-by: Liam Beckman --- tests/kubernetes/kubernetes_test.go | 22 +++++++++++----------- 1 file changed, 11 insertions(+), 11 deletions(-) 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) } From 206f414bfad516d78c333fe59dbaf3df242bdfbd Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Fri, 1 May 2026 17:00:07 -0700 Subject: [PATCH 14/16] fix: Add kubernetes package for managing Kubernetes resources Signed-off-by: Liam Beckman --- cmd/kubernetes/kubernetes.go | 114 +++++++++++++++++++++++++++++++++++ 1 file changed, 114 insertions(+) create mode 100644 cmd/kubernetes/kubernetes.go 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) + } +} From 957d200a878546a8f8f2ce281ab200d9e0a396e6 Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Fri, 1 May 2026 17:11:40 -0700 Subject: [PATCH 15/16] fix: K8s unit test Signed-off-by: Liam Beckman --- compute/kubernetes/backend.go | 106 ++++++++++++++++++---------------- 1 file changed, 57 insertions(+), 49 deletions(-) diff --git a/compute/kubernetes/backend.go b/compute/kubernetes/backend.go index 38aabaa24..05f199d3a 100644 --- a/compute/kubernetes/backend.go +++ b/compute/kubernetes/backend.go @@ -191,62 +191,68 @@ func (b *Backend) createResources(ctx context.Context, task *tes.Task, config *c Controller: &isController, } - // Create ConfigMap - b.log.Debug("creating Worker ConfigMap", "taskID", task.Id) - err = resources.CreateConfigMap(timeoutCtx, task.Id, config, b.client, b.log, ownerRef) - if err != nil { - _ = b.Cancel(context.Background(), task.Id) - b.log.Debug("creating Worker ConfigMap", "error", err) - return fmt.Errorf("creating Worker ConfigMap: %w", err) + // 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, ownerRef) + if err != nil { + _ = b.Cancel(context.Background(), task.Id) + return fmt.Errorf("creating Worker ConfigMap: %w", err) + } } - // Create ServiceAccount: - // - This should only be created if no such ServiceAccount with the same name exists - // - ServiceAccount will still always need to be added to Worker Job and Executor - // - External (user-managed) SAs are not owned by the Job — they outlive individual tasks - saName := fmt.Sprintf("funnel-worker-sa-%s-%s", config.Kubernetes.JobsNamespace, task.Id) - externalSA := false - if sa, exists := task.Tags["_WORKER_SA"]; exists && sa != "" { - saName = sa - externalSA = true - } + // 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) + 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{}) + // 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) - // 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) + // ServiceAccount does not exist, create it if err != nil { - _ = b.Cancel(context.Background(), task.Id) - return fmt.Errorf("creating Worker ServiceAccount: %w", err) + b.log.Debug("Error getting ServiceAccount:", "ServiceAccount", saName, "taskID", task.Id, "error", err) + b.log.Debug("Creating Worker ServiceAccount", "taskID", task.Id) + // 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) + } + } else { + b.log.Debug("ServiceAccount already exists, skipping creation", "ServiceAccount", saName, "taskID", task.Id) } - } else { - b.log.Debug("ServiceAccount already exists, skipping creation", "ServiceAccount", saName, "taskID", task.Id) } - // Create Role - b.log.Debug("creating Worker Role", "taskID", task.Id) - 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) + if config.Kubernetes.RoleTemplate != "" { + b.log.Debug("creating Worker Role", "taskID", task.Id) + 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) + } } - // Create RoleBinding - b.log.Debug("creating Worker RoleBinding", "taskID", task.Id) - 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) + if config.Kubernetes.RoleBindingTemplate != "" { + b.log.Debug("creating Worker RoleBinding", "taskID", task.Id) + 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) + } } // If the task has inputs, outputs, or declared volumes, create a PVC so @@ -287,9 +293,11 @@ func (b *Backend) cleanResources(ctx context.Context, taskId string) error { // 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 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 + 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 + } } } From ec64e00bcef1244de24038d1b34c1db17cbb4e86 Mon Sep 17 00:00:00 2001 From: Liam Beckman Date: Mon, 4 May 2026 10:47:16 -0700 Subject: [PATCH 16/16] fix: Single element quote test (TODO: check if OK) Signed-off-by: Liam Beckman --- tests/fixtures/quotes/single-element-single-quote.json | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/fixtures/quotes/single-element-single-quote.json b/tests/fixtures/quotes/single-element-single-quote.json index 08fd1a6a3..7aa6ec9c8 100644 --- a/tests/fixtures/quotes/single-element-single-quote.json +++ b/tests/fixtures/quotes/single-element-single-quote.json @@ -1,9 +1,9 @@ { - "name": "Multi-element exec — single quote in argument", + "name": "Single-element exec — single quote in argument", "executors": [ { "image": "alpine", - "command": ["echo Hello O'hare!"] + "command": ["echo \"Hello O'hare!\""] } ] }