Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails | 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 | } |
This file's history is long; its oldest lines are credited to the oldest commit read.