| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285 |
- // Code generated by smithy-go-codegen DO NOT EDIT.
- package s3
- import (
- "context"
- "fmt"
- "github.com/aws/aws-sdk-go-v2/aws"
- "github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream"
- "github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream/eventstreamapi"
- "github.com/aws/aws-sdk-go-v2/service/s3/types"
- smithy "github.com/aws/smithy-go"
- "github.com/aws/smithy-go/middleware"
- smithysync "github.com/aws/smithy-go/sync"
- smithyhttp "github.com/aws/smithy-go/transport/http"
- "io"
- "io/ioutil"
- "sync"
- )
- // SelectObjectContentEventStreamReader provides the interface for reading events
- // from a stream.
- //
- // The writer's Close method must allow multiple concurrent calls.
- type SelectObjectContentEventStreamReader interface {
- Events() <-chan types.SelectObjectContentEventStream
- Close() error
- Err() error
- }
- type selectObjectContentEventStreamReader struct {
- stream chan types.SelectObjectContentEventStream
- decoder *eventstream.Decoder
- eventStream io.ReadCloser
- err *smithysync.OnceErr
- payloadBuf []byte
- done chan struct{}
- closeOnce sync.Once
- }
- func newSelectObjectContentEventStreamReader(readCloser io.ReadCloser, decoder *eventstream.Decoder) *selectObjectContentEventStreamReader {
- w := &selectObjectContentEventStreamReader{
- stream: make(chan types.SelectObjectContentEventStream),
- decoder: decoder,
- eventStream: readCloser,
- err: smithysync.NewOnceErr(),
- done: make(chan struct{}),
- payloadBuf: make([]byte, 10*1024),
- }
- go w.readEventStream()
- return w
- }
- func (r *selectObjectContentEventStreamReader) Events() <-chan types.SelectObjectContentEventStream {
- return r.stream
- }
- func (r *selectObjectContentEventStreamReader) readEventStream() {
- defer r.Close()
- defer close(r.stream)
- for {
- r.payloadBuf = r.payloadBuf[0:0]
- decodedMessage, err := r.decoder.Decode(r.eventStream, r.payloadBuf)
- if err != nil {
- if err == io.EOF {
- return
- }
- select {
- case <-r.done:
- return
- default:
- r.err.SetError(err)
- return
- }
- }
- event, err := r.deserializeEventMessage(&decodedMessage)
- if err != nil {
- r.err.SetError(err)
- return
- }
- select {
- case r.stream <- event:
- case <-r.done:
- return
- }
- }
- }
- func (r *selectObjectContentEventStreamReader) deserializeEventMessage(msg *eventstream.Message) (types.SelectObjectContentEventStream, error) {
- messageType := msg.Headers.Get(eventstreamapi.MessageTypeHeader)
- if messageType == nil {
- return nil, fmt.Errorf("%s event header not present", eventstreamapi.MessageTypeHeader)
- }
- switch messageType.String() {
- case eventstreamapi.EventMessageType:
- var v types.SelectObjectContentEventStream
- if err := awsRestxml_deserializeEventStreamSelectObjectContentEventStream(&v, msg); err != nil {
- return nil, err
- }
- return v, nil
- case eventstreamapi.ExceptionMessageType:
- return nil, awsRestxml_deserializeEventStreamExceptionSelectObjectContentEventStream(msg)
- case eventstreamapi.ErrorMessageType:
- errorCode := "UnknownError"
- errorMessage := errorCode
- if header := msg.Headers.Get(eventstreamapi.ErrorCodeHeader); header != nil {
- errorCode = header.String()
- }
- if header := msg.Headers.Get(eventstreamapi.ErrorMessageHeader); header != nil {
- errorMessage = header.String()
- }
- return nil, &smithy.GenericAPIError{
- Code: errorCode,
- Message: errorMessage,
- }
- default:
- mc := msg.Clone()
- return nil, &UnknownEventMessageError{
- Type: messageType.String(),
- Message: &mc,
- }
- }
- }
- func (r *selectObjectContentEventStreamReader) ErrorSet() <-chan struct{} {
- return r.err.ErrorSet()
- }
- func (r *selectObjectContentEventStreamReader) Close() error {
- r.closeOnce.Do(r.safeClose)
- return r.Err()
- }
- func (r *selectObjectContentEventStreamReader) safeClose() {
- close(r.done)
- r.eventStream.Close()
- }
- func (r *selectObjectContentEventStreamReader) Err() error {
- return r.err.Err()
- }
- func (r *selectObjectContentEventStreamReader) Closed() <-chan struct{} {
- return r.done
- }
- type awsRestxml_deserializeOpEventStreamSelectObjectContent struct {
- LogEventStreamWrites bool
- LogEventStreamReads bool
- }
- func (*awsRestxml_deserializeOpEventStreamSelectObjectContent) ID() string {
- return "OperationEventStreamDeserializer"
- }
- func (m *awsRestxml_deserializeOpEventStreamSelectObjectContent) HandleDeserialize(ctx context.Context, in middleware.DeserializeInput, next middleware.DeserializeHandler) (
- out middleware.DeserializeOutput, metadata middleware.Metadata, err error,
- ) {
- defer func() {
- if err == nil {
- return
- }
- m.closeResponseBody(out)
- }()
- logger := middleware.GetLogger(ctx)
- request, ok := in.Request.(*smithyhttp.Request)
- if !ok {
- return out, metadata, fmt.Errorf("unknown transport type: %T", in.Request)
- }
- _ = request
- out, metadata, err = next.HandleDeserialize(ctx, in)
- if err != nil {
- return out, metadata, err
- }
- deserializeOutput, ok := out.RawResponse.(*smithyhttp.Response)
- if !ok {
- return out, metadata, fmt.Errorf("unknown transport type: %T", out.RawResponse)
- }
- _ = deserializeOutput
- output, ok := out.Result.(*SelectObjectContentOutput)
- if out.Result != nil && !ok {
- return out, metadata, fmt.Errorf("unexpected output result type: %T", out.Result)
- } else if out.Result == nil {
- output = &SelectObjectContentOutput{}
- out.Result = output
- }
- eventReader := newSelectObjectContentEventStreamReader(
- deserializeOutput.Body,
- eventstream.NewDecoder(func(options *eventstream.DecoderOptions) {
- options.Logger = logger
- options.LogMessages = m.LogEventStreamReads
- }),
- )
- defer func() {
- if err == nil {
- return
- }
- _ = eventReader.Close()
- }()
- output.eventStream = NewSelectObjectContentEventStream(func(stream *SelectObjectContentEventStream) {
- stream.Reader = eventReader
- })
- go output.eventStream.waitStreamClose()
- return out, metadata, nil
- }
- func (*awsRestxml_deserializeOpEventStreamSelectObjectContent) closeResponseBody(out middleware.DeserializeOutput) {
- if resp, ok := out.RawResponse.(*smithyhttp.Response); ok && resp != nil && resp.Body != nil {
- _, _ = io.Copy(ioutil.Discard, resp.Body)
- _ = resp.Body.Close()
- }
- }
- func addEventStreamSelectObjectContentMiddleware(stack *middleware.Stack, options Options) error {
- if err := stack.Deserialize.Insert(&awsRestxml_deserializeOpEventStreamSelectObjectContent{
- LogEventStreamWrites: options.ClientLogMode.IsRequestEventMessage(),
- LogEventStreamReads: options.ClientLogMode.IsResponseEventMessage(),
- }, "OperationDeserializer", middleware.Before); err != nil {
- return err
- }
- return nil
- }
- // UnknownEventMessageError provides an error when a message is received from the stream,
- // but the reader is unable to determine what kind of message it is.
- type UnknownEventMessageError struct {
- Type string
- Message *eventstream.Message
- }
- // Error retruns the error message string.
- func (e *UnknownEventMessageError) Error() string {
- return "unknown event stream message type, " + e.Type
- }
- func setSafeEventStreamClientLogMode(o *Options, operation string) {
- switch operation {
- case "SelectObjectContent":
- toggleEventStreamClientLogMode(o, false, true)
- return
- default:
- return
- }
- }
- func toggleEventStreamClientLogMode(o *Options, request, response bool) {
- mode := o.ClientLogMode
- if request && mode.IsRequestWithBody() {
- mode.ClearRequestWithBody()
- mode |= aws.LogRequest
- }
- if response && mode.IsResponseWithBody() {
- mode.ClearResponseWithBody()
- mode |= aws.LogResponse
- }
- o.ClientLogMode = mode
- }
|