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

677 lines25,896 bytesCodeBlame
1//! 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
13use g1t_contracts::events::IssueEvent;
14use g1t_contracts::repos::{GetByIdArgs, Repo, RepoPath};
15use g1t_contracts::time::rfc3339;
16use g1t_contracts::work::*;
17use g1t_contracts::{FailureCode, Outcome, User, new_id};
18use g1t_kit::now_ms;
19use serde::Deserialize;
20use sha2::{Digest, Sha256};
21use worker::Result;
22use worker::wasm_bindgen::JsValue;
23
24use crate::rows::user;
25use crate::{MAX_TITLE_CHARS, Work, valid_title};
26
27const MAX_BRIEF_CHARS: usize = 8_000;
28const MAX_PLANNED_ISSUES: usize = 12;
29const MAX_BODY_CHARS: usize = 20_000;
30const PLAN_PAGE: u32 = 20;
31/// Planning that has not reported back in this long has failed.
32const PLANNING_MINUTES: u64 = 20;
33/// How many g1t agents make changes in one repository at once. The rest of
34/// a plan waits its turn, which also keeps a large plan from flooding the
35/// sandboxes that checks and reviews need.
36const MAX_AGENTS_AT_WORK: usize = 6;
37
38#[derive(Deserialize)]
39struct PlanRow {
40 id: String,
41 repo_id: String,
42 brief: String,
43 status: PlanStatus,
44 summary: Option<String>,
45 issues: String,
46 error: Option<String>,
47 token_hash: String,
48 author_id: String,
49 author_name: String,
50 created_at: String,
51 finished_at: Option<String>,
52}
53
54impl From<PlanRow> for Plan {
55 fn from(row: PlanRow) -> Self {
56 // A sandbox that never reported is not still planning.
57 let oldest = rfc3339(now_ms().saturating_sub(PLANNING_MINUTES * 60 * 1000));
58 let abandoned = row.status == PlanStatus::Planning && row.created_at < oldest;
59 Plan {
60 id: row.id,
61 repo_id: row.repo_id,
62 brief: row.brief,
63 status: if abandoned {
64 PlanStatus::Failed
65 } else {
66 row.status
67 },
68 summary: row.summary.unwrap_or_default(),
69 issues: serde_json::from_str(&row.issues).unwrap_or_default(),
70 error: row
71 .error
72 .or_else(|| abandoned.then(|| "The planner did not report back.".to_owned())),
73 author: user(row.author_id, row.author_name),
74 created_at: row.created_at,
75 finished_at: row.finished_at,
76 progress: Vec::new(),
77 exchanges: Vec::new(),
78 }
79 }
80}
81
82#[derive(Deserialize)]
83struct LatestPull {
84 number: u32,
85 agent: String,
86 status: String,
87 stage: Option<String>,
88 stage_detail: Option<String>,
89}
90
91/// An issue waiting for a g1t agent, as selected.
92#[derive(Deserialize)]
93struct QueuedRow {
94 repo_id: String,
95 number: u32,
96 queued_by: String,
97}
98
99fn hash(token: &str) -> String {
100 hex::encode(Sha256::digest(token.as_bytes()))
101}
102
103/// Tidies what an agent proposed into something that can be opened as it
104/// is: bounded, with titles that are valid and dependencies that point only
105/// at earlier issues.
106pub(crate) fn tidy(proposed: Vec<PlannedIssue>) -> Vec<PlannedIssue> {
107 let mut issues: Vec<PlannedIssue> = Vec::new();
108 for issue in proposed.into_iter().take(MAX_PLANNED_ISSUES) {
109 // Nothing is dropped, so that the positions later issues depend on
110 // stay what the agent meant.
111 let title = match valid_title(&issue.title) {
112 Ok(title) => title.to_owned(),
113 Err(_) => match issue.title.trim() {
114 "" => "Untitled change".to_owned(),
115 long => long.chars().take(MAX_TITLE_CHARS).collect(),
116 },
117 };
118 let position = issues.len() as u32 + 1;
119 let mut depends_on: Vec<u32> = issue
120 .depends_on
121 .into_iter()
122 .filter(|earlier| (1..position).contains(earlier))
123 .collect();
124 depends_on.sort_unstable();
125 depends_on.dedup();
126 issues.push(PlannedIssue {
127 title,
128 body: issue.body.trim().chars().take(MAX_BODY_CHARS).collect(),
129 labels: normalize_labels(&issue.labels).unwrap_or_default(),
130 checks: issue
131 .checks
132 .into_iter()
133 .map(|check| check.trim().to_owned())
134 .filter(|check| !check.is_empty())
135 .take(10)
136 .collect(),
137 files: issue.files.into_iter().take(40).collect(),
138 depends_on,
139 number: None,
140 });
141 }
142 issues
143}
144
145impl Work {
146 async fn plan_row(&self, id: &str) -> Result<Option<PlanRow>> {
147 self.db
148 .prepare("SELECT * FROM plans WHERE id = ?")
149 .bind(&[id.into()])?
150 .first::<PlanRow>(None)
151 .await
152 }
153
154 /// Records an outcome to plan for, and returns what a sandbox needs to
155 /// plan it. Members of the repository's workspace only: a plan becomes
156 /// issues and agents at work, which the workspace pays for.
157 pub(crate) async fn start_plan(&self, a: StartPlanArgs) -> Result<Outcome<PlanJob>> {
158 let repo = match self.repo(&a.repo, &Some(a.actor.clone())).await? {
159 Outcome::Ok(repo) => repo,
160 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
161 };
162 if !a.actor.verified || !a.actor.is_member(&repo.namespace) {
163 return Ok(members_only());
164 }
165 let brief: String = a.brief.trim().chars().take(MAX_BRIEF_CHARS).collect();
166 if brief.is_empty() {
167 return Ok(Outcome::fail(
168 FailureCode::Invalid,
169 "Say what you want to be true when the work is done.",
170 ));
171 }
172 let now = now_ms();
173 let id = new_id("pln", now);
174 let mut bytes = [0u8; 32];
175 getrandom::getrandom(&mut bytes).expect("no source of randomness");
176 let token = hex::encode(bytes);
177 self.db
178 .prepare(
179 "INSERT INTO plans
180 (id, repo_id, brief, token_hash, author_id, author_name, created_at)
181 VALUES (?, ?, ?, ?, ?, ?, ?)",
182 )
183 .bind(&[
184 id.as_str().into(),
185 repo.id.as_str().into(),
186 brief.as_str().into(),
187 hash(&token).into(),
188 a.actor.id.as_str().into(),
189 a.actor.username.as_str().into(),
190 rfc3339(now).into(),
191 ])?
192 .run()
193 .await?;
194 Ok(Outcome::Ok(PlanJob {
195 plan_id: id,
196 token,
197 brief,
198 repo: RepoPath {
199 namespace: repo.namespace,
200 name: repo.name,
201 },
202 }))
203 }
204
205 /// Records the plan a sandbox's agent wrote, or why it could not write
206 /// one. The plan's token is the only credential.
207 pub(crate) async fn report_plan(&self, a: ReportPlanArgs) -> Result<Outcome<bool>> {
208 let row = self
209 .plan_row(&a.plan_id)
210 .await?
211 .filter(|row| row.token_hash == hash(&a.token));
212 let Some(row) = row else {
213 return Ok(Outcome::fail(FailureCode::NotFound, "Plan not found."));
214 };
215 if row.status != PlanStatus::Planning {
216 return Ok(Outcome::fail(
217 FailureCode::Conflict,
218 "This plan has already been reported.",
219 ));
220 }
221 let issues = tidy(a.issues);
222 let error = a.error.or_else(|| {
223 issues
224 .is_empty()
225 .then(|| "The planner proposed no issues.".to_owned())
226 });
227 self.db
228 .prepare(
229 "UPDATE plans SET status = ?, summary = ?, issues = ?, error = ?, finished_at = ?
230 WHERE id = ? AND status = 'planning'",
231 )
232 .bind(&[
233 if error.is_some() { "failed" } else { "ready" }.into(),
234 a.summary.trim().into(),
235 serde_json::to_string(&issues)?.into(),
236 error.as_deref().map_or(JsValue::NULL, JsValue::from),
237 rfc3339(now_ms()).into(),
238 row.id.into(),
239 ])?
240 .run()
241 .await?;
242 Ok(Outcome::Ok(true))
243 }
244
245 pub(crate) async fn get_plan(&self, a: PlanArgs) -> Result<Outcome<Plan>> {
246 let repo = match self.repo(&a.repo, &a.viewer).await? {
247 Outcome::Ok(repo) => repo,
248 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
249 };
250 if !a
251 .viewer
252 .is_some_and(|viewer| viewer.is_member(&repo.namespace))
253 {
254 return Ok(members_only());
255 }
256 Ok(
257 match self
258 .plan_row(&a.id)
259 .await?
260 .filter(|row| row.repo_id == repo.id)
261 {
262 Some(row) => {
263 let mut plan: Plan = row.into();
264 if plan.status == PlanStatus::Applied {
265 plan.progress = self.progress(&repo.id, &plan).await?;
266 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?;
272 }
273 Outcome::Ok(plan)
274 }
275 None => Outcome::fail(FailureCode::NotFound, "Plan not found."),
276 },
277 )
278 }
279
280 /// 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
351 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 };
356 if !a
357 .viewer
358 .is_some_and(|viewer| viewer.is_member(&repo.namespace))
359 {
360 return Ok(members_only());
361 }
362 let rows = self
363 .db
364 .prepare("SELECT * FROM plans WHERE repo_id = ? ORDER BY id DESC LIMIT ?")
365 .bind(&[repo.id.into(), PLAN_PAGE.into()])?
366 .all()
367 .await?
368 .results::<PlanRow>()?;
369 Ok(Outcome::Ok(rows.into_iter().map(Plan::from).collect()))
370 }
371
372 /// Opens a plan's issues. Each depends on the issues the plan said it
373 /// does, by their new numbers. With `assign`, every one is queued for a
374 /// g1t agent: those that depend on nothing are ready at once, and the
375 /// rest as what they depend on merges.
376 pub(crate) async fn apply_plan(&self, a: ApplyPlanArgs) -> Result<Outcome<Plan>> {
377 let repo = match self.repo(&a.repo, &Some(a.actor.clone())).await? {
378 Outcome::Ok(repo) => repo,
379 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
380 };
381 if !a.actor.verified || !a.actor.is_member(&repo.namespace) {
382 return Ok(members_only());
383 }
384 let Some(row) = self
385 .plan_row(&a.id)
386 .await?
387 .filter(|row| row.repo_id == repo.id)
388 else {
389 return Ok(Outcome::fail(FailureCode::NotFound, "Plan not found."));
390 };
391 // Only whoever flips it from ready to applied opens the issues.
392 let claimed = self
393 .db
394 .prepare(
395 "UPDATE plans SET status = 'applied' WHERE id = ? AND status = 'ready'
396 RETURNING id AS value",
397 )
398 .bind(&[row.id.as_str().into()])?
399 .first::<crate::rows::ValueRow>(None)
400 .await?;
401 if claimed.is_none() {
402 return Ok(Outcome::fail(
403 FailureCode::Conflict,
404 "This plan is not waiting to be applied.",
405 ));
406 }
407
408 let mut plan: Plan = row.into();
409 // What the person kept, in the plan's order. A dependency on an
410 // issue they dropped is dropped with it.
411 let kept: Vec<usize> = match &a.keep {
412 Some(positions) => (0..plan.issues.len())
413 .filter(|index| positions.contains(&(*index as u32 + 1)))
414 .collect(),
415 None => (0..plan.issues.len()).collect(),
416 };
417 let queued_by = a
418 .assign
419 .then(|| serde_json::to_string(&a.actor))
420 .transpose()?;
421 let mut numbers: Vec<Option<u32>> = vec![None; plan.issues.len()];
422 for index in kept {
423 let planned = &plan.issues[index];
424 let blocked_by: Vec<u32> = planned
425 .depends_on
426 .iter()
427 .filter_map(|position| numbers.get(*position as usize - 1).copied().flatten())
428 .collect();
429 let number = self.next_number(&repo.id).await?;
430 let now = now_ms();
431 let timestamp = rfc3339(now);
432 let issue_id = new_id("iss", now);
433 self.db
434 .prepare(
435 "INSERT INTO issues
436 (id, repo_id, number, title, body, labels, checks, blocked_by, queued_by,
437 author_id, author_name, created_at, updated_at)
438 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
439 )
440 .bind(&[
441 issue_id.as_str().into(),
442 repo.id.as_str().into(),
443 number.into(),
444 planned.title.as_str().into(),
445 planned.body.as_str().into(),
446 serde_json::to_string(&planned.labels)?.into(),
447 serde_json::to_string(&planned.checks)?.into(),
448 serde_json::to_string(&blocked_by)?.into(),
449 queued_by.as_deref().map_or(JsValue::NULL, JsValue::from),
450 a.actor.id.as_str().into(),
451 a.actor.username.as_str().into(),
452 timestamp.as_str().into(),
453 timestamp.as_str().into(),
454 ])?
455 .run()
456 .await?;
457 if a.assign {
458 self.note(
459 &repo.id,
460 number,
461 (&a.actor.id, &a.actor.username),
462 &if blocked_by.is_empty() {
463 "queued this for g1t-agent".to_owned()
464 } else {
465 format!(
466 "queued this for g1t-agent, to start once {} {} merged",
467 blocked_by
468 .iter()
469 .map(|number| format!("#{number}"))
470 .collect::<Vec<_>>()
471 .join(", "),
472 if blocked_by.len() == 1 { "has" } else { "have" }
473 )
474 },
475 )
476 .await?;
477 }
478 self.publish(
479 "issue.opened",
480 &repo.id,
481 &a.actor,
482 IssueEvent {
483 issue_id,
484 repo_id: repo.id.clone(),
485 number,
486 title: Some(planned.title.clone()),
487 ..IssueEvent::default()
488 },
489 )
490 .await?;
491 numbers[index] = Some(number);
492 }
493 for (issue, number) in plan.issues.iter_mut().zip(&numbers) {
494 issue.number = *number;
495 }
496 self.db
497 .prepare("UPDATE plans SET issues = ? WHERE id = ?")
498 .bind(&[
499 serde_json::to_string(&plan.issues)?.into(),
500 plan.id.as_str().into(),
501 ])?
502 .run()
503 .await?;
504 plan.status = PlanStatus::Applied;
505 Ok(Outcome::Ok(plan))
506 }
507
508 /// Queues an issue for a g1t agent, or takes it out of the queue.
509 pub(crate) async fn queue_issue(&self, a: QueueIssueArgs) -> Result<Outcome<bool>> {
510 let issue = match self.manageable_issue(&a.actor, &a.repo, a.number).await? {
511 Outcome::Ok(issue) => issue,
512 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
513 };
514 let queued_by = a
515 .queued
516 .then(|| serde_json::to_string(&a.actor))
517 .transpose()?;
518 self.db
519 .prepare("UPDATE issues SET queued_by = ? WHERE id = ?")
520 .bind(&[
521 queued_by.as_deref().map_or(JsValue::NULL, JsValue::from),
522 issue.id.as_str().into(),
523 ])?
524 .run()
525 .await?;
526 Ok(Outcome::Ok(true))
527 }
528
529 /// Takes the issues that are waiting for a g1t agent and can be given
530 /// one now: open, queued, with nobody working on them, and with
531 /// everything they depend on closed. Oldest first, in one repository or
532 /// in all, and no more than leaves each repository with
533 /// `MAX_AGENTS_AT_WORK` agents making changes at once.
534 ///
535 /// Taking an issue takes it out of the queue, in one statement, so two
536 /// callers cannot both start an agent on it. A caller that then cannot
537 /// start one puts it back with `queue_issue`.
538 pub(crate) async fn ready_issues(&self, a: ReadyIssuesArgs) -> Result<Vec<ReadyIssue>> {
539 let rows = self
540 .db
541 .prepare(
542 "SELECT repo_id, number, queued_by FROM issues
543 WHERE state = 'open' AND queued_by IS NOT NULL
544 AND (?1 IS NULL OR repo_id = ?1)
545 AND NOT EXISTS (
546 SELECT 1 FROM pulls
547 WHERE pulls.issue_id = issues.id AND pulls.status IN ('draft', 'open'))
548 AND NOT EXISTS (
549 SELECT 1 FROM json_each(issues.blocked_by) AS blocker
550 JOIN issues AS earlier
551 ON earlier.repo_id = issues.repo_id AND earlier.number = blocker.value
552 WHERE earlier.state = 'open')
553 ORDER BY number LIMIT 50",
554 )
555 .bind(&[a.repo_id.map_or(JsValue::NULL, JsValue::from)])?
556 .all()
557 .await?
558 .results::<QueuedRow>()?;
559 let mut ready = Vec::new();
560 let mut room: HashMap<String, usize> = HashMap::new();
561 for row in rows {
562 let Ok(actor) = serde_json::from_str::<User>(&row.queued_by) else {
563 continue;
564 };
565 // How many more agents this repository has room for.
566 if !room.contains_key(&row.repo_id) {
567 let working = self
568 .db
569 .prepare(
570 "SELECT count(*) AS n FROM pulls
571 WHERE repo_id = ? AND managed = 1 AND status = 'draft'",
572 )
573 .bind(&[row.repo_id.as_str().into()])?
574 .first::<crate::rows::NumberRow>(None)
575 .await?
576 .map_or(0, |row| row.n as usize);
577 room.insert(
578 row.repo_id.clone(),
579 MAX_AGENTS_AT_WORK.saturating_sub(working),
580 );
581 }
582 let left = room.get_mut(&row.repo_id).expect("just inserted");
583 if *left == 0 {
584 continue;
585 }
586 let taken = self
587 .db
588 .prepare(
589 "UPDATE issues SET queued_by = NULL
590 WHERE repo_id = ? AND number = ? AND queued_by IS NOT NULL
591 RETURNING id AS value",
592 )
593 .bind(&[row.repo_id.as_str().into(), row.number.into()])?
594 .first::<crate::rows::ValueRow>(None)
595 .await?;
596 if taken.is_none() {
597 continue;
598 }
599 *left -= 1;
600 // Whoever queued it could see the repository then; where it is
601 // now is asked as them.
602 let repo: Outcome<Repo> = g1t_kit::call(
603 &self.repos,
604 "get_by_id",
605 &GetByIdArgs {
606 id: row.repo_id.clone(),
607 viewer: Some(actor.clone()),
608 },
609 )
610 .await?;
611 if let Outcome::Ok(repo) = repo {
612 ready.push(ReadyIssue {
613 repo: RepoPath {
614 namespace: repo.namespace,
615 name: repo.name,
616 },
617 number: row.number,
618 actor,
619 });
620 }
621 }
622 Ok(ready)
623 }
624}
625
626fn members_only<T>() -> Outcome<T> {
627 Outcome::fail(
628 FailureCode::Forbidden,
629 "Only members of the repository's workspace can plan work for it.",
630 )
631}
632
633#[cfg(test)]
634mod tests {
635 use super::*;
636
637 fn proposed(title: &str, depends_on: &[u32]) -> PlannedIssue {
638 PlannedIssue {
639 title: title.to_owned(),
640 body: " What to do. ".to_owned(),
641 labels: vec!["Feature".to_owned()],
642 checks: vec![" cargo test ".to_owned(), String::new()],
643 files: vec!["src/lib.rs".to_owned()],
644 depends_on: depends_on.to_vec(),
645 number: None,
646 }
647 }
648
649 #[test]
650 fn a_proposal_is_trimmed_and_normalised() {
651 let issues = tidy(vec![proposed(" Add a flag ", &[])]);
652 assert_eq!(issues[0].title, "Add a flag");
653 assert_eq!(issues[0].body, "What to do.");
654 assert_eq!(issues[0].labels, ["feature"]);
655 assert_eq!(issues[0].checks, ["cargo test"]);
656 }
657
658 #[test]
659 fn an_issue_can_only_depend_on_earlier_ones() {
660 let issues = tidy(vec![
661 proposed("First", &[2]),
662 proposed("Second", &[1, 1, 2, 9]),
663 proposed("Third", &[2, 1]),
664 ]);
665 assert!(issues[0].depends_on.is_empty());
666 assert_eq!(issues[1].depends_on, [1]);
667 assert_eq!(issues[2].depends_on, [1, 2]);
668 }
669
670 #[test]
671 fn a_plan_is_bounded() {
672 let many = (0..30)
673 .map(|i| proposed(&format!("Issue {i}"), &[]))
674 .collect();
675 assert_eq!(tidy(many).len(), MAX_PLANNED_ISSUES);
676 }
677}