flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/work/src/plans.rs

704 lines27,793 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Agents as a team: lifecycle, merge queue, billing and a new shell1//! Planning: an outcome someone wrote, turned by an agent into issues and
2//! the order they have to land in.
3//!
4//! This service keeps the plan. The runner starts the sandbox in which an
5//! agent reads the repository and writes it, and the sandbox reports back
6//! with the plan's one-time token. A person then reads the proposal and
7//! applies it, which opens the issues. Issues that depend on others are
8//! held back and handed to a g1t agent when what they depend on has
9//! merged, so that each starts from the result of the last.
10
11use std::collections::HashMap;
12
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look13use g1t_contracts::access::Capability;
Agents as a team: lifecycle, merge queue, billing and a new shell14use g1t_contracts::events::IssueEvent;
15use g1t_contracts::repos::{GetByIdArgs, Repo, RepoPath};
16use g1t_contracts::time::rfc3339;
17use g1t_contracts::work::*;
18use g1t_contracts::{FailureCode, Outcome, User, new_id};
19use g1t_kit::now_ms;
20use serde::Deserialize;
21use sha2::{Digest, Sha256};
22use worker::Result;
23use worker::wasm_bindgen::JsValue;
24
25use crate::rows::user;
26use crate::{MAX_TITLE_CHARS, Work, valid_title};
27
28const MAX_BRIEF_CHARS: usize = 8_000;
29const MAX_PLANNED_ISSUES: usize = 12;
30const MAX_BODY_CHARS: usize = 20_000;
31const PLAN_PAGE: u32 = 20;
32/// Planning that has not reported back in this long has failed.
33const PLANNING_MINUTES: u64 = 20;
34/// How many g1t agents make changes in one repository at once. The rest of
35/// a plan waits its turn, which also keeps a large plan from flooding the
36/// sandboxes that checks and reviews need.
37const MAX_AGENTS_AT_WORK: usize = 6;
38
39#[derive(Deserialize)]
40struct PlanRow {
41 id: String,
42 repo_id: String,
43 brief: String,
44 status: PlanStatus,
45 summary: Option<String>,
46 issues: String,
47 error: Option<String>,
48 token_hash: String,
49 author_id: String,
50 author_name: String,
51 created_at: String,
52 finished_at: Option<String>,
53}
54
55impl From<PlanRow> for Plan {
56 fn from(row: PlanRow) -> Self {
57 // A sandbox that never reported is not still planning.
58 let oldest = rfc3339(now_ms().saturating_sub(PLANNING_MINUTES * 60 * 1000));
59 let abandoned = row.status == PlanStatus::Planning && row.created_at < oldest;
60 Plan {
61 id: row.id,
62 repo_id: row.repo_id,
63 brief: row.brief,
64 status: if abandoned {
65 PlanStatus::Failed
66 } else {
67 row.status
68 },
69 summary: row.summary.unwrap_or_default(),
70 issues: serde_json::from_str(&row.issues).unwrap_or_default(),
71 error: row
72 .error
73 .or_else(|| abandoned.then(|| "The planner did not report back.".to_owned())),
74 author: user(row.author_id, row.author_name),
75 created_at: row.created_at,
76 finished_at: row.finished_at,
Dark gray base with lavender as an accent, and a live outcome view for plans77 progress: Vec::new(),
Agents ask each other, hand each other work, and answer78 exchanges: Vec::new(),
Agents as a team: lifecycle, merge queue, billing and a new shell79 }
80 }
81}
82
Dark gray base with lavender as an accent, and a live outcome view for plans83#[derive(Deserialize)]
84struct LatestPull {
85 number: u32,
86 agent: String,
87 status: String,
88 stage: Option<String>,
89 stage_detail: Option<String>,
90}
91
Agents as a team: lifecycle, merge queue, billing and a new shell92/// An issue waiting for a g1t agent, as selected.
93#[derive(Deserialize)]
94struct QueuedRow {
95 repo_id: String,
96 number: u32,
97 queued_by: String,
98}
99
100fn hash(token: &str) -> String {
101 hex::encode(Sha256::digest(token.as_bytes()))
102}
103
104/// Tidies what an agent proposed into something that can be opened as it
105/// is: bounded, with titles that are valid and dependencies that point only
106/// at earlier issues.
107pub(crate) fn tidy(proposed: Vec<PlannedIssue>) -> Vec<PlannedIssue> {
108 let mut issues: Vec<PlannedIssue> = Vec::new();
109 for issue in proposed.into_iter().take(MAX_PLANNED_ISSUES) {
110 // Nothing is dropped, so that the positions later issues depend on
111 // stay what the agent meant.
112 let title = match valid_title(&issue.title) {
113 Ok(title) => title.to_owned(),
114 Err(_) => match issue.title.trim() {
115 "" => "Untitled change".to_owned(),
116 long => long.chars().take(MAX_TITLE_CHARS).collect(),
117 },
118 };
119 let position = issues.len() as u32 + 1;
120 let mut depends_on: Vec<u32> = issue
121 .depends_on
122 .into_iter()
123 .filter(|earlier| (1..position).contains(earlier))
124 .collect();
125 depends_on.sort_unstable();
126 depends_on.dedup();
127 issues.push(PlannedIssue {
128 title,
129 body: issue.body.trim().chars().take(MAX_BODY_CHARS).collect(),
130 labels: normalize_labels(&issue.labels).unwrap_or_default(),
Fast pages, required checks on the branch, self-hosted runners, honest incidents131 done: issue
132 .done
Agents as a team: lifecycle, merge queue, billing and a new shell133 .into_iter()
Fast pages, required checks on the branch, self-hosted runners, honest incidents134 .map(|item| item.trim().to_owned())
135 .filter(|item| !item.is_empty())
Agents as a team: lifecycle, merge queue, billing and a new shell136 .take(10)
137 .collect(),
138 files: issue.files.into_iter().take(40).collect(),
139 depends_on,
140 number: None,
141 });
142 }
143 issues
144}
145
146impl Work {
147 async fn plan_row(&self, id: &str) -> Result<Option<PlanRow>> {
148 self.db
149 .prepare("SELECT * FROM plans WHERE id = ?")
150 .bind(&[id.into()])?
151 .first::<PlanRow>(None)
152 .await
153 }
154
155 /// Records an outcome to plan for, and returns what a sandbox needs to
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look156 /// plan it. Needs the Write role: a plan becomes issues and agents at
157 /// work, which the workspace pays for.
Agents as a team: lifecycle, merge queue, billing and a new shell158 pub(crate) async fn start_plan(&self, a: StartPlanArgs) -> Result<Outcome<PlanJob>> {
159 let repo = match self.repo(&a.repo, &Some(a.actor.clone())).await? {
160 Outcome::Ok(repo) => repo,
161 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
162 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look163 if let Outcome::Fail(failure) = crate::retired::writable(&repo) {
164 return Ok(Outcome::Fail(failure));
Agents as a team: lifecycle, merge queue, billing and a new shell165 }
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look166 if let Some(refused) = may_plan(&a.actor, &repo) {
167 return Ok(refused);
168 }
Agents as a team: lifecycle, merge queue, billing and a new shell169 let brief: String = a.brief.trim().chars().take(MAX_BRIEF_CHARS).collect();
170 if brief.is_empty() {
171 return Ok(Outcome::fail(
172 FailureCode::Invalid,
173 "Say what you want to be true when the work is done.",
174 ));
175 }
176 let now = now_ms();
177 let id = new_id("pln", now);
178 let mut bytes = [0u8; 32];
179 getrandom::getrandom(&mut bytes).expect("no source of randomness");
180 let token = hex::encode(bytes);
181 self.db
182 .prepare(
183 "INSERT INTO plans
184 (id, repo_id, brief, token_hash, author_id, author_name, created_at)
185 VALUES (?, ?, ?, ?, ?, ?, ?)",
186 )
187 .bind(&[
188 id.as_str().into(),
189 repo.id.as_str().into(),
190 brief.as_str().into(),
191 hash(&token).into(),
192 a.actor.id.as_str().into(),
193 a.actor.username.as_str().into(),
194 rfc3339(now).into(),
195 ])?
196 .run()
197 .await?;
198 Ok(Outcome::Ok(PlanJob {
199 plan_id: id,
200 token,
201 brief,
202 repo: RepoPath {
203 namespace: repo.namespace,
204 name: repo.name,
205 },
206 }))
207 }
208
209 /// Records the plan a sandbox's agent wrote, or why it could not write
210 /// one. The plan's token is the only credential.
211 pub(crate) async fn report_plan(&self, a: ReportPlanArgs) -> Result<Outcome<bool>> {
212 let row = self
213 .plan_row(&a.plan_id)
214 .await?
215 .filter(|row| row.token_hash == hash(&a.token));
216 let Some(row) = row else {
217 return Ok(Outcome::fail(FailureCode::NotFound, "Plan not found."));
218 };
219 if row.status != PlanStatus::Planning {
220 return Ok(Outcome::fail(
221 FailureCode::Conflict,
222 "This plan has already been reported.",
223 ));
224 }
225 let issues = tidy(a.issues);
226 let error = a.error.or_else(|| {
227 issues
228 .is_empty()
229 .then(|| "The planner proposed no issues.".to_owned())
230 });
231 self.db
232 .prepare(
233 "UPDATE plans SET status = ?, summary = ?, issues = ?, error = ?, finished_at = ?
234 WHERE id = ? AND status = 'planning'",
235 )
236 .bind(&[
237 if error.is_some() { "failed" } else { "ready" }.into(),
238 a.summary.trim().into(),
239 serde_json::to_string(&issues)?.into(),
240 error.as_deref().map_or(JsValue::NULL, JsValue::from),
241 rfc3339(now_ms()).into(),
242 row.id.into(),
243 ])?
244 .run()
245 .await?;
246 Ok(Outcome::Ok(true))
247 }
248
249 pub(crate) async fn get_plan(&self, a: PlanArgs) -> Result<Outcome<Plan>> {
250 let repo = match self.repo(&a.repo, &a.viewer).await? {
251 Outcome::Ok(repo) => repo,
252 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
253 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look254 // The repos service found it for the viewer, so they can read it,
255 // and whoever can read a repository can read its plans.
Agents as a team: lifecycle, merge queue, billing and a new shell256 Ok(
257 match self
258 .plan_row(&a.id)
259 .await?
260 .filter(|row| row.repo_id == repo.id)
261 {
Dark gray base with lavender as an accent, and a live outcome view for plans262 Some(row) => {
263 let mut plan: Plan = row.into();
264 if plan.status == PlanStatus::Applied {
265 plan.progress = self.progress(&repo.id, &plan).await?;
Agents ask each other, hand each other work, and answer266 let numbers: Vec<u32> = plan
267 .progress
268 .iter()
269 .flat_map(|item| std::iter::once(item.number).chain(item.pull))
270 .collect();
271 plan.exchanges = self.exchanges(&repo.id, &numbers).await?;
Dark gray base with lavender as an accent, and a live outcome view for plans272 }
273 Outcome::Ok(plan)
274 }
Agents as a team: lifecycle, merge queue, billing and a new shell275 None => Outcome::fail(FailureCode::NotFound, "Plan not found."),
276 },
277 )
278 }
279
Dark gray base with lavender as an accent, and a live outcome view for plans280 /// Where each issue an applied plan opened stands now, all at once.
281 async fn progress(&self, repo_id: &str, plan: &Plan) -> Result<Vec<IssueProgress>> {
282 let numbers: Vec<u32> = plan.issues.iter().filter_map(|issue| issue.number).collect();
283 let issues = futures_util::future::try_join_all(
284 numbers.iter().map(|number| self.issue(repo_id, *number)),
285 )
286 .await?;
287 let pulls = futures_util::future::try_join_all(numbers.iter().map(|number| async move {
288 self.db
289 .prepare(
290 "SELECT number, agent, status, stage, stage_detail FROM pulls
291 WHERE repo_id = ? AND issue_number = ? AND status != 'closed'
292 ORDER BY number DESC LIMIT 1",
293 )
294 .bind(&[repo_id.into(), (*number).into()])?
295 .first::<LatestPull>(None)
296 .await
297 }))
298 .await?;
299 let open: std::collections::HashSet<u32> = issues
300 .iter()
301 .flatten()
302 .filter(|issue| issue.state == State::Open)
303 .map(|issue| issue.number)
304 .collect();
305 Ok(issues
306 .into_iter()
307 .zip(pulls)
308 .filter_map(|(issue, pull)| {
309 let issue = issue?;
310 let blocked_by: Vec<u32> = issue
311 .blocked_by
312 .iter()
313 .copied()
314 .filter(|number| open.contains(number))
315 .collect();
316 let (state, detail) = if issue.state == State::Closed {
317 match (issue.reason, issue.resolved_by) {
318 (Some(IssueReason::Completed), Some(by)) => {
319 ("landed".to_owned(), format!("Landed with #{by}."))
320 }
321 _ => ("closed".to_owned(), "Closed without landing.".to_owned()),
322 }
323 } else if let Some(pull) = pull.as_ref().filter(|pull| pull.status != "merged") {
324 (
325 pull.stage.clone().unwrap_or_else(|| {
326 if pull.status == "draft" { "working" } else { "ready" }.to_owned()
327 }),
328 pull.stage_detail.clone().unwrap_or_default(),
329 )
330 } else if !blocked_by.is_empty() {
331 let named: Vec<String> = blocked_by.iter().map(|n| format!("#{n}")).collect();
332 ("blocked".to_owned(), format!("Waiting for {} to land.", named.join(", ")))
333 } else if issue.queued {
334 ("waiting".to_owned(), "Waiting for an agent to be free.".to_owned())
335 } else {
336 ("open".to_owned(), "Nobody is working on it.".to_owned())
337 };
338 Some(IssueProgress {
339 number: issue.number,
340 title: issue.title,
341 state,
342 detail,
343 blocked_by,
344 pull: pull.as_ref().map(|pull| pull.number),
345 agent: pull.map(|pull| pull.agent).or(issue.agent),
346 })
347 })
348 .collect())
349 }
350
Agents as a team: lifecycle, merge queue, billing and a new shell351 pub(crate) async fn list_plans(&self, a: ViewArgs) -> Result<Outcome<Vec<Plan>>> {
352 let repo = match self.repo(&a.repo, &a.viewer).await? {
353 Outcome::Ok(repo) => repo,
354 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
355 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look356 // Found for the viewer: they can read it, and so its plans.
Agents as a team: lifecycle, merge queue, billing and a new shell357 let rows = self
358 .db
359 .prepare("SELECT * FROM plans WHERE repo_id = ? ORDER BY id DESC LIMIT ?")
360 .bind(&[repo.id.into(), PLAN_PAGE.into()])?
361 .all()
362 .await?
363 .results::<PlanRow>()?;
364 Ok(Outcome::Ok(rows.into_iter().map(Plan::from).collect()))
365 }
366
367 /// Opens a plan's issues. Each depends on the issues the plan said it
368 /// does, by their new numbers. With `assign`, every one is queued for a
369 /// g1t agent: those that depend on nothing are ready at once, and the
370 /// rest as what they depend on merges.
371 pub(crate) async fn apply_plan(&self, a: ApplyPlanArgs) -> Result<Outcome<Plan>> {
372 let repo = match self.repo(&a.repo, &Some(a.actor.clone())).await? {
373 Outcome::Ok(repo) => repo,
374 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
375 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look376 if let Outcome::Fail(failure) = crate::retired::writable(&repo) {
377 return Ok(Outcome::Fail(failure));
378 }
379 if let Some(refused) = may_plan(&a.actor, &repo) {
380 return Ok(refused);
Agents as a team: lifecycle, merge queue, billing and a new shell381 }
382 let Some(row) = self
383 .plan_row(&a.id)
384 .await?
385 .filter(|row| row.repo_id == repo.id)
386 else {
387 return Ok(Outcome::fail(FailureCode::NotFound, "Plan not found."));
388 };
389 // Only whoever flips it from ready to applied opens the issues.
390 let claimed = self
391 .db
392 .prepare(
393 "UPDATE plans SET status = 'applied' WHERE id = ? AND status = 'ready'
394 RETURNING id AS value",
395 )
396 .bind(&[row.id.as_str().into()])?
397 .first::<crate::rows::ValueRow>(None)
398 .await?;
399 if claimed.is_none() {
400 return Ok(Outcome::fail(
401 FailureCode::Conflict,
402 "This plan is not waiting to be applied.",
403 ));
404 }
405
406 let mut plan: Plan = row.into();
407 // What the person kept, in the plan's order. A dependency on an
408 // issue they dropped is dropped with it.
409 let kept: Vec<usize> = match &a.keep {
410 Some(positions) => (0..plan.issues.len())
411 .filter(|index| positions.contains(&(*index as u32 + 1)))
412 .collect(),
413 None => (0..plan.issues.len()).collect(),
414 };
415 let queued_by = a
416 .assign
417 .then(|| serde_json::to_string(&a.actor))
418 .transpose()?;
419 let mut numbers: Vec<Option<u32>> = vec![None; plan.issues.len()];
g1t is the stored author of what it opens; the person who asked is requested_by and keeps the author's rights420 // Whoever applies the plan opens its issues; g1t's agent applying
421 // one opens them as g1t, for the person it works for.
422 let (author, requested_by) = authorship(&a.actor, false);
Agents as a team: lifecycle, merge queue, billing and a new shell423 for index in kept {
424 let planned = &plan.issues[index];
425 let blocked_by: Vec<u32> = planned
426 .depends_on
427 .iter()
428 .filter_map(|position| numbers.get(*position as usize - 1).copied().flatten())
429 .collect();
430 let number = self.next_number(&repo.id).await?;
431 let now = now_ms();
432 let timestamp = rfc3339(now);
433 let issue_id = new_id("iss", now);
434 self.db
435 .prepare(
436 "INSERT INTO issues
437 (id, repo_id, number, title, body, labels, checks, blocked_by, queued_by,
g1t is the stored author of what it opens; the person who asked is requested_by and keeps the author's rights438 author_id, author_name, requested_by_id, requested_by_name, created_at, updated_at)
439 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
Agents as a team: lifecycle, merge queue, billing and a new shell440 )
441 .bind(&[
442 issue_id.as_str().into(),
443 repo.id.as_str().into(),
444 number.into(),
445 planned.title.as_str().into(),
Fast pages, required checks on the branch, self-hosted runners, honest incidents446 // What done looks like, in words, for the agent and
447 // reviewers; the branch's required checks gate the merge.
448 with_definition_of_done(&planned.body, &planned.done).into(),
Agents as a team: lifecycle, merge queue, billing and a new shell449 serde_json::to_string(&planned.labels)?.into(),
Fast pages, required checks on the branch, self-hosted runners, honest incidents450 "[]".into(),
Agents as a team: lifecycle, merge queue, billing and a new shell451 serde_json::to_string(&blocked_by)?.into(),
452 queued_by.as_deref().map_or(JsValue::NULL, JsValue::from),
g1t is the stored author of what it opens; the person who asked is requested_by and keeps the author's rights453 author.id.as_str().into(),
454 author.username.as_str().into(),
455 requested_by.as_ref().map_or(JsValue::NULL, |user| user.id.as_str().into()),
456 requested_by.as_ref().map_or(JsValue::NULL, |user| user.username.as_str().into()),
Agents as a team: lifecycle, merge queue, billing and a new shell457 timestamp.as_str().into(),
458 timestamp.as_str().into(),
459 ])?
460 .run()
461 .await?;
462 if a.assign {
463 self.note(
464 &repo.id,
465 number,
466 (&a.actor.id, &a.actor.username),
467 &if blocked_by.is_empty() {
g1t is one name: its agent's work, commits and comments show as @g1t, and nobody can claim g1t or g1t-agent468 "queued this for g1t".to_owned()
Agents as a team: lifecycle, merge queue, billing and a new shell469 } else {
470 format!(
g1t is one name: its agent's work, commits and comments show as @g1t, and nobody can claim g1t or g1t-agent471 "queued this for g1t, to start once {} {} merged",
Agents as a team: lifecycle, merge queue, billing and a new shell472 blocked_by
473 .iter()
474 .map(|number| format!("#{number}"))
475 .collect::<Vec<_>>()
476 .join(", "),
477 if blocked_by.len() == 1 { "has" } else { "have" }
478 )
479 },
480 )
481 .await?;
482 }
483 self.publish(
484 "issue.opened",
485 &repo.id,
486 &a.actor,
487 IssueEvent {
488 issue_id,
489 repo_id: repo.id.clone(),
490 number,
g1t is the stored author of what it opens; the person who asked is requested_by and keeps the author's rights491 author: Some((&author).into()),
492 requested_by: requested_by.as_ref().map(Into::into),
Agents as a team: lifecycle, merge queue, billing and a new shell493 title: Some(planned.title.clone()),
494 ..IssueEvent::default()
495 },
496 )
497 .await?;
498 numbers[index] = Some(number);
499 }
500 for (issue, number) in plan.issues.iter_mut().zip(&numbers) {
501 issue.number = *number;
502 }
503 self.db
504 .prepare("UPDATE plans SET issues = ? WHERE id = ?")
505 .bind(&[
506 serde_json::to_string(&plan.issues)?.into(),
507 plan.id.as_str().into(),
508 ])?
509 .run()
510 .await?;
511 plan.status = PlanStatus::Applied;
512 Ok(Outcome::Ok(plan))
513 }
514
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look515 /// Queues an issue for a g1t agent, which needs the Write role, or
516 /// takes it out of the queue, as whoever may change the issue.
Agents as a team: lifecycle, merge queue, billing and a new shell517 pub(crate) async fn queue_issue(&self, a: QueueIssueArgs) -> Result<Outcome<bool>> {
518 let issue = match self.manageable_issue(&a.actor, &a.repo, a.number).await? {
519 Outcome::Ok(issue) => issue,
520 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
521 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look522 if a.queued {
523 let repo = match self.repo(&a.repo, &Some(a.actor.clone())).await? {
524 Outcome::Ok(repo) => repo,
525 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
526 };
527 if let Outcome::Fail(failure) = crate::allowed(Some(&a.actor), &repo, Capability::Run) {
528 return Ok(Outcome::Fail(failure));
529 }
530 }
Agents as a team: lifecycle, merge queue, billing and a new shell531 let queued_by = a
532 .queued
533 .then(|| serde_json::to_string(&a.actor))
534 .transpose()?;
535 self.db
536 .prepare("UPDATE issues SET queued_by = ? WHERE id = ?")
537 .bind(&[
538 queued_by.as_deref().map_or(JsValue::NULL, JsValue::from),
539 issue.id.as_str().into(),
540 ])?
541 .run()
542 .await?;
543 Ok(Outcome::Ok(true))
544 }
545
546 /// Takes the issues that are waiting for a g1t agent and can be given
547 /// one now: open, queued, with nobody working on them, and with
548 /// everything they depend on closed. Oldest first, in one repository or
549 /// in all, and no more than leaves each repository with
550 /// `MAX_AGENTS_AT_WORK` agents making changes at once.
551 ///
552 /// Taking an issue takes it out of the queue, in one statement, so two
553 /// callers cannot both start an agent on it. A caller that then cannot
554 /// start one puts it back with `queue_issue`.
555 pub(crate) async fn ready_issues(&self, a: ReadyIssuesArgs) -> Result<Vec<ReadyIssue>> {
556 let rows = self
557 .db
558 .prepare(
559 "SELECT repo_id, number, queued_by FROM issues
560 WHERE state = 'open' AND queued_by IS NOT NULL
561 AND (?1 IS NULL OR repo_id = ?1)
562 AND NOT EXISTS (
563 SELECT 1 FROM pulls
564 WHERE pulls.issue_id = issues.id AND pulls.status IN ('draft', 'open'))
565 AND NOT EXISTS (
566 SELECT 1 FROM json_each(issues.blocked_by) AS blocker
567 JOIN issues AS earlier
568 ON earlier.repo_id = issues.repo_id AND earlier.number = blocker.value
569 WHERE earlier.state = 'open')
570 ORDER BY number LIMIT 50",
571 )
572 .bind(&[a.repo_id.map_or(JsValue::NULL, JsValue::from)])?
573 .all()
574 .await?
575 .results::<QueuedRow>()?;
576 let mut ready = Vec::new();
577 let mut room: HashMap<String, usize> = HashMap::new();
578 for row in rows {
579 let Ok(actor) = serde_json::from_str::<User>(&row.queued_by) else {
580 continue;
581 };
582 // How many more agents this repository has room for.
583 if !room.contains_key(&row.repo_id) {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look584 // An archived or deleted repository's issues stay queued,
585 // untaken, for when it is writable again.
586 if !self.repo_active(&row.repo_id).await? {
587 room.insert(row.repo_id.clone(), 0);
588 continue;
589 }
Agents as a team: lifecycle, merge queue, billing and a new shell590 let working = self
591 .db
592 .prepare(
593 "SELECT count(*) AS n FROM pulls
594 WHERE repo_id = ? AND managed = 1 AND status = 'draft'",
595 )
596 .bind(&[row.repo_id.as_str().into()])?
597 .first::<crate::rows::NumberRow>(None)
598 .await?
599 .map_or(0, |row| row.n as usize);
600 room.insert(
601 row.repo_id.clone(),
602 MAX_AGENTS_AT_WORK.saturating_sub(working),
603 );
604 }
605 let left = room.get_mut(&row.repo_id).expect("just inserted");
606 if *left == 0 {
607 continue;
608 }
609 let taken = self
610 .db
611 .prepare(
612 "UPDATE issues SET queued_by = NULL
613 WHERE repo_id = ? AND number = ? AND queued_by IS NOT NULL
614 RETURNING id AS value",
615 )
616 .bind(&[row.repo_id.as_str().into(), row.number.into()])?
617 .first::<crate::rows::ValueRow>(None)
618 .await?;
619 if taken.is_none() {
620 continue;
621 }
622 *left -= 1;
623 // Whoever queued it could see the repository then; where it is
624 // now is asked as them.
625 let repo: Outcome<Repo> = g1t_kit::call(
626 &self.repos,
627 "get_by_id",
628 &GetByIdArgs {
629 id: row.repo_id.clone(),
630 viewer: Some(actor.clone()),
631 },
632 )
633 .await?;
634 if let Outcome::Ok(repo) = repo {
635 ready.push(ReadyIssue {
636 repo: RepoPath {
637 namespace: repo.namespace,
638 name: repo.name,
639 },
640 number: row.number,
641 actor,
642 });
643 }
644 }
645 Ok(ready)
646 }
647}
648
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look649/// Refuses whoever may not plan work in `repo`: the Write role, verified.
650fn may_plan<T>(actor: &User, repo: &Repo) -> Option<Outcome<T>> {
651 if !actor.verified {
652 return Some(Outcome::fail(FailureCode::Forbidden, crate::UNVERIFIED));
653 }
654 match crate::allowed(Some(actor), repo, Capability::Run) {
655 Outcome::Ok(()) => None,
656 Outcome::Fail(failure) => Some(Outcome::Fail(failure)),
657 }
Agents as a team: lifecycle, merge queue, billing and a new shell658}
659
660#[cfg(test)]
661mod tests {
662 use super::*;
663
664 fn proposed(title: &str, depends_on: &[u32]) -> PlannedIssue {
665 PlannedIssue {
666 title: title.to_owned(),
667 body: " What to do. ".to_owned(),
668 labels: vec!["Feature".to_owned()],
Fast pages, required checks on the branch, self-hosted runners, honest incidents669 done: vec![" The flag is documented in --help. ".to_owned(), String::new()],
Agents as a team: lifecycle, merge queue, billing and a new shell670 files: vec!["src/lib.rs".to_owned()],
671 depends_on: depends_on.to_vec(),
672 number: None,
673 }
674 }
675
676 #[test]
677 fn a_proposal_is_trimmed_and_normalised() {
678 let issues = tidy(vec![proposed(" Add a flag ", &[])]);
679 assert_eq!(issues[0].title, "Add a flag");
680 assert_eq!(issues[0].body, "What to do.");
681 assert_eq!(issues[0].labels, ["feature"]);
Fast pages, required checks on the branch, self-hosted runners, honest incidents682 assert_eq!(issues[0].done, ["The flag is documented in --help."]);
Agents as a team: lifecycle, merge queue, billing and a new shell683 }
684
685 #[test]
686 fn an_issue_can_only_depend_on_earlier_ones() {
687 let issues = tidy(vec![
688 proposed("First", &[2]),
689 proposed("Second", &[1, 1, 2, 9]),
690 proposed("Third", &[2, 1]),
691 ]);
692 assert!(issues[0].depends_on.is_empty());
693 assert_eq!(issues[1].depends_on, [1]);
694 assert_eq!(issues[2].depends_on, [1, 2]);
695 }
696
697 #[test]
698 fn a_plan_is_bounded() {
699 let many = (0..30)
700 .map(|i| proposed(&format!("Issue {i}"), &[]))
701 .collect();
702 assert_eq!(tidy(many).len(), MAX_PLANNED_ISSUES);
703 }
704}