返回 DeepSeek-Reasonix
rebuild_subgraph.go
根目录 / internal / boot / rebuild_subgraph.go
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
328 lines GO