Skip to main content

camel_master/
supervision.rs

1use std::sync::Arc;
2use std::time::Duration;
3
4use async_trait::async_trait;
5use camel_api::CamelError;
6use camel_component_api::{Consumer, ConsumerContext, is_retryable_camel_error};
7use tokio::time::{interval, timeout};
8use tokio_util::sync::CancellationToken;
9use tracing::{error, warn};
10
11use crate::consumer::{DelegateState, MasterConsumer};
12use crate::leadership::{ReconcileContext, reconcile_event, stop_delegate};
13
14const DELEGATE_RETRY_INTERVAL: Duration = Duration::from_millis(200);
15
16#[async_trait]
17impl Consumer for MasterConsumer {
18    async fn start(&mut self, context: ConsumerContext) -> Result<(), CamelError> {
19        if self.leadership_task.is_some() {
20            return Ok(());
21        }
22
23        let handle = self
24            .platform_service
25            .leadership()
26            .start(&self.lock_name)
27            .await
28            .map_err(|e| {
29                CamelError::EndpointCreationFailed(format!("failed to start leader election: {e}"))
30            })?;
31
32        let lock_name = self.lock_name.clone();
33        let delegate_uri = self.delegate_uri.clone();
34        let delegate_component = Arc::clone(&self.delegate_component);
35        let metrics = Arc::clone(&self.metrics);
36        let platform_service = Arc::clone(&self.platform_service);
37        let sender = context.sender();
38        let parent_cancel = context.cancel_token();
39        let route_id = context.route_id().to_string();
40        let drain_timeout = self.drain_timeout;
41        let reconnect = self.reconnect.clone();
42        let runtime = Arc::clone(&self.runtime);
43        let mut events = handle.events.clone();
44
45        let stop_token = CancellationToken::new();
46        let stop_token_loop = stop_token.clone();
47        let leadership_handle = handle;
48        let leader_epoch = leadership_handle.leader_epoch_arc();
49
50        let task = tokio::spawn(async move {
51            let mut state = DelegateState::Inactive;
52            let mut is_leading = false;
53            let mut delegate_attempts = 0u32;
54            let mut retry_tick = interval(DELEGATE_RETRY_INTERVAL);
55
56            let rctx = ReconcileContext {
57                lock_name: &lock_name,
58                delegate_component: &delegate_component,
59                delegate_uri: &delegate_uri,
60                route_id,
61                sender: &sender,
62                parent_cancel: &parent_cancel,
63                drain_timeout,
64                metrics: &metrics,
65                platform_service: &platform_service,
66                runtime: Arc::clone(&runtime),
67                leader_epoch: Arc::clone(&leader_epoch),
68            };
69
70            let initial_event = { events.borrow().clone() };
71            if let Some(initial_event) = initial_event {
72                is_leading = matches!(&initial_event, camel_api::LeadershipEvent::StartedLeading);
73                if is_leading {
74                    delegate_attempts = 0;
75                }
76                if let Err(err) = reconcile_event(initial_event, &mut state, &rctx).await {
77                    // log-policy: system-broken
78                    error!(lock = %lock_name, "master delegate error: {err}");
79                    return Err(err);
80                }
81            }
82
83            loop {
84                tokio::select! {
85                    _ = stop_token_loop.cancelled() => {
86                        break;
87                    }
88                    _ = context.cancelled() => {
89                        break;
90                    }
91                    changed = events.changed() => {
92                        if changed.is_err() {
93                            break;
94                        }
95                        let event = { events.borrow().clone() };
96                        if let Some(event) = event {
97                            let was_leading = is_leading;
98                            is_leading = matches!(&event, camel_api::LeadershipEvent::StartedLeading);
99                            if !was_leading && is_leading {
100                                delegate_attempts = 0;
101                            }
102                            if let Err(err) = reconcile_event(event, &mut state, &rctx).await {
103                                // log-policy: system-broken
104                                error!(lock = %lock_name, "master delegate error: {err}");
105                                return Err(err);
106                            }
107                        }
108                    }
109                    _ = retry_tick.tick() => {
110                        if matches!(&state, DelegateState::Active { handle, .. } if handle.is_finished())
111                            && let Err(err) = stop_delegate(&mut state, drain_timeout).await
112                        {
113                            // log-policy: system-broken
114                            error!(lock = %lock_name, "master delegate task failed: {err}");
115                            return Err(err);
116                        }
117
118                        if is_leading && matches!(state, DelegateState::Inactive) {
119                            // Manual retry loop (not retry_async) because:
120                            // - The retry logic is embedded inside a periodic
121                            //   retry_tick.tick() handler; the outer select! runs
122                            //   every DELEGATE_RETRY_INTERVAL regardless, so the
123                            //   delay is applied as an additive sleep on top of
124                            //   the tick interval, not as a replacement for it.
125                            // - reconcile_event() requires &mut state, and the
126                            //   inter-attempt logic checks handle.is_finished()
127                            //   before retrying — both require state access
128                            //   between iterations that retry_async cannot provide.
129                            // - Classifies errors (rc-i1z): permanent → fail-fast,
130                            //   transient → retry with backoff.
131                            // Use NetworkRetryPolicy for bounded retries.
132                            // delegate_attempts tracks the next zero-based attempt index.
133                            if !reconnect.should_retry(delegate_attempts) {
134                                warn!(
135                                    lock = %lock_name,
136                                    attempts = delegate_attempts,
137                                    "delegate start exceeded max attempts, stopping consumer"
138                                );
139                                break;
140                            }
141                            // Apply backoff delay for retries (skip first attempt).
142                            if delegate_attempts > 0 {
143                                let delay = reconnect.delay_for(delegate_attempts - 1);
144                                if delay > DELEGATE_RETRY_INTERVAL {
145                                    tokio::select! {
146                                        _ = stop_token_loop.cancelled() => break,
147                                        _ = tokio::time::sleep(delay.saturating_sub(DELEGATE_RETRY_INTERVAL)) => {}
148                                    }
149                                }
150                            }
151                            delegate_attempts = delegate_attempts.saturating_add(1);
152                            if let Err(err) = reconcile_event(
153                                camel_api::LeadershipEvent::StartedLeading,
154                                &mut state,
155                                &rctx,
156                            )
157                            .await {
158                                if is_retryable_camel_error(&err) {
159                                    // log-policy: system-broken
160                                    error!(
161                                        lock = %lock_name,
162                                        error = %err,
163                                        attempt = delegate_attempts,
164                                        "master delegate transient error, will retry"
165                                    );
166                                    // Don't return — let the next tick attempt retry.
167                                } else {
168                                    // log-policy: system-broken
169                                    error!(
170                                        lock = %lock_name,
171                                        error = %err,
172                                        "master delegate permanent error, terminating"
173                                    );
174                                    return Err(err);
175                                }
176                            }
177                        }
178                    }
179                }
180            }
181
182            stop_delegate(&mut state, drain_timeout).await?;
183            let _ = timeout(drain_timeout, leadership_handle.step_down()).await;
184            Ok::<(), CamelError>(())
185        });
186
187        self.stop_token = Some(stop_token);
188        self.leadership_task = Some(task);
189
190        Ok(())
191    }
192
193    async fn stop(&mut self) -> Result<(), CamelError> {
194        if let Some(token) = self.stop_token.take() {
195            token.cancel();
196        }
197
198        if let Some(handle) = self.leadership_task.take() {
199            if handle.is_finished() {
200                match timeout(self.drain_timeout, handle).await {
201                    Ok(Ok(Ok(()))) => {}
202                    Ok(Ok(Err(err))) => return Err(err),
203                    Ok(Err(e)) => {
204                        return Err(CamelError::ProcessorError(format!(
205                            "leadership task join failed: {e}"
206                        )));
207                    }
208                    Err(_) => {
209                        return Err(CamelError::ProcessorError(
210                            "leadership task join timed out".to_string(),
211                        ));
212                    }
213                }
214                return Ok(());
215            }
216
217            // Abort first so the task is guaranteed to stop; then await with
218            // a timeout as a safety-net in case abort takes a moment to land.
219            handle.abort();
220            match timeout(self.drain_timeout, handle).await {
221                Ok(Ok(Ok(()))) => {}
222                Ok(Ok(Err(err))) => return Err(err),
223                Ok(Err(e)) if e.is_panic() => {
224                    // log-policy: system-broken
225                    error!(lock = %self.lock_name, error = %e, "leadership task panicked");
226                }
227                Ok(Err(e)) => {
228                    warn!(lock = %self.lock_name, error = %e, "leadership task cancelled");
229                }
230                Err(_) => {
231                    warn!("master leadership loop shutdown timed out after abort");
232                }
233            }
234        }
235
236        Ok(())
237    }
238
239    fn background_task_handle(
240        &mut self,
241    ) -> Option<tokio::task::JoinHandle<Result<(), CamelError>>> {
242        self.leadership_task.take()
243    }
244}