(ctx context.Context, stream streaming.Stream)
| 141 | } |
| 142 | |
| 143 | func ReceiveStream(ctx context.Context, stream streaming.Stream) io.Reader { |
| 144 | r, w := io.Pipe() |
| 145 | go func() { |
| 146 | defer stream.Close() |
| 147 | var window int32 |
| 148 | for { |
| 149 | var werr error |
| 150 | if window < windowSize { |
| 151 | update := &transferapi.WindowUpdate{ |
| 152 | Update: windowSize, |
| 153 | } |
| 154 | anyType, err := typeurl.MarshalAny(update) |
| 155 | if err != nil { |
| 156 | w.CloseWithError(fmt.Errorf("failed to marshal window update: %w", err)) |
| 157 | return |
| 158 | } |
| 159 | // check window update error after recv, stream may be complete |
| 160 | if werr = stream.Send(anyType); werr == nil { |
| 161 | window += windowSize |
| 162 | } else if errors.Is(werr, io.EOF) { |
| 163 | // TODO: Why does send return EOF here |
| 164 | werr = nil |
| 165 | } |
| 166 | } |
| 167 | anyType, err := stream.Recv() |
| 168 | if err != nil { |
| 169 | if errors.Is(err, io.EOF) || errors.Is(err, context.Canceled) { |
| 170 | err = nil |
| 171 | } else { |
| 172 | err = fmt.Errorf("received failed: %w", err) |
| 173 | } |
| 174 | w.CloseWithError(err) |
| 175 | return |
| 176 | } else if werr != nil { |
| 177 | // Try receive before erroring out |
| 178 | w.CloseWithError(fmt.Errorf("failed to send window update: %w", werr)) |
| 179 | return |
| 180 | } |
| 181 | i, err := typeurl.UnmarshalAny(anyType) |
| 182 | if err != nil { |
| 183 | w.CloseWithError(fmt.Errorf("failed to unmarshal received object: %w", err)) |
| 184 | return |
| 185 | } |
| 186 | switch v := i.(type) { |
| 187 | case *transferapi.Data: |
| 188 | n, err := w.Write(v.Data) |
| 189 | if err != nil { |
| 190 | w.CloseWithError(fmt.Errorf("failed to unmarshal received object: %w", err)) |
| 191 | // Close will error out sender |
| 192 | return |
| 193 | } |
| 194 | window = window - int32(n) |
| 195 | // TODO: Handle error case |
| 196 | default: |
| 197 | log.G(ctx).Warnf("Ignoring unknown stream object of type %T", i) |
| 198 | continue |
| 199 | } |
| 200 | } |
searching dependent graphs…