PR #3362
Sections
Review

fix(backend): recover dynamic agent routes stuck in starting with no owner

main ← feature/fix-dynamic-route-st-jsp 18 files +682 −94 PR #3362 ↗

Dynamic routes that claimed a generation but never reached a terminal status no longer stay stuck at starting. The engine now marks successful launches active and failed or orphaned launches action_required with a status-fenced CAS, and startup sweeps orphaned starting rows.

Why this change

A dynamic route claims generation N as starting before the agent process starts. If the launch fails after the claim or the process crashes before the async start callback, no timer or event moves the row. The route stays at starting with no owner and the UI shows no recovery action.

What it does

Architecture, end to end

Select claims starting. The async process callback moves it to active. Any failure or orphan sweep moves it to action_required. The projection mirrors the durable row.

flowchart LR
  Select[Engine.Select] --> Starting[(dynamic_route_states: starting)]
  Starting --> ActiveCheck{process start?}
  ActiveCheck -- success --> MarkActive[MarkActive] --> Active[(active)]
  ActiveCheck -- failure --> MarkReq[MarkActionRequired] --> ActionReq[(action_required)]
  Starting -- orphan sweep --> ActionReq
  Starting -- declined fallback --> ActionReq
  Active --> MirrorActive[mirrorDynamicRouteProjection]
  ActionReq --> MirrorReq[mirrorDynamicRouteProjection]
  MirrorActive --> Session[(task_sessions)]
  MirrorReq --> Session

Key code changes

Drag to pan. Use the + and − buttons to zoom. Click a node to open the full code. The arrows show how the parts interact.

drag to pan · +/− to zoom · click a node for details

Engine lifecycle: starting to active or action_requiredapps/backend/internal/agent/runtime/dynamic/engine.go ↗
func (e *Engine) MarkActive(ctx context.Context, sessionID string, expectedGeneration int64) error
Click for details →

The engine now owns the only transition that produces active, and a fenced transition that produces action_required from starting.

Status constants
const (
  routeStatusStarting       = "starting"
  routeStatusActive         = "active"
  routeStatusActionRequired = "action_required"
)
MarkActive
func (e *Engine) MarkActive(ctx context.Context, sessionID string, expectedGeneration int64) error {
  e.mu.Lock()
  defer e.mu.Unlock()
  state, exists, err := e.loadStateLocked(ctx, sessionID)
  if err != nil {
    return err
  }
  if !exists {
    return ErrRouteStateNotFound
  }
  if state.Generation != expectedGeneration {
    return ErrStaleGeneration
  }
  if state.Status != routeStatusStarting {
    return nil
  }
  expectedStatus := state.Status
  state.Status = routeStatusActive
  state.UpdatedAt = e.now()
  if err := e.persistSameGeneration(ctx, expectedGeneration, expectedStatus, state); err != nil {
    return err
  }
  e.states[sessionID] = state
  return nil
}
MarkActionRequired
func (e *Engine) MarkActionRequired(ctx context.Context, sessionID string, expectedGeneration int64, reason string) (RouteDecision, error) {
  e.mu.Lock()
  defer e.mu.Unlock()
  state, exists, err := e.loadStateLocked(ctx, sessionID)
  if err != nil {
    return RouteDecision{}, err
  }
  if !exists {
    return RouteDecision{}, ErrRouteStateNotFound
  }
  if state.Generation != expectedGeneration {
    return RouteDecision{}, ErrStaleGeneration
  }
  if state.Status == routeStatusStarting {
    expectedStatus := state.Status
    state.Status = routeStatusActionRequired
    state.UpdatedAt = e.now()
    if err := e.persistSameGeneration(ctx, expectedGeneration, expectedStatus, state); err != nil {
      return RouteDecision{}, err
    }
    e.states[sessionID] = state
  }
  return RouteDecision{
    SessionID: sessionID, LogicalProfileID: state.LogicalProfileID,
    ExecutionProfileID: state.ExecutionProfileID, Generation: state.Generation,
    ProfileVersion: state.ProfileVersion, Reason: reason, Status: state.Status,
  }, nil
}
func (r *Repository) ClaimRouteStateFrom(ctx context.Context, expectedGeneration int64, expectedStatus string, state RouteState) (bool, error)
Click for details →

