weida 0.1.0-alpha.1

QUIC-native messaging framework: runtime, native QUIC transport, Req/Rep, Push/Pull, Pub/Sub, PAIR, SURVEY and BUS
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
//! Publisher-side subscription registry and fan-out.
//!
//! A publisher does not know its subscribers: they arrive as SUBSCRIBE frames
//! on whatever connections the listener has accepted. This module owns that
//! mapping and the one genuinely new selection policy in Phase 3 — fan-out
//! (`docs/ARCHITECTURE.md`, pattern primitive P3). Everything else a published
//! message needs is the ordinary one-way transfer path.
//!
//! **Delivery is per-subscriber best effort.** Each subscriber has a byte
//! budget; a message that does not fit is dropped for that subscriber, counted,
//! and the publisher moves on. A slow consumer therefore costs the publisher
//! nothing but its own messages (master doc §17: telemetry fan-out tolerates
//! drops and cannot tolerate a stalled producer).

use std::collections::{HashMap, HashSet};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, RwLock};

use bytes::Bytes;
use tokio::sync::mpsc::OwnedPermit;
use tokio::sync::{Semaphore, mpsc};
use weida_core::{Error, Limits, TraceContext};
use weida_protocol::header::OrderingMode;
use weida_protocol::{DataHeader, filter};

use crate::conn::ConnHandle;
use crate::ordering::Sequencer;
use crate::transfer::write_data_preamble;
use crate::transport::SendHalf;
use weida_protocol::codes;

/// Messages one subscriber's writer task may hold. The real bound is the byte
/// budget; this only keeps the channel itself from growing without limit when
/// messages are tiny.
const WRITER_QUEUE: usize = 1024;

/// One message queued for one subscriber.
///
/// The DATA header is *not* pre-encoded and shared: each copy carries the
/// subscriber's own path. The payload is a `Bytes`, so the fan-out shares one
/// allocation regardless of subscriber count.
struct PubMsg {
    topic: Arc<str>,
    payload: Bytes,
    /// The trace context the caller propagated, if any: weida mints none
    /// ([0028](../../../docs/decisions/0028-trace-propagation-is-the-callers.md)).
    trace: Option<TraceContext>,
    /// The producer's sequence number for this topic, assigned once per
    /// published message. `None` unless `PerProducer` ordering is negotiated.
    sequence: Option<u64>,
}

/// What one subscriber's writer is told to do.
///
/// A whole message is one item, which is what [`SubRegistry::publish`] sends.
/// A **streamed** publish ([`SubRegistry::open`], B-064) is `Begin`, then a
/// `Chunk` per piece, then `Finish` — so a payload the publisher never
/// materializes still becomes one stream per subscriber, and the per-chunk
/// byte budget is what bounds this side rather than the payload's size.
enum PubItem {
    /// A message whose payload is already in memory: open, write, finish.
    Whole(PubMsg),
    /// Opens a stream for a streamed transfer and writes its DATA header,
    /// which carries no `content_len` because nobody knows it yet.
    Begin { id: u64, head: PubMsg },
    /// The next piece of the streamed transfer `id`.
    Chunk { id: u64, payload: Bytes },
    /// The end of `id`: FIN, and the receipt parked for the drain.
    Finish { id: u64 },
    /// `id` will not be completed — a chunk did not fit this subscriber's
    /// budget, or the publisher dropped the handle. The stream is reset, so
    /// the subscriber discards a partial payload instead of waiting for a FIN
    /// that will never arrive.
    Abort { id: u64 },
}

/// One subscribing connection's state for one publisher path.
struct SubEntry {
    /// `quinn::Connection::stable_id`, the identity a SUBSCRIBE arrives with.
    conn_id: usize,
    filters: HashSet<String>,
    tx: mpsc::Sender<PubItem>,
    /// Payload bytes this subscriber may hold queued. Permits are taken by
    /// `publish` and returned by the writer once the bytes are on the wire.
    budget: Arc<Semaphore>,
}

/// Why a published copy was dropped.
///
/// Three causes, because they are three different failures: the first two
/// are the subscriber not keeping up, the third is the *local* transport
/// having no connection to carry the copy
/// ([decisions/0012](../../../docs/decisions/0012-local-connection-grouping.md)
/// §4.4).
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub(crate) enum DropCause {
    /// The subscriber's byte budget (`Limits::subscriber_buffer_bytes`) had
    /// no room for the payload.
    SubscriberBudget,
    /// The subscriber's writer queue was full of messages.
    SubscriberQueue,
    /// The subscriber had parked no connection for this copy: a socket
    /// transport's fan-out rides connections the subscriber parks, and the
    /// pool was empty.
    NoParkedConnection,
}

