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
|
use std::future;
use axum::extract::State;
use axum_extra::extract::Query;
use futures::stream::{self, StreamExt as _};
use crate::{
event::{routes::get, Sequenced as _},
test::fixtures::{self, future::Expect as _},
};
#[tokio::test]
async fn resume() {
// Set up the environment
let app = fixtures::scratch_app().await;
let channel = fixtures::channel::create(&app, &fixtures::now()).await;
let sender = fixtures::login::create(&app, &fixtures::now()).await;
let initial_message = fixtures::message::send(&app, &channel, &sender, &fixtures::now()).await;
let later_messages = vec![
fixtures::message::send(&app, &channel, &sender, &fixtures::now()).await,
fixtures::message::send(&app, &channel, &sender, &fixtures::now()).await,
];
// Call the endpoint
let subscriber = fixtures::identity::create(&app, &fixtures::now()).await;
let resume_at = {
// First subscription
let get::Response(events) = get::handler(
State(app.clone()),
subscriber.clone(),
None,
Query::default(),
)
.await
.expect("subscribe never fails");
let event = events
.filter_map(fixtures::event::message)
.filter_map(fixtures::event::message::sent)
.filter(|event| future::ready(event.message == initial_message))
.next()
.expect_some("delivered event for initial message")
.await;
event.sequence()
};
// Resume after disconnect
let get::Response(resumed) = get::handler(
State(app),
subscriber,
Some(resume_at.into()),
Query::default(),
)
.await
.expect("subscribe never fails");
// Verify final events
let mut events = resumed
.filter_map(fixtures::event::message)
.filter_map(fixtures::event::message::sent)
.zip(stream::iter(later_messages));
while let Some((event, message)) = events.next().expect_ready("event ready").await {
assert_eq!(message, event.message);
}
}
// This test verifies a real bug I hit developing the vector-of-sequences
// approach to resuming events. A small omission caused the event IDs in a
// resumed stream to _omit_ channels that were in the original stream until
// those channels also appeared in the resumed stream.
//
// Clients would see something like
// * In the original stream, Cfoo=5,Cbar=8
// * In the resumed stream, Cfoo=6 (no Cbar sequence number)
//
// Disconnecting and reconnecting a second time, using event IDs from that
// initial period of the first resume attempt, would then cause the second
// resume attempt to restart all other channels from the beginning, and not
// from where the first disconnection happened.
//
// As we have switched to a single global event sequence number, this scenario
// can no longer arise, but this test is preserved because the actual behaviour
// _is_ a valid way for clients to behave, and should work. We might as well
// keep testing it.
#[tokio::test]
async fn serial_resume() {
// Set up the environment
let app = fixtures::scratch_app().await;
let sender = fixtures::login::create(&app, &fixtures::now()).await;
let channel_a = fixtures::channel::create(&app, &fixtures::now()).await;
let channel_b = fixtures::channel::create(&app, &fixtures::now()).await;
// Call the endpoint
let subscriber = fixtures::identity::create(&app, &fixtures::now()).await;
let resume_at = {
let initial_messages = [
fixtures::message::send(&app, &channel_a, &sender, &fixtures::now()).await,
fixtures::message::send(&app, &channel_b, &sender, &fixtures::now()).await,
];
// First subscription
let get::Response(events) = get::handler(
State(app.clone()),
subscriber.clone(),
None,
Query::default(),
)
.await
.expect("subscribe never fails");
// Check for expected events
let events = events
.filter_map(fixtures::event::message)
.filter_map(fixtures::event::message::sent)
.zip(stream::iter(initial_messages))
.collect::<Vec<_>>()
.expect_ready("zipping a finite list of events is ready immediately")
.await;
assert!(events
.iter()
.all(|(event, message)| message == &event.message));
let (event, _) = events.last().expect("this vec is non-empty");
// Take the last one's resume point
event.sequence()
};
// Resume after disconnect
let resume_at = {
let resume_messages = [
// Note that channel_b does not appear here. The buggy behaviour
// would be masked if channel_b happened to send a new message
// into the resumed event stream.
fixtures::message::send(&app, &channel_a, &sender, &fixtures::now()).await,
fixtures::message::send(&app, &channel_a, &sender, &fixtures::now()).await,
];
// Second subscription
let get::Response(events) = get::handler(
State(app.clone()),
subscriber.clone(),
Some(resume_at.into()),
Query::default(),
)
.await
.expect("subscribe never fails");
// Check for expected events
let events = events
.filter_map(fixtures::event::message)
.filter_map(fixtures::event::message::sent)
.zip(stream::iter(resume_messages))
.collect::<Vec<_>>()
.expect_ready("zipping a finite list of events is ready immediately")
.await;
assert!(events
.iter()
.all(|(event, message)| message == &event.message));
let (event, _) = events.last().expect("this vec is non-empty");
// Take the last one's resume point
event.sequence()
};
// Resume after disconnect a second time
{
// At this point, we can send on either channel and demonstrate the
// problem. The resume point should before both of these messages, but
// after _all_ prior messages.
let final_messages = [
fixtures::message::send(&app, &channel_a, &sender, &fixtures::now()).await,
fixtures::message::send(&app, &channel_b, &sender, &fixtures::now()).await,
];
// Third subscription
let get::Response(events) = get::handler(
State(app.clone()),
subscriber.clone(),
Some(resume_at.into()),
Query::default(),
)
.await
.expect("subscribe never fails");
// Check for expected events
let events = events
.filter_map(fixtures::event::message)
.filter_map(fixtures::event::message::sent)
.zip(stream::iter(final_messages))
.collect::<Vec<_>>()
.expect_ready("zipping a finite list of events is ready immediately")
.await;
assert!(events
.iter()
.all(|(event, message)| message == &event.message));
};
}
|