Same-generation updates now fence on both generation and status so a late recovery cannot overwrite an active route.

GenerationStatusClaimer
// GenerationStatusClaimer updates one generation only while its status still
// matches the caller's observation. This closes the same-generation race
// between a launch becoming active and a concurrent recovery transition.
type GenerationStatusClaimer interface {
  ClaimRouteStateFrom(context.Context, int64, string, RouteState) (bool, error)
}
ClaimRouteStateFrom
func (r *Repository) ClaimRouteStateFrom(ctx context.Context, expectedGeneration int64, expectedStatus string, state RouteState) (bool, error) {
  if isTransientRouteSession(state.SessionID) {
    return true, nil
  }
  result, err := r.db.ExecContext(ctx, r.db.Rebind(`
    UPDATE dynamic_route_states
    SET logical_profile_id = ?, execution_profile_id = ?,
      route_generation = ?, profile_version = ?, state = ?, continuation_json = ?, policy_state_json = ?, updated_at = ?
    WHERE session_id = ? AND route_generation = ? AND state = ?
  `), state.LogicalProfileID, state.ExecutionProfileID,
    state.Generation, state.ProfileVersion, state.Status, state.ContinuationJSON, state.PolicyStateJSON, state.UpdatedAt,
    state.SessionID, expectedGeneration, expectedStatus)
  if err != nil {
    return false, err
  }
  rows, err := result.RowsAffected()
  return rows == 1, err
}
ListStartingRouteStates
func (r *Repository) ListStartingRouteStates(ctx context.Context) ([]RouteState, error) {
  rows, err := r.ro.QueryContext(ctx, r.ro.Rebind(`
    SELECT session_id, logical_profile_id, execution_profile_id,
      route_generation, profile_version, state, continuation_json, policy_state_json, updated_at
    FROM dynamic_route_states
    WHERE state = ? ORDER BY updated_at ASC
  `), "starting")
  if err != nil {
    return nil, err
  }
  defer func() { _ = rows.Close() }()
  states := make([]RouteState, 0)
  for rows.Next() {
    var state RouteState
    if err := rows.Scan(&state.SessionID, &state.LogicalProfileID, &state.ExecutionProfileID,
      &state.Generation, &state.ProfileVersion, &state.Status,
      &state.ContinuationJSON, &state.PolicyStateJSON, &state.UpdatedAt); err != nil {
      return nil, err
    }
    states = append(states, state)
  }
  return states, rows.Err()
}
Resolver surface for active and action_requiredapps/backend/internal/agent/runtime/dynamic_resolver.go ↗
func (r *ProfileExecutionResolver) MarkRouteActive(ctx context.Context, sessionID string, expectedGeneration int64) error
Click for details →

The orchestrator calls the resolver to complete the starting phase after the async process callback.

MarkRouteActive
func (r *ProfileExecutionResolver) MarkRouteActive(ctx context.Context, sessionID string, expectedGeneration int64) error {
  if r.engine == nil {
    return errors.New("dynamic profile execution is not configured")
  }
  return r.engine.MarkActive(ctx, sessionID, expectedGeneration)
}
MarkRouteActionRequired
func (r *ProfileExecutionResolver) MarkRouteActionRequired(ctx context.Context, sessionID string, expectedGeneration int64, reason string) (RouteDecision, error) {
  if r.engine == nil {
    return RouteDecision{}, errors.New("dynamic profile execution is not configured")
  }
  return r.engine.MarkActionRequired(ctx, sessionID, expectedGeneration, reason)
}
Orchestrator mirroring and async callbacksapps/backend/internal/orchestrator/dynamic_launch.go ↗
func (s *Service) markDynamicRouteActive(ctx context.Context, sessionID string, generation int64)
Click for details →

