| 1 | package boot |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "fmt" |
| 6 | "time" |
| 7 | |
| 8 | "reasonix/internal/config" |
| 9 | "reasonix/internal/control" |
| 10 | "reasonix/internal/extension" |
| 11 | "reasonix/internal/extension/dispatch" |
| 12 | "reasonix/internal/extension/protocol" |
| 13 | "reasonix/internal/extension/sidecar" |
| 14 | "reasonix/internal/extension/uihub" |
| 15 | "reasonix/internal/provider" |
| 16 | ) |
| 17 | |
| 18 | // tryRebuildSubgraph patches narrow plans without BuildRuntime (fail-atomic). |
| 19 | // Callers must skip Close when BuildResult.ReusedController is set. |
| 20 | func tryRebuildSubgraph(ctx context.Context, old *control.Controller, previous *BuildResult, opts Options, m runtimeMigration) (res *BuildResult, handled bool, err error) { |
| 21 | if previous == nil || previous.Snapshot == nil || old == nil { |
| 22 | return nil, false, nil |
| 23 | } |
| 24 | start := time.Now() |
| 25 | from := opts.Graph |
| 26 | if previous.Plan != nil && previous.Plan.Graph != nil { |
| 27 | from = previous.Plan.Graph |
| 28 | } |
| 29 | graphStart := time.Now() |
| 30 | to, gerr := buildRuntimeGraph(config.ReasonixHomeDir(), nil) |
| 31 | extension.DefaultLifecycleMetrics.ObserveGraphBuild(time.Since(graphStart)) |
| 32 | if gerr != nil { |
| 33 | return nil, false, nil |
| 34 | } |
| 35 | diffStart := time.Now() |
| 36 | plan := extension.DiffRuntimePlan(from, to, opts.Generation, 0) |
| 37 | extension.DefaultLifecycleMetrics.ObservePlanDiff(time.Since(diffStart)) |
| 38 | if plan == nil { |
| 39 | return nil, false, nil |
| 40 | } |
| 41 | |
| 42 | switch plan.Kind { |
| 43 | case extension.SubgraphNone: |
| 44 | extension.DefaultLifecycleMetrics.NoOpRebuilds.Add(1) |
| 45 | case extension.SubgraphInterceptorOnly, extension.SubgraphUIOnly, extension.SubgraphProviderOnly, extension.SubgraphMCPOnly: |
| 46 | extension.DefaultLifecycleMetrics.SubgraphRebuilds.Add(1) |
| 47 | default: |
| 48 | extension.DefaultLifecycleMetrics.FullRebuilds.Add(1) |
| 49 | return nil, false, nil |
| 50 | } |
| 51 | |
| 52 | gen := nextRuntimeGeneration() |
| 53 | plan.ToGeneration = gen |
| 54 | plan.Graph = to |
| 55 | |
| 56 | // Checkpoint bindings for fail-atomic restore. |
| 57 | prevDispatcher := previous.Dispatcher |
| 58 | prevResolver := previous.ProviderResolver |
| 59 | prevUI := previous.ExtensionUI |
| 60 | prevUISession := controllerSessionID(previous.Controller) |
| 61 | prevUIGen := previous.Snapshot.Generation() |
| 62 | res = &BuildResult{ |
| 63 | Controller: previous.Controller, |
| 64 | Snapshot: previous.Snapshot.WithGeneration(gen), |
| 65 | Runtime: extension.NewRuntimeSet(gen), |
| 66 | Owner: opts.Owner, |
| 67 | Extensions: previous.Extensions, |
| 68 | Dispatcher: previous.Dispatcher, |
| 69 | ExtensionUI: previous.ExtensionUI, |
| 70 | ProviderResolver: previous.ProviderResolver, |
| 71 | BaseProviderResolver: previous.BaseProviderResolver, |
| 72 | Assembly: previous.Assembly, |
| 73 | SkillWatchService: previous.SkillWatchService, |
| 74 | Plan: plan, |
| 75 | ReusedController: true, |
| 76 | } |
| 77 | session := protocol.SessionContext{ |
| 78 | SessionID: controllerSessionID(previous.Controller), |
| 79 | WorkspaceRoot: previous.Controller.WorkspaceRoot(), |
| 80 | Generation: gen, |
| 81 | } |
| 82 | if session.SessionID == "" { |
| 83 | session.SessionID = "session" |
| 84 | } |
| 85 | if session.WorkspaceRoot == "" { |
| 86 | session.WorkspaceRoot = "." |
| 87 | } |
| 88 | |
| 89 | oldMgr := previous.Extensions |
| 90 | fail := func(stageErr error) (*BuildResult, bool, error) { |
| 91 | // Restore controller to pre-patch bindings (no partial commit should remain). |
| 92 | restoreControllerBindings(previous.Controller, prevDispatcher, prevResolver, prevUI, prevUISession, prevUIGen, oldMgr) |
| 93 | if res.Runtime != nil { |
| 94 | _ = res.Runtime.Close() |
| 95 | } |
| 96 | if res.Extensions != nil && res.Extensions != oldMgr { |
| 97 | res.Extensions.RollbackPlanStart(oldMgr) |
| 98 | } |
| 99 | return nil, true, stageErr |
| 100 | } |
| 101 | |
| 102 | var patchErr error |
| 103 | switch plan.Kind { |
| 104 | case extension.SubgraphNone: |
| 105 | if res.Extensions != nil { |
| 106 | _ = res.Runtime.Track(extension.Effect{ |
| 107 | ID: "sidecar-manager-adopted", Owner: "boot", Class: extension.Cancelable, |
| 108 | Dispose: func(context.Context) error { return nil }, |
| 109 | }) |
| 110 | } |
| 111 | case extension.SubgraphUIOnly, extension.SubgraphInterceptorOnly, extension.SubgraphProviderOnly, extension.SubgraphMCPOnly: |
| 112 | // Stage only: does not mutate controller dispatcher/resolver/UI. |
| 113 | patchErr = stageSidecarSubgraph(ctx, res, oldMgr, plan, session, gen) |
| 114 | } |
| 115 | if patchErr != nil { |
| 116 | return fail(patchErr) |
| 117 | } |
| 118 | |
| 119 | if err := awaitSidecarsReady(ctx, res.Extensions); err != nil { |
| 120 | return fail(err) |
| 121 | } |
| 122 | |
| 123 | // Commit after staging + ready. If this panics mid-way, fail restores. |
| 124 | if err := commitControllerExtPatch(res, session, gen); err != nil { |
| 125 | return fail(err) |
| 126 | } |
| 127 | |
| 128 | _ = m |
| 129 | attachPlanAndStatus(res, from, to, opts.Generation, previous.Snapshot) |
| 130 | |
| 131 | if prevGen := previous.Snapshot.Generation(); prevGen != 0 && prevGen != gen { |
| 132 | registerControllerDrainCancel(res.Owner, prevGen, old) |
| 133 | if !res.ReusedController { |
| 134 | if host := old.Host(); host != nil { |
| 135 | h := host |
| 136 | res.Owner.Gate.RegisterDrainCancel(prevGen, func() { h.CancelInFlightMCP() }) |
| 137 | } |
| 138 | } |
| 139 | } |
| 140 | finishRebuildPublish(res, nil, start) |
| 141 | if oldMgr != nil && res.Extensions != oldMgr && res.Plan != nil { |
| 142 | drainStart := time.Now() |
| 143 | oldMgr.DrainPlan(res.Plan) |
| 144 | extension.DefaultLifecycleMetrics.ObserveDrain(time.Since(drainStart)) |
| 145 | } |
| 146 | return res, true, nil |
| 147 | } |
| 148 | |
| 149 | func restoreControllerBindings(ctrl *control.Controller, disp *dispatch.Dispatcher, resolver provider.Resolver, ui *uihub.Hub, uiSession string, uiGen uint64, oldMgr *sidecar.Manager) { |
| 150 | if ctrl == nil { |
| 151 | return |
| 152 | } |
| 153 | if disp != nil { |
| 154 | ctrl.ReplaceExtensions(disp) |
| 155 | } |
| 156 | ctrl.SetProviderResolver(resolver) |
| 157 | if ui != nil { |
| 158 | ui.BindGeneration(uiSession, uiGen) |
| 159 | if oldMgr != nil { |
| 160 | bindExtensionUI(ui, oldMgr, func(string) {}) |
| 161 | } |
| 162 | ctrl.SetExtensionUI(ui) |
| 163 | } |
| 164 | } |
| 165 | |
| 166 | func (res *BuildResult) ensureRuntime(gen uint64) *extension.RuntimeSet { |
| 167 | if res == nil { |
| 168 | return nil |
| 169 | } |
| 170 | if res.Runtime == nil || res.Runtime.Generation() != gen { |
| 171 | res.Runtime = extension.NewRuntimeSet(gen) |
| 172 | } |
| 173 | return res.Runtime |
| 174 | } |
| 175 | |
| 176 | // stageSidecarSubgraph prepares manager/snapshot/dispatcher/resolver without |
| 177 | // mutating live controller bindings. BindGeneration and stream-router install |
| 178 | // wait for commit; staged-generation host/ui/* is dropped until then. |
| 179 | func stageSidecarSubgraph(ctx context.Context, res *BuildResult, oldMgr *sidecar.Manager, plan *extension.RuntimePlan, session protocol.SessionContext, gen uint64) error { |
| 180 | home := config.ReasonixHomeDir() |
| 181 | var ui sidecar.UIHandler |
| 182 | if res.ExtensionUI != nil { |
| 183 | ui = res.ExtensionUI |
| 184 | } |
| 185 | mgr, _, err := sidecar.StartPackagesWithPlan(ctx, home, session, ui, oldMgr, plan) |
| 186 | if err != nil { |
| 187 | if mgr != nil { |
| 188 | mgr.RollbackPlanStart(oldMgr) |
| 189 | } |
| 190 | return err |
| 191 | } |
| 192 | res.Extensions = mgr |
| 193 | rs := res.ensureRuntime(gen) |
| 194 | closeOnDispose := mgr != oldMgr |
| 195 | if err := rs.Track(extension.Effect{ |
| 196 | ID: "sidecar-manager", |
| 197 | Owner: "boot", |
| 198 | Component: "extension-runtimes", |
| 199 | Class: extension.Cancelable, |
| 200 | Dispose: func(context.Context) error { |
| 201 | if closeOnDispose { |
| 202 | return mgr.Close() |
| 203 | } |
| 204 | return nil |
| 205 | }, |
| 206 | }); err != nil { |
| 207 | if closeOnDispose { |
| 208 | mgr.RollbackPlanStart(oldMgr) |
| 209 | } |
| 210 | return err |
| 211 | } |
| 212 | _ = extension.TrackUIHub(rs.Scope(), gen) |
| 213 | for _, client := range mgr.Clients() { |
| 214 | _ = extension.TrackEventSubscription(rs.Scope(), "sidecar:"+client.PluginID(), func() error { return nil }) |
| 215 | } |
| 216 | if res.Snapshot != nil { |
| 217 | res.Snapshot = res.Snapshot.WithLiveContributions(gen, mgr.Contributions()) |
| 218 | } |
| 219 | if err := stageDispatcher(res); err != nil { |
| 220 | return err |
| 221 | } |
| 222 | base := res.BaseProviderResolver |
| 223 | if base == nil { |
| 224 | base = res.ProviderResolver |
| 225 | } |
| 226 | var claims map[extension.Slot]extension.ContributionSource |
| 227 | if res.Snapshot != nil { |
| 228 | claims = res.Snapshot.Replacements() |
| 229 | } |
| 230 | merged, merr := mergeSidecarProviders(base, mgr, claims, res.Owner) |
| 231 | if merr != nil { |
| 232 | return merr |
| 233 | } |
| 234 | if merged != nil { |
| 235 | res.ProviderResolver = merged |
| 236 | } |
| 237 | return nil |
| 238 | } |
| 239 | |
| 240 | func stageDispatcher(res *BuildResult) error { |
| 241 | if res == nil || res.Snapshot == nil { |
| 242 | return nil |
| 243 | } |
| 244 | clients := sidecarClientResolver(res.Extensions) |
| 245 | required := requiredRuntimeSet(res.Extensions) |
| 246 | res.Dispatcher = dispatch.New(res.Snapshot.InterceptorChain(), res.Snapshot.Replacements(), clients, required, dispatch.Options{}) |
| 247 | return nil |
| 248 | } |
| 249 | |
| 250 | func commitControllerExtPatch(res *BuildResult, session protocol.SessionContext, gen uint64) error { |
| 251 | if res == nil || res.Controller == nil { |
| 252 | return nil |
| 253 | } |
| 254 | if res.Dispatcher != nil { |
| 255 | res.Controller.ReplaceExtensions(res.Dispatcher) |
| 256 | } |
| 257 | if res.ProviderResolver != nil { |
| 258 | res.Controller.SetProviderResolver(res.ProviderResolver) |
| 259 | // Install stream routers only at commit so stage/ready failure never |
| 260 | // leaves Unchanged clients pointing at a discarded generation resolver. |
| 261 | installSidecarStreamRouters(res.Extensions, res.ProviderResolver) |
| 262 | } |
| 263 | if res.ExtensionUI != nil { |
| 264 | if res.Extensions != nil { |
| 265 | bindExtensionUI(res.ExtensionUI, res.Extensions, func(string) {}) |
| 266 | } |
| 267 | res.ExtensionUI.BindGeneration(session.SessionID, gen) |
| 268 | res.Controller.SetExtensionUI(res.ExtensionUI) |
| 269 | } |
| 270 | return nil |
| 271 | } |
| 272 | |
| 273 | func awaitSidecarsReady(ctx context.Context, mgr *sidecar.Manager) error { |
| 274 | if mgr == nil { |
| 275 | return extension.AwaitReady(ctx, nil) |
| 276 | } |
| 277 | ready := make(chan struct{}) |
| 278 | go func() { |
| 279 | _ = mgr.Clients() |
| 280 | close(ready) |
| 281 | }() |
| 282 | return extension.AwaitReady(ctx, ready) |
| 283 | } |
| 284 | |
| 285 | func controllerSessionID(c *control.Controller) string { |
| 286 | if c == nil { |
| 287 | return "" |
| 288 | } |
| 289 | if p := c.SessionPath(); p != "" { |
| 290 | return p |
| 291 | } |
| 292 | return c.WorkspaceRoot() |
| 293 | } |
| 294 | |
| 295 | func registerControllerDrainCancel(owner *extension.RuntimeOwner, gen uint64, ctrl *control.Controller) { |
| 296 | if gen == 0 || ctrl == nil { |
| 297 | return |
| 298 | } |
| 299 | if owner == nil { |
| 300 | owner = extension.RuntimeOwnerOrDefault(nil) |
| 301 | } |
| 302 | owner.Gate.RegisterDrainCancel(gen, func() { |
| 303 | if ctrl.RuntimeGeneration() == gen || ctrl.RuntimeGeneration() == 0 { |
| 304 | ctrl.Cancel() |
| 305 | } |
| 306 | }) |
| 307 | } |
| 308 | |
| 309 | func finishRebuildPublish(res *BuildResult, drainMgr *sidecar.Manager, start time.Time) { |
| 310 | publishBuildResult(res) |
| 311 | if drainMgr != nil && res != nil && res.Plan != nil && !res.ReusedController { |
| 312 | drainStart := time.Now() |
| 313 | drainMgr.DrainPlan(res.Plan) |
| 314 | extension.DefaultLifecycleMetrics.ObserveDrain(time.Since(drainStart)) |
| 315 | } |
| 316 | if res != nil && res.Runtime != nil && res.Snapshot != nil { |
| 317 | _ = res.Runtime.Track(extension.Effect{ |
| 318 | ID: fmt.Sprintf("rebuild-publish-%d", res.Snapshot.Generation()), |
| 319 | Owner: "boot", |
| 320 | Class: extension.Irreversible, |
| 321 | Dispose: func(context.Context) error { |
| 322 | return nil |
| 323 | }, |
| 324 | }) |
| 325 | } |
| 326 | extension.DefaultLifecycleMetrics.ObserveActivate(time.Since(start)) |
| 327 | } |
| 328 |