| 1 | //! Fitting events into queue messages, and batches into `sendBatch` calls. |
| 2 | //! |
| 3 | //! A queue takes up to 100 messages and 256 KB in one `sendBatch`, and up |
| 4 | //! to 128 KB in one message. Batches are cut by count and by size, with |
| 5 | //! room to spare for how the runtime encodes them. An event too large for |
| 6 | //! one message has its long text shortened and is marked `truncated`: |
| 7 | //! whoever needs the whole text reads it from the service that owns it. |
| 8 | |
| 9 | use g1t_contracts::events::Event; |
| 10 | use serde_json::Value; |
| 11 | |
| 12 | /// The most messages in one `sendBatch`. |
| 13 | pub const MAX_BATCH_MESSAGES: usize = 100; |
| 14 | /// The most bytes of JSON in one `sendBatch`, under the queue's 256 KB. |
| 15 | pub const MAX_BATCH_BYTES: usize = 200 * 1024; |
| 16 | /// The most bytes of JSON in one message, under the queue's 128 KB. |
| 17 | pub const MAX_MESSAGE_BYTES: usize = 96 * 1024; |
| 18 | /// The key set in `data` on an event whose long text was shortened. |
| 19 | pub const TRUNCATED: &str = "truncated"; |
| 20 | /// How long text is cut to, longest first, until the event fits. Ids and |
| 21 | /// names are shorter than the last, so they are never cut. |
| 22 | const CUTS: [usize; 3] = [4096, 512, 128]; |
| 23 | |
| 24 | /// An event's size as JSON. |
| 25 | pub fn size(event: &Event) -> usize { |
| 26 | serde_json::to_string(event).map_or(usize::MAX, |json| json.len()) |
| 27 | } |
| 28 | |
| 29 | /// Shortens `event`'s long strings until it fits in one message, and marks |
| 30 | /// it truncated if anything was cut. Ids, numbers and short fields stay, so |
| 31 | /// every consumer still knows what the event is about. |
| 32 | pub fn fit(event: &mut Event) { |
| 33 | if size(event) <= MAX_MESSAGE_BYTES { |
| 34 | return; |
| 35 | } |
| 36 | for cut in CUTS { |
| 37 | shorten(&mut event.data, cut); |
| 38 | mark(&mut event.data); |
| 39 | if size(event) <= MAX_MESSAGE_BYTES { |
| 40 | return; |
| 41 | } |
| 42 | } |
| 43 | // Still too large (a long list, say): only the top level's short |
| 44 | // values are kept. |
| 45 | let kept = match &event.data { |
| 46 | Value::Object(map) => map |
| 47 | .iter() |
| 48 | .filter(|(_, value)| matches!(value, Value::String(_) | Value::Number(_) | Value::Bool(_) | Value::Null)) |
| 49 | .map(|(key, value)| (key.clone(), value.clone())) |
| 50 | .collect(), |
| 51 | _ => serde_json::Map::new(), |
| 52 | }; |
| 53 | event.data = Value::Object(kept); |
| 54 | mark(&mut event.data); |
| 55 | } |
| 56 | |
| 57 | fn shorten(value: &mut Value, cut: usize) { |
| 58 | match value { |
| 59 | Value::String(text) if text.len() > cut => { |
| 60 | let end = (0..=cut).rev().find(|end| text.is_char_boundary(*end)).unwrap_or(0); |
| 61 | text.truncate(end); |
| 62 | } |
| 63 | Value::Array(items) => items.iter_mut().for_each(|item| shorten(item, cut)), |
| 64 | Value::Object(map) => map.values_mut().for_each(|item| shorten(item, cut)), |
| 65 | _ => {} |
| 66 | } |
| 67 | } |
| 68 | |
| 69 | fn mark(data: &mut Value) { |
| 70 | if let Value::Object(map) = data { |
| 71 | map.insert(TRUNCATED.to_owned(), Value::Bool(true)); |
| 72 | } |
| 73 | } |
| 74 | |
| 75 | /// `events` cut into `sendBatch` calls, in order: no more than |
| 76 | /// [`MAX_BATCH_MESSAGES`] and [`MAX_BATCH_BYTES`] each. Events are expected |
| 77 | /// to have been [`fit`]. |
| 78 | pub fn chunks<'a>(events: &[&'a Event]) -> Vec<Vec<&'a Event>> { |
| 79 | let mut chunks: Vec<Vec<&Event>> = Vec::new(); |
| 80 | let mut bytes = 0; |
| 81 | for event in events { |
| 82 | let event_bytes = size(event); |
| 83 | let full = chunks |
| 84 | .last() |
| 85 | .is_none_or(|chunk| chunk.len() >= MAX_BATCH_MESSAGES || bytes + event_bytes > MAX_BATCH_BYTES); |
| 86 | if full { |
| 87 | chunks.push(Vec::new()); |
| 88 | bytes = 0; |
| 89 | } |
| 90 | chunks.last_mut().expect("just pushed").push(event); |
| 91 | bytes += event_bytes; |
| 92 | } |
| 93 | chunks |
| 94 | } |
| 95 | |
| 96 | #[cfg(test)] |
| 97 | mod tests { |
| 98 | use super::*; |
| 99 | use serde_json::json; |
| 100 | |
| 101 | fn event(data: Value) -> Event { |
| 102 | Event { |
| 103 | id: "evt_1".into(), |
| 104 | kind: "comment.edited".into(), |
| 105 | source: "work".into(), |
| 106 | time: "2026-10-08T00:00:00Z".into(), |
| 107 | repo_id: Some("rep_1".into()), |
| 108 | actor: Some("usr_1".into()), |
| 109 | data, |
| 110 | } |
| 111 | } |
| 112 | |
| 113 | #[test] |
| 114 | fn a_small_event_is_left_as_it_is() { |
| 115 | let mut small = event(json!({ "number": 4, "changes": { "body": { "from": "hello" } } })); |
| 116 | fit(&mut small); |
| 117 | assert_eq!(small.data, json!({ "number": 4, "changes": { "body": { "from": "hello" } } })); |
| 118 | } |
| 119 | |
| 120 | #[test] |
| 121 | fn a_large_event_keeps_its_shape_and_ids_with_its_text_shortened() { |
| 122 | let mut large = event(json!({ |
| 123 | "commentId": "cmt_1", |
| 124 | "number": 4, |
| 125 | "comment": { "body": "é".repeat(100_000), "author": "ana" }, |
| 126 | })); |
| 127 | fit(&mut large); |
| 128 | assert!(size(&large) <= MAX_MESSAGE_BYTES); |
| 129 | assert_eq!(large.data[TRUNCATED], true); |
| 130 | assert_eq!(large.data["commentId"], "cmt_1"); |
| 131 | assert_eq!(large.data["number"], 4); |
| 132 | // Still an object, so whoever reads its fields still can. |
| 133 | assert_eq!(large.data["comment"]["author"], "ana"); |
| 134 | assert!(large.data["comment"]["body"].as_str().unwrap().len() <= 4096); |
| 135 | } |
| 136 | |
| 137 | #[test] |
| 138 | fn an_event_too_large_even_shortened_keeps_its_top_level_values() { |
| 139 | let many: Vec<Value> = (0..20_000).map(|n| json!({ "n": n })).collect(); |
| 140 | let mut huge = event(json!({ "pullId": "pul_1", "files": many })); |
| 141 | fit(&mut huge); |
| 142 | assert!(size(&huge) <= MAX_MESSAGE_BYTES); |
| 143 | assert_eq!(huge.data, json!({ "pullId": "pul_1", "truncated": true })); |
| 144 | } |
| 145 | |
| 146 | #[test] |
| 147 | fn batches_are_cut_by_count_and_by_size() { |
| 148 | let small: Vec<Event> = (0..250).map(|_| event(json!({ "number": 1 }))).collect(); |
| 149 | let refs: Vec<&Event> = small.iter().collect(); |
| 150 | let counted: Vec<usize> = chunks(&refs).iter().map(Vec::len).collect(); |
| 151 | assert_eq!(counted, [100, 100, 50]); |
| 152 | |
| 153 | let big: Vec<Event> = (0..5).map(|_| event(json!({ "text": "x".repeat(80 * 1024) }))).collect(); |
| 154 | let refs: Vec<&Event> = big.iter().collect(); |
| 155 | let cut = chunks(&refs); |
| 156 | assert_eq!(cut.iter().map(Vec::len).collect::<Vec<_>>(), [2, 2, 1]); |
| 157 | for chunk in cut { |
| 158 | assert!(chunk.iter().map(|event| size(event)).sum::<usize>() <= MAX_BATCH_BYTES); |
| 159 | } |
| 160 | assert!(chunks(&[]).is_empty()); |
| 161 | } |
| 162 | } |