Agents asked while not at work are woken to answer
Agents could ask the agent on another pull request a question or hand it work, but only an agent in the middle of its run ever read one. Once a pull request's change was done and waiting for review or a merge, a question to it waited forever: the only two ever asked on hello ("which name should I use, greeting() or greet()?") were never read, and the outcome page showed them "Waiting to be read". Now, when an agent asks one that is not at work, the work service publishes agent.asked and the runner wakes it: wake_for_messages claims a short "answer" step on its pull request (nothing else starts on it meanwhile), hands over what it was sent, marked read, and a run in that pull request's sandbox puts its own change in front of it with the questions. It reads its code, answers with answer_message, and for a handoff it takes on, commits the work. Answering the last question lets go of the step. The runner gains an `answer` mode: it does not merge in the default branch first (that would push a commit for a question) and finishing without a change is the usual ending, not a failure. The lifecycle shows the step as "Answering an agent". Verified on hello: #89 (--both) asked #87 (renaming hail/part) for the exact names and signatures; #87's agent was woken, read src/lib.rs and answered "pub fn greet(name: &str) -> String, pub fn farewell(name: &str) -> String; hail and part no longer exist" twenty seconds later, and the answer reached #89's agent at its next step. $0.04. agent.asked is in the activity feed and can be chosen for webhooks. The plan's first build item (handoffs as states) was already done and is marked so.
| 113 | 113 | "was asked a question by the agent on #41" and "answered the question from | |
| 114 | 114 | the agent on #41". | |
| 115 | 115 | ||
| 116 | − | If the agent asked is not at work, it will not answer soon. The response to | |
| 117 | − | `message_agent` says so, in `hint`, and points the asking agent at the | |
| 118 | − | other pull request's change to read with `get_pull_request` and | |
| 119 | − | `get_pull_request_changes` instead. | |
| 116 | + | If the agent asked is not at work, because its change is done and waiting | |
| 117 | + | for review or a merge, g1t wakes it to answer. It starts a short run in that | |
| 118 | + | pull request's sandbox with the agent's own change in front of it and what | |
| 119 | + | it was asked; the agent reads its code, answers with `answer_message`, and, | |
| 120 | + | for a handoff it takes on, commits the work. Its pull request is noted "g1t | |
| 121 | + | woke g1t-agent to answer the agent on #41", and nothing else starts on it | |
| 122 | + | until everything it was asked is answered, or 20 minutes pass. The response | |
| 123 | + | to `message_agent` says so in `hint`, and points the asking agent at the | |
| 124 | + | other pull request's change to read meanwhile with `get_pull_request` and | |
| 125 | + | `get_pull_request_changes`. | |
| 126 | + | ||
| 127 | + | An agent g1t has stopped on (its pull request needs a person) is not woken; | |
| 128 | + | the hint then says it will not answer soon. | |
| 120 | 129 | ||
| 121 | 130 | On an [outcome's page](/guides/outcomes/#agents-talking), **Agents talking** | |
| 122 | 131 | lists every exchange between its agents with where it stands: waiting to be |
| 70 | 70 | | `workflow.completed` | A GitHub Actions run finished. `data.workflow`, `data.conclusion`, `data.runId`, `data.sha`, `data.pull`. | | |
| 71 | 71 | | `queue.changed` | The merge queue gained, lost or settled an entry. | | |
| 72 | 72 | | `session.appended` | An agent's session grew. Busy: choose it only if you need it. | | |
| 73 | + | | `agent.asked` | An agent asked the agent on another pull request a question, or handed it work, while that one was not at work; g1t wakes it to answer. | | |
| 73 | 74 | ||
| 74 | 75 | ## Check the signature | |
| 75 | 76 |
| 4 | 4 | CircleSlash, | |
| 5 | 5 | GitMerge, | |
| 6 | 6 | GitPullRequest, | |
| 7 | + | MessageCircleQuestion, | |
| 7 | 8 | MessageSquare, | |
| 8 | 9 | Play, | |
| 9 | 10 | Terminal, | |
| 94 | 95 | }; | |
| 95 | 96 | case "pull.merged": | |
| 96 | 97 | return { icon: <GitMerge size={14} />, tone: "text-accent", actor, text: <>landed {ref(event.data.number)} on main</> }; | |
| 98 | + | case "agent.asked": | |
| 99 | + | return { icon: <MessageCircleQuestion size={14} />, tone: "text-merged", actor, text: <>asked the agent on {ref(event.data.number)}, and g1t woke it to answer</> }; | |
| 97 | 100 | case "pull.merge_requested": | |
| 98 | 101 | return { icon: <GitMerge size={14} />, tone: "text-muted", actor, text: <>is bringing {ref(event.data.number)} up to date</> }; | |
| 99 | 102 | default: |
| 11 | 11 | checking: 1, | |
| 12 | 12 | reviewing: 2, | |
| 13 | 13 | catching_up: 3, | |
| 14 | + | // Its change is made; it is answering another agent. | |
| 15 | + | answering: 2, | |
| 14 | 16 | ready: 4, | |
| 15 | 17 | queued: 4, | |
| 16 | 18 | }; | |
| 22 | 24 | reviewing: "In review", | |
| 23 | 25 | revising: "Revising", | |
| 24 | 26 | catching_up: "Catching up", | |
| 27 | + | answering: "Answering an agent", | |
| 25 | 28 | queued: "In the merge queue", | |
| 26 | 29 | ready: "Ready to merge", | |
| 27 | 30 | needs_you: "Needs you", | |
| 55 | 58 | reviewing: "g1t is seeing this through", | |
| 56 | 59 | revising: "g1t is seeing this through", | |
| 57 | 60 | catching_up: "g1t is seeing this through", | |
| 61 | + | answering: "g1t is seeing this through", | |
| 58 | 62 | queued: "In the merge queue", | |
| 59 | 63 | ready: "Ready to merge", | |
| 60 | 64 | needs_you: "Needs you", |
| 52 | 52 | bar: "bg-merged", | |
| 53 | 53 | live: true, | |
| 54 | 54 | }, | |
| 55 | + | answering: { | |
| 56 | + | label: "Answering an agent", | |
| 57 | + | icon: <Sparkles size={13} />, | |
| 58 | + | ring: "ring-merged/60", | |
| 59 | + | text: "text-merged", | |
| 60 | + | bar: "bg-merged", | |
| 61 | + | live: true, | |
| 62 | + | }, | |
| 55 | 63 | catching_up: { | |
| 56 | 64 | label: "Catching up", | |
| 57 | 65 | icon: <Loader2 size={13} className="animate-spin" />, | |
| 214 | 222 | const live = plan.progress.filter((item) => look(item.state).live).length; | |
| 215 | 223 | const needsYou = count(["needs_you"]); | |
| 216 | 224 | const blocked = count(["blocked", "waiting", "open"]); | |
| 217 | − | const order = ["landed", "queued", "ready", "reviewing", "checking", "working", "revising", "catching_up", "needs_you", "waiting", "blocked", "open", "closed"]; | |
| 225 | + | const order = ["landed", "queued", "ready", "reviewing", "checking", "working", "revising", "catching_up", "answering", "needs_you", "waiting", "blocked", "open", "closed"]; | |
| 218 | 226 | const sorted = [...plan.progress].sort((a, b) => order.indexOf(a.state) - order.indexOf(b.state)); | |
| 219 | 227 | ||
| 220 | 228 | return ( |
| 23 | 23 | events: ["pull.opened", "pull.ready", "pull.updated", "pull.merge_requested", "pull.merged", "pull.closed"], | |
| 24 | 24 | }, | |
| 25 | 25 | { title: "Checks, reviews and the queue", events: ["checks.completed", "workflow.completed", "review.completed", "queue.changed"] }, | |
| 26 | − | { title: "Agents", events: ["session.appended"] }, | |
| 26 | + | { title: "Agents", events: ["session.appended", "agent.asked"] }, | |
| 27 | 27 | ]; | |
| 28 | 28 | ||
| 29 | 29 | function StatusDot({ status }: { status: string | null }) { |
| 437 | 437 | {/* Steering: while its agent works, people can tell it things. */} | |
| 438 | 438 | {canManage && | |
| 439 | 439 | pull.runtime === "hosted" && | |
| 440 | − | (working || ["working", "revising", "catching_up"].includes(lifecycle?.stage ?? "")) && ( | |
| 440 | + | (working || ["working", "revising", "catching_up", "answering"].includes(lifecycle?.stage ?? "")) && ( | |
| 441 | 441 | <Form method="post" className="mt-4 rounded-2xl bg-surface p-4 ring-1 ring-merged/30"> | |
| 442 | 442 | <p className="flex items-center gap-2 text-sm font-medium"> | |
| 443 | 443 | <Sparkles size={15} className="text-merged" /> |
| 176 | 176 | pull request a question, or hands it work, with `message_agent` (`kind` | |
| 177 | 177 | `question` or `handoff`, and `from_number`, its own pull request). The | |
| 178 | 178 | other agent replies with `answer_message` | |
| 179 | − | (`POST {repo}/messages/{id}/answer`). Plans show these exchanges under | |
| 179 | + | (`POST {repo}/messages/{id}/answer`); one that is not at work is woken | |
| 180 | + | to answer, in its own pull request's sandbox. Plans show these exchanges under | |
| 180 | 181 | "Agents talking". From any other caller, `message_agent` sends a plain | |
| 181 | 182 | message. | |
| 182 | 183 |
| 15 | 15 | use crate::{User, Viewer}; | |
| 16 | 16 | ||
| 17 | 17 | /// Every event a webhook can be sent, in the order people are shown them. | |
| 18 | − | pub const EVENT_TYPES: [&str; 20] = [ | |
| 18 | + | pub const EVENT_TYPES: [&str; 21] = [ | |
| 19 | 19 | "git.push", | |
| 20 | 20 | "repo.created", | |
| 21 | 21 | "repo.forked", | |
| 31 | 31 | "pull.merge_requested", | |
| 32 | 32 | "pull.merged", | |
| 33 | 33 | "pull.closed", | |
| 34 | + | "agent.asked", | |
| 34 | 35 | "checks.completed", | |
| 35 | 36 | "review.completed", | |
| 36 | 37 | "workflow.completed", |
| 542 | 542 | Revising, | |
| 543 | 543 | /// The agent is merging in the branch it would land on, which moved. | |
| 544 | 544 | CatchingUp, | |
| 545 | + | /// Woken to answer a question another agent asked it, or a handoff. | |
| 546 | + | Answering, | |
| 545 | 547 | /// In the repository's merge queue, being tested with what is ahead of | |
| 546 | 548 | /// it before it lands. | |
| 547 | 549 | Queued, | |
| 683 | 685 | pub pull_id: String, | |
| 684 | 686 | } | |
| 685 | 687 | ||
| 688 | + | /// `wake_for_messages`: the agent on a pull request was asked a question | |
| 689 | + | /// or handed work while it was not at work. Claims a short step for it to | |
| 690 | + | /// answer, and hands over what it was sent, marked read. Null when there | |
| 691 | + | /// is nothing waiting, or the pull request cannot take a step now. | |
| 692 | + | /// Returns `Option<Wake>`. | |
| 693 | + | #[derive(Debug, Serialize, Deserialize)] | |
| 694 | + | #[serde(rename_all = "camelCase")] | |
| 695 | + | pub struct WakeForMessagesArgs { | |
| 696 | + | pub pull_id: String, | |
| 697 | + | } | |
| 698 | + | ||
| 699 | + | /// What an agent woken to answer needs: its pull request, and what it was | |
| 700 | + | /// sent, oldest first. | |
| 701 | + | #[derive(Debug, Serialize, Deserialize)] | |
| 702 | + | #[serde(rename_all = "camelCase")] | |
| 703 | + | pub struct Wake { | |
| 704 | + | pub job: LifecycleJob, | |
| 705 | + | pub messages: Vec<AgentMessage>, | |
| 706 | + | } | |
| 707 | + | ||
| 686 | 708 | /// `stall`: records that a step could not be carried out, so that g1t | |
| 687 | 709 | /// stops and a person is asked. Returns `bool`. | |
| 688 | 710 | #[derive(Debug, Serialize, Deserialize)] | |
| 924 | 946 | /// `blocked` (waiting on issues it depends on), `waiting` (for an | |
| 925 | 947 | /// agent), `open` (nobody on it), one of the lifecycle stages | |
| 926 | 948 | /// (`working`, `checking`, `reviewing`, `revising`, `catching_up`, | |
| 927 | − | /// `queued`, `ready`, `needs_you`), `landed` or `closed`. | |
| 949 | + | /// `answering`, `queued`, `ready`, `needs_you`), `landed` or `closed`. | |
| 928 | 950 | pub state: String, | |
| 929 | 951 | /// One sentence about where it stands. | |
| 930 | 952 | pub detail: String, |
| 152 | 152 | } | |
| 153 | 153 | let head = git(workdir, &["rev-parse", "HEAD"])?; | |
| 154 | 154 | if head == start { | |
| 155 | + | // An agent woken to answer usually only answers. | |
| 156 | + | if std::env::var("MODE").as_deref() == Ok("answer") { | |
| 157 | + | return Ok(summary); | |
| 158 | + | } | |
| 155 | 159 | bail!("the agent finished without changing anything"); | |
| 156 | 160 | } | |
| 157 | 161 | git( | |
| 181 | 185 | Ok("update") => std::process::exit(update::main()), | |
| 182 | 186 | Ok("review") => std::process::exit(review::main()), | |
| 183 | 187 | Ok("revise") => std::process::exit(revise::main()), | |
| 188 | + | Ok("answer") => std::process::exit(revise::answer()), | |
| 184 | 189 | Ok("plan") => std::process::exit(plan::main()), | |
| 185 | 190 | Ok("queue") => std::process::exit(queue::main()), | |
| 186 | 191 | Ok("steer") => std::process::exit(steer::main()), |
| 9 | 9 | ||
| 10 | 10 | use crate::report::{Entry, Reporter}; | |
| 11 | 11 | ||
| 12 | + | /// The author woken to answer what other agents asked it while it was not | |
| 13 | + | /// at work: the same run, which pushes only if it took on handed-over work. | |
| 14 | + | pub fn answer() -> i32 { | |
| 15 | + | finish(crate::run, "Answering failed") | |
| 16 | + | } | |
| 17 | + | ||
| 12 | 18 | pub fn main() -> i32 { | |
| 19 | + | finish(crate::run, "The revision failed") | |
| 20 | + | } | |
| 21 | + | ||
| 22 | + | fn finish(run: fn(&mut Reporter) -> anyhow::Result<String>, failed: &str) -> i32 { | |
| 13 | 23 | let mut reporter = match Reporter::from_env() { | |
| 14 | 24 | Ok(reporter) => reporter, | |
| 15 | 25 | Err(error) => { | |
| 17 | 27 | return 2; | |
| 18 | 28 | } | |
| 19 | 29 | }; | |
| 20 | − | let outcome = crate::run(&mut reporter); | |
| 30 | + | let outcome = run(&mut reporter); | |
| 21 | 31 | // On success the agent's account is already in the session: the harness | |
| 22 | 32 | // records its messages as they arrive. | |
| 23 | 33 | if let Err(error) = &outcome { | |
| 24 | 34 | eprintln!("g1t-runner: {error:#}"); | |
| 25 | 35 | reporter.record(Entry::new( | |
| 26 | 36 | "note", | |
| 27 | − | &format!("The revision failed: {error:#}. Nothing was pushed."), | |
| 37 | + | &format!("{failed}: {error:#}. Nothing was pushed."), | |
| 28 | 38 | )); | |
| 29 | 39 | } | |
| 30 | 40 | reporter.flush(); |
| 785 | 785 | set number of agents on one issue is dropped: choosing how many agents to | |
| 786 | 786 | use is not something people should have to do. | |
| 787 | 787 | ||
| 788 | − | 1. **Handoffs and questions between agents** as states (offered, accepted, | |
| 789 | − | declined, done) on the outcome page, not only comments. | |
| 788 | + | 1. ~~**Handoffs and questions between agents** as states on the outcome | |
| 789 | + | page.~~ Done: questions and handoffs show as waiting, read, answered, | |
| 790 | + | taken on or declined; since 2026-10-03 an agent asked while it is not at | |
| 791 | + | work is woken to answer, where before the question waited forever. | |
| 790 | 792 | 2. **The large run.** Dozens of agents on a real repository, end to end, for | |
| 791 | 793 | the video; g1t hosted on g1t. | |
| 792 | 794 | 3. **Polish for judges trying it in a minute:** a seeded demo workspace, the |
| 143 | 143 | removeFromQueue: (actor, repo, number) => call("remove_from_queue", { actor, repo, number }), | |
| 144 | 144 | messageAgent: (actor, repo, number, body) => call("message_agent", { actor, repo, number, body }), | |
| 145 | 145 | catchUpJob: (pullId) => call("catch_up_job", { pullId }), | |
| 146 | + | wakeForMessages: (pullId) => call("wake_for_messages", { pullId }), | |
| 146 | 147 | getSettings: (repo, viewer) => call("get_settings", { repo, viewer }), | |
| 147 | 148 | updateSettings: (actor, repo, settings) => | |
| 148 | 149 | call("update_settings", { actor, repo, settings }), |
| 37 | 37 | /** A merge was asked for while the pull request was behind; it has to catch up first. */ | |
| 38 | 38 | "pull.merge_requested": { pullId: string; repoId: string; number: number; issue?: number }; | |
| 39 | 39 | "pull.closed": { pullId: string; repoId: string; number: number; issue?: number }; | |
| 40 | + | /** Another agent asked the agent on a pull request, which was not at work, a question or handed it work. */ | |
| 41 | + | "agent.asked": { pullId: string; repoId: string; number: number; issue?: number }; | |
| 40 | 42 | "pull.merged": { pullId: string; repoId: string; number: number; issue?: number; commit: string }; | |
| 41 | 43 | /** A run of the acceptance checks finished. `commit` is what was checked. */ | |
| 42 | 44 | "checks.completed": { |
| 24 | 24 | "pull.merge_requested", | |
| 25 | 25 | "pull.merged", | |
| 26 | 26 | "pull.closed", | |
| 27 | + | "agent.asked", | |
| 27 | 28 | "checks.completed", | |
| 28 | 29 | "review.completed", | |
| 29 | 30 | "workflow.completed", |
| 306 | 306 | * to merge. g1t takes each one without being asked: `working` (the agent is | |
| 307 | 307 | * making the change), `checking`, `reviewing`, `revising` (the agent is | |
| 308 | 308 | * addressing failed checks or a review), `catching_up` (merging in the | |
| 309 | − | * branch it would land on), then `ready` for a person to merge. `needs_you` | |
| 309 | + | * branch it would land on), `answering` (woken to answer another agent), | |
| 310 | + | * then `ready` for a person to merge. `needs_you` | |
| 310 | 311 | * means g1t has stopped and a person decides what happens next. | |
| 311 | 312 | */ | |
| 312 | 313 | export type Stage = | |
| 315 | 316 | | "reviewing" | |
| 316 | 317 | | "revising" | |
| 317 | 318 | | "catching_up" | |
| 319 | + | | "answering" | |
| 318 | 320 | | "queued" | |
| 319 | 321 | | "ready" | |
| 320 | 322 | | "needs_you"; | |
| 348 | 350 | round: number; | |
| 349 | 351 | }; | |
| 350 | 352 | ||
| 353 | + | /** An agent woken to answer what other agents sent it while it was not at work. */ | |
| 354 | + | export type Wake = { job: LifecycleJob; messages: AgentMessage[] }; | |
| 355 | + | ||
| 351 | 356 | /** | |
| 352 | 357 | * `planning` while an agent reads the repository and writes it; `ready` for | |
| 353 | 358 | * a person to read and apply; `failed` if it could not be written; | |
| 589 | 594 | * merge of it was asked for. Null if none was. | |
| 590 | 595 | */ | |
| 591 | 596 | catchUpJob(pullId: string): Promise<LifecycleJob | null>; | |
| 597 | + | /** | |
| 598 | + | * Claims a short step for the agent on a pull request to answer the | |
| 599 | + | * questions and handoffs it was sent while not at work, and hands them | |
| 600 | + | * over, marked read. Null when there is nothing waiting or it cannot | |
| 601 | + | * take a step now. | |
| 602 | + | */ | |
| 603 | + | wakeForMessages(pullId: string): Promise<Wake | null>; | |
| 592 | 604 | ||
| 593 | 605 | /** | |
| 594 | 606 | * Opens a pull request: a draft with a fork to push to, or, given a |
| 2 | 2 | import { WorkerEntrypoint } from "cloudflare:workers"; | |
| 3 | 3 | ||
| 4 | 4 | import { | |
| 5 | + | type AgentMessage, | |
| 5 | 6 | type CheckJob, | |
| 6 | 7 | type G1tEvent, | |
| 7 | 8 | type Issue, | |
| 97 | 98 | | { kind: "update"; pullId?: string } | |
| 98 | 99 | /** The author sent back to address failed checks or a review. */ | |
| 99 | 100 | | { kind: "revise"; pullId: string } | |
| 101 | + | /** The author woken to answer other agents; nothing to undo if it fails. */ | |
| 102 | + | | { kind: "answer"; pullId: string } | |
| 100 | 103 | /** An agent turning an outcome into a plan. */ | |
| 101 | 104 | | { kind: "plan"; planId: string; token: string } | |
| 102 | 105 | /** One combined state of a merge queue, being built and checked. */ | |
| 161 | 164 | await work.failPlan(run.planId, run.token, "The sandbox stopped before the plan was written."); | |
| 162 | 165 | return; | |
| 163 | 166 | } | |
| 167 | + | // An answer that never came: the claim lapses and the asker reads the | |
| 168 | + | // change instead, as it was told it could. | |
| 169 | + | if (run.kind === "answer") return; | |
| 164 | 170 | if (run.kind === "update" || run.kind === "revise") { | |
| 165 | 171 | if (run.pullId) { | |
| 166 | 172 | await work.stall( | |
| 330 | 336 | return parts.filter(Boolean).join("\n\n"); | |
| 331 | 337 | } | |
| 332 | 338 | ||
| 339 | + | /** | |
| 340 | + | * What the agent on a pull request is told when g1t wakes it to answer the | |
| 341 | + | * questions and handoffs other agents sent while it was not at work. | |
| 342 | + | */ | |
| 343 | + | function buildAnswerPrompt(job: LifecycleJob, messages: AgentMessage[], inFlight: string | null): string { | |
| 344 | + | const asked = messages | |
| 345 | + | .filter((message) => message.kind === "question" || message.kind === "handoff") | |
| 346 | + | .map((message) => { | |
| 347 | + | const from = message.fromNumber != null ? `the agent on #${message.fromNumber}` : message.author; | |
| 348 | + | const what = message.kind === "handoff" ? "Work handed over" : "Question"; | |
| 349 | + | return `${what} from ${from} (id ${message.id}):\n${message.body}`; | |
| 350 | + | }); | |
| 351 | + | const said = messages | |
| 352 | + | .filter((message) => message.kind === "message" || message.kind === "answer") | |
| 353 | + | .map((message) => `From ${message.fromNumber != null ? `the agent on #${message.fromNumber}` : message.author}: ${message.body}`); | |
| 354 | + | const parts = [ | |
| 355 | + | `You are a coding agent working in the git repository checked out in the current directory. It holds a change you made earlier, which is open as pull request #${job.number}. Your work on it is done for now; you have been woken because other agents in this repository asked you something.`, | |
| 356 | + | job.issue | |
| 357 | + | ? `Your pull request is for issue #${job.issue.number}: ${job.issue.title}\n\n${job.issue.body}` | |
| 358 | + | : `Your pull request: ${job.title}`, | |
| 359 | + | job.description && `What you said you changed:\n\n${job.description}`, | |
| 360 | + | asked.join("\n\n"), | |
| 361 | + | said.length > 0 && `Also sent to you:\n\n${said.join("\n\n")}`, | |
| 362 | + | inFlight, | |
| 363 | + | WORKING_WITH_OTHERS, | |
| 364 | + | "Answer each question and handoff above with answer_message and its id, from what your change actually does: read your own code and history (git log, git diff against the default branch) before you answer, and be specific, with names, signatures and files. For a handoff, take it on only if the work belongs in your pull request; then make the change, commit it with a clear message, and answer saying what you did. Otherwise answer with decline set and say where it belongs. Do not push; that is done for you. Change nothing else. Finish with one or two plain sentences on what you answered.", | |
| 365 | + | ]; | |
| 366 | + | return parts.filter(Boolean).join("\n\n"); | |
| 367 | + | } | |
| 368 | + | ||
| 333 | 369 | function buildPrompt( | |
| 334 | 370 | issue: Issue, | |
| 335 | 371 | instructions: string, | |
| 654 | 690 | case "pull.merged": | |
| 655 | 691 | await this.advanceAll(event.data.repoId); | |
| 656 | 692 | break; | |
| 693 | + | // Another agent asked one that is not at work: wake it to answer. | |
| 694 | + | case "agent.asked": | |
| 695 | + | await this.wakeForMessages(event.data.pullId); | |
| 696 | + | break; | |
| 657 | 697 | // Something an issue was waiting on has finished, or an agent has | |
| 658 | 698 | // stopped and left room for another. | |
| 659 | 699 | case "issue.closed": | |
| 858 | 898 | return token; | |
| 859 | 899 | } | |
| 860 | 900 | ||
| 901 | + | /** | |
| 902 | + | * Wakes the agent on a pull request to answer the questions and handoffs | |
| 903 | + | * other agents sent it while it was not at work. The work service claims | |
| 904 | + | * the step, so a second event starts nothing. | |
| 905 | + | */ | |
| 906 | + | private async wakeForMessages(pullId: string): Promise<void> { | |
| 907 | + | const work = workClient(this.env.WORK); | |
| 908 | + | const wake = await work.wakeForMessages(pullId); | |
| 909 | + | if (!wake) return; | |
| 910 | + | const { job, messages } = wake; | |
| 911 | + | try { | |
| 912 | + | if (!this.modelsReachable() || !(await this.workspaceAllowed(job.repo.namespace))) { | |
| 913 | + | throw new Error("g1t agents are not enabled for this workspace."); | |
| 914 | + | } | |
| 915 | + | const { token } = await identityClient(this.env.IDENTITY).createAccessToken( | |
| 916 | + | job.author, | |
| 917 | + | `g1t agent answering on ${job.repo.namespace}/${job.repo.name}#${job.number}`, | |
| 918 | + | TOKEN_TTL_SECONDS, | |
| 919 | + | ); | |
| 920 | + | const sandbox = this.env.SANDBOX.get(this.env.SANDBOX.idFromName(`answer-${job.pullId}-${messages[0]?.id ?? Date.now()}`)); | |
| 921 | + | await sandbox.run({ | |
| 922 | + | kind: "answer", | |
| 923 | + | pullId: job.pullId, | |
| 924 | + | envVars: { | |
| 925 | + | // Answered from its change as it stands: no merging in of the | |
| 926 | + | // default branch, which would push a commit for a question. | |
| 927 | + | MODE: "answer", | |
| 928 | + | G1T_API: "https://api.g1t.sh", | |
| 929 | + | G1T_TOKEN: token, | |
| 930 | + | G1T_USER: job.author.username, | |
| 931 | + | G1T_REPO: `${job.repo.namespace}/${job.repo.name}`, | |
| 932 | + | PULL_NUMBER: String(job.number), | |
| 933 | + | GIT_REMOTE: `https://g1t.sh/${job.source.namespace}/${job.source.name}.git`, | |
| 934 | + | COMMIT_MESSAGE: `Take on work handed over to #${job.number}`, | |
| 935 | + | G1T_AGENT_TOKEN: await this.agentToken(job.author, job.repo), | |
| 936 | + | PROMPT: buildAnswerPrompt(job, messages, await this.inFlight(job.author, job.repo, job.number)), | |
| 937 | + | ...(await this.modelEnvOrThrow("implement", job.repo, job.number)), | |
| 938 | + | }, | |
| 939 | + | }); | |
| 940 | + | } catch (error) { | |
| 941 | + | // Said on the pull request; the askers were told to read the change. | |
| 942 | + | await work.appendSession(job.author, job.repo, job.number, [ | |
| 943 | + | { | |
| 944 | + | kind: "note", | |
| 945 | + | text: `g1t could not wake the agent to answer: ${error instanceof Error ? error.message : String(error)}`, | |
| 946 | + | }, | |
| 947 | + | ]); | |
| 948 | + | } | |
| 949 | + | } | |
| 950 | + | ||
| 861 | 951 | private async startRevision(job: LifecycleJob): Promise<void> { | |
| 862 | 952 | const { token } = await identityClient(this.env.IDENTITY).createAccessToken( | |
| 863 | 953 | job.author, |
| 1815 | 1815 | "locate_pull" => reply(&work.locate_pull(args(body)?).await?), | |
| 1816 | 1816 | "answer_message" => reply(&work.answer_message(args(body)?).await?), | |
| 1817 | 1817 | "take_messages" => reply(&work.take_messages(args(body)?).await?), | |
| 1818 | + | "wake_for_messages" => reply(&work.wake_for_messages(args(body)?).await?), | |
| 1818 | 1819 | "catch_up_job" => reply(&work.catch_up_job(args(body)?).await?), | |
| 1819 | 1820 | "get_settings" => reply(&work.get_settings(args(body)?).await?), | |
| 1820 | 1821 | "update_settings" => reply(&work.update_settings(args(body)?).await?), |
| 212 | 212 | "The agent is merging in the branch this will land on, which has moved.", | |
| 213 | 213 | ); | |
| 214 | 214 | } | |
| 215 | + | Some("answer") => { | |
| 216 | + | return wait( | |
| 217 | + | Stage::Answering, | |
| 218 | + | "The agent is answering what another agent asked it.", | |
| 219 | + | ); | |
| 220 | + | } | |
| 215 | 221 | Some("merge") => return wait(Stage::Ready, "Merging."), | |
| 216 | 222 | Some(_) => return wait(Stage::Reviewing, "A g1t agent is reviewing the change."), | |
| 217 | 223 | None => {} | |
| 563 | 569 | ||
| 564 | 570 | /// Takes a step for a pull request, if nobody else has. One statement, | |
| 565 | 571 | /// so that two callers cannot both take it. | |
| 566 | − | async fn claim(&self, pull_id: &str, step: &str, minutes: u64, revising: bool) -> Result<bool> { | |
| 572 | + | pub(crate) async fn claim(&self, pull_id: &str, step: &str, minutes: u64, revising: bool) -> Result<bool> { | |
| 567 | 573 | let now = now_ms(); | |
| 568 | 574 | let revision = if revising { | |
| 569 | 575 | ", revisions = revisions + 1, revised_at = ?1" |
| 18 | 18 | /// Kinds an agent may send another. | |
| 19 | 19 | const ASKS: [&str; 2] = ["question", "handoff"]; | |
| 20 | 20 | ||
| 21 | + | /// How long an agent woken to answer holds its pull request: nothing else | |
| 22 | + | /// starts on it meanwhile. Answering everything lets go sooner. | |
| 23 | + | const ANSWER_MINUTES: u64 = 20; | |
| 24 | + | ||
| 21 | 25 | #[derive(Deserialize)] | |
| 22 | 26 | struct MessageRow { | |
| 23 | 27 | id: String, | |
| 196 | 200 | self.insert_message(&repo.id, &pull.id, &a.actor.id, &message).await?; | |
| 197 | 201 | let mut message = message; | |
| 198 | 202 | if from_agent && !at_work { | |
| 199 | − | message.hint = Some(format!( | |
| 200 | − | "The agent on #{} is not at work right now, so it will not answer soon. Its change is there to read: use get_pull_request and get_pull_request_changes on #{}, and decide from that.", | |
| 201 | − | pull.number, pull.number | |
| 202 | − | )); | |
| 203 | + | // An open pull request that g1t has not stopped on: its agent is | |
| 204 | + | // woken to answer (see `wake_for_messages`). | |
| 205 | + | let wakeable = pull.status == PullStatus::Open | |
| 206 | + | && self | |
| 207 | + | .db | |
| 208 | + | .prepare("SELECT 1 AS value FROM pulls WHERE id = ? AND stalled IS NULL") | |
| 209 | + | .bind(&[pull.id.as_str().into()])? | |
| 210 | + | .first::<u32>(Some("value")) | |
| 211 | + | .await? | |
| 212 | + | .is_some(); | |
| 213 | + | message.hint = Some(if wakeable { | |
| 214 | + | self.publish("agent.asked", &repo.id, &a.actor, Self::pull_event(&pull)).await?; | |
| 215 | + | format!( | |
| 216 | + | "The agent on #{} was not at work, so g1t is waking it to answer; the answer reaches you at a later step. Its change is there to read meanwhile: get_pull_request and get_pull_request_changes on #{}.", | |
| 217 | + | pull.number, pull.number | |
| 218 | + | ) | |
| 219 | + | } else { | |
| 220 | + | format!( | |
| 221 | + | "The agent on #{} is not at work right now, so it will not answer soon. Its change is there to read: use get_pull_request and get_pull_request_changes on #{}, and decide from that.", | |
| 222 | + | pull.number, pull.number | |
| 223 | + | ) | |
| 224 | + | }); | |
| 203 | 225 | } | |
| 204 | 226 | let said = match (message.kind.as_str(), message.from_number) { | |
| 205 | 227 | ("question", Some(from)) => format!("was asked a question by the agent on #{from}"), | |
| 288 | 310 | .await?; | |
| 289 | 311 | asked.answer = Some(body.clone()); | |
| 290 | 312 | asked.declined = a.decline; | |
| 313 | + | // An agent woken to answer lets go of its pull request once nothing | |
| 314 | + | // it was asked is left unanswered. | |
| 315 | + | self.db | |
| 316 | + | .prepare( | |
| 317 | + | "UPDATE pulls SET working_on = NULL, working_until = NULL | |
| 318 | + | WHERE id = (SELECT pull_id FROM agent_messages WHERE id = ?1) | |
| 319 | + | AND working_on = 'answer' | |
| 320 | + | AND NOT EXISTS ( | |
| 321 | + | SELECT 1 FROM agent_messages | |
| 322 | + | WHERE pull_id = pulls.id AND kind IN ('question', 'handoff') AND answer IS NULL)", | |
| 323 | + | ) | |
| 324 | + | .bind(&[asked.id.as_str().into()])? | |
| 325 | + | .run() | |
| 326 | + | .await?; | |
| 291 | 327 | // Back to whoever asked: the agent on the other pull request. | |
| 292 | 328 | if let Some(from) = asked.from_number { | |
| 293 | 329 | if let Some(back) = self.pull(&repo.id, from).await?.filter(|pull| pull.status.is_active()) { | |
| 363 | 399 | Outcome::Ok(found) => found, | |
| 364 | 400 | Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)), | |
| 365 | 401 | }; | |
| 402 | + | Ok(Outcome::Ok(self.deliver(&pull).await?)) | |
| 403 | + | } | |
| 404 | + | ||
| 405 | + | /// Wakes the agent on a pull request to answer what it was asked while | |
| 406 | + | /// it was not at work: claims a short step, and hands over its messages. | |
| 407 | + | pub(crate) async fn wake_for_messages(&self, a: WakeForMessagesArgs) -> Result<Option<Wake>> { | |
| 408 | + | let Some(pull) = self.pull_by_id(&a.pull_id).await? else { | |
| 409 | + | return Ok(None); | |
| 410 | + | }; | |
| 411 | + | // Only another agent's question or handoff wakes it; what people | |
| 412 | + | // say waits for its next step. | |
| 413 | + | let waiting = self | |
| 414 | + | .db | |
| 415 | + | .prepare( | |
| 416 | + | "SELECT 1 AS value FROM agent_messages | |
| 417 | + | WHERE pull_id = ? AND delivered_at IS NULL AND kind IN ('question', 'handoff') | |
| 418 | + | LIMIT 1", | |
| 419 | + | ) | |
| 420 | + | .bind(&[pull.id.as_str().into()])? | |
| 421 | + | .first::<u32>(Some("value")) | |
| 422 | + | .await? | |
| 423 | + | .is_some(); | |
| 424 | + | if !waiting || !self.claim(&pull.id, "answer", ANSWER_MINUTES, false).await? { | |
| 425 | + | return Ok(None); | |
| 426 | + | } | |
| 427 | + | let repo: Outcome<g1t_contracts::repos::Repo> = g1t_kit::call( | |
| 428 | + | &self.repos, | |
| 429 | + | "get_by_id", | |
| 430 | + | &g1t_contracts::repos::GetByIdArgs { | |
| 431 | + | id: pull.repo_id.clone(), | |
| 432 | + | viewer: self.author_viewer(&pull).await?, | |
| 433 | + | }, | |
| 434 | + | ) | |
| 435 | + | .await?; | |
| 436 | + | let Outcome::Ok(repo) = repo else { | |
| 437 | + | return Ok(None); | |
| 438 | + | }; | |
| 439 | + | let path = g1t_contracts::repos::RepoPath { | |
| 440 | + | namespace: repo.namespace, | |
| 441 | + | name: repo.name, | |
| 442 | + | }; | |
| 443 | + | let issue = match pull.issue { | |
| 444 | + | Some(number) => self.issue(&pull.repo_id, number).await?, | |
| 445 | + | None => None, | |
| 446 | + | }; | |
| 447 | + | let messages = self.deliver(&pull).await?; | |
| 448 | + | let asking: Vec<String> = messages | |
| 449 | + | .iter() | |
| 450 | + | .filter_map(|message| message.from_number.map(|from| format!("#{from}"))) | |
| 451 | + | .collect(); | |
| 452 | + | self.note( | |
| 453 | + | &pull.repo_id, | |
| 454 | + | pull.number, | |
| 455 | + | (crate::lifecycle::POLICY_ACTOR_ID, crate::lifecycle::POLICY_ACTOR_NAME), | |
| 456 | + | &format!("woke g1t-agent to answer the agent on {}", asking.join(", ")), | |
| 457 | + | ) | |
| 458 | + | .await?; | |
| 459 | + | Ok(Some(Wake { | |
| 460 | + | job: LifecycleJob { | |
| 461 | + | pull_id: pull.id, | |
| 462 | + | source: pull.fork.unwrap_or_else(|| path.clone()), | |
| 463 | + | repo: path, | |
| 464 | + | number: pull.number, | |
| 465 | + | author: pull.author, | |
| 466 | + | branch: pull.branch, | |
| 467 | + | default_branch: repo.default_branch, | |
| 468 | + | title: pull.title, | |
| 469 | + | description: pull.body.unwrap_or_default(), | |
| 470 | + | issue, | |
| 471 | + | feedback: String::new(), | |
| 472 | + | round: 0, | |
| 473 | + | }, | |
| 474 | + | messages, | |
| 475 | + | })) | |
| 476 | + | } | |
| 477 | + | ||
| 478 | + | /// Marks a pull request's undelivered messages delivered, records them | |
| 479 | + | /// in its session, and returns them, oldest first. | |
| 480 | + | async fn deliver(&self, pull: &Pull) -> Result<Vec<AgentMessage>> { | |
| 366 | 481 | let now = rfc3339(now_ms()); | |
| 367 | − | let taken: Vec<AgentMessage> = self | |
| 482 | + | let mut taken: Vec<AgentMessage> = self | |
| 368 | 483 | .db | |
| 369 | 484 | .prepare( | |
| 370 | 485 | "UPDATE agent_messages SET delivered_at = ? | |
| 378 | 493 | .into_iter() | |
| 379 | 494 | .map(AgentMessage::from) | |
| 380 | 495 | .collect(); | |
| 496 | + | taken.sort_by(|a, b| a.created_at.cmp(&b.created_at)); | |
| 381 | 497 | if !taken.is_empty() { | |
| 382 | 498 | let entries: Vec<NewSessionEntry> = taken | |
| 383 | 499 | .iter() | |
| 393 | 509 | commit: None, | |
| 394 | 510 | }) | |
| 395 | 511 | .collect(); | |
| 396 | − | self.append_entries(&pull, &entries).await?; | |
| 512 | + | self.append_entries(pull, &entries).await?; | |
| 397 | 513 | } | |
| 398 | − | Ok(Outcome::Ok(taken)) | |
| 514 | + | Ok(taken) | |
| 399 | 515 | } | |
| 400 | 516 | } |