Skip to main content

photon_backend/
executor_services.rs

1//! Shared runtime services for handler delivery (constructed at Photon build time).
2
3use std::sync::Arc;
4
5use crate::checkpoint::CheckpointCoalescer;
6use crate::delivery::{DlqSink, WorkerPool};
7use crate::retention::{
8    PartitionReclaim, RetentionDeps, RetentionHook, RetentionPolicy, RetentionReclaimer,
9};
10use crate::storage::StoragePort;
11
12/// Shared runtime services for handler delivery.
13pub struct ExecutorServices {
14    /// Bounded handler concurrency pool.
15    pub worker_pool: Arc<WorkerPool>,
16    /// Coalesced checkpoint writer.
17    pub checkpoint_coalescer: Arc<CheckpointCoalescer>,
18    /// Dead-letter sink for failed deliveries.
19    pub dlq: Arc<DlqSink>,
20    /// Storage retention reclaim coordinator.
21    pub retention_reclaimer: Arc<RetentionReclaimer>,
22}
23
24impl ExecutorServices {
25    /// Construct delivery services bound to the storage port.
26    #[allow(clippy::needless_pass_by_value)] // Arc-by-value is the public ownership API
27    pub fn new(
28        port: Arc<dyn StoragePort>,
29        policy: RetentionPolicy,
30        hook: Option<Arc<dyn RetentionHook>>,
31    ) -> Self {
32        let coalescer = Arc::new(CheckpointCoalescer::new(Arc::clone(&port)));
33        let dlq = Arc::new(DlqSink::new());
34        let reclaimer = RetentionReclaimer::new(RetentionDeps {
35            port: Arc::clone(&port),
36            coalescer: Arc::clone(&coalescer),
37            dlq: Arc::clone(&dlq),
38            policy,
39            hook,
40        });
41        coalescer.attach_reclaimer(Arc::clone(&reclaimer) as Arc<dyn PartitionReclaim>);
42        Self {
43            worker_pool: Arc::new(WorkerPool::from_env()),
44            checkpoint_coalescer: coalescer,
45            dlq,
46            retention_reclaimer: reclaimer,
47        }
48    }
49}