Skip to content

Latest commit

 

History

History
328 lines (250 loc) · 9.08 KB

File metadata and controls

328 lines (250 loc) · 9.08 KB

Phase 2 Report: Kubernetes Job Integration

Project: RCDiags Temporal Integration – Phase 1 PoC
Date: 2026-05-19
Status: ✅ COMPLETE


Executive Summary

Phase 2 Kubernetes Job integration has been successfully completed. Temporal activities now execute real Kubernetes Jobs via the Kubernetes Python client. Job lifecycle handling (create → poll → complete/fail → cleanup) validated. Parallel execution maintained from Phase 1. Failure handling confirmed. All Phase 2 deliverables met as defined in worker instructions.


Objectives (from Worker Instructions)

Extend the Temporal PoC to integrate real Kubernetes Jobs as activity execution:

  • Temporal → Kubernetes interaction
  • Activity implementation via real Jobs
  • Job lifecycle handling (create → poll → complete/fail)

Scope: Replace sleep activity with Kubernetes-backed activity, create simple reusable job templates, validate success + failure paths


Implementation Details

1. Kubernetes Job Templates

Directory: k8s-jobs/

success-job.yaml:

  • Sleeps 3 seconds
  • Exits 0 (success)
  • Container: busybox:1.36

fail-job.yaml:

  • Sleeps 2 seconds
  • Exits 1 (failure)
  • Container: busybox:1.36

slow-job.yaml:

  • Sleeps 10 seconds
  • Exits 0 (success)
  • Container: busybox:1.36

All jobs use simple busybox containers with shell commands for minimal complexity.

2. Activity: K8s Job Runner

File: app/activities/k8s_job_activity.py

Functionality:

  • Creates Kubernetes Job via Kubernetes Python client
  • Polls job status until completion or failure
  • Detects completion or failure
  • Returns result (success / failure)
  • Unique job name per invocation (UUID-based)
  • Cleanup job after completion (Foreground deletion)

Key Implementation Details:

  • Kubernetes client imported inside activity to avoid Temporal workflow sandbox restrictions
  • Uses config.load_kube_config() for local kind cluster access
  • Polls with 1-second intervals, 120-second max timeout
  • Returns format: {job-name}:{status}

3. Workflow Update

File: app/workflows/simple_workflow.py

Changes:

  • Replaced sleep_activity with k8s_job_activity
  • Runs 3 jobs in parallel using asyncio.gather():
    • 1 success job
    • 1 slow job
    • 1 fail job
  • Increased timeout to 120 seconds per activity (for slow job)

4. Worker Update

File: app/worker.py

Changes:

  • Registered k8s_job_activity alongside sleep_activity
  • Maintains single task queue: "test-queue"

Validation Results

1. Kubernetes Integration

✅ Jobs visible via kubectl get jobs: Jobs created and executed
✅ Pods created and completed: Jobs spawned pods that completed
✅ Cleanup working: No orphan jobs or pods remaining after completion

kubectl Output (after workflow completion):

kubectl get jobs: No resources found in default namespace.
kubectl get pods: No resources found in default namespace.

Analysis: Jobs were created, executed, and cleaned up as designed. No resources remain.

2. Parallel Execution

✅ Jobs start at nearly same time: All 3 jobs created within ~60ms

Worker Logs:

- success-2396c884 created: 2026-05-19 14:42:14.289375
- slow-27d08dd4 created: 2026-05-19 14:42:14.318786
- fail-ae1ab2d4 created: 2026-05-19 14:42:14.348830

Analysis: All jobs created within ~60ms of each other (nearly identical start times). Parallel execution validated.

3. Failure Handling

✅ Failed job detected correctly: fail-ae1ab2d4 detected as failed
✅ Workflow does not crash unexpectedly: Workflow continued despite failure

Worker Logs:

- fail-ae1ab2d4 failed: 2026-05-19 14:42:31.440968
- fail-ae1ab2d4 deleted: 2026-05-19 14:42:31.450734

Analysis: Failed job detected and reported correctly. Job cleaned up after failure. Workflow handled failure gracefully.


Evidence

1. kubectl Output

File: phase_2_kubectl_output.txt

kubectl get jobs (after workflow completion):
No resources found in default namespace.

kubectl get pods (after workflow completion):
No resources found in default namespace.

2. Logs

File: phase_2_worker_logs.txt

Job Creation Timestamps:

  • success-2396c884 created: 2026-05-19 14:42:14.289375
  • slow-27d08dd4 created: 2026-05-19 14:42:14.318786
  • fail-ae1ab2d4 created: 2026-05-19 14:42:14.348830

