package streamgate import ( "context" "errors" "sort" "sync" "time" ) var ( // ErrRecoveryCycleInProgress rejects overlapping recovery cycles for the // same request. Host calls remain strictly serialized by the coordinator. ErrRecoveryCycleInProgress = errors.New("streamgate: recovery cycle already in progress") // ErrRecoveryNoCurrentAttempt indicates there is no provider transport // ownership to close before recovery. ErrRecoveryNoCurrentAttempt = errors.New("streamgate: recovery current attempt is not bound") // ErrRecoveryArbitrationIneligible indicates the arbiter did not select one // validated recovery intent. ErrRecoveryArbitrationIneligible = errors.New("streamgate: arbitration result is not recoverable") // ErrRecoveryAbortFailed indicates current provider ownership could not be // closed. No preparation, rebuild, budget consume, or dispatch follows. ErrRecoveryAbortFailed = errors.New("streamgate: recovery attempt abort failed") // ErrRecoveryPreparerUnavailable indicates a selected preparer id has no // registered request-local implementation. ErrRecoveryPreparerUnavailable = errors.New("streamgate: recovery plan preparer is unavailable") // ErrRecoveryPreparationFailed indicates a bounded preparer invocation // failed without exposing its potentially raw error. ErrRecoveryPreparationFailed = errors.New("streamgate: recovery plan preparation failed") // ErrRecoveryPreparationDeadline indicates the bounded preparer deadline // elapsed. It is distinct from caller cancellation. ErrRecoveryPreparationDeadline = errors.New("streamgate: recovery plan preparation deadline exceeded") // ErrRecoveryPreparationInvalid indicates the preparer returned a directive // that does not satisfy the selected immutable plan. ErrRecoveryPreparationInvalid = errors.New("streamgate: recovery plan preparation result is invalid") // ErrRecoveryPreparationSnapshotRelease indicates request-local preparation // state could not be released. Rebuild and dispatch are not attempted. ErrRecoveryPreparationSnapshotRelease = errors.New("streamgate: recovery preparation snapshot release failed") // ErrRecoveryRebuildFailed sanitizes a host rebuilder failure. ErrRecoveryRebuildFailed = errors.New("streamgate: recovery request rebuild failed") // ErrRecoveryRebuildInvalid indicates the rebuilt draft could not be bound // to the immutable recovery plan. ErrRecoveryRebuildInvalid = errors.New("streamgate: recovery rebuilt request is invalid") // ErrRecoveryDispatchFailed sanitizes dispatcher errors and invalid returned // bindings. Fault usage remains consumed because an outbound attempt began. ErrRecoveryDispatchFailed = errors.New("streamgate: recovery attempt dispatch failed") ) // RecoveryCoordinatorOptions binds request-local policy and host seams. The // policy, usage, ingress snapshot, preparer registry, and current attempt are // defensively snapshotted at construction. Concrete request/auth data remains // behind host-owned references. type RecoveryCoordinatorOptions struct { Policy RecoveryPolicySnapshot Usage RecoveryUsageSnapshot RequestSnapshot RecoveryRequestSnapshotRef CurrentBinding AttemptBinding Rebuilder RequestRebuilder Dispatcher AttemptDispatcher Preparers map[string]RecoveryPlanPreparer PreparationTimeout time.Duration ObservationSink ObservationSink } // RecoveryCycleInput carries the immutable data needed to turn one selected // ArbitrationResult.RecoveryIntent into one plan and at most one dispatch. // ConsumerID completes the selected arbiter filter/rule contributor identity. // A non-nil PreparationSnapshot transfers release responsibility to Execute. type RecoveryCycleInput struct { Arbitration ArbitrationResult PlanID string IdempotencyKey string ConsumerID string CommitState CommitState CallerCanceled bool HasToolSideEffect bool HasCompleteToolCall bool ManagedTrajectory *ManagedTrajectoryBudget TerminalReason TerminalReason PreparerID string PreparationSnapshot RecoveryPreparationSnapshot RequiredCapabilities []string FailureCauses FailureCauseChain } // RecoveryCycleResult is an immutable, raw-free account of one cycle. A // failed cycle may still expose its selected plan and consumed usage so the // host can produce one terminal result without retrying the same plan. type RecoveryCycleResult struct { plan RecoveryPlan hasPlan bool binding AttemptBinding hasBinding bool usage RecoveryUsageSnapshot failureCauses FailureCauseChain previousAttemptClosed bool } // PreviousAttemptClosed returns true if the previous attempt was closed/aborted. func (r RecoveryCycleResult) PreviousAttemptClosed() bool { return r.previousAttemptClosed } // Plan returns the immutable selected/finalized plan when plan creation // succeeded. func (r RecoveryCycleResult) Plan() (RecoveryPlan, bool) { if !r.hasPlan { return RecoveryPlan{}, false } return r.plan, true } // Binding returns the single successfully dispatched attempt binding. func (r RecoveryCycleResult) Binding() (AttemptBinding, bool) { if !r.hasBinding { return AttemptBinding{}, false } return r.binding, true } // UsageSnapshot returns the coordinator ledger after this cycle. A failed // outbound dispatch consumes fault usage; failures before dispatch do not. func (r RecoveryCycleResult) UsageSnapshot() RecoveryUsageSnapshot { return copyRecoveryUsageSnapshot(r.usage) } // FailureCauses returns a defensive bounded sanitized cause chain. func (r RecoveryCycleResult) FailureCauses() FailureCauseChain { return r.failureCauses.Copy() } func (r RecoveryCycleResult) withPlan(plan RecoveryPlan) RecoveryCycleResult { r.plan = plan r.hasPlan = true return r } func (r RecoveryCycleResult) withFailure(code string) RecoveryCycleResult { if r.failureCauses.Len() >= MaxFailureCauses { return r } cause, err := NewFailureCause("recovery_coordinator", code, "", "", "") if err != nil { return r } chain, err := r.failureCauses.Append(cause) if err == nil { r.failureCauses = chain } return r } // RecoveryCoordinator owns the mutable request-local recovery ledger and // current provider binding. It never serializes endpoint errors and never // participates in filter evaluation. type RecoveryCoordinator struct { mu sync.Mutex busy bool policy RecoveryPolicySnapshot usage RecoveryUsageSnapshot requestSnapshot RecoveryRequestSnapshotRef currentBinding AttemptBinding hasCurrentBinding bool rebuilder RequestRebuilder dispatcher AttemptDispatcher preparers map[string]RecoveryPlanPreparer preparationTimeout time.Duration usedPlanIDs map[string]struct{} preparedPlanKeys map[string]struct{} sequencer *ObservationSequencer } // NewRecoveryCoordinator creates one request-local coordinator. func NewRecoveryCoordinator(options RecoveryCoordinatorOptions) (*RecoveryCoordinator, error) { if err := options.Policy.Validate(); err != nil { return nil, err } if err := options.Usage.ValidateAgainst(options.Policy); err != nil { return nil, err } if err := options.RequestSnapshot.Validate(); err != nil { return nil, err } if err := options.CurrentBinding.Validate(); err != nil { return nil, err } if options.Rebuilder == nil { return nil, errors.New("streamgate: recovery coordinator rebuilder is required") } if options.Dispatcher == nil { return nil, errors.New("streamgate: recovery coordinator dispatcher is required") } if options.PreparationTimeout < 0 { return nil, errors.New("streamgate: recovery coordinator preparation timeout must not be negative") } if len(options.Preparers) > 0 && options.PreparationTimeout <= 0 { return nil, errors.New("streamgate: recovery coordinator preparation timeout must be positive") } preparers := make(map[string]RecoveryPlanPreparer, len(options.Preparers)) for id, preparer := range options.Preparers { if err := validateStableTokenRequired("preparerID", id); err != nil { return nil, err } if preparer == nil { return nil, errors.New("streamgate: recovery coordinator preparer must not be nil") } preparers[id] = preparer } return &RecoveryCoordinator{ policy: copyRecoveryPolicySnapshot(options.Policy), usage: copyRecoveryUsageSnapshot(options.Usage), requestSnapshot: options.RequestSnapshot, currentBinding: options.CurrentBinding, hasCurrentBinding: true, rebuilder: options.Rebuilder, dispatcher: options.Dispatcher, preparers: preparers, preparationTimeout: options.PreparationTimeout, usedPlanIDs: make(map[string]struct{}), preparedPlanKeys: make(map[string]struct{}), sequencer: NewObservationSequencer(options.ObservationSink, nil), }, nil } // SetObservationSequencer sets or overrides the coordinator's ObservationSequencer. func (c *RecoveryCoordinator) SetObservationSequencer(seq *ObservationSequencer) { c.mu.Lock() defer c.mu.Unlock() c.sequencer = seq } // UsageSnapshot returns the current request-local fault dispatch ledger. func (c *RecoveryCoordinator) UsageSnapshot() RecoveryUsageSnapshot { c.mu.Lock() defer c.mu.Unlock() return copyRecoveryUsageSnapshot(c.usage) } // CurrentBinding returns the currently owned provider attempt, if any. func (c *RecoveryCoordinator) CurrentBinding() (AttemptBinding, bool) { c.mu.Lock() defer c.mu.Unlock() if !c.hasCurrentBinding { return AttemptBinding{}, false } return c.currentBinding, true } // UsedPlanIDs returns a stable sorted copy of the request-local plan ledger. func (c *RecoveryCoordinator) UsedPlanIDs() []string { c.mu.Lock() defer c.mu.Unlock() return c.usedPlanIDsLocked() } func (c *RecoveryCoordinator) usedPlanIDsLocked() []string { ids := make([]string, 0, len(c.usedPlanIDs)) for id := range c.usedPlanIDs { ids = append(ids, id) } sort.Strings(ids) return ids } func (c *RecoveryCoordinator) beginCycle(planID string) (AttemptBinding, RecoveryUsageSnapshot, []string, error) { c.mu.Lock() defer c.mu.Unlock() if c.busy { return AttemptBinding{}, RecoveryUsageSnapshot{}, nil, ErrRecoveryCycleInProgress } if _, used := c.usedPlanIDs[planID]; used { return AttemptBinding{}, RecoveryUsageSnapshot{}, nil, ErrRecoveryPlanReentry } if !c.hasCurrentBinding { return AttemptBinding{}, RecoveryUsageSnapshot{}, nil, ErrRecoveryNoCurrentAttempt } c.busy = true return c.currentBinding, copyRecoveryUsageSnapshot(c.usage), c.usedPlanIDsLocked(), nil } func (c *RecoveryCoordinator) finishCycle() { c.mu.Lock() c.busy = false c.mu.Unlock() } func (c *RecoveryCoordinator) markPlanUsed(planID string) { c.mu.Lock() c.usedPlanIDs[planID] = struct{}{} c.mu.Unlock() } func (c *RecoveryCoordinator) clearCurrentBinding(attemptID string) { c.mu.Lock() defer c.mu.Unlock() if c.hasCurrentBinding && c.currentBinding.AttemptID() == attemptID { c.currentBinding = AttemptBinding{} c.hasCurrentBinding = false } } func (c *RecoveryCoordinator) consumeUsage(usage RecoveryUsageSnapshot) { c.mu.Lock() c.usage = copyRecoveryUsageSnapshot(usage) c.mu.Unlock() } func (c *RecoveryCoordinator) installBinding(binding AttemptBinding) { c.mu.Lock() c.currentBinding = binding c.hasCurrentBinding = true c.mu.Unlock() } func (c *RecoveryCoordinator) claimPreparation(plan RecoveryPlan) bool { key := plan.PlanID() + "\x00" + plan.IdempotencyKey() c.mu.Lock() defer c.mu.Unlock() if _, exists := c.preparedPlanKeys[key]; exists { return false } c.preparedPlanKeys[key] = struct{}{} return true } func (c *RecoveryCoordinator) preparer(id string) RecoveryPlanPreparer { c.mu.Lock() defer c.mu.Unlock() return c.preparers[id] } // Execute performs exactly one recovery cycle. Ordering is fixed as: // eligibility/plan -> abort -> optional one-shot prepare -> rebuild/finalize -> // pre-dispatch usage consume -> single dispatch -> immutable binding install. func (c *RecoveryCoordinator) Execute(ctx context.Context, input RecoveryCycleInput) (result RecoveryCycleResult, retErr error) { result.failureCauses = input.FailureCauses.Copy() result.usage = c.UsageSnapshot() c.mu.Lock() seq := c.sequencer c.mu.Unlock() correlation, configGen, attemptID := "default", "default", "default" attemptTarget, _ := NewObservationAttemptTarget("default", "default", "default", "default") var epochID uint64 if ctx != nil { if obsCtx, ok := ObservationContextFromContext(ctx); ok { if obsCtx.Sequencer != nil { seq = obsCtx.Sequencer } if obsCtx.StableCorrelation != "" { correlation = obsCtx.StableCorrelation } if obsCtx.ConfigGeneration != "" { configGen = obsCtx.ConfigGeneration } if obsCtx.AttemptID != "" { attemptID = obsCtx.AttemptID } if err := obsCtx.AttemptTarget.Validate(); err == nil { attemptTarget = obsCtx.AttemptTarget } epochID = obsCtx.EpochID } } var snapshotReleased bool releasePreparationSnapshot := func() error { if input.PreparationSnapshot == nil || snapshotReleased { return nil } snapshotReleased = true if err := input.PreparationSnapshot.Release(); err != nil { return ErrRecoveryPreparationSnapshotRelease } return nil } defer func() { // Every result path below finalizes the snapshot before it emits a // terminal observation. This is only a safety net for future paths. _ = releasePreparationSnapshot() }() emitPlanRejected := func(failed RecoveryCycleResult) { if seq == nil { return } _, _ = seq.Emit(context.Background(), FilterObservationInput{ Kind: ObservationKindRecoveryPlanRejected, StableCorrelation: correlation, ConfigGeneration: configGen, AttemptID: attemptID, AttemptTarget: attemptTarget, EpochID: epochID, CommitState: input.CommitState, Causes: failed.FailureCauses(), OccurredAt: time.Now(), }) } finishPlanRejection := func(failed RecoveryCycleResult, cause error) (RecoveryCycleResult, error) { if err := releasePreparationSnapshot(); err != nil { failed = failed.withFailure("preparation_snapshot_release_failed") cause = errors.Join(cause, err) } emitPlanRejected(failed) return failed, cause } if err := input.Arbitration.Validate(); err != nil || input.Arbitration.Action() != ArbitrationActionRecover { return finishPlanRejection(result.withFailure("arbitration_ineligible"), ErrRecoveryArbitrationIneligible) } intent := input.Arbitration.RecoveryIntent() if intent == nil { return finishPlanRejection(result.withFailure("arbitration_ineligible"), ErrRecoveryArbitrationIneligible) } if err := validateStableTokenRequired("planID", input.PlanID); err != nil { return finishPlanRejection(result.withFailure("plan_ineligible"), err) } if ctx == nil { return finishPlanRejection(result.withFailure("caller_canceled"), ErrRecoveryCallerCanceled) } current, usage, usedPlanIDs, err := c.beginCycle(input.PlanID) if err != nil { if errors.Is(err, ErrRecoveryCycleInProgress) && attemptID == "default" { c.mu.Lock() if c.hasCurrentBinding { attemptID = c.currentBinding.AttemptID() } c.mu.Unlock() } code := "cycle_ineligible" if errors.Is(err, ErrRecoveryCycleInProgress) { code = "cycle_in_progress" } else if errors.Is(err, ErrRecoveryPlanReentry) { code = "plan_reentry" } else if errors.Is(err, ErrRecoveryNoCurrentAttempt) { code = "current_attempt_missing" } return finishPlanRejection(result.withFailure(code), err) } defer c.finishCycle() if attemptID == "default" { attemptID = current.AttemptID() } if attemptTarget.ModelGroup() == "default" { attemptTarget, _ = NewObservationAttemptTarget("default", current.Model(), current.Provider(), current.ExecutionPath()) } emitRecoveryObs := func(kind ObservationKind, plan RecoveryPlan, observedAttemptID string, observedTarget ObservationAttemptTarget, preparerInfo *ObservationPreparerInfo, causes FailureCauseChain) { if seq == nil { return } sharedAttemptID := "" if plan.ResumeMode() == RecoveryResumeModeContinueStream { sharedAttemptID = current.AttemptID() } recInfo, err := NewObservationRecoveryInfo(plan.PlanID(), plan.Strategy(), plan.ResumeMode(), sharedAttemptID) if err != nil { return } _, _ = seq.Emit(ctx, FilterObservationInput{ Kind: kind, StableCorrelation: correlation, ConfigGeneration: configGen, AttemptID: observedAttemptID, AttemptTarget: observedTarget, EpochID: epochID, CommitState: input.CommitState, Recovery: &recInfo, Preparer: preparerInfo, Causes: causes, OccurredAt: time.Now(), }) } callerCanceled := input.CallerCanceled || ctx.Err() != nil var preparationRef string if input.PreparerID != "" { if input.PreparationSnapshot == nil { return finishPlanRejection(result.withFailure("preparation_snapshot_missing"), ErrRecoveryPreparationInvalid) } preparationRef = input.PreparationSnapshot.SnapshotRef() } else if input.PreparationSnapshot != nil { return finishPlanRejection(result.withFailure("preparer_missing"), ErrRecoveryPreparationInvalid) } contributor, err := NewRecoveryContributor(input.ConsumerID, input.Arbitration.FilterID(), input.Arbitration.RuleID()) if err != nil { return finishPlanRejection(result.withFailure("plan_ineligible"), err) } plan, err := NewRecoveryPlan(RecoveryEligibilityInput{ PlanID: input.PlanID, IdempotencyKey: input.IdempotencyKey, Intent: *intent, Contributors: []RecoveryContributor{contributor}, CommitState: input.CommitState, CallerCanceled: callerCanceled, HasToolSideEffect: input.HasToolSideEffect, HasCompleteToolCall: input.HasCompleteToolCall, PreviouslyUsedPlanIDs: usedPlanIDs, Policy: c.policy, Usage: usage, ManagedTrajectory: input.ManagedTrajectory, TerminalReason: input.TerminalReason, PreparerID: input.PreparerID, PreparationSnapshotRef: preparationRef, RequiredCapabilities: input.RequiredCapabilities, FailureCauses: input.FailureCauses, }) if err != nil { return finishPlanRejection(result.withFailure("plan_ineligible"), err) } result = result.withPlan(plan) c.markPlanUsed(plan.PlanID()) emitRecoveryObs(ObservationKindRecoveryPlanSelected, plan, attemptID, attemptTarget, nil, result.FailureCauses()) // Closing current provider ownership is the first host side effect. if err := current.Controller().AbortAttempt(ctx); err != nil { result = result.withFailure("attempt_abort_failed") retErr := error(ErrRecoveryAbortFailed) if releaseErr := releasePreparationSnapshot(); releaseErr != nil { result = result.withFailure("preparation_snapshot_release_failed") retErr = errors.Join(retErr, releaseErr) } emitRecoveryObs(ObservationKindRecoveryAttemptAbortFailed, plan, attemptID, attemptTarget, nil, result.FailureCauses()) return result, retErr } result.previousAttemptClosed = true c.clearCurrentBinding(current.AttemptID()) emitRecoveryObs(ObservationKindRecoveryAttemptAborted, plan, attemptID, attemptTarget, nil, result.FailureCauses()) if ctx.Err() != nil { return result.withFailure("caller_canceled"), ErrRecoveryCallerCanceled } if plan.RequiresPreparation() { preparer := c.preparer(plan.PreparerID()) emitPreparationFailure := func(code string, deadlineOutcome ObservationDeadlineOutcome, retErr error) (RecoveryCycleResult, error) { result = result.withFailure(code) if err := releasePreparationSnapshot(); err != nil { result = result.withFailure("preparation_snapshot_release_failed") retErr = errors.Join(retErr, err) } preparerInfo, infoErr := NewObservationPreparerInfo(plan.PreparerID(), RecoveryPreparationRequired, deadlineOutcome) if infoErr == nil { emitRecoveryObs(ObservationKindRecoveryPrepared, plan, attemptID, attemptTarget, &preparerInfo, result.FailureCauses()) } return result, retErr } if preparer == nil || c.preparationTimeout <= 0 { return emitPreparationFailure("preparer_unavailable", ObservationDeadlineOutcomeFailed, ErrRecoveryPreparerUnavailable) } if input.PreparationSnapshot.SnapshotRef() != plan.PreparationSnapshotRef() { return emitPreparationFailure("preparation_snapshot_mismatch", ObservationDeadlineOutcomeFailed, ErrRecoveryPreparationInvalid) } if !c.claimPreparation(plan) { return emitPreparationFailure("preparation_reentry", ObservationDeadlineOutcomeFailed, ErrRecoveryPlanReentry) } prepareCtx, cancel := context.WithTimeout(ctx, c.preparationTimeout) directive, prepareErr := preparer.PrepareRecoveryPlan(prepareCtx, plan, input.PreparationSnapshot) prepareCtxErr := prepareCtx.Err() cancel() if prepareErr != nil { if errors.Is(prepareCtxErr, context.DeadlineExceeded) { return emitPreparationFailure("preparation_deadline", ObservationDeadlineOutcomeTimeout, ErrRecoveryPreparationDeadline) } if errors.Is(prepareCtxErr, context.Canceled) || ctx.Err() != nil { return emitPreparationFailure("caller_canceled", ObservationDeadlineOutcomeFailed, ErrRecoveryCallerCanceled) } return emitPreparationFailure("preparation_failed", ObservationDeadlineOutcomeFailed, ErrRecoveryPreparationFailed) } if prepareCtxErr != nil { if errors.Is(prepareCtxErr, context.DeadlineExceeded) { return emitPreparationFailure("preparation_deadline", ObservationDeadlineOutcomeTimeout, ErrRecoveryPreparationDeadline) } return emitPreparationFailure("caller_canceled", ObservationDeadlineOutcomeFailed, ErrRecoveryCallerCanceled) } preparedPlan, err := plan.WithPreparedDirective(directive) if err != nil { return emitPreparationFailure("preparation_invalid", ObservationDeadlineOutcomeFailed, ErrRecoveryPreparationInvalid) } plan = preparedPlan result = result.withPlan(plan) if err := releasePreparationSnapshot(); err != nil { return emitPreparationFailure("preparation_snapshot_release_failed", ObservationDeadlineOutcomeFailed, err) } preparerInfo, _ := NewObservationPreparerInfo(plan.PreparerID(), RecoveryPreparationPrepared, ObservationDeadlineOutcomeCompleted) emitRecoveryObs(ObservationKindRecoveryPrepared, plan, attemptID, attemptTarget, &preparerInfo, result.FailureCauses()) } if ctx.Err() != nil { return result.withFailure("caller_canceled"), ErrRecoveryCallerCanceled } draft, err := c.rebuilder.RebuildRequest(ctx, c.requestSnapshot, plan) if err != nil { result = result.withFailure("request_rebuild_failed") emitRecoveryObs(ObservationKindRecoveryRebuildFailed, plan, attemptID, attemptTarget, nil, result.FailureCauses()) return result, ErrRecoveryRebuildFailed } finalPlan, request, err := plan.FinalizeRebuiltRequest(draft) if err != nil { if errors.Is(err, ErrManagedTrajectoryExhausted) { return result.withFailure("managed_trajectory_exhausted"), ErrManagedTrajectoryExhausted } result = result.withFailure("rebuilt_request_invalid") emitRecoveryObs(ObservationKindRecoveryRebuildFailed, plan, attemptID, attemptTarget, nil, result.FailureCauses()) return result, ErrRecoveryRebuildInvalid } plan = finalPlan result = result.withPlan(plan) emitRecoveryObs(ObservationKindRecoveryRebuilt, plan, attemptID, attemptTarget, nil, result.FailureCauses()) if ctx.Err() != nil { return result.withFailure("caller_canceled"), ErrRecoveryCallerCanceled } usageAfterDispatch, err := plan.UsageAfterDispatch() if err != nil { return result.withFailure("dispatch_budget_ineligible"), err } // This write is intentionally immediately before the outbound call. It is // not rolled back when the dispatcher reports failure. c.consumeUsage(usageAfterDispatch) result.usage = copyRecoveryUsageSnapshot(usageAfterDispatch) binding, err := c.dispatcher.DispatchAttempt(ctx, request) if err != nil { result = result.withFailure("attempt_dispatch_failed") emitRecoveryObs(ObservationKindRecoveryDispatchFailed, plan, attemptID, attemptTarget, nil, result.FailureCauses()) return result, ErrRecoveryDispatchFailed } if err := binding.Validate(); err != nil { result = result.withFailure("attempt_binding_invalid") emitRecoveryObs(ObservationKindRecoveryDispatchFailed, plan, attemptID, attemptTarget, nil, result.FailureCauses()) return result, ErrRecoveryDispatchFailed } reboundTarget, targetErr := NewObservationAttemptTarget(attemptTarget.ModelGroup(), binding.Model(), binding.Provider(), binding.ExecutionPath()) if targetErr != nil { _ = binding.Controller().AbortAttempt(ctx) result = result.withFailure("attempt_target_invalid") emitRecoveryObs(ObservationKindRecoveryDispatchFailed, plan, attemptID, attemptTarget, nil, result.FailureCauses()) return result, ErrRecoveryDispatchFailed } c.installBinding(binding) result.binding = binding result.hasBinding = true emitRecoveryObs(ObservationKindRecoveryDispatched, plan, binding.AttemptID(), reboundTarget, nil, result.FailureCauses()) return result, nil }