Skip to main content

risingwave_meta/manager/sink_coordination/
coordinator_worker.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
16use std::fmt::Debug;
17use std::future::{Future, poll_fn};
18use std::pin::pin;
19use std::sync::LazyLock;
20use std::task::Poll;
21use std::time::{Duration, Instant};
22
23use anyhow::anyhow;
24use await_tree::InstrumentAwait;
25use futures::future::{Either, pending, select};
26use futures::pin_mut;
27use itertools::Itertools;
28use risingwave_common::bail;
29use risingwave_common::bitmap::Bitmap;
30use risingwave_common::util::retry::exponential_backoff;
31use risingwave_connector::connector_common::IcebergSinkCompactionUpdate;
32use risingwave_connector::dispatch_sink;
33use risingwave_connector::sink::catalog::SinkId;
34use risingwave_connector::sink::{
35    Sink, SinkCommitCoordinator, SinkCommittedEpochSubscriber, SinkError, SinkParam, build_sink,
36};
37use risingwave_meta_model::pending_sink_state::SinkState;
38use risingwave_pb::connector_service::{SinkMetadata, coordinate_request};
39use risingwave_pb::stream_plan::PbSinkSchemaChange;
40use sea_orm::DatabaseConnection;
41use thiserror_ext::AsReport;
42use tokio::select;
43use tokio::sync::Semaphore;
44use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender};
45use tokio::time::sleep;
46use tokio_retry::strategy::jitter;
47use tonic::Status;
48use tracing::{error, warn};
49
50use crate::manager::exactly_once_util::{
51    clean_aborted_records, commit_and_prune_epoch, list_sink_states_ordered_by_epoch,
52    persist_pre_commit_metadata,
53};
54use crate::manager::sink_coordination::handle::SinkWriterCoordinationHandle;
55
56// Keep coordinator initialization below the default MetaStore pool size (10), leaving
57// connections available for recovery and other Meta services.
58const MAX_CONCURRENT_COORDINATOR_INITIALIZATIONS: usize = 8;
59static COORDINATOR_INIT_SEMAPHORE: LazyLock<Semaphore> =
60    LazyLock::new(|| Semaphore::new(MAX_CONCURRENT_COORDINATOR_INITIALIZATIONS));
61
62async fn run_future_with_periodic_fn<F: Future>(
63    future: F,
64    interval: Duration,
65    mut f: impl FnMut(),
66) -> F::Output {
67    pin_mut!(future);
68    loop {
69        match select(&mut future, pin!(sleep(interval))).await {
70            Either::Left((output, _)) => {
71                break output;
72            }
73            Either::Right(_) => f(),
74        }
75    }
76}
77
78type HandleId = usize;
79
80#[derive(Default)]
81struct AligningRequests<R> {
82    requests: Vec<R>,
83    handle_ids: HashSet<HandleId>,
84    committed_bitmap: Option<Bitmap>, // lazy-initialized on first request
85}
86
87impl<R> AligningRequests<R> {
88    fn add_new_request(
89        &mut self,
90        handle_id: HandleId,
91        request: R,
92        vnode_bitmap: &Bitmap,
93    ) -> anyhow::Result<()>
94    where
95        R: Debug,
96    {
97        let committed_bitmap = self
98            .committed_bitmap
99            .get_or_insert_with(|| Bitmap::zeros(vnode_bitmap.len()));
100        assert_eq!(committed_bitmap.len(), vnode_bitmap.len());
101
102        let check_bitmap = (&*committed_bitmap) & vnode_bitmap;
103        if check_bitmap.count_ones() > 0 {
104            return Err(anyhow!(
105                "duplicate vnode {:?}. request vnode: {:?}, prev vnode: {:?}. pending request: {:?}, request: {:?}",
106                check_bitmap.iter_ones().collect_vec(),
107                vnode_bitmap,
108                committed_bitmap,
109                self.requests,
110                request
111            ));
112        }
113        *committed_bitmap |= vnode_bitmap;
114        self.requests.push(request);
115        assert!(self.handle_ids.insert(handle_id));
116        Ok(())
117    }
118
119    fn aligned(&self) -> bool {
120        self.committed_bitmap.as_ref().is_some_and(|b| b.all())
121    }
122}
123
124type RetryBackoffFuture = std::pin::Pin<Box<tokio::time::Sleep>>;
125type RetryBackoffStrategy = impl Iterator<Item = RetryBackoffFuture> + Send + 'static;
126
127struct TwoPhaseCommitHandler {
128    db: DatabaseConnection,
129    sink_id: SinkId,
130    curr_hummock_committed_epoch: u64,
131    job_committed_epoch_rx: UnboundedReceiver<u64>,
132    last_committed_epoch: Option<u64>,
133    pending_epochs: VecDeque<(u64, Option<Vec<u8>>, Option<PbSinkSchemaChange>)>,
134    prepared_epochs: VecDeque<(u64, Option<Vec<u8>>, Option<PbSinkSchemaChange>)>,
135    backoff_state: Option<(RetryBackoffFuture, RetryBackoffStrategy)>,
136}
137
138impl TwoPhaseCommitHandler {
139    fn new(
140        db: DatabaseConnection,
141        sink_id: SinkId,
142        initial_hummock_committed_epoch: u64,
143        job_committed_epoch_rx: UnboundedReceiver<u64>,
144        last_committed_epoch: Option<u64>,
145    ) -> Self {
146        Self {
147            db,
148            sink_id,
149            curr_hummock_committed_epoch: initial_hummock_committed_epoch,
150            job_committed_epoch_rx,
151            last_committed_epoch,
152            pending_epochs: VecDeque::new(),
153            prepared_epochs: VecDeque::new(),
154            backoff_state: None,
155        }
156    }
157
158    #[define_opaque(RetryBackoffStrategy)]
159    fn get_retry_backoff_strategy() -> RetryBackoffStrategy {
160        exponential_backoff(Duration::from_millis(10), 10, Duration::from_secs(60))
161            .map(jitter)
162            .map(|delay| Box::pin(tokio::time::sleep(delay)))
163    }
164
165    async fn next_to_commit(
166        &mut self,
167    ) -> anyhow::Result<(u64, Option<Vec<u8>>, Option<PbSinkSchemaChange>)> {
168        loop {
169            let wait_backoff = async {
170                if self.prepared_epochs.is_empty() {
171                    pending::<()>().await;
172                } else if let Some((backoff_fut, _)) = &mut self.backoff_state {
173                    backoff_fut.await;
174                }
175            };
176
177            select! {
178                _ = wait_backoff => {
179                    let item = self.prepared_epochs.front().cloned().expect("non-empty");
180                    return Ok(item);
181                }
182
183                recv_epoch = self.job_committed_epoch_rx.recv() => {
184                    let Some(recv_epoch) = recv_epoch else {
185                        return Err(anyhow!(
186                            "Hummock committed epoch sender closed unexpectedly"
187                        ));
188                    };
189                    self.curr_hummock_committed_epoch = recv_epoch;
190                    while let Some((epoch, metadata, schema_change)) = self.pending_epochs.pop_front_if(|(epoch, _, _)| *epoch <= recv_epoch) {
191                        if let Some((last_epoch, _, _)) = self.prepared_epochs.back() {
192                            assert!(epoch > *last_epoch, "prepared epochs must be in increasing order");
193                        }
194                        self.prepared_epochs.push_back((epoch, metadata, schema_change));
195                    }
196                }
197            }
198        }
199    }
200
201    fn push_new_item(
202        &mut self,
203        epoch: u64,
204        metadata: Option<Vec<u8>>,
205        schema_change: Option<PbSinkSchemaChange>,
206    ) {
207        if epoch > self.curr_hummock_committed_epoch {
208            if let Some((last_epoch, _, _)) = self.pending_epochs.back() {
209                assert!(
210                    epoch > *last_epoch,
211                    "pending epochs must be in increasing order"
212                );
213            }
214            self.pending_epochs
215                .push_back((epoch, metadata, schema_change));
216        } else {
217            assert!(self.pending_epochs.is_empty());
218            if let Some((last_epoch, _, _)) = self.prepared_epochs.back() {
219                assert!(
220                    epoch > *last_epoch,
221                    "prepared epochs must be in increasing order"
222                );
223            }
224            self.prepared_epochs
225                .push_back((epoch, metadata, schema_change));
226        }
227    }
228
229    async fn ack_committed(&mut self, epoch: u64) -> anyhow::Result<()> {
230        self.backoff_state = None;
231        let (last_epoch, _, _) = self.prepared_epochs.pop_front().expect("non-empty");
232        assert_eq!(last_epoch, epoch);
233
234        commit_and_prune_epoch(&self.db, self.sink_id, epoch, self.last_committed_epoch).await?;
235        self.last_committed_epoch = Some(epoch);
236        Ok(())
237    }
238
239    fn failed_committed(&mut self, epoch: u64, err: SinkError) {
240        assert_eq!(self.prepared_epochs.front().expect("non-empty").0, epoch,);
241        if let Some((prev_fut, strategy)) = &mut self.backoff_state {
242            let new_fut = strategy.next().expect("infinite");
243            *prev_fut = new_fut;
244        } else {
245            let mut strategy = Self::get_retry_backoff_strategy();
246            let backoff_fut = strategy.next().expect("infinite");
247            self.backoff_state = Some((backoff_fut, strategy));
248        }
249        tracing::error!(
250            error = %err.as_report(),
251            %self.sink_id,
252            "failed to commit epoch {}, Retrying after backoff",
253            epoch,
254        );
255    }
256
257    fn is_empty(&self) -> bool {
258        self.pending_epochs.is_empty() && self.prepared_epochs.is_empty()
259    }
260
261    /// Whether there exists an uncommitted schema change.
262    ///
263    /// Per current design, if a `schema_change` exists, it should be attached to the latest
264    /// uncommitted item across `pending_epochs` and `prepared_epochs`.
265    fn has_uncommitted_schema_change(&self) -> bool {
266        if let Some((_, _, schema_change)) = self.pending_epochs.back() {
267            schema_change.is_some()
268        } else if let Some((_, _, schema_change)) = self.prepared_epochs.back() {
269            schema_change.is_some()
270        } else {
271            false
272        }
273    }
274}
275
276struct CoordinationHandleManager {
277    param: SinkParam,
278    writer_handles: HashMap<HandleId, SinkWriterCoordinationHandle>,
279    next_handle_id: HandleId,
280    request_rx: UnboundedReceiver<SinkWriterCoordinationHandle>,
281}
282
283impl CoordinationHandleManager {
284    fn start(
285        &mut self,
286        log_store_rewind_start_epoch: Option<u64>,
287        handle_ids: impl IntoIterator<Item = HandleId>,
288    ) -> anyhow::Result<()> {
289        for handle_id in handle_ids {
290            let handle = self
291                .writer_handles
292                .get_mut(&handle_id)
293                .ok_or_else(|| anyhow!("failed to find handle {} to start", handle_id,))?;
294            handle.start(log_store_rewind_start_epoch).map_err(|_| {
295                anyhow!(
296                    "failed to start {:?} for handle {}",
297                    log_store_rewind_start_epoch,
298                    handle_id
299                )
300            })?;
301        }
302        Ok(())
303    }
304
305    fn ack_aligned_initial_epoch(&mut self, aligned_initial_epoch: u64) -> anyhow::Result<()> {
306        for (handle_id, handle) in &mut self.writer_handles {
307            handle
308                .ack_aligned_initial_epoch(aligned_initial_epoch)
309                .map_err(|_| {
310                    anyhow!(
311                        "failed to ack aligned initial epoch {:?} for handle {}",
312                        aligned_initial_epoch,
313                        handle_id
314                    )
315                })?;
316        }
317        Ok(())
318    }
319
320    fn ack_commit(
321        &mut self,
322        epoch: u64,
323        handle_ids: impl IntoIterator<Item = HandleId>,
324    ) -> anyhow::Result<()> {
325        for handle_id in handle_ids {
326            let handle = self.writer_handles.get_mut(&handle_id).ok_or_else(|| {
327                anyhow!(
328                    "failed to find handle {} when acknowledging the commit for epoch {}",
329                    handle_id,
330                    epoch
331                )
332            })?;
333            handle.ack_commit(epoch).map_err(|_| {
334                anyhow!(
335                    "failed to acknowledge the commit for epoch {} on handle {}",
336                    epoch,
337                    handle_id
338                )
339            })?;
340        }
341        Ok(())
342    }
343
344    async fn next_request_inner(
345        writer_handles: &mut HashMap<HandleId, SinkWriterCoordinationHandle>,
346    ) -> anyhow::Result<(HandleId, coordinate_request::Msg)> {
347        poll_fn(|cx| {
348            for (handle_id, handle) in writer_handles.iter_mut() {
349                if let Poll::Ready(result) = handle.poll_next_request(cx) {
350                    return Poll::Ready(result.map(|request| (*handle_id, request)));
351                }
352            }
353            Poll::Pending
354        })
355        .await
356    }
357}
358
359enum CoordinationHandleManagerEvent {
360    NewHandle,
361    UpdateVnodeBitmap,
362    Stop,
363    CommitRequest {
364        epoch: u64,
365        metadata: SinkMetadata,
366        schema_change: Option<PbSinkSchemaChange>,
367    },
368    AlignInitialEpoch(u64),
369}
370
371impl CoordinationHandleManagerEvent {
372    fn name(&self) -> &'static str {
373        match self {
374            CoordinationHandleManagerEvent::NewHandle => "NewHandle",
375            CoordinationHandleManagerEvent::UpdateVnodeBitmap => "UpdateVnodeBitmap",
376            CoordinationHandleManagerEvent::Stop => "Stop",
377            CoordinationHandleManagerEvent::CommitRequest { .. } => "CommitRequest",
378            CoordinationHandleManagerEvent::AlignInitialEpoch(_) => "AlignInitialEpoch",
379        }
380    }
381}
382
383impl CoordinationHandleManager {
384    async fn next_event(&mut self) -> anyhow::Result<(HandleId, CoordinationHandleManagerEvent)> {
385        select! {
386            handle = self.request_rx.recv() => {
387                let handle = handle.ok_or_else(|| anyhow!("end of writer request stream"))?;
388                if handle.param() != &self.param {
389                    warn!(prev_param = ?self.param, new_param = ?handle.param(), "sink param mismatch");
390                }
391                let handle_id = self.next_handle_id;
392                self.next_handle_id += 1;
393                self.writer_handles.insert(handle_id, handle);
394                Ok((handle_id, CoordinationHandleManagerEvent::NewHandle))
395            }
396            result = Self::next_request_inner(&mut self.writer_handles) => {
397                let (handle_id, request) = result?;
398                let event = match request {
399                    coordinate_request::Msg::CommitRequest(request) => {
400                        CoordinationHandleManagerEvent::CommitRequest {
401                            epoch: request.epoch,
402                            metadata: request.metadata.ok_or_else(|| anyhow!("empty sink metadata"))?,
403                            schema_change: request.schema_change,
404                        }
405                    }
406                    coordinate_request::Msg::AlignInitialEpochRequest(epoch) => {
407                        CoordinationHandleManagerEvent::AlignInitialEpoch(epoch)
408                    }
409                    coordinate_request::Msg::UpdateVnodeRequest(_) => {
410                        CoordinationHandleManagerEvent::UpdateVnodeBitmap
411                    }
412                    coordinate_request::Msg::Stop(_) => {
413                        CoordinationHandleManagerEvent::Stop
414                    }
415                    coordinate_request::Msg::StartRequest(_) => {
416                        unreachable!("should have been handled");
417                    }
418                };
419                Ok((handle_id, event))
420            }
421        }
422    }
423
424    fn vnode_bitmap(&self, handle_id: HandleId) -> &Bitmap {
425        self.writer_handles[&handle_id].vnode_bitmap()
426    }
427
428    fn stop_handle(&mut self, handle_id: HandleId) -> anyhow::Result<()> {
429        self.writer_handles
430            .remove(&handle_id)
431            .expect("should exist")
432            .stop()
433    }
434
435    async fn wait_init_handles(&mut self) -> anyhow::Result<HashSet<HandleId>> {
436        assert!(self.writer_handles.is_empty());
437        let mut init_requests = AligningRequests::default();
438        while !init_requests.aligned() {
439            let (handle_id, event) = self.next_event().await?;
440            let unexpected_event = match event {
441                CoordinationHandleManagerEvent::NewHandle => {
442                    init_requests.add_new_request(handle_id, (), self.vnode_bitmap(handle_id))?;
443                    continue;
444                }
445                event => event.name(),
446            };
447            return Err(anyhow!(
448                "expect new handle during init, but got {}",
449                unexpected_event
450            ));
451        }
452        Ok(init_requests.handle_ids)
453    }
454
455    async fn alter_parallelisms(
456        &mut self,
457        altered_handles: impl Iterator<Item = HandleId>,
458    ) -> anyhow::Result<HashSet<HandleId>> {
459        let mut requests = AligningRequests::default();
460        for handle_id in altered_handles {
461            requests.add_new_request(handle_id, (), self.vnode_bitmap(handle_id))?;
462        }
463        let mut remaining_handles: HashSet<_> = self
464            .writer_handles
465            .keys()
466            .filter(|handle_id| !requests.handle_ids.contains(handle_id))
467            .cloned()
468            .collect();
469        while !remaining_handles.is_empty() || !requests.aligned() {
470            let (handle_id, event) = self.next_event().await?;
471            match event {
472                CoordinationHandleManagerEvent::NewHandle => {
473                    requests.add_new_request(handle_id, (), self.vnode_bitmap(handle_id))?;
474                }
475                CoordinationHandleManagerEvent::UpdateVnodeBitmap => {
476                    assert!(remaining_handles.remove(&handle_id));
477                    requests.add_new_request(handle_id, (), self.vnode_bitmap(handle_id))?;
478                }
479                CoordinationHandleManagerEvent::Stop => {
480                    assert!(remaining_handles.remove(&handle_id));
481                    self.stop_handle(handle_id)?;
482                }
483                CoordinationHandleManagerEvent::CommitRequest { epoch, .. } => {
484                    bail!(
485                        "receive commit request on epoch {} from handle {} during alter parallelism",
486                        epoch,
487                        handle_id
488                    );
489                }
490                CoordinationHandleManagerEvent::AlignInitialEpoch(epoch) => {
491                    bail!(
492                        "receive AlignInitialEpoch on epoch {} from handle {} during alter parallelism",
493                        epoch,
494                        handle_id
495                    );
496                }
497            }
498        }
499        Ok(requests.handle_ids)
500    }
501}
502
503/// Represents the coordinator worker's state machine for handling schema changes.
504///
505/// - `Running`: Normal operation, handles can be started immediately
506/// - `WaitingForFlushed`: Waiting for all pending two-phase commits to complete before starting new handles. This
507///   ensures new sink executors load the correct schema.
508enum CoordinatorWorkerState {
509    Running,
510    WaitingForFlushed(HashSet<HandleId>),
511}
512
513pub struct CoordinatorWorker {
514    handle_manager: CoordinationHandleManager,
515    /// Last epoch whose commit has been acknowledged to sink writers.
516    ///
517    /// On recovery, pending sink states are treated as already acknowledged to writers, so this is
518    /// initialized to the latest pending epoch if any. Otherwise, it starts from the latest
519    /// committed epoch persisted in `pending_sink_state`.
520    last_writer_acked_epoch: Option<u64>,
521    curr_state: CoordinatorWorkerState,
522}
523
524enum CoordinatorWorkerEvent {
525    HandleManagerEvent(HandleId, CoordinationHandleManagerEvent),
526    ReadyToCommit(u64, Option<Vec<u8>>, Option<PbSinkSchemaChange>),
527}
528
529impl CoordinatorWorker {
530    pub async fn run(
531        param: SinkParam,
532        request_rx: UnboundedReceiver<SinkWriterCoordinationHandle>,
533        db: DatabaseConnection,
534        subscriber: SinkCommittedEpochSubscriber,
535        iceberg_compact_stat_sender: UnboundedSender<IcebergSinkCompactionUpdate>,
536    ) {
537        let sink = match build_sink(param.clone()) {
538            Ok(sink) => sink,
539            Err(e) => {
540                error!(
541                    error = %e.as_report(),
542                    "unable to build sink with param {:?}",
543                    param
544                );
545                return;
546            }
547        };
548
549        dispatch_sink!(sink, sink, {
550            let coordinator =
551                match Box::pin(sink.new_coordinator(Some(iceberg_compact_stat_sender))).await {
552                    Ok(coordinator) => coordinator,
553                    Err(e) => {
554                        error!(
555                            error = %e.as_report(),
556                            "unable to build coordinator with param {:?}",
557                            param
558                        );
559                        return;
560                    }
561                };
562            Self::execute_coordinator(db, param, request_rx, coordinator, subscriber).await
563        });
564    }
565
566    pub async fn execute_coordinator(
567        db: DatabaseConnection,
568        param: SinkParam,
569        request_rx: UnboundedReceiver<SinkWriterCoordinationHandle>,
570        coordinator: SinkCommitCoordinator,
571        subscriber: SinkCommittedEpochSubscriber,
572    ) {
573        let mut worker = CoordinatorWorker {
574            handle_manager: CoordinationHandleManager {
575                param,
576                writer_handles: HashMap::new(),
577                next_handle_id: 0,
578                request_rx,
579            },
580            last_writer_acked_epoch: None,
581            curr_state: CoordinatorWorkerState::Running,
582        };
583
584        if let Err(e) = worker.run_coordination(db, coordinator, subscriber).await {
585            for handle in worker.handle_manager.writer_handles.into_values() {
586                handle.abort(Status::internal(format!(
587                    "failed to run coordination: {:?}",
588                    e.as_report()
589                )))
590            }
591        }
592    }
593
594    async fn try_handle_init_requests(
595        &mut self,
596        pending_handle_ids: &HashSet<HandleId>,
597        two_phase_handler: &mut TwoPhaseCommitHandler,
598    ) -> anyhow::Result<()> {
599        assert!(matches!(self.curr_state, CoordinatorWorkerState::Running));
600        if two_phase_handler.has_uncommitted_schema_change() {
601            // Delay handling init requests until all pending epochs are flushed.
602            self.curr_state = CoordinatorWorkerState::WaitingForFlushed(pending_handle_ids.clone());
603        } else {
604            self.handle_init_requests_impl(pending_handle_ids.clone())
605                .await?;
606        }
607        Ok(())
608    }
609
610    async fn handle_init_requests_impl(
611        &mut self,
612        pending_handle_ids: impl IntoIterator<Item = HandleId>,
613    ) -> anyhow::Result<()> {
614        let log_store_rewind_start_epoch = self.last_writer_acked_epoch;
615        self.handle_manager
616            .start(log_store_rewind_start_epoch, pending_handle_ids)?;
617        if log_store_rewind_start_epoch.is_none() {
618            let mut align_requests = AligningRequests::default();
619            while !align_requests.aligned() {
620                let (handle_id, event) = self.handle_manager.next_event().await?;
621                match event {
622                    CoordinationHandleManagerEvent::AlignInitialEpoch(initial_epoch) => {
623                        align_requests.add_new_request(
624                            handle_id,
625                            initial_epoch,
626                            self.handle_manager.vnode_bitmap(handle_id),
627                        )?;
628                    }
629                    other => {
630                        return Err(anyhow!("expect AlignInitialEpoch but got {}", other.name()));
631                    }
632                }
633            }
634            let aligned_initial_epoch = align_requests
635                .requests
636                .into_iter()
637                .max()
638                .expect("non-empty");
639            self.handle_manager
640                .ack_aligned_initial_epoch(aligned_initial_epoch)?;
641        }
642        Ok(())
643    }
644
645    async fn next_event(
646        &mut self,
647        two_phase_handler: &mut TwoPhaseCommitHandler,
648    ) -> anyhow::Result<CoordinatorWorkerEvent> {
649        if let CoordinatorWorkerState::WaitingForFlushed(pending_handle_ids) = &self.curr_state
650            && two_phase_handler.is_empty()
651        {
652            let pending_handle_ids = pending_handle_ids.clone();
653            self.handle_init_requests_impl(pending_handle_ids).await?;
654            self.curr_state = CoordinatorWorkerState::Running;
655        }
656
657        select! {
658            next_handle_event = self.handle_manager.next_event() => {
659                let (handle_id, event) = next_handle_event?;
660                Ok(CoordinatorWorkerEvent::HandleManagerEvent(handle_id, event))
661            }
662
663            next_item_to_commit = two_phase_handler.next_to_commit() => {
664                let (epoch, metadata, schema_change) = next_item_to_commit?;
665                Ok(CoordinatorWorkerEvent::ReadyToCommit(epoch, metadata, schema_change))
666            }
667        }
668    }
669
670    async fn run_coordination(
671        &mut self,
672        db: DatabaseConnection,
673        mut coordinator: SinkCommitCoordinator,
674        subscriber: SinkCommittedEpochSubscriber,
675    ) -> anyhow::Result<()> {
676        let sink_id = self.handle_manager.param.sink_id;
677
678        let coordinator_init_permit = COORDINATOR_INIT_SEMAPHORE
679            .acquire()
680            .instrument_await("acquire_sink_coordinator_init_permit")
681            .await?;
682        let mut two_phase_handler = self
683            .init_state_from_store(&db, sink_id, subscriber, &mut coordinator)
684            .await?;
685        match &mut coordinator {
686            SinkCommitCoordinator::SinglePhase(coordinator) => coordinator.init().await?,
687            SinkCommitCoordinator::TwoPhase(coordinator) => coordinator.init().await?,
688        }
689        drop(coordinator_init_permit);
690
691        let mut running_handles = self.handle_manager.wait_init_handles().await?;
692        self.try_handle_init_requests(&running_handles, &mut two_phase_handler)
693            .await?;
694
695        let mut pending_epochs: BTreeMap<u64, AligningRequests<_>> = BTreeMap::new();
696        let mut pending_new_handles = vec![];
697        loop {
698            let event = self.next_event(&mut two_phase_handler).await?;
699            let (handle_id, epoch, commit_request) = match event {
700                CoordinatorWorkerEvent::HandleManagerEvent(handle_id, event) => match event {
701                    CoordinationHandleManagerEvent::NewHandle => {
702                        pending_new_handles.push(handle_id);
703                        continue;
704                    }
705                    CoordinationHandleManagerEvent::UpdateVnodeBitmap => {
706                        running_handles = self
707                            .handle_manager
708                            .alter_parallelisms(pending_new_handles.drain(..).chain([handle_id]))
709                            .await?;
710                        self.try_handle_init_requests(&running_handles, &mut two_phase_handler)
711                            .await?;
712                        continue;
713                    }
714                    CoordinationHandleManagerEvent::Stop => {
715                        self.handle_manager.stop_handle(handle_id)?;
716                        running_handles = self
717                            .handle_manager
718                            .alter_parallelisms(pending_new_handles.drain(..))
719                            .await?;
720                        self.try_handle_init_requests(&running_handles, &mut two_phase_handler)
721                            .await?;
722
723                        continue;
724                    }
725                    CoordinationHandleManagerEvent::CommitRequest {
726                        epoch,
727                        metadata,
728                        schema_change,
729                    } => (handle_id, epoch, (metadata, schema_change)),
730                    CoordinationHandleManagerEvent::AlignInitialEpoch(_) => {
731                        bail!("receive AlignInitialEpoch after initialization")
732                    }
733                },
734                CoordinatorWorkerEvent::ReadyToCommit(epoch, metadata, schema_change) => {
735                    let start_time = Instant::now();
736                    let commit_fut = async {
737                        match &mut coordinator {
738                            SinkCommitCoordinator::SinglePhase(coordinator) => {
739                                assert!(metadata.is_none());
740                                if let Some(schema_change) = schema_change {
741                                    coordinator
742                                        .commit_schema_change(epoch, schema_change)
743                                        .instrument_await(Self::commit_span(
744                                            "single_phase_schema_change",
745                                            sink_id,
746                                            epoch,
747                                        ))
748                                        .await?;
749                                }
750                            }
751                            SinkCommitCoordinator::TwoPhase(coordinator) => {
752                                if let Some(metadata) = metadata {
753                                    coordinator
754                                        .commit_data(epoch, metadata)
755                                        .instrument_await(Self::commit_span(
756                                            "two_phase_commit_data",
757                                            sink_id,
758                                            epoch,
759                                        ))
760                                        .await?;
761                                }
762                                if let Some(schema_change) = schema_change {
763                                    coordinator
764                                        .commit_schema_change(epoch, schema_change)
765                                        .instrument_await(Self::commit_span(
766                                            "two_phase_commit_schema_change",
767                                            sink_id,
768                                            epoch,
769                                        ))
770                                        .await?;
771                                }
772                            }
773                        }
774                        Ok(())
775                    };
776                    let commit_res =
777                        run_future_with_periodic_fn(commit_fut, Duration::from_secs(5), || {
778                            warn!(
779                                elapsed = ?start_time.elapsed(),
780                                %sink_id,
781                                "committing"
782                            );
783                        })
784                        .await;
785
786                    match commit_res {
787                        Ok(_) => {
788                            two_phase_handler.ack_committed(epoch).await?;
789                        }
790                        Err(e) => {
791                            two_phase_handler.failed_committed(epoch, e);
792                        }
793                    }
794
795                    continue;
796                }
797            };
798            if !running_handles.contains(&handle_id) {
799                bail!(
800                    "receiving commit request from non-running handle {}, running handles: {:?}",
801                    handle_id,
802                    running_handles
803                );
804            }
805            pending_epochs.entry(epoch).or_default().add_new_request(
806                handle_id,
807                commit_request,
808                self.handle_manager.vnode_bitmap(handle_id),
809            )?;
810            if pending_epochs
811                .first_key_value()
812                .expect("non-empty")
813                .1
814                .aligned()
815            {
816                let (epoch, commit_requests) = pending_epochs.pop_first().expect("non-empty");
817                let mut metadatas = Vec::with_capacity(commit_requests.requests.len());
818                let mut requests = commit_requests.requests.into_iter();
819                let (first_metadata, first_schema_change) = requests.next().expect("non-empty");
820                metadatas.push(first_metadata);
821                for (metadata, schema_change) in requests {
822                    if first_schema_change != schema_change {
823                        return Err(anyhow!(
824                            "got different schema change {:?} to prev schema change {:?}",
825                            schema_change,
826                            first_schema_change
827                        ));
828                    }
829                    metadatas.push(metadata);
830                }
831
832                match &mut coordinator {
833                    SinkCommitCoordinator::SinglePhase(coordinator) => {
834                        if !metadatas.is_empty() {
835                            let start_time = Instant::now();
836                            run_future_with_periodic_fn(
837                                coordinator.commit_data(epoch, metadatas).instrument_await(
838                                    Self::commit_span("single_phase_commit_data", sink_id, epoch),
839                                ),
840                                Duration::from_secs(5),
841                                || {
842                                    warn!(
843                                        elapsed = ?start_time.elapsed(),
844                                        %sink_id,
845                                        "committing"
846                                    );
847                                },
848                            )
849                            .await
850                            .map_err(|e| anyhow!(e))?;
851                        }
852                        if first_schema_change.is_some() {
853                            persist_pre_commit_metadata(
854                                &db,
855                                sink_id as _,
856                                epoch,
857                                None,
858                                first_schema_change.as_ref(),
859                            )
860                            .await?;
861                            two_phase_handler.push_new_item(epoch, None, first_schema_change);
862                        }
863                    }
864                    SinkCommitCoordinator::TwoPhase(coordinator) => {
865                        let commit_metadata = coordinator
866                            .pre_commit(epoch, metadatas, first_schema_change.clone())
867                            .instrument_await(Self::commit_span(
868                                "two_phase_pre_commit",
869                                sink_id,
870                                epoch,
871                            ))
872                            .await?;
873                        // Persist every acknowledged epoch, even when there is no metadata or
874                        // schema change. Writers may truncate the epoch as soon as they receive
875                        // the acknowledgement, so recovery must retain the same progress. The
876                        // commit handler treats a `None`/`None` item as an external no-op before
877                        // marking it committed, and recovery re-enqueues it through the same path.
878                        persist_pre_commit_metadata(
879                            &db,
880                            sink_id as _,
881                            epoch,
882                            commit_metadata.clone(),
883                            first_schema_change.as_ref(),
884                        )
885                        .await?;
886                        two_phase_handler.push_new_item(
887                            epoch,
888                            commit_metadata,
889                            first_schema_change,
890                        );
891                    }
892                }
893
894                self.handle_manager
895                    .ack_commit(epoch, commit_requests.handle_ids)?;
896                self.last_writer_acked_epoch = Some(epoch);
897            }
898        }
899    }
900
901    /// Return `TwoPhaseCommitHandler` initialized from the persisted state in the meta store.
902    async fn init_state_from_store(
903        &mut self,
904        db: &DatabaseConnection,
905        sink_id: SinkId,
906        subscriber: SinkCommittedEpochSubscriber,
907        coordinator: &mut SinkCommitCoordinator,
908    ) -> anyhow::Result<TwoPhaseCommitHandler> {
909        let ordered_metadata = list_sink_states_ordered_by_epoch(db, sink_id as _).await?;
910
911        let mut metadata_iter = ordered_metadata.into_iter().peekable();
912        let last_committed_epoch = metadata_iter
913            .next_if(|(_, state, _, _)| matches!(state, SinkState::Committed))
914            .map(|(epoch, _, _, _)| epoch);
915
916        let pending_items = metadata_iter
917            .peeking_take_while(|(_, state, _, _)| matches!(state, SinkState::Pending))
918            .map(|(epoch, _, metadata, schema_change)| (epoch, metadata, schema_change))
919            .collect_vec();
920        self.last_writer_acked_epoch = pending_items
921            .last()
922            .map(|(epoch, _, _)| *epoch)
923            .or(last_committed_epoch);
924
925        let mut aborted_epochs = vec![];
926
927        for (epoch, state, metadata, _) in metadata_iter {
928            match state {
929                SinkState::Aborted => {
930                    if let Some(metadata) = metadata
931                        && let SinkCommitCoordinator::TwoPhase(coordinator) = coordinator
932                    {
933                        coordinator.abort(epoch, metadata).await;
934                    }
935                    aborted_epochs.push(epoch);
936                }
937                other => {
938                    unreachable!(
939                        "unexpected state {:?} after pending items at epoch {}",
940                        other, epoch
941                    );
942                }
943            }
944        }
945
946        // Records for all aborted epochs and previously committed epochs are no longer needed.
947        clean_aborted_records(db, sink_id, aborted_epochs).await?;
948
949        let (initial_hummock_committed_epoch, job_committed_epoch_rx) = subscriber(sink_id).await?;
950        let mut two_phase_handler = TwoPhaseCommitHandler::new(
951            db.clone(),
952            sink_id,
953            initial_hummock_committed_epoch,
954            job_committed_epoch_rx,
955            last_committed_epoch,
956        );
957
958        for (epoch, metadata, schema_change) in pending_items {
959            two_phase_handler.push_new_item(epoch, metadata, schema_change);
960        }
961
962        Ok(two_phase_handler)
963    }
964
965    fn commit_span(stage: &str, sink_id: SinkId, epoch: u64) -> await_tree::Span {
966        await_tree::span!("sink_coord_{stage} (sink_id {sink_id}, epoch {epoch})").long_running()
967    }
968}