Skip to main content

gix_protocol/fetch/arguments/
io.rs

1#[crate::bisync::only_async]
2use futures_lite::io::AsyncWriteExt;
3#[crate::bisync::only_sync]
4use std::io::Write;
5
6#[crate::bisync::only_async]
7use gix_transport::client::{
8    self,
9    async_io::{ExtendedBufRead, Transport, TransportV2Ext},
10};
11#[crate::bisync::only_sync]
12use gix_transport::client::{
13    self,
14    blocking_io::{ExtendedBufRead, Transport, TransportV2Ext},
15};
16
17use crate::{Command, fetch::Arguments};
18
19impl Arguments {
20    /// Send fetch arguments to the server, and indicate this is the end of negotiations only if `add_done_argument` is present.
21    #[crate::bisync::bisync]
22    pub async fn send<'a, T: Transport + 'a>(
23        &mut self,
24        transport: &'a mut T,
25        add_done_argument: bool,
26    ) -> Result<Box<dyn ExtendedBufRead<'a> + Unpin + 'a>, client::Error> {
27        if self.haves.is_empty() {
28            assert!(add_done_argument, "If there are no haves, is_done must be true.");
29        }
30        match self.version {
31            gix_transport::Protocol::V0 | gix_transport::Protocol::V1 => {
32                let (on_into_read, retained_state) = self.prepare_v1(
33                    transport.connection_persists_across_multiple_requests(),
34                    add_done_argument,
35                )?;
36                let mut line_writer = transport.request(
37                    client::WriteMode::OneLfTerminatedLinePerWriteCall,
38                    on_into_read,
39                    self.trace,
40                )?;
41                let had_args = !self.args.is_empty();
42                for arg in self.args.drain(..) {
43                    line_writer.write_all(&arg).await?;
44                }
45                if had_args {
46                    line_writer.write_message(client::MessageKind::Flush).await?;
47                }
48                for line in self.haves.drain(..) {
49                    line_writer.write_all(&line).await?;
50                }
51                if let Some(next_args) = retained_state {
52                    self.args = next_args;
53                }
54                Ok(line_writer.into_read().await?)
55            }
56            gix_transport::Protocol::V2 => {
57                let retained_state = self.args.clone();
58                self.args.append(&mut self.haves);
59                if add_done_argument {
60                    self.args.push("done".into());
61                }
62                transport
63                    .invoke(
64                        Command::Fetch.as_str(),
65                        self.features.iter().filter(|(_, v)| v.is_some()).cloned(),
66                        Some(std::mem::replace(&mut self.args, retained_state).into_iter()),
67                        self.trace,
68                    )
69                    .await
70            }
71        }
72    }
73}