/// What a publisher dropped on one topic, by cause.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct TopicDrops {
    /// The topic the dropped copies were published on.
    pub topic: Arc<str>,
    /// Copies dropped for an exhausted subscriber byte budget.
    pub subscriber_budget: u64,
    /// Copies dropped for a full subscriber queue.
    pub subscriber_queue: u64,
    /// Copies dropped because the subscriber had parked no connection.
    pub no_parked_connection: u64,
}

impl TopicDrops {
    /// All three causes summed.
    pub fn total(&self) -> u64 {
        self.subscriber_budget + self.subscriber_queue + self.no_parked_connection
    }
}

/// Three counters, one per cause.
#[derive(Default)]
struct Causes {
    budget: AtomicU64,
    queue: AtomicU64,
    no_parked: AtomicU64,
}

impl Causes {
    fn record(&self, cause: DropCause) {
        let counter = match cause {
            DropCause::SubscriberBudget => &self.budget,
            DropCause::SubscriberQueue => &self.queue,
            DropCause::NoParkedConnection => &self.no_parked,
        };
        counter.fetch_add(1, Ordering::Relaxed);
    }

    fn snapshot(&self, topic: &Arc<str>) -> TopicDrops {
        TopicDrops {
            topic: Arc::clone(topic),
            subscriber_budget: self.budget.load(Ordering::Relaxed),
            subscriber_queue: self.queue.load(Ordering::Relaxed),
            no_parked_connection: self.no_parked.load(Ordering::Relaxed),
        }
    }
}

/// The drops on one publisher path: a total, and a count per topic and
/// cause.
///
/// The per-topic table is bounded by `Limits::max_sequence_scopes`, the same
/// ceiling the sequencer's per-topic table has and for the same reason: the
/// topics are the publishing application's, but a table nobody bounds is a
/// table that grows for the life of the process. At the cap a new topic's
/// drops count in the total and nowhere else, which is what the sequencer
/// does with a new scope.
struct DropTable {
    total: AtomicU64,
    per_topic: RwLock<HashMap<Arc<str>, Causes>>,
    max_topics: usize,
}

impl DropTable {
    fn new(max_topics: usize) -> DropTable {
        DropTable {
            total: AtomicU64::new(0),
            per_topic: RwLock::new(HashMap::new()),
            max_topics,
        }
    }

    /// Counts one dropped copy of a message on `topic`.
    ///
    /// A read lock on the common path — the topic is known — and a write
    /// lock only for a topic's first drop. Drops are the exceptional path of
    /// fan-out, so neither is on the path of a message that gets through.
    fn record(&self, topic: &Arc<str>, cause: DropCause) {
        self.total.fetch_add(1, Ordering::Relaxed);
        {
            let table = self.per_topic.read().expect("drop table poisoned");
            if let Some(causes) = table.get(topic) {
                causes.record(cause);
                return;
            }
        }
        let mut table = self.per_topic.write().expect("drop table poisoned");
        if let Some(causes) = table.get(topic) {
            causes.record(cause);
        } else if table.len() < self.max_topics {
            let causes = Causes::default();
            causes.record(cause);
            table.insert(Arc::clone(topic), causes);
        }
    }

    fn on_topic(&self, topic: &str) -> Option<TopicDrops> {
        let table = self.per_topic.read().expect("drop table poisoned");
        table
            .get_key_value(topic)
            .map(|(topic, causes)| causes.snapshot(topic))
    }

    fn by_topic(&self) -> Vec<TopicDrops> {
        let table = self.per_topic.read().expect("drop table poisoned");
        table
            .iter()
            .map(|(topic, causes)| causes.snapshot(topic))
            .collect()
    }
}

/// Everything registered for one publisher path.
struct PathState {
    subs: Vec<SubEntry>,
    /// Copies dropped for a subscriber that could not take them, by topic
    /// and cause.
    drops: Arc<DropTable>,
}

/// Subscriptions for every publisher path served by one listener.
pub(crate) struct SubRegistry {
    paths: RwLock<HashMap<Arc<str>, PathState>>,
    /// Filters held per connection, summed over paths, for `max_subscriptions`.
    per_conn: RwLock<HashMap<usize, usize>>,
    limits: Limits,
    /// Numbers published messages per topic. The number belongs to the
    /// *message*, not to a subscriber's copy, which is what makes a dropped
    /// copy visible as a gap: the survivors keep the numbers the lost ones
    /// would have had.
    sequencer: Sequencer,
    /// Names one streamed publish, so a writer can tell the chunks of two
    /// concurrent ones apart. Local and never on the wire: what identifies a
    /// transfer to a subscriber is the stream it arrives on.
    next_stream: AtomicU64,
}

