(&mut self, req: Request)
| 163 | } |
| 164 | |
| 165 | fn call(&mut self, req: Request) -> Self::Future { |
| 166 | // Must make sure this is a cheap clone, required so the app can be moved into the async boxed future. |
| 167 | // See https://tokio.rs/blog/2021-05-14-inventing-the-service-trait |
| 168 | // The alternative is to perform the operation synchronously right here, |
| 169 | // but if we use the `tower_abci::buffer4::Worker` that means nothing else |
| 170 | // get processed during that time. |
| 171 | let app = self.0.clone(); |
| 172 | |
| 173 | // Another trick to avoid any subtle bugs is the mem::replace. |
| 174 | // See https://github.com/tower-rs/tower/issues/547 |
| 175 | let app: A = std::mem::replace(&mut self.0, app); |
| 176 | |
| 177 | // Because this is async, make sure the `Consensus` service is wrapped in a concurrency limiting Tower layer. |
| 178 | let res = async move { |
| 179 | let res = match req { |
| 180 | Request::Echo(r) => Response::Echo(log_error(app.echo(r).await)?), |
| 181 | Request::Info(r) => Response::Info(log_error(app.info(r).await)?), |
| 182 | Request::InitChain(r) => Response::InitChain(log_error(app.init_chain(r).await)?), |
| 183 | Request::Query(r) => Response::Query(log_error(app.query(r).await)?), |
| 184 | Request::CheckTx(r) => Response::CheckTx(log_error(app.check_tx(r).await)?), |
| 185 | Request::PrepareProposal(r) => { |
| 186 | Response::PrepareProposal(log_error(app.prepare_proposal(r).await)?) |
| 187 | } |
| 188 | Request::ProcessProposal(r) => { |
| 189 | Response::ProcessProposal(log_error(app.process_proposal(r).await)?) |
| 190 | } |
| 191 | Request::ExtendVote(r) => { |
| 192 | Response::ExtendVote(log_error(app.extend_vote(r).await)?) |
| 193 | } |
| 194 | Request::VerifyVoteExtension(r) => { |
| 195 | Response::VerifyVoteExtension(log_error(app.verify_vote_extension(r).await)?) |
| 196 | } |
| 197 | Request::FinalizeBlock(r) => { |
| 198 | Response::FinalizeBlock(log_error(app.finalize_block(r).await)?) |
| 199 | } |
| 200 | Request::Commit => Response::Commit(log_error(app.commit().await)?), |
| 201 | Request::ListSnapshots => { |
| 202 | Response::ListSnapshots(log_error(app.list_snapshots().await)?) |
| 203 | } |
| 204 | Request::OfferSnapshot(r) => { |
| 205 | Response::OfferSnapshot(log_error(app.offer_snapshot(r).await)?) |
| 206 | } |
| 207 | Request::LoadSnapshotChunk(r) => { |
| 208 | Response::LoadSnapshotChunk(log_error(app.load_snapshot_chunk(r).await)?) |
| 209 | } |
| 210 | Request::ApplySnapshotChunk(r) => { |
| 211 | Response::ApplySnapshotChunk(log_error(app.apply_snapshot_chunk(r).await)?) |
| 212 | } |
| 213 | // Some clients still send explicit Flush requests; never panic on protocol input. |
| 214 | Request::Flush => Response::Flush, |
| 215 | }; |
| 216 | Ok(res) |
| 217 | }; |
| 218 | res.boxed() |
| 219 | } |
| 220 | } |
| 221 | |
| 222 | fn log_error<T>(res: AbciResult<T>) -> AbciResult<T> { |
nothing calls this directly
no test coverage detected