Skip to content

fix(plugins/ray): honor user task error verdicts - #8083

Open
1fanwang wants to merge 2 commits into
flyteorg:masterfrom
1fanwang:1fanwang/ray-v1-task-error-verdict
Open

1fanwang wants to merge 2 commits into
flyteorg:masterfrom
1fanwang:1fanwang/ray-v1-task-error-verdict

Conversation

@1fanwang

Copy link
Copy Markdown
Contributor

Tracking issue

Related to #7153.

Why are the changes needed?

When a Ray task uses FLYTE_FAIL_ON_ERROR=true, the SDK records a user exception in error.pb and exits 1, so KubeRay reports Failed/AppFailed. The v1 Ray plugin currently classifies even a non-recoverable USER error as a retryable SYSTEM failure. This loses the SDK's user-error verdict.

What changes were proposed in this pull request?

After KubeRay reports failure, use the existing output reader to honor readable USER errors: non-recoverable errors become permanent USER failures, and recoverable errors become retryable USER failures. Both request cleanup.

SYSTEM errors, absent writers, and missing, corrupt or oversized error files retain the existing SYSTEM retry-and-cleanup result. Metadata and read failures are logged. No flag or API is added.

How was this patch tested?

Added file-backed regressions for user/system verdicts and missing, corrupt and oversized error documents.

Labels

fixed

Setup process

The SDK/KubeRay capture and plugin replay setup are included under Testing Done.

Testing Done

The tasks ran with Flytekit 1.16.28, flytekitplugins-ray 1.16.28 and Ray 2.46.0 under KubeRay 1.5.1 on a Kind cluster running Kubernetes 1.35. Both tasks explicitly enable FLYTE_FAIL_ON_ERROR=true. KubeRay produced the failed statuses naturally, and the SDK wrote the captured protobufs.

This is SDK/KubeRay-to-plugin component integration. The Go driver uses the repository fixture for execution metadata, then calls production GetTaskPhase with the captured RayJob and the real filesystem-backed output reader. Status, storage and classification are not mocked. This evidence does not include a full Flyte Admin/Propeller workflow or an observed multi-attempt scheduler run. cleanup=true is the returned cleanup request, not an observation of resource deletion.

Revision Source
Before b2be54c
After 2e0379a

The classifier replays used Go 1.26.2. Set EVIDENCE_ROOT to the absolute directory containing the captures, and run each classifier command from the corresponding checkout's flyteplugins/ directory with the driver below available in the Ray package.

Scenario Command Returned verdict
Before, permanent user error RAY_EVIDENCE_ROOT="$EVIDENCE_ROOT" RAY_EVIDENCE_CASE=permanent go test ./go/tasks/plugins/k8s/ray -run TestLiveRayErrorVerdict -count=1 -v Retryable SYSTEM; cleanup requested
After, permanent user error RAY_EVIDENCE_ROOT="$EVIDENCE_ROOT" RAY_EVIDENCE_CASE=permanent go test ./go/tasks/plugins/k8s/ray -run TestLiveRayErrorVerdict -count=1 -v Permanent USER; cleanup requested
Before, recoverable user error RAY_EVIDENCE_ROOT="$EVIDENCE_ROOT" RAY_EVIDENCE_CASE=recoverable go test ./go/tasks/plugins/k8s/ray -run TestLiveRayErrorVerdict -count=1 -v Retryable SYSTEM; cleanup requested
After, recoverable user error RAY_EVIDENCE_ROOT="$EVIDENCE_ROOT" RAY_EVIDENCE_CASE=recoverable go test ./go/tasks/plugins/k8s/ray -run TestLiveRayErrorVerdict -count=1 -v Retryable USER; cleanup requested

Captured SDK failures

The decoded permanent error document contains:

jq '{origin,kind,code}' permanent/error.json
{
  "origin": "USER",
  "kind": "NON_RECOVERABLE",
  "code": "USER:RuntimeError"
}
jq -r '.message' permanent/error.json

The message ends with this verbatim excerpt:

Message:

    ValueError: public Ray task failure

The recoverable task produced:

jq '{origin,kind,code}' recoverable/error.json
{
  "origin": "USER",
  "kind": "RECOVERABLE",
  "code": "USER:Recoverable"
}
jq -r '.message' recoverable/error.json
Message:

    FlyteRecoverableException: USER:Recoverable: error=public recoverable Ray task failure

