pr_01m47d24b0e6n91zwymwxg0vpx/services/work/src/messages.rs

524 lines20,848 bytesCodeBlame
1//! 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.
5
6use g1t_contracts::access::Capability;
7use 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
19/// Kinds an agent may send another.
20const ASKS: [&str; 2] = ["question", "handoff"];
21
22/// 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
26#[derive(Deserialize)]
27struct MessageRow {
28 id: String,
29 #[allow(dead_code)]
30 pull_id: String,
31 author_name: String,
32 body: String,
33 created_at: String,
34 delivered_at: Option<String>,
35 kind: String,
36 from_number: Option<u32>,
37 to_number: Option<u32>,
38 answer: Option<String>,
39 declined: u32,
40}
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,
50 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,
56 }
57 }
58}
59
60/// 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
68impl Work {
69 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
97 /// 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>> {
99 Ok(self
100 .db
101 .prepare("SELECT * FROM agent_messages WHERE pull_id = ? ORDER BY created_at, id")
102 .bind(&[pull_id.into()])?
103 .all()
104 .await?
105 .results::<MessageRow>()?
106 .into_iter()
107 .map(AgentMessage::from)
108 .collect())
109 }
110
111 pub(crate) async fn message_agent(&self, a: MessageAgentArgs) -> Result<Outcome<AgentMessage>> {
112 let viewer = Some(a.actor.clone());
113 let (repo, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
114 Outcome::Ok(found) => found,
115 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
116 };
117 if let Outcome::Fail(failure) = crate::retired::writable(&repo) {
118 return Ok(Outcome::Fail(failure));
119 }
120 let from_agent = a.actor.kind == PrincipalKind::Agent;
121 let kind = match (&a.kind, from_agent) {
122 (Some(kind), true) if ASKS.contains(&kind.as_str()) => kind.clone(),
123 (None, true) => "question".to_owned(),
124 (_, true) => {
125 return Ok(Outcome::fail(
126 FailureCode::Invalid,
127 "An agent sends a question or a handoff.",
128 ));
129 }
130 (_, false) => "message".to_owned(),
131 };
132 if from_agent && a.from_number.is_none() {
133 return Ok(Outcome::fail(
134 FailureCode::Invalid,
135 "Say which pull request you are working on, as from_number; it is where the answer goes.",
136 ));
137 }
138 if !a.actor.verified {
139 return Ok(Outcome::fail(FailureCode::Forbidden, crate::UNVERIFIED));
140 }
141 // Its author may always steer it; anyone else puts compute to work.
142 if pull.author.id != a.actor.id
143 && let Outcome::Fail(failure) = crate::allowed(Some(&a.actor), &repo, Capability::Run)
144 {
145 return Ok(Outcome::Fail(failure));
146 }
147 if !pull.status.is_active() {
148 return Ok(Outcome::fail(
149 FailureCode::Conflict,
150 "This pull request is no longer being worked on.",
151 ));
152 }
153 let body = a.body.trim();
154 if body.is_empty() {
155 return Ok(Outcome::fail(FailureCode::Invalid, "Write a message."));
156 }
157 let body: String = body.chars().take(MAX_MESSAGE_CHARS).collect();
158 // An agent may name its issue rather than its pull request: the
159 // answer goes to the issue's pull request that is still open.
160 let from_number = match (from_agent, a.from_number) {
161 (true, Some(from)) => match self.pull(&repo.id, from).await? {
162 Some(_) => Some(from),
163 None => self.db
164 .prepare(
165 "SELECT number AS value FROM pulls
166 WHERE repo_id = ? AND issue_number = ? AND status IN ('draft', 'open')
167 ORDER BY number DESC LIMIT 1",
168 )
169 .bind(&[repo.id.as_str().into(), from.into()])?
170 .first::<u32>(Some("value"))
171 .await?
172 .or(Some(from)),
173 },
174 _ => None,
175 };
176 if from_agent && from_number == Some(pull.number) {
177 return Ok(Outcome::fail(FailureCode::Invalid, "That is your own pull request."));
178 }
179 // Whether the agent asked is at work now, to read it soon.
180 let at_work = pull.status == PullStatus::Draft
181 || self
182 .db
183 .prepare(
184 "SELECT 1 AS value FROM pulls
185 WHERE id = ? AND working_on = 'revision' AND working_until > ?",
186 )
187 .bind(&[pull.id.as_str().into(), rfc3339(now_ms()).into()])?
188 .first::<u32>(Some("value"))
189 .await?
190 .is_some();
191 let now = now_ms();
192 let message = AgentMessage {
193 id: new_id("msg", now),
194 author: a.actor.username.clone(),
195 body,
196 created_at: rfc3339(now),
197 delivered_at: None,
198 kind,
199 from_number,
200 to_number: pull.number,
201 answer: None,
202 declined: false,
203 hint: None,
204 };
205 self.insert_message(&repo.id, &pull.id, &a.actor.id, &message).await?;
206 let mut message = message;
207 if from_agent && !at_work {
208 // An open pull request that g1t has not stopped on: its agent is
209 // woken to answer (see `wake_for_messages`).
210 let wakeable = pull.status == PullStatus::Open
211 && self
212 .db
213 .prepare("SELECT 1 AS value FROM pulls WHERE id = ? AND stalled IS NULL")
214 .bind(&[pull.id.as_str().into()])?
215 .first::<u32>(Some("value"))
216 .await?
217 .is_some();
218 message.hint = Some(if wakeable {
219 self.publish("agent.asked", &repo.id, &a.actor, Self::pull_event(&pull)).await?;
220 format!(
221 "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 #{}.",
222 pull.number, pull.number
223 )
224 } else {
225 format!(
226 "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.",
227 pull.number, pull.number
228 )
229 });
230 }
231 let said = match (message.kind.as_str(), message.from_number) {
232 ("question", Some(from)) => format!("was asked a question by the agent on #{from}"),
233 ("handoff", Some(from)) => format!("was handed work by the agent on #{from}"),
234 _ => "sent the agent a message".to_owned(),
235 };
236 self.note(&repo.id, pull.number, (a.actor.id.as_str(), a.actor.username.as_str()), &said)
237 .await?;
238 Ok(Outcome::Ok(message))
239 }
240
241 async fn insert_message(
242 &self,
243 repo_id: &str,
244 pull_id: &str,
245 author_id: &str,
246 message: &AgentMessage,
247 ) -> Result<()> {
248 self.db
249 .prepare(
250 "INSERT INTO agent_messages
251 (id, pull_id, repo_id, author_id, author_name, body, created_at, kind,
252 from_number, to_number)
253 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
254 )
255 .bind(&[
256 message.id.as_str().into(),
257 pull_id.into(),
258 repo_id.into(),
259 author_id.into(),
260 message.author.as_str().into(),
261 message.body.as_str().into(),
262 message.created_at.as_str().into(),
263 message.kind.as_str().into(),
264 message
265 .from_number
266 .map_or(worker::wasm_bindgen::JsValue::NULL, |n| n.into()),
267 message.to_number.into(),
268 ])?
269 .run()
270 .await?;
271 Ok(())
272 }
273
274 /// The agent asked answers: the reply is recorded on the question and
275 /// sent back to the asking agent as a message of its own.
276 pub(crate) async fn answer_message(&self, a: AnswerMessageArgs) -> Result<Outcome<AgentMessage>> {
277 let viewer = Some(a.actor.clone());
278 let repo = match self.repo(&a.repo, &viewer).await? {
279 Outcome::Ok(repo) => repo,
280 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
281 };
282 if let Outcome::Fail(failure) = crate::retired::writable(&repo) {
283 return Ok(Outcome::Fail(failure));
284 }
285 if let Outcome::Fail(failure) = crate::allowed(Some(&a.actor), &repo, Capability::Run) {
286 return Ok(Outcome::Fail(failure));
287 }
288 let row = self
289 .db
290 .prepare("SELECT * FROM agent_messages WHERE id = ? AND repo_id = ?")
291 .bind(&[a.id.as_str().into(), repo.id.as_str().into()])?
292 .first::<MessageRow>(None)
293 .await?;
294 let Some(row) = row else {
295 return Ok(Outcome::fail(FailureCode::NotFound, "No such message."));
296 };
297 let mut asked = AgentMessage::from(row);
298 if !ASKS.contains(&asked.kind.as_str()) {
299 return Ok(Outcome::fail(FailureCode::Invalid, "Only a question or a handoff is answered."));
300 }
301 if asked.answer.is_some() {
302 return Ok(Outcome::fail(FailureCode::Conflict, "It has been answered already."));
303 }
304 let body: String = a.body.trim().chars().take(MAX_MESSAGE_CHARS).collect();
305 if body.is_empty() {
306 return Ok(Outcome::fail(FailureCode::Invalid, "Write an answer."));
307 }
308 let now = now_ms();
309 self.db
310 .prepare("UPDATE agent_messages SET answer = ?, answered_at = ?, declined = ? WHERE id = ?")
311 .bind(&[
312 body.as_str().into(),
313 rfc3339(now).into(),
314 u32::from(a.decline).into(),
315 asked.id.as_str().into(),
316 ])?
317 .run()
318 .await?;
319 asked.answer = Some(body.clone());
320 asked.declined = a.decline;
321 // An agent woken to answer lets go of its pull request once nothing
322 // it was asked is left unanswered.
323 self.db
324 .prepare(
325 "UPDATE pulls SET working_on = NULL, working_until = NULL
326 WHERE id = (SELECT pull_id FROM agent_messages WHERE id = ?1)
327 AND working_on = 'answer'
328 AND NOT EXISTS (
329 SELECT 1 FROM agent_messages
330 WHERE pull_id = pulls.id AND kind IN ('question', 'handoff') AND answer IS NULL)",
331 )
332 .bind(&[asked.id.as_str().into()])?
333 .run()
334 .await?;
335 // Back to whoever asked: the agent on the other pull request.
336 if let Some(from) = asked.from_number {
337 if let Some(back) = self.pull(&repo.id, from).await?.filter(|pull| pull.status.is_active()) {
338 let reply = AgentMessage {
339 id: new_id("msg", now),
340 author: a.actor.username.clone(),
341 body: if a.decline { format!("Declined: {body}") } else { body },
342 created_at: rfc3339(now),
343 delivered_at: None,
344 kind: "answer".to_owned(),
345 from_number: Some(asked.to_number),
346 to_number: from,
347 answer: None,
348 declined: a.decline,
349 hint: None,
350 };
351 self.insert_message(&repo.id, &back.id, &a.actor.id, &reply).await?;
352 }
353 let said = if asked.kind == "handoff" {
354 if a.decline { "declined the handoff from" } else { "took on the handoff from" }
355 } else {
356 "answered the question from"
357 };
358 self.note(
359 &repo.id,
360 asked.to_number,
361 (a.actor.id.as_str(), a.actor.username.as_str()),
362 &format!("{said} the agent on #{from}"),
363 )
364 .await?;
365 }
366 Ok(Outcome::Ok(asked))
367 }
368
369 /// Questions, handoffs and answers between the agents on these pull
370 /// requests, newest first.
371 pub(crate) async fn exchanges(&self, repo_id: &str, numbers: &[u32]) -> Result<Vec<AgentMessage>> {
372 if numbers.is_empty() {
373 return Ok(Vec::new());
374 }
375 let rows = self
376 .db
377 .prepare(
378 "SELECT * FROM agent_messages
379 WHERE repo_id = ? AND kind IN ('question', 'handoff')
380 ORDER BY created_at DESC LIMIT 50",
381 )
382 .bind(&[repo_id.into()])?
383 .all()
384 .await?
385 .results::<MessageRow>()?;
386 Ok(rows
387 .into_iter()
388 .map(AgentMessage::from)
389 .filter(|message| {
390 numbers.contains(&message.to_number)
391 || message.from_number.is_some_and(|from| numbers.contains(&from))
392 })
393 .collect())
394 }
395
396 /// The undelivered messages, marked delivered and recorded in the
397 /// session, for the agent's sandbox.
398 pub(crate) async fn take_messages(&self, a: TakeMessagesArgs) -> Result<Outcome<Vec<AgentMessage>>> {
399 if a.actor.kind != PrincipalKind::Agent {
400 return Ok(Outcome::fail(
401 FailureCode::Forbidden,
402 "Only g1t's agents take messages.",
403 ));
404 }
405 let viewer = Some(a.actor.clone());
406 let (_, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
407 Outcome::Ok(found) => found,
408 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
409 };
410 Ok(Outcome::Ok(self.deliver(&pull).await?))
411 }
412
413 /// Wakes the agent on a pull request to answer what it was asked while
414 /// it was not at work: claims a short step, and hands over its messages.
415 pub(crate) async fn wake_for_messages(&self, a: WakeForMessagesArgs) -> Result<Option<Wake>> {
416 let Some(pull) = self.pull_by_id(&a.pull_id).await? else {
417 return Ok(None);
418 };
419 // Only another agent's question or handoff wakes it; what people
420 // say waits for its next step.
421 let waiting = self
422 .db
423 .prepare(
424 "SELECT 1 AS value FROM agent_messages
425 WHERE pull_id = ? AND delivered_at IS NULL AND kind IN ('question', 'handoff')
426 LIMIT 1",
427 )
428 .bind(&[pull.id.as_str().into()])?
429 .first::<u32>(Some("value"))
430 .await?
431 .is_some();
432 if !waiting || !self.claim(&pull.id, "answer", ANSWER_MINUTES, false).await? {
433 return Ok(None);
434 }
435 let repo: Outcome<g1t_contracts::repos::Repo> = g1t_kit::call(
436 &self.repos,
437 "get_by_id",
438 &g1t_contracts::repos::GetByIdArgs {
439 id: pull.repo_id.clone(),
440 viewer: self.author_viewer(&pull).await?,
441 },
442 )
443 .await?;
444 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
445 return Ok(None);
446 };
447 let path = g1t_contracts::repos::RepoPath {
448 namespace: repo.namespace,
449 name: repo.name,
450 };
451 let issue = match pull.issue {
452 Some(number) => self.issue(&pull.repo_id, number).await?,
453 None => None,
454 };
455 let messages = self.deliver(&pull).await?;
456 let asking: Vec<String> = messages
457 .iter()
458 .filter_map(|message| message.from_number.map(|from| format!("#{from}")))
459 .collect();
460 self.note(
461 &pull.repo_id,
462 pull.number,
463 (crate::lifecycle::POLICY_ACTOR_ID, crate::lifecycle::POLICY_ACTOR_NAME),
464 &format!("woke g1t-agent to answer the agent on {}", asking.join(", ")),
465 )
466 .await?;
467 Ok(Some(Wake {
468 job: LifecycleJob {
469 pull_id: pull.id,
470 source: pull.fork.unwrap_or_else(|| path.clone()),
471 repo: path,
472 number: pull.number,
473 author: pull.author,
474 branch: pull.branch,
475 default_branch: repo.default_branch,
476 title: pull.title,
477 description: pull.body.unwrap_or_default(),
478 issue,
479 feedback: String::new(),
480 round: 0,
481 },
482 messages,
483 }))
484 }
485
486 /// Marks a pull request's undelivered messages delivered, records them
487 /// in its session, and returns them, oldest first.
488 async fn deliver(&self, pull: &Pull) -> Result<Vec<AgentMessage>> {
489 let now = rfc3339(now_ms());
490 let mut taken: Vec<AgentMessage> = self
491 .db
492 .prepare(
493 "UPDATE agent_messages SET delivered_at = ?
494 WHERE pull_id = ? AND delivered_at IS NULL
495 RETURNING *",
496 )
497 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
498 .all()
499 .await?
500 .results::<MessageRow>()?
501 .into_iter()
502 .map(AgentMessage::from)
503 .collect();
504 taken.sort_by(|a, b| a.created_at.cmp(&b.created_at));
505 if !taken.is_empty() {
506 let entries: Vec<NewSessionEntry> = taken
507 .iter()
508 .map(|message| NewSessionEntry {
509 kind: SessionEntryKind::Prompt,
510 text: match message.kind.as_str() {
511 "question" => format!("Question from {} ({}): {}", sender(message), message.id, message.body),
512 "handoff" => format!("Work handed over by {} ({}): {}", sender(message), message.id, message.body),
513 "answer" => format!("Answer from {}: {}", sender(message), message.body),
514 _ => format!("Message from {}: {}", message.author, message.body),
515 },
516 tool: None,
517 commit: None,
518 })
519 .collect();
520 self.append_entries(pull, &entries).await?;
521 }
522 Ok(taken)
523 }
524}