返回 DeepSeek-Reasonix
model_settings_source.go
根目录 / internal / serve / model_settings_source.go
1 package serve
2
3 import (
4 "bytes"
5 "context"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "io"
10 "net/http"
11 "net/url"
12 "strconv"
13 "strings"
14 "time"
15
16 "reasonix/internal/config"
17 "reasonix/internal/control"
18 )
19
20 var modelSettingsSourceClient = &http.Client{Transport: &http.Transport{Proxy: nil}, Timeout: 20 * time.Second}
21
22 func requestModelSettingsSource(ctx context.Context, settings *config.ModelRuntimeSettings, request config.ModelSettingsSourceRequest) (config.ModelSettingsSourceResponse, error) {
23 var response config.ModelSettingsSourceResponse
24 endpoint, err := url.Parse(settings.ProxyURL)
25 if err != nil || endpoint.Scheme != "http" || (endpoint.Hostname() != "127.0.0.1" && endpoint.Hostname() != "::1") || endpoint.User != nil {
26 return response, fmt.Errorf("model settings source requires the local credential tunnel")
27 }
28 endpoint.Path, endpoint.RawQuery, endpoint.Fragment = "/model-settings-source", "", ""
29 body, err := json.Marshal(request)
30 if err != nil {
31 return response, err
32 }
33 req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint.String(), bytes.NewReader(body))
34 if err != nil {
35 return response, err
36 }
37 req.Header.Set("Authorization", "Bearer "+settings.SourceToken)
38 req.Header.Set("Content-Type", "application/json")
39 res, err := modelSettingsSourceClient.Do(req)
40 if err != nil {
41 return response, fmt.Errorf("Desktop model settings source is unavailable")
42 }
43 defer res.Body.Close()
44 if res.StatusCode != http.StatusOK {
45 return response, fmt.Errorf("Desktop cannot apply saved model settings (status %d); check its available models and connection", res.StatusCode)
46 }
47 if err := json.NewDecoder(io.LimitReader(res.Body, 4<<20)).Decode(&response); err != nil || response.Version != 1 || response.Revision == "" {
48 return response, fmt.Errorf("invalid Desktop model settings acknowledgement")
49 }
50 return response, nil
51 }
52
53 func sourceModelRef(settings *config.ModelRuntimeSettings, remoteRef string) string {
54 for source, target := range settings.References {
55 if target == remoteRef && strings.Contains(source, "/") {
56 return source
57 }
58 }
59 return remoteRef
60 }
61
62 // New HTTP turns and autonomous FIFO dispatch share this boundary. A source
63 // failure leaves the existing controller and queued work intact. Approvals,
64 // steers, and children remain on their already accepted runtime.
65 func (s *Server) refreshRunModelSettingsLocked(ctx context.Context) error {
66 if err := s.cachedModelApplicationFailureLocked(ctx); err != nil {
67 return err
68 }
69 err := s.refreshModelSettingsOwnerLocked(ctx, modelSettingsRuntimeOwner{
70 current: s.ctl, settings: &s.managedModels, offerID: &s.modelSettingsOfferID, apply: s.switchModelLocked,
71 })
72 if err != nil && !control.ModelReplacementBlocked(s.ctl()) {
73 s.recordModelApplicationFailureLocked(ctx, err)
74 }
75 return err
76 }
77
78 type modelSettingsRuntimeOwner struct {
79 current func() control.SessionAPI
80 settings **config.ModelRuntimeSettings
81 offerID *string
82 apply func(context.Context, string) error
83 }
84
85 func (s *Server) refreshModelSettingsOwnerLocked(ctx context.Context, owner modelSettingsRuntimeOwner) error {
86 for {
87 if err := ctx.Err(); err != nil {
88 return err
89 }
90 current := owner.current()
91 settings := (*owner.settings)
92 var offered *config.ModelRuntimeSettings
93 ref := current.ModelRef()
94 var sourceRequest config.ModelSettingsSourceRequest
95 if settings != nil && settings.SourceToken != "" {
96 offerID, err := config.NewModelSettingsOfferID()
97 if err != nil {
98 return err
99 }
100 endpoint, err := url.Parse(settings.ProxyURL)
101 if err != nil {
102 return fmt.Errorf("invalid model settings tunnel")
103 }
104 port, _ := strconv.Atoi(endpoint.Port())
105 status := s.modelSettingsStatusLocked()
106 sourceRequest = config.ModelSettingsSourceRequest{Mode: "prepare", OfferID: offerID, PreviousOfferID: (*owner.offerID), Model: sourceModelRef(settings, ref), AppliedRevision: settings.Revision, RemotePort: port, OwnedRevisions: status.OwnedRevisions, UnversionedOwners: status.UnversionedOwners}
107 (*owner.offerID) = offerID
108 sourceRequest.ModelSettingsOwnership = status.ModelSettingsOwnership
109 response, err := requestModelSettingsSource(ctx, settings, sourceRequest)
110 if err != nil {
111 return err
112 }
113 if response.Settings != nil {
114 if err := validateModelSettingsSourceOffer(offerID, response); err != nil {
115 return err
116 }
117 offered, ref = response.Settings, response.Ref
118 } else if response.Revision != settings.Revision {
119 return fmt.Errorf("changed model settings require a complete resolver")
120 }
121 }
122 needsApply := offered != nil && (offered.Revision != settings.Revision || ref != current.ModelRef())
123 if snapshot, ok := current.(interface {
124 ModelSettingsState() (string, string, error)
125 }); ok {
126 applied, desired, err := snapshot.ModelSettingsState()
127 if err != nil {
128 return err
129 }
130 needsApply = needsApply || applied != desired
131 }
132 var applyErr error
133 if needsApply {
134 if offered != nil {
135 (*owner.settings) = offered
136 }
137 applyErr = owner.apply(ctx, ref)
138 if applyErr != nil {
139 (*owner.settings) = settings
140 }
141 }
142 if settings != nil && settings.SourceToken != "" {
143 ackSettings := settings
144 if offered != nil {
145 ackSettings = offered
146 }
147 status := s.modelSettingsStatusLocked()
148 sourceRequest.Mode = "finish"
149 sourceRequest.ModelSettingsOwnership = status.ModelSettingsOwnership
150 sourceRequest.PreviousOfferID = ""
151 sourceRequest.OwnedRevisions, sourceRequest.UnversionedOwners = status.OwnedRevisions, status.UnversionedOwners
152 ack, ackErr := requestModelSettingsSource(ctx, ackSettings, sourceRequest)
153 if ackErr == nil {
154 (*owner.offerID) = ""
155 }
156 if applyErr != nil {
157 if sourceCandidateOvertaken(applyErr, ackErr, offered, ack.Revision) {
158 continue // discard an overtaken candidate before publishing it
159 }
160 return fmt.Errorf("saved model settings could not be applied: %w", applyErr)
161 }
162 if ackErr != nil {
163 return ackErr
164 }
165 if ack.Revision != (*owner.settings).Revision {
166 continue // a save overtook the candidate; no new run was admitted
167 }
168 }
169 if applyErr != nil {
170 return applyErr
171 }
172 if changed, err := runtimeModelSettingsChanged(owner.current()); err != nil {
173 return err
174 } else if changed {
175 continue
176 }
177 return nil
178 }
179 }
180
181 func sourceCandidateOvertaken(applyErr, ackErr error, offered *config.ModelRuntimeSettings, acknowledged string) bool {
182 return errors.Is(applyErr, control.ErrModelChoiceStale) && ackErr == nil && offered != nil && acknowledged != offered.Revision
183 }
184
185 func validateModelSettingsSourceOffer(id string, response config.ModelSettingsSourceResponse) error {
186 if response.Settings.OfferID != id || response.Settings.Revision != response.Revision || response.Ref == "" || response.Settings.SourceToken == "" {
187 return fmt.Errorf("model settings offer does not match its request")
188 }
189 return nil
190 }
191
192 func runtimeModelSettingsChanged(ctrl control.SessionAPI) (bool, error) {
193 snapshot, ok := ctrl.(interface {
194 ModelSettingsState() (string, string, error)
195 })
196 if !ok {
197 return false, nil
198 }
199 applied, desired, err := snapshot.ModelSettingsState()
200 return applied != desired, err
201 }
202
203 func (s *Server) beforeInboxDispatch(ctrl *control.Controller) (func(), error) {
204 if s.ctl() != ctrl {
205 return s.beforeDetachedInboxDispatch(ctrl)
206 }
207 s.bindMu.Lock()
208 if s.ctl() != ctrl {
209 s.bindMu.Unlock()
210 return nil, control.ErrInboxRuntimeUnpublished
211 }
212 if ctrl.Running() {
213 s.bindMu.Unlock()
214 return nil, control.ErrTurnRunning
215 }
216 ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
217 err := s.refreshRunModelSettingsLocked(ctx)
218 cancel()
219 if err != nil {
220 s.bindMu.Unlock()
221 return nil, err
222 }
223 if current := s.ctl(); current != ctrl {
224 s.bindMu.Unlock()
225 if replacement, ok := current.(*control.Controller); ok {
226 replacement.NotifyInboxRuntimeReady()
227 }
228 return nil, control.ErrInboxRuntimeUnpublished
229 }
230 return s.bindMu.Unlock, nil
231 }
232
232 lines GO