smoldot_light/json_rpc_service/
statement.rs

1// Smoldot
2// Copyright (C) 2019-2022  Parity Technologies (UK) Ltd.
3// SPDX-License-Identifier: GPL-3.0-or-later WITH Classpath-exception-2.0
4
5// This program is free software: you can redistribute it and/or modify
6// it under the terms of the GNU General Public License as published by
7// the Free Software Foundation, either version 3 of the License, or
8// (at your option) any later version.
9
10// This program is distributed in the hope that it will be useful,
11// but WITHOUT ANY WARRANTY; without even the implied warranty of
12// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
13// GNU General Public License for more details.
14
15// You should have received a copy of the GNU General Public License
16// along with this program.  If not, see <http://www.gnu.org/licenses/>.
17
18use crate::network_service::{self, BroadcastStatementResult};
19use alloc::{format, string::String, vec::Vec};
20use core::{num::NonZero, time::Duration};
21use smoldot::json_rpc::methods::{HexString, InvalidReason, StatementSubmitResult, TopicFilter};
22use smoldot::json_rpc::parse;
23use smoldot::network::codec;
24
25/// Configuration for the Statement Store protocol.
26#[derive(Debug, Clone)]
27pub struct StatementProtocolConfig {
28    /// Per-subscription LRU cache size used for deduplicating delivered statements.
29    max_seen_statements: NonZero<usize>,
30    false_positive_rate: f64,
31    bloom_seed: u128,
32    affinity_update_interval: Duration,
33}
34
35impl StatementProtocolConfig {
36    pub fn new(
37        max_seen_statements: NonZero<usize>,
38        false_positive_rate: f64,
39        bloom_seed: u128,
40        affinity_update_interval: Duration,
41    ) -> Self {
42        assert!(
43            false_positive_rate.is_finite()
44                && false_positive_rate > 0.0
45                && false_positive_rate < 1.0
46        );
47        assert!(!affinity_update_interval.is_zero());
48        StatementProtocolConfig {
49            max_seen_statements,
50            false_positive_rate,
51            bloom_seed,
52            affinity_update_interval,
53        }
54    }
55
56    pub fn max_seen_statements(&self) -> NonZero<usize> {
57        self.max_seen_statements
58    }
59
60    pub fn false_positive_rate(&self) -> f64 {
61        self.false_positive_rate
62    }
63
64    pub fn bloom_seed(&self) -> u128 {
65        self.bloom_seed
66    }
67
68    pub fn affinity_update_interval(&self) -> Duration {
69        self.affinity_update_interval
70    }
71}
72
73/// JSON-RPC error code answering a submission the statement store couldn't process.
74///
75/// Matches polkadot-sdk.
76pub const STATEMENT_STORE_ERROR_CODE: i64 = 7001;
77
78/// Submission failure reported as a JSON-RPC error
79#[derive(Debug, Clone, PartialEq, Eq)]
80pub enum StatementSubmitError {
81    /// The submitted bytes don't decode into a statement.
82    InvalidEncoding,
83    /// The statement is valid but reached none of the gossip-connected peers. A light
84    /// client keeps no store, so a statement nobody received is permanently lost.
85    NotSent { connected: usize },
86}
87
88impl StatementSubmitError {
89    /// Builds the JSON-RPC error response answering this failure.
90    ///
91    /// # Panic
92    ///
93    /// Panics if `request_id_json` isn't valid JSON.
94    pub fn to_json_rpc_error(&self, request_id_json: &str) -> String {
95        // Both messages carry polkadot-sdk's prefix. It answers every statement-store failure with
96        // the same code, so the message is all that tells these two apart.
97        let message = match self {
98            StatementSubmitError::InvalidEncoding => {
99                String::from("Statement store error: Error decoding statement")
100            }
101            StatementSubmitError::NotSent { connected: 0 } => String::from(
102                "Statement store error: No connected peers to broadcast the statement to",
103            ),
104            // Connected peers can still decline to queue, so the count separates the two.
105            StatementSubmitError::NotSent { connected } => format!(
106                "Statement store error: none of the {connected} connected peers accepted the \
107                 statement"
108            ),
109        };
110
111        parse::build_error_response(
112            request_id_json,
113            parse::ErrorResponse::ApplicationDefined(STATEMENT_STORE_ERROR_CODE, &message),
114            None,
115        )
116    }
117}
118
119/// Validates a SCALE-encoded statement and broadcasts it to the network.
120///
121/// The checks run in the order polkadot-sdk's `Store::submit` applies them — expiry, then size,
122/// then proof — so that a client submitting a statement failing several of them is told the same
123/// reason a full node would give. Checks needing a local store or chain state are skipped.
124///
125/// The `broadcast` closure is only called if the statement is valid.
126pub async fn validate_and_broadcast_statement<F, Fut>(
127    encoded: &[u8],
128    now_from_unix_epoch: Duration,
129    broadcast: F,
130) -> Result<StatementSubmitResult, StatementSubmitError>
131where
132    F: FnOnce(Vec<u8>) -> Fut,
133    Fut: core::future::Future<Output = BroadcastStatementResult>,
134{
135    let Ok(statement) = codec::decode_statement(encoded) else {
136        return Err(StatementSubmitError::InvalidEncoding);
137    };
138
139    if is_expired(&statement, now_from_unix_epoch) {
140        return Ok(StatementSubmitResult::Invalid(
141            InvalidReason::AlreadyExpired,
142        ));
143    }
144
145    if encoded.len() > codec::MAX_STATEMENT_SIZE {
146        return Ok(StatementSubmitResult::Invalid(
147            InvalidReason::EncodingTooLarge {
148                submitted_size: encoded.len(),
149                max_size: codec::MAX_STATEMENT_SIZE,
150            },
151        ));
152    }
153
154    if statement.proof.is_none() {
155        return Ok(StatementSubmitResult::Invalid(InvalidReason::NoProof));
156    }
157
158    let broadcasted = broadcast(encoded.to_vec()).await;
159    if broadcasted.sent == 0 {
160        return Err(StatementSubmitError::NotSent {
161            connected: broadcasted.total,
162        });
163    }
164
165    Ok(StatementSubmitResult::New)
166}
167
168/// Whether a statement received from a peer may be delivered to subscriptions: the checks of
169/// `statement_submit` that need no store, expiry and presence of a proof.
170pub(super) fn is_deliverable(statement: &codec::Statement, now_from_unix_epoch: Duration) -> bool {
171    !is_expired(statement, now_from_unix_epoch) && statement.proof.is_some()
172}
173
174/// The most significant 32 bits of `expiry` hold a UNIX timestamp in seconds. A statement expiring
175/// exactly now is already expired, as in polkadot-sdk.
176fn is_expired(statement: &codec::Statement, now_from_unix_epoch: Duration) -> bool {
177    now_from_unix_epoch.as_secs() >= statement.expiry >> 32
178}
179
180pub(super) struct StatementSubscription {
181    topic_filter: TopicFilter,
182    seen: Option<lru::LruCache<[u8; 32], (), fnv::FnvBuildHasher>>,
183}
184
185impl StatementSubscription {
186    pub(super) fn new(topic_filter: TopicFilter, max_seen: Option<NonZero<usize>>) -> Self {
187        Self {
188            topic_filter,
189            seen: max_seen
190                .map(|cap| lru::LruCache::with_hasher(cap, fnv::FnvBuildHasher::default())),
191        }
192    }
193
194    pub(super) fn accept(&mut self, hash: &[u8; 32], statement: &codec::Statement) -> bool {
195        if !self.topic_filter.matches(&statement.topics) {
196            return false;
197        }
198        if let Some(seen) = &mut self.seen {
199            if seen.put(*hash, ()).is_some() {
200                return false;
201            }
202        }
203        true
204    }
205}
206
207/// Set of active statement subscriptions together with a reverse index mapping each topic to the
208/// subscriptions that reference it.
209///
210/// The reverse index lets statement matching scale with the number of subscriptions that share a
211/// topic with the incoming statement, rather than with the total number of subscriptions.
212pub(super) struct StatementSubscriptions {
213    /// Maps subscription ID to its state.
214    subscriptions: hashbrown::HashMap<String, StatementSubscription, fnv::FnvBuildHasher>,
215
216    /// Reverse index: maps a topic to the IDs of all subscriptions whose filter references it.
217    /// Only populated for `MatchAny`/`MatchAll` filters with a non-empty topic list.
218    by_topic: hashbrown::HashMap<
219        [u8; 32],
220        hashbrown::HashSet<String, fnv::FnvBuildHasher>,
221        fnv::FnvBuildHasher,
222    >,
223
224    /// IDs of subscriptions that match every statement irrespective of its topics: either
225    /// `TopicFilter::Any`, or a `TopicFilter::MatchAll` whose topic list is empty.
226    wildcard: hashbrown::HashSet<String, fnv::FnvBuildHasher>,
227}
228
229impl StatementSubscriptions {
230    pub(super) fn with_capacity(capacity: usize) -> Self {
231        Self {
232            subscriptions: hashbrown::HashMap::with_capacity_and_hasher(
233                capacity,
234                Default::default(),
235            ),
236            by_topic: hashbrown::HashMap::with_hasher(Default::default()),
237            wildcard: hashbrown::HashSet::with_hasher(Default::default()),
238        }
239    }
240
241    pub(super) fn is_empty(&self) -> bool {
242        self.subscriptions.is_empty()
243    }
244
245    /// Inserts a new subscription and updates the reverse index.
246    pub(super) fn insert(
247        &mut self,
248        id: String,
249        topic_filter: TopicFilter,
250        max_seen: Option<NonZero<usize>>,
251    ) {
252        match &topic_filter {
253            TopicFilter::Any => {
254                self.wildcard.insert(id.clone());
255            }
256            // An empty `MatchAll` filter matches every statement.
257            TopicFilter::MatchAll(topics) if topics.is_empty() => {
258                self.wildcard.insert(id.clone());
259            }
260            TopicFilter::MatchAll(topics) | TopicFilter::MatchAny(topics) => {
261                for topic in topics {
262                    self.by_topic
263                        .entry(*topic)
264                        .or_insert_with(|| hashbrown::HashSet::with_hasher(Default::default()))
265                        .insert(id.clone());
266                }
267            }
268        }
269
270        self.subscriptions
271            .insert(id, StatementSubscription::new(topic_filter, max_seen));
272    }
273
274    /// Removes a subscription and cleans up the reverse index. Returns whether it existed.
275    pub(super) fn remove(&mut self, id: &str) -> bool {
276        let Some(sub) = self.subscriptions.remove(id) else {
277            return false;
278        };
279
280        match &sub.topic_filter {
281            TopicFilter::Any => {
282                self.wildcard.remove(id);
283            }
284            TopicFilter::MatchAll(topics) if topics.is_empty() => {
285                self.wildcard.remove(id);
286            }
287            TopicFilter::MatchAll(topics) | TopicFilter::MatchAny(topics) => {
288                for topic in topics {
289                    if let Some(ids) = self.by_topic.get_mut(topic) {
290                        ids.remove(id);
291                        if ids.is_empty() {
292                            self.by_topic.remove(topic);
293                        }
294                    }
295                }
296            }
297        }
298
299        true
300    }
301
302    pub(super) fn shrink_to_fit(&mut self) {
303        self.subscriptions.shrink_to_fit();
304        for ids in self.by_topic.values_mut() {
305            ids.shrink_to_fit();
306        }
307        self.by_topic.shrink_to_fit();
308        self.wildcard.shrink_to_fit();
309    }
310
311    /// Matches a batch of statements against the subscriptions.
312    ///
313    /// Returns, for every subscription that accepts at least one statement, the list of re-encoded
314    /// matching statements. Uses the reverse index to only consider subscriptions that either match
315    /// everything or share a topic with the statement; the precise per-subscription filter and
316    /// deduplication is then applied via [`StatementSubscription::accept`].
317    pub(super) fn matching(
318        &mut self,
319        statements: &[([u8; 32], codec::Statement)],
320    ) -> Vec<(String, Vec<HexString>)> {
321        // Disjoint borrows: `subscriptions` is mutated while `by_topic`/`wildcard` are only read.
322        let Self {
323            subscriptions,
324            by_topic,
325            wildcard,
326        } = self;
327
328        // Subscription ID -> its matching re-encoded statements.
329        let mut out: hashbrown::HashMap<&str, Vec<HexString>, fnv::FnvBuildHasher> =
330            hashbrown::HashMap::with_hasher(Default::default());
331        // Reused across statements to avoid reallocating.
332        let mut candidates: hashbrown::HashSet<&str, fnv::FnvBuildHasher> =
333            hashbrown::HashSet::with_hasher(Default::default());
334
335        for (hash, statement) in statements {
336            candidates.clear();
337            candidates.extend(wildcard.iter().map(String::as_str));
338            for topic in &statement.topics {
339                if let Some(ids) = by_topic.get(topic) {
340                    candidates.extend(ids.iter().map(String::as_str));
341                }
342            }
343
344            // Re-encoded lazily on first match and reused for every matching subscription.
345            let mut encoded: Option<HexString> = None;
346            for id in &candidates {
347                let sub = subscriptions
348                    .get_mut(*id)
349                    .expect("`candidates` is a subset of `subscriptions`; qed");
350                if sub.accept(hash, statement) {
351                    let encoded = encoded.get_or_insert_with(|| {
352                        HexString(
353                            codec::encode_statement(statement)
354                                .expect("re-encoding a decoded statement always succeeds; qed"),
355                        )
356                    });
357                    out.entry(*id).or_default().push(encoded.clone());
358                }
359            }
360        }
361
362        out.into_iter()
363            .map(|(id, matching)| (String::from(id), matching))
364            .collect()
365    }
366
367    pub(super) fn build_combined_affinity_filter(
368        &self,
369        config: &StatementProtocolConfig,
370    ) -> network_service::AffinityFilter {
371        let mut all_topics: Vec<&[u8; 32]> = Vec::new();
372
373        for sub in self.subscriptions.values() {
374            match &sub.topic_filter {
375                TopicFilter::Any => {
376                    return network_service::AffinityFilter::match_all(config.bloom_seed());
377                }
378                TopicFilter::MatchAll(topics) | TopicFilter::MatchAny(topics) => {
379                    all_topics.extend(topics.iter());
380                }
381            }
382        }
383
384        network_service::AffinityFilter::from_topics(
385            all_topics.into_iter(),
386            config.bloom_seed(),
387            config.false_positive_rate(),
388        )
389    }
390}
391
392#[cfg(test)]
393mod tests {
394    use super::*;
395    use alloc::string::ToString as _;
396    use core::time::Duration;
397    use futures_lite::future::block_on;
398
399    const SEED: u128 = 0x5EED_5EED_5EED_5EED_5EED_5EED_5EED_5EED;
400    const FPR: f64 = 0.01;
401
402    fn test_config() -> StatementProtocolConfig {
403        StatementProtocolConfig::new(
404            NonZero::new(128).unwrap(),
405            FPR,
406            SEED,
407            Duration::from_secs(1),
408        )
409    }
410
411    fn make_subscriptions(
412        entries: Vec<(&str, TopicFilter, Option<NonZero<usize>>)>,
413    ) -> StatementSubscriptions {
414        let mut subs = StatementSubscriptions::with_capacity(entries.len());
415        for (id, filter, max_seen) in entries {
416            subs.insert(id.to_string(), filter, max_seen);
417        }
418        subs
419    }
420
421    fn statement_with_topics(topics: Vec<[u8; 32]>) -> codec::Statement {
422        codec::Statement {
423            proof: None,
424            decryption_key: None,
425            expiry: 42,
426            channel: None,
427            topics,
428            data: None,
429        }
430    }
431
432    const NOW: Duration = Duration::from_secs(1_000);
433
434    /// Expiration timestamp, in the most significant 32 bits, later than [`NOW`].
435    const FUTURE_EXPIRY: u64 = 2_000 << 32;
436
437    fn encoded_statement(with_proof: bool, expiry: u64, data: Option<Vec<u8>>) -> Vec<u8> {
438        codec::encode_statement(&codec::Statement {
439            proof: with_proof.then(|| codec::Proof::Ed25519 {
440                signature: [0; 64],
441                signer: [0; 32],
442            }),
443            decryption_key: None,
444            expiry,
445            channel: None,
446            topics: Vec::new(),
447            data,
448        })
449        .unwrap()
450    }
451
452    #[test]
453    fn validate_and_broadcast_invalid_encoding() {
454        let result = block_on(validate_and_broadcast_statement(
455            &[0xff, 0xff],
456            NOW,
457            |_| async { unreachable!() },
458        ));
459        assert_eq!(result, Err(StatementSubmitError::InvalidEncoding));
460    }
461
462    #[test]
463    fn validate_and_broadcast_already_expired() {
464        // The statement also has no proof: the expiry check runs first.
465        let encoded = encoded_statement(false, 500 << 32, None);
466        let result = block_on(validate_and_broadcast_statement(&encoded, NOW, |_| async {
467            unreachable!()
468        }));
469        assert_eq!(
470            result,
471            Ok(StatementSubmitResult::Invalid(
472                InvalidReason::AlreadyExpired
473            ))
474        );
475    }
476
477    #[test]
478    fn validate_and_broadcast_expiry_equal_to_now_is_expired() {
479        let encoded = encoded_statement(true, NOW.as_secs() << 32, None);
480        let result = block_on(validate_and_broadcast_statement(&encoded, NOW, |_| async {
481            unreachable!()
482        }));
483        assert_eq!(
484            result,
485            Ok(StatementSubmitResult::Invalid(
486                InvalidReason::AlreadyExpired
487            ))
488        );
489    }
490
491    #[test]
492    fn validate_and_broadcast_encoding_too_large() {
493        // The statement also has no proof: the size check runs before the proof check.
494        let encoded = encoded_statement(false, FUTURE_EXPIRY, Some(vec![0; 1024 * 1024]));
495        assert!(encoded.len() > codec::MAX_STATEMENT_SIZE);
496        let result = block_on(validate_and_broadcast_statement(&encoded, NOW, |_| async {
497            unreachable!()
498        }));
499        assert_eq!(
500            result,
501            Ok(StatementSubmitResult::Invalid(
502                InvalidReason::EncodingTooLarge {
503                    submitted_size: encoded.len(),
504                    max_size: codec::MAX_STATEMENT_SIZE,
505                }
506            ))
507        );
508    }
509
510    #[test]
511    fn validate_and_broadcast_no_proof() {
512        let encoded = encoded_statement(false, FUTURE_EXPIRY, None);
513        let result = block_on(validate_and_broadcast_statement(&encoded, NOW, |_| async {
514            unreachable!()
515        }));
516        assert_eq!(
517            result,
518            Ok(StatementSubmitResult::Invalid(InvalidReason::NoProof))
519        );
520    }
521
522    fn decoded_statement(with_proof: bool, expiry: u64) -> codec::Statement {
523        codec::decode_statement(&encoded_statement(with_proof, expiry, None)).unwrap()
524    }
525
526    #[test]
527    fn is_deliverable_requires_proof_and_future_expiry() {
528        assert!(is_deliverable(&decoded_statement(true, FUTURE_EXPIRY), NOW));
529        assert!(!is_deliverable(
530            &decoded_statement(false, FUTURE_EXPIRY),
531            NOW
532        ));
533        assert!(!is_deliverable(&decoded_statement(true, 500 << 32), NOW));
534        // Expiring exactly now counts as expired, as in `validate_and_broadcast_statement`.
535        assert!(!is_deliverable(
536            &decoded_statement(true, NOW.as_secs() << 32),
537            NOW
538        ));
539    }
540
541    #[test]
542    fn validate_and_broadcast_no_peers() {
543        let encoded = encoded_statement(true, FUTURE_EXPIRY, None);
544        let result = block_on(validate_and_broadcast_statement(&encoded, NOW, |_| async {
545            BroadcastStatementResult { sent: 0, total: 0 }
546        }));
547        assert_eq!(result, Err(StatementSubmitError::NotSent { connected: 0 }));
548    }
549
550    #[test]
551    fn submit_errors_carry_the_statement_store_code() {
552        // Both failures answer with polkadot-sdk's single statement-store code, telling themselves
553        // apart by message alone, exactly as it does.
554        assert_eq!(
555            StatementSubmitError::InvalidEncoding.to_json_rpc_error("7"),
556            r#"{"jsonrpc":"2.0","id":7,"error":{"code":7001,"message":"Statement store error: Error decoding statement"}}"#
557        );
558        assert_eq!(
559            StatementSubmitError::NotSent { connected: 0 }.to_json_rpc_error("7"),
560            r#"{"jsonrpc":"2.0","id":7,"error":{"code":7001,"message":"Statement store error: No connected peers to broadcast the statement to"}}"#
561        );
562        // Connected peers that took nothing must not read as no peers.
563        assert_eq!(
564            StatementSubmitError::NotSent { connected: 5 }.to_json_rpc_error("7"),
565            r#"{"jsonrpc":"2.0","id":7,"error":{"code":7001,"message":"Statement store error: none of the 5 connected peers accepted the statement"}}"#
566        );
567    }
568
569    #[test]
570    fn validate_and_broadcast_reaching_no_peer_is_not_new() {
571        // Gossip-connected peers whose statement substream is missing or whose queue is full leave
572        // the statement unsent. Answering `new` would tell the client it was published.
573        let encoded = encoded_statement(true, FUTURE_EXPIRY, None);
574        let result = block_on(validate_and_broadcast_statement(&encoded, NOW, |_| async {
575            BroadcastStatementResult { sent: 0, total: 5 }
576        }));
577        assert_eq!(result, Err(StatementSubmitError::NotSent { connected: 5 }));
578    }
579
580    #[test]
581    fn validate_and_broadcast_new() {
582        let encoded = encoded_statement(true, FUTURE_EXPIRY, None);
583        let result = block_on(validate_and_broadcast_statement(&encoded, NOW, |_| async {
584            BroadcastStatementResult { sent: 3, total: 5 }
585        }));
586        assert_eq!(result, Ok(StatementSubmitResult::New));
587    }
588
589    #[test]
590    fn build_combined_affinity_empty_subscriptions() {
591        let config = test_config();
592        let subs = make_subscriptions(vec![]);
593        let filter = subs.build_combined_affinity_filter(&config);
594
595        // Empty subscription set: no topics are ever in the filter.
596        assert!(!filter.contains(&[1u8; 32]));
597        // A statement with no topics (broadcast) still matches.
598        let broadcast: &[&[u8; 32]] = &[];
599        assert!(filter.matches_statement(broadcast));
600    }
601
602    #[test]
603    fn build_combined_affinity_any_filter_matches_everything() {
604        let config = test_config();
605        let subs = make_subscriptions(vec![("s", TopicFilter::Any, None)]);
606        let filter = subs.build_combined_affinity_filter(&config);
607
608        // TopicFilter::Any returns the broadcast `match_all` filter: every topic matches.
609        assert!(filter.contains(&[1u8; 32]));
610        assert!(filter.contains(&[99u8; 32]));
611        let t = [7u8; 32];
612        assert!(filter.matches_statement(&[&t]));
613    }
614
615    #[test]
616    fn build_combined_affinity_match_any_union() {
617        let config = test_config();
618        let t1 = [1u8; 32];
619        let t2 = [2u8; 32];
620        let subs = make_subscriptions(vec![
621            ("a", TopicFilter::match_any(vec![t1]).unwrap(), None),
622            ("b", TopicFilter::match_any(vec![t2]).unwrap(), None),
623        ]);
624        let filter = subs.build_combined_affinity_filter(&config);
625
626        assert!(filter.contains(&t1));
627        assert!(filter.contains(&t2));
628    }
629
630    #[test]
631    fn accept_fresh_statement_passes() {
632        let t1 = [1u8; 32];
633        let mut sub =
634            StatementSubscription::new(TopicFilter::match_any(vec![t1]).unwrap(), NonZero::new(8));
635        let stmt = statement_with_topics(vec![t1]);
636        assert!(sub.accept(&[0xbb; 32], &stmt));
637    }
638
639    #[test]
640    fn accept_duplicate_returns_false() {
641        let mut sub = StatementSubscription::new(TopicFilter::Any, NonZero::new(8));
642        let stmt = statement_with_topics(vec![]);
643        let hash = [0xcc; 32];
644        assert!(sub.accept(&hash, &stmt));
645        assert!(!sub.accept(&hash, &stmt));
646    }
647
648    #[test]
649    fn accept_lru_eviction_allows_resubmit() {
650        let mut sub = StatementSubscription::new(TopicFilter::Any, NonZero::new(2));
651        let stmt = statement_with_topics(vec![]);
652        let h_a = [0xa; 32];
653        let h_b = [0xb; 32];
654        let h_c = [0xc; 32];
655
656        assert!(sub.accept(&h_a, &stmt));
657        assert!(sub.accept(&h_b, &stmt));
658        // Inserting a third eviction-capacity 2 item evicts h_a (oldest).
659        assert!(sub.accept(&h_c, &stmt));
660        // h_a was evicted: it is accepted again as if fresh.
661        assert!(sub.accept(&h_a, &stmt));
662    }
663
664    #[test]
665    fn dedup_is_per_subscription() {
666        let mut sub_a = StatementSubscription::new(TopicFilter::Any, NonZero::new(8));
667        let mut sub_b = StatementSubscription::new(TopicFilter::Any, NonZero::new(8));
668        let stmt = statement_with_topics(vec![]);
669        let hash = [0xee; 32];
670
671        assert!(sub_a.accept(&hash, &stmt));
672        assert!(!sub_a.accept(&hash, &stmt));
673        // Same hash on a different subscription is still fresh: caches are independent.
674        assert!(sub_b.accept(&hash, &stmt));
675    }
676
677    /// Builds a `(hash, statement)` batch entry from a list of topics.
678    fn batch_entry(hash: u8, topics: Vec<[u8; 32]>) -> ([u8; 32], codec::Statement) {
679        ([hash; 32], statement_with_topics(topics))
680    }
681
682    /// Collects the IDs of all subscriptions that matched at least once.
683    fn matched_ids(matches: &[(String, Vec<HexString>)]) -> Vec<String> {
684        let mut ids: Vec<String> = matches.iter().map(|(id, _)| id.clone()).collect();
685        ids.sort();
686        ids
687    }
688
689    #[test]
690    fn matching_match_any_only_returns_relevant_subscriptions() {
691        let t1 = [1u8; 32];
692        let t2 = [2u8; 32];
693        let mut subs = make_subscriptions(vec![
694            ("a", TopicFilter::match_any(vec![t1]).unwrap(), None),
695            ("b", TopicFilter::match_any(vec![t2]).unwrap(), None),
696        ]);
697
698        // A statement carrying only `t1` must match `a` and not `b`.
699        let matches = subs.matching(&[batch_entry(0xaa, vec![t1])]);
700        assert_eq!(matched_ids(&matches), vec!["a".to_string()]);
701
702        // A statement with an unrelated topic matches nothing.
703        let matches = subs.matching(&[batch_entry(0xbb, vec![[9u8; 32]])]);
704        assert!(matches.is_empty());
705    }
706
707    #[test]
708    fn matching_wildcard_filters_match_every_statement() {
709        // `Any` and an empty `MatchAll` both match every statement, with or without topics.
710        let mut subs = make_subscriptions(vec![
711            ("any", TopicFilter::Any, None),
712            ("all", TopicFilter::match_all(vec![]).unwrap(), None),
713        ]);
714
715        let matches = subs.matching(&[batch_entry(0x01, vec![[7u8; 32]])]);
716        assert_eq!(
717            matched_ids(&matches),
718            vec!["all".to_string(), "any".to_string()]
719        );
720
721        let matches = subs.matching(&[batch_entry(0x02, vec![])]);
722        assert_eq!(
723            matched_ids(&matches),
724            vec!["all".to_string(), "any".to_string()]
725        );
726    }
727
728    #[test]
729    fn matching_match_all_requires_every_topic() {
730        let t1 = [1u8; 32];
731        let t2 = [2u8; 32];
732        let mut subs = make_subscriptions(vec![(
733            "all",
734            TopicFilter::match_all(vec![t1, t2]).unwrap(),
735            None,
736        )]);
737
738        // A statement carrying only one of the required topics is a candidate via the reverse
739        // index but must be rejected by the precise re-check.
740        let matches = subs.matching(&[batch_entry(0xaa, vec![t1])]);
741        assert!(matches.is_empty());
742
743        // A statement carrying both topics matches.
744        let matches = subs.matching(&[batch_entry(0xbb, vec![t1, t2])]);
745        assert_eq!(matched_ids(&matches), vec!["all".to_string()]);
746    }
747
748    #[test]
749    fn matching_empty_match_any_never_matches() {
750        let mut subs = make_subscriptions(vec![(
751            "none",
752            TopicFilter::match_any(vec![]).unwrap(),
753            None,
754        )]);
755
756        let matches = subs.matching(&[batch_entry(0x01, vec![[1u8; 32]])]);
757        assert!(matches.is_empty());
758        let matches = subs.matching(&[batch_entry(0x02, vec![])]);
759        assert!(matches.is_empty());
760    }
761
762    #[test]
763    fn matching_dedup_applies_across_batches() {
764        let t1 = [1u8; 32];
765        let mut subs = make_subscriptions(vec![(
766            "a",
767            TopicFilter::match_any(vec![t1]).unwrap(),
768            NonZero::new(8),
769        )]);
770
771        let entry = batch_entry(0xaa, vec![t1]);
772        let matches = subs.matching(&[entry.clone()]);
773        assert_eq!(matches.len(), 1);
774        assert_eq!(matches[0].1.len(), 1);
775
776        // The same statement hash is deduplicated and produces no further notification.
777        let matches = subs.matching(&[entry]);
778        assert!(matches.is_empty());
779    }
780
781    #[test]
782    fn matching_groups_multiple_statements_per_subscription() {
783        let t1 = [1u8; 32];
784        let mut subs =
785            make_subscriptions(vec![("a", TopicFilter::match_any(vec![t1]).unwrap(), None)]);
786
787        let matches = subs.matching(&[batch_entry(0x01, vec![t1]), batch_entry(0x02, vec![t1])]);
788        assert_eq!(matches.len(), 1);
789        assert_eq!(matches[0].0, "a");
790        assert_eq!(matches[0].1.len(), 2);
791    }
792
793    #[test]
794    fn remove_cleans_reverse_index() {
795        let t1 = [1u8; 32];
796        let mut subs =
797            make_subscriptions(vec![("a", TopicFilter::match_any(vec![t1]).unwrap(), None)]);
798
799        assert!(subs.remove("a"));
800        assert!(!subs.remove("a"));
801        assert!(subs.is_empty());
802        // The topic entry must have been cleaned up, so a matching statement finds nothing.
803        let matches = subs.matching(&[batch_entry(0xaa, vec![t1])]);
804        assert!(matches.is_empty());
805    }
806}