weida 0.1.0-alpha.2

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
//! The blocking facade, from threads with no executor (B-194).
//!
//! These are plain `#[test]` functions: no `#[tokio::test]` anywhere, which
//! is the point — a synchronous caller has no reactor, and everything below
//! runs on the reactor `weida::blocking::Runtime` owns and never shows.
//!
//! One round trip per pattern, because what needs proving is that the facade
//! drives the asynchronous surface correctly, not that the patterns work: they
//! have their own suites. The fourth test is the one that earns its place
//! twice over — calling a blocking method from inside a Tokio runtime is the
//! one mistake that deadlocks instead of failing, and the facade refuses it.
#![cfg(all(feature = "blocking", feature = "generate"))]

use std::time::Duration;

use weida::blocking::Runtime;
use weida::{
    Acknowledgement, CursorLevel, Error, Identity, ReportMode, Reported, RuntimeConfig,
    TransferMeta, Trust,
};

/// A payload ceiling for every receive: the facade has no default, on purpose.
const CAP: usize = 1024 * 1024;

/// Binds a server on loopback and returns the runtime, the binding and the
/// address a client dials.
fn served(path: &str) -> (Runtime, weida::blocking::Binding, String) {
    let runtime = Runtime::new(RuntimeConfig::default()).expect("an owned reactor");
    let identity = Identity::generate().expect("identity");
    let fingerprint = identity.fingerprint().expect("fingerprint");
    let binding = runtime
        .bind_quic("127.0.0.1:0".parse().expect("addr"), identity)
        .expect("bind");
    let url = format!("weida://{fingerprint}@{}{path}", binding.local_addr());
    (runtime, binding, url)
}

#[test]
fn a_request_and_its_reply_cross_between_two_threads() {
    let (server, binding, url) = served("/echo");
    let replier = binding.replier("/echo").expect("replier");

    // The replier runs on its own thread, blocking on `accept`, which is how
    // a synchronous server is written.
    let answering = std::thread::spawn(move || {
        let request = replier.accept(CAP).expect("accept");
        assert_eq!(request.message().payload, b"ping");
        // The metadata is the asynchronous surface's, passed through: this
        // client presented no identity, so the peer is `None` — which is the
        // documented answer and not a facade artefact.
        assert!(request.message().meta.peer.is_none());
        assert_eq!(request.message().meta.endpoint.as_deref(), Some("/echo"));
        request.reply(b"pong").expect("reply");
    });

    let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
    let requester = client.requester(Trust::by_address());
    requester.connect(&url).expect("connect");
    assert_eq!(requester.request(b"ping", CAP).expect("request"), b"pong");

    answering.join().expect("the replier thread");
    client.shutdown().expect("client shutdown");
    server.shutdown().expect("server shutdown");
}

#[test]
fn a_push_reaches_a_puller_and_the_receipt_comes_back() {
    let (server, binding, url) = served("/ingest");
    let puller = binding.puller("/ingest").expect("puller");

    let draining = std::thread::spawn(move || {
        let first = puller.recv(CAP).expect("recv");
        assert_eq!(first.payload, b"one");
        let second = puller.recv(CAP).expect("recv");
        assert_eq!(second.payload, b"two");
    });

    let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
    let pusher = client.pusher(Trust::by_address());
    pusher.connect(&url).expect("connect");
    // `send` returns when the peer's transport holds the bytes, which is what
    // a synchronous caller has to be able to check.
    pusher.send(b"one").expect("send");
    pusher.send(b"two").expect("send");

    draining.join().expect("the puller thread");
    client.shutdown().expect("client shutdown");
    server.shutdown().expect("server shutdown");
}