impl SubRegistry {
    pub(crate) fn new(limits: Limits, ordering: OrderingMode) -> SubRegistry {
        SubRegistry {
            paths: RwLock::new(HashMap::new()),
            per_conn: RwLock::new(HashMap::new()),
            limits,
            sequencer: Sequencer::new(ordering),
            next_stream: AtomicU64::new(0),
        }
    }

    /// Records interest of `ctx` in `filter` on `path`.
    ///
    /// The path need not be registered by a publisher yet: a subscriber may
    /// arrive first, and the subscription is bounded by `max_subscriptions`
    /// either way. Duplicate filters are idempotent.
    ///
    /// Returns [`Error::LimitExceeded`] when the connection is already holding
    /// `max_subscriptions` filters; the caller closes the connection, because
    /// SUBSCRIBE carries no transfer id to answer with an ERROR frame.
    pub(crate) fn subscribe(
        &self,
        path: &str,
        ctx: &ConnHandle,
        filter: String,
    ) -> Result<(), Error> {
        let conn_id = ctx.conn.stable_id();
        let mut paths = self.paths.write().expect("subscription lock poisoned");
        let mut per_conn = self.per_conn.write().expect("subscription lock poisoned");

        let state = paths.entry(Arc::from(path)).or_insert_with(|| PathState {
            subs: Vec::new(),
            drops: Arc::new(DropTable::new(self.limits.max_sequence_scopes)),
        });
        if let Some(entry) = state.subs.iter_mut().find(|e| e.conn_id == conn_id) {
            if entry.filters.contains(&filter) {
                return Ok(());
            }
            let count = per_conn.entry(conn_id).or_insert(0);
            if *count >= self.limits.max_subscriptions {
                return Err(Error::LimitExceeded);
            }
            *count += 1;
            entry.filters.insert(filter);
            return Ok(());
        }

        let count = per_conn.entry(conn_id).or_insert(0);
        if *count >= self.limits.max_subscriptions {
            return Err(Error::LimitExceeded);
        }
        *count += 1;

        let (tx, rx) = mpsc::channel(WRITER_QUEUE);
        let budget = Arc::new(Semaphore::new(self.limits.subscriber_buffer_bytes));
        let drops = Arc::clone(&state.drops);
        state.subs.push(SubEntry {
            conn_id,
            filters: HashSet::from([filter]),
            tx,
            budget: Arc::clone(&budget),
        });
        // One writer task per (connection, path): it serializes this
        // subscriber's messages, which is what makes delivery FIFO per
        // subscriber even though each message rides its own QUIC stream.
        ctx.exec
            .spawn(writer(Arc::clone(ctx), Arc::from(path), rx, budget, drops));
        Ok(())
    }

    /// Counts one subscription against `max_subscriptions` without recording
    /// a filter.
    ///
    /// For a subscription this registry does not serve: an L2 queue keeps its
    /// consumers itself, but the bound is per connection over *all* paths, so
    /// the count has to stay in one place or it bounds nothing
    /// (`docs/PROTOCOL.md` §10).
    pub(crate) fn reserve(&self, conn_id: usize) -> Result<(), Error> {
        let mut per_conn = self.per_conn.write().expect("subscription lock poisoned");
        let count = per_conn.entry(conn_id).or_insert(0);
        if *count >= self.limits.max_subscriptions {
            return Err(Error::LimitExceeded);
        }
        *count += 1;
        Ok(())
    }

    /// Gives one reserved subscription back. An unknown connection is ignored:
    /// UNSUBSCRIBE is idempotent.
    pub(crate) fn release(&self, conn_id: usize) {
        let mut per_conn = self.per_conn.write().expect("subscription lock poisoned");
        decrement(&mut per_conn, conn_id, 1);
    }

    /// Withdraws one filter. An unknown filter, path or connection is ignored:
    /// UNSUBSCRIBE is idempotent and races legitimately with teardown.
    pub(crate) fn unsubscribe(&self, path: &str, conn_id: usize, filter: &str) {
        let mut paths = self.paths.write().expect("subscription lock poisoned");
        let mut per_conn = self.per_conn.write().expect("subscription lock poisoned");
        let Some(state) = paths.get_mut(path) else {
            return;
        };
        let Some(index) = state.subs.iter().position(|e| e.conn_id == conn_id) else {
            return;
        };
        if !state.subs[index].filters.remove(filter) {
            return;
        }
        decrement(&mut per_conn, conn_id, 1);
        if state.subs[index].filters.is_empty() {
            // The writer task ends when its channel closes.
            state.subs.swap_remove(index);
        }
    }

