Skip to content

Commit

Merge branch 'worktree-agent-ad8a36dfcd4176015' into spend-guardrails

syntaqxcommitted Parents7aafb03db3e845Browse files
41 files+1038−1350/41 viewed
+7−0
10631063 or a comment made with it runs nothing, so a workflow cannot set itself
10641064 off. `workflow_dispatch` and [`repository_dispatch`](#repository-dispatch)
10651065 are the exceptions, for a workflow that means to start another.
1066+- It **never puts g1t to work**. A comment it posts that mentions
1067+ `@g1t` starts nothing, and it cannot assign an issue or a plan to g1t,
1068+ queue one for it, hand it work or ask it for a review. Otherwise a
1069+ workflow that asks g1t to fix a failing check would run again on g1t's
1070+ push, and ask again, without end. A step that should put g1t to work
1071+ uses a token of a person's own, stored as a
1072+ [secret](/guides/secrets-and-variables/).
10661073
10671074 `permissions:` goes at the top of the workflow, for every job, or on a job,
10681075 which then ignores the workflow's. Once either is written, every permission
+15−0
184184
185185 Each push sends only what the one before did not.
186186
187+### Pushes of many branches or tags
188+
189+Every branch and tag in a push is stored. Each one is also announced as a
190+`git.push` event, which starts workflows, mirrors the repository and
191+calls webhooks, except in a push of many:
192+
193+| A push of | What is announced |
194+| --- | --- |
195+| Up to 3 tags | Each tag |
196+| More than 3 tags (`git push --tags`, say) | None of the tags |
197+| Up to 1,000 branches | Each branch |
198+| More than 1,000 branches | Only the default branch, if it moved |
199+
200+To have tags start workflows, push them 3 or fewer at a time.
201+
187202 ### When the store is busy
188203
189204 If Cloudflare Artifacts is rate limiting g1t or not answering, g1t tries
+4−2
4848
4949 `data` holds what the event is about: ids and numbers to fetch the rest with
5050 the [API](/reference/api/). `actor` is null for something g1t did by
51−itself.
51+itself. An event whose `data` would be larger than about 96 KB has its
52+long text, such as a comment's or release's body, shortened, and
53+`data.truncated` is `true`: fetch the whole text from the API.
5254
5355 With these headers:
5456
6466
6567 | Event | When |
6668 | --- | --- |
67−| `git.push` | A branch moved. `data.ref`, `data.after`, `data.default_branch`. |
69+| `git.push` | A branch moved. `data.ref`, `data.after`, `data.default_branch`. A push of more than 3 tags, or more than 1,000 branches, sends fewer; see [pushes of many branches or tags](/guides/git/#pushes-of-many-branches-or-tags). |
6870 | `repo.created`, `repo.forked` | A repository was made, or forked for a pull request. |
6971 | `repo.updated` | Its description, website, topics, protection or visibility changed. |
7072 | `repo.visibility_changed` | It was made public or private. |
+10−1
414414 (`ops@g1t.sh`), in URLs, in package scopes (`@g1t/platform`) and in longer
415415 names (`@g1t-bot`) are ignored. Matching ignores case. Agents mentioning
416416 `@g1t` start nothing, so
417−agents cannot set each other to work this way.
417+agents cannot set each other to work this way. Neither does a comment made
418+with a workflow job's own token (`G1T_TOKEN`), so a workflow cannot set
419+off g1t whose push sets off the workflow again; see
420+[the job's token](/guides/actions/#the-jobs-token).
418421
419422 **Who can.** People with the Write [role](/guides/access-and-roles/) or higher on the
420423 repository, members or not. Anyone else who mentions it gets a short reply
428431
429432 Each comment starts one run at most; to ask again, write a new comment.
430433
434+A mention sends g1t back to a pull request it made even after it stopped
435+there, or used up the repository's revisions. At most 10 mentions in a day
436+send it back to the same pull request; past that, g1t replies that it
437+has reached the most it takes, and the next one works a day after the
438+first of those 10.
439+
431440 ## The label rule
432441
433442 Under a project's **Settings → Agents**, someone with the Maintain role or
+1−0
3939 pub mod security;
4040 pub mod teams;
4141 pub mod security_suite;
42+pub mod subscribers;
4243 pub mod time;
4344 pub mod tokens;
4445 pub mod updates;
+308−0
1+//! Which events each subscribing service's queue is sent.
2+//!
3+//! The events service passes each batch on to one queue per subscriber
4+//! (`SUBSCRIBER_*` bindings). Each is sent only the types its consumer acts
5+//! on, so an event costs a queue write for the services that read it, not
6+//! for all of them. A consumer that starts acting on a new type adds it
7+//! here in the same change, or it never hears of it; the tests below and
8+//! beside each consumer check what they can.
9+//!
10+//! A pattern is an exact type (`git.push`), a family (`pull.*`, every type
11+//! that starts `pull.`), or `*`, every type. A binding not named here is
12+//! sent everything, so a new subscriber is never left out.
13+
14+/// What the shared helpers in `g1t_kit` act on (`rename`, `transfer`,
15+/// `lifecycle`, `deleted`, `user_deleted`): rows that follow a workspace
16+/// or repository, or go with it. Rare, so every subscriber is sent them.
17+const LIFECYCLE: [&str; 8] = [
18+ "workspace.renamed",
19+ "workspace.deleted",
20+ "repo.renamed",
21+ "repo.transferred",
22+ "repo.deleted",
23+ "repo.restored",
24+ "repo.purged",
25+ "user.deleted",
26+];
27+
28+/// Each subscriber's binding and the types it is sent beside [`LIFECYCLE`].
29+pub const ROUTES: &[(&str, &[&str])] = &[
30+ // Memory captured from merges and comments, mentions, pull request
31+ // bases, runs stopped in archived repositories (work/src/lib.rs queue).
32+ (
33+ "SUBSCRIBER_WORK",
34+ &[
35+ "git.push",
36+ "branch.renamed",
37+ "repo.created",
38+ "repo.archived",
39+ "repo.default_branch_changed",
40+ "pull.merged",
41+ "comment.created",
42+ ],
43+ ),
44+ // Pull requests g1t is seeing through, mentions, queued issues and the
45+ // merge queue (runner/src/index.ts queue).
46+ (
47+ "SUBSCRIBER_RUNNER",
48+ &[
49+ "pull.opened",
50+ "pull.ready",
51+ "pull.updated",
52+ "pull.mergeability",
53+ "pull.mergecheck",
54+ "pull.merge_requested",
55+ "pull.merged",
56+ "pull.closed",
57+ "checks.completed",
58+ "review.completed",
59+ "queue.changed",
60+ "comment.created",
61+ "issue.opened",
62+ "issue.updated",
63+ "issue.closed",
64+ "agent.asked",
65+ "repo.archived",
66+ ],
67+ ),
68+ // Push mirrors and write-back to linked trackers (integrations/src).
69+ ("SUBSCRIBER_INTEGRATIONS", &["git.push", "issue.closed", "pull.opened"]),
70+ // A hook may ask for any type, or for every one.
71+ ("SUBSCRIBER_WEBHOOKS", &["*"]),
72+ // Every event a workflow can run on (g1t_actions::events::github_events),
73+ // and runs stopped in archived repositories.
74+ (
75+ "SUBSCRIBER_ACTIONS",
76+ &[
77+ "git.push",
78+ "pull.*",
79+ "issue.*",
80+ "comment.*",
81+ "release.*",
82+ "deployment.created",
83+ "deployment_status.created",
84+ "review.completed",
85+ "workflow.completed",
86+ "repo.archived",
87+ ],
88+ ),
89+ // Previews and production deploys (deployments/src/index.ts queue).
90+ (
91+ "SUBSCRIBER_DEPLOYMENTS",
92+ &[
93+ "git.push",
94+ "branch.renamed",
95+ "repo.default_branch_changed",
96+ "pull.opened",
97+ "pull.ready",
98+ "pull.updated",
99+ "pull.reopened",
100+ "pull.closed",
101+ "pull.merged",
102+ "workspace.deleting",
103+ "workspace.restored",
104+ ],
105+ ),
106+ // Projects follow their repositories, and wake on activity
107+ // (projects/src/index.ts, activity.ts).
108+ (
109+ "SUBSCRIBER_PROJECTS",
110+ &[
111+ "git.push",
112+ "repo.*",
113+ "issue.opened",
114+ "issue.closed",
115+ "issue.reopened",
116+ "pull.opened",
117+ "pull.updated",
118+ "pull.merged",
119+ "pull.closed",
120+ "review.completed",
121+ "comment.created",
122+ "deployment.succeeded",
123+ "deployment.failed",
124+ "package.published",
125+ "package.deleted",
126+ "package.version_deleted",
127+ "workspace.deleting",
128+ "workspace.restored",
129+ ],
130+ ),
131+ // Pull request working copies, workspaces' repositories, contributors
132+ // shown as ghost (repos/src/lib.rs handle_event).
133+ (
134+ "SUBSCRIBER_REPOS",
135+ &[
136+ "pull.merged",
137+ "pull.closed",
138+ "pull.reopened",
139+ "workspace.deleting",
140+ "workspace.restored",
141+ "user.deleting",
142+ "user.restored",
143+ ],
144+ ),
145+ // Usage rows that follow a renamed workspace or repository.
146+ ("SUBSCRIBER_BILLING", &[]),
147+ // Secret and dependency scans, and version update pull requests
148+ // (security/src/lib.rs on_event).
149+ (
150+ "SUBSCRIBER_SECURITY",
151+ &[
152+ "git.push",
153+ "pull.opened",
154+ "pull.updated",
155+ "pull.ready",
156+ "pull.merged",
157+ "pull.closed",
158+ "checks.completed",
159+ "comment.created",
160+ "repo.created",
161+ "repo.visibility_changed",
162+ ],
163+ ),
164+ // What agents are told about a repository (context/src/index.ts).
165+ (
166+ "SUBSCRIBER_CONTEXT",
167+ &[
168+ "git.push",
169+ "issue.opened",
170+ "issue.updated",
171+ "issue.closed",
172+ "pull.ready",
173+ "pull.merged",
174+ "memory.changed",
175+ ],
176+ ),
177+ // The site-wide search index (search/src/index.rs).
178+ ("SUBSCRIBER_SEARCH", &["git.push", "repo.*", "issue.*", "pull.*", "user.*", "workspace.*"]),
179+ // Composer packages read from repositories, and packages that follow
180+ // their workspace, team or repository (packages/src/lib.rs on_event).
181+ (
182+ "SUBSCRIBER_PACKAGES",
183+ &[
184+ "git.push",
185+ "repo.visibility_changed",
186+ "workspace.deleting",
187+ "workspace.restored",
188+ "team.edited",
189+ "team.deleted",
190+ ],
191+ ),
192+];
193+
194+/// Whether `pattern` (see the module's notes) matches the type `kind`.
195+fn matches(pattern: &str, kind: &str) -> bool {
196+ match pattern.strip_suffix('*') {
197+ Some(prefix) => kind.starts_with(prefix),
198+ None => pattern == kind,
199+ }
200+}
201+
202+/// Whether the subscriber behind `binding` is sent events of type `kind`.
203+pub fn routed(binding: &str, kind: &str) -> bool {
204+ match ROUTES.iter().find(|(name, _)| *name == binding) {
205+ Some((_, patterns)) => {
206+ LIFECYCLE.contains(&kind) || patterns.iter().any(|pattern| matches(pattern, kind))
207+ }
208+ None => true,
209+ }
210+}
211+
212+#[cfg(test)]
213+mod tests {
214+ use super::*;
215+
216+ /// Types published on the bus that hooks are not offered.
217+ const UNOFFERED: [&str; 18] = [
218+ "pull.mergecheck",
219+ "pull.mergeability",
220+ "deployment.review_requested",
221+ "memory.changed",
222+ "abuse.flagged",
223+ "invite.created",
224+ "invite.redeemed",
225+ "waitlist.requested",
226+ "user.updated",
227+ "user.deleting",
228+ "user.restored",
229+ "user.deleted",
230+ "user.email_added",
231+ "workspace.updated",
232+ "workspace.renamed",
233+ "workspace.deleting",
234+ "workspace.restored",
235+ "workspace.deleted",
236+ ];
237+
238+ fn published(kind: &str) -> bool {
239+ crate::webhooks::EVENT_TYPES.contains(&kind) || UNOFFERED.contains(&kind)
240+ }
241+
242+ #[test]
243+ fn every_exact_route_names_a_type_that_is_published() {
244+ for (binding, patterns) in ROUTES {
245+ for pattern in patterns.iter().chain(LIFECYCLE.iter()) {
246+ if pattern.ends_with('*') {
247+ let prefix = pattern.trim_end_matches('*');
248+ assert!(
249+ prefix.is_empty()
250+ || crate::webhooks::EVENT_TYPES.iter().chain(UNOFFERED.iter()).any(|kind| kind.starts_with(prefix)),
251+ "{binding}: {pattern} matches nothing"
252+ );
253+ } else {
254+ assert!(published(pattern), "{binding}: {pattern} is not a type anything publishes");
255+ }
256+ }
257+ }
258+ }
259+
260+ #[test]
261+ fn every_route_is_a_subscriber_binding_named_once() {
262+ let mut names: Vec<&str> = ROUTES.iter().map(|(name, _)| *name).collect();
263+ assert!(names.iter().all(|name| name.starts_with("SUBSCRIBER_")));
264+ let count = names.len();
265+ names.sort_unstable();
266+ names.dedup();
267+ assert_eq!(names.len(), count);
268+ assert_eq!(count, 13);
269+ }
270+
271+ #[test]
272+ fn subscribers_hear_what_they_act_on_and_not_the_rest() {
273+ assert!(routed("SUBSCRIBER_WEBHOOKS", "secret_scanning_alert.created"));
274+ assert!(routed("SUBSCRIBER_WEBHOOKS", "anything.new"));
275+ assert!(routed("SUBSCRIBER_ACTIONS", "pull.labeled"));
276+ assert!(routed("SUBSCRIBER_ACTIONS", "release.published"));
277+ assert!(!routed("SUBSCRIBER_ACTIONS", "check_run.completed"));
278+ assert!(routed("SUBSCRIBER_RUNNER", "comment.created"));
279+ assert!(!routed("SUBSCRIBER_RUNNER", "status.created"));
280+ assert!(routed("SUBSCRIBER_BILLING", "workspace.renamed"));
281+ assert!(!routed("SUBSCRIBER_BILLING", "git.push"));
282+ assert!(routed("SUBSCRIBER_SEARCH", "user.updated"));
283+ assert!(routed("SUBSCRIBER_REPOS", "user.deleting"));
284+ assert!(!routed("SUBSCRIBER_REPOS", "git.push"));
285+ // Every subscriber follows what moves or removes a repository.
286+ for (binding, _) in ROUTES {
287+ for kind in LIFECYCLE {
288+ assert!(routed(binding, kind), "{binding} misses {kind}");
289+ }
290+ }
291+ }
292+
293+ #[test]
294+ fn a_subscriber_not_in_the_table_hears_everything() {
295+ assert!(routed("SUBSCRIBER_NEW", "git.push"));
296+ assert!(routed("SUBSCRIBER_NEW", "status.created"));
297+ }
298+
299+ #[test]
300+ fn a_family_matches_only_its_own_types() {
301+ assert!(matches("pull.*", "pull.opened"));
302+ assert!(!matches("pull.*", "pulls.opened"));
303+ assert!(!matches("pull.*", "pull"));
304+ assert!(matches("*", "x"));
305+ assert!(matches("git.push", "git.push"));
306+ assert!(!matches("git.push", "git.pushed"));
307+ }
308+}
+3−0
5858 "d1": { "database": "g1t-events", "migrations": "migrations" },
5959 "stage": "core",
6060 "secrets": [],
61+ "setup": [
62+ "The dead-letter queue every queue consumer sends what it gave up on to, before any unit that names it deploys: npx wrangler queues create g1t-events-dlq"
63+ ],
6164 "self_host": "run"
6265 },
6366 "identity": {
+19−1
611611 `resources.kv`, then the ids in the configs.
612612 - Queues: `npx wrangler queues create <queue>` for each queue in the
613613 manifest: `g1t-events`, `g1t-events-<service>` for every subscriber,
614− `g1t-search-jobs`, `g1t-context-jobs`.
614+ `g1t-search-jobs`, `g1t-context-jobs`, and the dead-letter queue
615+ `g1t-events-dlq` (below).
615616 - R2: `npx wrangler r2 bucket create g1t-screenshots`,
616617 `npx wrangler r2 bucket create g1t-actions-cache` (with its two
617618 lifecycle rules, above), and `npx wrangler r2 bucket create g1t-git-packs`,
631632 that does not exist yet may be refused; deploy that service first with
632633 `--only`.
633634
635+## The dead-letter queue
636+
637+Every queue consumer (`g1t-events`, each `g1t-events-<service>`,
638+`g1t-search-jobs`, `g1t-context-jobs`) names `max_retries` and the
639+dead-letter queue `g1t-events-dlq`, so a message that keeps failing stops
640+after its retries instead of being retried for ever. A consumer bound to a
641+queue that does not exist fails to deploy, so create it once per account,
642+before the first deploy that names it:
643+
644+```sh
645+npx wrangler queues create g1t-events-dlq
646+```
647+
648+Nothing consumes it: read what landed there with `npx wrangler queues
649+info g1t-events-dlq`, fix the cause, and replay by hand if needed. The
650+manifest lists each unit's dead-letter queues under `queues.dead_letter`.
651+
634652 ## OIDC tokens for workflow jobs
635653
636654 The API is the issuer of workflow jobs' OIDC tokens,
+5−1
471471 }
472472 const out = stack.units.map(({ config, ...unit }) => ({
473473 ...unit,
474− queues: { produces: (config?.queues?.producers ?? []).map((q) => q.queue), consumes: (config?.queues?.consumers ?? []).map((q) => q.queue) },
474+ queues: {
475+ produces: (config?.queues?.producers ?? []).map((q) => q.queue),
476+ consumes: (config?.queues?.consumers ?? []).map((q) => q.queue),
477+ dead_letter: [...new Set((config?.queues?.consumers ?? []).map((q) => q.dead_letter_queue).filter(Boolean))],
478+ },
475479 kv: (config?.kv_namespaces ?? []).map((kv) => stack.resources.kv?.[kv.id] ?? kv.id),
476480 r2: (config?.r2_buckets ?? []).map((b) => b.bucket_name),
477481 vectorize: (config?.vectorize ?? []).map((v) => v.index_name),
+2−0
3333 echo "== Dispatch namespace and event queue"
3434 w dispatch-namespace list 2>/dev/null | grep -q "$NAMESPACE" || w dispatch-namespace create "$NAMESPACE"
3535 w queues list 2>/dev/null | grep -q g1t-events-deployments || w queues create g1t-events-deployments
36+# Where every consumer sends what it gave up on (docs/DEPLOYING.md).
37+w queues list 2>/dev/null | grep -q g1t-events-dlq || w queues create g1t-events-dlq
3638
3739 echo "== DNS record and the service's token"
3840 if [ -n "${CLOUDFLARE_API_KEY:-}" ] && [ -n "${CLOUDFLARE_EMAIL:-}" ]; then
+12−0
11141114 mod tests {
11151115 use super::*;
11161116 use g1t_contracts::work::{Pull, g1t_author};
1117+
1118+ #[test]
1119+ fn every_event_a_workflow_runs_on_is_sent_to_this_queue() {
1120+ for kind in g1t_contracts::webhooks::EVENT_TYPES {
1121+ if !github_events(kind).is_empty() {
1122+ assert!(
1123+ g1t_contracts::subscribers::routed("SUBSCRIBER_ACTIONS", kind),
1124+ "{kind} starts workflows but is not routed to actions (g1t_contracts::subscribers)"
1125+ );
1126+ }
1127+ }
1128+ }
11171129 use g1t_contracts::{Membership, PrincipalKind};
11181130
11191131 fn repo() -> Repo {
+1−1
3737 "r2_buckets": [{ "binding": "ACTIONS_CACHE", "bucket_name": "g1t-actions-cache" }],
3838 // Every event on the bus: what starts workflows, and pushes that change them.
3939 "queues": {
40− "consumers": [{ "queue": "g1t-events-actions", "max_batch_size": 10, "max_batch_timeout": 1, "max_retries": 3 }]
40+ "consumers": [{ "queue": "g1t-events-actions", "max_batch_size": 10, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
4141 },
4242 // Schedules, jobs waiting for room, and jobs whose runner went quiet.
4343 "triggers": { "crons": ["* * * * *"] },
+1−1
3333 // Events from the bus: a workspace renamed moves its billing to the
3434 // new slug (src/rename.rs).
3535 "queues": {
36− "consumers": [{ "queue": "g1t-events-billing", "max_batch_size": 20, "max_batch_timeout": 1 }]
36+ "consumers": [{ "queue": "g1t-events-billing", "max_batch_size": 20, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
3737 },
3838 "vars": {
3939 // What is added to a model run's cost, in percent: g1t's overhead for
+2−2
3737 // jobs, one project of a backfill each, so no request does them all.
3838 "producers": [{ "binding": "JOBS", "queue": "g1t-context-jobs" }],
3939 "consumers": [
40− { "queue": "g1t-events-context", "max_batch_size": 20, "max_batch_timeout": 2 },
41− { "queue": "g1t-context-jobs", "max_batch_size": 1, "max_batch_timeout": 1, "max_retries": 2 }
40+ { "queue": "g1t-events-context", "max_batch_size": 20, "max_batch_timeout": 2, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" },
41+ { "queue": "g1t-context-jobs", "max_batch_size": 1, "max_batch_timeout": 1, "max_retries": 2, "dead_letter_queue": "g1t-events-dlq" }
4242 ]
4343 },
4444 "observability": { "enabled": true }
+1−1
3434 // Pull requests opened, pushed to, closed and merged; pushes to the
3535 // default branch.
3636 "queues": {
37− "consumers": [{ "queue": "g1t-events-deployments", "max_batch_size": 20, "max_batch_timeout": 1 }]
37+ "consumers": [{ "queue": "g1t-events-deployments", "max_batch_size": 20, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
3838 },
3939 // Builds that died, usage counts, idle previews, apps of ended plans,
4040 // and charging a month that is over.
+10−0
1+-- Which subscriber queues an event already reached, written only when
2+-- passing a batch on failed part way: its retry skips the queues that
3+-- have it. Rows older than a day are removed by the daily sweep.
4+CREATE TABLE fanout_sent (
5+ event_id TEXT PRIMARY KEY,
6+ -- JSON array of SUBSCRIBER_* binding names.
7+ bindings TEXT NOT NULL,
8+ created_at TEXT NOT NULL
9+);
10+CREATE INDEX fanout_sent_by_age ON fanout_sent (created_at);
+162−0
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+}
+131−27
22 //! and its durable log.
33 //!
44 //! Publishing puts events on a queue and returns. The queue consumer writes
5−//! them to the log and passes each batch on to every subscriber's own
6−//! queue, so a slow or failing subscriber holds up nobody else.
5+//! them to the log and passes each batch on to each subscriber's own
6+//! queue, so a slow or failing subscriber holds up nobody else. Each
7+//! subscriber is sent only the types it acts on
8+//! (`g1t_contracts::subscribers`), in batches the queues take (fanout.rs).
79 //!
810 //! It keeps two more things beside the log: the audit log (audit.rs) and
911 //! each person's inbox (inbox.rs), written as events arrive, with who
1416 //! arguments.
1517
1618 mod audit;
19+mod fanout;
1720 mod inbox;
1821 mod subscriptions;
1922
2225 use g1t_contracts::time::rfc3339;
2326 use g1t_kit::{args, js, now_ms, reply, rpc_method};
2427 use serde::Deserialize;
28+use std::collections::{HashMap, HashSet};
2529 use worker::js_sys::{Array, Object};
2630 use worker::wasm_bindgen::{JsCast, JsValue};
2731 use worker::{Context, D1Database, Env, MessageBatch, Request, Response, Result, event};
3943 "check_suite.completed",
4044 "check_suite.rerequested",
4145 ];
42−/// Every binding whose name starts with this is a queue that receives all
43−/// events: one per subscribing service.
46+/// Every binding whose name starts with this is a queue that receives
47+/// events: one per subscribing service, sent the types it routes.
4448 const SUBSCRIBER_PREFIX: &str = "SUBSCRIBER_";
49+/// How long a record of which queues a batch reached is kept for its retries.
50+const FANOUT_KEEP_MS: u64 = 24 * 60 * 60 * 1000;
4551
4652 #[derive(Deserialize)]
4753 struct EventRow {
7480 value.as_deref().map_or(JsValue::NULL, JsValue::from)
7581 }
7682
77−/// Sends `events` to a queue binding, one message each.
78−async fn send(queue: &JsValue, events: &[Event]) -> Result<()> {
79− let messages = Array::new();
80− for event in events {
81− let message = Object::new();
82− js::set(&message, "body", &js::to_js(event)?);
83− messages.push(&message);
83+#[derive(Deserialize)]
84+struct FanoutRow {
85+ event_id: String,
86+ /// JSON array of binding names.
87+ bindings: String,
88+}
89+
90+/// Sends `events` to a queue binding, one message each, in as many
91+/// `sendBatch` calls as the queue's limits need (fanout.rs). Events are
92+/// expected to have been fit to a message.
93+async fn send(queue: &JsValue, events: &[&Event]) -> Result<()> {
94+ for chunk in fanout::chunks(events) {
95+ let messages = Array::new();
96+ for event in chunk {
97+ let message = Object::new();
98+ js::set(&message, "body", &js::to_js(event)?);
99+ messages.push(&message);
100+ }
101+ js::call(queue, "sendBatch", &[messages.into()]).await?;
84102 }
85− js::call(queue, "sendBatch", &[messages.into()]).await?;
86103 Ok(())
87104 }
88105
101118 let events: Vec<Event> = a
102119 .events
103120 .into_iter()
104− .map(|event| Event {
105− id: new_id("evt", now),
106− kind: event.kind,
107− source: event.source,
108− time: rfc3339(now),
109− repo_id: event.repo_id,
110− actor: event.actor,
111− data: event.data,
121+ .map(|event| {
122+ let mut event = Event {
123+ id: new_id("evt", now),
124+ kind: event.kind,
125+ source: event.source,
126+ time: rfc3339(now),
127+ repo_id: event.repo_id,
128+ actor: event.actor,
129+ data: event.data,
130+ };
131+ // Too large for one message: its long text is shortened.
132+ fanout::fit(&mut event);
133+ event
112134 })
113135 .collect();
136+ let events: Vec<&Event> = events.iter().collect();
114137 send(&js::binding(&self.env, "BUS")?, &events).await
115138 }
116139
138161 Ok(rows.into_iter().map(Event::from).collect())
139162 }
140163
141− /// Writes a batch from the bus to the log, hands it to every
142− /// subscriber, then tells the people it concerns (inbox.rs).
164+ /// Writes a batch from the bus to the log, notes who now follows what
165+ /// (inbox.rs), hands each subscriber the events it routes, then tells
166+ /// the people it concerns.
167+ ///
168+ /// Everything before the hand-off is safe to repeat. A retried batch
169+ /// skips the queues that already have it: when a hand-off fails part
170+ /// way, which queues each event reached is written down first.
143171 async fn deliver(&self, events: &[Event]) -> Result<()> {
144− let mut statements = Vec::with_capacity(events.len());
172+ let mut statements = Vec::with_capacity(events.len() + 1);
145173 for event in events {
146174 statements.push(
147175 self.db
161189 ])?,
162190 );
163191 }
164− self.db.batch(statements).await?;
192+ let ids: Vec<&str> = events.iter().map(|event| event.id.as_str()).collect();
193+ // Read in the same round trip: nearly always nothing.
194+ statements.push(
195+ self.db
196+ .prepare("SELECT event_id, bindings FROM fanout_sent WHERE event_id IN (SELECT value FROM json_each(?))")
197+ .bind(&[serde_json::to_string(&ids)?.into()])?,
198+ );
199+ let results = self.db.batch(statements).await?;
200+ let mut sent: HashMap<String, HashSet<String>> = HashMap::new();
201+ if let Some(rows) = results.last() {
202+ for row in rows.results::<FanoutRow>()? {
203+ let bindings: Vec<String> = serde_json::from_str(&row.bindings).unwrap_or_default();
204+ sent.insert(row.event_id, bindings.into_iter().collect());
205+ }
206+ }
165207 audit::follow_renames(&self.db, events).await?;
208+ // Before the hand-off, so its failing does not send the batch to
209+ // every queue again.
210+ inbox::follow(&self.db, events).await?;
166211
167212 let bindings: &JsValue = self.env.as_ref();
213+ let mut reached: Vec<String> = Vec::new();
214+ let mut failed = None;
168215 for name in Object::keys(bindings.unchecked_ref::<Object>()).iter() {
169216 let Some(name) = name
170217 .as_string()
172219 else {
173220 continue;
174221 };
175− send(&js::binding(&self.env, &name)?, events).await?;
222+ let routed: Vec<&Event> = events
223+ .iter()
224+ .filter(|event| g1t_contracts::subscribers::routed(&name, &event.kind))
225+ .filter(|event| !sent.get(&event.id).is_some_and(|reached| reached.contains(&name)))
226+ .collect();
227+ if routed.is_empty() {
228+ continue;
229+ }
230+ match send(&js::binding(&self.env, &name)?, &routed).await {
231+ Ok(()) => reached.push(name),
232+ Err(error) => {
233+ worker::console_error!("events: passing {} events to {name} failed: {error}", routed.len());
234+ failed = Some(error);
235+ }
236+ }
237+ }
238+ if let Some(error) = failed {
239+ self.note_reached(events, &sent, &reached).await?;
240+ return Err(error);
176241 }
177− inbox::follow(&self.db, events).await?;
178242 let (work, repos, identity) = (
179243 self.env.service("WORK")?,
180244 self.env.service("REPOS")?,
188252 inbox::deliver(&self.db, &sources, events).await;
189253 Ok(())
190254 }
255+
256+ /// Writes down which queues each event of a batch has reached, with
257+ /// those it had reached before, for the batch's retry.
258+ async fn note_reached(
259+ &self,
260+ events: &[Event],
261+ sent: &HashMap<String, HashSet<String>>,
262+ reached: &[String],
263+ ) -> Result<()> {
264+ let now = rfc3339(now_ms());
265+ let mut statements = Vec::with_capacity(events.len());
266+ for event in events {
267+ let mut bindings: Vec<&str> = reached.iter().map(String::as_str).collect();
268+ if let Some(before) = sent.get(&event.id) {
269+ bindings.extend(before.iter().map(String::as_str));
270+ }
271+ bindings.sort_unstable();
272+ bindings.dedup();
273+ statements.push(
274+ self.db
275+ .prepare(
276+ "INSERT INTO fanout_sent (event_id, bindings, created_at) VALUES (?, ?, ?)
277+ ON CONFLICT (event_id) DO UPDATE SET bindings = excluded.bindings",
278+ )
279+ .bind(&[event.id.as_str().into(), serde_json::to_string(&bindings)?.into(), now.as_str().into()])?,
280+ );
281+ }
282+ self.db.batch(statements).await?;
283+ Ok(())
284+ }
191285 }
192286
193287 #[event(fetch)]
234328 }
235329
236330 /// Events from the bus. A batch that fails is retried whole, which is safe
237−/// because both the log and subscribers ignore an event they have seen.
331+/// because the log ignores an event it has seen, and queues that already
332+/// have it are skipped (`deliver`). After its retries it goes to the
333+/// dead-letter queue (wrangler.jsonc).
238334 #[event(queue)]
239335 async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
240336 let events = Events {
269365 let min_days = days("AUDIT_MIN_DAYS", audit::DEFAULT_MIN_DAYS).min(max_days);
270366 let Ok(db) = env.d1("DB") else { return };
271367 let now = now_ms();
368+ let cutoff = rfc3339(now.saturating_sub(FANOUT_KEEP_MS));
369+ let purged = match db.prepare("DELETE FROM fanout_sent WHERE created_at < ?").bind(&[cutoff.into()]) {
370+ Ok(statement) => statement.run().await.map(|_| ()),
371+ Err(error) => Err(error),
372+ };
373+ if let Err(error) = purged {
374+ worker::console_error!("could not remove old fan-out records: {error}");
375+ }
272376 match inbox::purge(&db, now).await {
273377 Ok(removed) if removed > 0 => worker::console_log!("removed {removed} old inbox items"),
274378 Ok(_) => {}
+2−2
2222 { "binding": "BUS", "queue": "g1t-events" },
2323 // One queue per subscribing service, so each consumes, retries
2424 // and scales on its own. Bindings named SUBSCRIBER_* receive
25− // every event.
25+ // the event types g1t_contracts::subscribers routes to them.
2626 { "binding": "SUBSCRIBER_WORK", "queue": "g1t-events-work" },
2727 { "binding": "SUBSCRIBER_RUNNER", "queue": "g1t-events-runner" },
2828 { "binding": "SUBSCRIBER_INTEGRATIONS", "queue": "g1t-events-integrations" },
3737 { "binding": "SUBSCRIBER_SEARCH", "queue": "g1t-events-search" },
3838 { "binding": "SUBSCRIBER_PACKAGES", "queue": "g1t-events-packages" }
3939 ],
40− "consumers": [{ "queue": "g1t-events", "max_batch_size": 100, "max_batch_timeout": 1 }]
40+ "consumers": [{ "queue": "g1t-events", "max_batch_size": 100, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
4141 },
4242 // Each workspace's audit log keeps as many days as billing says (7 free,
4343 // 90 on the plan, or what staff set), asked for over this binding.
+28−0
10731073 }
10741074
10751075 /// For the queue: pushes on g1t that GitHub follows.
1076+/// Whether `event` should push its repository out to GitHub, given the
1077+/// repositories `pushed` out already for this batch: a sync sends every
1078+/// ref, so a push of many refs, or many pushes in one batch, is one sync.
1079+pub fn first_push_in_batch(pushed: &mut std::collections::HashSet<String>, event: &Event) -> bool {
1080+ match event.repo_id.as_deref() {
1081+ Some(repo_id) if event.kind == "git.push" => pushed.insert(repo_id.to_owned()),
1082+ _ => true,
1083+ }
1084+}
1085+
10761086 pub async fn on_event(env: &Env, event: &Event) -> Result<()> {
10771087 if event.kind != "git.push" || !AppConfig::from_env(env).configured() {
10781088 return Ok(());
10851095 use super::*;
10861096
10871097 #[test]
1098+ fn a_batch_pushes_each_repository_out_once() {
1099+ let event = |kind: &str, repo: &str| Event {
1100+ id: "evt_1".into(),
1101+ kind: kind.into(),
1102+ source: "repos".into(),
1103+ time: "2026-10-08T00:00:00Z".into(),
1104+ repo_id: Some(repo.into()),
1105+ actor: None,
1106+ data: serde_json::json!({}),
1107+ };
1108+ let mut pushed = std::collections::HashSet::new();
1109+ assert!(first_push_in_batch(&mut pushed, &event("git.push", "rep_1")));
1110+ assert!(!first_push_in_batch(&mut pushed, &event("git.push", "rep_1")));
1111+ assert!(first_push_in_batch(&mut pushed, &event("git.push", "rep_2")));
1112+ assert!(first_push_in_batch(&mut pushed, &event("issue.opened", "rep_1")));
1113+ }
1114+
1115+ #[test]
10881116 fn times_are_read_as_github_writes_them() {
10891117 assert_eq!(parse_time("1970-01-01T00:00:00Z"), Some(0));
10901118 assert_eq!(parse_time("2016-07-11T22:14:10Z"), Some(1_468_275_250_000));
+4−1
15801580 #[event(queue)]
15811581 async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
15821582 let service = Integrations::new(&env)?;
1583+ let mut pushed = std::collections::HashSet::new();
15831584 for message in batch.messages()? {
15841585 // A workspace renamed: its rows move to the slug it has now.
15851586 if g1t_kit::rename::on_event(&env, &env.d1("DB")?, message.body(), rename::STATEMENTS).await? {
15981599 }
15991600 // A repository purged: what was kept for it goes.
16001601 rename::on_purged(&env.d1("DB")?, message.body()).await?;
1601− github::on_event(&env, message.body()).await?;
1602+ if github::first_push_in_batch(&mut pushed, message.body()) {
1603+ github::on_event(&env, message.body()).await?;
1604+ }
16021605 service.on_event(message.body()).await?;
16031606 message.ack();
16041607 }
+1−1
2828 ],
2929 // Work starting on, and finishing, what an outside system is waiting on.
3030 "queues": {
31− "consumers": [{ "queue": "g1t-events-integrations", "max_batch_size": 20, "max_batch_timeout": 1 }]
31+ "consumers": [{ "queue": "g1t-events-integrations", "max_batch_size": 20, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
3232 },
3333 "vars": {
3434 "API_URL": "https://api.g1t.sh",
+10−2
4747 use g1t_contracts::repos::{GetArgs, Repo, RepoPath};
4848 use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
4949 use g1t_kit::{args, now_ms, reply, rpc_method};
50−use worker::{Context, Env, Fetcher, MessageBatch, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
50+use worker::{Context, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
5151
5252 use access::{Action, LinkedTo, Target};
5353 use db::{Db, PackageRow};
944944 #[event(queue)]
945945 async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
946946 let packages = Packages::from_env(&env)?;
947+ // Each event is acknowledged or retried on its own, so one that fails
948+ // does not run the rest of its batch again.
947949 for message in batch.messages()? {
948− packages.on_event(&env, message.body()).await?;
950+ match packages.on_event(&env, message.body()).await {
951+ Ok(_) => message.ack(),
952+ Err(error) => {
953+ worker::console_error!("packages: event {} failed: {error}", message.body().id);
954+ message.retry();
955+ }
956+ }
949957 }
950958 Ok(())
951959 }
+1−1
7070 // Workspaces renamed or deleted, repositories that change visibility,
7171 // are renamed, move or are purged.
7272 "queues": {
73− "consumers": [{ "queue": "g1t-events-packages", "max_batch_size": 20, "max_batch_timeout": 1 }]
73+ "consumers": [{ "queue": "g1t-events-packages", "max_batch_size": 20, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
7474 },
7575 "observability": { "enabled": true }
7676 }
+1−1
2626 ],
2727 // A repository created is given its own project.
2828 "queues": {
29− "consumers": [{ "queue": "g1t-events-projects", "max_batch_size": 20, "max_batch_timeout": 1 }]
29+ "consumers": [{ "queue": "g1t-events-projects", "max_batch_size": 20, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
3030 },
3131 "observability": { "enabled": true }
3232 }
+15−0
379379 use super::*;
380380
381381 #[test]
382+ fn every_event_that_settles_a_copy_is_sent_to_this_queue() {
383+ let users = ["user.deleting", "user.restored"];
384+ for kind in g1t_contracts::webhooks::EVENT_TYPES.iter().chain(users.iter()) {
385+ let acts = pull_change(kind).is_some()
386+ || crate::stats::shown_differently(kind, &serde_json::json!({ "username": "ana" })).is_some();
387+ if acts {
388+ assert!(g1t_contracts::subscribers::routed("SUBSCRIBER_REPOS", kind), "{kind} is not routed to repos");
389+ }
390+ }
391+ for kind in ["workspace.renamed", "workspace.deleting", "workspace.restored", "workspace.deleted"] {
392+ assert!(g1t_contracts::subscribers::routed("SUBSCRIBER_REPOS", kind));
393+ }
394+ }
395+
396+ #[test]
382397 fn a_copy_is_removed_days_after_its_pull_request_settles() {
383398 let settled = 1_791_936_000_000; // 2026-10-14T00:00:00Z
384399 assert_eq!(retire_after(settled, 7), "2026-10-21T00:00:00.000Z");
+149−71
6161
6262 use serde::Serialize;
6363 use worker::{
64− Context, Env, Fetcher, MessageBatch, Method, Request, Response, Result, ScheduleContext, ScheduledEvent,
64+ Context, Env, Fetcher, MessageBatch, MessageExt, Method, Request, Response, Result, ScheduleContext, ScheduledEvent,
6565 event,
6666 };
6767
21882188 // not fit a repo per pull request, so the front end reports pushes
21892189 // itself: one event for each branch that moved.
21902190 let stored = self.store.open(&store_key(&repo)).await?;
2191− for pushed in &pushed {
2191+ let announced = announced_refs(&pushed, &repo.default_branch);
2192+ if announced.len() < pushed.len() {
2193+ worker::console_log!(
2194+ "push to {}: {} of {} refs announced",
2195+ repo.name,
2196+ announced.len(),
2197+ pushed.len()
2198+ );
2199+ }
2200+ for pushed in announced {
21922201 // The store can refuse one ref and accept another, so each
21932202 // branch is checked against where it actually is. A tag the
21942203 // store cannot read back is taken as pushed.
22192228 }
22202229 }
22212230
2231+/// The most tags one push announces: a push of more (`git push --tags`
2232+/// into a new repository, say) announces none of them, so it starts no
2233+/// workflows, mirror syncs or package reads, one per tag.
2234+const MAX_PUSH_TAG_EVENTS: usize = 3;
2235+/// The most branches one push announces. A push of more announces only
2236+/// the default branch, if it moved, which the rest of g1t reads from.
2237+const MAX_PUSH_BRANCH_EVENTS: usize = 1000;
2238+
2239+/// The refs of a push that are announced with a `git.push` event each.
2240+/// Every ref is stored whatever this says; only the events are capped.
2241+fn announced_refs<'a>(pushed: &'a [git_http::Pushed], default_branch: &str) -> Vec<&'a git_http::Pushed> {
2242+ let (branches, tags): (Vec<&git_http::Pushed>, Vec<&git_http::Pushed>) =
2243+ pushed.iter().partition(|pushed| pushed.branch().is_some());
2244+ let mut announced = if branches.len() > MAX_PUSH_BRANCH_EVENTS {
2245+ branches.into_iter().filter(|pushed| pushed.branch() == Some(default_branch)).collect()
2246+ } else {
2247+ branches
2248+ };
2249+ if tags.len() <= MAX_PUSH_TAG_EVENTS {
2250+ announced.extend(tags);
2251+ }
2252+ announced
2253+}
2254+
22222255 /// A push the store accepted, to be recorded once git has its answer.
22232256 struct PushDone {
22242257 repo: Repo,
26682701 handled
26692702 }
26702703
2704+/// Each event is acknowledged or retried on its own, so one that fails is
2705+/// tried again without the others before and after it running twice.
26712706 async fn handle_events(batch: &MessageBatch<Event>, env: &Env, registry: &Registry, identity: &Fetcher) -> Result<()> {
26722707 for message in batch.messages()? {
2673− let event = message.body();
2674− // A pull request merged, closed or reopened: its working copy is
2675− // kept or let go (forks.rs).
2676− if let Some(change) = forks::pull_change(&event.kind) {
2677− let Some(pull_id) = forks::pull_id_of(&event.data) else {
2678− worker::console_error!("{} {} names no pull request", event.kind, event.id);
2679− continue;
2680− };
2681− let repos = service(env)?;
2682− match change {
2683− forks::PullChange::Settled => repos.pull_settled(&pull_id).await?,
2684− forks::PullChange::Reopened => repos.pull_reopened(&pull_id).await?,
2708+ match handle_event(message.body(), env, registry, identity).await {
2709+ Ok(()) => message.ack(),
2710+ Err(error) => {
2711+ worker::console_error!("repos: event {} failed: {error}", message.body().id);
2712+ message.retry();
26852713 }
2686− continue;
26872714 }
2688− // A workspace deleted, restored or purged: its repositories go with
2689− // it, come back with it, or are purged with it (lifecycle.rs).
2690− if event.kind == "workspace.deleting" {
2691− match serde_json::from_value::<WorkspaceDeleting>(event.data.clone()) {
2692− Ok(deleting) => service(env)?.delete_with_workspace(&deleting, &protected_workspaces(env)).await?,
2693− Err(_) => worker::console_error!("workspace.deleting {} could not be read", event.id),
2694− }
2695− continue;
2715+ }
2716+ Ok(())
2717+}
2718+
2719+async fn handle_event(event: &Event, env: &Env, registry: &Registry, identity: &Fetcher) -> Result<()> {
2720+ // A pull request merged, closed or reopened: its working copy is
2721+ // kept or let go (forks.rs).
2722+ if let Some(change) = forks::pull_change(&event.kind) {
2723+ let Some(pull_id) = forks::pull_id_of(&event.data) else {
2724+ worker::console_error!("{} {} names no pull request", event.kind, event.id);
2725+ return Ok(());
2726+ };
2727+ let repos = service(env)?;
2728+ match change {
2729+ forks::PullChange::Settled => repos.pull_settled(&pull_id).await?,
2730+ forks::PullChange::Reopened => repos.pull_reopened(&pull_id).await?,
26962731 }
2697− if event.kind == "workspace.restored" {
2698− match serde_json::from_value::<WorkspaceRestored>(event.data.clone()) {
2699− Ok(restored) => service(env)?.restore_with_workspace(&restored).await?,
2700− Err(_) => worker::console_error!("workspace.restored {} could not be read", event.id),
2701− }
2702− continue;
2732+ return Ok(());
2733+ }
2734+ // A workspace deleted, restored or purged: its repositories go with
2735+ // it, come back with it, or are purged with it (lifecycle.rs).
2736+ if event.kind == "workspace.deleting" {
2737+ match serde_json::from_value::<WorkspaceDeleting>(event.data.clone()) {
2738+ Ok(deleting) => service(env)?.delete_with_workspace(&deleting, &protected_workspaces(env)).await?,
2739+ Err(_) => worker::console_error!("workspace.deleting {} could not be read", event.id),
27032740 }
2704− if event.kind == "workspace.deleted" {
2705− match serde_json::from_value::<WorkspaceDeleted>(event.data.clone()) {
2706− Ok(deleted) => service(env)?.purge_workspace(&deleted, &protected_workspaces(env)).await?,
2707− Err(_) => worker::console_error!("workspace.deleted {} could not be read", event.id),
2708− }
2709− continue;
2741+ return Ok(());
2742+ }
2743+ if event.kind == "workspace.restored" {
2744+ match serde_json::from_value::<WorkspaceRestored>(event.data.clone()) {
2745+ Ok(restored) => service(env)?.restore_with_workspace(&restored).await?,
2746+ Err(_) => worker::console_error!("workspace.restored {} could not be read", event.id),
27102747 }
2711− // An account deleted or restored: kept contributors that name it
2712− // (or ghost, for a restore) are worked out again, so it shows as
2713− // ghost for its 30 days, and as itself again if restored.
2714− if let Some(needle) = stats::shown_differently(&event.kind, &event.data) {
2715− let db = env.d1("DB")?;
2716− if let Err(error) = stats::rework_naming(&db, &needle).await {
2717− worker::console_error!("{} {}: contributors not marked to be counted again: {error}", event.kind, event.id);
2718− }
2719− continue;
2748+ return Ok(());
2749+ }
2750+ if event.kind == "workspace.deleted" {
2751+ match serde_json::from_value::<WorkspaceDeleted>(event.data.clone()) {
2752+ Ok(deleted) => service(env)?.purge_workspace(&deleted, &protected_workspaces(env)).await?,
2753+ Err(_) => worker::console_error!("workspace.deleted {} could not be read", event.id),
27202754 }
2721− if event.kind != "workspace.renamed" {
2722− continue;
2755+ return Ok(());
2756+ }
2757+ // An account deleted or restored: kept contributors that name it
2758+ // (or ghost, for a restore) are worked out again, so it shows as
2759+ // ghost for its 30 days, and as itself again if restored.
2760+ if let Some(needle) = stats::shown_differently(&event.kind, &event.data) {
2761+ let db = env.d1("DB")?;
2762+ if let Err(error) = stats::rework_naming(&db, &needle).await {
2763+ worker::console_error!("{} {}: contributors not marked to be counted again: {error}", event.kind, event.id);
27232764 }
2724− let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
2725− worker::console_error!("workspace.renamed {} could not be read", event.id);
2726− continue;
2727− };
2728− let names: HashMap<String, String> = g1t_kit::call(
2729− identity,
2730− "usernames",
2731− &g1t_contracts::identity::UsernamesArgs {
2732− ids: vec![renamed.workspace_id.clone()],
2733− },
2734− )
2765+ return Ok(());
2766+ }
2767+ if event.kind != "workspace.renamed" {
2768+ return Ok(());
2769+ }
2770+ let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
2771+ worker::console_error!("workspace.renamed {} could not be read", event.id);
2772+ return Ok(());
2773+ };
2774+ let names: HashMap<String, String> = g1t_kit::call(
2775+ identity,
2776+ "usernames",
2777+ &g1t_contracts::identity::UsernamesArgs {
2778+ ids: vec![renamed.workspace_id.clone()],
2779+ },
2780+ )
2781+ .await?;
2782+ let current = names
2783+ .get(&renamed.workspace_id)
2784+ .cloned()
2785+ .unwrap_or_else(|| renamed.to.clone());
2786+ let left = registry
2787+ .rename_namespace(&renamed.stale_slugs(&current), &current)
27352788 .await?;
2736− let current = names
2737− .get(&renamed.workspace_id)
2738− .cloned()
2739− .unwrap_or_else(|| renamed.to.clone());
2740− let left = registry
2741− .rename_namespace(&renamed.stale_slugs(&current), &current)
2742− .await?;
2743− if left > 0 {
2744− worker::console_error!(
2745− "{left} repositories stayed under {} or {}: {current} already has repositories of the same names",
2746− renamed.from,
2747− renamed.to
2748− );
2749− }
2789+ if left > 0 {
2790+ worker::console_error!(
2791+ "{left} repositories stayed under {} or {}: {current} already has repositories of the same names",
2792+ renamed.from,
2793+ renamed.to
2794+ );
27502795 }
27512796 Ok(())
27522797 }
28282873 }
28292874
28302875 #[cfg(test)]
2876+mod announced_refs_tests {
2877+ use super::*;
2878+
2879+ fn pushed(git_ref: String) -> git_http::Pushed {
2880+ git_http::Pushed { git_ref, before: None, after: "abc".into() }
2881+ }
2882+
2883+ fn refs(announced: Vec<&git_http::Pushed>) -> Vec<&str> {
2884+ announced.into_iter().map(|pushed| pushed.git_ref.as_str()).collect()
2885+ }
2886+
2887+ #[test]
2888+ fn a_few_tags_are_announced_and_many_are_not() {
2889+ let few: Vec<_> = (1..=3).map(|n| pushed(format!("refs/tags/v{n}"))).chain([pushed("refs/heads/main".into())]).collect();
2890+ assert_eq!(refs(announced_refs(&few, "main")), ["refs/heads/main", "refs/tags/v1", "refs/tags/v2", "refs/tags/v3"]);
2891+ let many: Vec<_> = (1..=10_000).map(|n| pushed(format!("refs/tags/v{n}"))).chain([pushed("refs/heads/main".into())]).collect();
2892+ assert_eq!(refs(announced_refs(&many, "main")), ["refs/heads/main"]);
2893+ }
2894+
2895+ #[test]
2896+ fn past_the_branch_cap_only_the_default_branch_is_announced() {
2897+ let at_cap: Vec<_> = (0..MAX_PUSH_BRANCH_EVENTS).map(|n| pushed(format!("refs/heads/b{n}"))).collect();
2898+ assert_eq!(announced_refs(&at_cap, "main").len(), MAX_PUSH_BRANCH_EVENTS);
2899+ let over: Vec<_> = (0..=MAX_PUSH_BRANCH_EVENTS)
2900+ .map(|n| pushed(format!("refs/heads/b{n}")))
2901+ .chain([pushed("refs/heads/main".into())])
2902+ .collect();
2903+ assert_eq!(refs(announced_refs(&over, "main")), ["refs/heads/main"]);
2904+ assert!(announced_refs(&over, "trunk").is_empty());
2905+ }
2906+}
2907+
2908+#[cfg(test)]
28312909 mod push_to_create_tests {
28322910 use super::*;
28332911
+1−1
122122 // A workspace renamed moves its repositories to the new slug; one
123123 // deleted purges the repositories it left in Recently deleted.
124124 "queues": {
125− "consumers": [{ "queue": "g1t-events-repos", "max_batch_size": 20, "max_batch_timeout": 1 }]
125+ "consumers": [{ "queue": "g1t-events-repos", "max_batch_size": 20, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
126126 },
127127 "observability": { "enabled": true }
128128 }
+6−1
9191 import { BACKUP_MINUTES, backupEnv, backupPace, backupSandboxName } from "./backup";
9292 import { type ProjectSurroundings, readableSurroundings } from "./surroundings";
9393 import { holdCredentials, pushGrant, remotePath, revokeCredentials, runCredential } from "./credentials";
94−import { buildMentionPrompt, describeThread, handleMention, planMention } from "./mentions";
94+import { buildMentionPrompt, describeThread, handleMention, jobTokenRefusal, planMention } from "./mentions";
9595 import { instructionsFor, repoInstructions, withBlock } from "./repo-instructions";
9696 import { cancelTask, enqueueTask, handedOverStep, selfHostedRoute, taskEnv, taskRepo } from "./self-hosted";
9797 import {
24382438 * the actor is a member or a collaborator.
24392439 */
24402440 private async refusal(actor: User, repo: RepoPath): Promise<Result<never> | null> {
2441+ // Agent compute is never started by a workflow job's token.
2442+ const byJob = jobTokenRefusal(actor);
2443+ if (byJob) return fail("forbidden", byJob);
24412444 const closed = await this.closedRepo(actor, repo);
24422445 if (closed) return closed;
24432446 if (!(await this.workspaceAllowed(repo.namespace))) {
28402843
28412844 async delegate(actor: User, repo: RepoPath, input: DelegateInput): Promise<Result<Delegated>> {
28422845 // Who may put agents to work here is settled before anything is opened.
2846+ const byJob = jobTokenRefusal(actor);
2847+ if (byJob) return fail("forbidden", byJob);
28432848 const closed = await this.closedRepo(actor, repo);
28442849 if (closed) return closed;
28452850 if (!actor || !(await this.repoAllows(actor, repo, "run"))) return fail("forbidden", needs("run"));
+9−1
33
44 import type { LifecycleJob, MentionJob, Result } from "@g1t/contracts";
55
6−import { type MentionPorts, buildMentionPrompt, handleMention, planMention } from "./mentions.ts";
6+import { type MentionPorts, buildMentionPrompt, handleMention, jobTokenRefusal, planMention } from "./mentions.ts";
77
88 const ana = { id: "usr_ana", username: "ana", verified: true, workspaces: [{ slug: "acme", role: "member" as const }] };
99
162162 };
163163 return { ports, replies, started, recorded };
164164 }
165+
166+test("a workflow job's token never puts g1t to work; a person's own token can", () => {
167+ const job = { ...ana, token: { scopes: ["repo"], job: { run_id: "run_9", job_id: "job_1" } } };
168+ assert.match(jobTokenRefusal(job) ?? "", /job's token/);
169+ assert.equal(jobTokenRefusal({ ...ana, token: { scopes: ["repo"] } }), null);
170+ assert.equal(jobTokenRefusal(ana), null);
171+ assert.equal(jobTokenRefusal(null), null);
172+});
+11−0
1111 */
1212 import type { Comment, LifecycleJob, MentionJob, MentionsApi, Pull, RepoPath, Result, User } from "@g1t/contracts";
1313
14+/**
15+ * Why a workflow job's token (`G1T_TOKEN`) may not put g1t to work, or
16+ * null for anyone else. A workflow that could summon an agent on a failing
17+ * check would start one whose push runs the workflow again, without end.
18+ */
19+export function jobTokenRefusal(actor: User | null | undefined): string | null {
20+ return actor?.token?.job
21+ ? "A workflow job's token (G1T_TOKEN) cannot put g1t to work. A person, or a token of their own, can."
22+ : null;
23+}
24+
1425 /** What a mention leads to. */
1526 export type MentionPlan =
1627 | { kind: "not_member" }
+1−1
6767 "triggers": { "crons": ["*/5 * * * *"] },
6868 // Events it reacts to: a pull request ready for review, or its head moving.
6969 "queues": {
70− "consumers": [{ "queue": "g1t-events-runner", "max_batch_size": 20, "max_batch_timeout": 1 }]
70+ "consumers": [{ "queue": "g1t-events-runner", "max_batch_size": 20, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
7171 },
7272 "vars": {
7373 // Each workspace chooses where its agents' model spend goes: its own
+3−3
2929 // of the backfill. One at a time, so each stays within its caps.
3030 "producers": [{ "binding": "JOBS", "queue": "g1t-search-jobs" }],
3131 "consumers": [
32− // Every event, to keep the index current.
33− { "queue": "g1t-events-search", "max_batch_size": 20, "max_batch_timeout": 2 },
34− { "queue": "g1t-search-jobs", "max_batch_size": 1, "max_batch_timeout": 1, "max_retries": 3 }
32+ // The events that change what is indexed, to keep it current.
33+ { "queue": "g1t-events-search", "max_batch_size": 20, "max_batch_timeout": 2, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" },
34+ { "queue": "g1t-search-jobs", "max_batch_size": 1, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }
3535 ]
3636 },
3737 "observability": { "enabled": true }
+1−1
3636 // failing checks (security updates move on), new repositories (their
3737 // history is scanned once) and renamed workspaces.
3838 "queues": {
39− "consumers": [{ "queue": "g1t-events-security", "max_batch_size": 10, "max_batch_timeout": 5 }]
39+ "consumers": [{ "queue": "g1t-events-security", "max_batch_size": 10, "max_batch_timeout": 5, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
4040 },
4141 // Continues history scans a page at a time, and reads every
4242 // repository's dependencies again once a day (every 30 minutes); runs
+1−1
2323 ],
2424 // Every event on the bus, to deliver to the webhooks that want it.
2525 "queues": {
26− "consumers": [{ "queue": "g1t-events-webhooks", "max_batch_size": 20, "max_batch_timeout": 1 }]
26+ "consumers": [{ "queue": "g1t-events-webhooks", "max_batch_size": 20, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
2727 },
2828 // Retries that are due, and deliveries old enough to forget.
2929 "triggers": { "crons": ["* * * * *"] },
+5−0
1+-- When a mention sent g1t back to its pull request, so that no more than
2+-- a day's ceiling of them do (settings.rs MAX_MENTION_REVISIONS_PER_DAY).
3+ALTER TABLE agent_mentions ADD COLUMN revised_at TEXT;
4+CREATE INDEX agent_mentions_revised ON agent_mentions (pull_id, revised_at)
5+ WHERE revised_at IS NOT NULL;
+3−0
490490 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
491491 check!(writable(&repo));
492492 check!(allowed(Some(&a.actor), &repo, Capability::Run));
493+ if let Some(refused) = mentions::refuse_job_token(&a.actor) {
494+ return Ok(refused);
495+ }
493496 self.open_issue(OpenIssueArgs {
494497 actor: a.actor,
495498 repo: a.repo,
+80−9
2323 use worker::wasm_bindgen::JsValue;
2424
2525 use crate::Work;
26+use crate::settings::MAX_MENTION_REVISIONS_PER_DAY;
2627 use crate::lifecycle::{POLICY_ACTOR_ID, POLICY_ACTOR_NAME, made_by_g1t};
2728 use crate::reviews::{AGENT_ID, AGENT_NAME};
2829 use crate::rows::{NumberRow, ValueRow};
388389 Ok(Some(label.to_lowercase()))
389390 }
390391
392+/// Whether `actor` mentioning `@g1t` may set it to work. g1t and other
393+/// agents may not, so no agent can set another to work; neither may a
394+/// workflow job's token (`G1T_TOKEN`), or a workflow that comments on a
395+/// failing check would start an agent whose push runs it again, and so on
396+/// without end.
397+pub(crate) fn may_summon(actor: &User) -> bool {
398+ actor.kind != PrincipalKind::Agent
399+ && !actor.is_system()
400+ && actor.id != AGENT_ID
401+ && g1t_contracts::events::job_run_of(actor).is_none()
402+}
403+
404+/// Why a workflow job's token may not put g1t to work (`may_summon`).
405+pub(crate) const JOB_TOKEN_REFUSED: &str =
406+ "A workflow job's token (G1T_TOKEN) cannot put g1t to work. A person, or a token of their own, can.";
407+
408+/// A refusal for an actor acting with a workflow job's token, which may
409+/// not start agents by any path: queueing an issue, assigning a plan's or
410+/// handing one over.
411+pub(crate) fn refuse_job_token<T>(actor: &User) -> Option<Outcome<T>> {
412+ g1t_contracts::events::job_run_of(actor).map(|_| Outcome::fail(FailureCode::Forbidden, JOB_TOKEN_REFUSED))
413+}
414+
415+/// The start of the day `mention_revision` counts back over, from `now`.
416+fn day_before(now: u64) -> String {
417+ rfc3339(now.saturating_sub(24 * 60 * 60 * 1000))
418+}
419+
391420 impl Work {
392421 /// Records a comment's mention of `@g1t`, if it makes one, for
393− /// the runner to take when it hears of the comment. g1t mentioning
394− /// itself is not recorded, so no agent can set another to work.
422+ /// the runner to take when it hears of the comment. What agents and
423+ /// workflow jobs say is not recorded (`may_summon`).
395424 pub(crate) async fn note_mention(
396425 &self,
397426 actor: &User,
400429 comment: &Comment,
401430 pull_id: Option<&str>,
402431 ) -> Result<()> {
403− if actor.kind == PrincipalKind::Agent
404− || actor.is_system()
405− || actor.id == AGENT_ID
406− || !mentions_agent(&comment.body)
407− {
432+ if !may_summon(actor) || !mentions_agent(&comment.body) {
408433 return Ok(());
409434 }
410435 // Whether it may set the agent to work: mentioning spends compute.
571596 return Ok(Outcome::fail(FailureCode::NotFound, "Pull request not found."));
572597 };
573598 let base = pull.base_branch(&repo.default_branch).to_owned();
574− // A person asking outranks a stop and the limit on revisions.
599+ // A person asking outranks a stop and the limit on revisions, but
600+ // not this ceiling: each revision spends compute, and nothing
601+ // should send g1t back to one pull request without end.
602+ let now = now_ms();
603+ let today = self
604+ .db
605+ .prepare("SELECT count(*) AS n FROM agent_mentions WHERE pull_id = ? AND revised_at >= ?")
606+ .bind(&[pull.id.as_str().into(), day_before(now).into()])?
607+ .first::<NumberRow>(None)
608+ .await?
609+ .map_or(0, |row| row.n);
610+ if today >= MAX_MENTION_REVISIONS_PER_DAY {
611+ return Ok(Outcome::fail(
612+ FailureCode::Limit,
613+ format!(
614+ "g1t has been sent back to this pull request {MAX_MENTION_REVISIONS_PER_DAY} times in the last day, the most it takes. Mention it again tomorrow, or push the change yourself."
615+ ),
616+ ));
617+ }
575618 let was_stalled = self.is_stalled(&pull.id).await?;
576619 self.db
577620 .prepare("UPDATE pulls SET stalled = NULL WHERE id = ?")
587630 "g1t is already taking a step on this pull request.",
588631 ));
589632 }
633+ self.db
634+ .prepare("UPDATE agent_mentions SET revised_at = ? WHERE comment_id = ?")
635+ .bind(&[rfc3339(now).into(), row.comment_id.as_str().into()])?
636+ .run()
637+ .await?;
590638 let round = self
591639 .db
592640 .prepare("SELECT revisions FROM pulls WHERE id = ?")
787835 // the repository is public does not matter here.
788836 let target = access::RepoRef { id: &issue.repo_id, namespace: &path.namespace, private: true };
789837 if !actor.verified
790− || actor.kind == PrincipalKind::Agent
838+ || !may_summon(actor)
791839 || !access::can(Some(actor), target, Capability::Run)
792840 {
793841 return Ok(());
819867 use super::*;
820868
821869 #[test]
870+ fn only_people_and_their_own_tokens_may_summon_g1t() {
871+ let mut ana = User { id: "usr_ana".into(), username: "ana".into(), verified: true, ..User::default() };
872+ assert!(may_summon(&ana));
873+ assert!(refuse_job_token::<()>(&ana).is_none());
874+ ana.token = Some(Box::new(g1t_contracts::scopes::TokenAccess::default()));
875+ assert!(may_summon(&ana), "a person's own token speaks for them");
876+ ana.token = Some(Box::new(g1t_contracts::scopes::TokenAccess {
877+ job: Some(g1t_contracts::scopes::JobToken { run_id: "run_9".into(), job_id: "job_1".into(), pull_requests: true }),
878+ ..Default::default()
879+ }));
880+ assert!(!may_summon(&ana), "a workflow's comment would set off the workflow again");
881+ assert!(matches!(refuse_job_token::<()>(&ana), Some(Outcome::Fail(_))));
882+ let g1t = User { id: AGENT_ID.into(), username: AGENT_NAME.into(), ..User::default() };
883+ assert!(!may_summon(&g1t));
884+ }
885+
886+ #[test]
887+ fn a_day_of_mention_revisions_counts_back_from_now() {
888+ assert_eq!(day_before(2 * 24 * 60 * 60 * 1000), rfc3339(24 * 60 * 60 * 1000));
889+ assert_eq!(day_before(5), rfc3339(0));
890+ }
891+
892+ #[test]
822893 fn a_mention_is_found_whatever_its_case() {
823894 assert!(mentions_agent("@g1t take this"));
824895 assert!(mentions_agent("@G1T take this"));
+8−0
379379 if let Some(refused) = may_plan(&a.actor, &repo) {
380380 return Ok(refused);
381381 }
382+ if a.assign
383+ && let Some(refused) = crate::mentions::refuse_job_token(&a.actor)
384+ {
385+ return Ok(refused);
386+ }
382387 let Some(row) = self
383388 .plan_row(&a.id)
384389 .await?
520525 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
521526 };
522527 if a.queued {
528+ if let Some(refused) = crate::mentions::refuse_job_token(&a.actor) {
529+ return Ok(refused);
530+ }
523531 let repo = match self.repo(&a.repo, &Some(a.actor.clone())).await? {
524532 Outcome::Ok(repo) => repo,
525533 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
+3−0
1717
1818 const MAX_REQUIRED_APPROVALS: u32 = 6;
1919 const MAX_REVISIONS: u32 = 5;
20+/// The most times in a day people's mentions may send g1t back to one pull
21+/// request. They outrank `MAX_REVISIONS` and a stall, but not this.
22+pub(crate) const MAX_MENTION_REVISIONS_PER_DAY: u32 = 10;
2023
2124 #[derive(Deserialize)]
2225 struct SettingsRow {
+1−1
2424 { "binding": "ACTIONS", "service": "g1t-actions" }
2525 ],
2626 "queues": {
27− "consumers": [{ "queue": "g1t-events-work", "max_batch_size": 100, "max_batch_timeout": 1 }]
27+ "consumers": [{ "queue": "g1t-events-work", "max_batch_size": 100, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }]
2828 },
2929 "observability": { "enabled": true }
3030 }