Both captured RayJob status messages report entrypoint exit 1:

jq -r '.status.message | split("\n")[0]' permanent-rayjob.json
jq -r '.status.message | split("\n")[0]' recoverable-rayjob.json

Each command prints:

Job entrypoint command failed with exit code 1, last available logs (truncated to 20,000 chars):

Permanent user error

Before:

RAY_EVIDENCE_ROOT="$EVIDENCE_ROOT" RAY_EVIDENCE_CASE=permanent go test ./go/tasks/plugins/k8s/ray -run TestLiveRayErrorVerdict -count=1 -v
case=permanent ray_state=Failed ray_reason=AppFailed sdk_origin=USER sdk_recoverable=false phase=PhaseRetryableFailure phase_origin=SYSTEM cleanup=true
exit_code=1

After:

RAY_EVIDENCE_ROOT="$EVIDENCE_ROOT" RAY_EVIDENCE_CASE=permanent go test ./go/tasks/plugins/k8s/ray -run TestLiveRayErrorVerdict -count=1 -v
case=permanent ray_state=Failed ray_reason=AppFailed sdk_origin=USER sdk_recoverable=false phase=PhasePermanentFailure phase_origin=USER cleanup=true
exit_code=0

Recoverable user error

Before:

RAY_EVIDENCE_ROOT="$EVIDENCE_ROOT" RAY_EVIDENCE_CASE=recoverable go test ./go/tasks/plugins/k8s/ray -run TestLiveRayErrorVerdict -count=1 -v
case=recoverable ray_state=Failed ray_reason=AppFailed sdk_origin=USER sdk_recoverable=true phase=PhaseRetryableFailure phase_origin=SYSTEM cleanup=true
exit_code=1

After:

RAY_EVIDENCE_ROOT="$EVIDENCE_ROOT" RAY_EVIDENCE_CASE=recoverable go test ./go/tasks/plugins/k8s/ray -run TestLiveRayErrorVerdict -count=1 -v
case=recoverable ray_state=Failed ray_reason=AppFailed sdk_origin=USER sdk_recoverable=true phase=PhaseRetryableFailure phase_origin=USER cleanup=true
exit_code=0

These are verbatim result excerpts. The baseline driver exits 1 because its assertions expect the corrected USER verdict; the printed phase shows the baseline's SYSTEM classification.

From-scratch capture and replay

The following is a recipe for a new disposable cluster, not an additional clean run performed for this description. The observed SDK executions and classifier results are above.

In an empty working directory, save the Dockerfile, requirements, task module and two RayJob manifests shown below. Build the image before creating the dedicated kubeconfig:

export EVIDENCE_ROOT="$PWD"
docker build -t flyte-ray-proof:bc85 .
kind create cluster --name flyte-ray-proof --image kindest/node:v1.35.0 --kubeconfig "$EVIDENCE_ROOT/kubeconfig"
export KUBECONFIG="$EVIDENCE_ROOT/kubeconfig"
kind load docker-image flyte-ray-proof:bc85 --name flyte-ray-proof
kubectl create namespace flyte-proof
helm repo add kuberay https://ray-project.github.io/kuberay-helm/
helm repo update kuberay
helm upgrade --install kuberay-operator kuberay/kuberay-operator --version 1.5.1 --namespace flyte-proof
kubectl -n flyte-proof rollout status deployment/kuberay-operator --timeout=300s
kubectl apply -f rayjob.yaml -f rayjob-recoverable.yaml

Capture the statuses and SDK-written files without changing either:

for scenario in permanent recoverable; do
  job="flyte-${scenario}-failure"
  kubectl -n flyte-proof wait --for=jsonpath='{.status.jobDeploymentStatus}'=Failed "rayjob/$job" --timeout=600s
  kubectl -n flyte-proof get "rayjob/$job" -o json > "${scenario}-rayjob.json"
  cluster="$(jq -r '.status.rayClusterName' "${scenario}-rayjob.json")"
  head="$(kubectl -n flyte-proof get pods -l "ray.io/cluster=$cluster,ray.io/node-type=head" -o jsonpath='{.items[0].metadata.name}')"
  mkdir -p "$scenario"
  kubectl -n flyte-proof exec "$head" -c ray-head -- cat "/outputs/$scenario/error.pb" > "$scenario/error.pb"
