返回 DeepSeek-Reasonix
request_observation.go
根目录 / internal / provider / request_observation.go
1 package provider
2
3 import (
4 "context"
5 "errors"
6 "io"
7 "net/http"
8 "sync"
9 "sync/atomic"
10 "time"
11 )
12
13 // RequestObservation is content-free, process-local transport evidence. Body
14 // bytes include SSE comments/heartbeats, not just model output. It deliberately
15 // excludes URL userinfo/query/fragment, headers and request/response bodies.
16 type RequestObservation struct {
17 ID uint64 `json:"id"`
18 StartedAt time.Time `json:"startedAt"`
19 ObservedAt time.Time `json:"observedAt"`
20 Phase string `json:"phase"`
21 LastPhase string `json:"lastPhase,omitempty"`
22 ConnectedAt time.Time `json:"connectedAt,omitempty"`
23 WrittenAt time.Time `json:"writtenAt,omitempty"`
24 HeadersAt time.Time `json:"headersAt,omitempty"`
25 FirstBodyAt time.Time `json:"firstBodyAt,omitempty"`
26 LastBodyAt time.Time `json:"lastBodyAt,omitempty"`
27 FinishedAt time.Time `json:"finishedAt,omitempty"`
28 Status int `json:"status,omitempty"`
29 BodyBytes int64 `json:"bodyBytes"`
30 Method string `json:"method,omitempty"`
31 Host string `json:"host,omitempty"`
32 RequestPath string `json:"requestPath,omitempty"`
33 RequestBytes int64 `json:"requestBytes,omitempty"`
34 Network string `json:"network,omitempty"`
35 DialAddress string `json:"dialAddress,omitempty"`
36 RemoteAddress string `json:"remoteAddress,omitempty"`
37 HTTPProtocol string `json:"httpProtocol,omitempty"`
38 ConnectionReused bool `json:"connectionReused,omitempty"`
39 TransportCode string `json:"transportCode,omitempty"`
40 }
41
42 type requestObserverKey struct{}
43
44 var observedRequestID atomic.Uint64
45
46 // WithRequestObserver attaches a synchronous, concurrency-safe observer. The
47 // callback must be short and must not call back into the observed transport.
48 // It changes no provider-visible bytes or timeout/retry behavior.
49 func WithRequestObserver(ctx context.Context, observe func(RequestObservation)) context.Context {
50 return context.WithValue(ctx, requestObserverKey{}, observe)
51 }
52
53 type requestObservationState struct {
54 mu sync.Mutex
55 value RequestObservation
56 observe func(RequestObservation)
57 ctx context.Context
58 }
59
60 func observeRequest(ctx context.Context) (context.Context, *requestObservationState) {
61 observe, _ := ctx.Value(requestObserverKey{}).(func(RequestObservation))
62 if observe == nil {
63 return ctx, nil
64 }
65 now := time.Now().UTC()
66 s := &requestObservationState{ctx: ctx, observe: observe, value: RequestObservation{ID: observedRequestID.Add(1), StartedAt: now}}
67 s.update("request_started", nil)
68 return s.traceContext(ctx), s
69 }
70
71 func (s *requestObservationState) update(phase string, edit func(*RequestObservation)) {
72 if s == nil {
73 return
74 }
75 s.mu.Lock()
76 defer s.mu.Unlock()
77 if !s.value.FinishedAt.IsZero() {
78 return
79 }
80 s.value.ObservedAt = time.Now().UTC()
81 previousPhase := s.value.Phase
82 s.value.Phase = phase
83 if edit != nil {
84 edit(&s.value)
85 }
86 if !s.value.FinishedAt.IsZero() {
87 s.value.LastPhase = previousPhase
88 }
89 s.observe(s.value)
90 }
91
92 func (s *requestObservationState) finish(err error, phase string) {
93 if s == nil {
94 return
95 }
96 if errors.Is(s.ctx.Err(), context.Canceled) || errors.Is(err, context.Canceled) {
97 phase = "canceled"
98 } else if errors.Is(s.ctx.Err(), context.DeadlineExceeded) || errors.Is(err, context.DeadlineExceeded) {
99 phase = "deadline_exceeded"
100 }
101 s.update(phase, func(v *RequestObservation) {
102 v.FinishedAt = v.ObservedAt
103 v.TransportCode = HTTP2TransportCode(err)
104 })
105 }
106
107 func (s *requestObservationState) response(resp *http.Response) {
108 if s == nil {
109 return
110 }
111 s.update("headers_received", func(v *RequestObservation) {
112 v.HeadersAt, v.Status, v.HTTPProtocol = v.ObservedAt, resp.StatusCode, resp.Proto
113 })
114 resp.Body = &observedResponseBody{ReadCloser: resp.Body, state: s}
115 }
116
117 type observedResponseBody struct {
118 io.ReadCloser
119 state *requestObservationState
120 }
121
122 func (b *observedResponseBody) Read(p []byte) (int, error) {
123 n, err := b.ReadCloser.Read(p)
124 if n > 0 {
125 b.state.update("body_received", func(v *RequestObservation) {
126 if v.FirstBodyAt.IsZero() {
127 v.FirstBodyAt = v.ObservedAt
128 }
129 v.LastBodyAt = v.ObservedAt
130 v.BodyBytes += int64(n)
131 })
132 }
133 if err != nil {
134 phase := "body_error"
135 if errors.Is(err, io.EOF) {
136 phase = "body_eof"
137 }
138 b.state.finish(err, phase)
139 }
140 return n, err
141 }
142
143 func (b *observedResponseBody) Close() error {
144 err := b.ReadCloser.Close()
145 b.state.finish(err, "body_closed")
146 return err
147 }
148
148 lines GO