pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/messages.rs

516 lines20,446 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::time::rfc3339;
7use g1t_contracts::work::*;
8use g1t_contracts::{FailureCode, Outcome, PrincipalKind, new_id};
9use g1t_kit::now_ms;
10use serde::Deserialize;
11use worker::Result;
12
13use crate::Work;
14
15/// The longest message an agent is sent.
16const MAX_MESSAGE_CHARS: usize = 4000;
17
18/// Kinds an agent may send another.
19const ASKS: [&str; 2] = ["question", "handoff"];
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.
23const ANSWER_MINUTES: u64 = 20;
24
25#[derive(Deserialize)]
26struct MessageRow {
27 id: String,
28 #[allow(dead_code)]
29 pull_id: String,
30 author_name: String,
31 body: String,
32 created_at: String,
33 delivered_at: Option<String>,
34 kind: String,
35 from_number: Option<u32>,
36 to_number: Option<u32>,
37 answer: Option<String>,
38 declined: u32,
39}
40
41impl From<MessageRow> for AgentMessage {
42 fn from(row: MessageRow) -> Self {
43 AgentMessage {
44 id: row.id,
45 author: row.author_name,
46 body: row.body,
47 created_at: row.created_at,
48 delivered_at: row.delivered_at,
49 kind: row.kind,
50 from_number: row.from_number,
51 to_number: row.to_number.unwrap_or_default(),
52 answer: row.answer,
53 declined: row.declined != 0,
54 hint: None,
55 }
56 }
57}
58
59/// Who sent a message, in a sentence's words.
60fn sender(message: &AgentMessage) -> String {
61 match message.from_number {
62 Some(number) => format!("the agent on #{number}"),
63 None => message.author.clone(),
64 }
65}
66
67impl Work {
68 pub(crate) async fn locate_pull(&self, a: LocatePullArgs) -> Result<Outcome<LocatedPull>> {
69 let missing = || Outcome::fail(FailureCode::NotFound, "Pull request not found.");
70 let Some(pull) = self.pull_by_id(&a.id).await? else {
71 return Ok(missing());
72 };
73 let repo: Outcome<g1t_contracts::repos::Repo> = g1t_kit::call(
74 &self.repos,
75 "get_by_id",
76 &g1t_contracts::repos::GetByIdArgs {
77 id: pull.repo_id.clone(),
78 viewer: a.viewer,
79 },
80 )
81 .await?;
82 let Outcome::Ok(repo) = repo else {
83 return Ok(missing());
84 };
85 Ok(Outcome::Ok(LocatedPull {
86 repo: g1t_contracts::repos::RepoPath {
87 namespace: repo.namespace,
88 name: repo.name,
89 },
90 number: pull.number,
91 title: pull.title,
92 status: pull.status,
93 }))
94 }
95
96 /// Every message sent to the agent on a pull request, oldest first.
97 pub(crate) async fn messages(&self, pull_id: &str) -> Result<Vec<AgentMessage>> {
98 Ok(self
99 .db
100 .prepare("SELECT * FROM agent_messages WHERE pull_id = ? ORDER BY created_at, id")
101 .bind(&[pull_id.into()])?
102 .all()
103 .await?
104 .results::<MessageRow>()?
105 .into_iter()
106 .map(AgentMessage::from)
107 .collect())
108 }
109
110 pub(crate) async fn message_agent(&self, a: MessageAgentArgs) -> Result<Outcome<AgentMessage>> {
111 let viewer = Some(a.actor.clone());
112 let (repo, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
113 Outcome::Ok(found) => found,
114 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
115 };
116 let from_agent = a.actor.kind == PrincipalKind::Agent;
117 let kind = match (&a.kind, from_agent) {
118 (Some(kind), true) if ASKS.contains(&kind.as_str()) => kind.clone(),
119 (None, true) => "question".to_owned(),
120 (_, true) => {
121 return Ok(Outcome::fail(
122 FailureCode::Invalid,
123 "An agent sends a question or a handoff.",
124 ));
125 }
126 (_, false) => "message".to_owned(),
127 };
128 if from_agent && a.from_number.is_none() {
129 return Ok(Outcome::fail(
130 FailureCode::Invalid,
131 "Say which pull request you are working on, as from_number; it is where the answer goes.",
132 ));
133 }
134 if !a.actor.verified
135 || (pull.author.id != a.actor.id && !a.actor.is_member(&repo.namespace))
136 {
137 return Ok(Outcome::fail(
138 FailureCode::Forbidden,
139 "Only the pull request's author and members of the workspace can message its agent.",
140 ));
141 }
142 if !pull.status.is_active() {
143 return Ok(Outcome::fail(
144 FailureCode::Conflict,
145 "This pull request is no longer being worked on.",
146 ));
147 }
148 let body = a.body.trim();
149 if body.is_empty() {
150 return Ok(Outcome::fail(FailureCode::Invalid, "Write a message."));
151 }
152 let body: String = body.chars().take(MAX_MESSAGE_CHARS).collect();
153 // An agent may name its issue rather than its pull request: the
154 // answer goes to the issue's pull request that is still open.
155 let from_number = match (from_agent, a.from_number) {
156 (true, Some(from)) => match self.pull(&repo.id, from).await? {
157 Some(_) => Some(from),
158 None => self.db
159 .prepare(
160 "SELECT number AS value FROM pulls
161 WHERE repo_id = ? AND issue_number = ? AND status IN ('draft', 'open')
162 ORDER BY number DESC LIMIT 1",
163 )
164 .bind(&[repo.id.as_str().into(), from.into()])?
165 .first::<u32>(Some("value"))
166 .await?
167 .or(Some(from)),
168 },
169 _ => None,
170 };
171 if from_agent && from_number == Some(pull.number) {
172 return Ok(Outcome::fail(FailureCode::Invalid, "That is your own pull request."));
173 }
174 // Whether the agent asked is at work now, to read it soon.
175 let at_work = pull.status == PullStatus::Draft
176 || self
177 .db
178 .prepare(
179 "SELECT 1 AS value FROM pulls
180 WHERE id = ? AND working_on = 'revision' AND working_until > ?",
181 )
182 .bind(&[pull.id.as_str().into(), rfc3339(now_ms()).into()])?
183 .first::<u32>(Some("value"))
184 .await?
185 .is_some();
186 let now = now_ms();
187 let message = AgentMessage {
188 id: new_id("msg", now),
189 author: a.actor.username.clone(),
190 body,
191 created_at: rfc3339(now),
192 delivered_at: None,
193 kind,
194 from_number,
195 to_number: pull.number,
196 answer: None,
197 declined: false,
198 hint: None,
199 };
200 self.insert_message(&repo.id, &pull.id, &a.actor.id, &message).await?;
201 let mut message = message;
202 if from_agent && !at_work {
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 });
225 }
226 let said = match (message.kind.as_str(), message.from_number) {
227 ("question", Some(from)) => format!("was asked a question by the agent on #{from}"),
228 ("handoff", Some(from)) => format!("was handed work by the agent on #{from}"),
229 _ => "sent the agent a message".to_owned(),
230 };
231 self.note(&repo.id, pull.number, (a.actor.id.as_str(), a.actor.username.as_str()), &said)
232 .await?;
233 Ok(Outcome::Ok(message))
234 }
235
236 async fn insert_message(
237 &self,
238 repo_id: &str,
239 pull_id: &str,
240 author_id: &str,
241 message: &AgentMessage,
242 ) -> Result<()> {
243 self.db
244 .prepare(
245 "INSERT INTO agent_messages
246 (id, pull_id, repo_id, author_id, author_name, body, created_at, kind,
247 from_number, to_number)
248 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
249 )
250 .bind(&[
251 message.id.as_str().into(),
252 pull_id.into(),
253 repo_id.into(),
254 author_id.into(),
255 message.author.as_str().into(),
256 message.body.as_str().into(),
257 message.created_at.as_str().into(),
258 message.kind.as_str().into(),
259 message
260 .from_number
261 .map_or(worker::wasm_bindgen::JsValue::NULL, |n| n.into()),
262 message.to_number.into(),
263 ])?
264 .run()
265 .await?;
266 Ok(())
267 }
268
269 /// The agent asked answers: the reply is recorded on the question and
270 /// sent back to the asking agent as a message of its own.
271 pub(crate) async fn answer_message(&self, a: AnswerMessageArgs) -> Result<Outcome<AgentMessage>> {
272 let viewer = Some(a.actor.clone());
273 let repo = match self.repo(&a.repo, &viewer).await? {
274 Outcome::Ok(repo) => repo,
275 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
276 };
277 if !a.actor.is_member(&repo.namespace) {
278 return Ok(Outcome::fail(FailureCode::Forbidden, "Only members and g1t's agents answer."));
279 }
280 let row = self
281 .db
282 .prepare("SELECT * FROM agent_messages WHERE id = ? AND repo_id = ?")
283 .bind(&[a.id.as_str().into(), repo.id.as_str().into()])?
284 .first::<MessageRow>(None)
285 .await?;
286 let Some(row) = row else {
287 return Ok(Outcome::fail(FailureCode::NotFound, "No such message."));
288 };
289 let mut asked = AgentMessage::from(row);
290 if !ASKS.contains(&asked.kind.as_str()) {
291 return Ok(Outcome::fail(FailureCode::Invalid, "Only a question or a handoff is answered."));
292 }
293 if asked.answer.is_some() {
294 return Ok(Outcome::fail(FailureCode::Conflict, "It has been answered already."));
295 }
296 let body: String = a.body.trim().chars().take(MAX_MESSAGE_CHARS).collect();
297 if body.is_empty() {
298 return Ok(Outcome::fail(FailureCode::Invalid, "Write an answer."));
299 }
300 let now = now_ms();
301 self.db
302 .prepare("UPDATE agent_messages SET answer = ?, answered_at = ?, declined = ? WHERE id = ?")
303 .bind(&[
304 body.as_str().into(),
305 rfc3339(now).into(),
306 u32::from(a.decline).into(),
307 asked.id.as_str().into(),
308 ])?
309 .run()
310 .await?;
311 asked.answer = Some(body.clone());
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?;
327 // Back to whoever asked: the agent on the other pull request.
328 if let Some(from) = asked.from_number {
329 if let Some(back) = self.pull(&repo.id, from).await?.filter(|pull| pull.status.is_active()) {
330 let reply = AgentMessage {
331 id: new_id("msg", now),
332 author: a.actor.username.clone(),
333 body: if a.decline { format!("Declined: {body}") } else { body },
334 created_at: rfc3339(now),
335 delivered_at: None,
336 kind: "answer".to_owned(),
337 from_number: Some(asked.to_number),
338 to_number: from,
339 answer: None,
340 declined: a.decline,
341 hint: None,
342 };
343 self.insert_message(&repo.id, &back.id, &a.actor.id, &reply).await?;
344 }
345 let said = if asked.kind == "handoff" {
346 if a.decline { "declined the handoff from" } else { "took on the handoff from" }
347 } else {
348 "answered the question from"
349 };
350 self.note(
351 &repo.id,
352 asked.to_number,
353 (a.actor.id.as_str(), a.actor.username.as_str()),
354 &format!("{said} the agent on #{from}"),
355 )
356 .await?;
357 }
358 Ok(Outcome::Ok(asked))
359 }
360
361 /// Questions, handoffs and answers between the agents on these pull
362 /// requests, newest first.
363 pub(crate) async fn exchanges(&self, repo_id: &str, numbers: &[u32]) -> Result<Vec<AgentMessage>> {
364 if numbers.is_empty() {
365 return Ok(Vec::new());
366 }
367 let rows = self
368 .db
369 .prepare(
370 "SELECT * FROM agent_messages
371 WHERE repo_id = ? AND kind IN ('question', 'handoff')
372 ORDER BY created_at DESC LIMIT 50",
373 )
374 .bind(&[repo_id.into()])?
375 .all()
376 .await?
377 .results::<MessageRow>()?;
378 Ok(rows
379 .into_iter()
380 .map(AgentMessage::from)
381 .filter(|message| {
382 numbers.contains(&message.to_number)
383 || message.from_number.is_some_and(|from| numbers.contains(&from))
384 })
385 .collect())
386 }
387
388 /// The undelivered messages, marked delivered and recorded in the
389 /// session, for the agent's sandbox.
390 pub(crate) async fn take_messages(&self, a: TakeMessagesArgs) -> Result<Outcome<Vec<AgentMessage>>> {
391 if a.actor.kind != PrincipalKind::Agent {
392 return Ok(Outcome::fail(
393 FailureCode::Forbidden,
394 "Only g1t's agents take messages.",
395 ));
396 }
397 let viewer = Some(a.actor.clone());
398 let (_, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
399 Outcome::Ok(found) => found,
400 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
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>> {
481 let now = rfc3339(now_ms());
482 let mut taken: Vec<AgentMessage> = self
483 .db
484 .prepare(
485 "UPDATE agent_messages SET delivered_at = ?
486 WHERE pull_id = ? AND delivered_at IS NULL
487 RETURNING *",
488 )
489 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
490 .all()
491 .await?
492 .results::<MessageRow>()?
493 .into_iter()
494 .map(AgentMessage::from)
495 .collect();
496 taken.sort_by(|a, b| a.created_at.cmp(&b.created_at));
497 if !taken.is_empty() {
498 let entries: Vec<NewSessionEntry> = taken
499 .iter()
500 .map(|message| NewSessionEntry {
501 kind: SessionEntryKind::Prompt,
502 text: match message.kind.as_str() {
503 "question" => format!("Question from {} ({}): {}", sender(message), message.id, message.body),
504 "handoff" => format!("Work handed over by {} ({}): {}", sender(message), message.id, message.body),
505 "answer" => format!("Answer from {}: {}", sender(message), message.body),
506 _ => format!("Message from {}: {}", message.author, message.body),
507 },
508 tool: None,
509 commit: None,
510 })
511 .collect();
512 self.append_entries(pull, &entries).await?;
513 }
514 Ok(taken)
515 }
516}