    /// Drops every subscription held by one connection. Called when the
    /// connection closes, so a churn of peers cannot grow the registry.
    pub(crate) fn remove_connection(&self, conn_id: usize) {
        let mut paths = self.paths.write().expect("subscription lock poisoned");
        let mut per_conn = self.per_conn.write().expect("subscription lock poisoned");
        for state in paths.values_mut() {
            if let Some(index) = state.subs.iter().position(|e| e.conn_id == conn_id) {
                state.subs.swap_remove(index);
            }
        }
        per_conn.remove(&conn_id);
    }

    /// Fans `payload` out to every subscriber of `path` whose filter matches
    /// `topic`, and returns how many it was enqueued to.
    ///
    /// Synchronous and non-blocking by design: a subscriber that cannot take
    /// the bytes right now loses this message rather than stalling everyone
    /// else. `want` is the payload length in permits.
    pub(crate) fn publish(
        &self,
        path: &str,
        topic: &str,
        payload: Bytes,
        trace: Option<TraceContext>,
        want: u32,
    ) -> usize {
        let paths = self.paths.read().expect("subscription lock poisoned");
        let Some(state) = paths.get(path) else {
            return 0;
        };

        let topic: Arc<str> = Arc::from(topic);
        // Before fan-out, so every copy of one message carries one number and
        // a copy dropped below leaves a hole rather than renumbering.
        let sequence = self.sequencer.next(&topic);
        let mut sent = 0usize;
        for entry in &state.subs {
            if !entry.filters.iter().any(|f| filter::matches(&topic, f)) {
                continue;
            }
            let Ok(permit) = entry.budget.try_acquire_many(want) else {
                state.drops.record(&topic, DropCause::SubscriberBudget);
                tracing::debug!(path, %topic, "subscriber budget exhausted; message dropped");
                continue;
            };
            let msg = PubMsg {
                topic: Arc::clone(&topic),
                payload: payload.clone(),
                trace,
                sequence,
            };
            match entry.tx.try_send(PubItem::Whole(msg)) {
                Ok(()) => {
                    // The writer returns the permits once the bytes are gone.
                    permit.forget();
                    sent += 1;
                }
                Err(_) => {
                    // Dropping `permit` returns the budget immediately.
                    state.drops.record(&topic, DropCause::SubscriberQueue);
                    tracing::debug!(path, %topic, "subscriber queue full; message dropped");
                }
            }
        }
        sent
    }

    /// Opens a streamed publish: one stream per matched subscriber, written
    /// chunk by chunk and never materialized whole.
    ///
    /// The subscriber set is fixed here, at `open`, because a stream is a
    /// stream: a subscriber that arrives mid-payload would receive a fragment
    /// and could not be told where it started. It gets the next message.
    pub(crate) fn open(&self, path: &str, topic: &str, trace: Option<TraceContext>) -> FanOut {
        let paths = self.paths.read().expect("subscription lock poisoned");
        let topic: Arc<str> = Arc::from(topic);
        let id = self.next_stream.fetch_add(1, Ordering::Relaxed);
        let Some(state) = paths.get(path) else {
            return FanOut::empty(id, topic);
        };
        // One number for the message, as `publish` does: a subscriber whose
        // copy is aborted below leaves a hole rather than renumbering.
        let sequence = self.sequencer.next(&topic);
        let mut targets = Vec::new();
        for entry in &state.subs {
            if !entry.filters.iter().any(|f| filter::matches(&topic, f)) {
                continue;
            }
            // The permit for the ending — `Finish` or `Abort` — is taken
            // before the first chunk, so ending a transfer can never fail for
            // want of queue room. Without it a subscriber whose queue filled
            // mid-payload would hold an open stream waiting for a FIN nobody
            // could enqueue.
            let Ok(ending) = entry.tx.clone().try_reserve_owned() else {
                state.drops.record(&topic, DropCause::SubscriberQueue);
                continue;
            };
            let head = PubMsg {
                topic: Arc::clone(&topic),
                payload: Bytes::new(),
                trace,
                sequence,
            };
            if entry.tx.try_send(PubItem::Begin { id, head }).is_err() {
                state.drops.record(&topic, DropCause::SubscriberQueue);
                continue;
            }
            targets.push(Target {
                tx: entry.tx.clone(),
                budget: Arc::clone(&entry.budget),
                ending: Some(ending),
            });
        }
        FanOut {
            id,
            topic,
            targets,
            drops: Some(Arc::clone(&state.drops)),
            finished: false,
        }
    }

