From 82ecff062c5de57f5a16d1c68408b617e787ce80 Mon Sep 17 00:00:00 2001 From: amh1k Date: Thu, 24 Sep 2026 10:24:09 +0500 Subject: [PATCH 1/3] feat(coordinator): add Job status polling and owned Pod discovery Signed-off-by: amh1k --- backend/internal/coordinator/jobs.go | 119 +++++++++++- backend/internal/coordinator/jobs_test.go | 212 ++++++++++++++++++++++ 2 files changed, 330 insertions(+), 1 deletion(-) diff --git a/backend/internal/coordinator/jobs.go b/backend/internal/coordinator/jobs.go index b9fa415..f89acd1 100644 --- a/backend/internal/coordinator/jobs.go +++ b/backend/internal/coordinator/jobs.go @@ -14,6 +14,7 @@ import ( corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/types" ) @@ -21,6 +22,7 @@ const ( runnerContainerName = "pgbench-runner" benchmarkJobNamePrefix = "plugin-bench-pgbench-" credentialCleanupTimeout = 10 * time.Second + jobPollInterval = time.Second ) // JobRef identifies the Job created for one benchmark run. @@ -29,11 +31,127 @@ type JobRef struct { UID types.UID } +// waitForJob polls a created Job until Kubernetes reports completion or +// failure. Log collection and resource cleanup happen in later lifecycle +// steps. +func (c *Coordinator) waitForJob(ctx context.Context, reference JobRef) error { + return c.waitForJobAtInterval(ctx, reference, jobPollInterval) +} + +func (c *Coordinator) waitForJobAtInterval(ctx context.Context, reference JobRef, interval time.Duration) error { + if c == nil { + return errors.New("coordinator is nil") + } + if ctx == nil { + return errors.New("job polling context is nil") + } + if c.kubeClient == nil { + return ErrKubernetesClientRequired + } + if err := c.config.Validate(); err != nil { + return fmt.Errorf("invalid coordinator configuration: %w", err) + } + if strings.TrimSpace(reference.Name) == "" { + return errors.New("benchmark Job name is required") + } + if interval <= 0 { + return errors.New("job polling interval must be positive") + } + + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + if err := ctx.Err(); err != nil { + return err + } + + job, err := c.kubeClient.BatchV1().Jobs(c.config.WorkloadNamespace).Get(ctx, reference.Name, metav1.GetOptions{}) + if err != nil { + if ctxErr := ctx.Err(); ctxErr != nil { + return ctxErr + } + return fmt.Errorf("get benchmark Job %q: %w", reference.Name, err) + } + if reference.UID != "" && job.UID != "" && job.UID != reference.UID { + return fmt.Errorf("benchmark Job %q UID does not match the created Job", reference.Name) + } + + for _, condition := range job.Status.Conditions { + if condition.Status != corev1.ConditionTrue { + continue + } + switch condition.Type { + case batchv1.JobComplete: + return nil + case batchv1.JobFailed: + return fmt.Errorf("benchmark Job %q failed", reference.Name) + } + } + + select { + case <-ctx.Done(): + return ctx.Err() + case <-ticker.C: + } + } +} + type benchmarkResources struct { secret SecretRef job JobRef } +// findJobPod uses the Kubernetes Job label to narrow the search, then checks +// controller ownership so a same-named Job from another run cannot match. +// Multiple owned Pods are an error: choosing one arbitrarily could hide output. +func (c *Coordinator) findJobPod(ctx context.Context, reference JobRef) (*corev1.Pod, error) { + if c == nil { + return nil, errors.New("coordinator is nil") + } + if ctx == nil { + return nil, errors.New("Pod discovery context is nil") + } + if c.kubeClient == nil { + return nil, ErrKubernetesClientRequired + } + if err := c.config.Validate(); err != nil { + return nil, fmt.Errorf("invalid coordinator configuration: %w", err) + } + if strings.TrimSpace(reference.Name) == "" || reference.UID == "" { + return nil, errors.New("benchmark Job name and UID are required") + } + if err := ctx.Err(); err != nil { + return nil, err + } + + pods, err := c.kubeClient.CoreV1().Pods(c.config.WorkloadNamespace).List(ctx, metav1.ListOptions{ + LabelSelector: labels.Set{batchv1.JobNameLabel: reference.Name}.String(), + }) + if ctxErr := ctx.Err(); ctxErr != nil { + return nil, ctxErr + } + if err != nil { + return nil, fmt.Errorf("list Pods for benchmark Job %q: %w", reference.Name, err) + } + var found *corev1.Pod + for i := range pods.Items { + pod := &pods.Items[i] + owner := metav1.GetControllerOf(pod) + if owner == nil || owner.APIVersion != batchv1.SchemeGroupVersion.String() || + owner.Kind != "Job" || owner.Name != reference.Name || owner.UID != reference.UID { + continue + } + if found != nil { + return nil, fmt.Errorf("multiple Pods found for benchmark Job %q", reference.Name) + } + found = pod + } + if found == nil { + return nil, fmt.Errorf("no Pod found for benchmark Job %q", reference.Name) + } + return found, nil +} + // createBenchmarkResources creates the Secret before the Job that consumes // it. If Job creation fails, the Secret is removed before the error returns. // Monitoring and normal run cleanup are handled by the execution lifecycle. @@ -71,7 +189,6 @@ func (c *Coordinator) createBenchmarkResources(ctx context.Context, connection C return benchmarkResources{secret: secret, job: job}, nil } - // Cleanup must survive caller cancellation but still have a deadline. cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), credentialCleanupTimeout) defer cancel() cleanupErr := c.deleteCredentialSecret(cleanupCtx, secret) diff --git a/backend/internal/coordinator/jobs_test.go b/backend/internal/coordinator/jobs_test.go index 55e5c46..bb34ee9 100644 --- a/backend/internal/coordinator/jobs_test.go +++ b/backend/internal/coordinator/jobs_test.go @@ -149,6 +149,218 @@ func TestSecretCreationFailurePreventsJobCreation(t *testing.T) { } } +func TestWaitForJobReturnsWhenComplete(t *testing.T) { + job := &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{Name: "job-1", Namespace: "plugin-bench-workloads", UID: "job-uid"}, + Status: batchv1.JobStatus{Conditions: []batchv1.JobCondition{{ + Type: batchv1.JobComplete, + Status: corev1.ConditionTrue, + }}}, + } + client := fake.NewSimpleClientset(job) + coordinator, err := New(validConfig(), client) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + if err := coordinator.waitForJobAtInterval(context.Background(), JobRef{Name: job.Name, UID: job.UID}, time.Millisecond); err != nil { + t.Fatalf("waitForJob() error = %v", err) + } +} + +func TestWaitForJobReturnsFailure(t *testing.T) { + job := &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{Name: "job-1", Namespace: "plugin-bench-workloads", UID: "job-uid"}, + Status: batchv1.JobStatus{Conditions: []batchv1.JobCondition{{ + Type: batchv1.JobFailed, + Status: corev1.ConditionTrue, + }}}, + } + client := fake.NewSimpleClientset(job) + coordinator, err := New(validConfig(), client) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + err = coordinator.waitForJobAtInterval(context.Background(), JobRef{Name: job.Name, UID: job.UID}, time.Millisecond) + if err == nil || !strings.Contains(err.Error(), `benchmark Job "job-1" failed`) { + t.Fatalf("error = %v, want Job failure", err) + } +} + +func TestWaitForJobPollsUntilTerminalState(t *testing.T) { + client := fake.NewSimpleClientset() + getCount := 0 + client.PrependReactor("get", "jobs", func(ktesting.Action) (bool, runtime.Object, error) { + getCount++ + job := &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Name: "job-1", Namespace: "plugin-bench-workloads", UID: "job-uid"}} + if getCount == 1 { + job.Status.Conditions = []batchv1.JobCondition{{Type: batchv1.JobComplete, Status: corev1.ConditionFalse}} + } + if getCount == 2 { + job.Status.Conditions = []batchv1.JobCondition{{Type: batchv1.JobComplete, Status: corev1.ConditionTrue}} + } + return true, job, nil + }) + coordinator, err := New(validConfig(), client) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + if err := coordinator.waitForJobAtInterval(context.Background(), JobRef{Name: "job-1", UID: "job-uid"}, time.Millisecond); err != nil { + t.Fatalf("waitForJob() error = %v", err) + } + if getCount != 2 { + t.Fatalf("Job was fetched %d times, want 2", getCount) + } +} + +func TestWaitForJobPreservesCancellation(t *testing.T) { + client := fake.NewSimpleClientset() + ctx, cancel := context.WithCancel(context.Background()) + client.PrependReactor("get", "jobs", func(ktesting.Action) (bool, runtime.Object, error) { + cancel() + return true, &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Name: "job-1", Namespace: "plugin-bench-workloads", UID: "job-uid"}}, nil + }) + coordinator, err := New(validConfig(), client) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + err = coordinator.waitForJobAtInterval(ctx, JobRef{Name: "job-1", UID: "job-uid"}, time.Millisecond) + if !errors.Is(err, context.Canceled) { + t.Fatalf("error = %v, want context.Canceled", err) + } +} + +func TestWaitForJobRejectsUIDMismatch(t *testing.T) { + client := fake.NewSimpleClientset(&batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{Name: "job-1", Namespace: "plugin-bench-workloads", UID: "different-uid"}, + }) + coordinator, err := New(validConfig(), client) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + err = coordinator.waitForJobAtInterval(context.Background(), JobRef{Name: "job-1", UID: "created-uid"}, time.Millisecond) + if err == nil || !strings.Contains(err.Error(), "UID does not match") { + t.Fatalf("error = %v, want UID mismatch", err) + } +} + +func TestFindJobPod(t *testing.T) { + ref := JobRef{Name: "benchmark", UID: "job-uid"} + newPod := func(name string) *corev1.Pod { + return &corev1.Pod{ObjectMeta: metav1.ObjectMeta{ + Name: name, Namespace: validConfig().WorkloadNamespace, + Labels: map[string]string{batchv1.JobNameLabel: ref.Name}, + OwnerReferences: []metav1.OwnerReference{{ + APIVersion: "batch/v1", Kind: "Job", Name: ref.Name, + UID: ref.UID, Controller: boolPtr(true), + }}, + }} + } + for _, test := range []struct { + name string + mutate func(*corev1.Pod) + wantMatch bool + }{ + {"matching owner", func(*corev1.Pod) {}, true}, + {"old Job UID", func(p *corev1.Pod) { p.OwnerReferences[0].UID = "old" }, false}, + {"no owner", func(p *corev1.Pod) { p.OwnerReferences = nil }, false}, + {"not controller", func(p *corev1.Pod) { p.OwnerReferences[0].Controller = boolPtr(false) }, false}, + {"wrong kind", func(p *corev1.Pod) { p.OwnerReferences[0].Kind = "ReplicaSet" }, false}, + {"wrong owner name", func(p *corev1.Pod) { p.OwnerReferences[0].Name = "other" }, false}, + {"wrong API", func(p *corev1.Pod) { p.OwnerReferences[0].APIVersion = "apps/v1" }, false}, + {"wrong label", func(p *corev1.Pod) { p.Labels[batchv1.JobNameLabel] = "other" }, false}, + {"wrong namespace", func(p *corev1.Pod) { p.Namespace = "other" }, false}, + } { + t.Run(test.name, func(t *testing.T) { + pod := newPod("runner") + test.mutate(pod) + client := fake.NewSimpleClientset(pod) + c, err := New(validConfig(), client) + if err != nil { + t.Fatal(err) + } + got, err := c.findJobPod(context.Background(), ref) + if test.wantMatch { + if err != nil || got == nil || got.Name != pod.Name { + t.Fatalf("Pod = %v, error = %v", got, err) + } + } else if err == nil || got != nil || !strings.Contains(err.Error(), "no Pod found") { + t.Fatalf("expected no matching Pod, got %v, %v", got, err) + } + }) + } + t.Run("multiple matches", func(t *testing.T) { + c, err := New(validConfig(), fake.NewSimpleClientset(newPod("one"), newPod("two"))) + if err != nil { + t.Fatal(err) + } + got, err := c.findJobPod(context.Background(), ref) + if got != nil || err == nil || !strings.Contains(err.Error(), "multiple Pods") { + t.Fatalf("expected ambiguity error, got %v, %v", got, err) + } + }) + t.Run("empty list", func(t *testing.T) { + c, err := New(validConfig(), fake.NewSimpleClientset()) + if err != nil { + t.Fatal(err) + } + got, err := c.findJobPod(context.Background(), ref) + if got != nil || err == nil || !strings.Contains(err.Error(), "no Pod found") { + t.Fatalf("expected missing Pod error, got %v, %v", got, err) + } + }) +} + +func TestFindJobPodErrors(t *testing.T) { + for _, test := range []string{"API failure", "cancellation during list", "already canceled", "expired deadline", "missing UID"} { + t.Run(test, func(t *testing.T) { + client := fake.NewSimpleClientset() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + ref := JobRef{Name: "benchmark", UID: "job-uid"} + cause := errors.New("API unavailable") + switch test { + case "already canceled": + cancel() + cause = context.Canceled + case "expired deadline": + var stop context.CancelFunc + ctx, stop = context.WithDeadline(context.Background(), time.Now().Add(-time.Second)) + defer stop() + cause = context.DeadlineExceeded + case "missing UID": + ref.UID = "" + case "cancellation during list": + cause = context.Canceled + } + client.PrependReactor("list", "pods", func(ktesting.Action) (bool, runtime.Object, error) { + if test == "cancellation during list" { + cancel() + } + return true, nil, cause + }) + c, err := New(validConfig(), client) + if err != nil { + t.Fatal(err) + } + pod, err := c.findJobPod(ctx, ref) + if pod != nil || err == nil { + t.Fatal("expected error and no Pod") + } + if test != "missing UID" && !errors.Is(err, cause) { + t.Fatalf("error = %v, want %v", err, cause) + } + if (test == "already canceled" || test == "expired deadline" || test == "missing UID") && len(client.Actions()) != 0 { + t.Fatal("invalid input or expired context must prevent API calls") + } + }) + } +} + func TestCreateBenchmarkJob(t *testing.T) { client := fake.NewSimpleClientset() client.PrependReactor("create", "jobs", func(action ktesting.Action) (bool, runtime.Object, error) { From a15c98995aced7b007388b0f6000a7c34c30b118 Mon Sep 17 00:00:00 2001 From: amh1k Date: Thu, 24 Sep 2026 11:56:22 +0500 Subject: [PATCH 2/3] feat(coordinator): collect runner logs and clean up benchmark resources Signed-off-by: amh1k --- backend/internal/coordinator/jobs.go | 155 ++++++++++++++++++- backend/internal/coordinator/jobs_test.go | 179 +++++++++++++++++++++- 2 files changed, 326 insertions(+), 8 deletions(-) diff --git a/backend/internal/coordinator/jobs.go b/backend/internal/coordinator/jobs.go index f89acd1..4e4a996 100644 --- a/backend/internal/coordinator/jobs.go +++ b/backend/internal/coordinator/jobs.go @@ -6,12 +6,14 @@ import ( "encoding/hex" "errors" "fmt" + "io" "strconv" "strings" "time" batchv1 "k8s.io/api/batch/v1" corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" @@ -19,10 +21,11 @@ import ( ) const ( - runnerContainerName = "pgbench-runner" - benchmarkJobNamePrefix = "plugin-bench-pgbench-" - credentialCleanupTimeout = 10 * time.Second - jobPollInterval = time.Second + runnerContainerName = "pgbench-runner" + benchmarkJobNamePrefix = "plugin-bench-pgbench-" + resourceCleanupTimeout = 10 * time.Second + jobPollInterval = time.Second + maxRunnerLogBytes = 1 << 20 ) // JobRef identifies the Job created for one benchmark run. @@ -152,6 +155,77 @@ func (c *Coordinator) findJobPod(ctx context.Context, reference JobRef) (*corev1 return found, nil } +type podLogReader func(context.Context, string, *corev1.PodLogOptions) (io.ReadCloser, error) + +// collectRunnerLogs reads the runner container's combined Pod log stream and +// limits the returned output to maxRunnerLogBytes. +func (c *Coordinator) collectRunnerLogs(ctx context.Context, pod *corev1.Pod) (string, bool, error) { + if c == nil { + return "", false, errors.New("coordinator is nil") + } + if ctx == nil { + return "", false, errors.New("log collection context is nil") + } + if c.kubeClient == nil { + return "", false, ErrKubernetesClientRequired + } + if err := c.config.Validate(); err != nil { + return "", false, fmt.Errorf("invalid coordinator configuration: %w", err) + } + if pod == nil || strings.TrimSpace(pod.Name) == "" { + return "", false, errors.New("Pod name is required") + } + + reader := func(ctx context.Context, name string, options *corev1.PodLogOptions) (io.ReadCloser, error) { + return c.kubeClient.CoreV1().Pods(c.config.WorkloadNamespace).GetLogs(name, options).Stream(ctx) + } + return collectRunnerLogsWithReader(ctx, pod, reader) +} + +func collectRunnerLogsWithReader(ctx context.Context, pod *corev1.Pod, readLogs podLogReader) (string, bool, error) { + if ctx == nil { + return "", false, errors.New("log collection context is nil") + } + if pod == nil || strings.TrimSpace(pod.Name) == "" { + return "", false, errors.New("Pod name is required") + } + if readLogs == nil { + return "", false, errors.New("Pod log reader is required") + } + if err := ctx.Err(); err != nil { + return "", false, err + } + + stream, err := readLogs(ctx, pod.Name, &corev1.PodLogOptions{Container: runnerContainerName}) + if err != nil { + if ctxErr := ctx.Err(); ctxErr != nil { + return "", false, ctxErr + } + return "", false, fmt.Errorf("open logs for Pod %q: %w", pod.Name, err) + } + if stream == nil { + return "", false, fmt.Errorf("open logs for Pod %q: empty log stream", pod.Name) + } + + data, readErr := io.ReadAll(io.LimitReader(stream, maxRunnerLogBytes+1)) + closeErr := stream.Close() + if readErr != nil { + if ctxErr := ctx.Err(); ctxErr != nil { + return "", false, ctxErr + } + return "", false, fmt.Errorf("read logs for Pod %q: %w", pod.Name, readErr) + } + if closeErr != nil { + return "", false, fmt.Errorf("close logs for Pod %q: %w", pod.Name, closeErr) + } + + truncated := int64(len(data)) > maxRunnerLogBytes + if truncated { + data = data[:maxRunnerLogBytes] + } + return string(data), truncated, nil +} + // createBenchmarkResources creates the Secret before the Job that consumes // it. If Job creation fails, the Secret is removed before the error returns. // Monitoring and normal run cleanup are handled by the execution lifecycle. @@ -189,15 +263,82 @@ func (c *Coordinator) createBenchmarkResources(ctx context.Context, connection C return benchmarkResources{secret: secret, job: job}, nil } - cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), credentialCleanupTimeout) - defer cancel() - cleanupErr := c.deleteCredentialSecret(cleanupCtx, secret) + cleanupErr := c.cleanupBenchmarkResources(ctx, benchmarkResources{secret: secret}) if cleanupErr != nil { return resources, errors.Join(err, fmt.Errorf("cleanup credential Secret after Job creation failure: %w", cleanupErr)) } return resources, err } +// deleteBenchmarkJob removes a Job created for a run. The UID precondition +// prevents deleting a different Job that was recreated with the same name. +// Foreground propagation asks Kubernetes to remove the owned Pod first. +func (c *Coordinator) deleteBenchmarkJob(ctx context.Context, reference JobRef) error { + if c == nil { + return errors.New("coordinator is nil") + } + if ctx == nil { + return errors.New("Job deletion context is nil") + } + if c.kubeClient == nil { + return ErrKubernetesClientRequired + } + if err := c.config.Validate(); err != nil { + return fmt.Errorf("invalid coordinator configuration: %w", err) + } + if strings.TrimSpace(reference.Name) == "" { + return errors.New("benchmark Job name is required") + } + if reference.UID == "" { + return errors.New("benchmark Job UID is required") + } + if err := ctx.Err(); err != nil { + return err + } + + propagation := metav1.DeletePropagationForeground + err := c.kubeClient.BatchV1().Jobs(c.config.WorkloadNamespace).Delete(ctx, reference.Name, metav1.DeleteOptions{ + Preconditions: &metav1.Preconditions{UID: &reference.UID}, + PropagationPolicy: &propagation, + }) + if apierrors.IsNotFound(err) { + return nil + } + if err != nil { + return fmt.Errorf("delete benchmark Job: %w", err) + } + return nil +} + +// cleanupBenchmarkResources removes the resources created for one run using +// a context that survives caller cancellation but has a bounded lifetime. +func (c *Coordinator) cleanupBenchmarkResources(ctx context.Context, resources benchmarkResources) error { + if c == nil { + return errors.New("coordinator is nil") + } + if ctx == nil { + return errors.New("resource cleanup context is nil") + } + cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), resourceCleanupTimeout) + defer cancel() + return c.cleanupBenchmarkResourcesWithContext(cleanupCtx, resources) +} + +func (c *Coordinator) cleanupBenchmarkResourcesWithContext(ctx context.Context, resources benchmarkResources) error { + var cleanupErrors []error + if resources.job.Name != "" { + if err := c.deleteBenchmarkJob(ctx, resources.job); err != nil { + cleanupErrors = append(cleanupErrors, fmt.Errorf("cleanup benchmark Job: %w", err)) + } + } + if resources.secret.Name != "" { + if err := c.deleteCredentialSecret(ctx, resources.secret); err != nil { + cleanupErrors = append(cleanupErrors, fmt.Errorf("cleanup credential Secret: %w", err)) + } + } + return errors.Join(cleanupErrors...) +} + // createBenchmarkJob builds and creates one runner Job in the configured // workload namespace. It does not wait for completion or read logs. func (c *Coordinator) createBenchmarkJob(ctx context.Context, secretRef SecretRef, options Options, labels map[string]string) (JobRef, error) { diff --git a/backend/internal/coordinator/jobs_test.go b/backend/internal/coordinator/jobs_test.go index bb34ee9..c6d094f 100644 --- a/backend/internal/coordinator/jobs_test.go +++ b/backend/internal/coordinator/jobs_test.go @@ -3,6 +3,7 @@ package coordinator import ( "context" "errors" + "io" "strings" "testing" "time" @@ -43,6 +44,20 @@ type cleanupContextSecrets struct { check func(context.Context) } +type failingLogStream struct { + err error + closed bool +} + +func (s *failingLogStream) Read([]byte) (int, error) { + return 0, s.err +} + +func (s *failingLogStream) Close() error { + s.closed = true + return nil +} + func (s cleanupContextSecrets) Delete(ctx context.Context, name string, options metav1.DeleteOptions) error { s.check(ctx) return s.SecretInterface.Delete(ctx, name, options) @@ -64,7 +79,7 @@ func TestResourceCleanupSurvivesCancellationWithDeadline(t *testing.T) { t.Fatal("cleanup inherited caller cancellation") } deadline, ok := cleanupCtx.Deadline() - if !ok || time.Until(deadline) <= 0 || time.Until(deadline) > credentialCleanupTimeout { + if !ok || time.Until(deadline) <= 0 || time.Until(deadline) > resourceCleanupTimeout { t.Fatal("cleanup must have a bounded future deadline") } }}) @@ -248,6 +263,99 @@ func TestWaitForJobRejectsUIDMismatch(t *testing.T) { } } +func TestDeleteBenchmarkJobUsesUIDAndForegroundPropagation(t *testing.T) { + client := fake.NewSimpleClientset(&batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{Name: "job-1", Namespace: validConfig().WorkloadNamespace, UID: "job-uid"}, + }) + client.PrependReactor("delete", "jobs", func(action ktesting.Action) (bool, runtime.Object, error) { + deleteAction := action.(ktesting.DeleteAction) + options := deleteAction.GetDeleteOptions() + if options.Preconditions == nil || options.Preconditions.UID == nil || *options.Preconditions.UID != "job-uid" { + t.Fatal("Job deletion must specify the recorded UID") + } + if options.PropagationPolicy == nil || *options.PropagationPolicy != metav1.DeletePropagationForeground { + t.Fatal("Job deletion must use foreground propagation") + } + return false, nil, nil + }) + coordinator, err := New(validConfig(), client) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + if err := coordinator.deleteBenchmarkJob(context.Background(), JobRef{Name: "job-1", UID: "job-uid"}); err != nil { + t.Fatalf("deleteBenchmarkJob() error = %v", err) + } + if _, err := client.BatchV1().Jobs(validConfig().WorkloadNamespace).Get(context.Background(), "job-1", metav1.GetOptions{}); err == nil { + t.Fatal("Job still exists after deletion") + } +} + +func TestDeleteBenchmarkJobTreatsMissingJobAsSuccess(t *testing.T) { + coordinator, err := New(validConfig(), fake.NewSimpleClientset()) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + if err := coordinator.deleteBenchmarkJob(context.Background(), JobRef{Name: "job-1", UID: "job-uid"}); err != nil { + t.Fatalf("missing Job should be treated as cleaned up: %v", err) + } +} + +func TestCleanupBenchmarkResourcesDeletesJobBeforeSecret(t *testing.T) { + client := fake.NewSimpleClientset( + &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Name: "job-1", Namespace: validConfig().WorkloadNamespace, UID: "job-uid"}}, + &corev1.Secret{ObjectMeta: metav1.ObjectMeta{Name: "secret-1", Namespace: validConfig().WorkloadNamespace, UID: "secret-uid"}}, + ) + var actions []string + client.PrependReactor("delete", "jobs", func(ktesting.Action) (bool, runtime.Object, error) { + actions = append(actions, "job") + return false, nil, nil + }) + client.PrependReactor("delete", "secrets", func(ktesting.Action) (bool, runtime.Object, error) { + actions = append(actions, "secret") + return false, nil, nil + }) + coordinator, err := New(validConfig(), client) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + err = coordinator.cleanupBenchmarkResources(context.Background(), benchmarkResources{ + job: JobRef{Name: "job-1", UID: "job-uid"}, + secret: SecretRef{Name: "secret-1", UID: "secret-uid"}, + }) + if err != nil { + t.Fatalf("cleanupBenchmarkResources() error = %v", err) + } + if strings.Join(actions, ",") != "job,secret" { + t.Fatalf("cleanup order = %v, want [job secret]", actions) + } +} + +func TestCleanupBenchmarkResourcesPreservesDeleteErrors(t *testing.T) { + client := fake.NewSimpleClientset() + jobErr, secretErr := errors.New("Job deletion failed"), errors.New("Secret deletion failed") + client.PrependReactor("delete", "jobs", func(ktesting.Action) (bool, runtime.Object, error) { + return true, nil, jobErr + }) + client.PrependReactor("delete", "secrets", func(ktesting.Action) (bool, runtime.Object, error) { + return true, nil, secretErr + }) + coordinator, err := New(validConfig(), client) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + err = coordinator.cleanupBenchmarkResources(context.Background(), benchmarkResources{ + job: JobRef{Name: "job-1", UID: "job-uid"}, + secret: SecretRef{Name: "secret-1", UID: "secret-uid"}, + }) + if !errors.Is(err, jobErr) || !errors.Is(err, secretErr) { + t.Fatalf("error = %v, want both cleanup errors", err) + } +} + func TestFindJobPod(t *testing.T) { ref := JobRef{Name: "benchmark", UID: "job-uid"} newPod := func(name string) *corev1.Pod { @@ -361,6 +469,75 @@ func TestFindJobPodErrors(t *testing.T) { } } +func TestCollectRunnerLogs(t *testing.T) { + pod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "runner-pod"}} + var requestedName string + var requestedContainer string + + output, truncated, err := collectRunnerLogsWithReader(context.Background(), pod, func(_ context.Context, name string, options *corev1.PodLogOptions) (io.ReadCloser, error) { + requestedName = name + requestedContainer = options.Container + return io.NopCloser(strings.NewReader("transaction output\n")), nil + }) + if err != nil || truncated || output != "transaction output\n" { + t.Fatalf("output = %q, truncated = %v, error = %v", output, truncated, err) + } + if requestedName != pod.Name || requestedContainer != runnerContainerName { + t.Fatalf("requested Pod/container = %q/%q", requestedName, requestedContainer) + } +} + +func TestCollectRunnerLogsUsesKubernetesClient(t *testing.T) { + coordinator, err := New(validConfig(), fake.NewSimpleClientset()) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + output, truncated, err := coordinator.collectRunnerLogs(context.Background(), &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "runner-pod"}}) + if err != nil || truncated || output != "fake logs" { + t.Fatalf("output = %q, truncated = %v, error = %v", output, truncated, err) + } +} + +func TestCollectRunnerLogsTruncatesOutput(t *testing.T) { + pod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "runner-pod"}} + largeOutput := strings.Repeat("x", int(maxRunnerLogBytes)+1) + + output, truncated, err := collectRunnerLogsWithReader(context.Background(), pod, func(context.Context, string, *corev1.PodLogOptions) (io.ReadCloser, error) { + return io.NopCloser(strings.NewReader(largeOutput)), nil + }) + if err != nil { + t.Fatalf("collectRunnerLogsWithReader() error = %v", err) + } + if !truncated || len(output) != int(maxRunnerLogBytes) || output != largeOutput[:maxRunnerLogBytes] { + t.Fatalf("output length = %d, truncated = %v", len(output), truncated) + } +} + +func TestCollectRunnerLogsClosesStreamAfterReadError(t *testing.T) { + pod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "runner-pod"}} + stream := &failingLogStream{err: errors.New("log stream failed")} + output, truncated, err := collectRunnerLogsWithReader(context.Background(), pod, func(context.Context, string, *corev1.PodLogOptions) (io.ReadCloser, error) { + return stream, nil + }) + if err == nil || !strings.Contains(err.Error(), "log stream failed") || output != "" || truncated || !stream.closed { + t.Fatalf("output = %q, truncated = %v, closed = %v, error = %v", output, truncated, stream.closed, err) + } +} + +func TestCollectRunnerLogsPreservesCancellation(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + called := false + output, truncated, err := collectRunnerLogsWithReader(ctx, &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "runner-pod"}}, func(context.Context, string, *corev1.PodLogOptions) (io.ReadCloser, error) { + called = true + return nil, nil + }) + if !errors.Is(err, context.Canceled) || output != "" || truncated || called { + t.Fatalf("output = %q, truncated = %v, called = %v, error = %v", output, truncated, called, err) + } +} + func TestCreateBenchmarkJob(t *testing.T) { client := fake.NewSimpleClientset() client.PrependReactor("create", "jobs", func(action ktesting.Action) (bool, runtime.Object, error) { From 06854c5a584b019f8607ebb30f9070a1a7d067c0 Mon Sep 17 00:00:00 2001 From: amh1k Date: Thu, 24 Sep 2026 13:41:49 +0500 Subject: [PATCH 3/3] feat(coordinator): integrate benchmark runs with Jobs and Secrets Signed-off-by: amh1k --- backend/internal/coordinator/jobs.go | 33 ++++- backend/internal/coordinator/jobs_test.go | 59 ++++++++ backend/internal/coordinator/runs.go | 49 ++++++- backend/internal/coordinator/runs_test.go | 167 ++++++++++++++++++++++ 4 files changed, 298 insertions(+), 10 deletions(-) create mode 100644 backend/internal/coordinator/runs_test.go diff --git a/backend/internal/coordinator/jobs.go b/backend/internal/coordinator/jobs.go index 4e4a996..88df241 100644 --- a/backend/internal/coordinator/jobs.go +++ b/backend/internal/coordinator/jobs.go @@ -18,6 +18,7 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/util/wait" ) const ( @@ -307,11 +308,28 @@ func (c *Coordinator) deleteBenchmarkJob(ctx context.Context, reference JobRef) if err != nil { return fmt.Errorf("delete benchmark Job: %w", err) } + // Foreground deletion keeps the Job until its blocking dependents are gone. + // A successful DELETE only acknowledges the request; wait for confirmation. + err = wait.PollUntilContextCancel(ctx, jobPollInterval, true, func(ctx context.Context) (bool, error) { + job, err := c.kubeClient.BatchV1().Jobs(c.config.WorkloadNamespace).Get(ctx, reference.Name, metav1.GetOptions{}) + if apierrors.IsNotFound(err) { + return true, nil + } + if err != nil { + return false, err + } + if job.UID != reference.UID { + return false, errors.New("Job UID changed before deletion could be confirmed") + } + return false, nil + }) + if err != nil { + return fmt.Errorf("wait for benchmark Job deletion: %w", err) + } return nil } -// cleanupBenchmarkResources removes the resources created for one run using -// a context that survives caller cancellation but has a bounded lifetime. +// cleanupBenchmarkResources removes the resources created for one run func (c *Coordinator) cleanupBenchmarkResources(ctx context.Context, resources benchmarkResources) error { if c == nil { return errors.New("coordinator is nil") @@ -327,7 +345,16 @@ func (c *Coordinator) cleanupBenchmarkResources(ctx context.Context, resources b func (c *Coordinator) cleanupBenchmarkResourcesWithContext(ctx context.Context, resources benchmarkResources) error { var cleanupErrors []error if resources.job.Name != "" { - if err := c.deleteBenchmarkJob(ctx, resources.job); err != nil { + jobCtx := ctx + cancel := func() {} + if deadline, ok := ctx.Deadline(); ok && resources.secret.Name != "" { + // Reserve half the remaining cleanup budget for Secret deletion, + // even if Kubernetes cannot finish deleting the Job in time. + jobCtx, cancel = context.WithTimeout(ctx, time.Until(deadline)/2) + } + err := c.deleteBenchmarkJob(jobCtx, resources.job) + cancel() + if err != nil { cleanupErrors = append(cleanupErrors, fmt.Errorf("cleanup benchmark Job: %w", err)) } } diff --git a/backend/internal/coordinator/jobs_test.go b/backend/internal/coordinator/jobs_test.go index c6d094f..a6d3cdf 100644 --- a/backend/internal/coordinator/jobs_test.go +++ b/backend/internal/coordinator/jobs_test.go @@ -10,6 +10,7 @@ import ( batchv1 "k8s.io/api/batch/v1" corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" @@ -356,6 +357,64 @@ func TestCleanupBenchmarkResourcesPreservesDeleteErrors(t *testing.T) { } } +func TestCleanupWaitsForJobDeletionBeforeSecret(t *testing.T) { + client := fake.NewSimpleClientset() + client.PrependReactor("delete", "jobs", func(ktesting.Action) (bool, runtime.Object, error) { + return true, nil, nil // The API acknowledged deletion, but it is still pending. + }) + reads := 0 + client.PrependReactor("get", "jobs", func(ktesting.Action) (bool, runtime.Object, error) { + reads++ + if reads == 1 { + return true, &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Name: "job", UID: "job-uid"}}, nil + } + return true, nil, apierrors.NewNotFound(batchv1.Resource("jobs"), "job") + }) + secretDeleted := false + client.PrependReactor("delete", "secrets", func(ktesting.Action) (bool, runtime.Object, error) { + if reads < 2 { + t.Fatal("Secret deleted before Job deletion was confirmed") + } + secretDeleted = true + return true, nil, nil + }) + c, err := New(validConfig(), client) + if err != nil { + t.Fatal(err) + } + err = c.cleanupBenchmarkResources(context.Background(), benchmarkResources{ + job: JobRef{Name: "job", UID: "job-uid"}, secret: SecretRef{Name: "secret", UID: "secret-uid"}, + }) + if err != nil || !secretDeleted { + t.Fatalf("error = %v, Secret deleted = %v", err, secretDeleted) + } +} + +func TestCleanupJobDeletionTimeoutStillAttemptsSecret(t *testing.T) { + client := fake.NewSimpleClientset() + client.PrependReactor("delete", "jobs", func(ktesting.Action) (bool, runtime.Object, error) { return true, nil, nil }) + client.PrependReactor("get", "jobs", func(ktesting.Action) (bool, runtime.Object, error) { + return true, &batchv1.Job{ObjectMeta: metav1.ObjectMeta{UID: "job-uid"}}, nil + }) + secretDeleted := false + client.PrependReactor("delete", "secrets", func(ktesting.Action) (bool, runtime.Object, error) { + secretDeleted = true + return true, nil, nil + }) + c, err := New(validConfig(), client) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + err = c.cleanupBenchmarkResourcesWithContext(ctx, benchmarkResources{ + job: JobRef{Name: "job", UID: "job-uid"}, secret: SecretRef{Name: "secret", UID: "secret-uid"}, + }) + if !errors.Is(err, context.DeadlineExceeded) || !secretDeleted { + t.Fatalf("error = %v, Secret deleted = %v", err, secretDeleted) + } +} + func TestFindJobPod(t *testing.T) { ref := JobRef{Name: "benchmark", UID: "job-uid"} newPod := func(name string) *corev1.Pod { diff --git a/backend/internal/coordinator/runs.go b/backend/internal/coordinator/runs.go index ee85ff2..c30385c 100644 --- a/backend/internal/coordinator/runs.go +++ b/backend/internal/coordinator/runs.go @@ -5,9 +5,10 @@ import ( "errors" "fmt" "strings" + "time" ) -var ErrRunNotImplemented = errors.New("benchmark execution is not implemented") +const failureLogTimeout = 10 * time.Second // Connection contains resolved database credentials. Do not log this value. // Database must explicitly identify the database used for benchmarking. @@ -77,11 +78,9 @@ type Result struct { OutputTruncated bool } -// Run validates one benchmark request. Kubernetes Secret and Job creation -// will be added when the execution components are implemented. -func (c *Coordinator) Run(ctx context.Context, connection Connection, options Options) (Result, error) { - var result Result - +// Run creates and executes one benchmark Job, returns the runner output, and +// removes the Job and temporary credential Secret before returning. +func (c *Coordinator) Run(ctx context.Context, connection Connection, options Options) (result Result, runErr error) { if c == nil { return result, errors.New("coordinator is nil") } @@ -101,7 +100,43 @@ func (c *Coordinator) Run(ctx context.Context, connection Connection, options Op return result, err } - return result, ErrRunNotImplemented + // Bound the coordinator's API calls and waiting, independently of the + // Kubernetes Job deadline. An earlier caller deadline still takes priority. + ctx, cancelExecution := context.WithTimeout(ctx, c.config.ExecutionTimeout) + defer cancelExecution() + + resources, err := c.createBenchmarkResources(ctx, connection, options, nil) + if err != nil { + return result, err + } + result.JobName = resources.job.Name + + defer func() { + if cleanupErr := c.cleanupBenchmarkResources(ctx, resources); cleanupErr != nil { + runErr = errors.Join(runErr, fmt.Errorf("cleanup benchmark resources: %w", cleanupErr)) + } + }() + + logCtx := ctx + if err := c.waitForJob(ctx, resources.job); err != nil { + runErr = fmt.Errorf("wait for benchmark Job %q: %w", resources.job.Name, err) + // Preserve diagnostics even when the execution context has expired. + // Both Pod discovery and log retrieval share this bounded budget. + var cancel context.CancelFunc + logCtx, cancel = context.WithTimeout(context.WithoutCancel(ctx), failureLogTimeout) + defer cancel() + } + + pod, err := c.findJobPod(logCtx, resources.job) + if err != nil { + return result, errors.Join(runErr, fmt.Errorf("find runner Pod for Job %q: %w", resources.job.Name, err)) + } + + result.Output, result.OutputTruncated, err = c.collectRunnerLogs(logCtx, pod) + if err != nil { + return result, errors.Join(runErr, fmt.Errorf("collect runner logs for Job %q: %w", resources.job.Name, err)) + } + return result, runErr } // DefaultOptions matches the existing pgbench runner defaults. diff --git a/backend/internal/coordinator/runs_test.go b/backend/internal/coordinator/runs_test.go new file mode 100644 index 0000000..0ece169 --- /dev/null +++ b/backend/internal/coordinator/runs_test.go @@ -0,0 +1,167 @@ +package coordinator + +import ( + "context" + "errors" + "strings" + "testing" + "time" + + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/client-go/kubernetes/fake" + typedcorev1 "k8s.io/client-go/kubernetes/typed/core/v1" + ktesting "k8s.io/client-go/testing" +) + +type diagnosticCore struct { + typedcorev1.CoreV1Interface + check func(context.Context) +} + +func (c diagnosticCore) Pods(namespace string) typedcorev1.PodInterface { + return diagnosticPods{PodInterface: c.CoreV1Interface.Pods(namespace), check: c.check} +} + +type diagnosticPods struct { + typedcorev1.PodInterface + check func(context.Context) +} + +func (p diagnosticPods) List(ctx context.Context, options metav1.ListOptions) (*corev1.PodList, error) { + p.check(ctx) + return p.PodInterface.List(ctx, options) +} + +type diagnosticClient struct { + *fake.Clientset + check func(context.Context) +} + +func (c diagnosticClient) CoreV1() typedcorev1.CoreV1Interface { + return diagnosticCore{CoreV1Interface: c.Clientset.CoreV1(), check: c.check} +} + +func TestRunCollectsLogsBeforeCleanupOnFailure(t *testing.T) { + for _, scenario := range []string{"failed job", "cancelled", "deadline", "missing pod", "execution timeout", "caller deadline"} { + t.Run(scenario, func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + config := validConfig() + if scenario == "execution timeout" || scenario == "caller deadline" { + config.ExecutionTimeout = time.Second + } + var callerDeadline time.Time + if scenario == "caller deadline" { + callerDeadline = time.Now().Add(100 * time.Millisecond) + var cancelDeadline context.CancelFunc + ctx, cancelDeadline = context.WithDeadline(ctx, callerDeadline) + defer cancelDeadline() + } + client := fake.NewSimpleClientset() + var jobName string + client.PrependReactor("create", "secrets", func(action ktesting.Action) (bool, runtime.Object, error) { + action.(ktesting.CreateAction).GetObject().(*corev1.Secret).UID = "secret-uid" + return false, nil, nil + }) + client.PrependReactor("create", "jobs", func(action ktesting.Action) (bool, runtime.Object, error) { + job := action.(ktesting.CreateAction).GetObject().(*batchv1.Job) + job.UID = "job-uid" + jobName = job.Name + return false, nil, nil + }) + var originalErr error + jobDeleted := false + client.PrependReactor("delete", "jobs", func(ktesting.Action) (bool, runtime.Object, error) { + jobDeleted = true + return false, nil, nil + }) + client.PrependReactor("get", "jobs", func(action ktesting.Action) (bool, runtime.Object, error) { + if jobDeleted { + return false, nil, nil + } + if scenario == "execution timeout" || scenario == "caller deadline" { + originalErr = context.DeadlineExceeded + // Leave the Job pending so Run must stop itself via its context. + return true, &batchv1.Job{ObjectMeta: metav1.ObjectMeta{UID: "job-uid"}}, nil + } + if scenario == "cancelled" { + cancel() + originalErr = context.Canceled + } else if scenario == "deadline" { + originalErr = context.DeadlineExceeded + } + if originalErr != nil { + return true, nil, originalErr + } + return true, &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{Name: action.(ktesting.GetAction).GetName(), UID: "job-uid"}, + Status: batchv1.JobStatus{Conditions: []batchv1.JobCondition{{Type: batchv1.JobFailed, Status: corev1.ConditionTrue}}}, + }, nil + }) + client.PrependReactor("list", "pods", func(ktesting.Action) (bool, runtime.Object, error) { + if scenario == "missing pod" { + return true, &corev1.PodList{}, nil + } + return true, &corev1.PodList{Items: []corev1.Pod{{ObjectMeta: metav1.ObjectMeta{ + Labels: map[string]string{batchv1.JobNameLabel: jobName}, + Name: "runner-pod", OwnerReferences: []metav1.OwnerReference{{APIVersion: "batch/v1", Kind: "Job", Name: jobName, UID: "job-uid", Controller: boolPtr(true)}}, + }}}}, nil + }) + checked := false + c, err := New(config, diagnosticClient{Clientset: client, check: func(logCtx context.Context) { + checked = true + deadline, ok := logCtx.Deadline() + if logCtx.Err() != nil || !ok || time.Until(deadline) <= 0 || time.Until(deadline) > failureLogTimeout { + t.Fatal("diagnostics need an active, bounded context") + } + }}) + if err != nil { + t.Fatal(err) + } + result, err := c.Run(ctx, validConnection(), DefaultOptions()) + if scenario == "execution timeout" && ctx.Err() != nil { + t.Fatal("execution timeout cancelled the caller context") + } + if scenario == "caller deadline" && time.Since(callerDeadline) > 500*time.Millisecond { + t.Fatal("Run did not honor the earlier caller deadline") + } + if err == nil || !checked { + t.Fatalf("error = %v, diagnostics attempted = %v", err, checked) + } + if originalErr != nil && !errors.Is(err, originalErr) { + t.Fatalf("lost original error: %v", err) + } + if originalErr == nil && !strings.Contains(err.Error(), "failed") { + t.Fatalf("lost Job failure: %v", err) + } + if scenario == "missing pod" { + if !strings.Contains(err.Error(), "no Pod found") { + t.Fatalf("lost discovery error: %v", err) + } + } else if result.Output != "fake logs" { + t.Fatalf("output = %q", result.Output) + } + logged, deletedJob, deletedSecret := false, false, false + for _, action := range client.Actions() { + if action.GetSubresource() == "log" { + logged = true + } + if action.Matches("delete", "jobs") { + deletedJob = true + if scenario != "missing pod" && !logged { + t.Fatal("Job deleted before collecting logs") + } + } + if action.Matches("delete", "secrets") { + deletedSecret = true + } + } + if !deletedJob || !deletedSecret { + t.Fatal("resources were not cleaned up") + } + }) + } +}