#[test]
fn a_published_message_reaches_a_blocking_subscriber() {
    let (server, binding, url) = served("/md");
    let publisher = binding.publisher("/md").expect("publisher");

    let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
    let subscriber = client.subscriber(Trust::by_address());
    subscriber.connect(&url).expect("connect");
    subscriber.subscribe("px.#").expect("subscribe");

    // Publishing never waits for a subscriber, so the publisher may run
    // before the subscription has arrived: retry until it has, exactly as the
    // asynchronous Pub/Sub tests do, because a sleep would only be a guess.
    let deadline = std::time::Instant::now() + Duration::from_secs(15);
    loop {
        assert!(
            std::time::Instant::now() < deadline,
            "the subscription never reached the publisher"
        );
        if publisher
            .publish("px.eur", &b"1.0812"[..])
            .expect("publish")
            > 0
        {
            break;
        }
        std::thread::sleep(Duration::from_millis(5));
    }

    let message = subscriber.recv(CAP).expect("recv");
    assert_eq!(message.payload, b"1.0812");
    assert_eq!(message.meta.topic.as_deref(), Some("px.eur"));
    // A filter that does not match gets nothing, which is the one claim a
    // facade could break by dropping the filter on the way down.
    assert_eq!(publisher.publish("fx.chf", &b"x"[..]).expect("publish"), 0);

    client.shutdown().expect("client shutdown");
    server.shutdown().expect("server shutdown");
}

/// PAIR over the facade, and the one rule a caller gets wrong: the **first**
/// peer is kept and a second is refused (B-244).
#[test]
fn a_pair_carries_both_directions_and_refuses_a_second_peer() {
    let (server, binding, url) = served("/link");
    let bound = binding.pair("/link").expect("bound pair");

    // The bound half runs on its own thread: a pair is symmetric, so both
    // ends block on `recv` and answer. It is **handed back** rather than
    // dropped: dropping a bound pair unregisters its route, so a first peer
    // that "still works" has to have something to talk to.
    let answering = std::thread::spawn(move || {
        let first = bound.recv(CAP).expect("recv");
        assert_eq!(first.payload, b"ping");
        bound.send(b"pong").expect("send back");
        bound
    });

    let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
    let paired = client.pair(Trust::by_address());
    paired.connect(&url).expect("connect");
    paired.send(b"ping").expect("send");
    assert_eq!(paired.recv(CAP).expect("recv").payload, b"pong");
    let bound = answering.join().expect("the bound thread");

    // A second peer: the connection is accepted — the refusal is per stream —
    // and the transfer is refused. Past the peer's 1 MiB stream window, so the
    // refusal cannot lose the race to the receipt
    // (`docs/decisions/0005-refusal-race.md`).
    let newcomer = Runtime::new(RuntimeConfig::default()).expect("second client runtime");
    let second = newcomer.pair(Trust::by_address());
    second
        .connect(&url)
        .expect("the refusal is per stream, not per connection");
    let refused = second
        .send(&vec![0x7au8; 2 * 1024 * 1024])
        .expect_err("a second peer is refused");
    assert!(
        matches!(refused, Error::LimitExceeded),
        "a capacity decision said out loud, got {refused:?}"
    );

    // And the first peer still **delivers**, which is the whole rule: ZeroMQ's
    // PAIR would have dropped it for the newcomer.
    paired.send(b"still mine").expect("the first peer is kept");
    assert_eq!(bound.recv(CAP).expect("recv").payload, b"still mine");

    newcomer.shutdown().expect("second client shutdown");
    client.shutdown().expect("client shutdown");
    server.shutdown().expect("server shutdown");
}