    /// Connections currently subscribed to `path`.
    pub(crate) fn subscriber_count(&self, path: &str) -> usize {
        self.paths
            .read()
            .expect("subscription lock poisoned")
            .get(path)
            .map_or(0, |s| s.subs.len())
    }

    /// Filters registered on `path`, summed over subscribers.
    pub(crate) fn filter_count(&self, path: &str) -> usize {
        self.paths
            .read()
            .expect("subscription lock poisoned")
            .get(path)
            .map_or(0, |s| s.subs.iter().map(|e| e.filters.len()).sum())
    }

    /// Messages dropped on `path` because a subscriber could not take them,
    /// summed over topics and causes.
    pub(crate) fn dropped(&self, path: &str) -> u64 {
        self.paths
            .read()
            .expect("subscription lock poisoned")
            .get(path)
            .map_or(0, |s| s.drops.total.load(Ordering::Relaxed))
    }

    /// The drops on `path` for one topic, or `None` if none was counted.
    pub(crate) fn dropped_on(&self, path: &str, topic: &str) -> Option<TopicDrops> {
        self.paths
            .read()
            .expect("subscription lock poisoned")
            .get(path)
            .and_then(|s| s.drops.on_topic(topic))
    }

    /// The drops on `path`, one entry per topic that lost a copy.
    pub(crate) fn drops(&self, path: &str) -> Vec<TopicDrops> {
        self.paths
            .read()
            .expect("subscription lock poisoned")
            .get(path)
            .map_or_else(Vec::new, |s| s.drops.by_topic())
    }
}

/// One subscriber a streamed publish is writing to.
struct Target {
    tx: mpsc::Sender<PubItem>,
    budget: Arc<Semaphore>,
    /// The reserved slot for this copy's `Finish` or `Abort`. Always `Some`
    /// until the transfer ends.
    ending: Option<OwnedPermit<PubItem>>,
}

/// A publish in progress: a payload written once and fanned out to one stream
/// per subscriber, without ever being held whole.
///
/// This is what [`Publisher::publish`](crate::Publisher::publish) cannot do:
/// that call takes a `Bytes` and refuses anything above
/// `Limits::subscriber_buffer_bytes`, because a message that large could not
/// be enqueued for anybody. Here the **chunk** is what the budget bounds, so
/// the payload is unbounded and a 33 MB frame is an ordinary publish
/// (B-064, `docs/requirements/zeughaus-video.md` request 1).
///
/// **The drop behaviour of [GUARANTEES](../../../docs/GUARANTEES.md) §6 is
/// per subscriber, not per publish.** A subscriber whose budget or queue
/// cannot take a chunk loses *this* transfer — its stream is reset, so it
/// never mistakes a partial payload for a whole one — and it is counted in
/// [`Publisher::drops`](crate::Publisher::drops) like any other fan-out drop.
/// Every other subscriber keeps receiving, and the publisher never waits for
/// the slowest one.
///
/// Dropping this handle without [`FanOut::finish`] aborts every copy, for the
/// same reason [`OutgoingTransfer`](crate::OutgoingTransfer) resets on drop.
pub struct FanOut {
    id: u64,
    topic: Arc<str>,
    targets: Vec<Target>,
    /// `None` only for a fan-out with no subscribers at all, which has
    /// nothing to count against.
    drops: Option<Arc<DropTable>>,
    finished: bool,
}

impl FanOut {
    fn empty(id: u64, topic: Arc<str>) -> FanOut {
        FanOut {
            id,
            topic,
            targets: Vec::new(),
            drops: None,
            finished: false,
        }
    }

    /// The topic this transfer is published on.
    pub fn topic(&self) -> &str {
        &self.topic
    }

    /// Subscribers still receiving this transfer.
    ///
    /// It only falls: a subscriber that loses a chunk is gone from this
    /// transfer, and one that subscribes while it is in flight receives the
    /// next message rather than half of this one.
    pub fn subscribers(&self) -> usize {
        self.targets.len()
    }

