返回 DeepSeek-Reasonix
remote_model_settings.go
根目录 / desktop / remote_model_settings.go
1 package main
2
3 import (
4 "bytes"
5 "context"
6 "crypto/sha256"
7 "encoding/hex"
8 "encoding/json"
9 "errors"
10 "fmt"
11 "io"
12 "net/http"
13 "strings"
14 "time"
15
16 "reasonix/internal/config"
17 )
18
19 type remoteModelSettingsStatus struct {
20 Application *ModelApplicationDetails `json:"application,omitempty"`
21 config.ModelSettingsOwnership
22 Version int `json:"version"`
23 Revision string `json:"revision"`
24 Model string `json:"model"`
25 SessionPath string `json:"sessionPath"`
26 OwnedRevisions []string `json:"ownedRevisions"`
27 UnversionedOwners bool `json:"unversionedOwners"`
28 }
29
30 func credentialProxyScope(host, workspace string) string {
31 sum := sha256.Sum256(fmt.Appendf(nil, "%d:%s%s", len(host), host, workspace))
32 return hex.EncodeToString(sum[:])
33 }
34
35 // Status refreshes already run when foreground/background work settles. Read
36 // ownership under the manager gate, then reject receipts overtaken by source
37 // refreshes using the Serve-wide ownership sequence.
38 func (a *App) refreshRemoteModelOwnership(ctx context.Context, tabID string, client *http.Client, generation uint64) {
39 a.remoteMu.Lock()
40 manager, ok := a.remoteRuntime.(*desktopRemoteManager)
41 a.remoteMu.Unlock()
42 if !ok {
43 return
44 }
45 a.remoteTabMu.Lock()
46 tab := a.remoteTabs[tabID]
47 if tab == nil || tab.client != client || tab.gen != generation || tab.settings.revision == "" {
48 a.remoteTabMu.Unlock()
49 return
50 }
51 host, workspace, path := tab.ref.HostID, tab.ref.Workspace, tab.routing.currentPath
52 a.remoteTabMu.Unlock()
53 managed := manager.managed(host)
54 if managed == nil {
55 return
56 }
57 managed.serveMu.Lock()
58 defer managed.serveMu.Unlock()
59 if !manager.isCurrent(host, managed) {
60 return
61 }
62 manager.mu.Lock()
63 server := managed.serves[workspace]
64 base := ""
65 if server != nil {
66 base = server.view.LocalURL
67 }
68 manager.mu.Unlock()
69 if base == "" {
70 return
71 }
72 status, err := remoteModelSettingsRequest(ctx, client, base, path, nil)
73 if err == nil && manager.isCurrent(host, managed) {
74 if !a.reconcileCredentialProxyGenerations(host, workspace, status) {
75 return
76 }
77 // Serve may have applied a new snapshot itself for a browser or queued
78 // turn. Reflect that acknowledgement on the same bound Desktop target.
79 if cfg, loadErr := config.LoadModelRuntimeSnapshot("."); loadErr == nil {
80 for _, provider := range cfg.Providers {
81 for _, modelID := range provider.ModelList() {
82 ref := provider.Name + "/" + modelID
83 if remoteSnapshotProviderName(provider.Name, modelID)+"/"+modelID != status.Model || cfg.ModelRuntimeFingerprint(ref) != status.Revision {
84 continue
85 }
86 a.remoteTabMu.Lock()
87 if current := a.remoteTabs[tabID]; current == tab && current.client == client && current.gen == generation && current.routing.currentPath == path && status.SessionPath == path {
88 current.model, current.settings.revision = ref, status.Revision
89 current.settings.generation, current.settings.sessionPath = generation, path
90 current.settings.failure, current.settings.failureRevision = "", ""
91 }
92 a.remoteTabMu.Unlock()
93 }
94 }
95 }
96 }
97 }
98
99 // Each model has its own immutable tunnel token. Provider names are stable
100 // across revisions; only the virtual credentials and resolver content change.
101 func remoteSnapshotProviderName(provider, model string) string {
102 sum := sha256.Sum256([]byte(model))
103 return provider + "-" + hex.EncodeToString(sum[:6])
104 }
105
106 func (a *App) buildRemoteModelSettings(host, workspace, model string, remotePort int, cfg *config.Config) (*config.ModelRuntimeSettings, string, error) {
107 offerID, err := config.NewModelSettingsOfferID()
108 if err != nil {
109 return nil, "", err
110 }
111 return a.buildRemoteModelSettingsOffer(host, workspace, model, remotePort, cfg, offerID)
112 }
113
114 func (a *App) buildRemoteModelSettingsOffer(host, workspace, model string, remotePort int, cfg *config.Config, offerID string) (_ *config.ModelRuntimeSettings, _ string, buildErr error) {
115 defer func() {
116 if buildErr != nil {
117 a.finishCredentialProxyOffer(host, workspace, offerID)
118 }
119 }()
120 revision := cfg.ModelRuntimeFingerprint(model)
121 bundle := &config.ModelRuntimeSettings{OfferID: offerID, Revision: revision, ProxyURL: fmt.Sprintf("http://127.0.0.1:%d", remotePort), Credentials: map[string]string{}, Preferences: cfg.RuntimeModelPreferences()}
122 refs := map[string]string{}
123 for _, base := range cfg.Providers {
124 if !base.Configured() || !modelProviderAccessAllowed(cfg.Desktop.ProviderAccess, base.Name) {
125 continue
126 }
127 for _, modelID := range base.ModelList() {
128 ref := base.Name + "/" + modelID
129 entry, ok := cfg.ResolveModel(ref)
130 if !ok {
131 continue
132 }
133 route, err := a.applyCredentialProxySnapshot(host, workspace, ref, cfg, revision, offerID)
134 if err != nil {
135 return nil, "", err
136 }
137 p := *entry
138 p.Name = remoteSnapshotProviderName(base.Name, modelID)
139 p.ModelsURL, p.BalanceURL, p.APIKeyEnv = "", "", ""
140 p.Headers, p.ExtraBody = nil, nil
141 p.NoProxy = true
142 p.Model, p.Default, p.Models = modelID, "", nil
143 bundle.Providers = append(bundle.Providers, p)
144 bundle.Credentials[p.Name] = route.token
145 if ref == model {
146 bundle.SourceToken = route.token
147 }
148 refs[ref] = p.Name + "/" + modelID
149 if modelID == base.DefaultModel() {
150 refs[base.Name] = refs[ref]
151 }
152 }
153 }
154 mapRef := func(ref string) string {
155 if mapped := refs[ref]; mapped != "" {
156 return mapped
157 }
158 return ref
159 }
160 p := &bundle.Preferences
161 p.PlannerModel, p.VisionModel, p.WebSearchModel = mapRef(p.PlannerModel), mapRef(p.VisionModel), mapRef(p.WebSearchModel)
162 p.GuardianModel, p.RecoveryModel, p.SubagentModel = mapRef(p.GuardianModel), mapRef(p.RecoveryModel), mapRef(p.SubagentModel)
163 p.SubagentModels = map[string]string{}
164 for name, ref := range cfg.Agent.SubagentModels {
165 p.SubagentModels[name] = mapRef(ref)
166 }
167 if refs[model] == "" {
168 return nil, "", fmt.Errorf("no available remote model in saved settings")
169 }
170 bundle.References = refs
171 return bundle, refs[model], nil
172 }
173
174 func (a *App) finishCredentialProxyOffer(host, workspace, offerID string) {
175 if offerID == "" {
176 return
177 }
178 a.credProxyMu.Lock()
179 p := a.credProxy
180 a.credProxyMu.Unlock()
181 if p == nil {
182 return
183 }
184 p.mu.Lock()
185 defer p.mu.Unlock()
186 scope := credentialProxyScope(host, workspace)
187 for _, route := range p.routes {
188 if route.scope == scope {
189 delete(route.holds, offerID)
190 }
191 }
192 }
193
194 // updateMu serializes reservations; releasing one only needs the route lock.
195 // Refuse excess candidates without evicting a runtime or an accepted request.
196 func (p *credentialProxy) validateModelSettingsOfferCapacity(up proxyUpstream) error {
197 if up.offerID == "" {
198 return nil
199 }
200 p.mu.Lock()
201 defer p.mu.Unlock()
202 offers := map[string]bool{}
203 for _, route := range p.routes {
204 if route.scope != up.scope {
205 continue
206 }
207 if route.holds[up.offerID] {
208 return nil
209 }
210 for id := range route.holds {
211 offers[id] = true
212 }
213 }
214 if len(offers) >= 64 {
215 return fmt.Errorf("too many unacknowledged model settings offers for this session")
216 }
217 return nil
218 }
219
220 // A managed Serve can refresh at its own run boundary, including durable
221 // inbox dispatch. Authentication is the currently owned immutable route token;
222 // the request cannot choose a local workspace or disclose a real credential.
223 func (a *App) serveModelSettingsSource(w http.ResponseWriter, r *http.Request, route *credProxyRoute) {
224 if r.Method != http.MethodPost {
225 http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
226 return
227 }
228 var request config.ModelSettingsSourceRequest
229 decoder := json.NewDecoder(http.MaxBytesReader(w, r.Body, 64<<10))
230 decoder.DisallowUnknownFields()
231 if err := decoder.Decode(&request); err != nil || len(request.OfferID) != 32 || (request.Mode != "prepare" && request.Mode != "finish" && request.Mode != "inspect") {
232 http.Error(w, "invalid model settings request", http.StatusBadRequest)
233 return
234 }
235 if _, err := hex.DecodeString(request.OfferID); err != nil {
236 http.Error(w, "invalid model settings offer", http.StatusBadRequest)
237 return
238 }
239 if request.Mode == "inspect" {
240 cfg, err := config.LoadModelRuntimeSnapshot(".")
241 if err != nil {
242 http.Error(w, "cannot verify saved settings", http.StatusServiceUnavailable)
243 return
244 }
245 response := config.ModelSettingsSourceResponse{Version: 1, Revision: cfg.ModelRuntimeFingerprint(route.ref)}
246 if route.modelSnapshot != nil {
247 if entry, ok := route.modelSnapshot.ResolveModel(route.ref); ok {
248 response.ConnectionTarget = config.SafeModelConnectionTarget(config.ProviderEffectiveRequestURL(entry))
249 }
250 }
251 if request.AppliedRevision != route.revision {
252 response.ContinuationUnavailable = "configuration ownership changed"
253 } else if err := config.ValidateModelRuntimeContinuation(route.modelSnapshot, cfg); err != nil {
254 response.ContinuationUnavailable = err.Error()
255 } else {
256 response.CanContinue = true
257 }
258 w.Header().Set("Content-Type", "application/json")
259 _ = json.NewEncoder(w).Encode(response)
260 return
261 }
262 offers := []string{request.PreviousOfferID}
263 if request.Mode == "finish" {
264 offers = append(offers, request.OfferID)
265 }
266 if !a.reconcileCredentialProxyGenerations(route.host, route.workspace, remoteModelSettingsStatus{ModelSettingsOwnership: request.ModelSettingsOwnership, Version: 1, OwnedRevisions: request.OwnedRevisions, UnversionedOwners: request.UnversionedOwners}, offers...) {
267 http.Error(w, "stale model settings ownership", http.StatusConflict)
268 return
269 }
270 cfg, err := config.LoadModelRuntimeSnapshot(".")
271 if err != nil {
272 http.Error(w, "cannot read saved model settings", http.StatusServiceUnavailable)
273 return
274 }
275 model := request.Model
276 if model == "" {
277 model = route.ref
278 }
279 model, err = resolveModelSettingsRuntime(cfg, model)
280 if err != nil {
281 http.Error(w, "choose an available model in Desktop Settings before starting another run", http.StatusConflict)
282 return
283 }
284 response := config.ModelSettingsSourceResponse{Version: 1, Revision: cfg.ModelRuntimeFingerprint(model)}
285 if request.Mode == "prepare" && request.AppliedRevision != response.Revision {
286 if request.RemotePort < 1 || request.RemotePort > 65535 {
287 http.Error(w, "invalid credential tunnel port", http.StatusBadRequest)
288 return
289 }
290 response.Settings, response.Ref, err = a.buildRemoteModelSettingsOffer(route.host, route.workspace, model, request.RemotePort, cfg, request.OfferID)
291 if err != nil {
292 http.Error(w, "cannot prepare saved model settings", http.StatusConflict)
293 return
294 }
295 if !a.reserveCredentialProxyInstall(route.host, route.workspace, request.OfferID, response.Revision, request.OwnershipIncarnation) {
296 a.finishCredentialProxyOffer(route.host, route.workspace, request.OfferID)
297 http.Error(w, "model settings owner was replaced", http.StatusConflict)
298 return
299 }
300 }
301 w.Header().Set("Content-Type", "application/json")
302 _ = json.NewEncoder(w).Encode(response)
303 }
304
305 const remoteModelSettingsUpgradeHint = "saved model settings require a newer remote Serve; upgrade or safely reconnect after its current work finishes"
306
307 type remoteModelSettingsRejection struct {
308 message string
309 // unsupported marks rejections that mean the remote Serve predates the
310 // model-settings protocol entirely (no /model-settings route at all).
311 unsupported bool
312 }
313
314 func (e *remoteModelSettingsRejection) Error() string { return e.message }
315
316 // isRemoteModelSettingsUnsupported reports whether err means the remote Serve
317 // cannot speak the model-settings protocol, as opposed to refusing a specific
318 // snapshot (ownership conflicts, busy turns, and other transient rejections).
319 func isRemoteModelSettingsUnsupported(err error) bool {
320 var rejected *remoteModelSettingsRejection
321 return errors.As(err, &rejected) && rejected.unsupported
322 }
323
324 func remoteModelSettingsRequest(ctx context.Context, client *http.Client, base, expectedPath string, body any) (remoteModelSettingsStatus, error) {
325 var result remoteModelSettingsStatus
326 method := http.MethodGet
327 var data []byte
328 if body != nil {
329 method = http.MethodPost
330 var err error
331 data, err = json.Marshal(body)
332 if err != nil {
333 return result, err
334 }
335 }
336 req, err := http.NewRequestWithContext(ctx, method, serveURL(base, "/model-settings"), bytes.NewReader(data))
337 if err != nil {
338 return result, err
339 }
340 req.Header.Set("Content-Type", "application/json")
341 if expectedPath != "" {
342 req.Header.Set(expectedSessionPathHeader, expectedPath)
343 }
344 resp, err := client.Do(req)
345 if err != nil {
346 return result, err
347 }
348 defer resp.Body.Close()
349 if resp.StatusCode == http.StatusNotFound || resp.StatusCode == http.StatusMethodNotAllowed {
350 return result, &remoteModelSettingsRejection{message: remoteModelSettingsUpgradeHint, unsupported: true}
351 }
352 if resp.StatusCode != http.StatusOK {
353 detail, _ := io.ReadAll(io.LimitReader(resp.Body, 2048))
354 return result, &remoteModelSettingsRejection{message: fmt.Sprintf("remote model settings status %d: %s", resp.StatusCode, strings.TrimSpace(string(detail)))}
355 }
356 payload, err := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
357 if err != nil {
358 return result, err
359 }
360 if err := json.Unmarshal(payload, &result); err != nil {
361 // Older Serves answer unknown routes with the HTML index; classify that
362 // document response as an unsupported model-settings protocol.
363 if trimmed := bytes.TrimSpace(payload); len(trimmed) > 0 && trimmed[0] == '<' {
364 return result, &remoteModelSettingsRejection{message: remoteModelSettingsUpgradeHint, unsupported: true}
365 }
366 return result, fmt.Errorf("decode remote model settings status: %w", err)
367 }
368 if result.Version != 1 {
369 return result, fmt.Errorf("remote Serve does not support immutable model settings")
370 }
371 return result, nil
372 }
373
374 // ensureRemoteModelSettings is the Desktop remote run-admission boundary.
375 // Approval/ask/steer replies bypass it because they belong to the accepted run.
376 // The returned generation scopes the admission (revision or legacy skip) to one
377 // tab generation; callers must re-admit before a target that no longer runs it.
378 func (a *App) ensureRemoteModelSettings(tabID string) (string, uint64, error) {
379 if !a.remoteTabLocalProxy(tabID) {
380 return "", 0, nil
381 }
382 for {
383 cfg, err := config.LoadModelRuntimeSnapshot(".")
384 if err != nil {
385 return "", 0, err
386 }
387 a.remoteTabMu.Lock()
388 tab := a.remoteTabs[tabID]
389 if tab == nil {
390 a.remoteTabMu.Unlock()
391 return "", 0, fmt.Errorf("remote session is no longer available")
392 }
393 model, applied, generation, path := tab.model, tab.settings.revision, tab.gen, tab.routing.currentPath
394 valid := tab.settings.generation == generation && tab.settings.sessionPath == path
395 // Generation 0 is an unrecorded verdict, not a legacy one: fresh and
396 // restored tabs run generation 0 before their first attachment.
397 unsupported := tab.settings.unsupportedGen != 0 && tab.settings.unsupportedGen == generation
398 a.remoteTabMu.Unlock()
399 if unsupported {
400 return "", generation, nil
401 }
402 if model == "" {
403 model = resolveNewSessionModel(cfg)
404 }
405 desired := cfg.ModelRuntimeFingerprint(model)
406 if valid && applied == desired {
407 return applied, generation, nil
408 }
409 next, err := resolveModelSettingsRuntime(cfg, model)
410 if err == nil {
411 err = a.SetRemoteTabModel(tabID, next)
412 }
413 if err != nil {
414 if isRemoteModelSettingsUnsupported(err) {
415 // A legacy Serve cannot apply snapshots; admit without a revision
416 // only when the verdict belongs to the current connection.
417 a.remoteTabMu.Lock()
418 current := a.remoteTabs[tabID]
419 recorded := current == tab && current.gen == generation && current.routing.currentPath == path
420 if recorded {
421 current.settings.unsupportedGen = generation
422 }
423 a.remoteTabMu.Unlock()
424 if !recorded {
425 // A reconnect replaced the probe's target mid-flight; the
426 // replacement Serve may speak the protocol, so re-probe it
427 // instead of admitting against the retired fence.
428 continue
429 }
430 return "", generation, nil
431 }
432 a.remoteTabMu.Lock()
433 if current := a.remoteTabs[tabID]; current == tab && current.gen == generation && current.routing.currentPath == path {
434 current.settings.failure = modelSettingsIssue("apply_failed", err).Message
435 current.settings.failureRevision = desired
436 }
437 a.remoteTabMu.Unlock()
438 if detail := a.remoteModelApplicationDetails(tabID); detail != nil {
439 return "", 0, &modelApplicationError{cause: err, details: detail}
440 }
441 return "", 0, fmt.Errorf("model settings were saved but the remote session could not apply them: %w", err)
442 }
443 a.remoteTabMu.Lock()
444 acknowledged := tab.settings.revision != "" && tab.settings.generation == tab.gen && tab.settings.sessionPath == tab.routing.currentPath
445 a.remoteTabMu.Unlock()
446 if !acknowledged {
447 return "", 0, fmt.Errorf("remote session did not acknowledge the saved model settings")
448 }
449 // Read the current file again: another save may have won during the
450 // remote build. Never admit a new request with that stale completion.
451 }
452 }
453
454 func (a *App) appendRemoteModelSettingsStatus(result *ModelSettingsResult) {
455 cfg, err := config.LoadModelRuntimeSnapshot(".")
456 if err != nil {
457 return
458 } // the global read failure is already in the result
459 a.remoteTabMu.Lock()
460 defer a.remoteTabMu.Unlock()
461 for _, tab := range a.remoteTabs {
462 if tab == nil {
463 continue
464 }
465 host, ok := cfg.RemoteHost(tab.ref.HostID)
466 if !ok || !host.CredentialProxyEnabled() {
467 continue
468 }
469 model := tab.model
470 if model == "" {
471 model = resolveNewSessionModel(cfg)
472 }
473 desired := cfg.ModelRuntimeFingerprint(model)
474 // Generation 0 is an unrecorded verdict: a not-yet-attached tab has not
475 // probed any Serve and must stay pending, not claim not_required.
476 if tab.settings.unsupportedGen != 0 && tab.settings.unsupportedGen == tab.gen {
477 // Legacy targets never apply snapshots, so report them as not required
478 // instead of leaving the settings receipt pending indefinitely.
479 result.Targets = append(result.Targets, ModelSettingsTarget{TabID: tab.id, Title: tab.topicTitle, Application: "not_required", AppliedRevision: tab.settings.revision, DesiredRevision: desired})
480 continue
481 }
482 state := "applied"
483 if tab.settings.revision != desired || tab.settings.generation != tab.gen || tab.settings.sessionPath != tab.routing.currentPath {
484 state = "pending"
485 blocked := tab.settings.details != nil && tab.settings.details.Code == "model_settings_pending"
486 if tab.settings.failure != "" && tab.settings.failureRevision == desired && !blocked {
487 state = "failed"
488 result.Application = "failed"
489 result.Issues = append(result.Issues, ModelSettingsIssue{Code: "remote_apply_failed", Message: tab.settings.failure})
490 } else if result.Application != "failed" {
491 result.Application = "pending"
492 }
493 }
494 target := ModelSettingsTarget{TabID: tab.id, Title: tab.topicTitle, Application: state, AppliedRevision: tab.settings.revision, DesiredRevision: desired}
495 if state != "applied" {
496 target.Details = tab.settings.details
497 }
498 result.Targets = append(result.Targets, target)
499 }
500 }
501
502 // remoteModelApplicationState is scoped to one acknowledged session binding.
503 type remoteModelApplicationState struct {
504 details *ModelApplicationDetails
505 detailsGeneration uint64
506 detailsPath string
507 revision string
508 failure string
509 failureRevision string
510 generation uint64
511 sessionPath string
512 // unsupportedGen records a generation whose Serve predates model-settings;
513 // a reconnect or replacement generation probes the protocol again.
514 unsupportedGen uint64
515 }
516
517 func applyRemoteModelSettingsSnapshot(ctx context.Context, client *http.Client, base, expectedPath, remoteRef string, bundle *config.ModelRuntimeSettings, status remoteModelSettingsStatus) (remoteModelSettingsStatus, error) {
518 var err error
519 if status.Revision != bundle.Revision || status.Model != remoteRef {
520 body := map[string]any{"version": 1, "ref": remoteRef, "settings": bundle}
521 status, err = remoteModelSettingsRequest(ctx, client, base, expectedPath, body)
522 if err != nil {
523 readCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
524 defer cancel()
525 observed, readErr := remoteModelSettingsRequest(readCtx, client, base, expectedPath, nil)
526 if readErr != nil {
527 return status, err
528 }
529 if observed.Revision != bundle.Revision || observed.Model != remoteRef {
530 return observed, err
531 }
532 status = observed
533 }
534 }
535 if status.Revision != bundle.Revision || status.Model != remoteRef {
536 return status, fmt.Errorf("remote model settings acknowledgement did not match the submitted snapshot")
537 }
538 return status, nil
539 }
540
540 lines GO