done

The manifests keep the Ray heads running so their output volumes remain available. The JSON below is only a readable rendering; the Go driver consumes the original protobuf:

docker run --rm -i -v "$EVIDENCE_ROOT/permanent:/evidence/permanent" -v "$EVIDENCE_ROOT/recoverable:/evidence/recoverable" -w /evidence flyte-ray-proof:bc85 python - <<'PY'
import json
from pathlib import Path
from flyteidl.core.errors_pb2 import ErrorDocument

for scenario in ("permanent", "recoverable"):
    document = ErrorDocument()
    document.ParseFromString(Path(scenario, "error.pb").read_bytes())
    error = document.error
    fields = error.DESCRIPTOR.fields_by_name
    result = {
        "origin": fields["origin"].enum_type.values_by_number[error.origin].name,
        "kind": fields["kind"].enum_type.values_by_number[error.kind].name,
        "code": error.code,
        "message": error.message,
    }
    Path(scenario, "error.json").write_text(json.dumps(result, indent=2) + "\n")
PY

Prepare separate Flyte source checkouts at the before/after revisions above. Save the complete Go driver below as flyteplugins/go/tasks/plugins/k8s/ray/live_evidence_test.go in each checkout. It is a reproduction aid outside the product diff. With Go 1.26 available and EVIDENCE_ROOT still exported, run the recorded commands from each checkout's flyteplugins/ directory.

Reproducer source: Dockerfile
FROM python:3.12-slim-bookworm
RUN apt-get update && apt-get install -y --no-install-recommends wget && apt-get clean
COPY requirements.txt /tmp/requirements.txt
RUN pip install --no-cache-dir -r /tmp/requirements.txt
WORKDIR /repro
COPY ray_tasks.py .
RUN python -c 'from pathlib import Path; from flyteidl.core.literals_pb2 import LiteralMap; Path("inputs.pb").write_bytes(LiteralMap().SerializeToString())'
ENV PYTHONUNBUFFERED=1

Dependency configuration, saved as requirements.txt:

ray[default]==2.46.0
flytekit==1.16.28
flytekitplugins-ray==1.16.28
Reproducer source: ray_tasks.py
from flytekit import task
from flytekit.core.constants import FLYTE_FAIL_ON_ERROR
from flytekit.exceptions.user import FlyteRecoverableException
from flytekitplugins.ray import RayJobConfig

ray_config: RayJobConfig = RayJobConfig(worker_node_config=[], address="auto")
environment: dict[str, str] = {FLYTE_FAIL_ON_ERROR: "true"}


@task(task_config=ray_config, environment=environment)
def permanent_failure() -> None:
    raise ValueError("public Ray task failure")


@task(task_config=ray_config, environment=environment)
def recoverable_failure() -> None:
    raise FlyteRecoverableException("public recoverable Ray task failure")
Reproducer source: rayjob.yaml
apiVersion: ray.io/v1
kind: RayJob
metadata:
  name: flyte-permanent-failure
  namespace: flyte-proof
spec:
  entrypoint: >-
    pyflyte-execute --inputs file:///repro/inputs.pb
    --output-prefix file:///outputs/permanent
    --raw-output-data-prefix file:///outputs/raw
    --resolver flytekit.core.python_auto_container.default_task_resolver
    -- task-module ray_tasks task-name permanent_failure
  shutdownAfterJobFinishes: false
  backoffLimit: 0
  rayClusterSpec:
    rayVersion: "2.46.0"
    headGroupSpec:
      rayStartParams:
        num-cpus: "1"
        object-store-memory: "100000000"
        dashboard-host: "0.0.0.0"
        disable-usage-stats: "true"
      template:
        spec:
          containers:
            - name: ray-head
              image: flyte-ray-proof:bc85
              imagePullPolicy: Never
              env:
                - name: FLYTE_FAIL_ON_ERROR
                  value: "true"
                - name: RAY_USAGE_STATS_ENABLED
                  value: "0"
              ports:
                - name: gcs-server
                  containerPort: 6379
                - name: dashboard
                  containerPort: 8265
                - name: client
                  containerPort: 10001
              resources:
                requests:
                  cpu: 200m
                  memory: 1Gi
                limits:
                  cpu: "1"
                  memory: 2Gi
              volumeMounts:
                - name: output
                  mountPath: /outputs
                - name: shared-memory
                  mountPath: /dev/shm
          volumes:
            - name: output
              emptyDir: {}
            - name: shared-memory
              emptyDir:
                medium: Memory
                sizeLimit: 256Mi
    workerGroupSpecs: []
  submitterPodTemplate:
    spec:
      restartPolicy: Never
      containers:
        - name: submitter
          image: flyte-ray-proof:bc85
          imagePullPolicy: Never
          resources:
            requests:
              cpu: 100m
              memory: 128Mi
            limits:
              cpu: "1"
              memory: 512Mi