    /// Writes the next chunk to every subscriber still receiving, waiting up
    /// to `limit` for one that has no room, and returns how many subscribers
    /// are left.
    ///
    /// **The bound is mandatory and finite, and it is the caller's**, for the
    /// reason [`Runtime::drain`](crate::Runtime::drain) takes one
    /// (`docs/decisions/0009-drain.md` §4.4): waiting on a subscriber with no
    /// deadline is how a publisher hangs on a peer's behaviour, and refusing
    /// to wait at all would make a payload larger than
    /// `Limits::subscriber_buffer_bytes` impossible to send to anyone —
    /// the publisher would outrun its own budget and abort every copy. A
    /// subscriber that frees room inside `limit` keeps the transfer; one that
    /// does not loses it, and only it.
    ///
    /// One `Bytes` allocation is shared by every copy, so a chunk costs one
    /// buffer regardless of subscriber count, and the payload is never held
    /// whole anywhere.
    ///
    /// Fails with [`Error::LimitExceeded`] for a chunk larger than
    /// `Limits::subscriber_buffer_bytes`: such a chunk could never be
    /// enqueued for anybody, and the point of this API is that the *payload*
    /// need not fit while a chunk does.
    pub async fn write_within(
        &mut self,
        chunk: impl Into<Bytes>,
        limit: std::time::Duration,
    ) -> Result<usize, Error> {
        let chunk = chunk.into();
        let want = u32::try_from(chunk.len()).map_err(|_| Error::LimitExceeded)?;
        let mut kept = Vec::with_capacity(self.targets.len());
        for mut target in std::mem::take(&mut self.targets) {
            let budget = Arc::clone(&target.budget);
            let acquired = tokio::select! {
                permit = tokio::time::timeout(limit, budget.acquire_many_owned(want)) => {
                    match permit {
                        Ok(Ok(permit)) => Some(permit),
                        // Timed out, or the semaphore is closed.
                        _ => None,
                    }
                }
                // The writer task is gone — the connection closed — so no
                // permit will ever come back. Without this arm the wait would
                // run to `limit` for a subscriber that cannot exist.
                () = target.tx.closed() => None,
            };
            let Some(permit) = acquired else {
                self.abort_one(&mut target, DropCause::SubscriberBudget);
                continue;
            };
            match self.enqueue(&target, &chunk) {
                Ok(()) => {
                    // The writer returns the permits once the bytes are gone.
                    permit.forget();
                    kept.push(target);
                }
                Err(()) => {
                    // Dropping the permit returns the budget immediately.
                    drop(permit);
                    self.abort_one(&mut target, DropCause::SubscriberQueue);
                }
            }
        }
        self.targets = kept;
        Ok(self.targets.len())
    }

    /// Writes the next chunk without ever waiting: a subscriber with no room
    /// right now loses the transfer.
    ///
    /// Fan-out's `Drop` from [GUARANTEES](../../../docs/GUARANTEES.md) §6 in
    /// its purest form, and the right call where a later chunk supersedes an
    /// earlier one — a video frame, a market snapshot — because a subscriber
    /// that cannot keep up should be waiting for the *next* transfer rather
    /// than holding this one up. A publisher streaming a payload that must
    /// arrive whole wants [`FanOut::write_within`].
    pub fn write_now(&mut self, chunk: impl Into<Bytes>) -> Result<usize, Error> {
        let chunk = chunk.into();
        let want = u32::try_from(chunk.len()).map_err(|_| Error::LimitExceeded)?;
        let mut kept = Vec::with_capacity(self.targets.len());
        for mut target in std::mem::take(&mut self.targets) {
            let Ok(permit) = target.budget.try_acquire_many(want) else {
                self.abort_one(&mut target, DropCause::SubscriberBudget);
                continue;
            };
            match self.enqueue(&target, &chunk) {
                Ok(()) => {
                    permit.forget();
                    kept.push(target);
                }
                Err(()) => {
                    drop(permit);
                    self.abort_one(&mut target, DropCause::SubscriberQueue);
                }
            }
        }
        self.targets = kept;
        Ok(self.targets.len())
    }

    /// Hands one chunk to one subscriber's writer. `Err` means the writer's
    /// queue is full or gone; the caller aborts that copy.
    fn enqueue(&self, target: &Target, chunk: &Bytes) -> Result<(), ()> {
        let item = PubItem::Chunk {
            id: self.id,
            payload: chunk.clone(),
        };
        target.tx.try_send(item).map_err(|_| ())
    }

    /// Ends the transfer, and returns how many subscribers received all of
    /// it as far as this side can tell.
    ///
    /// "As far as this side can tell" is the honest claim: the count is the
    /// subscribers whose every chunk was enqueued and whose FIN is queued
    /// behind them. A fan-out copy carries no receipt — nobody holds a
    /// `Delivery` for it — so the transport acknowledgement is awaited by the
    /// drain and by nothing else (`docs/decisions/0009-drain.md` §4.2).
    pub fn finish(mut self) -> usize {
        self.finished = true;
        let id = self.id;
        let delivered = self.targets.len();
        for target in &mut self.targets {
            if let Some(ending) = target.ending.take() {
                ending.send(PubItem::Finish { id });
            }
        }
        delivered
    }