The service mirrors the durable status onto task_sessions and settles routes from the async process start callbacks.

markDynamicRouteActive
func (s *Service) markDynamicRouteActive(ctx context.Context, sessionID string, generation int64) {
  if s.profileExecutionResolver == nil || sessionID == "" || generation <= 0 {
    return
  }
  if err := s.profileExecutionResolver.MarkRouteActive(ctx, sessionID, generation); err != nil {
    if !errors.Is(err, ErrStaleGeneration) && !errors.Is(err, ErrRouteStateNotFound) {
      s.logger.Warn("failed to mark dynamic route active", zap.String("session_id", sessionID), zap.Error(err))
    }
    return
  }
  session, err := s.repo.GetTaskSession(ctx, sessionID)
  if err != nil || session == nil || session.RouteGeneration != generation {
    return
  }
  s.mirrorDynamicRouteProjection(ctx, session, generation, dynamicRouteStatusActive, session.RouteReason)
}
handleAgentProcessStarted
func (s *Service) handleAgentProcessStarted(ctx context.Context, _, sessionID, agentExecutionID string) {
  if s.profileExecutionResolver == nil || sessionID == "" {
    return
  }
  session, err := s.repo.GetTaskSession(ctx, sessionID)
  if err != nil || session == nil || session.RouteGeneration <= 0 || session.ExecutionProfileID == "" {
    return
  }
  if session.AgentExecutionID != "" && agentExecutionID != "" && session.AgentExecutionID != agentExecutionID {
    return
  }
  if session.State != models.TaskSessionStateStarting && session.State != models.TaskSessionStateRunning {
    return
  }
  s.markDynamicRouteActive(ctx, sessionID, session.RouteGeneration)
}
Deferred guard in routeDynamicAgentFailure
generation := session.RouteGeneration
  handled := false
  defer func() {
    if !handled {
      s.markDynamicRouteActionRequired(ctx, session.ID, generation, reason)
    }
  }()
  if classified == nil || !classified.FallbackAllowed {
    return false
  }
func (s *Service) reconcileOrphanedDynamicStartingRoutes(ctx context.Context)
Click for details →

On restart the service lists every starting row and moves orphaned ones to action_required so the UI can recover.

Reconcile sweep
func (s *Service) reconcileOrphanedDynamicStartingRoutes(ctx context.Context) {
  if s.profileExecutionResolver == nil {
    return
  }
  lister, ok := s.repo.(dynamicStartingRouteLister)
  if !ok {
    return
  }
  states, err := lister.ListStartingRouteStates(ctx)
  if err != nil {
    s.logger.Warn("failed to list starting dynamic route states", zap.Error(err))
    return
  }
  for _, state := range states {
    s.reconcileOrphanedDynamicStartingRoute(ctx, state)
  }
}
Orphan check
func isOrphanableDynamicSessionState(state models.TaskSessionState) bool {
  return state == models.TaskSessionStateStarting || state == models.TaskSessionStateIdle
}
Guarded projection update
func (s *Service) mirrorDynamicRouteProjection(ctx context.Context, session *models.TaskSession, generation int64, status, reason string) {
  if session == nil || session.RouteGeneration != generation || status == "" {
    return
  }
  if projector, ok := s.repo.(dynamicRouteSessionProjector); ok {
    changed, updatedAt, err := projector.UpdateTaskSessionDynamicRouteIfCurrent(ctx, session.ID, generation, session.RouteState, status, reason)
    if err != nil {
      s.logger.Warn("failed to mirror dynamic route state to task session", zap.String("session_id", session.ID), zap.Error(err))
      return
    }
    if !changed {
      return
    }
    session.RouteState = status
    session.RouteReason = reason
    session.UpdatedAt = updatedAt
    s.publishTaskSessionStateChanged(ctx, session.TaskID, session.ID, oldState, session.State, session.ErrorMessage, &updatedAt, session)
    return
  }
}
Read the changes as a list

Engine lifecycle: starting to active or action_required

apps/backend/internal/agent/runtime/dynamic/engine.go

The engine now owns the only transition that produces active, and a fenced transition that produces action_required from starting.

Status constants
const (
  routeStatusStarting       = "starting"
  routeStatusActive         = "active"
  routeStatusActionRequired = "action_required"
)
MarkActive
func (e *Engine) MarkActive(ctx context.Context, sessionID string, expectedGeneration int64) error {
  e.mu.Lock()
  defer e.mu.Unlock()
  state, exists, err := e.loadStateLocked(ctx, sessionID)
  if err != nil {
    return err
  }
  if !exists {
    return ErrRouteStateNotFound
  }
  if state.Generation != expectedGeneration {
    return ErrStaleGeneration
  }
  if state.Status != routeStatusStarting {
    return nil
  }
  expectedStatus := state.Status
  state.Status = routeStatusActive
  state.UpdatedAt = e.now()
  if err := e.persistSameGeneration(ctx, expectedGeneration, expectedStatus, state); err != nil {
    return err
  }
  e.states[sessionID] = state
  return nil
}
MarkActionRequired
func (e *Engine) MarkActionRequired(ctx context.Context, sessionID string, expectedGeneration int64, reason string) (RouteDecision, error) {
  e.mu.Lock()
  defer e.mu.Unlock()
  state, exists, err := e.loadStateLocked(ctx, sessionID)
  if err != nil {
    return RouteDecision{}, err
  }
  if !exists {
    return RouteDecision{}, ErrRouteStateNotFound
  }
  if state.Generation != expectedGeneration {
    return RouteDecision{}, ErrStaleGeneration
  }
  if state.Status == routeStatusStarting {
    expectedStatus := state.Status
    state.Status = routeStatusActionRequired
    state.UpdatedAt = e.now()
    if err := e.persistSameGeneration(ctx, expectedGeneration, expectedStatus, state); err != nil {
      return RouteDecision{}, err
    }
    e.states[sessionID] = state
  }
  return RouteDecision{
    SessionID: sessionID, LogicalProfileID: state.LogicalProfileID,
    ExecutionProfileID: state.ExecutionProfileID, Generation: state.Generation,
    ProfileVersion: state.ProfileVersion, Reason: reason, Status: state.Status,
  }, nil
}

Status-fenced persistence

apps/backend/internal/task/repository/sqlite/dynamic_route.go

Same-generation updates now fence on both generation and status so a late recovery cannot overwrite an active route.

GenerationStatusClaimer
// GenerationStatusClaimer updates one generation only while its status still
// matches the caller's observation. This closes the same-generation race
// between a launch becoming active and a concurrent recovery transition.
type GenerationStatusClaimer interface {
  ClaimRouteStateFrom(context.Context, int64, string, RouteState) (bool, error)
}
ClaimRouteStateFrom
func (r *Repository) ClaimRouteStateFrom(ctx context.Context, expectedGeneration int64, expectedStatus string, state RouteState) (bool, error) {
  if isTransientRouteSession(state.SessionID) {
    return true, nil
  }
  result, err := r.db.ExecContext(ctx, r.db.Rebind(`
    UPDATE dynamic_route_states
    SET logical_profile_id = ?, execution_profile_id = ?,
      route_generation = ?, profile_version = ?, state = ?, continuation_json = ?, policy_state_json = ?, updated_at = ?
    WHERE session_id = ? AND route_generation = ? AND state = ?
  `), state.LogicalProfileID, state.ExecutionProfileID,
    state.Generation, state.ProfileVersion, state.Status, state.ContinuationJSON, state.PolicyStateJSON, state.UpdatedAt,
    state.SessionID, expectedGeneration, expectedStatus)
  if err != nil {
    return false, err
  }
  rows, err := result.RowsAffected()
  return rows == 1, err
}
ListStartingRouteStates
func (r *Repository) ListStartingRouteStates(ctx context.Context) ([]RouteState, error) {
  rows, err := r.ro.QueryContext(ctx, r.ro.Rebind(`
    SELECT session_id, logical_profile_id, execution_profile_id,
      route_generation, profile_version, state, continuation_json, policy_state_json, updated_at
    FROM dynamic_route_states
    WHERE state = ? ORDER BY updated_at ASC
  `), "starting")
  if err != nil {
    return nil, err
  }
  defer func() { _ = rows.Close() }()
  states := make([]RouteState, 0)
  for rows.Next() {
    var state RouteState
    if err := rows.Scan(&state.SessionID, &state.LogicalProfileID, &state.ExecutionProfileID,
      &state.Generation, &state.ProfileVersion, &state.Status,
      &state.ContinuationJSON, &state.PolicyStateJSON, &state.UpdatedAt); err != nil {
      return nil, err
    }
    states = append(states, state)
  }
  return states, rows.Err()
}