Reproducer source: rayjob-recoverable.yaml
apiVersion: ray.io/v1
kind: RayJob
metadata:
  name: flyte-recoverable-failure
  namespace: flyte-proof
spec:
  backoffLimit: 0
  entrypoint: pyflyte-execute --inputs file:///repro/inputs.pb --output-prefix file:///outputs/recoverable
    --raw-output-data-prefix file:///outputs/raw --resolver flytekit.core.python_auto_container.default_task_resolver
    -- task-module ray_tasks task-name recoverable_failure
  rayClusterSpec:
    headGroupSpec:
      rayStartParams:
        dashboard-host: 0.0.0.0
        disable-usage-stats: "true"
        num-cpus: "1"
        object-store-memory: "100000000"
      template:
        spec:
          containers:
          - env:
            - name: FLYTE_FAIL_ON_ERROR
              value: "true"
            - name: RAY_USAGE_STATS_ENABLED
              value: "0"
            image: flyte-ray-proof:bc85
            imagePullPolicy: Never
            name: ray-head
            ports:
            - containerPort: 6379
              name: gcs-server
            - containerPort: 8265
              name: dashboard
            - containerPort: 10001
              name: client
            resources:
              limits:
                cpu: "1"
                memory: 2Gi
              requests:
                cpu: 200m
                memory: 1Gi
            volumeMounts:
            - mountPath: /outputs
              name: output
            - mountPath: /dev/shm
              name: shared-memory
          volumes:
          - emptyDir: {}
            name: output
          - emptyDir:
              medium: Memory
              sizeLimit: 256Mi
            name: shared-memory
    rayVersion: 2.46.0
    workerGroupSpecs: []
  shutdownAfterJobFinishes: false
  submitterPodTemplate:
    spec:
      containers:
      - image: flyte-ray-proof:bc85
        imagePullPolicy: Never
        name: submitter
        resources:
          limits:
            cpu: "1"
            memory: 512Mi
          requests:
            cpu: 100m
            memory: 128Mi
      restartPolicy: Never
Reproducer source: live_evidence_test.go
package ray

import (
	"context"
	"encoding/json"
	"fmt"
	"os"
	"path/filepath"
	"testing"

	rayv1 "github.com/ray-project/kuberay/ray-operator/apis/ray/v1"
	"github.com/stretchr/testify/require"

	"github.com/flyteorg/flyte/flyteidl/gen/pb-go/flyteidl/core"
	pluginsCore "github.com/flyteorg/flyte/flyteplugins/go/tasks/pluginmachinery/core"
	"github.com/flyteorg/flyte/flyteplugins/go/tasks/pluginmachinery/ioutils"
	"github.com/flyteorg/flyte/flyteplugins/go/tasks/pluginmachinery/k8s"
	k8sMocks "github.com/flyteorg/flyte/flyteplugins/go/tasks/pluginmachinery/k8s/mocks"
	"github.com/flyteorg/flyte/flytestdlib/contextutils"
	"github.com/flyteorg/flyte/flytestdlib/promutils"
	"github.com/flyteorg/flyte/flytestdlib/promutils/labeled"
	"github.com/flyteorg/flyte/flytestdlib/storage"
	"github.com/flyteorg/stow/local"
)

