Skip to content
162 linesCodeBlameRaw

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-guardrails1//! 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
9use g1t_contracts::events::Event;
10use serde_json::Value;
11
12/// The most messages in one `sendBatch`.
13pub const MAX_BATCH_MESSAGES: usize = 100;
14/// The most bytes of JSON in one `sendBatch`, under the queue's 256 KB.
15pub const MAX_BATCH_BYTES: usize = 200 * 1024;
16/// The most bytes of JSON in one message, under the queue's 128 KB.
17pub const MAX_MESSAGE_BYTES: usize = 96 * 1024;
18/// The key set in `data` on an event whose long text was shortened.
19pub 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.
22const CUTS: [usize; 3] = [4096, 512, 128];
23
24/// An event's size as JSON.
25pub 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.
32pub 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
57fn 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
69fn 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`].
78pub 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)]
97mod 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.