| 91 | } |
| 92 | } |
| 93 | async fn work<T: tokio::io::AsyncWrite + Send + 'static>(&self, writer: T) { |
| 94 | use tokio::io::AsyncWriteExt; |
| 95 | let mut writer = pin!(writer); |
| 96 | loop { |
| 97 | while let Some(job) = self.pop() { |
| 98 | match job { |
| 99 | Job::Flush => { |
| 100 | if let Err(e) = writer.flush().await { |
| 101 | self.report_error(e); |
| 102 | return; |
| 103 | } |
| 104 | |
| 105 | tracing::debug!("worker marking flush complete"); |
| 106 | self.state().flush_pending = false; |
| 107 | } |
| 108 | |
| 109 | Job::Write(mut bytes) => { |
| 110 | tracing::debug!("worker writing: {bytes:?}"); |
| 111 | let len = bytes.len(); |
| 112 | match writer.write_all_buf(&mut bytes).await { |
| 113 | Err(e) => { |
| 114 | self.report_error(e); |
| 115 | return; |
| 116 | } |
| 117 | Ok(_) => { |
| 118 | self.state().write_budget += len; |
| 119 | } |
| 120 | } |
| 121 | } |
| 122 | } |
| 123 | |
| 124 | let waker = self.state().write_ready_changed.take(); |
| 125 | if let Some(waker) = waker { |
| 126 | waker.wake(); |
| 127 | } |
| 128 | } |
| 129 | self.new_work.notified().await; |
| 130 | } |
| 131 | } |
| 132 | } |
| 133 | |
| 134 | /// Provides a [`OutputStream`] impl from a [`tokio::io::AsyncWrite`] impl |