func TestLiveRayErrorVerdict(t *testing.T) {
	labeled.SetMetricKeys(contextutils.ExecIDKey)
	root := os.Getenv("RAY_EVIDENCE_ROOT")
	if root == "" {
		t.Skip("RAY_EVIDENCE_ROOT must point to captured live KubeRay and SDK output")
	}
	name := os.Getenv("RAY_EVIDENCE_CASE")
	require.Contains(t, []string{"permanent", "recoverable"}, name)
	data, err := os.ReadFile(filepath.Join(root, name+"-rayjob.json"))
	require.NoError(t, err)
	var job rayv1.RayJob
	require.NoError(t, json.Unmarshal(data, &job))
	require.Equal(t, rayv1.JobDeploymentStatusFailed, job.Status.JobDeploymentStatus)

	directory := filepath.Join(root, name)
	store, err := storage.NewDataStore(&storage.Config{
		Type:                  storage.TypeLocal,
		InitContainer:         directory,
		MultiContainerEnabled: true,
		Stow: storage.StowConfig{
			Kind:   local.Kind,
			Config: map[string]string{local.ConfigKeyPath: "/"},
		},
		Limits: storage.LimitsConfig{GetLimitMegabytes: 2},
	}, promutils.NewTestScope())
	require.NoError(t, err)
	ctx := context.Background()
	prefix := storage.DataReference("file://" + directory)
	paths := ioutils.NewCheckpointRemoteFilePaths(ctx, store, prefix, ioutils.NewRawOutputPaths(ctx, prefix), "")
	writer := ioutils.NewRemoteFileOutputWriter(ctx, store, paths)
	reader := ioutils.NewRemoteFileOutputReader(ctx, store, writer, 0)
	found, err := reader.IsError(ctx)
	require.NoError(t, err)
	require.True(t, found)
	verdict, err := reader.ReadError(ctx)
	require.NoError(t, err)
	require.Equal(t, core.ExecutionError_USER, verdict.Kind)

	pluginContext := newPluginContext(k8s.PluginState{}).(*k8sMocks.PluginContext)
	pluginContext.EXPECT().OutputWriter().Return(writer)
	pluginContext.EXPECT().DataStore().Return(store)
	phase, err := (rayJobResourceHandler{}).GetTaskPhase(ctx, pluginContext, &job)
	require.NoError(t, err)
	fmt.Printf("case=%s ray_state=%s ray_reason=%s sdk_origin=%s sdk_recoverable=%t phase=%s phase_origin=%s cleanup=%t\n",
		name, job.Status.JobDeploymentStatus, job.Status.Reason, verdict.Kind, verdict.IsRecoverable,
		phase.Phase(), phase.Err().Kind, phase.CleanupOnFailure())
	require.Equal(t, core.ExecutionError_USER, phase.Err().Kind)
	expectedPhase := pluginsCore.PhasePermanentFailure
	if name == "recoverable" {
		expectedPhase = pluginsCore.PhaseRetryableFailure
	}
	require.Equal(t, expectedPhase, phase.Phase())
	require.True(t, phase.CleanupOnFailure())
}

Screenshots

Not applicable; the status, error-document and classifier output above are the evidence.

Check all the applicable boxes

  • I updated the documentation accordingly. Not applicable; no configuration or public API was added.
  • All new and existing tests passed. This covers the Ray package, including race detection; the full repository suite was not run.
  • All commits are signed-off.

Related PRs

Stack

No dependent PRs. Git Town documentation.

Docs link

Not applicable; no documentation changes.

Signed-off-by: 1fanwang <1fannnw@gmail.com>
@github-actions github-actions Bot added the flyte label Sep 26, 2026
@codecov

codecov Bot commented Sep 26, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 57.33%. Comparing base (b2be54c) to head (bb54559).

Additional details and impacted files
@@            Coverage Diff             @@
##           master    #8083      +/-   ##
==========================================
+ Coverage   57.32%   57.33%   +0.01%     
==========================================
  Files         931      931              
  Lines       58315    58329      +14     
==========================================
+ Hits        33427    33444      +17     
+ Misses      21835    21832       -3     
  Partials     3053     3053              
Flag Coverage Δ
unittests-datacatalog 53.62% <ø> (ø)
unittests-flyteadmin 53.24% <ø> (ø)
unittests-flytecopilot 48.05% <ø> (ø)
unittests-flytectl 64.16% <ø> (+0.04%) ⬆️
unittests-flyteidl 76.63% <ø> (ø)
unittests-flyteplugins 60.64% <100.00%> (+0.05%) ⬆️
unittests-flytepropeller 53.84% <ø> (ø)
unittests-flytestdlib 64.41% <ø> (ø)

Flags with carried forward coverage won't be shown. Click here to find out more.

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

Signed-off-by: 1fanwang <1fannnw@gmail.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant