(ctx context.Context, request []byte)
| 243 | } |
| 244 | |
| 245 | func (h *Host) callHostStreamEmit(ctx context.Context, request []byte) ([]byte, error) { |
| 246 | var req rpcStreamEmitRequest |
| 247 | if errUnmarshal := json.Unmarshal(request, &req); errUnmarshal != nil { |
| 248 | return nil, fmt.Errorf("decode stream emit request: %w", errUnmarshal) |
| 249 | } |
| 250 | chunk := pluginapi.ExecutorStreamChunk{Payload: append([]byte(nil), req.Payload...)} |
| 251 | if req.Error != "" { |
| 252 | chunk.Err = fmt.Errorf("%s", req.Error) |
| 253 | } |
| 254 | if errEmit := h.streams.emit(ctx, req.StreamID, chunk); errEmit != nil { |
| 255 | return nil, errEmit |
| 256 | } |
| 257 | return marshalRPCResult(rpcEmptyResponse{}) |
| 258 | } |
| 259 | |
| 260 | func (h *Host) callHostStreamClose(request []byte) ([]byte, error) { |
| 261 | var req rpcStreamCloseRequest |
no test coverage detected