Resolver surface for active and action_required

apps/backend/internal/agent/runtime/dynamic_resolver.go

The orchestrator calls the resolver to complete the starting phase after the async process callback.

MarkRouteActive
func (r *ProfileExecutionResolver) MarkRouteActive(ctx context.Context, sessionID string, expectedGeneration int64) error {
  if r.engine == nil {
    return errors.New("dynamic profile execution is not configured")
  }
  return r.engine.MarkActive(ctx, sessionID, expectedGeneration)
}
MarkRouteActionRequired
func (r *ProfileExecutionResolver) MarkRouteActionRequired(ctx context.Context, sessionID string, expectedGeneration int64, reason string) (RouteDecision, error) {
  if r.engine == nil {
    return RouteDecision{}, errors.New("dynamic profile execution is not configured")
  }
  return r.engine.MarkActionRequired(ctx, sessionID, expectedGeneration, reason)
}

Orchestrator mirroring and async callbacks

apps/backend/internal/orchestrator/dynamic_launch.go

The service mirrors the durable status onto task_sessions and settles routes from the async process start callbacks.

markDynamicRouteActive
func (s *Service) markDynamicRouteActive(ctx context.Context, sessionID string, generation int64) {
  if s.profileExecutionResolver == nil || sessionID == "" || generation <= 0 {
    return
  }
  if err := s.profileExecutionResolver.MarkRouteActive(ctx, sessionID, generation); err != nil {
    if !errors.Is(err, ErrStaleGeneration) && !errors.Is(err, ErrRouteStateNotFound) {
      s.logger.Warn("failed to mark dynamic route active", zap.String("session_id", sessionID), zap.Error(err))
    }
    return
  }
  session, err := s.repo.GetTaskSession(ctx, sessionID)
  if err != nil || session == nil || session.RouteGeneration != generation {
    return
  }
  s.mirrorDynamicRouteProjection(ctx, session, generation, dynamicRouteStatusActive, session.RouteReason)
}
handleAgentProcessStarted
func (s *Service) handleAgentProcessStarted(ctx context.Context, _, sessionID, agentExecutionID string) {
  if s.profileExecutionResolver == nil || sessionID == "" {
    return
  }
  session, err := s.repo.GetTaskSession(ctx, sessionID)
  if err != nil || session == nil || session.RouteGeneration <= 0 || session.ExecutionProfileID == "" {
    return
  }
  if session.AgentExecutionID != "" && agentExecutionID != "" && session.AgentExecutionID != agentExecutionID {
    return
  }
  if session.State != models.TaskSessionStateStarting && session.State != models.TaskSessionStateRunning {
    return
  }
  s.markDynamicRouteActive(ctx, sessionID, session.RouteGeneration)
}
Deferred guard in routeDynamicAgentFailure
generation := session.RouteGeneration
  handled := false
  defer func() {
    if !handled {
      s.markDynamicRouteActionRequired(ctx, session.ID, generation, reason)
    }
  }()
  if classified == nil || !classified.FallbackAllowed {
    return false
  }

