durable-streams-server 0.2.0

Durable Streams protocol server in Rust, built with axum and tokio
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
mod common;

use common::{
    spawn_test_server, spawn_test_server_with_timeout, test_client, test_client_with_timeout,
    unique_stream_name,
};

/// Validates spec: 03-read-modes.md#long-poll-mode
///
/// When data exists at the requested offset, long-poll returns
/// 200 immediately with the data (same as catch-up).
#[tokio::test]
async fn test_long_poll_returns_immediately_when_data_exists() {
    let (base_url, _port) = spawn_test_server().await;
    let client = test_client();
    let stream_name = unique_stream_name();

    // Create stream and append data
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    client
        .post(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .body("hello")
        .send()
        .await
        .unwrap();

    // Long-poll with data available should return 200 immediately
    let response = client
        .get(format!(
            "{base_url}/v1/stream/{stream_name}?offset=-1&live=long-poll"
        ))
        .send()
        .await
        .unwrap();

    assert_eq!(response.status(), 200);

    let body = response.text().await.unwrap();
    assert_eq!(body, "hello");
}

/// Validates spec: 03-read-modes.md#long-poll-mode
///
/// Long-poll catches up first: if client is behind tail, returns
/// historical data immediately instead of waiting.
#[tokio::test]
async fn test_long_poll_catches_up_first() {
    let (base_url, _port) = spawn_test_server().await;
    let client = test_client();
    let stream_name = unique_stream_name();

    // Create stream and append multiple messages
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    client
        .post(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .body("first")
        .send()
        .await
        .unwrap();

    client
        .post(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .body("second")
        .send()
        .await
        .unwrap();

    // Long-poll from start should return all data immediately
    let response = client
        .get(format!(
            "{base_url}/v1/stream/{stream_name}?offset=-1&live=long-poll"
        ))
        .send()
        .await
        .unwrap();

    assert_eq!(response.status(), 200);

    let body = response.text().await.unwrap();
    assert_eq!(body, "firstsecond");
}

/// Validates spec: 03-read-modes.md#long-poll-mode
///
/// Long-poll waits for new data and returns 200 when it arrives.
#[tokio::test]
async fn test_long_poll_waits_and_returns_on_new_data() {
    let (base_url, _port) = spawn_test_server().await;
    let client = test_client();
    let stream_name = unique_stream_name();

    // Create stream
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    let url = format!("{base_url}/v1/stream/{stream_name}");
    let poll_url = format!("{url}?offset=-1&live=long-poll");
    let append_url = url.clone();

    // Start long-poll in background
    let poll_handle = tokio::spawn({
        let client = test_client();
        async move { client.get(&poll_url).send().await.unwrap() }
    });

    // Wait a moment then append data
    tokio::time::sleep(tokio::time::Duration::from_millis(200)).await;

    client
        .post(&append_url)
        .header("Content-Type", "text/plain")
        .body("arrived")
        .send()
        .await
        .unwrap();

    // Long-poll should return with the new data
    let response = poll_handle.await.unwrap();
    assert_eq!(response.status(), 200);

    let body = response.text().await.unwrap();
    assert_eq!(body, "arrived");
}

/// Validates spec: 03-read-modes.md#timeout
///
/// Long-poll returns 204 when timeout expires with no new data.
#[tokio::test]
async fn test_long_poll_timeout_returns_204() {
    let (base_url, _port) = spawn_test_server_with_timeout(std::time::Duration::from_secs(1)).await;
    let client = test_client();
    let stream_name = unique_stream_name();

    // Create empty stream
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    // Long-poll should timeout and return 204
    let start = std::time::Instant::now();
    let response = client
        .get(format!(
            "{base_url}/v1/stream/{stream_name}?offset=-1&live=long-poll"
        ))
        .send()
        .await
        .unwrap();
    let elapsed = start.elapsed();

    assert_eq!(response.status(), 204, "Expected 204 No Content on timeout");
    assert!(
        elapsed >= std::time::Duration::from_millis(900),
        "Should wait at least ~1s, waited {elapsed:?}"
    );
}

/// Validates spec: 03-read-modes.md#closed-stream-at-tail
///
/// When stream is closed and client is at tail, long-poll MUST
/// immediately return 204 with Stream-Closed: true (no waiting).
#[tokio::test]
async fn test_long_poll_closed_stream_at_tail_returns_204_immediately() {
    let (base_url, _port) = spawn_test_server().await;
    let client = test_client();
    let stream_name = unique_stream_name();

    // Create stream, append, and close
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    client
        .post(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .header("Stream-Closed", "true")
        .body("final")
        .send()
        .await
        .unwrap();

    // Get the next offset (tail position)
    let read_response = client
        .get(format!("{base_url}/v1/stream/{stream_name}?offset=-1"))
        .send()
        .await
        .unwrap();

    let next_offset = read_response
        .headers()
        .get("Stream-Next-Offset")
        .unwrap()
        .to_str()
        .unwrap()
        .to_string();

    // Consume the body to avoid connection issues
    let _ = read_response.text().await.unwrap();

    // Long-poll from tail of closed stream should return immediately
    let start = std::time::Instant::now();
    let response = client
        .get(format!(
            "{base_url}/v1/stream/{stream_name}?offset={next_offset}&live=long-poll"
        ))
        .send()
        .await
        .unwrap();
    let elapsed = start.elapsed();

    assert_eq!(response.status(), 204, "Expected 204 No Content");
    assert!(
        elapsed < std::time::Duration::from_millis(500),
        "Should return immediately for closed stream, took {elapsed:?}"
    );

    let closed = response
        .headers()
        .get("Stream-Closed")
        .expect("Missing Stream-Closed header")
        .to_str()
        .unwrap();
    assert_eq!(closed, "true");
}

/// Validates spec: 03-read-modes.md#long-poll-mode
///
/// Long-poll on non-existent stream returns 404.
#[tokio::test]
async fn test_long_poll_nonexistent_stream_returns_404() {
    let (base_url, _port) = spawn_test_server().await;
    let client = test_client();
    let stream_name = unique_stream_name();

    let response = client
        .get(format!(
            "{base_url}/v1/stream/{stream_name}?live=long-poll&offset=-1"
        ))
        .send()
        .await
        .unwrap();

    assert_eq!(response.status(), 404);
}

/// Validates spec: 03-read-modes.md#long-poll-mode
///
/// Invalid `live` parameter value returns 400.
#[tokio::test]
async fn test_long_poll_invalid_live_param_returns_400() {
    let (base_url, _port) = spawn_test_server().await;
    let client = test_client();
    let stream_name = unique_stream_name();

    // Create stream
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    let response = client
        .get(format!("{base_url}/v1/stream/{stream_name}?live=invalid"))
        .send()
        .await
        .unwrap();

    assert_eq!(response.status(), 400, "Expected 400 for invalid live mode");
}

/// Validates spec: 03-read-modes.md#stream-cursor
///
/// Long-poll responses include Stream-Cursor header.
#[tokio::test]
async fn test_long_poll_includes_stream_cursor() {
    let (base_url, _port) = spawn_test_server().await;
    let client = test_client();
    let stream_name = unique_stream_name();

    // Create stream and append data
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    client
        .post(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .body("data")
        .send()
        .await
        .unwrap();

    // Long-poll should include Stream-Cursor
    let response = client
        .get(format!(
            "{base_url}/v1/stream/{stream_name}?offset=-1&live=long-poll"
        ))
        .send()
        .await
        .unwrap();

    assert_eq!(response.status(), 200);

    let cursor = response
        .headers()
        .get("Stream-Cursor")
        .expect("Missing Stream-Cursor header on long-poll response")
        .to_str()
        .unwrap();

    // Cursor should be non-empty and digits-only (opaque monotonic counter)
    assert!(!cursor.is_empty(), "Stream-Cursor should not be empty");
    assert!(
        cursor.chars().all(|c| c.is_ascii_digit()),
        "Stream-Cursor should be digits only, got: {cursor}"
    );
}

/// Validates spec: 03-read-modes.md#long-poll-mode
///
/// 204 response includes correct headers: Stream-Next-Offset,
/// Stream-Up-To-Date, and Stream-Cursor.
#[tokio::test]
async fn test_long_poll_204_includes_correct_headers() {
    let (base_url, _port) = spawn_test_server_with_timeout(std::time::Duration::from_secs(1)).await;
    let client = test_client();
    let stream_name = unique_stream_name();

    // Create stream with some data
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    client
        .post(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .body("existing")
        .send()
        .await
        .unwrap();

    // Get next offset (tail position)
    let read_response = client
        .get(format!("{base_url}/v1/stream/{stream_name}?offset=-1"))
        .send()
        .await
        .unwrap();

    let next_offset = read_response
        .headers()
        .get("Stream-Next-Offset")
        .unwrap()
        .to_str()
        .unwrap()
        .to_string();
    let _ = read_response.text().await.unwrap();

    // Long-poll from tail, wait for timeout
    let response = client
        .get(format!(
            "{base_url}/v1/stream/{stream_name}?offset={next_offset}&live=long-poll"
        ))
        .send()
        .await
        .unwrap();

    assert_eq!(response.status(), 204);

    // Verify required headers
    let resp_next_offset = response
        .headers()
        .get("Stream-Next-Offset")
        .expect("204 must include Stream-Next-Offset")
        .to_str()
        .unwrap();
    assert_eq!(resp_next_offset, next_offset);

    let up_to_date = response
        .headers()
        .get("Stream-Up-To-Date")
        .expect("204 must include Stream-Up-To-Date")
        .to_str()
        .unwrap();
    assert_eq!(up_to_date, "true");

    let cursor = response
        .headers()
        .get("Stream-Cursor")
        .expect("204 must include Stream-Cursor")
        .to_str()
        .unwrap();
    assert!(!cursor.is_empty());

    // Stream-Closed should NOT be present (stream is open)
    assert!(
        response.headers().get("Stream-Closed").is_none(),
        "Stream-Closed should not be present on open stream"
    );
}

/// Validates spec: 03-read-modes.md#long-poll-mode
///
/// Long-poll response includes all standard read headers (`ETag`, etc.)
/// when returning data.
#[tokio::test]
async fn test_long_poll_200_includes_all_read_headers() {
    let (base_url, _port) = spawn_test_server().await;
    let client = test_client();
    let stream_name = unique_stream_name();

    // Create stream and append data
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    client
        .post(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .body("data")
        .send()
        .await
        .unwrap();

    let response = client
        .get(format!(
            "{base_url}/v1/stream/{stream_name}?offset=-1&live=long-poll"
        ))
        .send()
        .await
        .unwrap();

    assert_eq!(response.status(), 200);

    // All standard read headers must be present
    assert!(
        response.headers().get("Content-Type").is_some(),
        "Missing Content-Type"
    );
    assert!(
        response.headers().get("Stream-Next-Offset").is_some(),
        "Missing Stream-Next-Offset"
    );
    assert!(
        response.headers().get("Stream-Up-To-Date").is_some(),
        "Missing Stream-Up-To-Date"
    );
    assert!(response.headers().get("ETag").is_some(), "Missing ETag");
    assert!(
        response.headers().get("Stream-Cursor").is_some(),
        "Missing Stream-Cursor"
    );
    assert!(
        response.headers().get("Cache-Control").is_some(),
        "Missing Cache-Control"
    );
}

/// Validates spec: 03-read-modes.md#long-poll-mode
///
/// Long-poll wakes up and returns 204 with Stream-Closed when
/// stream is closed mid-wait.
#[tokio::test]
async fn test_long_poll_wakes_on_stream_close() {
    let (base_url, _port) = spawn_test_server().await;
    let client = test_client_with_timeout(10);
    let stream_name = unique_stream_name();

    // Create stream
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    let url = format!("{base_url}/v1/stream/{stream_name}");

    // Start long-poll in background
    let poll_handle = tokio::spawn({
        let poll_url = format!("{url}?offset=-1&live=long-poll");
        let client = test_client_with_timeout(10);
        async move { client.get(&poll_url).send().await.unwrap() }
    });

    // Wait a moment then close the stream
    tokio::time::sleep(tokio::time::Duration::from_millis(200)).await;

    client
        .post(&url)
        .header("Content-Type", "text/plain")
        .header("Stream-Closed", "true")
        .body("")
        .send()
        .await
        .unwrap();

    // Long-poll should wake up and return 204 with Stream-Closed
    let response = poll_handle.await.unwrap();
    assert_eq!(response.status(), 204);

    let closed = response
        .headers()
        .get("Stream-Closed")
        .expect("Should include Stream-Closed after close event")
        .to_str()
        .unwrap();
    assert_eq!(closed, "true");
}

/// Validates spec: 03-read-modes.md#long-poll-mode
///
/// Regression: long-poll from offset=now must deliver data that arrives
/// during the wait. The sentinel `now` resolves to the current tail on
/// the initial read; re-reads after wake must use the resolved concrete
/// offset, not the raw sentinel (which would always return empty).
#[tokio::test]
async fn test_long_poll_from_now_sentinel_delivers_data() {
    let (base_url, _port) = spawn_test_server().await;
    let client = test_client();
    let stream_name = unique_stream_name();

    // Create stream
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    let url = format!("{base_url}/v1/stream/{stream_name}");

    // Start long-poll from `now` in background
    let poll_handle = tokio::spawn({
        let poll_url = format!("{url}?offset=now&live=long-poll");
        let client = test_client();
        async move { client.get(&poll_url).send().await.unwrap() }
    });

    // Wait a moment then append data
    tokio::time::sleep(tokio::time::Duration::from_millis(200)).await;

    client
        .post(&url)
        .header("Content-Type", "text/plain")
        .body("after-now")
        .send()
        .await
        .unwrap();

    // Long-poll should wake and return 200 with the new data
    let response = poll_handle.await.unwrap();
    assert_eq!(
        response.status(),
        200,
        "Long-poll from now should return 200 when data arrives"
    );

    let body = response.text().await.unwrap();
    assert_eq!(body, "after-now");
}

/// Validates spec: 03-read-modes.md#long-poll-mode
///
/// Catch-up mode (no live param) is unaffected by long-poll changes;
/// it does NOT include Stream-Cursor header.
#[tokio::test]
async fn test_catch_up_mode_unchanged() {
    let (base_url, _port) = spawn_test_server().await;
    let client = test_client();
    let stream_name = unique_stream_name();

    // Create stream and append data
    client
        .put(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .send()
        .await
        .unwrap();

    client
        .post(format!("{base_url}/v1/stream/{stream_name}"))
        .header("Content-Type", "text/plain")
        .body("data")
        .send()
        .await
        .unwrap();

    // Regular catch-up GET (no live param)
    let response = client
        .get(format!("{base_url}/v1/stream/{stream_name}?offset=-1"))
        .send()
        .await
        .unwrap();

    assert_eq!(response.status(), 200);

    // Catch-up mode should NOT have Stream-Cursor
    assert!(
        response.headers().get("Stream-Cursor").is_none(),
        "Catch-up mode should not include Stream-Cursor"
    );

    let body = response.text().await.unwrap();
    assert_eq!(body, "data");
}