camel_master/
supervision.rs1use 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 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 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 error!(lock = %lock_name, "master delegate task failed: {err}");
115 return Err(err);
116 }
117
118 if is_leading && matches!(state, DelegateState::Inactive) {
119 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 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 error!(
161 lock = %lock_name,
162 error = %err,
163 attempt = delegate_attempts,
164 "master delegate transient error, will retry"
165 );
166 } else {
168 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 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 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}