/// SURVEY over the facade, and the one rule a caller gets wrong: the
/// **deadline** is the caller's and silence is a number (B-244).
#[test]
fn a_survey_collects_what_answers_and_counts_what_does_not() {
    let (server, binding, url) = served("/poll");
    let quiet_url = url.replace("/poll", "/quiet");
    let answering = binding.respondent("/poll").expect("respondent");
    let silent = binding.respondent("/quiet").expect("a second respondent");

    let answers = std::thread::spawn(move || {
        let question = answering.accept(CAP).expect("accept");
        assert_eq!(question.message().payload, b"who is there");
        question.reply(b"me").expect("reply");
    });
    // Accepts the question and never answers it: not unreachable, not
    // refusing, just silent — which is what a deadline exists for.
    let holding = std::thread::spawn(move || {
        let question = silent.accept(CAP).expect("accept");
        std::thread::sleep(Duration::from_secs(3));
        drop(question);
    });

    let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
    let surveyor = client.surveyor(Trust::by_address());
    surveyor.connect(&url).expect("connect the answering one");
    surveyor
        .connect(&quiet_url)
        .expect("connect the silent one");

    let survey = surveyor
        .survey(b"who is there", Duration::from_millis(500), CAP)
        .expect("the survey ran");
    assert_eq!(survey.asked, 2, "both respondents were asked: {survey:?}");
    assert_eq!(
        survey.replies,
        vec![b"me".to_vec()],
        "one answered: {survey:?}"
    );
    assert_eq!(
        survey.silent(),
        1,
        "and the other's silence is a number, not an error: {survey:?}"
    );

    answers.join().expect("the answering thread");
    client.shutdown().expect("client shutdown");
    server.shutdown().expect("server shutdown");
    holding.join().expect("the silent thread");
}

/// BUS over the facade, and the one rule a caller gets wrong: **never your
/// own message** (B-244).
#[test]
fn a_bus_message_reaches_every_other_member_and_never_the_sender() {
    let (first_rt, first_binding, first_url) = served("/bus1");
    let (second_rt, second_binding, second_url) = served("/bus2");
    let one = first_binding
        .bus("/bus1", Trust::by_address())
        .expect("first member");
    let two = second_binding
        .bus("/bus2", Trust::by_address())
        .expect("second member");
    one.connect(&second_url).expect("one dials two");
    two.connect(&first_url).expect("two dials one");

    let listening = std::thread::spawn(move || {
        let heard = two.recv(CAP).expect("recv");
        assert_eq!(heard.payload, b"hello all");
        two.send(b"and back").expect("send back");
        two
    });

    assert_eq!(one.send(b"hello all").expect("send"), 1, "one other member");
    assert_eq!(one.recv(CAP).expect("recv").payload, b"and back");

    // The sender never hears itself: with the exchange above complete, a
    // further send leaves nothing for this member to receive.
    let two = listening.join().expect("the second member's thread");
    one.send(b"mine alone").expect("send");
    assert_eq!(two.recv(CAP).expect("recv").payload, b"mine alone");
    drop(two);

    first_rt.shutdown().expect("first shutdown");
    second_rt.shutdown().expect("second shutdown");
}

/// The mistake that would otherwise be a hang: blocking a reactor worker.
#[test]
fn the_facade_refuses_to_block_a_reactor_thread() {
    let reactor = tokio::runtime::Builder::new_current_thread()
        .build()
        .expect("a runtime to be inside of");
    let refused = reactor.block_on(async { Runtime::new(RuntimeConfig::default()) });
    match refused {
        Err(Error::Runtime(message)) => {
            assert!(
                message.contains("deadlock"),
                "the refusal must say why: {message}"
            );
        }
        Ok(_) => panic!("a blocking runtime built inside a reactor would deadlock on first use"),
        Err(other) => panic!("{other:?}"),
    }

    // And the same for a call on an endpoint built outside: it is the *call*
    // that blocks, so the check belongs on every entry point rather than only
    // on the constructor.
    let outside = Runtime::new(RuntimeConfig::default()).expect("outside a reactor");
    let requester = outside.requester(Trust::by_address());
    let refused = reactor.block_on(async { requester.connect("weida://127.0.0.1:1/x") });
    assert!(
        matches!(refused, Err(Error::Runtime(_))),
        "{refused:?} must be refused rather than deadlock"
    );
    outside.shutdown().expect("shutdown");
}

