Skip to main content

aven_core/sync/
session.rs

1use std::collections::HashSet;
2use std::io::Write;
3use std::path::PathBuf;
4use std::time::Instant;
5
6use anyhow::{Context, Result, bail};
7use flate2::write::GzEncoder;
8
9use super::blob::{
10    MissingLocalBlob, attachment_blob_contract, contract_by_hash, missing_counts,
11    unique_blob_contracts,
12};
13use super::persistence::{ApplySyncPage, ClientSyncPage};
14use super::planner::{
15    PendingChange, TransferBudget, TransferObject, plan_change_prefix, plan_transfers,
16};
17use super::wire::{
18    BlobUploadContract, MAX_BLOB_TRANSFER_BYTES, MAX_BLOB_TRANSFER_OBJECTS, MAX_PULL_BATCH,
19    MAX_PUSH_BATCH, MissingBlobsRequest, MissingBlobsResponse, SyncResponse,
20    sync_server_url_is_valid, validate_blob_hashes,
21};
22use crate::attachments::lifecycle::LifecyclePolicy;
23use crate::db::Database;
24use crate::ids::{new_id, now};
25
26const GZIP_THRESHOLD: usize = 256;
27
28#[derive(Clone, PartialEq, Eq)]
29pub struct SyncRequestContext {
30    session_id: String,
31    request: usize,
32}
33
34#[derive(Clone, PartialEq, Eq)]
35pub struct SyncHttpHeader {
36    pub name: String,
37    pub value: String,
38}
39
40#[derive(Clone)]
41pub struct PreparedSyncRequest {
42    pub method: String,
43    pub url: String,
44    pub headers: Vec<SyncHttpHeader>,
45    pub body: Vec<u8>,
46    pub context: SyncRequestContext,
47}
48
49pub struct SyncHttpResponse {
50    pub status: u16,
51    pub headers: Vec<SyncHttpHeader>,
52    pub body: Vec<u8>,
53}
54
55#[derive(Debug, Clone, PartialEq, Eq)]
56pub struct SyncPageOutcome {
57    pub page: usize,
58    pub pushed: usize,
59    pub pulled: usize,
60    pub blob_uploaded: usize,
61    pub blob_uploaded_bytes: u64,
62    pub blob_downloaded: usize,
63    pub blob_downloaded_bytes: u64,
64    pub cursor: i64,
65    pub complete: bool,
66    pub has_more: bool,
67    pub local_more: bool,
68    pub request_bytes: usize,
69    pub request_wire_bytes: usize,
70    pub response_decoded_bytes: usize,
71    pub response_compression: String,
72    pub apply_ms: u128,
73}
74
75#[derive(Debug, Clone, PartialEq, Eq)]
76pub struct SyncSessionSummary {
77    pub pushed: i64,
78    pub pulled: usize,
79    pub blob_uploaded: usize,
80    pub blob_uploaded_bytes: u64,
81    pub blob_downloaded: usize,
82    pub blob_downloaded_bytes: u64,
83    pub blob_upload_remaining: usize,
84    pub blob_upload_remaining_bytes: u64,
85    pub blob_download_remaining: usize,
86    pub blob_download_remaining_bytes: u64,
87    pub cursor: i64,
88    pub complete: bool,
89    pub pages: usize,
90    pub request_bytes: usize,
91    pub request_wire_bytes: usize,
92    pub response_decoded_bytes: usize,
93    pub response_compression: String,
94    pub apply_ms: u128,
95}
96
97struct ActivePage {
98    page: ClientSyncPage,
99    contracts: Vec<BlobUploadContract>,
100    probed: bool,
101    missing: HashSet<String>,
102    uploads: Vec<BlobUploadContract>,
103    upload_index: usize,
104    confirmed: bool,
105    metadata_sent: bool,
106    downloads: Vec<MissingLocalBlob>,
107    download_index: usize,
108    has_more: bool,
109    page_pushed: usize,
110    page_pulled: usize,
111    page_uploaded: usize,
112    page_uploaded_bytes: u64,
113    page_downloaded: usize,
114    page_downloaded_bytes: u64,
115    transfer_budget_blocked: bool,
116}
117
118impl ActivePage {
119    fn new(mut page: ClientSyncPage) -> Result<Self> {
120        let mut hashes = HashSet::new();
121        let mut count = page.request.changes.len();
122        for (index, change) in page.request.changes.iter().enumerate() {
123            if let Some(contract) = attachment_blob_contract(change)?
124                && hashes.insert(contract.sha256)
125                && hashes.len() > MAX_BLOB_TRANSFER_OBJECTS
126            {
127                count = index;
128                break;
129            }
130        }
131        page.request.changes.truncate(count);
132        page.pending = count;
133        let contracts = unique_blob_contracts(&page.request.changes)?;
134        Ok(Self {
135            page,
136            contracts,
137            probed: false,
138            missing: HashSet::new(),
139            uploads: Vec::new(),
140            upload_index: 0,
141            confirmed: false,
142            metadata_sent: false,
143            downloads: Vec::new(),
144            download_index: 0,
145            has_more: false,
146            page_pushed: 0,
147            page_pulled: 0,
148            page_uploaded: 0,
149            page_uploaded_bytes: 0,
150            page_downloaded: 0,
151            page_downloaded_bytes: 0,
152            transfer_budget_blocked: false,
153        })
154    }
155}
156
157#[derive(Clone)]
158enum RequestKind {
159    Missing,
160    Upload {
161        lease_id: String,
162        contract: BlobUploadContract,
163    },
164    Confirm,
165    Metadata {
166        request_bytes: usize,
167    },
168    Download {
169        blob: MissingLocalBlob,
170    },
171}
172
173struct OutstandingRequest {
174    prepared: PreparedSyncRequest,
175    kind: RequestKind,
176}
177
178pub struct SyncSession {
179    database: Database,
180    server: String,
181    auth_token: Option<String>,
182    blob_dir: PathBuf,
183    lifecycle_policy: LifecyclePolicy,
184    page_budget: Option<usize>,
185    attempted_at: String,
186    session_id: String,
187    request_number: usize,
188    active: Option<ActivePage>,
189    outstanding: Option<OutstandingRequest>,
190    known_server_blobs: HashSet<String>,
191    transfer_budget: TransferBudget,
192    summary: SyncSessionSummary,
193    last_has_more: bool,
194    last_local_more: bool,
195    stopped: bool,
196}
197
198impl SyncSession {
199    pub async fn start(
200        database: Database,
201        server: String,
202        auth_token: Option<String>,
203        page_budget: Option<usize>,
204    ) -> Result<Self> {
205        let blob_dir = crate::attachments::default_blob_dir(database.path());
206        Self::start_with_attachment_storage(
207            database,
208            server,
209            auth_token,
210            page_budget,
211            blob_dir,
212            LifecyclePolicy::default(),
213        )
214        .await
215    }
216
217    pub async fn start_with_attachment_storage(
218        database: Database,
219        server: String,
220        auth_token: Option<String>,
221        page_budget: Option<usize>,
222        blob_dir: PathBuf,
223        lifecycle_policy: LifecyclePolicy,
224    ) -> Result<Self> {
225        let attempted_at = now();
226        database.begin_sync_attempt(attempted_at.clone()).await?;
227        if !sync_server_url_is_valid(&server) {
228            let message = "invalid sync server URL";
229            database.record_sync_error(message.to_string()).await?;
230            bail!("{message}");
231        }
232        Ok(Self {
233            database,
234            server: server.trim_end_matches('/').to_string(),
235            auth_token,
236            blob_dir,
237            lifecycle_policy,
238            page_budget,
239            attempted_at,
240            session_id: new_id(),
241            request_number: 0,
242            active: None,
243            outstanding: None,
244            known_server_blobs: HashSet::new(),
245            transfer_budget: TransferBudget {
246                objects: MAX_BLOB_TRANSFER_OBJECTS,
247                bytes: MAX_BLOB_TRANSFER_BYTES,
248                completed_objects: 0,
249            },
250            summary: SyncSessionSummary {
251                pushed: 0,
252                pulled: 0,
253                blob_uploaded: 0,
254                blob_uploaded_bytes: 0,
255                blob_downloaded: 0,
256                blob_downloaded_bytes: 0,
257                blob_upload_remaining: 0,
258                blob_upload_remaining_bytes: 0,
259                blob_download_remaining: 0,
260                blob_download_remaining_bytes: 0,
261                cursor: 0,
262                complete: false,
263                pages: 0,
264                request_bytes: 0,
265                request_wire_bytes: 0,
266                response_decoded_bytes: 0,
267                response_compression: "none".to_string(),
268                apply_ms: 0,
269            },
270            last_has_more: false,
271            last_local_more: false,
272            stopped: false,
273        })
274    }
275
276    pub async fn prepare_request(&mut self) -> Result<Option<PreparedSyncRequest>> {
277        if self.stopped {
278            return Ok(None);
279        }
280        if let Some(outstanding) = &self.outstanding {
281            return Ok(Some(outstanding.prepared.clone()));
282        }
283        if self.active.is_none() {
284            let page = self
285                .database
286                .prepare_client_sync_page(self.server.clone(), MAX_PUSH_BATCH, MAX_PULL_BATCH)
287                .await?;
288            self.active = Some(ActivePage::new(page)?);
289        }
290
291        let active = self.active.as_ref().expect("active sync page");
292        if !active.contracts.is_empty() && !active.probed {
293            return self.prepare_json_request(
294                "POST",
295                "/sync/blobs/missing",
296                &MissingBlobsRequest {
297                    blobs: active.contracts.clone(),
298                },
299                RequestKind::Missing,
300            );
301        }
302        if active.upload_index < active.uploads.len() {
303            let contract = active.uploads[active.upload_index].clone();
304            let upload = self
305                .database
306                .prepare_blob_upload(&self.blob_dir, &contract)
307                .await?;
308            let mut headers = vec![header("content-type", &contract.media_type)];
309            headers.push(header("x-aven-workspace-id", &contract.workspace_id));
310            headers.push(header("x-aven-byte-size", &contract.byte_size.to_string()));
311            headers.push(header("x-aven-width", &contract.width.to_string()));
312            headers.push(header("x-aven-height", &contract.height.to_string()));
313            return self.prepare_raw_request(
314                "PUT",
315                &format!("/sync/blobs/{}", contract.sha256),
316                headers,
317                upload.bytes,
318                RequestKind::Upload {
319                    lease_id: upload.lease_id,
320                    contract,
321                },
322            );
323        }
324        if !active.contracts.is_empty() && !active.confirmed {
325            return self.prepare_json_request(
326                "POST",
327                "/sync/blobs/missing",
328                &MissingBlobsRequest {
329                    blobs: active.contracts.clone(),
330                },
331                RequestKind::Confirm,
332            );
333        }
334        if !active.metadata_sent {
335            let body = serde_json::to_vec(&active.page.request).context("encode sync request")?;
336            let request_bytes = body.len();
337            let mut headers = vec![header("content-type", "application/json")];
338            let body = if request_bytes > GZIP_THRESHOLD {
339                headers.push(header("content-encoding", "gzip"));
340                gzip_encode(&body)?
341            } else {
342                body
343            };
344            return self.prepare_raw_request(
345                "POST",
346                "/sync",
347                headers,
348                body,
349                RequestKind::Metadata { request_bytes },
350            );
351        }
352        if active.download_index < active.downloads.len() {
353            let blob = active.downloads[active.download_index].clone();
354            return self.prepare_raw_request(
355                "GET",
356                &format!("/sync/blobs/{}", blob.sha256),
357                vec![header("accept-encoding", "identity")],
358                Vec::new(),
359                RequestKind::Download { blob },
360            );
361        }
362
363        self.finish_page().await?;
364        if self.stopped {
365            Ok(None)
366        } else {
367            Box::pin(self.prepare_request()).await
368        }
369    }
370
371    fn prepare_json_request<T: serde::Serialize>(
372        &mut self,
373        method: &str,
374        path: &str,
375        value: &T,
376        kind: RequestKind,
377    ) -> Result<Option<PreparedSyncRequest>> {
378        self.prepare_raw_request(
379            method,
380            path,
381            vec![header("content-type", "application/json")],
382            serde_json::to_vec(value)?,
383            kind,
384        )
385    }
386
387    fn prepare_raw_request(
388        &mut self,
389        method: &str,
390        path: &str,
391        mut headers: Vec<SyncHttpHeader>,
392        body: Vec<u8>,
393        kind: RequestKind,
394    ) -> Result<Option<PreparedSyncRequest>> {
395        if let Some(token) = &self.auth_token {
396            headers.push(header("authorization", &format!("Bearer {token}")));
397        }
398        self.request_number += 1;
399        let prepared = PreparedSyncRequest {
400            method: method.to_string(),
401            url: format!("{}{}", self.server, path),
402            headers,
403            body,
404            context: SyncRequestContext {
405                session_id: self.session_id.clone(),
406                request: self.request_number,
407            },
408        };
409        self.outstanding = Some(OutstandingRequest {
410            prepared: prepared.clone(),
411            kind,
412        });
413        Ok(Some(prepared))
414    }
415
416    pub async fn accept_response(
417        &mut self,
418        context: &SyncRequestContext,
419        response: SyncHttpResponse,
420    ) -> Result<SyncPageOutcome> {
421        {
422            let outstanding = self
423                .outstanding
424                .as_ref()
425                .context("no outstanding sync request")?;
426            validate_context(&outstanding.prepared.context, context)?;
427        }
428        if !(200..300).contains(&response.status) {
429            if response.status == 507 {
430                bail!("error attachment-quota-exceeded");
431            }
432            bail!("sync HTTP status {}", response.status);
433        }
434        let outstanding = self
435            .outstanding
436            .as_ref()
437            .expect("validated outstanding request");
438        let prepared_body_len = outstanding.prepared.body.len();
439        let request_kind = outstanding.kind.clone();
440        let response_bytes = response.body.len();
441        match request_kind {
442            RequestKind::Missing => {
443                let decoded: MissingBlobsResponse = serde_json::from_slice(&response.body)
444                    .context("decode missing blobs response")?;
445                validate_blob_hashes(&decoded.missing)?;
446                let active = self.active.as_mut().expect("active sync page");
447                let requested = active
448                    .contracts
449                    .iter()
450                    .map(|c| c.sha256.as_str())
451                    .collect::<HashSet<_>>();
452                if decoded
453                    .missing
454                    .iter()
455                    .any(|hash| !requested.contains(hash.as_str()))
456                {
457                    bail!("error invalid-blob-missing-response");
458                }
459                let missing = decoded.missing.into_iter().collect::<HashSet<_>>();
460                self.known_server_blobs.extend(
461                    active
462                        .contracts
463                        .iter()
464                        .filter(|contract| !missing.contains(&contract.sha256))
465                        .map(|contract| contract.sha256.clone()),
466                );
467                active.probed = true;
468                active.missing = missing;
469                let pending = active
470                    .page
471                    .request
472                    .changes
473                    .iter()
474                    .map(|change| {
475                        let blob = match attachment_blob_contract(change)? {
476                            Some(contract) if active.missing.contains(&contract.sha256) => {
477                                Some(TransferObject {
478                                    sha256: contract.sha256,
479                                    byte_size: u64::try_from(contract.byte_size)
480                                        .context("attachment bytes exceed u64")?,
481                                })
482                            }
483                            _ => None,
484                        };
485                        Ok(PendingChange { missing_blob: blob })
486                    })
487                    .collect::<Result<Vec<_>>>()?;
488                let plan = plan_change_prefix(&pending, self.transfer_budget);
489                active.transfer_budget_blocked = plan.change_count < pending.len();
490                active.page.request.changes.truncate(plan.change_count);
491                active.page.pending = plan.change_count;
492                active.contracts = unique_blob_contracts(&active.page.request.changes)?;
493                active.uploads = plan
494                    .transfers
495                    .iter()
496                    .map(|object| contract_by_hash(&active.contracts, &object.sha256))
497                    .collect::<Result<Vec<_>>>()?;
498            }
499            RequestKind::Upload { lease_id, contract } => {
500                self.database.finish_blob_upload(&lease_id).await?;
501                let bytes =
502                    u64::try_from(contract.byte_size).context("attachment bytes exceed u64")?;
503                self.known_server_blobs.insert(contract.sha256.clone());
504                self.transfer_budget.consume(&TransferObject {
505                    sha256: contract.sha256,
506                    byte_size: bytes,
507                });
508                let active = self.active.as_mut().expect("active sync page");
509                active.upload_index += 1;
510                active.page_uploaded += 1;
511                active.page_uploaded_bytes += bytes;
512                self.summary.blob_uploaded += 1;
513                self.summary.blob_uploaded_bytes += bytes;
514            }
515            RequestKind::Confirm => {
516                let decoded: MissingBlobsResponse = serde_json::from_slice(&response.body)
517                    .context("decode missing blobs response")?;
518                validate_blob_hashes(&decoded.missing)?;
519                if !decoded.missing.is_empty() {
520                    bail!("error attachment-blob-admission-missing");
521                }
522                self.active.as_mut().expect("active sync page").confirmed = true;
523            }
524            RequestKind::Metadata { request_bytes } => {
525                let decoded: SyncResponse =
526                    serde_json::from_slice(&response.body).context("decode sync response")?;
527                let active = self.active.as_mut().expect("active sync page");
528                let request = active.page.request.clone();
529                let pending = active.page.pending;
530                let has_more = decoded.has_more;
531                let cursor = decoded.cursor;
532                let apply_started = Instant::now();
533                let pulled = self
534                    .database
535                    .apply_client_sync_page(ApplySyncPage {
536                        request,
537                        response: decoded,
538                        attempted_at: self.attempted_at.clone(),
539                        previous_pushed: self.summary.pushed,
540                        previous_pulled: self.summary.pulled,
541                    })
542                    .await?;
543                let apply_ms = apply_started.elapsed().as_millis();
544                active.metadata_sent = true;
545                active.has_more = has_more;
546                active.page_pushed = pending;
547                active.page_pulled = pulled;
548                self.summary.pushed += pending as i64;
549                self.summary.pulled += pulled;
550                self.summary.cursor = cursor;
551                self.summary.request_bytes += request_bytes;
552                self.summary.request_wire_bytes += prepared_body_len;
553                self.summary.response_decoded_bytes += response_bytes;
554                self.summary.response_compression =
555                    header_value(&response.headers, "content-encoding")
556                        .unwrap_or("none")
557                        .to_string();
558                self.summary.apply_ms += apply_ms;
559                let missing = self.database.missing_local_blobs(&self.blob_dir).await?;
560                let objects = missing
561                    .iter()
562                    .map(|blob| {
563                        Ok(TransferObject {
564                            sha256: blob.sha256.clone(),
565                            byte_size: u64::try_from(blob.byte_size)
566                                .context("attachment bytes exceed u64")?,
567                        })
568                    })
569                    .collect::<Result<Vec<_>>>()?;
570                let plan = plan_transfers(&objects, self.transfer_budget);
571                active.transfer_budget_blocked |= plan.len() < objects.len();
572                let planned = plan
573                    .iter()
574                    .map(|o| o.sha256.as_str())
575                    .collect::<HashSet<_>>();
576                active.downloads = missing
577                    .into_iter()
578                    .filter(|b| planned.contains(b.sha256.as_str()))
579                    .collect();
580            }
581            RequestKind::Download { blob } => {
582                let expected =
583                    usize::try_from(blob.byte_size).context("attachment bytes exceed usize")?;
584                let content_length_valid = header_value(&response.headers, "content-length")
585                    .map(|length| length.parse::<usize>().ok() == Some(expected))
586                    .unwrap_or(true);
587                if !content_length_valid || response.body.len() != expected {
588                    bail!("error attachment-blob-remote-invalid");
589                }
590                self.database
591                    .store_downloaded_blob(
592                        &self.blob_dir,
593                        self.lifecycle_policy,
594                        &blob,
595                        response.body,
596                    )
597                    .await?;
598                let bytes = u64::try_from(blob.byte_size).context("attachment bytes exceed u64")?;
599                self.transfer_budget.consume(&TransferObject {
600                    sha256: blob.sha256,
601                    byte_size: bytes,
602                });
603                let active = self.active.as_mut().expect("active sync page");
604                active.download_index += 1;
605                active.page_downloaded += 1;
606                active.page_downloaded_bytes += bytes;
607                self.summary.blob_downloaded += 1;
608                self.summary.blob_downloaded_bytes += bytes;
609            }
610        }
611        self.outstanding = None;
612        let page_finished = self.active.as_ref().is_some_and(|active| {
613            active.metadata_sent && active.download_index >= active.downloads.len()
614        });
615        if page_finished {
616            let mut outcome = self.outcome();
617            self.finish_page().await?;
618            outcome.page = self.summary.pages;
619            outcome.complete = self.summary.complete;
620            outcome.has_more = self.last_has_more;
621            outcome.local_more = self.last_local_more;
622            return Ok(outcome);
623        }
624        Ok(self.outcome())
625    }
626
627    async fn finish_page(&mut self) -> Result<()> {
628        let active = self.active.take().expect("active sync page");
629        self.summary.pages += 1;
630        self.last_has_more = active.has_more;
631        let local_more = self.database.pending_sync_change_count().await? > 0;
632        self.last_local_more = local_more;
633        let downloads = self.database.missing_local_blobs(&self.blob_dir).await?;
634        let download_counts = missing_counts(&downloads);
635        self.summary.blob_download_remaining = download_counts.count as usize;
636        self.summary.blob_download_remaining_bytes = download_counts.bytes;
637        let uploads = self
638            .database
639            .pending_blob_contracts()
640            .await?
641            .into_iter()
642            .filter(|contract| !self.known_server_blobs.contains(&contract.sha256))
643            .collect::<Vec<_>>();
644        self.summary.blob_upload_remaining = uploads.len();
645        self.summary.blob_upload_remaining_bytes = uploads
646            .iter()
647            .map(|c| u64::try_from(c.byte_size).unwrap_or(0))
648            .sum();
649        let complete = !local_more && !active.has_more && downloads.is_empty();
650        self.summary.complete = complete;
651        let budget_exhausted = self
652            .page_budget
653            .is_some_and(|budget| self.summary.pages >= budget);
654        let transfer_stalled = active.transfer_budget_blocked
655            || self.transfer_budget.objects == 0
656            || self.transfer_budget.bytes == 0;
657        self.stopped = complete || budget_exhausted || transfer_stalled;
658        Ok(())
659    }
660
661    fn outcome(&self) -> SyncPageOutcome {
662        let active = self.active.as_ref();
663        SyncPageOutcome {
664            page: self.summary.pages + usize::from(active.is_some_and(|page| page.metadata_sent)),
665            pushed: active.map_or(0, |page| page.page_pushed),
666            pulled: active.map_or(0, |page| page.page_pulled),
667            blob_uploaded: active.map_or(0, |page| page.page_uploaded),
668            blob_uploaded_bytes: active.map_or(0, |page| page.page_uploaded_bytes),
669            blob_downloaded: active.map_or(0, |page| page.page_downloaded),
670            blob_downloaded_bytes: active.map_or(0, |page| page.page_downloaded_bytes),
671            cursor: self.summary.cursor,
672            complete: self.summary.complete,
673            has_more: active.map_or(self.last_has_more, |page| page.has_more),
674            local_more: self.last_local_more,
675            request_bytes: self.summary.request_bytes,
676            request_wire_bytes: self.summary.request_wire_bytes,
677            response_decoded_bytes: self.summary.response_decoded_bytes,
678            response_compression: self.summary.response_compression.clone(),
679            apply_ms: self.summary.apply_ms,
680        }
681    }
682
683    pub async fn fail_request(
684        &mut self,
685        context: &SyncRequestContext,
686        error: impl Into<String>,
687    ) -> Result<()> {
688        let outstanding = self
689            .outstanding
690            .as_ref()
691            .context("no outstanding sync request")?;
692        validate_context(&outstanding.prepared.context, context)?;
693        if let RequestKind::Upload { lease_id, .. } = &outstanding.kind {
694            self.database.finish_blob_upload(lease_id).await?;
695        }
696        self.database.record_sync_error(error.into()).await?;
697        self.outstanding = None;
698        self.stopped = true;
699        Ok(())
700    }
701
702    pub fn summary(&self) -> SyncSessionSummary {
703        self.summary.clone()
704    }
705}
706
707fn header(name: &str, value: &str) -> SyncHttpHeader {
708    SyncHttpHeader {
709        name: name.to_string(),
710        value: value.to_string(),
711    }
712}
713
714fn validate_context(expected: &SyncRequestContext, actual: &SyncRequestContext) -> Result<()> {
715    if expected != actual {
716        bail!("sync request context does not match the outstanding page");
717    }
718    Ok(())
719}
720
721fn header_value<'a>(headers: &'a [SyncHttpHeader], name: &str) -> Option<&'a str> {
722    headers
723        .iter()
724        .find(|header| header.name.eq_ignore_ascii_case(name))
725        .map(|header| header.value.as_str())
726}
727
728fn gzip_encode(body: &[u8]) -> Result<Vec<u8>> {
729    let mut encoder = GzEncoder::new(Vec::new(), flate2::Compression::default());
730    encoder
731        .write_all(body)
732        .context("gzip encode sync request body")?;
733    encoder.finish().context("finish sync request compression")
734}