Stream executes a streaming HTTP request and returns channel with stream events.
(ctx context.Context, provider string, req *HTTPRequest)
| 113 | |
| 114 | // Stream executes a streaming HTTP request and returns channel with stream events. |
| 115 | func (m *Manager) Stream(ctx context.Context, provider string, req *HTTPRequest) (<-chan StreamEvent, error) { |
| 116 | if req == nil { |
| 117 | return nil, fmt.Errorf("wsrelay: request is nil") |
| 118 | } |
| 119 | msg := Message{ID: uuid.NewString(), Type: MessageTypeHTTPReq, Payload: encodeRequest(req)} |
| 120 | respCh, err := m.Send(ctx, provider, msg) |
| 121 | if err != nil { |
| 122 | return nil, err |
| 123 | } |
| 124 | out := make(chan StreamEvent) |
| 125 | go func() { |
| 126 | defer close(out) |
| 127 | send := func(ev StreamEvent) bool { |
| 128 | if ctx == nil { |
| 129 | out <- ev |
| 130 | return true |
| 131 | } |
| 132 | select { |
| 133 | case <-ctx.Done(): |
| 134 | return false |
| 135 | case out <- ev: |
| 136 | return true |
| 137 | } |
| 138 | } |
| 139 | for { |
| 140 | select { |
| 141 | case <-ctx.Done(): |
| 142 | return |
| 143 | case msg, ok := <-respCh: |
| 144 | if !ok { |
| 145 | _ = send(StreamEvent{Err: errors.New("wsrelay: stream closed")}) |
| 146 | return |
| 147 | } |
| 148 | switch msg.Type { |
| 149 | case MessageTypeStreamStart: |
| 150 | resp := decodeResponse(msg.Payload) |
| 151 | if okSend := send(StreamEvent{Type: MessageTypeStreamStart, Status: resp.Status, Headers: resp.Headers}); !okSend { |
| 152 | return |
| 153 | } |
| 154 | case MessageTypeStreamChunk: |
| 155 | chunk := decodeChunk(msg.Payload) |
| 156 | if okSend := send(StreamEvent{Type: MessageTypeStreamChunk, Payload: chunk}); !okSend { |
| 157 | return |
| 158 | } |
| 159 | case MessageTypeStreamEnd: |
| 160 | _ = send(StreamEvent{Type: MessageTypeStreamEnd}) |
| 161 | return |
| 162 | case MessageTypeError: |
| 163 | _ = send(StreamEvent{Type: MessageTypeError, Err: decodeError(msg.Payload)}) |
| 164 | return |
| 165 | case MessageTypeHTTPResp: |
| 166 | resp := decodeResponse(msg.Payload) |
| 167 | _ = send(StreamEvent{Type: MessageTypeHTTPResp, Status: resp.Status, Headers: resp.Headers, Payload: resp.Body}) |
| 168 | return |
| 169 | default: |
| 170 | } |
| 171 | } |
| 172 | } |
no test coverage detected