Skip to main content

rustfs_targets/runtime/
adapter.rs

1// Copyright 2024 RustFS Team
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 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/// Shared runtime contract for target plugins.
64#[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/// Built-in in-process runtime adapter that preserves the current replay and
96/// activation behavior while presenting a stable runtime contract to callers.
97#[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    /// Opens queue stores only after the previous runtime generation has been
196    /// quiesced. Targets whose stores cannot be opened retain the established
197    /// fault-isolation behavior and are returned for lock-free shutdown.
198    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                // Vec's indexed parallel iterator preserves configuration
207                // order in collect, keeping failure summaries deterministic.
208                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        // Stop (and join) the old replay workers before installing the new set so
386        // no two workers ever drain the same store concurrently, then close the
387        // old targets. A close failure during reload is logged but does not abort
388        // the reload — the new configuration must still take effect.
389        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        // On explicit shutdown, propagate any close/flush failures instead of
425        // swallowing them: the runtime is still fully torn down, but the caller
426        // learns that a target could not be flushed/closed cleanly.
427        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    /// Builds the tempdir-backed, already-opened queue store the store-backed mock targets use.
570    /// The tempdir handle is dropped here on purpose, matching the previous in-module mock: these
571    /// tests never write through the store, they only need an openable handle.
572    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}