Startup orphan sweep for starting routes

apps/backend/internal/orchestrator/dynamic_policy_recovery.go

On restart the service lists every starting row and moves orphaned ones to action_required so the UI can recover.

Reconcile sweep
func (s *Service) reconcileOrphanedDynamicStartingRoutes(ctx context.Context) {
  if s.profileExecutionResolver == nil {
    return
  }
  lister, ok := s.repo.(dynamicStartingRouteLister)
  if !ok {
    return
  }
  states, err := lister.ListStartingRouteStates(ctx)
  if err != nil {
    s.logger.Warn("failed to list starting dynamic route states", zap.Error(err))
    return
  }
  for _, state := range states {
    s.reconcileOrphanedDynamicStartingRoute(ctx, state)
  }
}
Orphan check
func isOrphanableDynamicSessionState(state models.TaskSessionState) bool {
  return state == models.TaskSessionStateStarting || state == models.TaskSessionStateIdle
}
Guarded projection update
func (s *Service) mirrorDynamicRouteProjection(ctx context.Context, session *models.TaskSession, generation int64, status, reason string) {
  if session == nil || session.RouteGeneration != generation || status == "" {
    return
  }
  if projector, ok := s.repo.(dynamicRouteSessionProjector); ok {
    changed, updatedAt, err := projector.UpdateTaskSessionDynamicRouteIfCurrent(ctx, session.ID, generation, session.RouteState, status, reason)
    if err != nil {
      s.logger.Warn("failed to mirror dynamic route state to task session", zap.String("session_id", session.ID), zap.Error(err))
      return
    }
    if !changed {
      return
    }
    session.RouteState = status
    session.RouteReason = reason
    session.UpdatedAt = updatedAt
    s.publishTaskSessionStateChanged(ctx, session.TaskID, session.ID, oldState, session.State, session.ErrorMessage, &updatedAt, session)
    return
  }
}

Data and storage

The durable route row is the source of truth. The session row mirrors it for the UI.

FieldTypeNotes
dynamic_route_states.session_idtext PKone row per logical session
dynamic_route_states.route_generationintegermonotonic claim, fenced by CAS
dynamic_route_states.statetextstarting, active, action_required, retry_wait, waiting_for_reset, waiting
dynamic_route_states.policy_state_jsonjsondeadline and retry ordinal for recovery timers
dynamic_route_states.continuation_jsonjsonhandoff package for successor launch
task_sessions.route_statetextprojection of dynamic_route_states.state
task_sessions.route_generationintegerprojection of route_generation
task_sessions.route_reasontextwhy the route moved, shown in UI

Risk

6 / 10 Medium
1 low5 medium10 high

Why this score

  • Touches the generation CAS and adds a new status-fenced CAS. A wrong fence could hide a live route or resurrect a stale one.
  • Startup sweep runs before general session reconciliation. Wrong ordering would make the sweep a no-op for rows with executors_running.
  • Mirrors state to task_sessions with a narrow guarded update. A missed mirror leaves the UI stale while the durable row is correct.

Trade-offs and review notes

Where to look first

  1. Check persistSameGeneration and ClaimRouteStateFrom fencing: generation plus expected status must match, and stale generation returns ErrStaleGeneration.
  2. Check MarkActionRequired is a no-op when status is already active. Verify the test that a permission_denied failure on an active route does not demote it.
  3. Check reconcileOrphanedDynamicStartingRoutes ordering in Service.Start and that it loads via WithStateLoader so it sees durable rows, not the process cache.
  4. Check deferred guards in routeDynamicAgentFailure, LaunchDynamicRouteAction, and relaunchDynamicTaskAfterFailure cover every early return.
  5. Check mirrorDynamicRouteProjection uses UpdateTaskSessionDynamicRouteIfCurrent and publishes taskSessionStateChanged only when the row actually changed.