Class AsyncStreamingInputPlus

All Implemented Interfaces:
Closeable, DataInput, AutoCloseable, DataInputPlus, StreamingDataInputPlus

public class AsyncStreamingInputPlus extends RebufferingInputStream implements StreamingDataInputPlus
  • Constructor Details

    • AsyncStreamingInputPlus

      public AsyncStreamingInputPlus(io.netty.channel.Channel channel)
  • Method Details

    • append

      public boolean append(io.netty.buffer.ByteBuf buf) throws IllegalStateException
      Append a ByteBuf to the end of the einternal queue. Note: it's expected this method is invoked on the netty event loop.
      Throws:
      IllegalStateException
    • reBuffer

      protected void reBuffer() throws ClosedChannelException
      Implementations must implement this method to refill the buffer. They can expect the buffer to be empty when this method is invoked. Release open buffers and poll the queue for more data.

      This is invoked on a consuming thread (not the event loop) because if we block on the queue we can't fill it on the event loop (as that's where the buffers are coming from).

      Specified by:
      reBuffer in class RebufferingInputStream
      Throws:
      ClosedChannelException - when no further reading from this instance should occur. Implies this instance is closed.
    • consume

      public void consume(AsyncStreamingInputPlus.Consumer consumer, long length) throws IOException
      Consumes bytes in the stream until the given length
      Throws:
      IOException
    • unsafeAvailable

      public int unsafeAvailable()
      As long as this method is invoked on the consuming thread the returned value will be accurate.
    • isEmpty

      public boolean isEmpty()
    • close

      public void close()
      Note: This should invoked on the consuming thread.
      Specified by:
      close in interface AutoCloseable
      Specified by:
      close in interface Closeable
      Specified by:
      close in interface StreamingDataInputPlus
      Overrides:
      close in class InputStream
    • requestClosure

      public void requestClosure()
      Mark this stream as closed, but do not release any of the resources. Note: this is best to be called from the producer thread.
    • getAllocator

      public io.netty.buffer.ByteBufAllocator getAllocator()