    /// Aborts one subscriber's copy, counting the cause.
    fn abort_one(&self, target: &mut Target, cause: DropCause) {
        if let Some(drops) = &self.drops {
            drops.record(&self.topic, cause);
        }
        if let Some(ending) = target.ending.take() {
            ending.send(PubItem::Abort { id: self.id });
        }
        tracing::debug!(topic = %self.topic, ?cause, "streamed fan-out copy aborted");
    }
}

impl Drop for FanOut {
    fn drop(&mut self) {
        if self.finished {
            return;
        }
        // Abandoned mid-payload: reset every copy rather than leave a
        // subscriber waiting for a FIN.
        for target in &mut self.targets {
            if let Some(ending) = target.ending.take() {
                ending.send(PubItem::Abort { id: self.id });
            }
        }
    }
}

fn decrement(counts: &mut HashMap<usize, usize>, conn_id: usize, by: usize) {
    if let Some(count) = counts.get_mut(&conn_id) {
        *count = count.saturating_sub(by);
        if *count == 0 {
            counts.remove(&conn_id);
        }
    }
}

/// Writes one subscriber's messages, in order, each on its own uni stream.
///
/// Fan-out is one-way by construction: a published copy owes nothing back, so
/// it rides a unidirectional stream and the publisher never waits on it.
async fn writer(
    ctx: ConnHandle,
    path: Arc<str>,
    mut rx: mpsc::Receiver<PubItem>,
    budget: Arc<Semaphore>,
    drops: Arc<DropTable>,
) {
    // The streams of the streamed publishes currently in flight for this
    // subscriber. Bounded by what this process opens, never by the peer: a
    // remote party cannot make a publisher open a transfer.
    let mut streaming: HashMap<u64, SendHalf> = HashMap::new();
    loop {
        let item = tokio::select! {
            item = rx.recv() => match item {
                Some(item) => item,
                None => break,
            },
            // Without this a subscriber on a dead connection would sit here
            // forever holding an `Arc<ConnCtx>`.
            _ = ctx.conn.closed() => break,
        };
        match item {
            PubItem::Whole(msg) => {
                let len = msg.payload.len();
                let outcome = write_one(&ctx, &path, &msg).await;
                budget.add_permits(len);
                match outcome {
                    Ok(()) => {}
                    // The subscriber parked no connection for this copy. That
                    // is a drop of the copy, not a failure of the
                    // subscription: the pool refills and the next message may
                    // well go out
                    // ([decisions/0012](../../../docs/decisions/0012-local-connection-grouping.md)
                    // §4.4). Counted exactly like an exhausted byte budget.
                    Err(Error::NoParkedConnection) => {
                        drops.record(&msg.topic, DropCause::NoParkedConnection);
                        tracing::debug!(path = %path, "no parked connection; copy dropped");
                    }
                    Err(e) => {
                        tracing::debug!(path = %path, error = %e, "fan-out write failed; subscriber writer ending");
                        break;
                    }
                }
            }
            PubItem::Begin { id, head } => match begin_one(&ctx, &path, &head).await {
                Ok(stream) => {
                    streaming.insert(id, stream);
                }
                Err(Error::NoParkedConnection) => {
                    drops.record(&head.topic, DropCause::NoParkedConnection);
                    tracing::debug!(path = %path, "no parked connection; streamed copy dropped");
                }
                Err(e) => {
                    tracing::debug!(path = %path, error = %e, "fan-out open failed; subscriber writer ending");
                    break;
                }
            },
            PubItem::Chunk { id, payload } => {
                let len = payload.len();
                // A chunk for a transfer whose stream is gone — the open
                // failed, or a write did — is dropped here: the budget is
                // still returned, because the publisher charged it.
                if let Some(stream) = streaming.get_mut(&id)
                    && let Err(e) = stream.write_all(&payload).await
                {
                    tracing::debug!(path = %path, error = %e, "streamed fan-out write failed");
                    if let Some(mut stream) = streaming.remove(&id) {
                        stream.reset(codes::CANCELED);
                    }
                }
                budget.add_permits(len);
            }
            PubItem::Finish { id } => {
                if let Some(mut stream) = streaming.remove(&id)
                    && stream.finish().is_ok()
                    && ctx.parked.park(stream.stopped())
                {
                    ctx.shared.drain.evict();
                }
            }
            PubItem::Abort { id } => {
                if let Some(mut stream) = streaming.remove(&id) {
                    stream.reset(codes::CANCELED);
                }
            }
        }
    }
    // Whatever is still open was abandoned by the connection ending, not by
    // the publisher: reset it so the subscriber does not wait for a FIN.
    for (_, mut stream) in streaming {
        stream.reset(codes::CANCELED);
    }
    tracing::debug!(path = %path, "subscriber writer ended");
}

