Skip to content
339 linesCodeBlameRaw
1//! Repositories that are archived, deleted or purged, and branches
2//! renamed: what issues, pull requests and agents do about each.
3//!
4//! An archived repository is read-only: its issues and pull requests are
5//! locked and nothing new starts on it. So is a mirror standing by (see
6//! `g1t_contracts::mirrors`): its work happens where it is mirrored from,
7//! until someone takes over on g1t. A deleted one looks missing (repos
8//! hides it) and keeps its rows for a restore. A purged one is gone, and
9//! every row kept for it goes with it.
10
11use g1t_contracts::events::{BranchRenamed, Event, RepoArchived, RepoDeleted};
12use g1t_contracts::repos::{Repo, RepoStatus, StatusByIdArgs, archived_message};
13use g1t_contracts::{FailureCode, Outcome};
14use g1t_contracts::time::rfc3339;
15use g1t_kit::now_ms;
16use serde::{Deserialize, Serialize};
17use worker::Result;
18
19use crate::Work;
20
21/// `Ok` when `repo` may be changed; refused, saying why, when it is
22/// archived or a mirror that is not taken over.
23pub(crate) fn writable(repo: &Repo) -> Outcome<()> {
24 match repo.read_only_reason() {
25 Some(reason) => Outcome::fail(FailureCode::Forbidden, reason),
26 None => Outcome::Ok(()),
27 }
28}
29
30/// `Ok` unless `repo` is archived. For what stays g1t's own on a mirror:
31/// its settings and rules, and the statuses and checks reported on its
32/// commits (CI failover reports them).
33pub(crate) fn not_archived(repo: &Repo) -> Outcome<()> {
34 if repo.archived() {
35 Outcome::fail(FailureCode::Forbidden, archived_message(&repo.namespace, &repo.name))
36 } else {
37 Outcome::Ok(())
38 }
39}
40
41/// `repo` as found, unless it is archived: then refused, as [`writable`]
42/// says. For the steps g1t starts by itself, which stop on a refusal.
43pub(crate) fn unless_archived(repo: Outcome<Repo>) -> Outcome<Repo> {
44 match repo {
45 Outcome::Ok(repo) => match writable(&repo) {
46 Outcome::Ok(()) => Outcome::Ok(repo),
47 Outcome::Fail(failure) => Outcome::Fail(failure),
48 },
49 failed => failed,
50 }
51}
52
53/// `repo.purged`: every row kept for the repository, `?1` its id. Rows
54/// kept by pull request go first, while the pull requests still name them.
55pub(crate) const PURGED: &[&str] = &[
56 "DELETE FROM session_entries WHERE pull_id IN (SELECT id FROM pulls WHERE repo_id = ?1)",
57 "DELETE FROM check_runs WHERE pull_id IN (SELECT id FROM pulls WHERE repo_id = ?1)",
58 "DELETE FROM review_runs WHERE pull_id IN (SELECT id FROM pulls WHERE repo_id = ?1)",
59 "DELETE FROM agent_messages WHERE repo_id = ?1 OR pull_id IN (SELECT id FROM pulls WHERE repo_id = ?1)",
60 "DELETE FROM queue_entries WHERE repo_id = ?1",
61 "DELETE FROM agent_mentions WHERE repo_id = ?1",
62 "DELETE FROM agent_rules WHERE repo_id = ?1",
63 "DELETE FROM pull_confidence WHERE repo_id = ?1",
64 "DELETE FROM run_confidence WHERE repo_id = ?1",
65 "DELETE FROM confidence_rules WHERE repo_id = ?1",
66 "DELETE FROM agent_runs WHERE repo_id = ?1",
67 "DELETE FROM commit_statuses WHERE repo_id = ?1",
68 "DELETE FROM commit_check_annotations WHERE run_id IN (SELECT id FROM commit_check_runs WHERE repo_id = ?1)",
69 "DELETE FROM commit_check_runs WHERE repo_id = ?1",
70 "DELETE FROM commit_check_suites WHERE repo_id = ?1",
71 "DELETE FROM plans WHERE repo_id = ?1",
72 "DELETE FROM comments WHERE repo_id = ?1",
73 "DELETE FROM pulls WHERE repo_id = ?1",
74 "DELETE FROM issues WHERE repo_id = ?1",
75 "DELETE FROM labels WHERE repo_id = ?1",
76 "DELETE FROM milestones WHERE repo_id = ?1",
77 "DELETE FROM counters WHERE repo_id = ?1",
78 "DELETE FROM repo_settings WHERE repo_id = ?1",
79 "DELETE FROM memories WHERE scope = 'project' AND scope_key = ?1",
80 "DELETE FROM guardrails WHERE scope <> 'workspace' AND scope_key = ?1",
81];
82
83/// A repository renamed or transferred (see `g1t_kit::transfer`): runs
84/// waiting for an agent slot name it by path in their payload, at `$.repo`
85/// or `$.job.repo`; they follow it, into its workspace's queue now.
86pub(crate) const WAITS_MOVED: &[&str] = &[
87 "UPDATE agent_waits SET workspace = lower(?3),
88 payload = json_set(payload, '$.repo.namespace', ?3, '$.repo.name', ?6)
89 WHERE json_valid(payload)
90 AND lower(json_extract(payload, '$.repo.namespace')) = lower(?4)
91 AND lower(json_extract(payload, '$.repo.name')) = lower(?7)",
92 "UPDATE agent_waits SET workspace = lower(?3),
93 payload = json_set(payload, '$.job.repo.namespace', ?3, '$.job.repo.name', ?6)
94 WHERE json_valid(payload)
95 AND lower(json_extract(payload, '$.job.repo.namespace')) = lower(?4)
96 AND lower(json_extract(payload, '$.job.repo.name')) = lower(?7)",
97];
98
99/// What a run stopped because its repository was archived or deleted says
100/// as its last step; [`Work::runs_in_repo`] finds such runs by it.
101const STOPPED_STEP: &str = "Stopped: the repository was archived or deleted.";
102
103/// `repo.deleted` and `repo.archived`: the repository's agent runs that
104/// are queued or running stop, `?1` the step they end on, `?2` the time,
105/// `?3` its id.
106const STOP_RUNS: &str = "UPDATE agent_runs SET status = 'stopped', step = ?1, finished_at = ?2, updated_at = ?2
107 WHERE repo_id = ?3 AND status IN ('queued', 'running')";
108
109/// And runs waiting for an agent slot in it are dropped: `?1` its
110/// workspace, `?2` its name, as the payload names it (see [`WAITS_MOVED`]).
111const DROP_WAITS: &[&str] = &[
112 "DELETE FROM agent_waits WHERE json_valid(payload)
113 AND lower(json_extract(payload, '$.repo.namespace')) = lower(?1)
114 AND lower(json_extract(payload, '$.repo.name')) = lower(?2)",
115 "DELETE FROM agent_waits WHERE json_valid(payload)
116 AND lower(json_extract(payload, '$.job.repo.namespace')) = lower(?1)
117 AND lower(json_extract(payload, '$.job.repo.name')) = lower(?2)",
118];
119
120/// `runs_in_repo` (internal, for the runner): the agent runs of a
121/// repository whose sandboxes should be stopped: those queued or running,
122/// and those work stopped in the last hour because the repository was
123/// archived or deleted. Takes `{ "repoId" }`; returns `[RepoRun]`.
124#[derive(Deserialize)]
125#[serde(rename_all = "camelCase")]
126pub(crate) struct RunsInRepoArgs {
127 pub repo_id: String,
128}
129
130#[derive(Debug, Deserialize, Serialize)]
131#[serde(rename_all = "camelCase")]
132pub(crate) struct RepoRun {
133 pub run_id: String,
134 pub sandbox: String,
135 pub kind: String,
136 pub status: String,
137 pub pull_id: Option<String>,
138}
139
140/// The repository, by id and path, that `event` says was deleted or
141/// archived (not unarchived): its runs stop.
142fn stopping(event: &Event) -> Option<(String, String, String)> {
143 let data = event.data.clone();
144 match event.kind.as_str() {
145 "repo.deleted" => serde_json::from_value::<RepoDeleted>(data).ok().map(|e| (e.repo_id, e.namespace, e.name)),
146 "repo.archived" => serde_json::from_value::<RepoArchived>(data)
147 .ok()
148 .filter(|e| e.archived)
149 .map(|e| (e.repo_id, e.namespace, e.name)),
150 _ => None,
151 }
152}
153
154/// `branch.renamed`: open pull requests from `?2` in repository `?1`
155/// come from `?3` now. A fork carries its change on its own default
156/// branch, so only branches of the repository itself are named.
157const BRANCH_RENAMED: &str = "UPDATE pulls SET source_branch = ?3
158 WHERE repo_id = ?1 AND source_branch = ?2 AND fork_repo_id IS NULL
159 AND status IN ('draft', 'open')";
160
161/// `branch.renamed`: open pull requests into `?2` merge into `?3` now.
162/// Those into the default branch name none, and follow it as they are.
163const BASE_RENAMED: &str = "UPDATE pulls SET base_branch = ?3
164 WHERE repo_id = ?1 AND base_branch = ?2 AND status IN ('draft', 'open')";
165
166impl Work {
167 /// Whether work may start on the repository: neither archived nor
168 /// deleted. When repos cannot say, it is taken as active, and the
169 /// repos calls that follow (which hide deleted repositories) decide.
170 pub(crate) async fn repo_active(&self, repo_id: &str) -> Result<bool> {
171 let status: Result<RepoStatus> = g1t_kit::call(
172 &self.repos,
173 "status_by_id",
174 &StatusByIdArgs {
175 id: repo_id.to_owned(),
176 },
177 )
178 .await;
179 Ok(match status {
180 Ok(status) => status.active(),
181 Err(error) => {
182 worker::console_error!("status_by_id {repo_id}: {error}");
183 true
184 }
185 })
186 }
187
188 /// See [`RunsInRepoArgs`].
189 pub(crate) async fn runs_in_repo(&self, a: RunsInRepoArgs) -> Result<Vec<RepoRun>> {
190 let hour_ago = rfc3339(now_ms().saturating_sub(60 * 60 * 1000));
191 self.db
192 .prepare(
193 "SELECT id AS run_id, sandbox, kind, status, pull_id FROM agent_runs
194 WHERE repo_id = ?1 AND (status IN ('queued', 'running')
195 OR (status = 'stopped' AND step = ?2 AND finished_at >= ?3))
196 ORDER BY created_at LIMIT 200",
197 )
198 .bind(&[a.repo_id.as_str().into(), STOPPED_STEP.into(), hour_ago.as_str().into()])?
199 .all()
200 .await?
201 .results::<RepoRun>()
202 }
203
204 /// `repo.deleted` and `repo.archived`: its agent runs stop (the runner
205 /// stops their sandboxes, see [`RunsInRepoArgs`]) and runs waiting for
206 /// a slot in it are dropped. Rows are kept, for a restore.
207 async fn stop_runs_in(&self, repo_id: &str, namespace: &str, name: &str) -> Result<()> {
208 let now = rfc3339(now_ms());
209 let mut batch = vec![self
210 .db
211 .prepare(STOP_RUNS)
212 .bind(&[STOPPED_STEP.into(), now.as_str().into(), repo_id.into()])?];
213 if !namespace.is_empty() && !name.is_empty() {
214 for sql in DROP_WAITS {
215 batch.push(self.db.prepare(*sql).bind(&[namespace.into(), name.into()])?);
216 }
217 }
218 self.db.batch(batch).await?;
219 Ok(())
220 }
221
222 /// `repo.purged`, `repo.deleted`, `repo.archived` and `branch.renamed`.
223 /// Says whether `event` was one this handles completely.
224 pub(crate) async fn on_retired(&self, event: &Event) -> Result<bool> {
225 if g1t_kit::lifecycle::on_purged(&self.db, event, PURGED).await? {
226 return Ok(true);
227 }
228 if let Some((repo_id, namespace, name)) = stopping(event) {
229 self.stop_runs_in(&repo_id, &namespace, &name).await?;
230 return Ok(true);
231 }
232 if event.kind != "branch.renamed" {
233 return Ok(false);
234 }
235 let Ok(renamed) = serde_json::from_value::<BranchRenamed>(event.data.clone()) else {
236 worker::console_error!("branch.renamed {} could not be read", event.id);
237 return Ok(true);
238 };
239 if renamed.from == renamed.to {
240 return Ok(true);
241 }
242 let names = [
243 renamed.repo_id.as_str().into(),
244 renamed.from.as_str().into(),
245 renamed.to.as_str().into(),
246 ];
247 self.db
248 .batch(vec![
249 self.db.prepare(BRANCH_RENAMED).bind(&names)?,
250 self.db.prepare(BASE_RENAMED).bind(&names)?,
251 ])
252 .await?;
253 Ok(true)
254 }
255}
256
257#[cfg(test)]
258mod tests {
259 use super::*;
260
261 fn repo(archived_at: Option<&str>) -> Repo {
262 Repo {
263 id: "rep_1".into(),
264 namespace: "acme".into(),
265 name: "web".into(),
266 description: None,
267 is_private: false,
268 owner_id: "usr_1".into(),
269 default_branch: "main".into(),
270 fork_of: None,
271 protected: false,
272 created_at: "2026-10-05T00:00:00Z".into(),
273 topics: Vec::new(),
274 website: None,
275 archived_at: archived_at.map(str::to_owned),
276 mirror: None,
277 }
278 }
279
280 #[test]
281 fn an_archived_repository_refuses_changes() {
282 assert!(matches!(writable(&repo(None)), Outcome::Ok(())));
283 match writable(&repo(Some("2026-10-05T00:00:00Z"))) {
284 Outcome::Fail(failure) => {
285 assert_eq!(failure.code, FailureCode::Forbidden);
286 assert!(failure.message.contains("acme/web is archived"));
287 }
288 Outcome::Ok(()) => panic!("an archived repository was writable"),
289 }
290 }
291
292 #[test]
293 fn purging_takes_only_the_repository_id() {
294 for sql in PURGED {
295 assert_eq!(g1t_kit::transfer::parameters(sql), 1, "{sql}");
296 }
297 // Rows kept by pull request go before the pull requests do.
298 let pulls = PURGED.iter().position(|sql| sql.starts_with("DELETE FROM pulls")).unwrap();
299 for table in ["session_entries", "check_runs", "review_runs", "agent_messages"] {
300 let at = PURGED.iter().position(|sql| sql.contains(&format!("FROM {table} "))).unwrap();
301 assert!(at < pulls, "{table}");
302 }
303 }
304
305 fn event(kind: &str, data: serde_json::Value) -> Event {
306 Event {
307 id: "evt_1".into(),
308 kind: kind.into(),
309 source: "repos".into(),
310 time: "2026-10-05T00:00:00Z".into(),
311 repo_id: Some("rep_1".into()),
312 actor: None,
313 data,
314 }
315 }
316
317 #[test]
318 fn deleting_or_archiving_stops_runs_and_unarchiving_does_not() {
319 let named = Some(("rep_1".to_owned(), "acme".to_owned(), "web".to_owned()));
320 let data = serde_json::json!({ "repoId": "rep_1", "namespace": "acme", "name": "web", "archived": true });
321 assert_eq!(stopping(&event("repo.deleted", data.clone())), named);
322 assert_eq!(stopping(&event("repo.archived", data)), named);
323 let unarchived = serde_json::json!({ "repoId": "rep_1", "namespace": "acme", "name": "web", "archived": false });
324 assert_eq!(stopping(&event("repo.unarchived", unarchived.clone())), None);
325 assert_eq!(stopping(&event("repo.archived", unarchived)), None);
326 assert_eq!(STOP_RUNS.matches('?').count(), 4);
327 for sql in DROP_WAITS {
328 assert_eq!(g1t_kit::transfer::parameters(sql), 2, "{sql}");
329 }
330 }
331
332 #[test]
333 fn waiting_runs_follow_with_the_names_they_bind() {
334 for sql in WAITS_MOVED {
335 assert_eq!(g1t_kit::transfer::parameters(sql), 7, "{sql}");
336 }
337 assert_eq!(g1t_kit::transfer::parameters(BRANCH_RENAMED), 3);
338 }
339}