| 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 |
use base64::Engine; |
| 35 |
use chrono::{DateTime, Utc}; |
| 36 |
use serde::Serialize; |
| 37 |
use uuid::Uuid; |
| 38 |
|
| 39 |
|
| 40 |
pub const DEFAULT_PAGE_SIZE: u32 = 50; |
| 41 |
|
| 42 |
pub const MAX_PAGE_SIZE: u32 = 200; |
| 43 |
|
| 44 |
|
| 45 |
|
| 46 |
|
| 47 |
|
| 48 |
#[derive(Debug, Clone, Copy, PartialEq, Eq)] |
| 49 |
pub struct Cursor { |
| 50 |
pub ts: DateTime<Utc>, |
| 51 |
pub id: Uuid, |
| 52 |
} |
| 53 |
|
| 54 |
impl Cursor { |
| 55 |
pub fn new(ts: DateTime<Utc>, id: Uuid) -> Self { |
| 56 |
Self { ts, id } |
| 57 |
} |
| 58 |
|
| 59 |
|
| 60 |
|
| 61 |
|
| 62 |
pub fn encode(&self) -> String { |
| 63 |
let raw = format!("{}:{}", self.ts.timestamp_millis(), self.id); |
| 64 |
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(raw) |
| 65 |
} |
| 66 |
|
| 67 |
|
| 68 |
|
| 69 |
|
| 70 |
|
| 71 |
|
| 72 |
|
| 73 |
|
| 74 |
pub fn decode(token: &str) -> Option<Self> { |
| 75 |
let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD |
| 76 |
.decode(token) |
| 77 |
.ok()?; |
| 78 |
let raw = String::from_utf8(bytes).ok()?; |
| 79 |
let (millis, id) = raw.split_once(':')?; |
| 80 |
let ts = DateTime::<Utc>::from_timestamp_millis(millis.parse().ok()?)?; |
| 81 |
Some(Self::new(ts, id.parse().ok()?)) |
| 82 |
} |
| 83 |
} |
| 84 |
|
| 85 |
impl Serialize for Cursor { |
| 86 |
|
| 87 |
|
| 88 |
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> { |
| 89 |
serializer.serialize_str(&self.encode()) |
| 90 |
} |
| 91 |
} |
| 92 |
|
| 93 |
|
| 94 |
#[derive(Debug, Clone)] |
| 95 |
pub struct PageParams { |
| 96 |
pub after: Option<Cursor>, |
| 97 |
pub limit: u32, |
| 98 |
} |
| 99 |
|
| 100 |
impl PageParams { |
| 101 |
|
| 102 |
|
| 103 |
|
| 104 |
pub fn new(after_token: Option<&str>, limit: Option<u32>) -> Self { |
| 105 |
Self { |
| 106 |
after: after_token.and_then(Cursor::decode), |
| 107 |
limit: limit.unwrap_or(DEFAULT_PAGE_SIZE).clamp(1, MAX_PAGE_SIZE), |
| 108 |
} |
| 109 |
} |
| 110 |
|
| 111 |
|
| 112 |
|
| 113 |
pub fn fetch_limit(&self) -> i64 { |
| 114 |
i64::from(self.limit) + 1 |
| 115 |
} |
| 116 |
} |
| 117 |
|
| 118 |
|
| 119 |
#[derive(Debug, Serialize)] |
| 120 |
pub struct Page<T> { |
| 121 |
pub items: Vec<T>, |
| 122 |
pub next: Option<Cursor>, |
| 123 |
} |
| 124 |
|
| 125 |
impl<T> Page<T> { |
| 126 |
|
| 127 |
|
| 128 |
|
| 129 |
|
| 130 |
pub fn from_rows( |
| 131 |
mut rows: Vec<T>, |
| 132 |
params: &PageParams, |
| 133 |
cursor_of: impl Fn(&T) -> Cursor, |
| 134 |
) -> Self { |
| 135 |
let has_more = rows.len() as u32 > params.limit; |
| 136 |
rows.truncate(params.limit as usize); |
| 137 |
let next = if has_more { |
| 138 |
rows.last().map(&cursor_of) |
| 139 |
} else { |
| 140 |
None |
| 141 |
}; |
| 142 |
Page { items: rows, next } |
| 143 |
} |
| 144 |
|
| 145 |
|
| 146 |
pub fn next_token(&self) -> Option<String> { |
| 147 |
self.next.map(|c| c.encode()) |
| 148 |
} |
| 149 |
} |
| 150 |
|
| 151 |
#[cfg(test)] |
| 152 |
mod tests { |
| 153 |
use super::*; |
| 154 |
|
| 155 |
fn ts(millis: i64) -> DateTime<Utc> { |
| 156 |
DateTime::<Utc>::from_timestamp_millis(millis).unwrap() |
| 157 |
} |
| 158 |
|
| 159 |
#[test] |
| 160 |
fn cursor_round_trips_through_opaque_token() { |
| 161 |
let c = Cursor::new(ts(1_700_000_000_123), Uuid::from_u128(42)); |
| 162 |
let token = c.encode(); |
| 163 |
|
| 164 |
assert_ne!(token, "1700000000123:00000000-0000-0000-0000-00000000002a"); |
| 165 |
assert_eq!(Cursor::decode(&token), Some(c)); |
| 166 |
} |
| 167 |
|
| 168 |
#[test] |
| 169 |
fn malformed_cursor_decodes_to_none_not_error() { |
| 170 |
assert_eq!(Cursor::decode("not-base64!!!"), None); |
| 171 |
assert_eq!(Cursor::decode(""), None); |
| 172 |
|
| 173 |
let garbage = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode("no-colon"); |
| 174 |
assert_eq!(Cursor::decode(&garbage), None); |
| 175 |
} |
| 176 |
|
| 177 |
#[test] |
| 178 |
fn params_clamp_and_default_page_size() { |
| 179 |
assert_eq!(PageParams::new(None, None).limit, DEFAULT_PAGE_SIZE); |
| 180 |
assert_eq!(PageParams::new(None, Some(0)).limit, 1); |
| 181 |
assert_eq!(PageParams::new(None, Some(10_000)).limit, MAX_PAGE_SIZE); |
| 182 |
assert_eq!(PageParams::new(None, Some(25)).fetch_limit(), 26); |
| 183 |
} |
| 184 |
|
| 185 |
#[test] |
| 186 |
fn bad_cursor_token_starts_from_beginning() { |
| 187 |
let p = PageParams::new(Some("tampered"), Some(20)); |
| 188 |
assert!( |
| 189 |
p.after.is_none(), |
| 190 |
"a bad cursor must not error, just restart" |
| 191 |
); |
| 192 |
} |
| 193 |
|
| 194 |
#[test] |
| 195 |
fn from_rows_detects_next_page_and_trims_extra() { |
| 196 |
let params = PageParams::new(None, Some(3)); |
| 197 |
|
| 198 |
let rows: Vec<(i64, Uuid)> = (0..4).map(|i| (i, Uuid::from_u128(i as u128))).collect(); |
| 199 |
let page = Page::from_rows(rows, ¶ms, |(m, id)| Cursor::new(ts(*m), *id)); |
| 200 |
assert_eq!(page.items.len(), 3, "the extra probe row is trimmed"); |
| 201 |
|
| 202 |
assert_eq!(page.next, Some(Cursor::new(ts(2), Uuid::from_u128(2)))); |
| 203 |
} |
| 204 |
|
| 205 |
#[test] |
| 206 |
fn from_rows_last_page_has_no_next() { |
| 207 |
let params = PageParams::new(None, Some(5)); |
| 208 |
|
| 209 |
let rows: Vec<(i64, Uuid)> = (0..3).map(|i| (i, Uuid::from_u128(i as u128))).collect(); |
| 210 |
let page = Page::from_rows(rows, ¶ms, |(m, id)| Cursor::new(ts(*m), *id)); |
| 211 |
assert_eq!(page.items.len(), 3); |
| 212 |
assert_eq!(page.next, None); |
| 213 |
assert_eq!(page.next_token(), None); |
| 214 |
} |
| 215 |
} |
| 216 |
|