Skip to content

Latest commit

 

History

History
373 lines (281 loc) · 11 KB

File metadata and controls

373 lines (281 loc) · 11 KB

Phase 3 Report: Device-Level Workflow Modeling

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


Executive Summary

Phase 3 device-level workflow modeling has been successfully completed. Introduced hierarchical workflow structure (Rack → Device → Task) on top of existing Temporal + Kubernetes integration. Validated parallel execution across devices, sequential execution within devices, and mixed outcome handling. Result propagation standardized for future branching in Phase 4.


Objectives (from Worker Instructions)

Introduce device-level workflow modeling and structured orchestration patterns:

  • Correct hierarchical workflow design (Rack → Device → Task)
  • Sequential execution within device workflows
  • Parallel execution across device workflows
  • Structured result propagation (foundation for branching in Phase 4)

Scope: Device abstraction, DeviceWorkflow, RackWorkflow, result contract standardization


Implementation Details

1. Device Abstraction

Implementation: Simple dict-based device object

Fields:

  • device_id: str
  • device_type: str (e.g., "server")

Example:

device = {"device_id": "device-A", "device_type": "server"}

No platform logic required (minimal abstraction as specified).

2. Activity Result Contract

File: app/activities/k8s_job_activity.py

Standardized Return Type:

@dataclass
class JobResult:
    job_name: str
    status: str

Activity Returns: JobResult object instead of string

Rationale: Enables structured result aggregation in workflows and future branching logic.

3. Device Workflow

File: app/workflows/device_workflow.py

Input: device: dict, include_failure: bool = False

Execution (sequential):

  1. Health → success job (3s)
  2. Firmware → slow job (10s) or fail job (2s) if include_failure=True
  3. Config → success job (3s)

Return:

{
    "device_id": str,
    "tasks": [{"job_name": str, "status": "success|failed"}, ...],
    "status": "success|failed"
}

Key Design: Tasks execute strictly sequentially using await (not asyncio.gather).

4. Rack Workflow

File: app/workflows/rack_workflow.py

Input: devices: list[dict]

Execution (parallel):

  • Executes DeviceWorkflow for each device in parallel using asyncio.gather
  • Uses workflow.execute_child_workflow for proper Temporal parent-child relationship
  • First device includes failure (via include_failure=True), second does not

Return: list[dict] - aggregated device results

5. Worker Registration

File: app/worker.py

Registered Workflows:

  • SimpleWorkflow (Phase 1/2)
  • DeviceWorkflow (Phase 3)
  • RackWorkflow (Phase 3)

Registered Activities:

  • sleep_activity (Phase 1)
  • k8s_job_activity (Phase 2/3)

6. Client Trigger

File: app/client.py

Configuration:

  • 2 devices: device-A (with failure), device-B (without failure)
  • Triggers RackWorkflow with device list
  • Non-blocking execution

Validation Results

1. Parallel Device Execution

✅ Multiple devices start at the same time: Both devices started health jobs within ~24ms

Evidence (Worker Logs):

Device A health: 2026-05-19 15:10:46.644374
Device B health: 2026-05-19 15:10:46.668971
Difference: ~24ms

Analysis: Devices executed in parallel as designed. RackWorkflow correctly used asyncio.gather for parallel child workflow execution.

2. Sequential Task Execution

✅ Within a device, firmware starts only AFTER health completes

Evidence (Worker Logs):

Device A (with failure):

  • Health: 15:10:46.644374 → 15:10:53.689495 (completed)
  • Firmware: 15:10:53.791590 → 15:10:59.835950 (failed)
  • Config: 15:10:59.888115 → 15:11:06.939224 (completed)

Device B (without failure):

  • Health: 15:10:46.668971 → 15:10:53.708936 (completed)
  • Firmware: 15:10:53.751276 → 15:11:07.846447 (completed)
  • Config: 15:11:07.901618 → 15:11:14.957239 (completed)

Analysis: Within each device, tasks executed sequentially. Firmware task only started after health task completed. Config task only started after firmware task completed. DeviceWorkflow correctly used await for sequential execution.

3. Mixed Outcomes

✅ At least one failing task introduced: Device A firmware task failed

Evidence (Worker Logs):

Device A firmware (fail-78f7eb12): 15:10:53.791590 → 15:10:59.835950 (failed)
Device B firmware (slow-fa73da96): 15:10:53.751276 → 15:11:07.846447 (succeeded)

Analysis: Device A had a failing firmware task, Device B had all successful tasks. Failure correctly detected and reported. Device A continued execution after failure (config task still ran).


Evidence

1. Logs

File: phase_3_worker_logs.txt

Device-Level Start Times:

  • Device A health: 2026-05-19 15:10:46.644374
  • Device B health: 2026-05-19 15:10:46.668971

