返回 DeepSeek-Reasonix
httpexec.go
根目录 / internal / browser / httpexec.go
1 package browser
2
3 import (
4 "bytes"
5 "context"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "io"
10 "net/http"
11 "strings"
12 "sync"
13 "time"
14 )
15
16 const httpHealthTTL = 30 * time.Second
17
18 // httpExecutor is the JSON-over-HTTP client half of the contract. It never
19 // retries: a write whose reply was lost is reported as ErrUnknownOutcome, and
20 // the tools tell the model not to try again.
21 type httpExecutor struct {
22 endpoint string
23 token string
24 client *http.Client
25 now func() time.Time
26
27 mu sync.Mutex
28 healthyAt time.Time
29 }
30
31 // NewHTTPExecutor returns an Executor that forwards every call to the
32 // broker at endpoint with a bearer token. A nil client uses
33 // http.DefaultClient; the caller decides timeouts through ctx.
34 func NewHTTPExecutor(endpoint, token string, client *http.Client) Executor {
35 if client == nil {
36 client = http.DefaultClient
37 }
38 return &httpExecutor{
39 endpoint: strings.TrimRight(strings.TrimSpace(endpoint), "/"),
40 token: strings.TrimSpace(token),
41 client: client,
42 now: time.Now,
43 }
44 }
45
46 // Available reports whether a health probe succeeded within the last 30 s,
47 // probing again when the cache is cold or expired.
48 func (e *httpExecutor) Available(ctx context.Context) bool {
49 e.mu.Lock()
50 fresh := !e.healthyAt.IsZero() && e.now().Sub(e.healthyAt) < httpHealthTTL
51 e.mu.Unlock()
52 if fresh {
53 return true
54 }
55 req, err := e.newRequest(ctx, http.MethodGet, e.endpoint+httpHealthRoute, nil)
56 if err != nil {
57 return false
58 }
59 resp, err := e.client.Do(req)
60 if err != nil {
61 return false
62 }
63 _, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, 4<<10))
64 _ = resp.Body.Close()
65 if resp.StatusCode != http.StatusNoContent && resp.StatusCode != http.StatusOK {
66 return false
67 }
68 e.mu.Lock()
69 e.healthyAt = e.now()
70 e.mu.Unlock()
71 return true
72 }
73
74 func (e *httpExecutor) newRequest(ctx context.Context, method, url string, body []byte) (*http.Request, error) {
75 req, err := http.NewRequestWithContext(ctx, method, url, bytes.NewReader(body))
76 if err != nil {
77 return nil, err
78 }
79 req.Header.Set("Authorization", "Bearer "+e.token)
80 if body != nil {
81 req.Header.Set("Content-Type", "application/json")
82 }
83 if id := SessionFromContext(ctx); id != "" {
84 req.Header.Set(SessionHeader, id)
85 }
86 return req, nil
87 }
88
89 // call posts in as JSON to /v1/browser/<method> and decodes the reply into
90 // out. Transport failures (no HTTP reply at all) come back as errTransport
91 // so Act can turn them into ErrUnknownOutcome.
92 func (e *httpExecutor) call(ctx context.Context, method string, in, out any) error {
93 body, err := json.Marshal(in)
94 if err != nil {
95 return fmt.Errorf("browser broker: encode %s: %w", method, err)
96 }
97 req, err := e.newRequest(ctx, http.MethodPost, e.endpoint+httpRoutePrefix+method, body)
98 if err != nil {
99 return fmt.Errorf("browser broker: %s: %w", method, err)
100 }
101 resp, err := e.client.Do(req)
102 if err != nil {
103 return &transportError{method: method, err: err}
104 }
105 defer resp.Body.Close()
106 data, err := io.ReadAll(io.LimitReader(resp.Body, httpMaxResponseBytes+1))
107 if err != nil {
108 return &transportError{method: method, err: err}
109 }
110 if len(data) > httpMaxResponseBytes {
111 return fmt.Errorf("browser broker: %s: reply exceeds %d bytes", method, httpMaxResponseBytes)
112 }
113 if resp.StatusCode == http.StatusConflict {
114 return decodeWireError(method, data)
115 }
116 if method == "capability" && (resp.StatusCode == http.StatusNotFound || resp.StatusCode == http.StatusNotImplemented) {
117 return errCapabilityUnsupported
118 }
119 if resp.StatusCode != http.StatusOK {
120 return fmt.Errorf("browser broker: %s: status %d: %s", method, resp.StatusCode, wireMessage(data))
121 }
122 if out == nil {
123 return nil
124 }
125 if err := json.Unmarshal(data, out); err != nil {
126 return fmt.Errorf("browser broker: decode %s reply: %w", method, err)
127 }
128 return nil
129 }
130
131 var errCapabilityUnsupported = errors.New("capability_unsupported: remote browser enhancement unavailable")
132
133 func (e *httpExecutor) BrowserCapability(ctx context.Context, name string, args json.RawMessage) (json.RawMessage, error) {
134 var out json.RawMessage
135 var params struct {
136 Action string `json:"action"`
137 }
138 if err := json.Unmarshal(args, &params); err != nil {
139 return nil, err
140 }
141 in := struct {
142 Name string `json:"name"`
143 Args json.RawMessage `json:"args"`
144 }{name, args}
145 err := e.call(ctx, "capability", in, &out)
146 write := name == "pointer" || name == "viewport" && params.Action != "get" || name == "record" && params.Action != "status"
147 if write && err != nil && !errors.Is(err, errCapabilityUnsupported) && !errors.Is(err, ErrStaleReference) && !errors.Is(err, ErrTakenOver) && !errors.Is(err, ErrNoGrant) && !errors.Is(err, ErrUnknownOutcome) {
148 err = fmt.Errorf("%w: %w", ErrUnknownOutcome, err)
149 }
150 return out, err
151 }
152
153 type transportError struct {
154 method string
155 err error
156 }
157
158 func (t *transportError) Error() string { return "browser broker: " + t.method + ": " + t.err.Error() }
159 func (t *transportError) Unwrap() error { return t.err }
160
161 // wireMessage prefers the handler's message over a raw body dump.
162 func wireMessage(data []byte) string {
163 var we wireError
164 if err := json.Unmarshal(data, &we); err == nil && we.Message != "" {
165 return we.Message
166 }
167 return strings.TrimSpace(string(data))
168 }
169
170 func decodeWireError(method string, data []byte) error {
171 var we wireError
172 if err := json.Unmarshal(data, &we); err != nil || we.Error == "" {
173 return fmt.Errorf("browser broker: %s: status 409: %s", method, strings.TrimSpace(string(data)))
174 }
175 sentinel, ok := wireErrorCodes[we.Error]
176 if !ok {
177 return fmt.Errorf("browser broker: %s: %s: %s", method, we.Error, we.Message)
178 }
179 detail := strings.TrimPrefix(strings.TrimPrefix(we.Message, sentinel.Error()), ": ")
180 if detail == "" {
181 return sentinel
182 }
183 return fmt.Errorf("%w: %s", sentinel, detail)
184 }
185
186 func (e *httpExecutor) Tabs(ctx context.Context) ([]Tab, error) {
187 var out wireTabs
188 if err := e.call(ctx, "tabs", struct{}{}, &out); err != nil {
189 return nil, err
190 }
191 tabs := make([]Tab, 0, len(out.Tabs))
192 for _, t := range out.Tabs {
193 tabs = append(tabs, t.tab())
194 }
195 return tabs, nil
196 }
197
198 func (e *httpExecutor) Open(ctx context.Context, req OpenRequest) (Tab, error) {
199 var out wireTab
200 if err := e.write(ctx, "open", wireOpenRequest(req), &out); err != nil {
201 return Tab{}, err
202 }
203 return out.tab(), nil
204 }
205
206 func (e *httpExecutor) Navigate(ctx context.Context, req NavigateRequest) (Tab, error) {
207 var out wireTab
208 if err := e.write(ctx, "navigate", wireNavigateRequest(req), &out); err != nil {
209 return Tab{}, err
210 }
211 return out.tab(), nil
212 }
213
214 func (e *httpExecutor) Snapshot(ctx context.Context, req SnapshotRequest) (Snapshot, error) {
215 var out wireSnapshot
216 if err := e.call(ctx, "snapshot", wireSnapshotRequest(req), &out); err != nil {
217 return Snapshot{}, err
218 }
219 return Snapshot(out), nil
220 }
221
222 func (e *httpExecutor) Screenshot(ctx context.Context, req ScreenshotRequest) (Screenshot, error) {
223 var out wireScreenshot
224 if err := e.call(ctx, "screenshot", wireScreenshotRequest(req), &out); err != nil {
225 return Screenshot{}, err
226 }
227 return Screenshot(out), nil
228 }
229
230 // Act sends one reserved write. A reply that never arrived leaves the
231 // action's fate unknown, which is exactly ErrUnknownOutcome.
232 func (e *httpExecutor) Act(ctx context.Context, req ActRequest) (ActResult, error) {
233 var out wireActResult
234 if err := e.write(ctx, "act", toWireAct(req), &out); err != nil {
235 if errors.Is(err, ErrUnknownOutcome) {
236 return ActResult{Outcome: OutcomeUnknown}, err
237 }
238 return ActResult{}, err
239 }
240 res := ActResult(out)
241 if res.Outcome == "" {
242 res.Outcome = OutcomeNotExecuted
243 if res.Executed {
244 res.Outcome = OutcomeExecuted
245 }
246 }
247 return res, nil
248 }
249
250 func (e *httpExecutor) Downloads(ctx context.Context, req DownloadsRequest) ([]Download, error) {
251 var out wireDownloads
252 in := wireDownloadsRequest{TabID: req.TabID, WaitForMs: req.WaitFor.Milliseconds()}
253 if err := e.call(ctx, "downloads", in, &out); err != nil {
254 return nil, err
255 }
256 downloads := make([]Download, 0, len(out.Downloads))
257 for _, d := range out.Downloads {
258 downloads = append(downloads, Download(d))
259 }
260 return downloads, nil
261 }
262
263 func (e *httpExecutor) Close(ctx context.Context, req CloseRequest) error {
264 return e.write(ctx, "close", wireCloseRequest(req), nil)
265 }
266
267 // Once a write is handed to HTTP, only explicit refusal codes prove it did
268 // not run. Truncated/invalid replies and HTTP failures also leave it unknown.
269 func (e *httpExecutor) write(ctx context.Context, method string, in, out any) error {
270 err := e.call(ctx, method, in, out)
271 if err == nil || errors.Is(err, ErrStaleReference) || errors.Is(err, ErrTakenOver) || errors.Is(err, ErrNoGrant) || errors.Is(err, ErrUnknownOutcome) {
272 return err
273 }
274 return fmt.Errorf("%w: %s", ErrUnknownOutcome, err.Error())
275 }
276
276 lines GO