Skip to content

g1t/services/work/src/retired.rs

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