Skip to main content

fynd_core/derived/
manager.rs

1//! Computation manager for derived data.
2//!
3//! The ComputationManager:
4//! - Subscribes to MarketEvents from TychoFeed
5//! - Runs derived computations (spot prices, token prices, component depths) and updates the
6//!   DerivedData store
7//! - Provides read access to workers via shared store reference
8
9use std::{
10    sync::Arc,
11    time::{Duration, Instant, SystemTime, UNIX_EPOCH},
12};
13
14use async_trait::async_trait;
15use futures::future::join_all;
16use metrics::{counter, gauge, histogram};
17use rustc_hash::{FxHashMap, FxHashSet};
18use tokio::sync::{broadcast, RwLock};
19use tracing::{error, info, trace, warn};
20use tycho_simulation::tycho_common::models::Address;
21
22use crate::types::ComponentId;
23
24/// Information about which components changed in a market update.
25///
26/// Used to enable incremental computation - only recomputing derived data
27/// for components that actually changed.
28#[derive(Debug, Clone, Default)]
29pub struct ChangedComponents {
30    /// Newly added components with their token addresses.
31    pub added: FxHashMap<ComponentId, Vec<Address>>,
32    /// Components that were removed.
33    pub removed: Vec<ComponentId>,
34    /// Components whose state was updated (but not added/removed).
35    pub updated: Vec<ComponentId>,
36    /// If true, this represents a full recompute (startup/lag recovery).
37    pub is_full_recompute: bool,
38}
39
40impl ChangedComponents {
41    /// Returns a HashSet of all changed component IDs.
42    pub fn all_changed_ids(&self) -> FxHashSet<ComponentId> {
43        let mut all = FxHashSet::default();
44        all.extend(self.added.keys().cloned());
45        all.extend(self.removed.iter().cloned());
46        all.extend(self.updated.iter().cloned());
47        all
48    }
49}
50
51/// Coalesces a drained batch of [`MarketEvent`]s into a single incremental
52/// [`ChangedComponents`], applying net semantics: a component that is added
53/// then removed within the batch nets to removed; an add supersedes a prior
54/// update; a remove supersedes a prior add/update.
55///
56/// Returns `None` when the batch carries no net changes. The result always has
57/// `is_full_recompute: false` — this is the bounded lag-recovery path, never a
58/// whole-topology recompute.
59fn coalesce_market_events(events: &[MarketEvent]) -> Option<ChangedComponents> {
60    let mut added: FxHashMap<ComponentId, Vec<Address>> = FxHashMap::default();
61    let mut removed: FxHashSet<ComponentId> = FxHashSet::default();
62    let mut updated: FxHashSet<ComponentId> = FxHashSet::default();
63
64    for event in events {
65        match event {
66            MarketEvent::MarketUpdated {
67                added_components,
68                removed_components,
69                updated_components,
70            } => {
71                for (id, tokens) in added_components {
72                    removed.remove(id);
73                    updated.remove(id);
74                    added.insert(id.clone(), tokens.clone());
75                }
76                for id in removed_components {
77                    added.remove(id);
78                    updated.remove(id);
79                    removed.insert(id.clone());
80                }
81                for id in updated_components {
82                    if !added.contains_key(id) && !removed.contains(id) {
83                        updated.insert(id.clone());
84                    }
85                }
86            }
87        }
88    }
89
90    if added.is_empty() && removed.is_empty() && updated.is_empty() {
91        return None;
92    }
93    Some(ChangedComponents {
94        added,
95        removed: removed.into_iter().collect(),
96        updated: updated.into_iter().collect(),
97        is_full_recompute: false,
98    })
99}
100
101use super::{
102    computation::{ComputationId, ComputationRequirements, DerivedComputation},
103    computations::{ComponentDepthComputation, SpotPriceComputation, TokenGasPriceComputation},
104    error::ComputationError,
105    events::DerivedDataEvent,
106    registry::ErasedComputation,
107    store::DerivedData,
108};
109use crate::feed::{
110    events::{EventError, MarketEvent, MarketEventHandler},
111    market_data::MarketData,
112};
113
114/// Thread-safe handle to shared derived data store.
115pub type SharedDerivedDataRef = Arc<RwLock<DerivedData>>;
116
117/// Configuration for the default computation set built by [`ComputationManager::new`].
118#[derive(Debug, Clone)]
119pub struct ComputationManagerConfig {
120    /// Gas token address (e.g., WETH) for token price computation.
121    gas_token: Address,
122    /// Max hop count for token gas price computation.
123    max_hop: usize,
124    /// Slippage threshold for component depth computation (0.0 < threshold < 1.0).
125    depth_slippage_threshold: f64,
126    /// Overrides the token pricing pass's sell-loop budget; `None` keeps the computation's
127    /// default. The replay harness sets an effectively unbounded budget so integration tests
128    /// can assert exact priced-token counts.
129    pricing_pass_budget: Option<Duration>,
130    /// Overrides how many tokens one token-pricing pass may attempt; `None` keeps the
131    /// computation's default.
132    pricing_max_tokens_per_pass: Option<usize>,
133    /// Overrides how long after a token-pricing pass starts the next one may start; `None`
134    /// keeps the computation's default.
135    pricing_min_pass_interval: Option<Duration>,
136}
137
138impl ComputationManagerConfig {
139    /// Creates a new configuration with the given gas token.
140    pub fn new() -> Self {
141        Self::default()
142    }
143
144    /// Sets the slippage threshold for component depth computation.
145    pub fn with_depth_slippage_threshold(mut self, threshold: f64) -> Self {
146        self.depth_slippage_threshold = threshold;
147        self
148    }
149
150    /// Sets the max hop count for token gas price computation.
151    pub fn with_max_hop(mut self, hop_count: usize) -> Self {
152        self.max_hop = hop_count;
153        self
154    }
155
156    /// Overrides the wall-clock budget for the token pricing pass's sell loop.
157    pub fn with_pricing_pass_budget(mut self, pass_budget: Duration) -> Self {
158        self.pricing_pass_budget = Some(pass_budget);
159        self
160    }
161
162    /// Overrides how many tokens one token-pricing pass may attempt.
163    pub fn with_pricing_max_tokens_per_pass(mut self, max_tokens: usize) -> Self {
164        self.pricing_max_tokens_per_pass = Some(max_tokens);
165        self
166    }
167
168    /// Overrides how long after a token-pricing pass starts the next one may start.
169    pub fn with_pricing_min_pass_interval(mut self, interval: Duration) -> Self {
170        self.pricing_min_pass_interval = Some(interval);
171        self
172    }
173
174    /// Sets the gas token address.
175    pub fn with_gas_token(mut self, gas_token: Address) -> Self {
176        self.gas_token = gas_token;
177        self
178    }
179
180    /// Returns the gas token address.
181    pub fn gas_token(&self) -> &Address {
182        &self.gas_token
183    }
184
185    /// Returns the max hop count.
186    pub fn max_hop(&self) -> usize {
187        self.max_hop
188    }
189
190    /// Returns the depth slippage threshold.
191    pub fn depth_slippage_threshold(&self) -> f64 {
192        self.depth_slippage_threshold
193    }
194
195    /// Builds the token price computation this configuration describes.
196    pub(crate) fn build_token_price_computation(&self) -> TokenGasPriceComputation {
197        let mut token_prices = TokenGasPriceComputation::default()
198            .with_max_hops(self.max_hop)
199            .with_gas_token(self.gas_token.clone());
200        if let Some(pass_budget) = self.pricing_pass_budget {
201            token_prices = token_prices.with_pass_budget(pass_budget);
202        }
203        if let Some(max_tokens) = self.pricing_max_tokens_per_pass {
204            token_prices = token_prices.with_max_tokens_per_pass(max_tokens);
205        }
206        if let Some(interval) = self.pricing_min_pass_interval {
207            token_prices = token_prices.with_min_pass_interval(interval);
208        }
209        token_prices
210    }
211}
212
213impl Default for ComputationManagerConfig {
214    fn default() -> Self {
215        // FyndBuilder overrides max_hop with the deepest configured pool's max_hops, so that
216        // every quotable token is priceable; the default only serves manual construction and
217        // matches the default pool max_hops.
218        Self {
219            gas_token: Address::zero(20),
220            max_hop: crate::solver::defaults::POOL_MAX_HOPS,
221            depth_slippage_threshold: 0.01,
222            pricing_pass_budget: None,
223            pricing_max_tokens_per_pass: None,
224            pricing_min_pass_interval: None,
225        }
226    }
227}
228
229/// Manages derived data computations triggered by market events.
230pub struct ComputationManager {
231    /// Reference to shared market data (read access).
232    market_data: MarketData,
233    /// Shared derived data store (write access).
234    store: SharedDerivedDataRef,
235    /// Registered computations, driven in dependency-stage order each block.
236    computations: Vec<Box<dyn ErasedComputation>>,
237    /// Event broadcaster for derived data updates.
238    event_tx: broadcast::Sender<DerivedDataEvent>,
239}
240
241/// A dependency-ordered execution plan for the registered computations.
242struct ComputationSchedule {
243    /// Indices into `ComputationManager::computations`, grouped into stages run in order.
244    stages: Vec<Vec<usize>>,
245    /// Indices that could not be ordered because of a requirement cycle.
246    unscheduled: Vec<usize>,
247}
248
249impl ComputationManager {
250    /// Creates a manager with the default computation set: spot prices, token prices, and
251    /// component depths, run for every block by [`run`](Self::run).
252    ///
253    /// Returns the manager and a receiver for derived data events.
254    /// Workers can subscribe to the event sender via `event_sender()` to track
255    /// computation readiness.
256    pub fn new(
257        config: ComputationManagerConfig,
258        market_data: MarketData,
259    ) -> Result<(Self, broadcast::Receiver<DerivedDataEvent>), ComputationError> {
260        let (mut manager, event_rx) = Self::empty(market_data);
261        manager.register(SpotPriceComputation::new())?;
262        manager.register(config.build_token_price_computation())?;
263        manager.register(ComponentDepthComputation::new(config.depth_slippage_threshold())?)?;
264        Ok((manager, event_rx))
265    }
266
267    /// Creates a manager with no computations registered.
268    ///
269    /// [`new`](Self::new) builds on this to assemble the default computation set, and
270    /// tests drive a custom set through [`register`](Self::register).
271    pub(crate) fn empty(market_data: MarketData) -> (Self, broadcast::Receiver<DerivedDataEvent>) {
272        let (event_tx, event_rx) = broadcast::channel(64);
273        (
274            Self {
275                market_data,
276                store: DerivedData::new_shared(),
277                computations: Vec::new(),
278                event_tx,
279            },
280            event_rx,
281        )
282    }
283
284    /// Registers a computation to be driven each block.
285    ///
286    /// Registration order is preserved within a dependency stage; cross-stage order is
287    /// derived from each computation's
288    /// [`requirements`](crate::derived::computation::DerivedComputation::requirements).
289    ///
290    /// # Errors
291    ///
292    /// Returns [`ComputationError::DuplicateComputationId`] if a computation with the same
293    /// [`ID`](DerivedComputation::ID) is already registered.
294    pub(crate) fn register<C: DerivedComputation>(
295        &mut self,
296        computation: C,
297    ) -> Result<(), ComputationError> {
298        if self
299            .computations
300            .iter()
301            .any(|existing| existing.id() == C::ID)
302        {
303            return Err(ComputationError::DuplicateComputationId(C::ID));
304        }
305        self.computations
306            .push(Box::new(computation));
307        Ok(())
308    }
309
310    /// Returns a reference to the shared derived data store.
311    pub fn store(&self) -> SharedDerivedDataRef {
312        Arc::clone(&self.store)
313    }
314
315    /// Returns the event sender for workers to subscribe.
316    pub fn event_sender(&self) -> broadcast::Sender<DerivedDataEvent> {
317        self.event_tx.clone()
318    }
319
320    /// Runs the main loop until shutdown or channel close.
321    ///
322    /// **Note:** Consumes `self`. Call [`store()`](Self::store) before `run()` to retain access.
323    pub async fn run(
324        mut self,
325        mut event_rx: broadcast::Receiver<MarketEvent>,
326        mut shutdown_rx: broadcast::Receiver<()>,
327    ) {
328        info!("computation manager started");
329
330        loop {
331            tokio::select! {
332                biased;
333
334                _ = shutdown_rx.recv() => {
335                    info!("computation manager shutting down");
336                    break;
337                }
338
339                event_result = event_rx.recv() => {
340                    match event_result {
341                        Ok(event) => {
342                            if let Err(e) = self.handle_event(&event).await {
343                                warn!(error = ?e, "failed to handle market event");
344                            }
345                        }
346                        Err(broadcast::error::RecvError::Closed) => {
347                            info!("event channel closed, computation manager shutting down");
348                            break;
349                        }
350                        Err(broadcast::error::RecvError::Lagged(skipped)) => {
351                            warn!(
352                                skipped,
353                                "computation manager lagged; draining buffered events and \
354                                 recomputing changed components incrementally"
355                            );
356                            counter!("derived_manager_lag_recoveries_total").increment(1);
357                            counter!("derived_manager_lagged_events_total")
358                                .increment(skipped);
359                            self.recover_from_lag(&mut event_rx).await;
360                        }
361                    }
362                }
363            }
364        }
365    }
366
367    /// Runs all registered computations for the current block and updates the store.
368    ///
369    /// Computations run in dependency stages derived from their
370    /// [`requirements`](crate::derived::computation::DerivedComputation::requirements):
371    /// a stage runs concurrently and is written before the next stage starts, and a
372    /// computation whose requirement did not succeed this block is skipped and reported
373    /// as failed. Broadcasts a `DerivedDataEvent` per computation.
374    async fn compute_all(&self, changed: &ChangedComponents) {
375        let total_start = Instant::now();
376
377        // Get block info for tracking
378        let Some(block) = self
379            .market_data
380            .read()
381            .await
382            .last_updated()
383            .map(|b| b.number())
384        else {
385            warn!("market data has no last updated block, skipping computations");
386            return;
387        };
388
389        // Broadcast new block event
390        let _ = self
391            .event_tx
392            .send(DerivedDataEvent::NewBlock { block });
393
394        let nodes: Vec<(ComputationId, ComputationRequirements)> = self
395            .computations
396            .iter()
397            .map(|computation| (computation.id(), computation.requirements()))
398            .collect();
399        let schedule = build_schedule(&nodes);
400        for &idx in &schedule.unscheduled {
401            let computation_id = nodes[idx].0;
402            error!(computation = computation_id, "computation skipped: requirement cycle");
403            counter!(
404                "derived_computation_failures_total",
405                "computation" => computation_id,
406                "reason" => "cycle"
407            )
408            .increment(1);
409            let _ = self
410                .event_tx
411                .send(DerivedDataEvent::ComputationFailed { computation_id, block });
412        }
413
414        let mut succeeded: FxHashSet<ComputationId> = FxHashSet::default();
415        for stage in &schedule.stages {
416            // Split the stage into runnable computations and ones whose requirements did
417            // not hold this block; the latter are skipped and reported as failed.
418            let mut runnable = Vec::new();
419            {
420                let store = self.store.read().await;
421                for &idx in stage {
422                    let reqs = &nodes[idx].1;
423                    let fresh_ready = reqs
424                        .fresh_requirements()
425                        .iter()
426                        .all(|id| succeeded.contains(id));
427                    let stale_ready = reqs
428                        .stale_requirements()
429                        .iter()
430                        .all(|id| succeeded.contains(id) || store.output_block(id).is_some());
431                    if fresh_ready && stale_ready {
432                        runnable.push(idx);
433                    } else {
434                        let computation_id = nodes[idx].0;
435                        counter!(
436                            "derived_computation_failures_total",
437                            "computation" => computation_id,
438                            "reason" => "upstream_failed"
439                        )
440                        .increment(1);
441                        let _ = self
442                            .event_tx
443                            .send(DerivedDataEvent::ComputationFailed { computation_id, block });
444                    }
445                }
446            }
447
448            if runnable.is_empty() {
449                continue;
450            }
451
452            // Run this stage's computations concurrently; they read the store as needed.
453            let results = join_all(runnable.iter().map(|&idx| async move {
454                let start = Instant::now();
455                let result = self.computations[idx]
456                    .compute_erased(&self.market_data, &self.store, changed, block)
457                    .await;
458                (idx, result, start.elapsed())
459            }))
460            .await;
461
462            // Persist and report in stage order, taking the write lock once for the stage.
463            let mut store = self.store.write().await;
464            for (idx, result, elapsed) in results {
465                let computation_id = nodes[idx].0;
466                match result {
467                    Ok(write) => {
468                        (write.persist)(&mut store);
469                        histogram!(
470                            "derived_computation_duration_seconds",
471                            "computation" => computation_id
472                        )
473                        .record(elapsed.as_secs_f64());
474                        gauge!(
475                            "derived_last_success_timestamp_seconds",
476                            "computation" => computation_id
477                        )
478                        .set(unix_now_seconds());
479                        info!(
480                            computation = computation_id,
481                            failed = write.failed_items.len(),
482                            elapsed_ms = elapsed.as_millis(),
483                            "computation complete"
484                        );
485                        let _ = self
486                            .event_tx
487                            .send(DerivedDataEvent::ComputationComplete {
488                                computation_id,
489                                block,
490                                failed_items: write.failed_items,
491                            });
492                        succeeded.insert(computation_id);
493                    }
494                    Err(e) => {
495                        counter!(
496                            "derived_computation_failures_total",
497                            "computation" => computation_id,
498                            "reason" => "error"
499                        )
500                        .increment(1);
501                        warn!(
502                            error = ?e,
503                            computation = computation_id,
504                            elapsed_ms = elapsed.as_millis(),
505                            "computation failed"
506                        );
507                        let _ = self
508                            .event_tx
509                            .send(DerivedDataEvent::ComputationFailed { computation_id, block });
510                    }
511                }
512            }
513        }
514
515        info!(
516            block,
517            total_ms = total_start.elapsed().as_millis(),
518            "all derived computations complete"
519        );
520    }
521
522    ////// Recovers from a broadcast lag without a full-topology recompute.
523    ///
524    /// Drains the events still buffered in `event_rx` (returning the receiver to the live tail so
525    /// it cannot immediately re-lag), coalesces them into one incremental `ChangedComponents`,
526    /// and recomputes just that union.
527    ///
528    /// Components lost in the dropped window are not recomputed. Added and updated ones
529    /// self-correct on their next `MarketUpdated`; removed ones never reappear, leaving stale
530    /// `spot_prices`/`pool_depths` entries for the life of the process. Routing is unaffected:
531    /// derived data is only read per graph edge, and a removed component has no edges.
532    async fn recover_from_lag(&self, event_rx: &mut broadcast::Receiver<MarketEvent>) {
533        let mut drained = Vec::new();
534        loop {
535            match event_rx.try_recv() {
536                Ok(event) => drained.push(event),
537                Err(broadcast::error::TryRecvError::Empty) => break,
538                Err(broadcast::error::TryRecvError::Lagged(n)) => {
539                    counter!("derived_manager_lagged_events_total").increment(n);
540                    continue;
541                }
542                Err(broadcast::error::TryRecvError::Closed) => break,
543            }
544        }
545        if let Some(changed) = coalesce_market_events(&drained) {
546            self.compute_all(&changed).await;
547        }
548    }
549}
550
551/// Seconds since the Unix epoch, for freshness gauges consumed as `time() - <gauge>`.
552fn unix_now_seconds() -> f64 {
553    SystemTime::now()
554        .duration_since(UNIX_EPOCH)
555        .map(|elapsed| elapsed.as_secs_f64())
556        .unwrap_or(0.0)
557}
558
559/// Computes the dependency-ordered execution plan for `nodes` (id paired with its
560/// requirements).
561///
562/// Each node lands in a later stage than the `nodes` it requires; input order is
563/// preserved within a stage. Nodes caught in a requirement cycle cannot be ordered and
564/// are returned as `unscheduled`. A requirement naming an id absent from `nodes` does
565/// not affect ordering (it is left to the runtime readiness check).
566fn build_schedule(nodes: &[(ComputationId, ComputationRequirements)]) -> ComputationSchedule {
567    let ids: Vec<ComputationId> = nodes
568        .iter()
569        .map(|(id, _)| *id)
570        .collect();
571    let mut stage_of: Vec<Option<usize>> = vec![None; nodes.len()];
572
573    loop {
574        let mut progressed = false;
575        for (idx, (_, reqs)) in nodes.iter().enumerate() {
576            if stage_of[idx].is_some() {
577                continue;
578            }
579            let mut stage = 0;
580            let mut ready = true;
581            for dep in reqs
582                .fresh_requirements()
583                .iter()
584                .chain(reqs.stale_requirements().iter())
585            {
586                let Some(dep_idx) = ids.iter().position(|id| id == dep) else {
587                    continue;
588                };
589                match stage_of[dep_idx] {
590                    Some(dep_stage) => stage = stage.max(dep_stage + 1),
591                    None => {
592                        ready = false;
593                        break;
594                    }
595                }
596            }
597            if ready {
598                stage_of[idx] = Some(stage);
599                progressed = true;
600            }
601        }
602        if !progressed {
603            break;
604        }
605    }
606
607    let stage_count = stage_of
608        .iter()
609        .filter_map(|stage| *stage)
610        .max()
611        .map_or(0, |max| max + 1);
612    let mut stages = vec![Vec::new(); stage_count];
613    let mut unscheduled = Vec::new();
614    for (idx, stage) in stage_of.iter().enumerate() {
615        match stage {
616            Some(stage) => stages[*stage].push(idx),
617            None => unscheduled.push(idx),
618        }
619    }
620    ComputationSchedule { stages, unscheduled }
621}
622
623#[async_trait]
624impl MarketEventHandler for ComputationManager {
625    async fn handle_event(&mut self, event: &MarketEvent) -> Result<(), EventError> {
626        match event {
627            MarketEvent::MarketUpdated {
628                added_components,
629                removed_components,
630                updated_components,
631            } if !added_components.is_empty() ||
632                !removed_components.is_empty() ||
633                !updated_components.is_empty() =>
634            {
635                trace!(
636                    added = added_components.len(),
637                    removed = removed_components.len(),
638                    updated = updated_components.len(),
639                    "market updated, running incremental computations"
640                );
641
642                let changed = ChangedComponents {
643                    added: added_components.clone(),
644                    removed: removed_components.clone(),
645                    updated: updated_components.clone(),
646                    is_full_recompute: false,
647                };
648                self.compute_all(&changed).await;
649            }
650            _ => {
651                trace!("empty market update, skipping computations");
652            }
653        }
654
655        Ok(())
656    }
657}
658
659#[cfg(test)]
660mod tests {
661    use std::sync::{
662        atomic::{AtomicBool, Ordering},
663        Arc,
664    };
665
666    use tokio::sync::broadcast;
667
668    use super::*;
669    use crate::{
670        algorithm::test_utils::{component, setup_market_weighted, token, MockProtocolSim},
671        derived::computation::{ComputationOutput, FailedItem, FailedItemError},
672        feed::market_data::{MarketData, MarketState},
673        types::BlockInfo,
674    };
675
676    /// Drains all currently-pending events from a broadcast receiver into a Vec.
677    fn drain_events(rx: &mut broadcast::Receiver<DerivedDataEvent>) -> Vec<DerivedDataEvent> {
678        let mut events = vec![];
679        loop {
680            match rx.try_recv() {
681                Ok(e) => events.push(e),
682                Err(broadcast::error::TryRecvError::Empty) => break,
683                Err(broadcast::error::TryRecvError::Lagged(_)) => continue,
684                Err(broadcast::error::TryRecvError::Closed) => break,
685            }
686        }
687        events
688    }
689
690    // --- coalesce_market_events: net semantics over a drained batch (pure) ---------
691
692    #[test]
693    fn coalesce_empty_batch_returns_none() {
694        assert!(coalesce_market_events(&[]).is_none());
695    }
696
697    #[test]
698    fn coalesce_unions_added_and_updated_across_events() {
699        let eth = token(1, "ETH");
700        let usdc = token(2, "USDC");
701        let e1 = MarketEvent::MarketUpdated {
702            added_components: FxHashMap::from_iter([(
703                "eth_usdc".to_string(),
704                vec![eth.address.clone(), usdc.address.clone()],
705            )]),
706            removed_components: vec![],
707            updated_components: vec![],
708        };
709        let e2 = MarketEvent::MarketUpdated {
710            added_components: FxHashMap::default(),
711            removed_components: vec![],
712            updated_components: vec!["eth_usdc".to_string(), "dai_usdc".to_string()],
713        };
714        let c = coalesce_market_events(&[e1, e2]).expect("net changes present");
715        assert!(!c.is_full_recompute);
716        // eth_usdc was added, so it stays in `added` (not double-counted in `updated`)
717        assert!(c.added.contains_key("eth_usdc"));
718        assert!(!c
719            .updated
720            .contains(&"eth_usdc".to_string()));
721        // dai_usdc only ever appeared as updated
722        assert!(c
723            .updated
724            .contains(&"dai_usdc".to_string()));
725    }
726
727    #[test]
728    fn coalesce_add_then_remove_nets_to_removed() {
729        let eth = token(1, "ETH");
730        let usdc = token(2, "USDC");
731        let add = MarketEvent::MarketUpdated {
732            added_components: FxHashMap::from_iter([(
733                "eth_usdc".to_string(),
734                vec![eth.address.clone(), usdc.address.clone()],
735            )]),
736            removed_components: vec![],
737            updated_components: vec![],
738        };
739        let remove = MarketEvent::MarketUpdated {
740            added_components: FxHashMap::default(),
741            removed_components: vec!["eth_usdc".to_string()],
742            updated_components: vec![],
743        };
744        let c = coalesce_market_events(&[add, remove]).expect("net removal present");
745        assert!(!c.added.contains_key("eth_usdc"));
746        assert!(c
747            .removed
748            .contains(&"eth_usdc".to_string()));
749        assert!(!c
750            .updated
751            .contains(&"eth_usdc".to_string()));
752    }
753
754    #[tokio::test]
755    async fn lag_recovery_recomputes_incrementally_and_drains_to_tail() {
756        let eth = token(1, "ETH");
757        let usdc = token(2, "USDC");
758        let (market, _) = setup_market_weighted(vec![(
759            "eth_usdc",
760            &eth,
761            &usdc,
762            MockProtocolSim::new(2000.0).with_gas(0),
763        )]);
764        let config = ComputationManagerConfig::new().with_gas_token(eth.address.clone());
765        let (manager, _out_rx) = ComputationManager::new(config, market).unwrap();
766
767        // Capacity-2 input channel; send 5 without reading to force Lagged on recv.
768        let (tx, mut rx) = broadcast::channel::<MarketEvent>(2);
769        for _ in 0..5 {
770            tx.send(MarketEvent::MarketUpdated {
771                added_components: FxHashMap::from_iter([(
772                    "eth_usdc".to_string(),
773                    vec![eth.address.clone(), usdc.address.clone()],
774                )]),
775                removed_components: vec![],
776                updated_components: vec![],
777            })
778            .unwrap();
779        }
780        let err = rx
781            .recv()
782            .await
783            .expect_err("receiver must have lagged");
784        assert!(matches!(err, broadcast::error::RecvError::Lagged(_)));
785
786        manager.recover_from_lag(&mut rx).await;
787
788        // Recovery recomputed the coalesced change incrementally...
789        let store = manager.store();
790        let guard = store.read().await;
791        assert!(guard.spot_prices().is_some());
792        assert!(guard.token_prices().is_some());
793        assert!(guard.component_depths().is_some());
794        drop(guard);
795        // ...and the receiver is back at the live tail (buffer drained).
796        assert!(matches!(rx.try_recv(), Err(broadcast::error::TryRecvError::Empty)));
797    }
798
799    // --- build_schedule: dependency staging (pure) --------------------------------
800
801    #[test]
802    fn schedule_empty_has_no_stages() {
803        let schedule = build_schedule(&[]);
804        assert!(schedule.stages.is_empty());
805        assert!(schedule.unscheduled.is_empty());
806    }
807
808    #[test]
809    fn schedule_single_root_is_one_stage() {
810        let schedule = build_schedule(&[("a", ComputationRequirements::none())]);
811        assert_eq!(schedule.stages, vec![vec![0]]);
812        assert!(schedule.unscheduled.is_empty());
813    }
814
815    #[test]
816    fn schedule_independent_roots_share_one_stage() {
817        let schedule = build_schedule(&[
818            ("a", ComputationRequirements::none()),
819            ("b", ComputationRequirements::none()),
820        ]);
821        assert_eq!(schedule.stages, vec![vec![0, 1]]);
822        assert!(schedule.unscheduled.is_empty());
823    }
824
825    #[test]
826    fn schedule_chain_orders_into_successive_stages() {
827        // a <- b <- c
828        let schedule = build_schedule(&[
829            ("a", ComputationRequirements::none()),
830            ("b", ComputationRequirements::fresh(["a"])),
831            ("c", ComputationRequirements::fresh(["b"])),
832        ]);
833        assert_eq!(schedule.stages, vec![vec![0], vec![1], vec![2]]);
834        assert!(schedule.unscheduled.is_empty());
835    }
836
837    #[test]
838    fn schedule_diamond_places_join_after_both_parents() {
839        // a <- {b, c} <- d; mirrors fynd's spot -> {token, component} fan-out.
840        let schedule = build_schedule(&[
841            ("a", ComputationRequirements::none()),
842            ("b", ComputationRequirements::fresh(["a"])),
843            ("c", ComputationRequirements::fresh(["a"])),
844            ("d", ComputationRequirements::fresh(["b", "c"])),
845        ]);
846        assert_eq!(schedule.stages, vec![vec![0], vec![1, 2], vec![3]]);
847        assert!(schedule.unscheduled.is_empty());
848    }
849
850    #[test]
851    fn schedule_preserves_input_order_within_a_stage() {
852        let schedule = build_schedule(&[
853            ("a", ComputationRequirements::none()),
854            ("b", ComputationRequirements::fresh(["a"])),
855            ("c", ComputationRequirements::fresh(["a"])),
856        ]);
857        // b registered before c, so it comes first in the shared stage.
858        assert_eq!(schedule.stages, vec![vec![0], vec![1, 2]]);
859    }
860
861    #[test]
862    fn schedule_stale_requirement_orders_after_its_producer() {
863        let schedule = build_schedule(&[
864            ("a", ComputationRequirements::none()),
865            ("b", ComputationRequirements::stale(["a"])),
866        ]);
867        assert_eq!(schedule.stages, vec![vec![0], vec![1]]);
868    }
869
870    #[test]
871    fn schedule_requirement_on_unregistered_id_does_not_affect_ordering() {
872        // "ghost" is not registered, so "a" is treated as a root.
873        let schedule = build_schedule(&[("a", ComputationRequirements::fresh(["ghost"]))]);
874        assert_eq!(schedule.stages, vec![vec![0]]);
875        assert!(schedule.unscheduled.is_empty());
876    }
877
878    #[test]
879    fn schedule_two_node_cycle_is_unscheduled() {
880        let schedule = build_schedule(&[
881            ("a", ComputationRequirements::fresh(["b"])),
882            ("b", ComputationRequirements::fresh(["a"])),
883        ]);
884        assert!(schedule.stages.is_empty());
885        assert_eq!(schedule.unscheduled, vec![0, 1]);
886    }
887
888    #[test]
889    fn schedule_isolates_cycle_from_schedulable_nodes() {
890        // "root" schedules normally; "x" and "y" form a cycle and are unscheduled.
891        let schedule = build_schedule(&[
892            ("root", ComputationRequirements::none()),
893            ("x", ComputationRequirements::fresh(["y"])),
894            ("y", ComputationRequirements::fresh(["x"])),
895        ]);
896        assert_eq!(schedule.stages, vec![vec![0]]);
897        assert_eq!(schedule.unscheduled, vec![1, 2]);
898    }
899
900    #[test]
901    fn invalid_slippage_threshold_returns_error() {
902        let (market, _) = setup_market_weighted(vec![]);
903        let config = ComputationManagerConfig::new().with_depth_slippage_threshold(1.5);
904
905        let result = ComputationManager::new(config, market);
906        assert!(matches!(result, Err(ComputationError::InvalidConfiguration(_))));
907    }
908
909    #[tokio::test]
910    async fn handle_event_runs_computations_on_market_update() {
911        let eth = token(1, "ETH");
912        let usdc = token(2, "USDC");
913
914        let (market, _) = setup_market_weighted(vec![(
915            "eth_usdc",
916            &eth,
917            &usdc,
918            MockProtocolSim::new(2000.0).with_gas(0),
919        )]);
920
921        let config = ComputationManagerConfig::new().with_gas_token(eth.address.clone());
922        let (mut manager, _event_rx) = ComputationManager::new(config, market).unwrap();
923
924        let event = MarketEvent::MarketUpdated {
925            added_components: FxHashMap::from_iter([(
926                "eth_usdc".to_string(),
927                vec![eth.address.clone(), usdc.address.clone()],
928            )]),
929            removed_components: vec![],
930            updated_components: vec![],
931        };
932
933        manager
934            .handle_event(&event)
935            .await
936            .unwrap();
937
938        let store = manager.store();
939        let guard = store.read().await;
940        assert!(guard.spot_prices().is_some());
941        assert!(guard.component_depths().is_some());
942        assert!(guard.token_prices().is_some());
943    }
944
945    #[tokio::test]
946    async fn handle_event_skips_empty_update() {
947        let (market, _) = setup_market_weighted(vec![]);
948        let config = ComputationManagerConfig::new();
949        let (mut manager, _event_rx) = ComputationManager::new(config, market).unwrap();
950
951        let event = MarketEvent::MarketUpdated {
952            added_components: FxHashMap::default(),
953            removed_components: vec![],
954            updated_components: vec![],
955        };
956
957        manager
958            .handle_event(&event)
959            .await
960            .unwrap();
961
962        let store = manager.store();
963        let guard = store.read().await;
964        assert!(guard.token_prices().is_none());
965    }
966
967    #[tokio::test]
968    async fn run_shuts_down_on_signal() {
969        let (market, _) = setup_market_weighted(vec![]);
970        let config = ComputationManagerConfig::new();
971        let (manager, _event_rx) = ComputationManager::new(config, market).unwrap();
972
973        let (_event_tx, event_rx) = broadcast::channel::<MarketEvent>(16);
974        let (shutdown_tx, shutdown_rx) = broadcast::channel::<()>(1);
975
976        let handle = tokio::spawn(async move {
977            manager.run(event_rx, shutdown_rx).await;
978        });
979
980        shutdown_tx.send(()).unwrap();
981
982        tokio::time::timeout(tokio::time::Duration::from_secs(1), handle)
983            .await
984            .expect("manager should shutdown")
985            .expect("task should complete successfully");
986    }
987
988    // --- registry seam: custom computations driven through the manager -------------
989
990    #[derive(Clone, Debug, PartialEq)]
991    struct CounterOutput(u32);
992
993    /// A minimal computation that ignores market data and uses the default `persist`
994    /// (the path a downstream computation takes: store into the generic slot).
995    struct CounterComputation;
996
997    #[async_trait::async_trait]
998    impl DerivedComputation for CounterComputation {
999        type Output = CounterOutput;
1000        const ID: ComputationId = "counter";
1001
1002        async fn compute(
1003            &self,
1004            _market: &MarketData,
1005            _store: &SharedDerivedDataRef,
1006            _changed: &ChangedComponents,
1007        ) -> Result<ComputationOutput<Self::Output>, ComputationError> {
1008            Ok(ComputationOutput::success(CounterOutput(7)))
1009        }
1010    }
1011
1012    /// Builds a market carrying a `last_updated` block so `compute_all` runs.
1013    fn market_with_block() -> MarketData {
1014        let eth = token(1, "ETH");
1015        let usdc = token(2, "USDC");
1016        let (market, _) = setup_market_weighted(vec![(
1017            "eth_usdc",
1018            &eth,
1019            &usdc,
1020            MockProtocolSim::new(2000.0).with_gas(0),
1021        )]);
1022        market
1023    }
1024
1025    #[tokio::test]
1026    async fn registered_custom_computation_runs_and_persists_via_default_slot() {
1027        let (mut manager, mut event_rx) = ComputationManager::empty(market_with_block());
1028        manager
1029            .register(CounterComputation)
1030            .unwrap();
1031
1032        manager
1033            .compute_all(&ChangedComponents { is_full_recompute: true, ..Default::default() })
1034            .await;
1035
1036        let store = manager.store();
1037        let guard = store.read().await;
1038        assert_eq!(
1039            guard.output::<CounterOutput>(CounterComputation::ID),
1040            Some(&CounterOutput(7)),
1041            "default persist should write the output into the generic slot"
1042        );
1043        assert!(guard
1044            .output_block(CounterComputation::ID)
1045            .is_some());
1046
1047        let events = drain_events(&mut event_rx);
1048        assert!(
1049            events.iter().any(|e| matches!(
1050                e,
1051                DerivedDataEvent::ComputationComplete { computation_id: "counter", .. }
1052            )),
1053            "expected ComputationComplete(counter), got: {events:?}"
1054        );
1055    }
1056
1057    #[test]
1058    fn registering_duplicate_id_is_rejected() {
1059        let (mut manager, _event_rx) = ComputationManager::empty(market_with_block());
1060        manager
1061            .register(CounterComputation)
1062            .unwrap();
1063
1064        let result = manager.register(CounterComputation);
1065
1066        assert!(matches!(result, Err(ComputationError::DuplicateComputationId("counter"))));
1067    }
1068
1069    // --- exact event sequences (characterization) ---------------------------------
1070
1071    /// Reduces an event stream to `(kind, computation_id)` pairs for exact comparison.
1072    fn event_summary(events: &[DerivedDataEvent]) -> Vec<(&'static str, &'static str)> {
1073        events
1074            .iter()
1075            .map(|event| match event {
1076                DerivedDataEvent::NewBlock { .. } => ("new_block", ""),
1077                DerivedDataEvent::ComputationComplete { computation_id, .. } => {
1078                    ("complete", *computation_id)
1079                }
1080                DerivedDataEvent::ComputationFailed { computation_id, .. } => {
1081                    ("failed", *computation_id)
1082                }
1083            })
1084            .collect()
1085    }
1086
1087    /// Subscribes, runs one full-recompute pass, and returns the events it emitted.
1088    async fn run_full_recompute(manager: &ComputationManager) -> Vec<DerivedDataEvent> {
1089        let mut event_rx = manager.event_sender().subscribe();
1090        manager
1091            .compute_all(&ChangedComponents { is_full_recompute: true, ..Default::default() })
1092            .await;
1093        drain_events(&mut event_rx)
1094    }
1095
1096    /// Defines a market-independent test computation with a fixed id, requirements, result.
1097    macro_rules! test_computation {
1098        ($name:ident, $id:literal, $reqs:expr, $result:expr) => {
1099            struct $name;
1100
1101            #[async_trait::async_trait]
1102            impl DerivedComputation for $name {
1103                type Output = ();
1104                const ID: ComputationId = $id;
1105
1106                fn requirements(&self) -> ComputationRequirements {
1107                    $reqs
1108                }
1109
1110                async fn compute(
1111                    &self,
1112                    _market: &MarketData,
1113                    _store: &SharedDerivedDataRef,
1114                    _changed: &ChangedComponents,
1115                ) -> Result<ComputationOutput<Self::Output>, ComputationError> {
1116                    $result
1117                }
1118            }
1119        };
1120    }
1121
1122    test_computation!(
1123        RootOk,
1124        "root",
1125        ComputationRequirements::none(),
1126        Ok(ComputationOutput::success(()))
1127    );
1128    test_computation!(
1129        DepOnRoot,
1130        "dep",
1131        ComputationRequirements::fresh(["root"]),
1132        Ok(ComputationOutput::success(()))
1133    );
1134    test_computation!(
1135        SecondDepOnRoot,
1136        "dep2",
1137        ComputationRequirements::fresh(["root"]),
1138        Ok(ComputationOutput::success(()))
1139    );
1140    test_computation!(
1141        RootErr,
1142        "boom",
1143        ComputationRequirements::none(),
1144        Err(ComputationError::InvalidConfiguration("boom".to_string()))
1145    );
1146    test_computation!(
1147        DepOnBoom,
1148        "dep_boom",
1149        ComputationRequirements::fresh(["boom"]),
1150        Ok(ComputationOutput::success(()))
1151    );
1152    test_computation!(
1153        ThirdOnBoom,
1154        "third",
1155        ComputationRequirements::fresh(["dep_boom"]),
1156        Ok(ComputationOutput::success(()))
1157    );
1158    test_computation!(
1159        StaleDepOnFlaky,
1160        "stale_dep",
1161        ComputationRequirements::stale(["flaky"]),
1162        Ok(ComputationOutput::success(()))
1163    );
1164    test_computation!(
1165        GhostDependent,
1166        "needs_ghost",
1167        ComputationRequirements::fresh(["ghost"]),
1168        Ok(ComputationOutput::success(()))
1169    );
1170    test_computation!(
1171        PartialProducer,
1172        "partial",
1173        ComputationRequirements::none(),
1174        Ok(ComputationOutput::with_failures(
1175            (),
1176            vec![FailedItem { key: "x".to_string(), error: FailedItemError::MissingSpotPrice }]
1177        ))
1178    );
1179    test_computation!(
1180        DepOnPartial,
1181        "dep_partial",
1182        ComputationRequirements::fresh(["partial"]),
1183        Ok(ComputationOutput::success(()))
1184    );
1185
1186    /// A producer that succeeds while its flag is set and fails once it is cleared, so a
1187    /// later block can exercise the stale-dependency path (producer failed this block, but
1188    /// a prior-block value is still in the store).
1189    struct FlakyProducer {
1190        succeed: Arc<AtomicBool>,
1191    }
1192
1193    #[async_trait::async_trait]
1194    impl DerivedComputation for FlakyProducer {
1195        type Output = ();
1196        const ID: ComputationId = "flaky";
1197
1198        async fn compute(
1199            &self,
1200            _market: &MarketData,
1201            _store: &SharedDerivedDataRef,
1202            _changed: &ChangedComponents,
1203        ) -> Result<ComputationOutput<Self::Output>, ComputationError> {
1204            if self.succeed.load(Ordering::SeqCst) {
1205                Ok(ComputationOutput::success(()))
1206            } else {
1207                Err(ComputationError::InvalidConfiguration("flaky".to_string()))
1208            }
1209        }
1210    }
1211
1212    #[tokio::test]
1213    async fn events_follow_dependency_order_across_stages() {
1214        let (mut manager, _event_rx) = ComputationManager::empty(market_with_block());
1215        manager.register(RootOk).unwrap();
1216        manager.register(DepOnRoot).unwrap();
1217
1218        let events = run_full_recompute(&manager).await;
1219
1220        assert_eq!(
1221            event_summary(&events),
1222            vec![("new_block", ""), ("complete", "root"), ("complete", "dep")]
1223        );
1224    }
1225
1226    #[tokio::test]
1227    async fn events_preserve_registration_order_within_a_stage() {
1228        let (mut manager, _event_rx) = ComputationManager::empty(market_with_block());
1229        manager.register(RootOk).unwrap();
1230        manager.register(DepOnRoot).unwrap();
1231        manager
1232            .register(SecondDepOnRoot)
1233            .unwrap();
1234
1235        let events = run_full_recompute(&manager).await;
1236
1237        // root runs in stage 0; dep then dep2 share stage 1 in registration order.
1238        assert_eq!(
1239            event_summary(&events),
1240            vec![
1241                ("new_block", ""),
1242                ("complete", "root"),
1243                ("complete", "dep"),
1244                ("complete", "dep2"),
1245            ]
1246        );
1247    }
1248
1249    #[tokio::test]
1250    async fn failed_dependency_cascades_to_dependents() {
1251        let (mut manager, _event_rx) = ComputationManager::empty(market_with_block());
1252        manager.register(RootErr).unwrap();
1253        manager.register(DepOnBoom).unwrap();
1254
1255        let events = run_full_recompute(&manager).await;
1256
1257        // boom fails in stage 0; its dependent is skipped and reported failed.
1258        assert_eq!(
1259            event_summary(&events),
1260            vec![("new_block", ""), ("failed", "boom"), ("failed", "dep_boom")]
1261        );
1262    }
1263
1264    #[tokio::test]
1265    async fn computation_with_unregistered_requirement_is_skipped() {
1266        let (mut manager, _event_rx) = ComputationManager::empty(market_with_block());
1267        manager
1268            .register(GhostDependent)
1269            .unwrap();
1270
1271        let events = run_full_recompute(&manager).await;
1272
1273        // "ghost" is never registered, so its fresh dependent never runs.
1274        assert_eq!(event_summary(&events), vec![("new_block", ""), ("failed", "needs_ghost")]);
1275    }
1276
1277    #[tokio::test]
1278    async fn fresh_dependent_runs_when_producer_succeeds_partially() {
1279        let (mut manager, _event_rx) = ComputationManager::empty(market_with_block());
1280        manager
1281            .register(PartialProducer)
1282            .unwrap();
1283        manager.register(DepOnPartial).unwrap();
1284
1285        let events = run_full_recompute(&manager).await;
1286
1287        // A partial success (Ok with failed_items) still counts as succeeded, so the fresh
1288        // dependent runs -- the compatibility invariant with the old hardcoded flow.
1289        assert_eq!(
1290            event_summary(&events),
1291            vec![("new_block", ""), ("complete", "partial"), ("complete", "dep_partial"),]
1292        );
1293    }
1294
1295    #[tokio::test]
1296    async fn failure_cascade_propagates_through_three_levels() {
1297        let (mut manager, _event_rx) = ComputationManager::empty(market_with_block());
1298        manager.register(RootErr).unwrap();
1299        manager.register(DepOnBoom).unwrap();
1300        manager.register(ThirdOnBoom).unwrap();
1301
1302        let events = run_full_recompute(&manager).await;
1303
1304        // boom fails; dep_boom is skipped; third (needs dep_boom) is skipped transitively.
1305        assert_eq!(
1306            event_summary(&events),
1307            vec![
1308                ("new_block", ""),
1309                ("failed", "boom"),
1310                ("failed", "dep_boom"),
1311                ("failed", "third"),
1312            ]
1313        );
1314    }
1315
1316    #[tokio::test]
1317    async fn stale_dependency_runs_on_prior_value_after_producer_fails() {
1318        let succeed = Arc::new(AtomicBool::new(true));
1319        let (mut manager, _event_rx) = ComputationManager::empty(market_with_block());
1320        manager
1321            .register(FlakyProducer { succeed: Arc::clone(&succeed) })
1322            .unwrap();
1323        manager
1324            .register(StaleDepOnFlaky)
1325            .unwrap();
1326
1327        // Block 1: producer succeeds and its value is stored.
1328        let first = run_full_recompute(&manager).await;
1329        assert_eq!(
1330            event_summary(&first),
1331            vec![("new_block", ""), ("complete", "flaky"), ("complete", "stale_dep")]
1332        );
1333
1334        // Block 2: producer fails, but its prior-block value remains, so the stale
1335        // dependent still runs.
1336        succeed.store(false, Ordering::SeqCst);
1337        let second = run_full_recompute(&manager).await;
1338        assert_eq!(
1339            event_summary(&second),
1340            vec![("new_block", ""), ("failed", "flaky"), ("complete", "stale_dep")]
1341        );
1342    }
1343
1344    #[tokio::test]
1345    async fn test_spot_price_failure_cascade() {
1346        // Real fynd flow: a full recompute with no sim state makes spot prices fail outright.
1347        // Pool depths depend on them and fail with them. Token prices read no derived data,
1348        // so they still complete.
1349        let (manager, _event_rx) = ComputationManager::new(
1350            ComputationManagerConfig::new(),
1351            market_with_component_no_sim_state(),
1352        )
1353        .unwrap();
1354
1355        let events = run_full_recompute(&manager).await;
1356
1357        assert_eq!(
1358            event_summary(&events),
1359            vec![
1360                ("new_block", ""),
1361                ("failed", "spot_prices"),
1362                ("complete", "token_prices"),
1363                ("failed", "pool_depths"),
1364            ]
1365        );
1366    }
1367
1368    /// Creates a market with a component in topology but WITHOUT simulation state.
1369    ///
1370    /// Used to trigger `TotalFailure` in spot_price computation (full recompute with
1371    /// all components missing sim_state → succeeded == 0 → failure).
1372    fn market_with_component_no_sim_state() -> MarketData {
1373        let eth = token(1, "ETH");
1374        let usdc = token(2, "USDC");
1375        let component = component("component", &[eth.clone(), usdc.clone()]);
1376
1377        let mut market = MarketState::new();
1378        market.update_last_updated(BlockInfo::new(10, "0xhash".into(), 0));
1379        market.upsert_components(std::iter::once(component));
1380        // Note: no update_states() — simulation state is intentionally absent
1381        market.upsert_tokens([eth, usdc]);
1382        MarketData::new(std::sync::Arc::new(tokio::sync::RwLock::new(market)))
1383    }
1384
1385    /// Creates a market with two components: one with sim state (component succeeds) and one
1386    /// without (component fails). Used to trigger partial spot price failure.
1387    fn market_with_mixed_sim_states() -> MarketData {
1388        let eth = token(1, "ETH");
1389        let usdc = token(2, "USDC");
1390        let dai = token(3, "DAI");
1391
1392        let component1 = component("eth_usdc", &[eth.clone(), usdc.clone()]);
1393        let component2 = component("eth_dai", &[eth.clone(), dai.clone()]);
1394
1395        let mut market = MarketState::new();
1396        market.update_last_updated(BlockInfo::new(10, "0xhash".into(), 0));
1397        market.upsert_components([component1, component2]);
1398        // Only component1 has simulation state; component2 intentionally has none
1399        market
1400            .update_states([("eth_usdc".to_string(), Box::new(MockProtocolSim::new(2000.0)) as _)]);
1401        market.upsert_tokens([eth, usdc, dai]);
1402        MarketData::new(std::sync::Arc::new(tokio::sync::RwLock::new(market)))
1403    }
1404
1405    #[tokio::test]
1406    async fn test_spot_price_failure_broadcasts_computation_failed() {
1407        let market = market_with_component_no_sim_state();
1408        let config = ComputationManagerConfig::new();
1409        let (manager, mut event_rx) = ComputationManager::new(config, market).unwrap();
1410
1411        // Full recompute with components that have no sim_state → TotalFailure
1412        let changed = ChangedComponents { is_full_recompute: true, ..Default::default() };
1413        manager.compute_all(&changed).await;
1414
1415        let events = drain_events(&mut event_rx);
1416
1417        assert!(
1418            events.iter().any(|e| matches!(
1419                e,
1420                DerivedDataEvent::ComputationFailed { computation_id: "spot_prices", .. }
1421            )),
1422            "expected ComputationFailed(spot_prices) in events: {events:?}"
1423        );
1424    }
1425
1426    #[tokio::test]
1427    async fn run_shuts_down_on_channel_close() {
1428        let (market, _) = setup_market_weighted(vec![]);
1429        let config = ComputationManagerConfig::new();
1430        let (manager, _event_rx) = ComputationManager::new(config, market).unwrap();
1431
1432        let (event_tx, event_rx) = broadcast::channel::<MarketEvent>(16);
1433        let (_shutdown_tx, shutdown_rx) = broadcast::channel::<()>(1);
1434
1435        let handle = tokio::spawn(async move {
1436            manager.run(event_rx, shutdown_rx).await;
1437        });
1438
1439        drop(event_tx);
1440
1441        tokio::time::timeout(tokio::time::Duration::from_secs(1), handle)
1442            .await
1443            .expect("manager should shutdown on channel close")
1444            .expect("task should complete successfully");
1445    }
1446
1447    #[tokio::test]
1448    async fn partial_spot_price_failure_broadcasts_computation_complete() {
1449        // market_with_mixed_sim_states has component1 (with sim state) and component2 (without)
1450        // → spot price computation partially succeeds → ComputationComplete with failed_items
1451        let market = market_with_mixed_sim_states();
1452        let config = ComputationManagerConfig::new();
1453        let (manager, mut event_rx) = ComputationManager::new(config, market).unwrap();
1454
1455        let changed = ChangedComponents { is_full_recompute: true, ..Default::default() };
1456        manager.compute_all(&changed).await;
1457
1458        let events = drain_events(&mut event_rx);
1459
1460        // Should broadcast ComputationComplete (not ComputationFailed) because component1 succeeds
1461        assert!(
1462            events.iter().any(|e| matches!(
1463                e,
1464                DerivedDataEvent::ComputationComplete { computation_id: "spot_prices", .. }
1465            )),
1466            "expected ComputationComplete(spot_prices), got: {events:?}"
1467        );
1468        assert!(
1469            !events.iter().any(|e| matches!(
1470                e,
1471                DerivedDataEvent::ComputationFailed { computation_id: "spot_prices", .. }
1472            )),
1473            "should not broadcast ComputationFailed for partial failure"
1474        );
1475
1476        // The ComputationComplete event should carry the failed item for component2
1477        let complete = events.iter().find(|e| {
1478            matches!(e, DerivedDataEvent::ComputationComplete { computation_id: "spot_prices", .. })
1479        });
1480        if let Some(DerivedDataEvent::ComputationComplete { failed_items, .. }) = complete {
1481            assert!(
1482                !failed_items.is_empty(),
1483                "ComputationComplete should carry failed_items for component2"
1484            );
1485        }
1486
1487        // The store should persist the failure reason for the failed component.
1488        // market_with_mixed_sim_states uses token(1, "ETH") and token(3, "DAI") for component2.
1489        let eth = token(1, "ETH");
1490        let dai = token(3, "DAI");
1491        let store = manager.store();
1492        let guard = store.read().await;
1493        let key_eth_dai = ("eth_dai".to_string(), eth.address.clone(), dai.address.clone());
1494        let key_dai_eth = ("eth_dai".to_string(), dai.address.clone(), eth.address.clone());
1495        assert!(
1496            guard
1497                .spot_price_failure(&key_eth_dai)
1498                .is_some() ||
1499                guard
1500                    .spot_price_failure(&key_dai_eth)
1501                    .is_some(),
1502            "store should persist failure reason for eth_dai (missing sim state)"
1503        );
1504    }
1505
1506    // --- metrics ---------------------------------------------------------------------
1507
1508    /// Mirrors `CounterComputation`, but always fails, to exercise the failure-counter path.
1509    struct FailingComputation;
1510
1511    #[async_trait::async_trait]
1512    impl DerivedComputation for FailingComputation {
1513        type Output = ();
1514        const ID: ComputationId = "failing";
1515
1516        async fn compute(
1517            &self,
1518            _market: &MarketData,
1519            _store: &SharedDerivedDataRef,
1520            _changed: &ChangedComponents,
1521        ) -> Result<ComputationOutput<Self::Output>, ComputationError> {
1522            Err(ComputationError::InvalidConfiguration("always fails".to_string()))
1523        }
1524    }
1525
1526    /// Finds the debug value recorded for `name` carrying every label in `labels`.
1527    fn find_metric<'a>(
1528        recorded: &'a [(
1529            metrics_util::CompositeKey,
1530            Option<metrics::Unit>,
1531            Option<metrics::SharedString>,
1532            metrics_util::debugging::DebugValue,
1533        )],
1534        name: &str,
1535        labels: &[(&str, &str)],
1536    ) -> &'a metrics_util::debugging::DebugValue {
1537        recorded
1538            .iter()
1539            .find(|(key, _, _, _)| {
1540                key.key().name() == name &&
1541                    labels
1542                        .iter()
1543                        .all(|(label_key, label_value)| {
1544                            key.key()
1545                                .labels()
1546                                .any(|l| l.key() == *label_key && l.value() == *label_value)
1547                        })
1548            })
1549            .map(|(_, _, _, value)| value)
1550            .unwrap_or_else(|| panic!("missing {name}{labels:?}, got {recorded:?}"))
1551    }
1552
1553    #[test]
1554    fn compute_all_records_derived_metrics() {
1555        use metrics_util::debugging::DebugValue;
1556
1557        let recorder = metrics_util::debugging::DebuggingRecorder::new();
1558        let snapshotter = recorder.snapshotter();
1559        let rt = tokio::runtime::Builder::new_current_thread()
1560            .enable_all()
1561            .build()
1562            .expect("runtime builds");
1563
1564        metrics::with_local_recorder(&recorder, || {
1565            rt.block_on(async {
1566                let (mut manager, _event_rx) = ComputationManager::empty(market_with_block());
1567                manager
1568                    .register(CounterComputation)
1569                    .unwrap();
1570                manager
1571                    .compute_all(&ChangedComponents {
1572                        is_full_recompute: true,
1573                        ..Default::default()
1574                    })
1575                    .await;
1576            })
1577        });
1578
1579        let recorded = snapshotter.snapshot().into_vec();
1580        let recorded_names: Vec<(String, Vec<String>)> = recorded
1581            .iter()
1582            .map(|(key, _, _, _)| {
1583                (
1584                    key.key().name().to_string(),
1585                    key.key()
1586                        .labels()
1587                        .map(|l| format!("{}={}", l.key(), l.value()))
1588                        .collect(),
1589                )
1590            })
1591            .collect();
1592        for expected in
1593            ["derived_computation_duration_seconds", "derived_last_success_timestamp_seconds"]
1594        {
1595            assert!(
1596                recorded_names
1597                    .iter()
1598                    .any(|(name, labels)| name == expected &&
1599                        labels.contains(&"computation=counter".to_string())),
1600                "missing {expected}{{computation=counter}}, got {recorded_names:?}"
1601            );
1602        }
1603
1604        match find_metric(
1605            &recorded,
1606            "derived_last_success_timestamp_seconds",
1607            &[("computation", "counter")],
1608        ) {
1609            DebugValue::Gauge(value) => {
1610                assert!(value.0 > 1.7e9, "gauge value {} not a sane unix timestamp", value.0);
1611            }
1612            other => panic!("derived_last_success_timestamp_seconds is not a gauge: {other:?}"),
1613        }
1614
1615        match find_metric(
1616            &recorded,
1617            "derived_computation_duration_seconds",
1618            &[("computation", "counter")],
1619        ) {
1620            DebugValue::Histogram(samples) => {
1621                assert!(!samples.is_empty(), "expected at least one recorded duration sample");
1622            }
1623            other => panic!("derived_computation_duration_seconds is not a histogram: {other:?}"),
1624        }
1625    }
1626
1627    #[test]
1628    fn compute_all_records_failure_metric() {
1629        use metrics_util::debugging::DebugValue;
1630
1631        let recorder = metrics_util::debugging::DebuggingRecorder::new();
1632        let snapshotter = recorder.snapshotter();
1633        let rt = tokio::runtime::Builder::new_current_thread()
1634            .enable_all()
1635            .build()
1636            .expect("runtime builds");
1637
1638        metrics::with_local_recorder(&recorder, || {
1639            rt.block_on(async {
1640                let (mut manager, _event_rx) = ComputationManager::empty(market_with_block());
1641                manager
1642                    .register(FailingComputation)
1643                    .unwrap();
1644                manager
1645                    .compute_all(&ChangedComponents {
1646                        is_full_recompute: true,
1647                        ..Default::default()
1648                    })
1649                    .await;
1650            })
1651        });
1652
1653        let recorded = snapshotter.snapshot().into_vec();
1654        match find_metric(
1655            &recorded,
1656            "derived_computation_failures_total",
1657            &[("computation", "failing"), ("reason", "error")],
1658        ) {
1659            DebugValue::Counter(value) => {
1660                assert!(*value >= 1, "expected failure counter >= 1, got {value}");
1661            }
1662            other => panic!("derived_computation_failures_total is not a counter: {other:?}"),
1663        }
1664    }
1665}