| 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 |