Skip to main content

chorus_client/
grpc.rs

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
29/// Cached zonal write routing token, learned from a `BidiWriteObjectRedirectedError`
30/// and replayed in `x-goog-request-params` to land on the bucket's location.
31/// Shared across every replica produced by one factory.
32type 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    /// The lane's persistent append session. The live service rate limits
43    /// per-object *mutations*, and every fresh append open is one — so a
44    /// lane opens its session once (at takeover) and streams every window
45    /// through it as continuations. Flushes are not rate limited.
46    /// Clones share the session; the hot send and
47    /// progress paths load its immutable handle without taking the replacement
48    /// mutex.
49    session: Arc<SessionSlot>,
50}
51
52#[derive(Clone)]
53/// Reusable Google Storage v2 channel, authentication handle, and zonal routing
54/// token for one Rapid bucket.
55pub struct GrpcReplicaFactory {
56    zone: usize,
57    bucket: String,
58    auth: Option<BearerAuth>,
59    client: StorageClient<Channel>,
60    routing_token: RoutingToken,
61}
62
63/// A live appendable write stream: an ordered request sender plus a
64/// background reader translating flush acknowledgments into a watchable
65/// durable tail. Send-ahead lanes write through `tx` while waiting on
66/// `state` — sends never block on acknowledgments.
67struct 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)]
123/// Durable tail plus any coalesced stream-termination error published by the
124/// session reader.
125struct LaneProgress {
126    durable: i64,
127    finalized: Option<ReplicaSnapshot>,
128    error: Option<TransportError>,
129}
130
131/// Target wire-message size for packed appends. Large enough to amortize
132/// per-message overhead (protobuf, CRC field, HTTP/2 framing, server
133/// per-message processing). The hard ceiling, measured live, is
134/// the service's 4 MiB gRPC inbound message cap — a 2 MiB + 1 byte chunk is
135/// accepted, 4 MiB is rejected with ResourceExhausted; the proto's only
136/// documented chunk constant (MaxReadChunkBytes = 2 MiB) is read-side.
137const 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
169/// Per-response progress timeout for the persistent session. The session
170/// itself has no overall deadline (it lives for the segment), but a server
171/// that stops acknowledging flushed appends must fail the lane.
172const 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    /// Read one generation in full over `BidiReadObject`.
248    ///
249    /// Preferred over `ReadObject`, which stalls a fixed ~5.1s opening its
250    /// stream against these appendable objects — measured across two zones and
251    /// three objects from 26KB to 198KB, while `GetObject` on the same objects
252    /// took 35ms and the drain under 2ms. Bidi returned the identical byte
253    /// count in 53ms. Same bytes, same generation guard, same per-chunk CRC
254    /// check; only the API differs.
255    ///
256    /// Errors carry the provider code, so the surrounding retry and quorum
257    /// logic classifies a bidi failure exactly as it classified a `ReadObject`
258    /// one — transient codes retry, `NotFound` and `DataLoss` keep their
259    /// meaning to recovery.
260    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            // Offset and size zero requests the whole object.
271            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    /// Like [`Self::request`], but without the per-RPC deadline: used for the
359    /// persistent append session, whose stream must outlive any fixed
360    /// deadline. Progress is bounded per response instead.
361    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        // GCS v2 gRPC routes every object RPC by the bucket resource path
371        // carried in x-goog-request-params; without it the service rejects the
372        // call with INVALID_ARGUMENT. Once a zonal write redirect has supplied a
373        // routing token, replay it here so the stream lands on the right location.
374        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        // Retry classification is a protocol contract, not a direct mirror of
407        // tonic's codes. This mapper is shared by append and non-append RPCs:
408        // RESOURCE_EXHAUSTED is GCS throttling, not a permanent rejection:
409        // per-object mutation-rate and per-project quotas reset over time, so
410        // retrying the valid request with backoff can succeed.
411        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    /// If `status` is a zonal write redirect, cache its routing token for replay
457    /// and report `true` so the caller retries the same guarded open. The
458    /// redirect fields are intentionally ignored: conditional creates replay
459    /// their create precondition, takeovers replay their metageneration guard,
460    /// and resumes replay their existing handle.
461    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    /// Open one write stream and obtain its first response, transparently
472    /// replaying redirects that occur before any response is observed.
473    ///
474    /// Both one-shot writes and persistent append sessions use this driver, so
475    /// routing-token capture, redirect limits, opening deadlines, and the
476    /// no-midstream-replay rule have one implementation.
477    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                // The established session has no RPC deadline, so bound its
513                // opening handshake separately before returning the live stream.
514                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    /// Open (or resume) the persistent append session: send the first
615    /// message, wait for the opening acknowledgment, and return the live
616    /// session with the server-certified durable tail and session handle.
617    /// Redirects are captured and retried with the routing token.
618    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    /// The opening message for this object's append session.
742    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    /// Construct the candidate handle-free current-generation open. Generation
768    /// zero is wire-identical to an unset proto3 scalar; the live probe decides
769    /// whether GCS accepts it as a selector or rejects the missing generation.
770    #[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                // The append-open response does not provide identity fields
806                // needed by later resume/finalize paths. Resolve them lazily.
807                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                        // Appendable sessions stay open server-side; once the
928                        // caller has what it needs, abandon the stream instead
929                        // of waiting for the live service's idle expiry.
930                        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    /// Connect a zonal bucket with an optional static bearer token.
955    ///
956    /// The bucket must be a full v2 resource name such as
957    /// `projects/_/buckets/example-zone-a`. All replicas created by the factory
958    /// share the routing token learned from zonal write redirects. The normal
959    /// production endpoint is `https://storage.googleapis.com`; Cloud Storage
960    /// regional JSON/XML endpoints do not support this gRPC client.
961    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    /// Build a factory over an already-established channel. The deterministic
989    /// simulator uses this to route the production client over a simulated
990    /// transport (`Endpoint::connect_with_connector`); production callers use
991    /// `connect`/`connect_with_auth`.
992    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    /// Connect a zonal bucket using a shared static or refreshing auth handle.
1008    /// Cloned refreshing handles update existing clients through `ArcSwap`.
1009    /// Use full v2 bucket resource names and a gRPC-compatible endpoint such as
1010    /// `https://storage.googleapis.com`.
1011    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)]
1041/// One generation-zero append-open response observed by the live-GCS probe.
1042pub struct GenerationZeroOpenObservation {
1043    /// Authoritative tail derived from the opening response.
1044    pub persisted_size: i64,
1045    /// Object-resource size in the opening response, if GCS supplied a resource.
1046    pub resource_size: Option<i64>,
1047}
1048
1049#[cfg(feature = "probe-support")]
1050#[derive(Clone, Debug)]
1051/// Independent observations from the focused live-GCS generation-zero probe.
1052pub struct GenerationZeroTakeoverProbeResult {
1053    /// Number of bytes the probe requested GCS to persist.
1054    pub expected_size: i64,
1055    /// Result of creating an appendable object and flushing the payload.
1056    pub append: Result<i64, TransportError>,
1057    /// Present-object takeover result, absent only when creation or append failed.
1058    pub takeover: Option<Result<GenerationZeroOpenObservation, TransportError>>,
1059    /// Generation-zero open result for a never-created object.
1060    pub absent: Result<GenerationZeroOpenObservation, TransportError>,
1061    /// Cleanup result for the object that may have been created.
1062    pub cleanup: Result<(), TransportError>,
1063}
1064
1065#[cfg(feature = "probe-support")]
1066/// Exercise the production gRPC append path against one present and one absent
1067/// object, then delete the present object before returning.
1068///
1069/// Both takeovers send `AppendObjectSpec.generation = 0` through the same
1070/// redirect-aware `open_session` implementation used by recovery. Operations
1071/// are reported independently so a provider rejection on the present object
1072/// does not suppress the absent-name observation.
1073pub 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        // A finalized object's size is frozen and authoritative; open
1242        // objects keep persisted_size = 0 here (stats are tail-blind).
1243        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        // Content comes over BidiReadObject. ReadObject stalls a fixed ~5.1s
1262        // opening its stream against these appendable objects, independent of
1263        // zone, object and size, while bidi returned the same bytes in 53ms.
1264        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            // The only precondition this request carries is
1303            // `if_generation_match=0`; the live service reports the conflict
1304            // as FAILED_PRECONDITION where the protocol expects AlreadyExists.
1305            if error.code == TransportCode::FailedPrecondition {
1306                return Err(TransportError {
1307                    code: TransportCode::AlreadyExists,
1308                    ..error
1309                });
1310            }
1311            return Err(error);
1312        }
1313        // Appendable create responses report only persisted size, not the
1314        // generation and metageneration required to guard the takeover open.
1315        // A metadata-only read supplies those fields without reading the
1316        // empty object body.
1317        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        // A plain one-shot WriteObject: regional buckets reject appendable
1378        // creates, and the register is never appended to anyway.
1379        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                // As with segment creates: the only precondition here is
1400                // `if_generation_match=0`, reported as FAILED_PRECONDITION.
1401                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, // handle-free: this open IS the fence
1457            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        // Generation 0 addresses the live object; the register is never
1482        // deleted or recreated, so the metageneration precondition alone is
1483        // the CAS guard.
1484        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        // The coalesced group was packed once before replica dispatch. Each
1664        // lane builds only its protobuf envelopes and shallow-clones the
1665        // refcounted message bytes; CRC32C and byte concatenation are shared.
1666        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                    // Preserve a coalesced response and stream error in one
1699                    // observation. Protocol code publishes the durable offset
1700                    // before deciding whether to recover or fence the writer.
1701                    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        // A healthy lane finishes on the same stream that conditionally
1743        // created the object. This needs neither another RPC nor a generation
1744        // lookup; the returned resource supplies the identity used by the
1745        // seal-time metadata CAS.
1746        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        // If an idle or failed stream disappeared before finalization, stat
1756        // only to obtain a candidate identity, then require the original
1757        // write handle to resume that exact server-side session. A replacement
1758        // generation cannot be adopted through this path.
1759        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
1820/// Bounded number of consecutive zonal redirects honored per write open.
1821const MAX_WRITE_REDIRECTS: u32 = 5;
1822
1823/// Minimal view of `google.rpc.Status` for decoding the rich-error payload that
1824/// tonic exposes via `Status::details()` (the `grpc-status-details-bin` trailer).
1825#[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
1835/// Extract the routing token from a `BidiWriteObjectRedirectedError` attached to
1836/// an ABORTED status, if present. The feature-gated `probe_support` module
1837/// exports this for repository probes that must replay redirects exactly like
1838/// the production transport.
1839///
1840/// Normal users do not need this helper; the production transport consumes and
1841/// caches redirects automatically.
1842pub 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}