1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138
use core::{ pin::Pin, task::{Context, Poll}, }; use std::io::Result; use crate::{codec::Decode, util::PartialBuffer}; use futures_core::ready; use pin_project_lite::pin_project; use tokio_02::io::{AsyncBufRead, AsyncRead}; #[derive(Debug)] enum State { Decoding, Flushing, Done, Next, } pin_project! { #[derive(Debug)] pub struct Decoder<R, D: Decode> { #[pin] reader: R, decoder: D, state: State, multiple_members: bool, } } impl<R: AsyncBufRead, D: Decode> Decoder<R, D> { pub fn new(reader: R, decoder: D) -> Self { Self { reader, decoder, state: State::Decoding, multiple_members: false, } } pub fn get_ref(&self) -> &R { &self.reader } pub fn get_mut(&mut self) -> &mut R { &mut self.reader } pub fn get_pin_mut(self: Pin<&mut Self>) -> Pin<&mut R> { self.project().reader } pub fn into_inner(self) -> R { self.reader } pub fn multiple_members(&mut self, enabled: bool) { self.multiple_members = enabled; } fn do_poll_read( self: Pin<&mut Self>, cx: &mut Context<'_>, output: &mut PartialBuffer<&mut [u8]>, ) -> Poll<Result<()>> { let mut this = self.project(); loop { *this.state = match this.state { State::Decoding => { let input = ready!(this.reader.as_mut().poll_fill_buf(cx))?; if input.is_empty() { State::Flushing } else { let mut input = PartialBuffer::new(input); let done = this.decoder.decode(&mut input, output)?; let len = input.written().len(); this.reader.as_mut().consume(len); if done { State::Flushing } else { State::Decoding } } } State::Flushing => { if this.decoder.finish(output)? { if *this.multiple_members { this.decoder.reinit()?; State::Next } else { State::Done } } else { State::Flushing } } State::Done => State::Done, State::Next => { let input = ready!(this.reader.as_mut().poll_fill_buf(cx))?; if input.is_empty() { State::Done } else { State::Decoding } } }; if let State::Done = *this.state { return Poll::Ready(Ok(())); } if output.unwritten().is_empty() { return Poll::Ready(Ok(())); } } } } impl<R: AsyncBufRead, D: Decode> AsyncRead for Decoder<R, D> { fn poll_read( self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut [u8], ) -> Poll<Result<usize>> { if buf.is_empty() { return Poll::Ready(Ok(0)); } let mut output = PartialBuffer::new(buf); match self.do_poll_read(cx, &mut output)? { Poll::Pending if output.written().is_empty() => Poll::Pending, _ => Poll::Ready(Ok(output.written().len())), } } }