/// Opens one subscriber's stream for a streamed publish and writes its DATA
/// header.
///
/// No `content_len`: the key is optional at the decoder
/// (`docs/PROTOCOL.md` §6.2) and a streaming publisher does not know the
/// length. A subscriber therefore reads until FIN, which is what every
/// streamed transfer in weida does.
async fn begin_one(ctx: &ConnHandle, path: &str, head: &PubMsg) -> Result<SendHalf, Error> {
    let mut header = DataHeader::addressed(path);
    header.topic = Some(head.topic.to_string());
    header.traceparent = head.trace.map(|t| t.to_traceparent());
    header.sequence = head.sequence;
    let mut stream = ctx.open_uni().await?;
    write_data_preamble(&mut stream, &header).await?;
    Ok(stream)
}

async fn write_one(ctx: &ConnHandle, path: &str, msg: &PubMsg) -> Result<(), Error> {
    let mut header = DataHeader::addressed(path);
    header.topic = Some(msg.topic.to_string());
    header.content_len = Some(msg.payload.len() as u64);
    header.traceparent = msg.trace.map(|t| t.to_traceparent());
    // The number the publisher assigned to this *message* (`None` under
    // `core`): every subscriber's copy carries the same one, so a copy this
    // subscriber lost shows up as a hole in its own sequence.
    header.sequence = msg.sequence;

    let mut stream = ctx.open_uni().await?;
    write_data_preamble(&mut stream, &header).await?;
    stream.write_all(&msg.payload).await?;
    stream.finish()?;
    // A published copy is a finished transfer like any other, and nobody
    // holds a receipt for it: park it on this connection so a drain waits
    // for it (`docs/decisions/0009-drain.md` §4.2).
    if ctx.parked.park(stream.stopped()) {
        ctx.shared.drain.evict();
    }
    Ok(())
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn a_literal_filter_matches_whole_segments_only() {
        assert!(filter::matches("px.eur", "px.eur"));
        assert!(!filter::matches("px.eur", "px.eur.spot"));
        assert!(!filter::matches("fx.usd", "px.eur"));
        // The boundary a byte prefix could not see: `px.` used to select
        // everything under `px`, and `sensors.temp` used to select
        // `sensors.temperature`.
        assert!(!filter::matches("px.eur", "px."));
        assert!(!filter::matches("sensors.temperature", "sensors.temp"));
    }

    #[test]
    fn the_empty_filter_and_the_rest_wildcard_take_everything() {
        assert!(filter::matches("", ""));
        assert!(filter::matches("anything.at.all", ""));
        assert!(filter::matches("anything.at.all", "#"));
        assert!(filter::matches("", "#"));
    }

    #[test]
    fn one_segment_wildcard_matches_exactly_one() {
        assert!(filter::matches("px.eur", "px.*"));
        assert!(filter::matches("sensors.a.temp", "sensors.*.temp"));
        assert!(filter::matches("px.eur", "*.eur"));
        // Exactly one: neither none nor two.
        assert!(!filter::matches("px", "px.*"));
        assert!(!filter::matches("px.eur.spot", "px.*"));
    }

    #[test]
    fn the_rest_wildcard_matches_zero_or_more_trailing_segments() {
        assert!(filter::matches("px", "px.#"));
        assert!(filter::matches("px.eur", "px.#"));
        assert!(filter::matches("px.eur.spot", "px.#"));
        assert!(!filter::matches("fx", "px.#"));
        assert!(!filter::matches("pxx", "px.#"));
    }

    #[test]
    fn a_topic_is_never_a_pattern() {
        // `*` and `#` in a published topic are ordinary bytes.
        assert!(filter::matches("px.*", "px.*"));
        assert!(!filter::matches("px.*", "px.eur"));
        assert!(filter::matches("px.#", "px.#"));
        assert!(filter::matches("px.*", "*.*"));
    }

    #[test]
    fn empty_segments_match_only_empty_segments() {
        assert!(filter::matches("px.", "px."));
        assert!(!filter::matches("px.eur", "px."));
        assert!(filter::matches("px.", "px.*"));
    }

    #[test]
    fn matching_is_byte_exact_not_case_folded() {
        assert!(!filter::matches("PX.EUR", "px.*"));
        assert!(!filter::matches("px.eur", "PX.*"));
    }
}