g1t/services/work/src/plans.rs

583 lines21,802 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 }
77 }
78}
79
80/// An issue waiting for a g1t agent, as selected.
81#[derive(Deserialize)]
82struct QueuedRow {
83 repo_id: String,
84 number: u32,
85 queued_by: String,
86}
87
88fn hash(token: &str) -> String {
89 hex::encode(Sha256::digest(token.as_bytes()))
90}
91
92/// Tidies what an agent proposed into something that can be opened as it
93/// is: bounded, with titles that are valid and dependencies that point only
94/// at earlier issues.
95pub(crate) fn tidy(proposed: Vec<PlannedIssue>) -> Vec<PlannedIssue> {
96 let mut issues: Vec<PlannedIssue> = Vec::new();
97 for issue in proposed.into_iter().take(MAX_PLANNED_ISSUES) {
98 // Nothing is dropped, so that the positions later issues depend on
99 // stay what the agent meant.
100 let title = match valid_title(&issue.title) {
101 Ok(title) => title.to_owned(),
102 Err(_) => match issue.title.trim() {
103 "" => "Untitled change".to_owned(),
104 long => long.chars().take(MAX_TITLE_CHARS).collect(),
105 },
106 };
107 let position = issues.len() as u32 + 1;
108 let mut depends_on: Vec<u32> = issue
109 .depends_on
110 .into_iter()
111 .filter(|earlier| (1..position).contains(earlier))
112 .collect();
113 depends_on.sort_unstable();
114 depends_on.dedup();
115 issues.push(PlannedIssue {
116 title,
117 body: issue.body.trim().chars().take(MAX_BODY_CHARS).collect(),
118 labels: normalize_labels(&issue.labels).unwrap_or_default(),
119 checks: issue
120 .checks
121 .into_iter()
122 .map(|check| check.trim().to_owned())
123 .filter(|check| !check.is_empty())
124 .take(10)
125 .collect(),
126 files: issue.files.into_iter().take(40).collect(),
127 depends_on,
128 number: None,
129 });
130 }
131 issues
132}
133
134impl Work {
135 async fn plan_row(&self, id: &str) -> Result<Option<PlanRow>> {
136 self.db
137 .prepare("SELECT * FROM plans WHERE id = ?")
138 .bind(&[id.into()])?
139 .first::<PlanRow>(None)
140 .await
141 }
142
143 /// Records an outcome to plan for, and returns what a sandbox needs to
144 /// plan it. Members of the repository's workspace only: a plan becomes
145 /// issues and agents at work, which the workspace pays for.
146 pub(crate) async fn start_plan(&self, a: StartPlanArgs) -> Result<Outcome<PlanJob>> {
147 let repo = match self.repo(&a.repo, &Some(a.actor.clone())).await? {
148 Outcome::Ok(repo) => repo,
149 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
150 };
151 if !a.actor.verified || !a.actor.is_member(&repo.namespace) {
152 return Ok(members_only());
153 }
154 let brief: String = a.brief.trim().chars().take(MAX_BRIEF_CHARS).collect();
155 if brief.is_empty() {
156 return Ok(Outcome::fail(
157 FailureCode::Invalid,
158 "Say what you want to be true when the work is done.",
159 ));
160 }
161 let now = now_ms();
162 let id = new_id("pln", now);
163 let mut bytes = [0u8; 32];
164 getrandom::getrandom(&mut bytes).expect("no source of randomness");
165 let token = hex::encode(bytes);
166 self.db
167 .prepare(
168 "INSERT INTO plans
169 (id, repo_id, brief, token_hash, author_id, author_name, created_at)
170 VALUES (?, ?, ?, ?, ?, ?, ?)",
171 )
172 .bind(&[
173 id.as_str().into(),
174 repo.id.as_str().into(),
175 brief.as_str().into(),
176 hash(&token).into(),
177 a.actor.id.as_str().into(),
178 a.actor.username.as_str().into(),
179 rfc3339(now).into(),
180 ])?
181 .run()
182 .await?;
183 Ok(Outcome::Ok(PlanJob {
184 plan_id: id,
185 token,
186 brief,
187 repo: RepoPath {
188 namespace: repo.namespace,
189 name: repo.name,
190 },
191 }))
192 }
193
194 /// Records the plan a sandbox's agent wrote, or why it could not write
195 /// one. The plan's token is the only credential.
196 pub(crate) async fn report_plan(&self, a: ReportPlanArgs) -> Result<Outcome<bool>> {
197 let row = self
198 .plan_row(&a.plan_id)
199 .await?
200 .filter(|row| row.token_hash == hash(&a.token));
201 let Some(row) = row else {
202 return Ok(Outcome::fail(FailureCode::NotFound, "Plan not found."));
203 };
204 if row.status != PlanStatus::Planning {
205 return Ok(Outcome::fail(
206 FailureCode::Conflict,
207 "This plan has already been reported.",
208 ));
209 }
210 let issues = tidy(a.issues);
211 let error = a.error.or_else(|| {
212 issues
213 .is_empty()
214 .then(|| "The planner proposed no issues.".to_owned())
215 });
216 self.db
217 .prepare(
218 "UPDATE plans SET status = ?, summary = ?, issues = ?, error = ?, finished_at = ?
219 WHERE id = ? AND status = 'planning'",
220 )
221 .bind(&[
222 if error.is_some() { "failed" } else { "ready" }.into(),
223 a.summary.trim().into(),
224 serde_json::to_string(&issues)?.into(),
225 error.as_deref().map_or(JsValue::NULL, JsValue::from),
226 rfc3339(now_ms()).into(),
227 row.id.into(),
228 ])?
229 .run()
230 .await?;
231 Ok(Outcome::Ok(true))
232 }
233
234 pub(crate) async fn get_plan(&self, a: PlanArgs) -> Result<Outcome<Plan>> {
235 let repo = match self.repo(&a.repo, &a.viewer).await? {
236 Outcome::Ok(repo) => repo,
237 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
238 };
239 if !a
240 .viewer
241 .is_some_and(|viewer| viewer.is_member(&repo.namespace))
242 {
243 return Ok(members_only());
244 }
245 Ok(
246 match self
247 .plan_row(&a.id)
248 .await?
249 .filter(|row| row.repo_id == repo.id)
250 {
251 Some(row) => Outcome::Ok(row.into()),
252 None => Outcome::fail(FailureCode::NotFound, "Plan not found."),
253 },
254 )
255 }
256
257 pub(crate) async fn list_plans(&self, a: ViewArgs) -> Result<Outcome<Vec<Plan>>> {
258 let repo = match self.repo(&a.repo, &a.viewer).await? {
259 Outcome::Ok(repo) => repo,
260 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
261 };
262 if !a
263 .viewer
264 .is_some_and(|viewer| viewer.is_member(&repo.namespace))
265 {
266 return Ok(members_only());
267 }
268 let rows = self
269 .db
270 .prepare("SELECT * FROM plans WHERE repo_id = ? ORDER BY id DESC LIMIT ?")
271 .bind(&[repo.id.into(), PLAN_PAGE.into()])?
272 .all()
273 .await?
274 .results::<PlanRow>()?;
275 Ok(Outcome::Ok(rows.into_iter().map(Plan::from).collect()))
276 }
277
278 /// Opens a plan's issues. Each depends on the issues the plan said it
279 /// does, by their new numbers. With `assign`, every one is queued for a
280 /// g1t agent: those that depend on nothing are ready at once, and the
281 /// rest as what they depend on merges.
282 pub(crate) async fn apply_plan(&self, a: ApplyPlanArgs) -> Result<Outcome<Plan>> {
283 let repo = match self.repo(&a.repo, &Some(a.actor.clone())).await? {
284 Outcome::Ok(repo) => repo,
285 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
286 };
287 if !a.actor.verified || !a.actor.is_member(&repo.namespace) {
288 return Ok(members_only());
289 }
290 let Some(row) = self
291 .plan_row(&a.id)
292 .await?
293 .filter(|row| row.repo_id == repo.id)
294 else {
295 return Ok(Outcome::fail(FailureCode::NotFound, "Plan not found."));
296 };
297 // Only whoever flips it from ready to applied opens the issues.
298 let claimed = self
299 .db
300 .prepare(
301 "UPDATE plans SET status = 'applied' WHERE id = ? AND status = 'ready'
302 RETURNING id AS value",
303 )
304 .bind(&[row.id.as_str().into()])?
305 .first::<crate::rows::ValueRow>(None)
306 .await?;
307 if claimed.is_none() {
308 return Ok(Outcome::fail(
309 FailureCode::Conflict,
310 "This plan is not waiting to be applied.",
311 ));
312 }
313
314 let mut plan: Plan = row.into();
315 // What the person kept, in the plan's order. A dependency on an
316 // issue they dropped is dropped with it.
317 let kept: Vec<usize> = match &a.keep {
318 Some(positions) => (0..plan.issues.len())
319 .filter(|index| positions.contains(&(*index as u32 + 1)))
320 .collect(),
321 None => (0..plan.issues.len()).collect(),
322 };
323 let queued_by = a
324 .assign
325 .then(|| serde_json::to_string(&a.actor))
326 .transpose()?;
327 let mut numbers: Vec<Option<u32>> = vec![None; plan.issues.len()];
328 for index in kept {
329 let planned = &plan.issues[index];
330 let blocked_by: Vec<u32> = planned
331 .depends_on
332 .iter()
333 .filter_map(|position| numbers.get(*position as usize - 1).copied().flatten())
334 .collect();
335 let number = self.next_number(&repo.id).await?;
336 let now = now_ms();
337 let timestamp = rfc3339(now);
338 let issue_id = new_id("iss", now);
339 self.db
340 .prepare(
341 "INSERT INTO issues
342 (id, repo_id, number, title, body, labels, checks, blocked_by, queued_by,
343 author_id, author_name, created_at, updated_at)
344 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
345 )
346 .bind(&[
347 issue_id.as_str().into(),
348 repo.id.as_str().into(),
349 number.into(),
350 planned.title.as_str().into(),
351 planned.body.as_str().into(),
352 serde_json::to_string(&planned.labels)?.into(),
353 serde_json::to_string(&planned.checks)?.into(),
354 serde_json::to_string(&blocked_by)?.into(),
355 queued_by.as_deref().map_or(JsValue::NULL, JsValue::from),
356 a.actor.id.as_str().into(),
357 a.actor.username.as_str().into(),
358 timestamp.as_str().into(),
359 timestamp.as_str().into(),
360 ])?
361 .run()
362 .await?;
363 if a.assign {
364 self.note(
365 &repo.id,
366 number,
367 (&a.actor.id, &a.actor.username),
368 &if blocked_by.is_empty() {
369 "queued this for g1t-agent".to_owned()
370 } else {
371 format!(
372 "queued this for g1t-agent, to start once {} {} merged",
373 blocked_by
374 .iter()
375 .map(|number| format!("#{number}"))
376 .collect::<Vec<_>>()
377 .join(", "),
378 if blocked_by.len() == 1 { "has" } else { "have" }
379 )
380 },
381 )
382 .await?;
383 }
384 self.publish(
385 "issue.opened",
386 &repo.id,
387 &a.actor,
388 IssueEvent {
389 issue_id,
390 repo_id: repo.id.clone(),
391 number,
392 title: Some(planned.title.clone()),
393 ..IssueEvent::default()
394 },
395 )
396 .await?;
397 numbers[index] = Some(number);
398 }
399 for (issue, number) in plan.issues.iter_mut().zip(&numbers) {
400 issue.number = *number;
401 }
402 self.db
403 .prepare("UPDATE plans SET issues = ? WHERE id = ?")
404 .bind(&[
405 serde_json::to_string(&plan.issues)?.into(),
406 plan.id.as_str().into(),
407 ])?
408 .run()
409 .await?;
410 plan.status = PlanStatus::Applied;
411 Ok(Outcome::Ok(plan))
412 }
413
414 /// Queues an issue for a g1t agent, or takes it out of the queue.
415 pub(crate) async fn queue_issue(&self, a: QueueIssueArgs) -> Result<Outcome<bool>> {
416 let issue = match self.manageable_issue(&a.actor, &a.repo, a.number).await? {
417 Outcome::Ok(issue) => issue,
418 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
419 };
420 let queued_by = a
421 .queued
422 .then(|| serde_json::to_string(&a.actor))
423 .transpose()?;
424 self.db
425 .prepare("UPDATE issues SET queued_by = ? WHERE id = ?")
426 .bind(&[
427 queued_by.as_deref().map_or(JsValue::NULL, JsValue::from),
428 issue.id.as_str().into(),
429 ])?
430 .run()
431 .await?;
432 Ok(Outcome::Ok(true))
433 }
434
435 /// Takes the issues that are waiting for a g1t agent and can be given
436 /// one now: open, queued, with nobody working on them, and with
437 /// everything they depend on closed. Oldest first, in one repository or
438 /// in all, and no more than leaves each repository with
439 /// `MAX_AGENTS_AT_WORK` agents making changes at once.
440 ///
441 /// Taking an issue takes it out of the queue, in one statement, so two
442 /// callers cannot both start an agent on it. A caller that then cannot
443 /// start one puts it back with `queue_issue`.
444 pub(crate) async fn ready_issues(&self, a: ReadyIssuesArgs) -> Result<Vec<ReadyIssue>> {
445 let rows = self
446 .db
447 .prepare(
448 "SELECT repo_id, number, queued_by FROM issues
449 WHERE state = 'open' AND queued_by IS NOT NULL
450 AND (?1 IS NULL OR repo_id = ?1)
451 AND NOT EXISTS (
452 SELECT 1 FROM pulls
453 WHERE pulls.issue_id = issues.id AND pulls.status IN ('draft', 'open'))
454 AND NOT EXISTS (
455 SELECT 1 FROM json_each(issues.blocked_by) AS blocker
456 JOIN issues AS earlier
457 ON earlier.repo_id = issues.repo_id AND earlier.number = blocker.value
458 WHERE earlier.state = 'open')
459 ORDER BY number LIMIT 50",
460 )
461 .bind(&[a.repo_id.map_or(JsValue::NULL, JsValue::from)])?
462 .all()
463 .await?
464 .results::<QueuedRow>()?;
465 let mut ready = Vec::new();
466 let mut room: HashMap<String, usize> = HashMap::new();
467 for row in rows {
468 let Ok(actor) = serde_json::from_str::<User>(&row.queued_by) else {
469 continue;
470 };
471 // How many more agents this repository has room for.
472 if !room.contains_key(&row.repo_id) {
473 let working = self
474 .db
475 .prepare(
476 "SELECT count(*) AS n FROM pulls
477 WHERE repo_id = ? AND managed = 1 AND status = 'draft'",
478 )
479 .bind(&[row.repo_id.as_str().into()])?
480 .first::<crate::rows::NumberRow>(None)
481 .await?
482 .map_or(0, |row| row.n as usize);
483 room.insert(
484 row.repo_id.clone(),
485 MAX_AGENTS_AT_WORK.saturating_sub(working),
486 );
487 }
488 let left = room.get_mut(&row.repo_id).expect("just inserted");
489 if *left == 0 {
490 continue;
491 }
492 let taken = self
493 .db
494 .prepare(
495 "UPDATE issues SET queued_by = NULL
496 WHERE repo_id = ? AND number = ? AND queued_by IS NOT NULL
497 RETURNING id AS value",
498 )
499 .bind(&[row.repo_id.as_str().into(), row.number.into()])?
500 .first::<crate::rows::ValueRow>(None)
501 .await?;
502 if taken.is_none() {
503 continue;
504 }
505 *left -= 1;
506 // Whoever queued it could see the repository then; where it is
507 // now is asked as them.
508 let repo: Outcome<Repo> = g1t_kit::call(
509 &self.repos,
510 "get_by_id",
511 &GetByIdArgs {
512 id: row.repo_id.clone(),
513 viewer: Some(actor.clone()),
514 },
515 )
516 .await?;
517 if let Outcome::Ok(repo) = repo {
518 ready.push(ReadyIssue {
519 repo: RepoPath {
520 namespace: repo.namespace,
521 name: repo.name,
522 },
523 number: row.number,
524 actor,
525 });
526 }
527 }
528 Ok(ready)
529 }
530}
531
532fn members_only<T>() -> Outcome<T> {
533 Outcome::fail(
534 FailureCode::Forbidden,
535 "Only members of the repository's workspace can plan work for it.",
536 )
537}
538
539#[cfg(test)]
540mod tests {
541 use super::*;
542
543 fn proposed(title: &str, depends_on: &[u32]) -> PlannedIssue {
544 PlannedIssue {
545 title: title.to_owned(),
546 body: " What to do. ".to_owned(),
547 labels: vec!["Feature".to_owned()],
548 checks: vec![" cargo test ".to_owned(), String::new()],
549 files: vec!["src/lib.rs".to_owned()],
550 depends_on: depends_on.to_vec(),
551 number: None,
552 }
553 }
554
555 #[test]
556 fn a_proposal_is_trimmed_and_normalised() {
557 let issues = tidy(vec![proposed(" Add a flag ", &[])]);
558 assert_eq!(issues[0].title, "Add a flag");
559 assert_eq!(issues[0].body, "What to do.");
560 assert_eq!(issues[0].labels, ["feature"]);
561 assert_eq!(issues[0].checks, ["cargo test"]);
562 }
563
564 #[test]
565 fn an_issue_can_only_depend_on_earlier_ones() {
566 let issues = tidy(vec![
567 proposed("First", &[2]),
568 proposed("Second", &[1, 1, 2, 9]),
569 proposed("Third", &[2, 1]),
570 ]);
571 assert!(issues[0].depends_on.is_empty());
572 assert_eq!(issues[1].depends_on, [1]);
573 assert_eq!(issues[2].depends_on, [1, 2]);
574 }
575
576 #[test]
577 fn a_plan_is_bounded() {
578 let many = (0..30)
579 .map(|i| proposed(&format!("Issue {i}"), &[]))
580 .collect();
581 assert_eq!(tidy(many).len(), MAX_PLANNED_ISSUES);
582 }
583}