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
}
}