Completion Status:

  • success-2396c884 succeeded: 2026-05-19 14:42:30.393986
  • fail-ae1ab2d4 failed: 2026-05-19 14:42:31.440968
  • slow-27d08dd4 succeeded: 2026-05-19 14:42:35.440398

Cleanup:

  • All jobs deleted after completion
  • No orphan resources

3. Temporal UI Screenshots

UI Location: http://localhost:8088

Screenshots Required: (User to capture from browser preview)

  • Activity execution showing 3 parallel jobs
  • Success + failure results visible
  • Job timeline with overlapping execution

Browser Preview: Available at http://127.0.0.1:46103

4. Code Snippets

Activity Implementation:

@activity.defn
async def k8s_job_activity(job_type: str) -> str:
    """Activity that creates a Kubernetes Job, polls for completion, and cleans up."""
    
    # Import kubernetes inside activity to avoid workflow sandbox restrictions
    from kubernetes import client, config
    
    # Load Kubernetes config (for local development with kind)
    config.load_kube_config()
    
    # Create Kubernetes API client
    batch_v1 = client.BatchV1Api()
    
    # Generate unique job name
    job_name = f"{job_type}-{uuid.uuid4().hex[:8]}"
    
    # Define job spec based on type
    job_spec = {...}
    
    # Create the job
    job = batch_v1.create_namespaced_job(
        body=job_spec,
        namespace="default"
    )
    
    # Poll for job completion
    while attempt < max_attempts:
        job_status = batch_v1.read_namespaced_job_status(...)
        if job_status.status.succeeded:
            job_succeeded = True
            break
        elif job_status.status.failed:
            job_succeeded = False
            break
        await asyncio.sleep(1)
    
    # Cleanup job
    batch_v1.delete_namespaced_job(
        name=job_name,
        namespace="default",
        body=client.V1DeleteOptions(propagation_policy="Foreground")
    )
    
    return f"{job_name}:success" if job_succeeded else f"{job_name}:failed"

Job Spec (Success):

apiVersion: batch/v1
kind: Job
metadata:
  name: success-job
spec:
  template:
    spec:
      containers:
      - name: success
        image: busybox:1.36
        command: ["sh", "-c", "echo 'Starting success job'; sleep 3; echo 'Success job completed'; exit 0"]
      restartPolicy: Never

Workflow Changes:

@workflow.defn
class SimpleWorkflow:
    @workflow.run
    async def run(self) -> list[str]:
        """Workflow that executes 3 Kubernetes Jobs in parallel."""
        
        results = await asyncio.gather(
            workflow.execute_activity(
                k8s_job_activity,
                args=["success"],
                start_to_close_timeout=timedelta(seconds=120)
            ),
            workflow.execute_activity(
                k8s_job_activity,
                args=["slow"],
                start_to_close_timeout=timedelta(seconds=120)
            ),
            workflow.execute_activity(
                k8s_job_activity,
                args=["fail"],
                start_to_close_timeout=timedelta(seconds=120)
            )
        )
        return results

Issues Resolved

1. Workflow Sandbox Restriction

Problem: RestrictedWorkflowAccessError: Cannot access http.client.IncompleteRead.__mro_entries__ from inside a workflow
Cause: Kubernetes client library imported at module level in activity, causing workflow sandbox to attempt validation
Solution: Moved kubernetes import inside activity function (runs outside workflow sandbox)


Acceptance Criteria

Phase 2 is complete when:

  • ✅ Kubernetes Jobs executed via Temporal activities
  • ✅ Parallel job execution validated (jobs created within ~60ms)
  • ✅ Failure case correctly handled (fail job detected and reported)
  • ✅ No blocking or orphan jobs remain (all jobs cleaned up)

Running Services

Temporal Server

cd temporal-server
docker compose up -d

Kubernetes (kind)

kubectl cluster-info
kubectl get nodes
  • Context: kind-kind
  • Status: Ready

Worker

cd /home/mcawood/projects/temporal_poc
./venv/bin/python app/worker.py
  • Listens on task queue: test-queue

Client

cd /home/mcawood/projects/temporal_poc
./venv/bin/python app/client.py
  • Starts workflow non-blocking
  • Prints workflow ID

STOP CONDITION

Phase 2 COMPLETE - As defined in worker instructions:

  • ✅ Do NOT proceed to Phase 3
  • ✅ Phase 2 Report produced
  • ⏳ Awaiting supervisor review

Sign-off

Phase 2 Kubernetes Job integration is complete and verified. Temporal activities successfully execute real Kubernetes Jobs via the Kubernetes Python client. Job lifecycle handling validated. Parallel execution maintained. Failure handling confirmed. All jobs cleaned up after execution. Ready for supervisor review and approval.