Task-Level Ordering per Device:

  • Device A: health → firmware → config (sequential)
  • Device B: health → firmware → config (sequential)

Parallel Overlap Across Devices:

  • Both devices executing simultaneously
  • Device A health and Device B health overlapped
  • Device A firmware and Device B firmware overlapped

2. kubectl Output

File: phase_3_kubectl_output.txt

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

Analysis: Jobs created, executed, and cleaned up successfully. No orphan resources.

3. Temporal UI Screenshots

UI Location: http://localhost:8088

Screenshots Required: (User to capture from browser preview)

  • RackWorkflow with child DeviceWorkflow executions
  • DeviceWorkflow with child activity executions
  • Clear device-level grouping in workflow tree

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

4. Code Snippets

Activity Result Contract:

@dataclass
class JobResult:
    job_name: str
    status: str

@activity.defn
async def k8s_job_activity(job_type: str) -> JobResult:
    # ... job creation and polling logic ...
    if job_succeeded:
        return JobResult(job_name=job_name, status="success")
    else:
        return JobResult(job_name=job_name, status="failed")

DeviceWorkflow:

@workflow.defn
class DeviceWorkflow:
    @workflow.run
    async def run(self, device: dict, include_failure: bool = False) -> dict:
        device_id = device["device_id"]
        
        # Execute tasks sequentially (using await, not gather)
        health_result = await workflow.execute_activity(
            k8s_job_activity,
            args=["success"],
            start_to_close_timeout=timedelta(seconds=120)
        )
        
        if include_failure:
            firmware_result = await workflow.execute_activity(
                k8s_job_activity,
                args=["fail"],
                start_to_close_timeout=timedelta(seconds=120)
            )
        else:
            firmware_result = await workflow.execute_activity(
                k8s_job_activity,
                args=["slow"],
                start_to_close_timeout=timedelta(seconds=120)
            )
        
        config_result = await workflow.execute_activity(
            k8s_job_activity,
            args=["success"],
            start_to_close_timeout=timedelta(seconds=120)
        )
        
        # Aggregate results
        task_results = [health_result, firmware_result, config_result]
        
        # Determine overall status
        overall_status = "success"
        for result in task_results:
            if result.status == "failed":
                overall_status = "failed"
                break
        
        return {
            "device_id": device_id,
            "tasks": [{"job_name": r.job_name, "status": r.status} for r in task_results],
            "status": overall_status
        }

RackWorkflow:

@workflow.defn
class RackWorkflow:
    @workflow.run
    async def run(self, devices: list[dict]) -> list[dict]:
        # Execute device workflows in parallel using gather
        device_results = await asyncio.gather(*[
            workflow.execute_child_workflow(
                DeviceWorkflow.run,
                args=[device, idx == 0],  # First device gets failure
            )
            for idx, device in enumerate(devices)
        ])
        
        return device_results

Issues Resolved

1. Dataclass Serialization

Problem: TypeError: Expected value to be str, was <class 'dict'> when activity returned dict
Cause: Temporal's default payload converter had issues with dict return type
Solution: Created JobResult dataclass with proper type hints for serialization


Architecture Diagram

RackWorkflow (Parent)
 ├── DeviceWorkflow (Child - device-A)
 │     ├── k8s_job_activity (health)
 │     ├── k8s_job_activity (firmware - FAIL)
 │     └── k8s_job_activity (config)
 ├── DeviceWorkflow (Child - device-B)
 │     ├── k8s_job_activity (health)
 │     ├── k8s_job_activity (firmware - SUCCESS)
 │     └── k8s_job_activity (config)
 └── Return: [device-A result, device-B result]

Execution Semantics:

  • RackWorkflow executes DeviceWorkflow children in parallel
  • Each DeviceWorkflow executes activities sequentially
  • Results aggregated and returned to client

Acceptance Criteria

Phase 3 is complete when:

  • ✅ Device-level workflows execute in parallel (validated: ~24ms difference)
  • ✅ Tasks within device execute strictly sequentially (validated: firmware starts after health)
  • ✅ Failure propagates to device-level status correctly (validated: device-A status = failed)
  • ✅ Results are structured and usable for future branching (validated: JobResult dataclass, structured device results)

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
  • Registered workflows: SimpleWorkflow, DeviceWorkflow, RackWorkflow

Client

cd /home/mcawood/projects/temporal_poc
./venv/bin/python app/client.py
  • Triggers RackWorkflow with 2 devices
  • First device includes failure

STOP CONDITION

Phase 3 COMPLETE - As defined in worker instructions:

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

Sign-off

Phase 3 device-level workflow modeling is complete and verified. Hierarchical workflow structure (Rack → Device → Task) implemented on top of Temporal + Kubernetes integration. Parallel execution across devices validated. Sequential execution within devices validated. Mixed outcome handling validated. Result contract standardized for future branching in Phase 4. Ready for supervisor review and approval.