g1t/services/work/src/messages.rs

532 lines21,192 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 ask each other, hand each other work, and answer1//! Messages to agents at work. People steer the agent on a pull request;
2//! agents ask each other questions and hand each other work, and answer.
3//! Each reaches its agent at the agent's next step: the sandbox asks for
4//! undelivered messages after each tool call and before it stops.
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request5
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look6use g1t_contracts::access::Capability;
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request7use g1t_contracts::time::rfc3339;
8use g1t_contracts::work::*;
9use g1t_contracts::{FailureCode, Outcome, PrincipalKind, new_id};
10use g1t_kit::now_ms;
11use serde::Deserialize;
12use worker::Result;
13
14use crate::Work;
15
16/// The longest message an agent is sent.
17const MAX_MESSAGE_CHARS: usize = 4000;
18
Agents ask each other, hand each other work, and answer19/// Kinds an agent may send another.
20const ASKS: [&str; 2] = ["question", "handoff"];
21
Agents asked while not at work are woken to answer22/// How long an agent woken to answer holds its pull request: nothing else
23/// starts on it meanwhile. Answering everything lets go sooner.
24const ANSWER_MINUTES: u64 = 20;
25
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request26#[derive(Deserialize)]
27struct MessageRow {
28 id: String,
Agents ask each other, hand each other work, and answer29 #[allow(dead_code)]
30 pull_id: String,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request31 author_name: String,
32 body: String,
33 created_at: String,
34 delivered_at: Option<String>,
Agents ask each other, hand each other work, and answer35 kind: String,
36 from_number: Option<u32>,
37 to_number: Option<u32>,
38 answer: Option<String>,
39 declined: u32,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request40}
41
42impl From<MessageRow> for AgentMessage {
43 fn from(row: MessageRow) -> Self {
44 AgentMessage {
45 id: row.id,
46 author: row.author_name,
47 body: row.body,
48 created_at: row.created_at,
49 delivered_at: row.delivered_at,
Agents ask each other, hand each other work, and answer50 kind: row.kind,
51 from_number: row.from_number,
52 to_number: row.to_number.unwrap_or_default(),
53 answer: row.answer,
54 declined: row.declined != 0,
55 hint: None,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request56 }
57 }
58}
59
Agents ask each other, hand each other work, and answer60/// Who sent a message, in a sentence's words.
61fn sender(message: &AgentMessage) -> String {
62 match message.from_number {
63 Some(number) => format!("the agent on #{number}"),
64 None => message.author.clone(),
65 }
66}
67
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request68impl Work {
Record your own agent's sessions automatically69 pub(crate) async fn locate_pull(&self, a: LocatePullArgs) -> Result<Outcome<LocatedPull>> {
70 let missing = || Outcome::fail(FailureCode::NotFound, "Pull request not found.");
71 let Some(pull) = self.pull_by_id(&a.id).await? else {
72 return Ok(missing());
73 };
74 let repo: Outcome<g1t_contracts::repos::Repo> = g1t_kit::call(
75 &self.repos,
76 "get_by_id",
77 &g1t_contracts::repos::GetByIdArgs {
78 id: pull.repo_id.clone(),
79 viewer: a.viewer,
80 },
81 )
82 .await?;
83 let Outcome::Ok(repo) = repo else {
84 return Ok(missing());
85 };
86 Ok(Outcome::Ok(LocatedPull {
87 repo: g1t_contracts::repos::RepoPath {
88 namespace: repo.namespace,
89 name: repo.name,
90 },
91 number: pull.number,
92 title: pull.title,
93 status: pull.status,
94 }))
95 }
96
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request97 /// Every message sent to the agent on a pull request, oldest first.
98 pub(crate) async fn messages(&self, pull_id: &str) -> Result<Vec<AgentMessage>> {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily99 if let Some(found) = self.prefetched_pull(pull_id) {
100 return Ok(found
101 .rows::<MessageRow>(crate::prefetch::Slot::Messages)?
102 .into_iter()
103 .map(AgentMessage::from)
104 .collect());
105 }
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request106 Ok(self
107 .db
108 .prepare("SELECT * FROM agent_messages WHERE pull_id = ? ORDER BY created_at, id")
109 .bind(&[pull_id.into()])?
110 .all()
111 .await?
112 .results::<MessageRow>()?
113 .into_iter()
114 .map(AgentMessage::from)
115 .collect())
116 }
117
118 pub(crate) async fn message_agent(&self, a: MessageAgentArgs) -> Result<Outcome<AgentMessage>> {
119 let viewer = Some(a.actor.clone());
120 let (repo, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
121 Outcome::Ok(found) => found,
122 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
123 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look124 if let Outcome::Fail(failure) = crate::retired::writable(&repo) {
125 return Ok(Outcome::Fail(failure));
126 }
Agents ask each other, hand each other work, and answer127 let from_agent = a.actor.kind == PrincipalKind::Agent;
128 let kind = match (&a.kind, from_agent) {
129 (Some(kind), true) if ASKS.contains(&kind.as_str()) => kind.clone(),
130 (None, true) => "question".to_owned(),
131 (_, true) => {
132 return Ok(Outcome::fail(
133 FailureCode::Invalid,
134 "An agent sends a question or a handoff.",
135 ));
136 }
137 (_, false) => "message".to_owned(),
138 };
139 if from_agent && a.from_number.is_none() {
140 return Ok(Outcome::fail(
141 FailureCode::Invalid,
142 "Say which pull request you are working on, as from_number; it is where the answer goes.",
143 ));
144 }
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look145 if !a.actor.verified {
146 return Ok(Outcome::fail(FailureCode::Forbidden, crate::UNVERIFIED));
147 }
g1t is the stored author of what it opens; the person who asked is requested_by and keeps the author's rights148 // Its owner (whoever asked g1t for it, or its author) may always
149 // steer it; anyone else puts compute to work.
150 if !pull.is_owned_by(&a.actor.id)
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look151 && let Outcome::Fail(failure) = crate::allowed(Some(&a.actor), &repo, Capability::Run)
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request152 {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look153 return Ok(Outcome::Fail(failure));
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request154 }
155 if !pull.status.is_active() {
156 return Ok(Outcome::fail(
157 FailureCode::Conflict,
158 "This pull request is no longer being worked on.",
159 ));
160 }
161 let body = a.body.trim();
162 if body.is_empty() {
163 return Ok(Outcome::fail(FailureCode::Invalid, "Write a message."));
164 }
165 let body: String = body.chars().take(MAX_MESSAGE_CHARS).collect();
Agents ask each other, hand each other work, and answer166 // An agent may name its issue rather than its pull request: the
167 // answer goes to the issue's pull request that is still open.
168 let from_number = match (from_agent, a.from_number) {
169 (true, Some(from)) => match self.pull(&repo.id, from).await? {
170 Some(_) => Some(from),
171 None => self.db
172 .prepare(
173 "SELECT number AS value FROM pulls
174 WHERE repo_id = ? AND issue_number = ? AND status IN ('draft', 'open')
175 ORDER BY number DESC LIMIT 1",
176 )
177 .bind(&[repo.id.as_str().into(), from.into()])?
178 .first::<u32>(Some("value"))
179 .await?
180 .or(Some(from)),
181 },
182 _ => None,
183 };
184 if from_agent && from_number == Some(pull.number) {
185 return Ok(Outcome::fail(FailureCode::Invalid, "That is your own pull request."));
186 }
187 // Whether the agent asked is at work now, to read it soon.
188 let at_work = pull.status == PullStatus::Draft
189 || self
190 .db
191 .prepare(
192 "SELECT 1 AS value FROM pulls
193 WHERE id = ? AND working_on = 'revision' AND working_until > ?",
194 )
195 .bind(&[pull.id.as_str().into(), rfc3339(now_ms()).into()])?
196 .first::<u32>(Some("value"))
197 .await?
198 .is_some();
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request199 let now = now_ms();
200 let message = AgentMessage {
201 id: new_id("msg", now),
202 author: a.actor.username.clone(),
203 body,
204 created_at: rfc3339(now),
205 delivered_at: None,
Agents ask each other, hand each other work, and answer206 kind,
207 from_number,
208 to_number: pull.number,
209 answer: None,
210 declined: false,
211 hint: None,
212 };
213 self.insert_message(&repo.id, &pull.id, &a.actor.id, &message).await?;
214 let mut message = message;
215 if from_agent && !at_work {
Agents asked while not at work are woken to answer216 // An open pull request that g1t has not stopped on: its agent is
217 // woken to answer (see `wake_for_messages`).
218 let wakeable = pull.status == PullStatus::Open
219 && self
220 .db
221 .prepare("SELECT 1 AS value FROM pulls WHERE id = ? AND stalled IS NULL")
222 .bind(&[pull.id.as_str().into()])?
223 .first::<u32>(Some("value"))
224 .await?
225 .is_some();
226 message.hint = Some(if wakeable {
227 self.publish("agent.asked", &repo.id, &a.actor, Self::pull_event(&pull)).await?;
228 format!(
229 "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 #{}.",
230 pull.number, pull.number
231 )
232 } else {
233 format!(
234 "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.",
235 pull.number, pull.number
236 )
237 });
Agents ask each other, hand each other work, and answer238 }
239 let said = match (message.kind.as_str(), message.from_number) {
240 ("question", Some(from)) => format!("was asked a question by the agent on #{from}"),
241 ("handoff", Some(from)) => format!("was handed work by the agent on #{from}"),
242 _ => "sent the agent a message".to_owned(),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request243 };
Agents ask each other, hand each other work, and answer244 self.note(&repo.id, pull.number, (a.actor.id.as_str(), a.actor.username.as_str()), &said)
245 .await?;
246 Ok(Outcome::Ok(message))
247 }
248
249 async fn insert_message(
250 &self,
251 repo_id: &str,
252 pull_id: &str,
253 author_id: &str,
254 message: &AgentMessage,
255 ) -> Result<()> {
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request256 self.db
257 .prepare(
Agents ask each other, hand each other work, and answer258 "INSERT INTO agent_messages
259 (id, pull_id, repo_id, author_id, author_name, body, created_at, kind,
260 from_number, to_number)
261 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request262 )
263 .bind(&[
264 message.id.as_str().into(),
Agents ask each other, hand each other work, and answer265 pull_id.into(),
266 repo_id.into(),
267 author_id.into(),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request268 message.author.as_str().into(),
269 message.body.as_str().into(),
270 message.created_at.as_str().into(),
Agents ask each other, hand each other work, and answer271 message.kind.as_str().into(),
272 message
273 .from_number
274 .map_or(worker::wasm_bindgen::JsValue::NULL, |n| n.into()),
275 message.to_number.into(),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request276 ])?
277 .run()
278 .await?;
Agents ask each other, hand each other work, and answer279 Ok(())
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request280 }
281
Agents ask each other, hand each other work, and answer282 /// The agent asked answers: the reply is recorded on the question and
283 /// sent back to the asking agent as a message of its own.
284 pub(crate) async fn answer_message(&self, a: AnswerMessageArgs) -> Result<Outcome<AgentMessage>> {
285 let viewer = Some(a.actor.clone());
286 let repo = match self.repo(&a.repo, &viewer).await? {
287 Outcome::Ok(repo) => repo,
288 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
289 };
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look290 if let Outcome::Fail(failure) = crate::retired::writable(&repo) {
291 return Ok(Outcome::Fail(failure));
292 }
293 if let Outcome::Fail(failure) = crate::allowed(Some(&a.actor), &repo, Capability::Run) {
294 return Ok(Outcome::Fail(failure));
Agents ask each other, hand each other work, and answer295 }
296 let row = self
297 .db
298 .prepare("SELECT * FROM agent_messages WHERE id = ? AND repo_id = ?")
299 .bind(&[a.id.as_str().into(), repo.id.as_str().into()])?
300 .first::<MessageRow>(None)
301 .await?;
302 let Some(row) = row else {
303 return Ok(Outcome::fail(FailureCode::NotFound, "No such message."));
304 };
305 let mut asked = AgentMessage::from(row);
306 if !ASKS.contains(&asked.kind.as_str()) {
307 return Ok(Outcome::fail(FailureCode::Invalid, "Only a question or a handoff is answered."));
308 }
309 if asked.answer.is_some() {
310 return Ok(Outcome::fail(FailureCode::Conflict, "It has been answered already."));
311 }
312 let body: String = a.body.trim().chars().take(MAX_MESSAGE_CHARS).collect();
313 if body.is_empty() {
314 return Ok(Outcome::fail(FailureCode::Invalid, "Write an answer."));
315 }
316 let now = now_ms();
317 self.db
318 .prepare("UPDATE agent_messages SET answer = ?, answered_at = ?, declined = ? WHERE id = ?")
319 .bind(&[
320 body.as_str().into(),
321 rfc3339(now).into(),
322 u32::from(a.decline).into(),
323 asked.id.as_str().into(),
324 ])?
325 .run()
326 .await?;
327 asked.answer = Some(body.clone());
328 asked.declined = a.decline;
Agents asked while not at work are woken to answer329 // An agent woken to answer lets go of its pull request once nothing
330 // it was asked is left unanswered.
331 self.db
332 .prepare(
333 "UPDATE pulls SET working_on = NULL, working_until = NULL
334 WHERE id = (SELECT pull_id FROM agent_messages WHERE id = ?1)
335 AND working_on = 'answer'
336 AND NOT EXISTS (
337 SELECT 1 FROM agent_messages
338 WHERE pull_id = pulls.id AND kind IN ('question', 'handoff') AND answer IS NULL)",
339 )
340 .bind(&[asked.id.as_str().into()])?
341 .run()
342 .await?;
Agents ask each other, hand each other work, and answer343 // Back to whoever asked: the agent on the other pull request.
344 if let Some(from) = asked.from_number {
345 if let Some(back) = self.pull(&repo.id, from).await?.filter(|pull| pull.status.is_active()) {
346 let reply = AgentMessage {
347 id: new_id("msg", now),
348 author: a.actor.username.clone(),
349 body: if a.decline { format!("Declined: {body}") } else { body },
350 created_at: rfc3339(now),
351 delivered_at: None,
352 kind: "answer".to_owned(),
353 from_number: Some(asked.to_number),
354 to_number: from,
355 answer: None,
356 declined: a.decline,
357 hint: None,
358 };
359 self.insert_message(&repo.id, &back.id, &a.actor.id, &reply).await?;
360 }
361 let said = if asked.kind == "handoff" {
362 if a.decline { "declined the handoff from" } else { "took on the handoff from" }
363 } else {
364 "answered the question from"
365 };
366 self.note(
367 &repo.id,
368 asked.to_number,
369 (a.actor.id.as_str(), a.actor.username.as_str()),
370 &format!("{said} the agent on #{from}"),
371 )
372 .await?;
373 }
374 Ok(Outcome::Ok(asked))
375 }
376
377 /// Questions, handoffs and answers between the agents on these pull
378 /// requests, newest first.
379 pub(crate) async fn exchanges(&self, repo_id: &str, numbers: &[u32]) -> Result<Vec<AgentMessage>> {
380 if numbers.is_empty() {
381 return Ok(Vec::new());
382 }
383 let rows = self
384 .db
385 .prepare(
386 "SELECT * FROM agent_messages
387 WHERE repo_id = ? AND kind IN ('question', 'handoff')
388 ORDER BY created_at DESC LIMIT 50",
389 )
390 .bind(&[repo_id.into()])?
391 .all()
392 .await?
393 .results::<MessageRow>()?;
394 Ok(rows
395 .into_iter()
396 .map(AgentMessage::from)
397 .filter(|message| {
398 numbers.contains(&message.to_number)
399 || message.from_number.is_some_and(|from| numbers.contains(&from))
400 })
401 .collect())
402 }
403
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request404 /// The undelivered messages, marked delivered and recorded in the
405 /// session, for the agent's sandbox.
406 pub(crate) async fn take_messages(&self, a: TakeMessagesArgs) -> Result<Outcome<Vec<AgentMessage>>> {
407 if a.actor.kind != PrincipalKind::Agent {
408 return Ok(Outcome::fail(
409 FailureCode::Forbidden,
410 "Only g1t's agents take messages.",
411 ));
412 }
413 let viewer = Some(a.actor.clone());
414 let (_, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
415 Outcome::Ok(found) => found,
416 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
417 };
Agents asked while not at work are woken to answer418 Ok(Outcome::Ok(self.deliver(&pull).await?))
419 }
420
421 /// Wakes the agent on a pull request to answer what it was asked while
422 /// it was not at work: claims a short step, and hands over its messages.
423 pub(crate) async fn wake_for_messages(&self, a: WakeForMessagesArgs) -> Result<Option<Wake>> {
424 let Some(pull) = self.pull_by_id(&a.pull_id).await? else {
425 return Ok(None);
426 };
427 // Only another agent's question or handoff wakes it; what people
428 // say waits for its next step.
429 let waiting = self
430 .db
431 .prepare(
432 "SELECT 1 AS value FROM agent_messages
433 WHERE pull_id = ? AND delivered_at IS NULL AND kind IN ('question', 'handoff')
434 LIMIT 1",
435 )
436 .bind(&[pull.id.as_str().into()])?
437 .first::<u32>(Some("value"))
438 .await?
439 .is_some();
440 if !waiting || !self.claim(&pull.id, "answer", ANSWER_MINUTES, false).await? {
441 return Ok(None);
442 }
443 let repo: Outcome<g1t_contracts::repos::Repo> = g1t_kit::call(
444 &self.repos,
445 "get_by_id",
446 &g1t_contracts::repos::GetByIdArgs {
447 id: pull.repo_id.clone(),
g1t is the stored author of what it opens; the person who asked is requested_by and keeps the author's rights448 viewer: self.owner_viewer(&pull).await?,
Agents asked while not at work are woken to answer449 },
450 )
451 .await?;
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look452 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
Agents asked while not at work are woken to answer453 return Ok(None);
454 };
455 let path = g1t_contracts::repos::RepoPath {
456 namespace: repo.namespace,
457 name: repo.name,
458 };
459 let issue = match pull.issue {
460 Some(number) => self.issue(&pull.repo_id, number).await?,
461 None => None,
462 };
463 let messages = self.deliver(&pull).await?;
464 let asking: Vec<String> = messages
465 .iter()
466 .filter_map(|message| message.from_number.map(|from| format!("#{from}")))
467 .collect();
468 self.note(
469 &pull.repo_id,
470 pull.number,
471 (crate::lifecycle::POLICY_ACTOR_ID, crate::lifecycle::POLICY_ACTOR_NAME),
g1t is one name: its agent's work, commits and comments show as @g1t, and nobody can claim g1t or g1t-agent472 &format!("woke g1t to answer the agent on {}", asking.join(", ")),
Agents asked while not at work are woken to answer473 )
474 .await?;
475 Ok(Some(Wake {
476 job: LifecycleJob {
477 pull_id: pull.id,
478 source: pull.fork.unwrap_or_else(|| path.clone()),
479 repo: path,
480 number: pull.number,
g1t is the stored author of what it opens; the person who asked is requested_by and keeps the author's rights481 author: pull.requested_by.unwrap_or(pull.author),
Agents asked while not at work are woken to answer482 branch: pull.branch,
483 default_branch: repo.default_branch,
484 title: pull.title,
485 description: pull.body.unwrap_or_default(),
486 issue,
487 feedback: String::new(),
488 round: 0,
489 },
490 messages,
491 }))
492 }
493
494 /// Marks a pull request's undelivered messages delivered, records them
495 /// in its session, and returns them, oldest first.
496 async fn deliver(&self, pull: &Pull) -> Result<Vec<AgentMessage>> {
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request497 let now = rfc3339(now_ms());
Agents asked while not at work are woken to answer498 let mut taken: Vec<AgentMessage> = self
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request499 .db
500 .prepare(
501 "UPDATE agent_messages SET delivered_at = ?
502 WHERE pull_id = ? AND delivered_at IS NULL
503 RETURNING *",
504 )
505 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
506 .all()
507 .await?
508 .results::<MessageRow>()?
509 .into_iter()
510 .map(AgentMessage::from)
511 .collect();
Agents asked while not at work are woken to answer512 taken.sort_by(|a, b| a.created_at.cmp(&b.created_at));
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request513 if !taken.is_empty() {
514 let entries: Vec<NewSessionEntry> = taken
515 .iter()
516 .map(|message| NewSessionEntry {
517 kind: SessionEntryKind::Prompt,
Agents ask each other, hand each other work, and answer518 text: match message.kind.as_str() {
519 "question" => format!("Question from {} ({}): {}", sender(message), message.id, message.body),
520 "handoff" => format!("Work handed over by {} ({}): {}", sender(message), message.id, message.body),
521 "answer" => format!("Answer from {}: {}", sender(message), message.body),
522 _ => format!("Message from {}: {}", message.author, message.body),
523 },
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request524 tool: None,
525 commit: None,
526 })
527 .collect();
Agents asked while not at work are woken to answer528 self.append_entries(pull, &entries).await?;
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request529 }
Agents asked while not at work are woken to answer530 Ok(taken)
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request531 }
532}