Let the listener own the EphemeralRunner Running phase transition

The controller derived Status.Phase directly from the pod phase, so a runner
became Running as soon as its pod started, whether or not it had picked up a
job. That made Running mean "the pod is up" instead of "the runner is busy",
and it left the EphemeralRunnerSet scale-down path unable to tell an idle
runner from one that is executing a job.

The listener already knows when a job is assigned to a specific runner, so
move the transition there. HandleJobStarted now reads the runner first and
only promotes it to Running when it is not terminal (Failed, Succeeded or
Outdated) and not being deleted, then patches the phase alongside the job
fields it already writes. The listener role gains "get" on ephemeralrunners
for that read.

On the controller side updateRunStatusFromPod keeps publishing the initial
Pending phase while the pod is starting, and no longer promotes to Running.
Runners waiting for work now stay Pending, so scale-down picks them before
runners that are actually executing a job.

Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
This commit is contained in:
Nikola Jokic
2026-09-12 23:37:31 +02:00
co-authored by Copilot App
parent 9b4c03e312
commit 03a2cdbdec
6 changed files with 309 additions and 51 deletions
+39 -11
View File
@@ -123,6 +123,7 @@ func (w *Scaler) applyDefaults() error {
// It takes a context and a jobInfo parameter which contains the details of the started job.
// This update marks the ephemeral runner so that the controller would have more context
// about the ephemeral runner that should not be deleted when scaling down.
// It also transitions the phase to Running if the runner is not in a terminal state.
// It returns an error if there is any issue with updating the job information.
func (w *Scaler) HandleJobStarted(ctx context.Context, jobInfo *scaleset.JobStarted) error {
w.logger.Info("Updating job info for the runner",
@@ -137,23 +138,50 @@ func (w *Scaler) HandleJobStarted(ctx context.Context, jobInfo *scaleset.JobStar
w.dirty = true
// Fetch current EphemeralRunner to check phase and deletion status
currentRunner := &v1alpha1.EphemeralRunner{}
err := w.clientset.RESTClient().
Get().
Prefix("apis", v1alpha1.GroupVersion.Group, v1alpha1.GroupVersion.Version).
Namespace(w.config.EphemeralRunnerSetNamespace).
Resource("EphemeralRunners").
Name(jobInfo.RunnerName).
Do(ctx).
Into(currentRunner)
if err != nil {
if kerrors.IsNotFound(err) {
w.logger.Info("Ephemeral runner not found, skipping job info update", "runnerName", jobInfo.RunnerName)
return nil
}
return fmt.Errorf("failed to get ephemeral runner: %w", err)
}
original, err := json.Marshal(&v1alpha1.EphemeralRunner{})
if err != nil {
return fmt.Errorf("failed to marshal empty ephemeral runner: %w", err)
}
patch, err := json.Marshal(
&v1alpha1.EphemeralRunner{
Status: v1alpha1.EphemeralRunnerStatus{
JobRequestID: jobInfo.RunnerRequestID,
JobRepositoryName: fmt.Sprintf("%s/%s", jobInfo.OwnerName, jobInfo.RepositoryName),
JobID: jobInfo.JobID,
WorkflowRunID: jobInfo.WorkflowRunID,
JobWorkflowRef: jobInfo.JobWorkflowRef,
JobDisplayName: jobInfo.JobDisplayName,
},
// Build patch with job fields
patchRunner := &v1alpha1.EphemeralRunner{
Status: v1alpha1.EphemeralRunnerStatus{
JobRequestID: jobInfo.RunnerRequestID,
JobRepositoryName: fmt.Sprintf("%s/%s", jobInfo.OwnerName, jobInfo.RepositoryName),
JobID: jobInfo.JobID,
WorkflowRunID: jobInfo.WorkflowRunID,
JobWorkflowRef: jobInfo.JobWorkflowRef,
JobDisplayName: jobInfo.JobDisplayName,
},
)
}
// Only set Running phase if current phase is not terminal/failure and deletion is not in progress
if currentRunner.DeletionTimestamp == nil &&
currentRunner.Status.Phase != v1alpha1.EphemeralRunnerPhaseFailed &&
currentRunner.Status.Phase != v1alpha1.EphemeralRunnerPhaseSucceeded &&
currentRunner.Status.Phase != v1alpha1.EphemeralRunnerPhaseOutdated {
patchRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
}
patch, err := json.Marshal(patchRunner)
if err != nil {
return fmt.Errorf("failed to marshal ephemeral runner patch: %w", err)
}
+145
View File
@@ -2,13 +2,22 @@ package scaler
import (
"bytes"
"context"
"encoding/json"
"log/slog"
"math"
"net/http"
"net/http/httptest"
"strconv"
"testing"
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
"github.com/actions/scaleset"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
)
var discardLogger = slog.New(slog.DiscardHandler)
@@ -121,6 +130,142 @@ func TestEffectiveRateLimiterConfig_QuietAtInfoLevel(t *testing.T) {
}
}
func TestHandleJobStarted(t *testing.T) {
jobInfo := &scaleset.JobStarted{
RunnerName: "runner-1",
JobMessageBase: scaleset.JobMessageBase{
OwnerName: "actions",
RepositoryName: "actions-runner-controller",
JobID: "job-1",
WorkflowRunID: 456,
JobWorkflowRef: "actions/actions-runner-controller/.github/workflows/ci.yaml@refs/heads/main",
JobDisplayName: "build",
RunnerRequestID: 123,
},
}
t.Run("patches job fields and running phase together", func(t *testing.T) {
runner := newTestEphemeralRunner(jobInfo.RunnerName, "")
scaler, shutdown := newTestScaler(t, runner)
defer shutdown()
require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo))
assertJobStartedStatus(t, runner, jobInfo)
assert.Equal(t, v1alpha1.EphemeralRunnerPhaseRunning, runner.Status.Phase)
})
t.Run("repeated assignment remains idempotent", func(t *testing.T) {
runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhaseRunning)
scaler, shutdown := newTestScaler(t, runner)
defer shutdown()
require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo))
firstStatus := runner.Status
require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo))
assert.Equal(t, firstStatus, runner.Status)
assertJobStartedStatus(t, runner, jobInfo)
assert.Equal(t, v1alpha1.EphemeralRunnerPhaseRunning, runner.Status.Phase)
})
for _, phase := range []v1alpha1.EphemeralRunnerPhase{
v1alpha1.EphemeralRunnerPhaseFailed,
v1alpha1.EphemeralRunnerPhaseSucceeded,
v1alpha1.EphemeralRunnerPhaseOutdated,
} {
t.Run("preserves "+string(phase)+" phase while patching job fields", func(t *testing.T) {
runner := newTestEphemeralRunner(jobInfo.RunnerName, phase)
scaler, shutdown := newTestScaler(t, runner)
defer shutdown()
require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo))
assertJobStartedStatus(t, runner, jobInfo)
assert.Equal(t, phase, runner.Status.Phase)
})
}
t.Run("preserves deleting runner phase while patching job fields", func(t *testing.T) {
runner := newTestEphemeralRunner(jobInfo.RunnerName, v1alpha1.EphemeralRunnerPhasePending)
deletionTimestamp := metav1.Now()
runner.DeletionTimestamp = &deletionTimestamp
scaler, shutdown := newTestScaler(t, runner)
defer shutdown()
require.NoError(t, scaler.HandleJobStarted(context.Background(), jobInfo))
assertJobStartedStatus(t, runner, jobInfo)
assert.Equal(t, v1alpha1.EphemeralRunnerPhasePending, runner.Status.Phase)
})
}
func newTestEphemeralRunner(name string, phase v1alpha1.EphemeralRunnerPhase) *v1alpha1.EphemeralRunner {
return &v1alpha1.EphemeralRunner{
ObjectMeta: metav1.ObjectMeta{
Name: name,
Namespace: "default",
},
Status: v1alpha1.EphemeralRunnerStatus{
Phase: phase,
},
}
}
func newTestScaler(t *testing.T, runner *v1alpha1.EphemeralRunner) (*Scaler, func()) {
t.Helper()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
switch r.Method {
case http.MethodGet:
require.NoError(t, json.NewEncoder(w).Encode(runner))
case http.MethodPatch:
var patch v1alpha1.EphemeralRunner
require.NoError(t, json.NewDecoder(r.Body).Decode(&patch))
runner.Status.JobRequestID = patch.Status.JobRequestID
runner.Status.JobRepositoryName = patch.Status.JobRepositoryName
runner.Status.JobID = patch.Status.JobID
runner.Status.WorkflowRunID = patch.Status.WorkflowRunID
runner.Status.JobWorkflowRef = patch.Status.JobWorkflowRef
runner.Status.JobDisplayName = patch.Status.JobDisplayName
if patch.Status.Phase != "" {
runner.Status.Phase = patch.Status.Phase
}
require.NoError(t, json.NewEncoder(w).Encode(runner))
default:
http.Error(w, "unexpected method", http.StatusMethodNotAllowed)
}
}))
clientset, err := kubernetes.NewForConfig(&rest.Config{Host: server.URL})
require.NoError(t, err)
return &Scaler{
clientset: clientset,
config: Config{
EphemeralRunnerSetNamespace: runner.Namespace,
},
targetRunners: -1,
patchSeq: -1,
logger: discardLogger,
}, server.Close
}
func assertJobStartedStatus(t *testing.T, runner *v1alpha1.EphemeralRunner, jobInfo *scaleset.JobStarted) {
t.Helper()
assert.Equal(t, jobInfo.RunnerRequestID, runner.Status.JobRequestID)
assert.Equal(t, jobInfo.JobID, runner.Status.JobID)
assert.Equal(t, jobInfo.OwnerName+"/"+jobInfo.RepositoryName, runner.Status.JobRepositoryName)
assert.Equal(t, jobInfo.WorkflowRunID, runner.Status.WorkflowRunID)
assert.Equal(t, jobInfo.JobWorkflowRef, runner.Status.JobWorkflowRef)
assert.Equal(t, jobInfo.JobDisplayName, runner.Status.JobDisplayName)
}
func TestSetDesiredWorkerState_MinMaxDefaults(t *testing.T) {
newEmptyWorker := func() *Scaler {
return &Scaler{