Source file
src/net/http/clientconn.go
1
2
3
4
5 package http
6
7 import (
8 "context"
9 "errors"
10 "fmt"
11 "net"
12 "net/http/httptrace"
13 "net/url"
14 "sync"
15 )
16
17
18
19
20
21 type ClientConn struct {
22 cc genericClientConn
23
24 stateHookMu sync.Mutex
25 userStateHook func(*ClientConn)
26 stateHookRunning bool
27 lastAvailable int
28 lastInFlight int
29 lastClosed bool
30 }
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46 type newClientConner interface {
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76 NewClientConn(nc net.Conn, internalStateHook func()) (RoundTripper, error)
77 }
78
79
80
81
82
83 type genericClientConn interface {
84 Close() error
85 Err() error
86 RoundTrip(req *Request) (*Response, error)
87 Reserve() error
88 Release()
89 Available() int
90 InFlight() int
91 }
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113 func (t *Transport) NewClientConn(ctx context.Context, scheme, address string) (*ClientConn, error) {
114 t.nextProtoOnce.Do(t.onceSetNextProtoDefaults)
115
116 if t.h2Config != nil {
117
118
119 if cc, err := t.http2NewClientConnFromContext(ctx); err != errors.ErrUnsupported {
120 return cc, err
121 }
122 }
123
124 switch scheme {
125 case "http", "https":
126 default:
127 return nil, fmt.Errorf("net/http: invalid scheme %q", scheme)
128 }
129
130 host, port, err := net.SplitHostPort(address)
131 if err != nil {
132 return nil, err
133 }
134 if port == "" {
135 port = schemePort(scheme)
136 }
137
138 var proxyURL *url.URL
139 if t.Proxy != nil {
140
141 req := &Request{
142 ctx: ctx,
143 Method: "GET",
144 URL: &url.URL{
145 Scheme: scheme,
146 Host: host,
147 Path: "/",
148 },
149 Proto: "HTTP/1.1",
150 ProtoMajor: 1,
151 ProtoMinor: 1,
152 Header: make(Header),
153 Body: NoBody,
154 Host: host,
155 }
156 var err error
157 proxyURL, err = t.Proxy(req)
158 if err != nil {
159 return nil, err
160 }
161 }
162
163 cm := connectMethod{
164 targetScheme: scheme,
165 targetAddr: net.JoinHostPort(host, port),
166 proxyURL: proxyURL,
167 }
168
169
170
171
172
173
174
175
176
177 cc := &ClientConn{}
178 const isClientConn = true
179 pconn, err := t.dialConn(ctx, cm, isClientConn, cc.maybeRunStateHook)
180 if err != nil {
181 return nil, err
182 }
183
184
185
186
187 cc.stateHookMu.Lock()
188 defer cc.stateHookMu.Unlock()
189 if pconn.alt != nil {
190
191
192 gc, ok := pconn.alt.(genericClientConn)
193 if !ok {
194 return nil, errors.New("http: NewClientConn returned something that is not a ClientConn")
195 }
196 cc.cc = gc
197 cc.lastAvailable = gc.Available()
198 } else {
199
200 pconn.availch = make(chan struct{}, 1)
201 pconn.availch <- struct{}{}
202 cc.cc = http1ClientConn{pconn}
203 cc.lastAvailable = 1
204 }
205 return cc, nil
206 }
207
208
209
210 func (cc *ClientConn) Close() error {
211 defer cc.maybeRunStateHook()
212 return cc.cc.Close()
213 }
214
215
216
217
218 func (cc *ClientConn) Err() error {
219 return cc.cc.Err()
220 }
221
222 func validateClientConnRequest(req *Request) error {
223 if req.URL == nil {
224 return errors.New("http: nil Request.URL")
225 }
226 if req.Header == nil {
227 return errors.New("http: nil Request.Header")
228 }
229
230 if err := validateHeaders(req.Header); err != "" {
231 return fmt.Errorf("http: invalid header %s", err)
232 }
233
234 if err := validateHeaders(req.Trailer); err != "" {
235 return fmt.Errorf("http: invalid trailer %s", err)
236 }
237 if req.Method != "" && !validMethod(req.Method) {
238 return fmt.Errorf("http: invalid method %q", req.Method)
239 }
240 if req.URL.Host == "" {
241 return errors.New("http: no Host in request URL")
242 }
243 return nil
244 }
245
246
247
248
249
250
251
252
253
254 func (cc *ClientConn) RoundTrip(req *Request) (*Response, error) {
255 defer cc.maybeRunStateHook()
256 if req.URL == nil && req.Method == ":ping" {
257
258
259 pinger, ok := cc.cc.(interface {
260 Ping(context.Context) error
261 })
262 if !ok {
263 return nil, errors.New("http: ClientConn does not support PING")
264 }
265 return nil, pinger.Ping(req.Context())
266 }
267 if err := validateClientConnRequest(req); err != nil {
268 cc.Release()
269 return nil, err
270 }
271 return cc.cc.RoundTrip(req)
272 }
273
274
275
276
277 func (cc *ClientConn) Available() int {
278 return cc.cc.Available()
279 }
280
281
282
283
284 func (cc *ClientConn) InFlight() int {
285 return cc.cc.InFlight()
286 }
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301 func (cc *ClientConn) Reserve() error {
302 defer cc.maybeRunStateHook()
303 return cc.cc.Reserve()
304 }
305
306
307
308 func (cc *ClientConn) Release() {
309 defer cc.maybeRunStateHook()
310 cc.cc.Release()
311 }
312
313
314
315 func (cc *ClientConn) shouldRunStateHook(stopRunning bool) func(*ClientConn) {
316 cc.stateHookMu.Lock()
317 defer cc.stateHookMu.Unlock()
318 if cc.cc == nil {
319 return nil
320 }
321 if stopRunning {
322 cc.stateHookRunning = false
323 }
324 if cc.userStateHook == nil {
325 return nil
326 }
327 if cc.stateHookRunning {
328 return nil
329 }
330 var (
331 available = cc.Available()
332 inFlight = cc.InFlight()
333 closed = cc.Err() != nil
334 )
335 var hook func(*ClientConn)
336 if available > cc.lastAvailable || inFlight < cc.lastInFlight || closed != cc.lastClosed {
337 hook = cc.userStateHook
338 cc.stateHookRunning = true
339 }
340 cc.lastAvailable = available
341 cc.lastInFlight = inFlight
342 cc.lastClosed = closed
343 return hook
344 }
345
346 func (cc *ClientConn) maybeRunStateHook() {
347 hook := cc.shouldRunStateHook(false)
348 if hook == nil {
349 return
350 }
351
352
353
354
355
356 hook(cc)
357
358
359
360
361
362 hook = cc.shouldRunStateHook(true)
363 if hook != nil {
364 go func() {
365 for hook != nil {
366 hook(cc)
367 hook = cc.shouldRunStateHook(true)
368 }
369 }()
370 }
371 }
372
373
374
375
376
377
378
379
380
381
382
383
384
385 func (cc *ClientConn) SetStateHook(f func(*ClientConn)) {
386 cc.stateHookMu.Lock()
387 cc.userStateHook = f
388 cc.stateHookMu.Unlock()
389 cc.maybeRunStateHook()
390 }
391
392
393
394 type http1ClientConn struct {
395 pconn *persistConn
396 }
397
398 func (cc http1ClientConn) RoundTrip(req *Request) (*Response, error) {
399 ctx := req.Context()
400 trace := httptrace.ContextClientTrace(ctx)
401
402
403 ctx, cancel := context.WithCancelCause(req.Context())
404 if req.Cancel != nil {
405 go awaitLegacyCancel(ctx, cancel, req)
406 }
407
408 treq := &transportRequest{Request: req, trace: trace, ctx: ctx, cancel: cancel}
409 resp, err := cc.pconn.roundTrip(treq)
410 if err != nil {
411 return nil, err
412 }
413 resp.Request = req
414 return resp, nil
415 }
416
417 func (cc http1ClientConn) Close() error {
418 cc.pconn.close(errors.New("ClientConn closed"))
419 return nil
420 }
421
422 func (cc http1ClientConn) Err() error {
423 select {
424 case <-cc.pconn.closech:
425 return cc.pconn.closed
426 default:
427 return nil
428 }
429 }
430
431 func (cc http1ClientConn) Available() int {
432 cc.pconn.mu.Lock()
433 defer cc.pconn.mu.Unlock()
434 if cc.pconn.closed != nil || cc.pconn.reserved || cc.pconn.inFlight {
435 return 0
436 }
437 return 1
438 }
439
440 func (cc http1ClientConn) InFlight() int {
441 cc.pconn.mu.Lock()
442 defer cc.pconn.mu.Unlock()
443 if cc.pconn.closed == nil && (cc.pconn.reserved || cc.pconn.inFlight) {
444 return 1
445 }
446 return 0
447 }
448
449 func (cc http1ClientConn) Reserve() error {
450 cc.pconn.mu.Lock()
451 defer cc.pconn.mu.Unlock()
452 if cc.pconn.closed != nil {
453 return cc.pconn.closed
454 }
455 select {
456 case <-cc.pconn.availch:
457 default:
458 return errors.New("connection is unavailable")
459 }
460 cc.pconn.reserved = true
461 return nil
462 }
463
464 func (cc http1ClientConn) Release() {
465 cc.pconn.mu.Lock()
466 defer cc.pconn.mu.Unlock()
467 if cc.pconn.reserved {
468 select {
469 case cc.pconn.availch <- struct{}{}:
470 default:
471 panic("cannot release reservation")
472 }
473 cc.pconn.reserved = false
474 }
475 }
476
View as plain text