1use std::collections::{HashMap, HashSet};
2use std::sync::Arc;
3
4use arc_swap::{ArcSwap, ArcSwapOption};
5use async_trait::async_trait;
6use bytes::Bytes;
7use googleapis_tonic_google_storage_v2::google::storage::v2::{
8 bidi_write_object_request, bidi_write_object_response, storage_client::StorageClient,
9 write_object_request, write_object_response, AppendObjectSpec, BidiReadObjectRequest,
10 BidiReadObjectSpec, BidiWriteHandle, BidiWriteObjectRedirectedError, BidiWriteObjectRequest,
11 BidiWriteObjectResponse, ChecksummedData, DeleteObjectRequest, GetObjectRequest,
12 ListObjectsRequest, Object, ReadRange, UpdateObjectRequest, WriteObjectRequest,
13 WriteObjectSpec,
14};
15use prost::Message as _;
16use tokio::sync::{mpsc, watch, Mutex as SessionMutex};
17use tokio_stream::wrappers::ReceiverStream;
18use tonic::metadata::MetadataValue;
19use tonic::transport::{Channel, ClientTlsConfig, Endpoint};
20use tonic::{Code, Request, Status, Streaming};
21
22use crate::auth::BearerAuth;
23use crate::error::Error;
24use crate::transport::{
25 AppendToken, LaneDurableChange, ListedObject, PackedAppend, PackedAppendMessage, Replica,
26 ReplicaFactory, ReplicaSnapshot, TransportCode, TransportError,
27};
28
29type RoutingToken = Arc<ArcSwap<Option<Arc<String>>>>;
33
34#[derive(Clone)]
35struct GrpcReplica {
36 zone: usize,
37 bucket: String,
38 object: String,
39 auth: Option<BearerAuth>,
40 client: StorageClient<Channel>,
41 routing_token: RoutingToken,
42 session: Arc<SessionSlot>,
50}
51
52#[derive(Clone)]
53pub struct GrpcReplicaFactory {
56 zone: usize,
57 bucket: String,
58 auth: Option<BearerAuth>,
59 client: StorageClient<Channel>,
60 routing_token: RoutingToken,
61}
62
63struct AppendSession {
68 handle: Arc<AppendSessionHandle>,
69 reader: tokio::task::JoinHandle<()>,
70}
71
72struct AppendSessionHandle {
73 tx: mpsc::Sender<BidiWriteObjectRequest>,
74 state: tokio::sync::watch::Receiver<LaneProgress>,
75}
76
77struct RedirectAwareStream {
78 tx: Option<mpsc::Sender<BidiWriteObjectRequest>>,
79 responses: Streaming<BidiWriteObjectResponse>,
80 first: BidiWriteObjectResponse,
81 attempt: u32,
82}
83
84struct SessionWait {
85 timeout: Option<std::time::Duration>,
86 clear_on_ready: bool,
87 inspect_before_error: bool,
88 reader_ended: &'static str,
89 stalled: &'static str,
90}
91
92struct SessionSlot {
93 current: ArcSwapOption<AppendSessionHandle>,
94 owned: SessionMutex<Option<AppendSession>>,
95 shutdown: watch::Sender<bool>,
96}
97
98impl SessionSlot {
99 fn new() -> Self {
100 let (shutdown, _) = watch::channel(false);
101 Self {
102 current: ArcSwapOption::empty(),
103 owned: SessionMutex::new(None),
104 shutdown,
105 }
106 }
107}
108
109impl Drop for AppendSession {
110 fn drop(&mut self) {
111 self.reader.abort();
112 }
113}
114
115impl AppendSession {
116 async fn shutdown(mut self) {
117 self.reader.abort();
118 let _ = (&mut self.reader).await;
119 }
120}
121
122#[derive(Clone, Debug, Default)]
123struct LaneProgress {
126 durable: i64,
127 finalized: Option<ReplicaSnapshot>,
128 error: Option<TransportError>,
129}
130
131const WIRE_MESSAGE_TARGET_BYTES: usize = 262_144;
138
139pub(crate) fn pack_append(chunks: Vec<Bytes>) -> PackedAppend {
140 let total_len = chunks.iter().map(Bytes::len).sum::<usize>();
141 let mut packed =
142 bytes::BytesMut::with_capacity(total_len.min(WIRE_MESSAGE_TARGET_BYTES.saturating_mul(2)));
143 let mut messages = Vec::new();
144 let mut relative_offset = 0i64;
145 for data in &chunks {
146 if !packed.is_empty() && packed.len() + data.len() > WIRE_MESSAGE_TARGET_BYTES {
147 let content = packed.split().freeze();
148 let len = content.len() as i64;
149 messages.push(PackedAppendMessage {
150 relative_offset,
151 crc32c: crc32c::crc32c(&content),
152 content,
153 });
154 relative_offset += len;
155 }
156 packed.extend_from_slice(data);
157 }
158 if !packed.is_empty() {
159 let content = packed.freeze();
160 messages.push(PackedAppendMessage {
161 relative_offset,
162 crc32c: crc32c::crc32c(&content),
163 content,
164 });
165 }
166 PackedAppend::new(chunks, messages, total_len)
167}
168
169const SESSION_PROGRESS_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
173const RPC_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
174const MAX_LIST_PAGES: usize = 10_000;
175
176impl AppendSessionHandle {
177 async fn send(
178 self: &Arc<Self>,
179 replica: &GrpcReplica,
180 request: BidiWriteObjectRequest,
181 disconnected: &'static str,
182 ) -> Result<(), TransportError> {
183 match tokio::time::timeout(SESSION_PROGRESS_TIMEOUT, self.tx.send(request)).await {
184 Ok(Ok(())) => Ok(()),
185 Ok(Err(_)) => {
186 replica.clear_session_if(self).await;
187 Err(replica.error(TransportCode::Unavailable, disconnected))
188 }
189 Err(_) => {
190 replica.clear_session_if(self).await;
191 Err(replica.error(
192 TransportCode::DeadlineExceeded,
193 "append session request channel made no progress",
194 ))
195 }
196 }
197 }
198
199 async fn wait_for<T>(
200 self: &Arc<Self>,
201 replica: &GrpcReplica,
202 wait: SessionWait,
203 mut inspect: impl FnMut(&LaneProgress) -> Option<Result<T, TransportError>>,
204 ) -> Result<T, TransportError> {
205 let mut state = self.state.clone();
206 loop {
207 let progress = state.borrow_and_update().clone();
208 if wait.inspect_before_error {
209 if let Some(result) = inspect(&progress) {
210 if wait.clear_on_ready || result.is_err() {
211 replica.clear_session_if(self).await;
212 }
213 return result;
214 }
215 }
216 if let Some(error) = progress.error {
217 replica.clear_session_if(self).await;
218 return Err(error);
219 }
220 if !wait.inspect_before_error {
221 if let Some(result) = inspect(&progress) {
222 if wait.clear_on_ready || result.is_err() {
223 replica.clear_session_if(self).await;
224 }
225 return result;
226 }
227 }
228 let changed = match wait.timeout {
229 Some(timeout) => match tokio::time::timeout(timeout, state.changed()).await {
230 Ok(changed) => changed,
231 Err(_) => {
232 replica.clear_session_if(self).await;
233 return Err(replica.error(TransportCode::DeadlineExceeded, wait.stalled));
234 }
235 },
236 None => state.changed().await,
237 };
238 if changed.is_err() {
239 replica.clear_session_if(self).await;
240 return Err(replica.error(TransportCode::Unavailable, wait.reader_ended));
241 }
242 }
243 }
244}
245
246impl GrpcReplica {
247 async fn bidi_read(&self, generation: i64) -> Result<Vec<u8>, TransportError> {
261 let spec = BidiReadObjectSpec {
262 bucket: self.bucket.clone(),
263 object: self.object.clone(),
264 generation,
265 if_generation_match: Some(generation),
266 ..Default::default()
267 };
268 let open = BidiReadObjectRequest {
269 read_object_spec: Some(spec),
270 read_ranges: vec![ReadRange {
272 read_offset: 0,
273 read_length: 0,
274 read_id: 1,
275 }],
276 };
277 let request = self.request(tokio_stream::once(open))?;
278 let mut stream = self
279 .client
280 .clone()
281 .bidi_read_object(request)
282 .await
283 .map_err(|status| self.status(status))?
284 .into_inner();
285 let mut bytes = Vec::new();
286 while let Some(response) = stream
287 .message()
288 .await
289 .map_err(|status| self.status(status))?
290 {
291 for range in response.object_data_ranges {
292 if let Some(data) = range.checksummed_data {
293 if data
294 .crc32c
295 .is_some_and(|expected| expected != crc32c::crc32c(&data.content))
296 {
297 return Err(self.error(
298 TransportCode::DataLoss,
299 "BidiReadObject response CRC32C mismatch",
300 ));
301 }
302 bytes.extend_from_slice(&data.content);
303 }
304 }
305 }
306 Ok(bytes)
307 }
308
309 fn live_session(&self) -> Result<Arc<AppendSessionHandle>, TransportError> {
310 self.session.current.load_full().ok_or_else(|| {
311 self.error(
312 TransportCode::Unavailable,
313 "no live append session (resume required)",
314 )
315 })
316 }
317
318 fn replace_session_locked(
319 &self,
320 owned: &mut Option<AppendSession>,
321 session: Option<AppendSession>,
322 ) -> Option<AppendSession> {
323 let previous = owned.take();
324 self.session
325 .current
326 .store(session.as_ref().map(|session| Arc::clone(&session.handle)));
327 *owned = session;
328 previous
329 }
330
331 async fn replace_session(&self, session: Option<AppendSession>) {
332 let previous = {
333 let mut owned = self.session.owned.lock().await;
334 self.replace_session_locked(&mut owned, session)
335 };
336 if let Some(previous) = previous {
337 previous.shutdown().await;
338 }
339 }
340
341 async fn clear_session_if(&self, expected: &Arc<AppendSessionHandle>) {
342 let previous = {
343 let mut owned = self.session.owned.lock().await;
344 if owned
345 .as_ref()
346 .is_some_and(|session| Arc::ptr_eq(&session.handle, expected))
347 {
348 self.replace_session_locked(&mut owned, None)
349 } else {
350 None
351 }
352 };
353 if let Some(previous) = previous {
354 previous.shutdown().await;
355 }
356 }
357
358 fn request_no_deadline<T>(&self, value: T) -> Result<Request<T>, TransportError> {
362 let mut request = self.request(value)?;
363 request.metadata_mut().remove("grpc-timeout");
364 Ok(request)
365 }
366
367 fn request<T>(&self, value: T) -> Result<Request<T>, TransportError> {
368 let mut request = Request::new(value);
369 request.set_timeout(RPC_TIMEOUT);
370 let guard = self.routing_token.load();
375 let value = match (**guard).as_ref() {
376 Some(token) => format!("bucket={}&routing_token={}", self.bucket, token),
377 None => format!("bucket={}", self.bucket),
378 };
379 let params = MetadataValue::try_from(value).map_err(|_| {
380 self.error(
381 TransportCode::Internal,
382 "request params are not valid gRPC metadata",
383 )
384 })?;
385 request
386 .metadata_mut()
387 .insert("x-goog-request-params", params);
388 if let Some(auth) = &self.auth {
389 let value = auth
390 .authorization_header()
391 .map_err(|error| self.error(TransportCode::Unauthenticated, error.to_string()))?;
392 request.metadata_mut().insert("authorization", value);
393 }
394 Ok(request)
395 }
396
397 fn error(&self, code: TransportCode, message: impl Into<String>) -> TransportError {
398 TransportError {
399 zone: self.zone,
400 code,
401 message: message.into(),
402 }
403 }
404
405 fn status(&self, status: Status) -> TransportError {
406 let code = match status.code() {
412 Code::NotFound => TransportCode::NotFound,
413 Code::AlreadyExists => TransportCode::AlreadyExists,
414 Code::InvalidArgument => TransportCode::InvalidArgument,
415 Code::FailedPrecondition => TransportCode::FailedPrecondition,
416 Code::Aborted => TransportCode::Aborted,
417 Code::OutOfRange => TransportCode::OutOfRange,
418 Code::ResourceExhausted => TransportCode::ResourceExhausted,
419 Code::Unimplemented => TransportCode::Unimplemented,
420 Code::DataLoss => TransportCode::DataLoss,
421 Code::Unauthenticated => TransportCode::Unauthenticated,
422 Code::PermissionDenied => TransportCode::PermissionDenied,
423 Code::Unavailable => TransportCode::Unavailable,
424 Code::DeadlineExceeded => TransportCode::DeadlineExceeded,
425 _ => TransportCode::Internal,
426 };
427 self.error(code, status.message())
428 }
429
430 fn snapshot_from_object(&self, object: Object, bytes: Vec<u8>) -> ReplicaSnapshot {
431 let crc32c = object
432 .checksums
433 .as_ref()
434 .and_then(|checksums| checksums.crc32c);
435 ReplicaSnapshot {
436 zone: self.zone,
437 generation: object.generation,
438 metageneration: object.metageneration,
439 persisted_size: bytes.len() as i64,
440 finalized: object.finalize_time.is_some(),
441 crc32c,
442 metadata: object.metadata,
443 bytes,
444 }
445 }
446
447 fn stat_from_object(&self, object: Object) -> ReplicaSnapshot {
448 let size = object.size;
449 let mut snapshot = self.snapshot_from_object(object, Vec::new());
450 if snapshot.finalized {
451 snapshot.persisted_size = size;
452 }
453 snapshot
454 }
455
456 fn try_capture_redirect(&self, status: &Status) -> bool {
462 match redirect_routing_token(status) {
463 Some(token) => {
464 self.routing_token.store(Arc::new(Some(Arc::new(token))));
465 true
466 }
467 None => false,
468 }
469 }
470
471 async fn open_redirect_aware_stream(
478 &self,
479 requests: &[BidiWriteObjectRequest],
480 persistent: bool,
481 ) -> Result<Option<RedirectAwareStream>, TransportError> {
482 let mut attempt = 0u32;
483 let mut shutdown = self.session.shutdown.subscribe();
484 loop {
485 if *shutdown.borrow_and_update() {
486 return Err(self.error(
487 TransportCode::Unavailable,
488 "append session open cancelled by shutdown",
489 ));
490 }
491 attempt += 1;
492 let (tx, rx) = mpsc::channel(requests.len().max(64));
493 for request in requests {
494 if tx.send(request.clone()).await.is_err() {
495 return Err(self.error(TransportCode::Internal, "write stream channel closed"));
496 }
497 }
498 let request = if persistent {
499 self.request_no_deadline(ReceiverStream::new(rx))?
500 } else {
501 self.request(ReceiverStream::new(rx))?
502 };
503 let mut tx = if persistent {
504 Some(tx)
505 } else {
506 drop(tx);
507 None
508 };
509 let mut client = self.client.clone();
510 let mut opening = Box::pin(client.bidi_write_object(request));
511 let opened = if persistent {
512 tokio::select! {
515 biased;
516 changed = shutdown.changed() => {
517 tx.take();
518 if changed.is_ok() {
519 if let Ok(Ok(response)) =
520 tokio::time::timeout(SESSION_PROGRESS_TIMEOUT, &mut opening).await
521 {
522 let mut responses = response.into_inner();
523 let _ = tokio::time::timeout(
524 SESSION_PROGRESS_TIMEOUT,
525 responses.message(),
526 )
527 .await;
528 }
529 }
530 return Err(self.error(
531 TransportCode::Unavailable,
532 "append session open cancelled by shutdown",
533 ));
534 }
535 opened = tokio::time::timeout(SESSION_PROGRESS_TIMEOUT, &mut opening) => {
536 match opened {
537 Ok(opened) => opened,
538 Err(_) => {
539 return Err(self.error(
540 TransportCode::DeadlineExceeded,
541 "append session open timed out",
542 ));
543 }
544 }
545 }
546 }
547 } else {
548 opening.await
549 };
550 let mut responses = match opened {
551 Ok(response) => response.into_inner(),
552 Err(status) => {
553 if attempt <= MAX_WRITE_REDIRECTS && self.try_capture_redirect(&status) {
554 continue;
555 }
556 return Err(self.status(status));
557 }
558 };
559 let first = if persistent {
560 tokio::select! {
561 biased;
562 changed = shutdown.changed() => {
563 tx.take();
564 if changed.is_ok() {
565 let _ = tokio::time::timeout(
566 SESSION_PROGRESS_TIMEOUT,
567 responses.message(),
568 )
569 .await;
570 }
571 return Err(self.error(
572 TransportCode::Unavailable,
573 "append session open cancelled by shutdown",
574 ));
575 }
576 response = tokio::time::timeout(
577 SESSION_PROGRESS_TIMEOUT,
578 responses.message(),
579 ) => {
580 match response {
581 Ok(response) => response,
582 Err(_) => {
583 return Err(self.error(
584 TransportCode::DeadlineExceeded,
585 "append open made no progress",
586 ));
587 }
588 }
589 }
590 }
591 } else {
592 responses.message().await
593 };
594 match first {
595 Ok(Some(first)) => {
596 return Ok(Some(RedirectAwareStream {
597 tx,
598 responses,
599 first,
600 attempt,
601 }));
602 }
603 Ok(None) => return Ok(None),
604 Err(status)
605 if attempt <= MAX_WRITE_REDIRECTS && self.try_capture_redirect(&status) =>
606 {
607 continue;
608 }
609 Err(status) => return Err(self.status(status)),
610 }
611 }
612 }
613
614 async fn open_session(
619 &self,
620 first: BidiWriteObjectRequest,
621 ) -> Result<(AppendSession, i64, Option<Bytes>), TransportError> {
622 let (session, persisted_size, write_handle, _) = self.open_session_observed(first).await?;
623 Ok((session, persisted_size, write_handle))
624 }
625
626 async fn open_session_observed(
627 &self,
628 first: BidiWriteObjectRequest,
629 ) -> Result<(AppendSession, i64, Option<Bytes>, Option<i64>), TransportError> {
630 let Some(opened) = self
631 .open_redirect_aware_stream(std::slice::from_ref(&first), true)
632 .await?
633 else {
634 return Err(self.error(
635 TransportCode::Unavailable,
636 "append open closed without a response",
637 ));
638 };
639 let opening_has_status = opened.first.write_status.is_some();
640 let opening_resource_size = match opened.first.write_status.as_ref() {
641 Some(bidi_write_object_response::WriteStatus::Resource(resource)) => {
642 Some(resource.size)
643 }
644 _ => None,
645 };
646 let handle = opened
647 .first
648 .write_handle
649 .clone()
650 .map(|handle| handle.handle);
651 let mut progress = LaneProgress {
652 durable: if opening_has_status {
653 0
654 } else {
655 first.write_offset
656 },
657 finalized: None,
658 error: None,
659 };
660 self.fold_session_progress(&mut progress, opened.first);
661 let persisted = progress.durable;
662 let (state_tx, state_rx) = tokio::sync::watch::channel(progress);
663 let this = self.clone();
664 let reader = tokio::spawn(async move {
665 let mut responses = opened.responses;
666 loop {
667 match responses.message().await {
668 Ok(Some(response)) => {
669 state_tx.send_modify(|progress| {
670 this.fold_session_progress(progress, response);
671 });
672 }
673 Ok(None) => {
674 state_tx.send_modify(|progress| {
675 progress.error = Some(this.error(
676 TransportCode::Unavailable,
677 "append session closed by the service",
678 ));
679 });
680 return;
681 }
682 Err(status) => {
683 this.try_capture_redirect(&status);
684 let error = this.status(status);
685 state_tx.send_modify(|progress| {
686 progress.error = Some(error.clone());
687 });
688 return;
689 }
690 }
691 }
692 });
693 tracing::debug!(
694 zone = self.zone,
695 persisted_size = persisted,
696 has_write_handle = handle.is_some(),
697 attempt = opened.attempt,
698 "append session opened"
699 );
700 let tx = opened.tx.ok_or_else(|| {
701 self.error(
702 TransportCode::Internal,
703 "persistent write stream omitted its request sender",
704 )
705 })?;
706 Ok((
707 AppendSession {
708 handle: Arc::new(AppendSessionHandle {
709 tx,
710 state: state_rx,
711 }),
712 reader,
713 },
714 persisted,
715 handle,
716 opening_resource_size,
717 ))
718 }
719
720 fn fold_session_progress(
721 &self,
722 progress: &mut LaneProgress,
723 response: BidiWriteObjectResponse,
724 ) {
725 match response.write_status {
726 Some(bidi_write_object_response::WriteStatus::PersistedSize(size)) => {
727 progress.durable = progress.durable.max(size);
728 }
729 Some(bidi_write_object_response::WriteStatus::Resource(resource)) => {
730 let persisted_size = resource.size;
731 let snapshot = self.stat_from_object(resource);
732 progress.durable = progress.durable.max(persisted_size);
733 if snapshot.finalized {
734 progress.finalized = Some(snapshot);
735 }
736 }
737 None => {}
738 }
739 }
740
741 fn session_open_request(
743 &self,
744 generation: i64,
745 metageneration: i64,
746 write_handle: Option<Bytes>,
747 write_offset: i64,
748 ) -> BidiWriteObjectRequest {
749 BidiWriteObjectRequest {
750 first_message: Some(bidi_write_object_request::FirstMessage::AppendObjectSpec(
751 AppendObjectSpec {
752 bucket: self.bucket.clone(),
753 object: self.object.clone(),
754 generation,
755 if_metageneration_match: Some(metageneration),
756 write_handle: write_handle.map(|handle| BidiWriteHandle { handle }),
757 ..Default::default()
758 },
759 )),
760 write_offset,
761 flush: true,
762 state_lookup: true,
763 ..Default::default()
764 }
765 }
766
767 #[cfg(feature = "probe-support")]
771 fn current_session_open_request(&self) -> BidiWriteObjectRequest {
772 BidiWriteObjectRequest {
773 first_message: Some(bidi_write_object_request::FirstMessage::AppendObjectSpec(
774 AppendObjectSpec {
775 bucket: self.bucket.clone(),
776 object: self.object.clone(),
777 generation: 0,
778 ..Default::default()
779 },
780 )),
781 write_offset: 0,
782 flush: true,
783 state_lookup: true,
784 ..Default::default()
785 }
786 }
787
788 #[cfg(feature = "probe-support")]
789 async fn takeover_current_generation(
790 &self,
791 ) -> Result<(AppendToken, Option<i64>), TransportError> {
792 let (session, persisted_size, write_handle, opening_resource_size) = self
793 .open_session_observed(self.current_session_open_request())
794 .await?;
795 self.replace_session(Some(session)).await;
796 tracing::debug!(
797 zone = self.zone,
798 persisted_size,
799 has_write_handle = write_handle.is_some(),
800 "current append session takeover completed"
801 );
802 Ok((
803 AppendToken {
804 zone: self.zone,
805 generation: None,
808 metageneration: None,
809 persisted_size,
810 write_handle,
811 },
812 opening_resource_size,
813 ))
814 }
815
816 async fn resolve_token_identity(
817 &self,
818 token: &mut AppendToken,
819 ) -> Result<(i64, i64), TransportError> {
820 match (token.generation, token.metageneration) {
821 (Some(generation), Some(metageneration)) => Ok((generation, metageneration)),
822 (None, None) => {
823 let observed = self.stat().await?;
824 if observed.finalized {
825 return Err(self.error(
826 TransportCode::FailedPrecondition,
827 "append session object is already finalized",
828 ));
829 }
830 token.generation = Some(observed.generation);
831 token.metageneration = Some(observed.metageneration);
832 Ok((observed.generation, observed.metageneration))
833 }
834 _ => Err(self.error(
835 TransportCode::Internal,
836 "append token has incomplete generation identity",
837 )),
838 }
839 }
840
841 async fn finish_live_session(
842 &self,
843 write_offset: i64,
844 expected_generation: Option<i64>,
845 ) -> Result<ReplicaSnapshot, TransportError> {
846 let handle = self.live_session().map_err(|_| {
847 self.error(
848 TransportCode::Unavailable,
849 "no live append session to finalize",
850 )
851 })?;
852 let request = BidiWriteObjectRequest {
853 first_message: None,
854 write_offset,
855 finish_write: true,
856 ..Default::default()
857 };
858 handle
859 .send(
860 self,
861 request,
862 "append session disconnected while finalizing",
863 )
864 .await?;
865 handle
866 .wait_for(
867 self,
868 SessionWait {
869 timeout: Some(SESSION_PROGRESS_TIMEOUT),
870 clear_on_ready: true,
871 inspect_before_error: true,
872 reader_ended: "append session reader ended while finalizing",
873 stalled: "append session finalization made no progress",
874 },
875 |progress| {
876 let finalized = progress.finalized.clone()?;
877 Some(
878 if !finalized.finalized
879 || expected_generation
880 .is_some_and(|generation| finalized.generation != generation)
881 || finalized.persisted_size != write_offset
882 {
883 tracing::warn!(
884 expected_generation,
885 actual_generation = finalized.generation,
886 expected_size = write_offset,
887 actual_size = finalized.persisted_size,
888 finalized = finalized.finalized,
889 "finalized append response did not match the requested prefix"
890 );
891 Err(self.error(
892 TransportCode::DataLoss,
893 "finalized segment does not match the committed prefix",
894 ))
895 } else {
896 Ok(finalized)
897 },
898 )
899 },
900 )
901 .await
902 }
903
904 async fn drive_redirect_aware_stream<S>(
905 &self,
906 requests: Vec<BidiWriteObjectRequest>,
907 make_state: impl Fn() -> S,
908 mut step: impl FnMut(&mut S, BidiWriteObjectResponse),
909 complete: impl Fn(&S) -> bool,
910 ) -> (S, Option<TransportError>) {
911 let mut state = make_state();
912 let opened = match self.open_redirect_aware_stream(&requests, false).await {
913 Ok(Some(opened)) => opened,
914 Ok(None) => return (state, None),
915 Err(error) => return (state, Some(error)),
916 };
917 step(&mut state, opened.first);
918 if complete(&state) {
919 return (state, None);
920 }
921 let mut stream = opened.responses;
922 loop {
923 match stream.message().await {
924 Ok(Some(response)) => {
925 step(&mut state, response);
926 if complete(&state) {
927 return (state, None);
931 }
932 }
933 Ok(None) => return (state, None),
934 Err(status) => return (state, Some(self.status(status))),
935 }
936 }
937 }
938}
939
940impl GrpcReplicaFactory {
941 #[cfg(feature = "probe-support")]
942 fn probe_replica(&self, object: &str) -> GrpcReplica {
943 GrpcReplica {
944 zone: self.zone,
945 bucket: self.bucket.clone(),
946 object: object.to_string(),
947 auth: self.auth.clone(),
948 client: self.client.clone(),
949 routing_token: self.routing_token.clone(),
950 session: Arc::new(SessionSlot::new()),
951 }
952 }
953
954 pub async fn connect(
962 zone: usize,
963 endpoint: &str,
964 bucket: impl Into<String>,
965 bearer_token: Option<String>,
966 ) -> Result<Self, Error> {
967 let mut builder = Endpoint::from_shared(endpoint.to_string())
968 .map_err(|error| Error::Connection(error.to_string()))?;
969 builder = builder.connect_timeout(RPC_TIMEOUT);
970 if endpoint.starts_with("https://") {
971 builder = builder
972 .tls_config(ClientTlsConfig::new().with_webpki_roots())
973 .map_err(|error| Error::Connection(error.to_string()))?;
974 }
975 let channel = builder
976 .connect()
977 .await
978 .map_err(|error| Error::Connection(error.to_string()))?;
979 Ok(Self {
980 zone,
981 bucket: bucket.into(),
982 auth: bearer_token.map(BearerAuth::static_token),
983 client: StorageClient::new(channel),
984 routing_token: Arc::new(ArcSwap::from_pointee(None)),
985 })
986 }
987
988 pub fn from_channel(
993 zone: usize,
994 channel: Channel,
995 bucket: impl Into<String>,
996 bearer_token: Option<String>,
997 ) -> Self {
998 Self {
999 zone,
1000 bucket: bucket.into(),
1001 auth: bearer_token.map(BearerAuth::static_token),
1002 client: StorageClient::new(channel),
1003 routing_token: Arc::new(ArcSwap::from_pointee(None)),
1004 }
1005 }
1006
1007 pub async fn connect_with_auth(
1012 zone: usize,
1013 endpoint: &str,
1014 bucket: impl Into<String>,
1015 auth: BearerAuth,
1016 ) -> Result<Self, Error> {
1017 let mut builder = Endpoint::from_shared(endpoint.to_string())
1018 .map_err(|error| Error::Connection(error.to_string()))?;
1019 builder = builder.connect_timeout(RPC_TIMEOUT);
1020 if endpoint.starts_with("https://") {
1021 builder = builder
1022 .tls_config(ClientTlsConfig::new().with_webpki_roots())
1023 .map_err(|error| Error::Connection(error.to_string()))?;
1024 }
1025 let channel = builder
1026 .connect()
1027 .await
1028 .map_err(|error| Error::Connection(error.to_string()))?;
1029 Ok(Self {
1030 zone,
1031 bucket: bucket.into(),
1032 auth: Some(auth),
1033 client: StorageClient::new(channel),
1034 routing_token: Arc::new(ArcSwap::from_pointee(None)),
1035 })
1036 }
1037}
1038
1039#[cfg(feature = "probe-support")]
1040#[derive(Clone, Debug)]
1041pub struct GenerationZeroOpenObservation {
1043 pub persisted_size: i64,
1045 pub resource_size: Option<i64>,
1047}
1048
1049#[cfg(feature = "probe-support")]
1050#[derive(Clone, Debug)]
1051pub struct GenerationZeroTakeoverProbeResult {
1053 pub expected_size: i64,
1055 pub append: Result<i64, TransportError>,
1057 pub takeover: Option<Result<GenerationZeroOpenObservation, TransportError>>,
1059 pub absent: Result<GenerationZeroOpenObservation, TransportError>,
1061 pub cleanup: Result<(), TransportError>,
1063}
1064
1065#[cfg(feature = "probe-support")]
1066pub async fn probe_generation_zero_takeover(
1074 factory: &GrpcReplicaFactory,
1075 present_object: &str,
1076 absent_object: &str,
1077 payload: Bytes,
1078) -> Result<GenerationZeroTakeoverProbeResult, TransportError> {
1079 let expected_size = i64::try_from(payload.len()).map_err(|_| TransportError {
1080 zone: factory.zone,
1081 code: TransportCode::InvalidArgument,
1082 message: "probe payload length does not fit in i64".into(),
1083 })?;
1084 if expected_size == 0 {
1085 return Err(TransportError {
1086 zone: factory.zone,
1087 code: TransportCode::InvalidArgument,
1088 message: "probe payload must be non-empty".into(),
1089 });
1090 }
1091
1092 let present = factory.probe_replica(present_object);
1093 let append = async {
1094 let token = present.create_append_session(HashMap::new()).await?;
1095 if token.persisted_size != 0 {
1096 return Err(present.error(
1097 TransportCode::DataLoss,
1098 format!(
1099 "new appendable object opened at persisted_size={}, expected 0",
1100 token.persisted_size
1101 ),
1102 ));
1103 }
1104 present
1105 .lane_send(token.persisted_size, std::slice::from_ref(&payload))
1106 .await?;
1107 let change = present.lane_durable_change(token.persisted_size).await?;
1108 if let Some(error) = change.error {
1109 return Err(error);
1110 }
1111 Ok(change.persisted_size)
1112 }
1113 .await;
1114 let takeover = if append.is_ok() {
1115 Some(
1116 present
1117 .takeover_current_generation()
1118 .await
1119 .map(|(token, resource_size)| GenerationZeroOpenObservation {
1120 persisted_size: token.persisted_size,
1121 resource_size,
1122 }),
1123 )
1124 } else {
1125 None
1126 };
1127
1128 present.shutdown().await;
1129 let cleanup = match present.stat().await {
1130 Ok(snapshot) => present.delete(snapshot.generation).await,
1131 Err(error) if error.code == TransportCode::NotFound => Ok(()),
1132 Err(error) => Err(error),
1133 };
1134
1135 let absent_replica = factory.probe_replica(absent_object);
1136 let absent =
1137 absent_replica
1138 .takeover_current_generation()
1139 .await
1140 .map(|(token, resource_size)| GenerationZeroOpenObservation {
1141 persisted_size: token.persisted_size,
1142 resource_size,
1143 });
1144 absent_replica.shutdown().await;
1145
1146 Ok(GenerationZeroTakeoverProbeResult {
1147 expected_size,
1148 append,
1149 takeover,
1150 absent,
1151 cleanup,
1152 })
1153}
1154
1155#[async_trait]
1156impl ReplicaFactory for GrpcReplicaFactory {
1157 fn bucket_name(&self) -> &str {
1158 self.bucket.rsplit('/').next().unwrap_or_default()
1159 }
1160
1161 fn replica(&self, object: &str) -> Arc<dyn Replica> {
1162 Arc::new(GrpcReplica {
1163 zone: self.zone,
1164 bucket: self.bucket.clone(),
1165 object: object.to_string(),
1166 auth: self.auth.clone(),
1167 client: self.client.clone(),
1168 routing_token: self.routing_token.clone(),
1169 session: Arc::new(SessionSlot::new()),
1170 })
1171 }
1172
1173 async fn list(&self, prefix: &str) -> Result<Vec<ListedObject>, TransportError> {
1174 let replica = GrpcReplica {
1175 zone: self.zone,
1176 bucket: self.bucket.clone(),
1177 object: String::new(),
1178 auth: self.auth.clone(),
1179 client: self.client.clone(),
1180 routing_token: self.routing_token.clone(),
1181 session: Arc::new(SessionSlot::new()),
1182 };
1183 let mut page_token = String::new();
1184 let mut seen_page_tokens = HashSet::new();
1185 let mut listed = Vec::new();
1186 for _ in 0..MAX_LIST_PAGES {
1187 let request = replica.request(ListObjectsRequest {
1188 parent: self.bucket.clone(),
1189 page_size: 1000,
1190 page_token,
1191 prefix: prefix.to_string(),
1192 ..Default::default()
1193 })?;
1194 let response = self
1195 .client
1196 .clone()
1197 .list_objects(request)
1198 .await
1199 .map_err(|status| replica.status(status))?
1200 .into_inner();
1201 listed.extend(response.objects.into_iter().map(|object| ListedObject {
1202 zone: self.zone,
1203 name: object.name,
1204 generation: object.generation,
1205 finalized: object.finalize_time.is_some(),
1206 metadata: object.metadata,
1207 }));
1208 if response.next_page_token.is_empty() {
1209 return Ok(listed);
1210 }
1211 if !seen_page_tokens.insert(response.next_page_token.clone()) {
1212 return Err(
1213 replica.error(TransportCode::Internal, "ListObjects repeated a page token")
1214 );
1215 }
1216 page_token = response.next_page_token;
1217 }
1218 Err(replica.error(
1219 TransportCode::Internal,
1220 "ListObjects exceeded the pagination bound",
1221 ))
1222 }
1223}
1224
1225#[async_trait]
1226impl Replica for GrpcReplica {
1227 async fn stat(&self) -> Result<ReplicaSnapshot, TransportError> {
1228 let get = GetObjectRequest {
1229 bucket: self.bucket.clone(),
1230 object: self.object.clone(),
1231 ..Default::default()
1232 };
1233 let request = self.request(get)?;
1234 let object = self
1235 .client
1236 .clone()
1237 .get_object(request)
1238 .await
1239 .map_err(|status| self.status(status))?
1240 .into_inner();
1241 Ok(self.stat_from_object(object))
1244 }
1245
1246 async fn snapshot(&self) -> Result<ReplicaSnapshot, TransportError> {
1247 let get = GetObjectRequest {
1248 bucket: self.bucket.clone(),
1249 object: self.object.clone(),
1250 ..Default::default()
1251 };
1252 let request = self.request(get)?;
1253 let object = self
1254 .client
1255 .clone()
1256 .get_object(request)
1257 .await
1258 .map_err(|status| self.status(status))?
1259 .into_inner();
1260
1261 let bytes = self.bidi_read(object.generation).await?;
1265 Ok(self.snapshot_from_object(object, bytes))
1266 }
1267
1268 async fn create_appendable(
1269 &self,
1270 metadata: HashMap<String, String>,
1271 ) -> Result<ReplicaSnapshot, TransportError> {
1272 let object = Object {
1273 bucket: self.bucket.clone(),
1274 name: self.object.clone(),
1275 metadata: metadata.clone(),
1276 content_type: "application/vnd.chorus.records".into(),
1277 ..Default::default()
1278 };
1279 let request = BidiWriteObjectRequest {
1280 first_message: Some(bidi_write_object_request::FirstMessage::WriteObjectSpec(
1281 WriteObjectSpec {
1282 resource: Some(object),
1283 if_generation_match: Some(0),
1284 appendable: Some(true),
1285 ..Default::default()
1286 },
1287 )),
1288 write_offset: 0,
1289 flush: true,
1290 state_lookup: true,
1291 ..Default::default()
1292 };
1293 let (_, error) = self
1294 .drive_redirect_aware_stream(
1295 vec![request],
1296 || false,
1297 |seen, _| *seen = true,
1298 |seen| *seen,
1299 )
1300 .await;
1301 if let Some(error) = error {
1302 if error.code == TransportCode::FailedPrecondition {
1306 return Err(TransportError {
1307 code: TransportCode::AlreadyExists,
1308 ..error
1309 });
1310 }
1311 return Err(error);
1312 }
1313 self.stat().await
1318 }
1319
1320 async fn create_append_session(
1321 &self,
1322 metadata: HashMap<String, String>,
1323 ) -> Result<AppendToken, TransportError> {
1324 let object = Object {
1325 bucket: self.bucket.clone(),
1326 name: self.object.clone(),
1327 metadata,
1328 content_type: "application/vnd.chorus.records".into(),
1329 ..Default::default()
1330 };
1331 let request = BidiWriteObjectRequest {
1332 first_message: Some(bidi_write_object_request::FirstMessage::WriteObjectSpec(
1333 WriteObjectSpec {
1334 resource: Some(object),
1335 if_generation_match: Some(0),
1336 appendable: Some(true),
1337 ..Default::default()
1338 },
1339 )),
1340 write_offset: 0,
1341 flush: true,
1342 state_lookup: true,
1343 ..Default::default()
1344 };
1345 let (session, persisted_size, write_handle) =
1346 self.open_session(request).await.map_err(|error| {
1347 if error.code == TransportCode::FailedPrecondition {
1348 TransportError {
1349 code: TransportCode::AlreadyExists,
1350 ..error
1351 }
1352 } else {
1353 error
1354 }
1355 })?;
1356 self.replace_session(Some(session)).await;
1357 Ok(AppendToken {
1358 zone: self.zone,
1359 generation: None,
1360 metageneration: None,
1361 persisted_size,
1362 write_handle,
1363 })
1364 }
1365
1366 async fn create_register(
1367 &self,
1368 metadata: HashMap<String, String>,
1369 ) -> Result<ReplicaSnapshot, TransportError> {
1370 let object = Object {
1371 bucket: self.bucket.clone(),
1372 name: self.object.clone(),
1373 metadata,
1374 content_type: "application/vnd.chorus.manifest".into(),
1375 ..Default::default()
1376 };
1377 let request = WriteObjectRequest {
1380 first_message: Some(write_object_request::FirstMessage::WriteObjectSpec(
1381 WriteObjectSpec {
1382 resource: Some(object),
1383 if_generation_match: Some(0),
1384 ..Default::default()
1385 },
1386 )),
1387 write_offset: 0,
1388 finish_write: true,
1389 ..Default::default()
1390 };
1391 let request = self.request(tokio_stream::iter([request]))?;
1392 let response = self
1393 .client
1394 .clone()
1395 .write_object(request)
1396 .await
1397 .map_err(|status| {
1398 let error = self.status(status);
1399 if error.code == TransportCode::FailedPrecondition {
1402 TransportError {
1403 code: TransportCode::AlreadyExists,
1404 ..error
1405 }
1406 } else {
1407 error
1408 }
1409 })?
1410 .into_inner();
1411 let Some(write_object_response::WriteStatus::Resource(object)) = response.write_status
1412 else {
1413 return Err(self.error(
1414 TransportCode::Internal,
1415 "manifest create response omitted the object resource",
1416 ));
1417 };
1418 Ok(self.stat_from_object(object))
1419 }
1420
1421 async fn resume_tail(&self, token: &mut AppendToken) -> Result<i64, TransportError> {
1422 let (generation, metageneration) = self.resolve_token_identity(token).await?;
1423 let mut owned = self.session.owned.lock().await;
1424 let first = self.session_open_request(
1425 generation,
1426 metageneration,
1427 token.write_handle.clone(),
1428 token.persisted_size,
1429 );
1430 let (result, previous) = match self.open_session(first).await {
1431 Ok((session, persisted, _)) => {
1432 let previous = self.replace_session_locked(&mut owned, Some(session));
1433 tracing::debug!(
1434 zone = self.zone,
1435 persisted_size = persisted,
1436 "append session resumed"
1437 );
1438 (Ok(persisted), previous)
1439 }
1440 Err(error) => {
1441 let previous = self.replace_session_locked(&mut owned, None);
1442 (Err(error), previous)
1443 }
1444 };
1445 drop(owned);
1446 if let Some(previous) = previous {
1447 previous.shutdown().await;
1448 }
1449 result
1450 }
1451
1452 async fn takeover(&self, observed: &ReplicaSnapshot) -> Result<AppendToken, TransportError> {
1453 let first = self.session_open_request(
1454 observed.generation,
1455 observed.metageneration,
1456 None, observed.persisted_size,
1458 );
1459 let (session, persisted_size, write_handle) = self.open_session(first).await?;
1460 self.replace_session(Some(session)).await;
1461 tracing::debug!(
1462 zone = self.zone,
1463 persisted_size,
1464 has_write_handle = write_handle.is_some(),
1465 "append session takeover completed"
1466 );
1467 Ok(AppendToken {
1468 zone: self.zone,
1469 generation: Some(observed.generation),
1470 metageneration: Some(observed.metageneration),
1471 persisted_size,
1472 write_handle,
1473 })
1474 }
1475
1476 async fn update_register(
1477 &self,
1478 metageneration: i64,
1479 metadata: HashMap<String, String>,
1480 ) -> Result<ReplicaSnapshot, TransportError> {
1481 let request = UpdateObjectRequest {
1485 object: Some(Object {
1486 bucket: self.bucket.clone(),
1487 name: self.object.clone(),
1488 metadata,
1489 ..Default::default()
1490 }),
1491 if_metageneration_match: Some(metageneration),
1492 update_mask: Some(prost_types::FieldMask {
1493 paths: vec!["metadata".into()],
1494 }),
1495 ..Default::default()
1496 };
1497 let object = self
1498 .client
1499 .clone()
1500 .update_object(self.request(request)?)
1501 .await
1502 .map_err(|status| self.status(status))?
1503 .into_inner();
1504 Ok(self.stat_from_object(object))
1505 }
1506
1507 async fn replace_appendable(
1508 &self,
1509 observed: &ReplicaSnapshot,
1510 data: Bytes,
1511 metadata: HashMap<String, String>,
1512 ) -> Result<AppendToken, TransportError> {
1513 let object = Object {
1514 bucket: self.bucket.clone(),
1515 name: self.object.clone(),
1516 metadata: metadata.clone(),
1517 content_type: "application/vnd.chorus.records".into(),
1518 ..Default::default()
1519 };
1520 let request = BidiWriteObjectRequest {
1521 first_message: Some(bidi_write_object_request::FirstMessage::WriteObjectSpec(
1522 WriteObjectSpec {
1523 resource: Some(object),
1524 if_generation_match: Some(observed.generation),
1525 if_metageneration_match: Some(observed.metageneration),
1526 appendable: Some(true),
1527 ..Default::default()
1528 },
1529 )),
1530 write_offset: 0,
1531 data: Some(bidi_write_object_request::Data::ChecksummedData(
1532 ChecksummedData {
1533 crc32c: Some(crc32c::crc32c(&data)),
1534 content: data.clone(),
1535 },
1536 )),
1537 flush: true,
1538 state_lookup: true,
1539 ..Default::default()
1540 };
1541 let expected = data.len() as i64;
1542 let (persisted_size, error) = self
1543 .drive_redirect_aware_stream(
1544 vec![request],
1545 || None,
1546 |persisted_size, response| {
1547 if let Some(bidi_write_object_response::WriteStatus::PersistedSize(size)) =
1548 response.write_status
1549 {
1550 *persisted_size = Some(size);
1551 }
1552 },
1553 |persisted_size| persisted_size.is_some_and(|size| size >= expected),
1554 )
1555 .await;
1556 if let Some(error) = error {
1557 return Err(error);
1558 }
1559 if persisted_size != Some(expected) {
1560 return Err(self.error(
1561 TransportCode::DataLoss,
1562 format!("replacement persisted {persisted_size:?}, expected {expected}"),
1563 ));
1564 }
1565 let snapshot = self.snapshot().await?;
1566 if snapshot.bytes != data[..] || snapshot.metadata != metadata {
1567 return Err(self.error(
1568 TransportCode::DataLoss,
1569 "replacement generation failed verification",
1570 ));
1571 }
1572 Ok(AppendToken {
1573 zone: self.zone,
1574 generation: Some(snapshot.generation),
1575 metageneration: Some(snapshot.metageneration),
1576 persisted_size: expected,
1577 write_handle: None,
1578 })
1579 }
1580
1581 async fn append(
1582 &self,
1583 token: &AppendToken,
1584 write_offset: i64,
1585 data: Vec<u8>,
1586 ) -> Result<i64, TransportError> {
1587 let Some(generation) = token.generation else {
1588 return Err(self.error(
1589 TransportCode::Internal,
1590 "one-shot append requires a generation-bound token",
1591 ));
1592 };
1593 let request = BidiWriteObjectRequest {
1594 first_message: Some(bidi_write_object_request::FirstMessage::AppendObjectSpec(
1595 AppendObjectSpec {
1596 bucket: self.bucket.clone(),
1597 object: self.object.clone(),
1598 generation,
1599 if_metageneration_match: token.metageneration,
1600 write_handle: token
1601 .write_handle
1602 .clone()
1603 .map(|handle| BidiWriteHandle { handle }),
1604 ..Default::default()
1605 },
1606 )),
1607 write_offset,
1608 data: Some(bidi_write_object_request::Data::ChecksummedData(
1609 ChecksummedData {
1610 content: Bytes::from(data.clone()),
1611 crc32c: Some(crc32c::crc32c(&data)),
1612 },
1613 )),
1614 flush: true,
1615 state_lookup: true,
1616 ..Default::default()
1617 };
1618 let expected = write_offset + data.len() as i64;
1619 let (persisted_size, error) = self
1620 .drive_redirect_aware_stream(
1621 vec![request],
1622 || None,
1623 |persisted_size, response| {
1624 if let Some(bidi_write_object_response::WriteStatus::PersistedSize(size)) =
1625 response.write_status
1626 {
1627 *persisted_size = Some(size);
1628 }
1629 },
1630 |persisted_size| persisted_size.is_some_and(|size| size >= expected),
1631 )
1632 .await;
1633 if let Some(error) = error {
1634 return Err(error);
1635 }
1636 match persisted_size {
1637 Some(size) if size >= expected => Ok(size),
1638 Some(size) => Err(self.error(
1639 TransportCode::DataLoss,
1640 format!("flush persisted {size}, expected at least {expected}"),
1641 )),
1642 None => Err(self.error(TransportCode::Internal, "missing persisted-size response")),
1643 }
1644 }
1645
1646 async fn lane_send(&self, write_offset: i64, chunks: &[Bytes]) -> Result<(), TransportError> {
1647 let packed = pack_append(chunks.to_vec());
1648 self.lane_send_packed(write_offset, &packed).await
1649 }
1650
1651 async fn lane_send_packed(
1652 &self,
1653 write_offset: i64,
1654 packed: &PackedAppend,
1655 ) -> Result<(), TransportError> {
1656 if packed.is_empty() {
1657 return Err(self.error(
1658 TransportCode::Internal,
1659 "append lane cannot send an empty flush group",
1660 ));
1661 }
1662 let handle = self.live_session()?;
1663 let messages = packed.messages();
1667 let last_index = messages.len() - 1;
1668 for (index, message) in messages.iter().enumerate() {
1669 let last = index == last_index;
1670 let request = BidiWriteObjectRequest {
1671 first_message: None,
1672 write_offset: write_offset + message.relative_offset,
1673 data: Some(bidi_write_object_request::Data::ChecksummedData(
1674 ChecksummedData {
1675 crc32c: Some(message.crc32c),
1676 content: message.content.clone(),
1677 },
1678 )),
1679 flush: last,
1680 state_lookup: last,
1681 ..Default::default()
1682 };
1683 handle
1684 .send(self, request, "append session disconnected while sending")
1685 .await?;
1686 }
1687 Ok(())
1688 }
1689
1690 async fn lane_durable_change(&self, seen: i64) -> Result<LaneDurableChange, TransportError> {
1691 let handle = self.live_session()?;
1692 let change = handle
1693 .wait_for(
1694 self,
1695 SessionWait {
1696 timeout: None,
1697 clear_on_ready: false,
1698 inspect_before_error: true,
1702 reader_ended: "append session reader ended",
1703 stalled: "append session made no progress",
1704 },
1705 |progress| {
1706 (progress.durable > seen).then(|| {
1707 Ok(LaneDurableChange {
1708 persisted_size: progress.durable,
1709 error: progress.error.clone(),
1710 })
1711 })
1712 },
1713 )
1714 .await?;
1715 if change.error.is_some() {
1716 self.clear_session_if(&handle).await;
1717 }
1718 Ok(change)
1719 }
1720 async fn delete(&self, generation: i64) -> Result<(), TransportError> {
1721 let request = DeleteObjectRequest {
1722 bucket: self.bucket.clone(),
1723 object: self.object.clone(),
1724 generation,
1725 if_generation_match: Some(generation),
1726 ..Default::default()
1727 };
1728 let request = self.request(request)?;
1729 self.client
1730 .clone()
1731 .delete_object(request)
1732 .await
1733 .map_err(|status| self.status(status))?;
1734 Ok(())
1735 }
1736
1737 async fn finalize(
1738 &self,
1739 token: &mut AppendToken,
1740 write_offset: i64,
1741 ) -> Result<ReplicaSnapshot, TransportError> {
1742 if self.session.current.load().is_some() {
1747 let finalized = self
1748 .finish_live_session(write_offset, token.generation)
1749 .await?;
1750 token.generation = Some(finalized.generation);
1751 token.metageneration = Some(finalized.metageneration);
1752 return Ok(finalized);
1753 }
1754
1755 let (generation, _) = self.resolve_token_identity(token).await?;
1760 let request = BidiWriteObjectRequest {
1761 first_message: Some(bidi_write_object_request::FirstMessage::AppendObjectSpec(
1762 AppendObjectSpec {
1763 bucket: self.bucket.clone(),
1764 object: self.object.clone(),
1765 generation,
1766 if_metageneration_match: None,
1767 write_handle: token
1768 .write_handle
1769 .clone()
1770 .map(|handle| BidiWriteHandle { handle }),
1771 ..Default::default()
1772 },
1773 )),
1774 write_offset,
1775 finish_write: true,
1776 ..Default::default()
1777 };
1778 let (resource, error) = self
1779 .drive_redirect_aware_stream(
1780 vec![request],
1781 || None,
1782 |resource, response| {
1783 if let Some(bidi_write_object_response::WriteStatus::Resource(object)) =
1784 response.write_status
1785 {
1786 *resource = Some(object);
1787 }
1788 },
1789 Option::is_some,
1790 )
1791 .await;
1792 if let Some(error) = error {
1793 return Err(error);
1794 }
1795 let Some(resource) = resource else {
1796 return Err(self.error(
1797 TransportCode::Internal,
1798 "missing finalized resource response",
1799 ));
1800 };
1801 let finalized = self.stat_from_object(resource);
1802 if !finalized.finalized
1803 || finalized.generation != generation
1804 || finalized.persisted_size != write_offset
1805 {
1806 return Err(self.error(
1807 TransportCode::DataLoss,
1808 "finalized segment does not match the recovered prefix",
1809 ));
1810 }
1811 Ok(finalized)
1812 }
1813
1814 async fn shutdown(&self) {
1815 self.session.shutdown.send_replace(true);
1816 self.replace_session(None).await;
1817 }
1818}
1819
1820const MAX_WRITE_REDIRECTS: u32 = 5;
1822
1823#[derive(Clone, PartialEq, ::prost::Message)]
1826struct RichStatus {
1827 #[prost(int32, tag = "1")]
1828 code: i32,
1829 #[prost(string, tag = "2")]
1830 message: ::prost::alloc::string::String,
1831 #[prost(message, repeated, tag = "3")]
1832 details: ::prost::alloc::vec::Vec<::prost_types::Any>,
1833}
1834
1835pub fn redirect_routing_token(status: &Status) -> Option<String> {
1843 let details = status.details();
1844 if details.is_empty() {
1845 return None;
1846 }
1847 let rich = RichStatus::decode(details).ok()?;
1848 for any in rich.details {
1849 if any
1850 .type_url
1851 .ends_with("google.storage.v2.BidiWriteObjectRedirectedError")
1852 {
1853 if let Ok(redirect) = BidiWriteObjectRedirectedError::decode(any.value.as_slice()) {
1854 if let Some(token) = redirect.routing_token {
1855 return Some(token);
1856 }
1857 }
1858 }
1859 }
1860 None
1861}