1use super::{
16 OpenedActivation, PrepareTargetResult, PreparedActivation, ReplayEvent, ReplayWorkerManager, RuntimeActivation,
17 RuntimeStatusSnapshot, RuntimeTargetHealthSnapshot, TargetActivationFailure, TargetRuntimeManager, prepare_target,
18 start_replay_worker,
19};
20use crate::plugin::PluginEvent;
21use crate::{SharedTarget, Target, TargetError};
22use async_trait::async_trait;
23use rayon::prelude::*;
24use std::future::Future;
25use std::panic::{AssertUnwindSafe, catch_unwind};
26use std::pin::Pin;
27use std::sync::{Arc, LazyLock};
28use std::time::Duration;
29use tokio::sync::Semaphore;
30use tokio_util::sync::CancellationToken;
31
32type ReplayHook<E> = Arc<dyn Fn(ReplayEvent<E>) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync>;
33type ReplayStartObserver = Arc<dyn Fn(&str, bool) + Send + Sync>;
34
35const MAX_PARALLEL_STORE_OPENS: usize = 4;
36
37static STORE_OPEN_POOL: LazyLock<Result<rayon::ThreadPool, rayon::ThreadPoolBuildError>> = LazyLock::new(|| {
38 rayon::ThreadPoolBuilder::new()
39 .num_threads(MAX_PARALLEL_STORE_OPENS)
40 .thread_name(|index| format!("rustfs-target-store-open-{index}"))
41 .build()
42});
43
44enum StoreOpenOutcome<E>
45where
46 E: PluginEvent,
47{
48 Accepted(SharedTarget<E>),
49 Rejected { panicked: bool, target: SharedTarget<E> },
50}
51
52fn open_target_store<E>(target: SharedTarget<E>) -> StoreOpenOutcome<E>
53where
54 E: PluginEvent,
55{
56 match catch_unwind(AssertUnwindSafe(|| target.store().map(|store| store.open()))) {
57 Ok(None | Some(Ok(()))) => StoreOpenOutcome::Accepted(target),
58 Ok(Some(Err(_))) => StoreOpenOutcome::Rejected { panicked: false, target },
59 Err(_) => StoreOpenOutcome::Rejected { panicked: true, target },
60 }
61}
62
63#[async_trait]
65pub trait PluginRuntimeAdapter<E>: Send + Sync
66where
67 E: PluginEvent,
68{
69 async fn activate_with_replay(&self, targets: Vec<Box<dyn Target<E> + Send + Sync>>) -> RuntimeActivation<E>;
70
71 async fn replace_runtime_targets(
72 &self,
73 runtime: &mut TargetRuntimeManager<E>,
74 replay_workers: &mut ReplayWorkerManager,
75 activation: RuntimeActivation<E>,
76 ) -> Result<(), TargetError>;
77
78 async fn stop_replay_workers(&self, replay_workers: &mut ReplayWorkerManager);
79
80 fn snapshot_runtime_status(
81 &self,
82 runtime: &TargetRuntimeManager<E>,
83 replay_workers: &ReplayWorkerManager,
84 ) -> RuntimeStatusSnapshot;
85
86 async fn snapshot_runtime_health(&self, runtime: &TargetRuntimeManager<E>) -> Vec<RuntimeTargetHealthSnapshot>;
87
88 async fn shutdown(
89 &self,
90 runtime: &mut TargetRuntimeManager<E>,
91 replay_workers: &mut ReplayWorkerManager,
92 ) -> Result<(), TargetError>;
93}
94
95#[derive(Clone)]
98pub struct BuiltinPluginRuntimeAdapter<E>
99where
100 E: PluginEvent,
101{
102 replay_hook: ReplayHook<E>,
103 replay_start_observer: ReplayStartObserver,
104 replay_semaphore: Option<Arc<Semaphore>>,
105 batch_timeout: Duration,
106 idle_sleep: Duration,
107 stop_log_prefix: Arc<str>,
108}
109
110impl<E> BuiltinPluginRuntimeAdapter<E>
111where
112 E: PluginEvent,
113{
114 pub fn new(
115 replay_hook: ReplayHook<E>,
116 replay_start_observer: ReplayStartObserver,
117 replay_semaphore: Option<Arc<Semaphore>>,
118 batch_timeout: Duration,
119 idle_sleep: Duration,
120 stop_log_prefix: impl Into<Arc<str>>,
121 ) -> Self {
122 Self {
123 replay_hook,
124 replay_start_observer,
125 replay_semaphore,
126 batch_timeout,
127 idle_sleep,
128 stop_log_prefix: stop_log_prefix.into(),
129 }
130 }
131
132 pub async fn prepare_targets(&self, targets: Vec<Box<dyn Target<E> + Send + Sync>>) -> PreparedActivation<E> {
133 self.prepare_targets_inner(targets, None).await
134 }
135
136 pub async fn prepare_targets_cancellable(
137 &self,
138 targets: Vec<Box<dyn Target<E> + Send + Sync>>,
139 cancellation: &CancellationToken,
140 ) -> PreparedActivation<E> {
141 self.prepare_targets_inner(targets, Some(cancellation)).await
142 }
143
144 async fn prepare_targets_inner(
145 &self,
146 targets: Vec<Box<dyn Target<E> + Send + Sync>>,
147 cancellation: Option<&CancellationToken>,
148 ) -> PreparedActivation<E> {
149 let mut prepared = Vec::with_capacity(targets.len());
150 let mut failures = Vec::new();
151 let mut rejected_targets = Vec::new();
152 let mut targets = targets.into_iter();
153 while let Some(target) = targets.next() {
154 match prepare_target(target, cancellation).await {
155 PrepareTargetResult::Ready(target) => prepared.push(target),
156 PrepareTargetResult::Degraded { error, target } => {
157 drop(error);
158 tracing::warn!(
159 target_id = %target.id(),
160 reason = "initialization_failed",
161 "Target initialization failed during lifecycle preparation"
162 );
163 failures.push(TargetActivationFailure {
164 detail: format!("{}: initialization failed", target.id()),
165 });
166 prepared.push(target);
167 }
168 PrepareTargetResult::Failed { error, target } => {
169 drop(error);
170 let target_id = target.id().to_string();
171 tracing::warn!(
172 target_id,
173 reason = "initialization_failed",
174 "Target initialization failed during lifecycle preparation"
175 );
176 failures.push(TargetActivationFailure {
177 detail: format!("{target_id}: initialization failed"),
178 });
179 rejected_targets.push(Arc::from(target));
180 }
181 PrepareTargetResult::Cancelled(target) => {
182 prepared.push(Arc::from(target));
183 prepared.extend(targets.map(Arc::from));
184 break;
185 }
186 }
187 }
188 PreparedActivation {
189 failures,
190 rejected_targets,
191 targets: prepared,
192 }
193 }
194
195 pub fn open_prepared_stores(&self, prepared: PreparedActivation<E>) -> (OpenedActivation<E>, PreparedActivation<E>) {
199 let mut accepted = Vec::with_capacity(prepared.targets.len());
200 let mut failures = prepared.failures;
201 let mut rejected = prepared.rejected_targets;
202 let outcomes = if prepared.targets.len() < 2 {
203 prepared.targets.into_iter().map(open_target_store).collect()
204 } else {
205 match STORE_OPEN_POOL.as_ref() {
206 Ok(pool) => pool.install(|| prepared.targets.into_par_iter().map(open_target_store).collect::<Vec<_>>()),
209 Err(err) => {
210 tracing::warn!(error = %err, "Failed to create target store open pool; opening stores serially");
211 prepared.targets.into_iter().map(open_target_store).collect()
212 }
213 }
214 };
215 for outcome in outcomes {
216 match outcome {
217 StoreOpenOutcome::Accepted(target) => accepted.push(target),
218 StoreOpenOutcome::Rejected { panicked, target } => {
219 if panicked {
220 tracing::error!(
221 target_id = %target.id(),
222 reason = "store_open_panicked",
223 "Target queue store panicked while opening during runtime handoff"
224 );
225 } else {
226 tracing::error!(
227 target_id = %target.id(),
228 reason = "store_open_failed",
229 "Failed to open target queue store during runtime handoff"
230 );
231 }
232 failures.push(TargetActivationFailure {
233 detail: format!("{}: queue store open failed", target.id()),
234 });
235 rejected.push(target);
236 }
237 }
238 }
239 (
240 OpenedActivation { targets: accepted },
241 PreparedActivation {
242 failures,
243 rejected_targets: rejected,
244 targets: Vec::new(),
245 },
246 )
247 }
248
249 pub fn try_activate_prepared(&self, opened: OpenedActivation<E>) -> (RuntimeActivation<E>, PreparedActivation<E>) {
250 let mut replay_workers = ReplayWorkerManager::new();
251 let mut accepted = Vec::with_capacity(opened.targets.len());
252 let mut failures = Vec::new();
253 let mut rejected_targets = Vec::new();
254 for target in opened.targets {
255 let target_id = target.id().to_string();
256 let replay = catch_unwind(AssertUnwindSafe(|| {
257 target.store().filter(|_| target.is_enabled()).map(|store| {
258 start_replay_worker(
259 store.boxed_clone(),
260 Arc::clone(&target),
261 Arc::clone(&self.replay_hook),
262 self.replay_semaphore.clone(),
263 self.batch_timeout,
264 self.idle_sleep,
265 )
266 })
267 }));
268 let replay = match replay {
269 Ok(replay) => replay,
270 Err(_) => {
271 tracing::error!(
272 target_id,
273 reason = "replay_activation_panicked",
274 "Target replay activation panicked during runtime handoff"
275 );
276 failures.push(TargetActivationFailure {
277 detail: format!("{target_id}: replay activation failed"),
278 });
279 rejected_targets.push(target);
280 continue;
281 }
282 };
283 (self.replay_start_observer)(&target_id, replay.is_some());
284 if let Some((cancel_tx, join)) = replay {
285 replay_workers.insert_with_handle(target_id, cancel_tx, join);
286 }
287 accepted.push(target);
288 }
289
290 (
291 RuntimeActivation {
292 replay_workers,
293 targets: accepted,
294 },
295 PreparedActivation {
296 failures,
297 rejected_targets,
298 targets: Vec::new(),
299 },
300 )
301 }
302
303 #[doc(hidden)]
304 pub async fn prepare_dormant_compat_activation(
305 &self,
306 targets: Vec<Box<dyn Target<E> + Send + Sync>>,
307 ) -> RuntimeActivation<E> {
308 let PreparedActivation {
309 failures,
310 rejected_targets,
311 targets,
312 } = self.prepare_targets(targets).await;
313 let rejected = PreparedActivation {
314 failures,
315 rejected_targets,
316 targets: Vec::new(),
317 };
318 if let Err(err) = self.close_prepared(rejected).await {
319 tracing::warn!(error = %err, "Failed to close targets rejected while preparing compatibility activation");
320 }
321 RuntimeActivation {
322 replay_workers: ReplayWorkerManager::new(),
323 targets,
324 }
325 }
326
327 #[doc(hidden)]
328 pub fn start_dormant_compat_activation(
329 &self,
330 activation: RuntimeActivation<E>,
331 ) -> (RuntimeActivation<E>, PreparedActivation<E>, PreparedActivation<E>) {
332 let prepared = PreparedActivation {
333 failures: Vec::new(),
334 rejected_targets: Vec::new(),
335 targets: activation.targets,
336 };
337 let (opened, open_rejected) = self.open_prepared_stores(prepared);
338 let (activation, activation_rejected) = self.try_activate_prepared(opened);
339 (activation, open_rejected, activation_rejected)
340 }
341
342 #[doc(hidden)]
343 pub async fn close_compat_activation(&self, mut activation: RuntimeActivation<E>) -> Result<(), TargetError> {
344 let mut runtime = TargetRuntimeManager::new();
345 for target in activation.targets {
346 runtime.add_arc(target);
347 }
348 self.shutdown(&mut runtime, &mut activation.replay_workers).await
349 }
350
351 pub async fn close_prepared(&self, prepared: PreparedActivation<E>) -> Result<(), TargetError> {
352 let mut runtime = TargetRuntimeManager::new();
353 for target in prepared.targets.into_iter().chain(prepared.rejected_targets) {
354 runtime.add_arc(target);
355 }
356 let mut replay_workers = ReplayWorkerManager::new();
357 self.shutdown(&mut runtime, &mut replay_workers).await
358 }
359}
360
361#[async_trait]
362impl<E> PluginRuntimeAdapter<E> for BuiltinPluginRuntimeAdapter<E>
363where
364 E: PluginEvent,
365{
366 async fn activate_with_replay(&self, targets: Vec<Box<dyn Target<E> + Send + Sync>>) -> RuntimeActivation<E> {
367 let prepared = self.prepare_targets(targets).await;
368 let (opened, rejected) = self.open_prepared_stores(prepared);
369 if let Err(err) = self.close_prepared(rejected).await {
370 tracing::warn!(error = %err, "Failed to close targets whose queue stores could not be opened");
371 }
372 let (activation, rejected) = self.try_activate_prepared(opened);
373 if let Err(err) = self.close_prepared(rejected).await {
374 tracing::warn!(error = %err, "Failed to close targets rejected during replay activation");
375 }
376 activation
377 }
378
379 async fn replace_runtime_targets(
380 &self,
381 runtime: &mut TargetRuntimeManager<E>,
382 replay_workers: &mut ReplayWorkerManager,
383 activation: RuntimeActivation<E>,
384 ) -> Result<(), TargetError> {
385 self.stop_replay_workers(replay_workers).await;
390 let close_errors = runtime.clear_and_close().await;
391 if !close_errors.is_empty() {
392 tracing::warn!(failed_targets = close_errors.len(), "Some targets failed to close during runtime reload");
393 }
394
395 for target in activation.targets {
396 runtime.add_arc(target);
397 }
398
399 *replay_workers = activation.replay_workers;
400 Ok(())
401 }
402
403 async fn stop_replay_workers(&self, replay_workers: &mut ReplayWorkerManager) {
404 replay_workers.stop_all(&self.stop_log_prefix).await;
405 }
406
407 fn snapshot_runtime_status(
408 &self,
409 runtime: &TargetRuntimeManager<E>,
410 replay_workers: &ReplayWorkerManager,
411 ) -> RuntimeStatusSnapshot {
412 runtime.status_snapshot(replay_workers)
413 }
414
415 async fn snapshot_runtime_health(&self, runtime: &TargetRuntimeManager<E>) -> Vec<RuntimeTargetHealthSnapshot> {
416 runtime.health_snapshots().await
417 }
418
419 async fn shutdown(
420 &self,
421 runtime: &mut TargetRuntimeManager<E>,
422 replay_workers: &mut ReplayWorkerManager,
423 ) -> Result<(), TargetError> {
424 self.stop_replay_workers(replay_workers).await;
428 let close_errors = runtime.clear_and_close().await;
429 if !close_errors.is_empty() {
430 let detail = close_errors
431 .into_iter()
432 .map(|(target_id, _)| target_id)
433 .collect::<Vec<_>>()
434 .join("; ");
435 return Err(TargetError::Storage(format!("Failed to close {detail}")));
436 }
437 Ok(())
438 }
439}
440
441#[cfg(test)]
442mod tests {
443 use super::{BuiltinPluginRuntimeAdapter, MAX_PARALLEL_STORE_OPENS, PluginRuntimeAdapter};
444 use crate::store::{Key, QueueStore, Store};
445 use crate::target::QueuedPayload;
446 use crate::testkit::MockTarget;
447 use crate::{StoreError, Target};
448 use std::sync::{Arc, Condvar, Mutex};
449 use std::time::Duration;
450 use tempfile::tempdir;
451 use tokio::sync::Notify;
452 use tokio_util::sync::CancellationToken;
453
454 type TestStore = dyn Store<QueuedPayload, Error = StoreError, Key = Key> + Send + Sync;
455
456 type BeforeOpen = Arc<dyn Fn() + Send + Sync>;
457
458 #[derive(Clone)]
459 struct TestOpenStore {
460 before_clone: BeforeOpen,
461 before_open: BeforeOpen,
462 store: QueueStore<QueuedPayload>,
463 }
464
465 impl Store<QueuedPayload> for TestOpenStore {
466 type Error = StoreError;
467 type Key = Key;
468
469 fn open(&self) -> Result<(), Self::Error> {
470 (self.before_open)();
471 self.store.open()
472 }
473
474 fn put(&self, item: Arc<QueuedPayload>) -> Result<Self::Key, Self::Error> {
475 self.store.put(item)
476 }
477
478 fn put_multiple(&self, items: Vec<QueuedPayload>) -> Result<Self::Key, Self::Error> {
479 self.store.put_multiple(items)
480 }
481
482 fn put_raw(&self, data: &[u8]) -> Result<Self::Key, Self::Error> {
483 self.store.put_raw(data)
484 }
485
486 fn get(&self, key: &Self::Key) -> Result<QueuedPayload, Self::Error> {
487 self.store.get(key)
488 }
489
490 fn get_multiple(&self, key: &Self::Key) -> Result<Vec<QueuedPayload>, Self::Error> {
491 self.store.get_multiple(key)
492 }
493
494 fn get_raw(&self, key: &Self::Key) -> Result<Vec<u8>, Self::Error> {
495 self.store.get_raw(key)
496 }
497
498 fn del(&self, key: &Self::Key) -> Result<(), Self::Error> {
499 self.store.del(key)
500 }
501
502 fn delete(&self) -> Result<(), Self::Error> {
503 self.store.delete()
504 }
505
506 fn list(&self) -> Vec<Self::Key> {
507 self.store.list()
508 }
509
510 fn len(&self) -> usize {
511 self.store.len()
512 }
513
514 fn is_empty(&self) -> bool {
515 self.store.is_empty()
516 }
517
518 fn boxed_clone(&self) -> Box<dyn Store<QueuedPayload, Error = Self::Error, Key = Self::Key> + Send + Sync> {
519 (self.before_clone)();
520 Box::new(self.clone())
521 }
522 }
523
524 #[derive(Default)]
525 struct StoreOpenGate {
526 changed: Condvar,
527 state: Mutex<StoreOpenGateState>,
528 }
529
530 #[derive(Default)]
531 struct StoreOpenGateState {
532 active: usize,
533 max_active: usize,
534 released: bool,
535 }
536
537 impl StoreOpenGate {
538 fn enter(&self) {
539 let mut state = self.state.lock().unwrap_or_else(|err| err.into_inner());
540 state.active += 1;
541 state.max_active = state.max_active.max(state.active);
542 self.changed.notify_all();
543 while !state.released {
544 state = self.changed.wait(state).unwrap_or_else(|err| err.into_inner());
545 }
546 state.active -= 1;
547 }
548
549 fn wait_for_active(&self, expected: usize, timeout: Duration) -> bool {
550 let state = self.state.lock().unwrap_or_else(|err| err.into_inner());
551 let (state, _) = self
552 .changed
553 .wait_timeout_while(state, timeout, |state| state.max_active < expected)
554 .unwrap_or_else(|err| err.into_inner());
555 state.max_active >= expected
556 }
557
558 fn release(&self) {
559 let mut state = self.state.lock().unwrap_or_else(|err| err.into_inner());
560 state.released = true;
561 self.changed.notify_all();
562 }
563
564 fn max_active(&self) -> usize {
565 self.state.lock().unwrap_or_else(|err| err.into_inner()).max_active
566 }
567 }
568
569 fn opened_queue_store() -> Arc<TestStore> {
573 let dir = tempdir().expect("tempdir should be created for queue store tests");
574 let store = QueueStore::<QueuedPayload>::new(dir.path(), 16, ".queue");
575 store.open().expect("queue store should open");
576 Arc::new(store)
577 }
578
579 fn builtin_adapter() -> BuiltinPluginRuntimeAdapter<String> {
580 BuiltinPluginRuntimeAdapter::new(
581 Arc::new(|_event| Box::pin(async {})),
582 Arc::new(|_target_id, _has_replay| {}),
583 None,
584 Duration::from_millis(10),
585 Duration::from_millis(10),
586 "stopping test replay worker",
587 )
588 }
589
590 #[tokio::test]
591 async fn builtin_adapter_handles_empty_target_activation() {
592 let adapter = builtin_adapter();
593 let activation = adapter.activate_with_replay(Vec::new()).await;
594
595 assert!(activation.targets.is_empty());
596 assert!(activation.replay_workers.is_empty());
597 }
598
599 #[tokio::test]
600 async fn builtin_adapter_skips_non_store_target_when_init_fails() {
601 let adapter = builtin_adapter();
602 let target = MockTarget::new("primary", "webhook").with_init_failures(usize::MAX);
603
604 let activation = adapter.activate_with_replay(vec![Box::new(target)]).await;
605
606 assert!(activation.targets.is_empty());
607 assert!(activation.replay_workers.is_empty());
608 }
609
610 #[tokio::test]
611 async fn builtin_adapter_keeps_store_backed_target_when_init_fails() {
612 let adapter = builtin_adapter();
613 let target = MockTarget::new("primary", "webhook")
614 .with_init_failures(usize::MAX)
615 .with_store(opened_queue_store());
616
617 let activation = adapter.activate_with_replay(vec![Box::new(target)]).await;
618
619 assert_eq!(activation.targets.len(), 1);
620 assert_eq!(activation.replay_workers.len(), 1);
621 }
622
623 #[tokio::test]
624 async fn prepared_store_target_reports_init_failure_without_dropping_queue_runtime() {
625 let adapter = builtin_adapter();
626 let target = MockTarget::new("primary", "webhook")
627 .with_init_failures(usize::MAX)
628 .with_store(opened_queue_store());
629
630 let prepared = adapter.prepare_targets(vec![Box::new(target)]).await;
631 assert_eq!(prepared.targets.len(), 1);
632 assert!(
633 prepared
634 .failure_summary()
635 .is_some_and(|summary| summary.contains("initialization failed") && !summary.contains("forced init failure"))
636 );
637
638 let (opened, rejected) = adapter.open_prepared_stores(prepared);
639 assert!(rejected.failure_summary().is_some());
640 let (mut activation, activation_rejected) = adapter.try_activate_prepared(opened);
641 assert!(activation_rejected.failure_summary().is_none());
642 assert_eq!(activation.targets.len(), 1);
643 assert_eq!(activation.replay_workers.len(), 1);
644 activation.replay_workers.stop_all("stop degraded target replay worker").await;
645 }
646
647 #[tokio::test]
648 async fn cancellable_preparation_returns_current_and_remaining_targets_for_shutdown() {
649 let adapter = builtin_adapter();
650 let init_entered = Arc::new(Notify::new());
651 let first = MockTarget::new("first", "webhook").with_blocking_init(init_entered.clone());
652 let first_observer = first.clone();
653 let second = MockTarget::new("second", "webhook");
654 let second_observer = second.clone();
655 let cancellation = CancellationToken::new();
656 let prepare_adapter = adapter.clone();
657 let prepare_cancellation = cancellation.clone();
658 let prepare = tokio::spawn(async move {
659 prepare_adapter
660 .prepare_targets_cancellable(vec![Box::new(first), Box::new(second)], &prepare_cancellation)
661 .await
662 });
663
664 init_entered.notified().await;
665 cancellation.cancel();
666 let prepared = tokio::time::timeout(Duration::from_secs(1), prepare)
667 .await
668 .expect("cancellation should interrupt target initialization")
669 .expect("preparation task should finish");
670 assert_eq!(prepared.targets.len(), 2);
671
672 adapter
673 .close_prepared(prepared)
674 .await
675 .expect("cancelled targets should close");
676 assert_eq!(first_observer.close_call_count(), 1);
677 assert_eq!(second_observer.close_call_count(), 1);
678 assert_eq!(second_observer.init_call_count(), 0);
679 }
680
681 #[tokio::test]
682 async fn prepared_activation_opens_store_before_starting_replay() {
683 let adapter = builtin_adapter();
684 let dir = tempdir().expect("tempdir should be created");
685 let queue_path = dir.path().join("queue");
686 let target = MockTarget::new("primary", "webhook").with_store(Arc::new(QueueStore::<QueuedPayload>::new(
687 &queue_path,
688 16,
689 ".queue",
690 )));
691
692 let prepared = adapter.prepare_targets(vec![Box::new(target)]).await;
693 assert!(!queue_path.exists(), "dormant preparation must not open the queue store");
694
695 let (opened, rejected) = adapter.open_prepared_stores(prepared);
696 assert!(rejected.targets.is_empty());
697 assert!(queue_path.is_dir(), "handoff must open the queue store before activation");
698 let (mut activation, activation_rejected) = adapter.try_activate_prepared(opened);
699 assert!(activation_rejected.failure_summary().is_none());
700 assert_eq!(activation.replay_workers.len(), 1);
701 activation
702 .replay_workers
703 .stop_all("stop prepared activation test worker")
704 .await;
705 }
706
707 #[tokio::test]
708 async fn activation_closes_a_target_when_its_queue_store_cannot_open() {
709 let adapter = builtin_adapter();
710 let dir = tempdir().expect("tempdir should be created");
711 let invalid_base = dir.path().join("not-a-directory");
712 std::fs::write(&invalid_base, b"file").expect("invalid queue base should be created");
713 let target = MockTarget::new("primary", "webhook").with_store(Arc::new(QueueStore::<QueuedPayload>::new(
714 &invalid_base,
715 16,
716 ".queue",
717 )));
718 let observer = target.clone();
719
720 let activation = adapter.activate_with_replay(vec![Box::new(target)]).await;
721
722 assert!(activation.targets.is_empty());
723 assert!(activation.replay_workers.is_empty());
724 assert_eq!(observer.close_call_count(), 1);
725 }
726
727 #[tokio::test]
728 async fn prepared_stores_open_with_bounded_parallelism_and_stable_order() {
729 const TARGETS: usize = MAX_PARALLEL_STORE_OPENS * 2;
730
731 let adapter = builtin_adapter();
732 let dir = tempdir().expect("tempdir should be created");
733 let gate = Arc::new(StoreOpenGate::default());
734 let mut targets: Vec<Box<dyn Target<String> + Send + Sync>> = Vec::with_capacity(TARGETS);
735 let mut expected_ids = Vec::with_capacity(TARGETS);
736 for index in 0..TARGETS {
737 let open_gate = gate.clone();
738 let target = MockTarget::new(&format!("target-{index}"), "webhook").with_store(Arc::new(TestOpenStore {
739 before_clone: Arc::new(|| {}),
740 before_open: Arc::new(move || open_gate.enter()),
741 store: QueueStore::new(dir.path().join(index.to_string()), 16, ".queue"),
742 }));
743 expected_ids.push(target.target_id().to_string());
744 targets.push(Box::new(target));
745 }
746
747 let prepared = adapter.prepare_targets(targets).await;
748 let open_adapter = adapter.clone();
749 let opening = tokio::task::spawn_blocking(move || open_adapter.open_prepared_stores(prepared));
750 let wait_gate = gate.clone();
751 let reached_bound =
752 tokio::task::spawn_blocking(move || wait_gate.wait_for_active(MAX_PARALLEL_STORE_OPENS, Duration::from_secs(30)))
753 .await
754 .expect("store-open observer should not panic");
755 gate.release();
756 let (opened, rejected) = opening.await.expect("bounded store opens should not panic");
757 let opened_ids = opened
758 .targets
759 .iter()
760 .map(|target| target.id().to_string())
761 .collect::<Vec<_>>();
762
763 assert!(reached_bound, "store opens did not use the configured parallelism");
764 assert_eq!(gate.max_active(), MAX_PARALLEL_STORE_OPENS);
765 assert_eq!(opened_ids, expected_ids, "parallel store opens must preserve configuration order");
766 assert!(rejected.targets.is_empty());
767 assert!(rejected.failure_summary().is_none());
768 }
769
770 #[tokio::test]
771 async fn panicking_store_open_rejects_and_closes_only_that_target() {
772 let adapter = builtin_adapter();
773 let dir = tempdir().expect("tempdir should be created");
774 let target = MockTarget::new("panicking", "webhook").with_store(Arc::new(TestOpenStore {
775 before_clone: Arc::new(|| {}),
776 before_open: Arc::new(|| panic!("forced store open panic: do-not-expose-payload")),
777 store: QueueStore::new(dir.path(), 16, ".queue"),
778 }));
779 let observer = target.clone();
780
781 let prepared = adapter.prepare_targets(vec![Box::new(target)]).await;
782 let (opened, rejected) = adapter.open_prepared_stores(prepared);
783 let summary = rejected
784 .failure_summary()
785 .expect("panicking store should report a generic activation failure");
786 let (activation, activation_rejected) = adapter.try_activate_prepared(opened);
787
788 assert!(activation.targets.is_empty(), "a target without an open store must not become visible");
789 assert!(
790 activation.replay_workers.is_empty(),
791 "a rejected target must not publish without a replay worker"
792 );
793 assert!(activation_rejected.failure_summary().is_none());
794 assert!(summary.contains("queue store open failed"));
795 assert!(!summary.contains("do-not-expose-payload"));
796 adapter
797 .close_prepared(rejected)
798 .await
799 .expect("a target rejected after a store panic should close");
800 assert_eq!(observer.close_call_count(), 1);
801 }
802
803 #[tokio::test]
804 async fn panicking_store_clone_cannot_publish_target_without_replay_worker() {
805 let adapter = builtin_adapter();
806 let dir = tempdir().expect("tempdir should be created");
807 let target = MockTarget::new("panicking-clone", "webhook").with_store(Arc::new(TestOpenStore {
808 before_clone: Arc::new(|| panic!("forced store clone panic: do-not-expose-payload")),
809 before_open: Arc::new(|| {}),
810 store: QueueStore::new(dir.path(), 16, ".queue"),
811 }));
812 let observer = target.clone();
813
814 let prepared = adapter.prepare_targets(vec![Box::new(target)]).await;
815 let (opened, open_rejected) = adapter.open_prepared_stores(prepared);
816 assert!(open_rejected.failure_summary().is_none());
817 let (activation, rejected) = adapter.try_activate_prepared(opened);
818 let summary = rejected
819 .failure_summary()
820 .expect("panicking store clone should report a generic activation failure");
821
822 assert!(activation.targets.is_empty(), "a target without a replay worker must not become visible");
823 assert!(activation.replay_workers.is_empty());
824 assert!(summary.contains("replay activation failed"));
825 assert!(!summary.contains("do-not-expose-payload"));
826 adapter
827 .close_prepared(rejected)
828 .await
829 .expect("a target rejected during replay activation should close");
830 assert_eq!(observer.close_call_count(), 1);
831 }
832
833 #[tokio::test]
834 async fn builtin_adapter_shutdown_clears_runtime_and_replay_workers() {
835 let adapter = builtin_adapter();
836 let target = MockTarget::new("primary", "webhook");
837 let observer = target.clone();
838 let mut runtime = crate::runtime::TargetRuntimeManager::new();
839 let mut replay_workers = crate::runtime::ReplayWorkerManager::new();
840
841 let activation = adapter.activate_with_replay(vec![Box::new(target)]).await;
842 adapter
843 .replace_runtime_targets(&mut runtime, &mut replay_workers, activation)
844 .await
845 .expect("replace_runtime_targets should succeed");
846
847 assert_eq!(runtime.len(), 1);
848 assert_eq!(replay_workers.len(), 0);
849
850 adapter
851 .shutdown(&mut runtime, &mut replay_workers)
852 .await
853 .expect("shutdown should succeed");
854
855 assert!(runtime.is_empty());
856 assert!(replay_workers.is_empty());
857 assert_eq!(observer.close_call_count(), 1);
858 }
859}