/// The cursor surface, synchronously: a fire-and-forget producer that gets a
/// **verdict** without an exchange, and a staged receiver that gives it
/// (B-243).
#[test]
fn a_push_gets_a_verdict_from_a_cursor_with_no_exchange() {
    const ACCEPTED: CursorLevel = CursorLevel::Known(Acknowledgement::Accepted);
    const PROCESSED: CursorLevel = CursorLevel::Known(Acknowledgement::Processed);

    let (server, binding, url) = served("/work");
    let puller = binding.puller("/work").expect("puller");

    let receiving = std::thread::spawn(move || {
        let (message, reporter) = puller.recv_reporting(CAP).expect("recv");
        assert_eq!(message.payload, b"a unit of work");
        // The order is in the metadata the DATA header carried, which is what
        // makes a report the sender's request rather than this side's idea.
        assert_eq!(message.meta.report, vec![ACCEPTED, PROCESSED]);
        assert_eq!(message.meta.report_mode, ReportMode::Progress);
        assert!(message.meta.report_id.is_some());

        let mut reporter = reporter.expect("the sender ordered a report");
        assert_eq!(reporter.levels(), [ACCEPTED, PROCESSED]);
        reporter
            .report(ACCEPTED, message.payload.len() as u64)
            .expect("accepted");
        // A level nobody ordered is ignored rather than refused: the order
        // says what the sender wants to hear.
        reporter
            .report(CursorLevel::Known(Acknowledgement::Stored), 1)
            .expect("ignored");
        reporter
            .report(PROCESSED, message.payload.len() as u64)
            .expect("processed");
        reporter.finish().expect("finish the report");
    });

    let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
    let pusher = client.pusher(Trust::by_address());
    pusher.connect(&url).expect("connect");
    let body = b"a unit of work";
    let mut cursors = pusher
        .send_reporting(
            TransferMeta::default()
                .with_content_len(body.len() as u64)
                .with_report([PROCESSED, ACCEPTED]),
            body,
        )
        .expect("send")
        .expect("the metadata ordered a report");

    // The verdict arrives **after** the transport receipt `send_reporting`
    // already waited for, which is the whole point of a cursor: `Processed` is
    // not something a FIN can carry. The wait is bounded, because a
    // synchronous one has to be: a receiver that died would otherwise park
    // this thread on a connection nobody closes.
    let mut processed = None;
    while processed.is_none() {
        match cursors.changed(Duration::from_secs(15)).expect("a cursor") {
            Reported::Changed => {
                processed = cursors.snapshot().offset(PROCESSED);
            }
            Reported::Waiting => panic!("no cursor within the deadline"),
            Reported::Ended => panic!("the report ended without a verdict"),
        }
    }
    assert_eq!(processed, Some(body.len() as u64));
    assert_eq!(
        cursors.snapshot().offset(ACCEPTED),
        Some(body.len() as u64),
        "a cursor is absolute, so the earlier level is still readable"
    );
    // The level nobody ordered never arrived.
    assert_eq!(
        cursors
            .snapshot()
            .offset(CursorLevel::Known(Acknowledgement::Stored)),
        None
    );

    receiving.join().expect("the receiving thread");
    client.shutdown().expect("client shutdown");
    server.shutdown().expect("server shutdown");
}

/// A send that orders nothing has nothing to read, and a receiver of it has
/// nothing to report with (B-243).
#[test]
fn a_transfer_that_orders_no_report_hands_out_no_handles() {
    let (server, binding, url) = served("/plain");
    let puller = binding.puller("/plain").expect("puller");

    let receiving = std::thread::spawn(move || {
        let (message, reporter) = puller.recv_reporting(CAP).expect("recv");
        assert_eq!(message.payload, b"no report");
        assert!(message.meta.report.is_empty());
        assert!(message.meta.report_id.is_none());
        assert!(
            reporter.is_none(),
            "a reporter with nothing ordered would be a handle that writes to nobody"
        );
    });

    let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
    let pusher = client.pusher(Trust::by_address());
    pusher.connect(&url).expect("connect");
    assert!(
        pusher
            .send_reporting(TransferMeta::default(), b"no report")
            .expect("send")
            .is_none(),
        "nothing was ordered, so there is nothing to read"
    );

    receiving.join().expect("the receiving thread");
    client.shutdown().expect("client shutdown");
    server.shutdown().expect("server shutdown");
}