pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/messages.rs

400 lines15,580 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
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
Agents ask each other, hand each other work, and answer18/// Kinds an agent may send another.
19const ASKS: [&str; 2] = ["question", "handoff"];
20
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request21#[derive(Deserialize)]
22struct MessageRow {
23 id: String,
Agents ask each other, hand each other work, and answer24 #[allow(dead_code)]
25 pull_id: String,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request26 author_name: String,
27 body: String,
28 created_at: String,
29 delivered_at: Option<String>,
Agents ask each other, hand each other work, and answer30 kind: String,
31 from_number: Option<u32>,
32 to_number: Option<u32>,
33 answer: Option<String>,
34 declined: u32,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request35}
36
37impl From<MessageRow> for AgentMessage {
38 fn from(row: MessageRow) -> Self {
39 AgentMessage {
40 id: row.id,
41 author: row.author_name,
42 body: row.body,
43 created_at: row.created_at,
44 delivered_at: row.delivered_at,
Agents ask each other, hand each other work, and answer45 kind: row.kind,
46 from_number: row.from_number,
47 to_number: row.to_number.unwrap_or_default(),
48 answer: row.answer,
49 declined: row.declined != 0,
50 hint: None,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request51 }
52 }
53}
54
Agents ask each other, hand each other work, and answer55/// Who sent a message, in a sentence's words.
56fn sender(message: &AgentMessage) -> String {
57 match message.from_number {
58 Some(number) => format!("the agent on #{number}"),
59 None => message.author.clone(),
60 }
61}
62
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request63impl Work {
Record your own agent's sessions automatically64 pub(crate) async fn locate_pull(&self, a: LocatePullArgs) -> Result<Outcome<LocatedPull>> {
65 let missing = || Outcome::fail(FailureCode::NotFound, "Pull request not found.");
66 let Some(pull) = self.pull_by_id(&a.id).await? else {
67 return Ok(missing());
68 };
69 let repo: Outcome<g1t_contracts::repos::Repo> = g1t_kit::call(
70 &self.repos,
71 "get_by_id",
72 &g1t_contracts::repos::GetByIdArgs {
73 id: pull.repo_id.clone(),
74 viewer: a.viewer,
75 },
76 )
77 .await?;
78 let Outcome::Ok(repo) = repo else {
79 return Ok(missing());
80 };
81 Ok(Outcome::Ok(LocatedPull {
82 repo: g1t_contracts::repos::RepoPath {
83 namespace: repo.namespace,
84 name: repo.name,
85 },
86 number: pull.number,
87 title: pull.title,
88 status: pull.status,
89 }))
90 }
91
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request92 /// Every message sent to the agent on a pull request, oldest first.
93 pub(crate) async fn messages(&self, pull_id: &str) -> Result<Vec<AgentMessage>> {
94 Ok(self
95 .db
96 .prepare("SELECT * FROM agent_messages WHERE pull_id = ? ORDER BY created_at, id")
97 .bind(&[pull_id.into()])?
98 .all()
99 .await?
100 .results::<MessageRow>()?
101 .into_iter()
102 .map(AgentMessage::from)
103 .collect())
104 }
105
106 pub(crate) async fn message_agent(&self, a: MessageAgentArgs) -> Result<Outcome<AgentMessage>> {
107 let viewer = Some(a.actor.clone());
108 let (repo, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
109 Outcome::Ok(found) => found,
110 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
111 };
Agents ask each other, hand each other work, and answer112 let from_agent = a.actor.kind == PrincipalKind::Agent;
113 let kind = match (&a.kind, from_agent) {
114 (Some(kind), true) if ASKS.contains(&kind.as_str()) => kind.clone(),
115 (None, true) => "question".to_owned(),
116 (_, true) => {
117 return Ok(Outcome::fail(
118 FailureCode::Invalid,
119 "An agent sends a question or a handoff.",
120 ));
121 }
122 (_, false) => "message".to_owned(),
123 };
124 if from_agent && a.from_number.is_none() {
125 return Ok(Outcome::fail(
126 FailureCode::Invalid,
127 "Say which pull request you are working on, as from_number; it is where the answer goes.",
128 ));
129 }
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request130 if !a.actor.verified
131 || (pull.author.id != a.actor.id && !a.actor.is_member(&repo.namespace))
132 {
133 return Ok(Outcome::fail(
134 FailureCode::Forbidden,
135 "Only the pull request's author and members of the workspace can message its agent.",
136 ));
137 }
138 if !pull.status.is_active() {
139 return Ok(Outcome::fail(
140 FailureCode::Conflict,
141 "This pull request is no longer being worked on.",
142 ));
143 }
144 let body = a.body.trim();
145 if body.is_empty() {
146 return Ok(Outcome::fail(FailureCode::Invalid, "Write a message."));
147 }
148 let body: String = body.chars().take(MAX_MESSAGE_CHARS).collect();
Agents ask each other, hand each other work, and answer149 // An agent may name its issue rather than its pull request: the
150 // answer goes to the issue's pull request that is still open.
151 let from_number = match (from_agent, a.from_number) {
152 (true, Some(from)) => match self.pull(&repo.id, from).await? {
153 Some(_) => Some(from),
154 None => self.db
155 .prepare(
156 "SELECT number AS value FROM pulls
157 WHERE repo_id = ? AND issue_number = ? AND status IN ('draft', 'open')
158 ORDER BY number DESC LIMIT 1",
159 )
160 .bind(&[repo.id.as_str().into(), from.into()])?
161 .first::<u32>(Some("value"))
162 .await?
163 .or(Some(from)),
164 },
165 _ => None,
166 };
167 if from_agent && from_number == Some(pull.number) {
168 return Ok(Outcome::fail(FailureCode::Invalid, "That is your own pull request."));
169 }
170 // Whether the agent asked is at work now, to read it soon.
171 let at_work = pull.status == PullStatus::Draft
172 || self
173 .db
174 .prepare(
175 "SELECT 1 AS value FROM pulls
176 WHERE id = ? AND working_on = 'revision' AND working_until > ?",
177 )
178 .bind(&[pull.id.as_str().into(), rfc3339(now_ms()).into()])?
179 .first::<u32>(Some("value"))
180 .await?
181 .is_some();
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request182 let now = now_ms();
183 let message = AgentMessage {
184 id: new_id("msg", now),
185 author: a.actor.username.clone(),
186 body,
187 created_at: rfc3339(now),
188 delivered_at: None,
Agents ask each other, hand each other work, and answer189 kind,
190 from_number,
191 to_number: pull.number,
192 answer: None,
193 declined: false,
194 hint: None,
195 };
196 self.insert_message(&repo.id, &pull.id, &a.actor.id, &message).await?;
197 let mut message = message;
198 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 }
204 let said = match (message.kind.as_str(), message.from_number) {
205 ("question", Some(from)) => format!("was asked a question by the agent on #{from}"),
206 ("handoff", Some(from)) => format!("was handed work by the agent on #{from}"),
207 _ => "sent the agent a message".to_owned(),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request208 };
Agents ask each other, hand each other work, and answer209 self.note(&repo.id, pull.number, (a.actor.id.as_str(), a.actor.username.as_str()), &said)
210 .await?;
211 Ok(Outcome::Ok(message))
212 }
213
214 async fn insert_message(
215 &self,
216 repo_id: &str,
217 pull_id: &str,
218 author_id: &str,
219 message: &AgentMessage,
220 ) -> Result<()> {
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request221 self.db
222 .prepare(
Agents ask each other, hand each other work, and answer223 "INSERT INTO agent_messages
224 (id, pull_id, repo_id, author_id, author_name, body, created_at, kind,
225 from_number, to_number)
226 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request227 )
228 .bind(&[
229 message.id.as_str().into(),
Agents ask each other, hand each other work, and answer230 pull_id.into(),
231 repo_id.into(),
232 author_id.into(),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request233 message.author.as_str().into(),
234 message.body.as_str().into(),
235 message.created_at.as_str().into(),
Agents ask each other, hand each other work, and answer236 message.kind.as_str().into(),
237 message
238 .from_number
239 .map_or(worker::wasm_bindgen::JsValue::NULL, |n| n.into()),
240 message.to_number.into(),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request241 ])?
242 .run()
243 .await?;
Agents ask each other, hand each other work, and answer244 Ok(())
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request245 }
246
Agents ask each other, hand each other work, and answer247 /// The agent asked answers: the reply is recorded on the question and
248 /// sent back to the asking agent as a message of its own.
249 pub(crate) async fn answer_message(&self, a: AnswerMessageArgs) -> Result<Outcome<AgentMessage>> {
250 let viewer = Some(a.actor.clone());
251 let repo = match self.repo(&a.repo, &viewer).await? {
252 Outcome::Ok(repo) => repo,
253 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
254 };
255 if !a.actor.is_member(&repo.namespace) {
256 return Ok(Outcome::fail(FailureCode::Forbidden, "Only members and g1t's agents answer."));
257 }
258 let row = self
259 .db
260 .prepare("SELECT * FROM agent_messages WHERE id = ? AND repo_id = ?")
261 .bind(&[a.id.as_str().into(), repo.id.as_str().into()])?
262 .first::<MessageRow>(None)
263 .await?;
264 let Some(row) = row else {
265 return Ok(Outcome::fail(FailureCode::NotFound, "No such message."));
266 };
267 let mut asked = AgentMessage::from(row);
268 if !ASKS.contains(&asked.kind.as_str()) {
269 return Ok(Outcome::fail(FailureCode::Invalid, "Only a question or a handoff is answered."));
270 }
271 if asked.answer.is_some() {
272 return Ok(Outcome::fail(FailureCode::Conflict, "It has been answered already."));
273 }
274 let body: String = a.body.trim().chars().take(MAX_MESSAGE_CHARS).collect();
275 if body.is_empty() {
276 return Ok(Outcome::fail(FailureCode::Invalid, "Write an answer."));
277 }
278 let now = now_ms();
279 self.db
280 .prepare("UPDATE agent_messages SET answer = ?, answered_at = ?, declined = ? WHERE id = ?")
281 .bind(&[
282 body.as_str().into(),
283 rfc3339(now).into(),
284 u32::from(a.decline).into(),
285 asked.id.as_str().into(),
286 ])?
287 .run()
288 .await?;
289 asked.answer = Some(body.clone());
290 asked.declined = a.decline;
291 // Back to whoever asked: the agent on the other pull request.
292 if let Some(from) = asked.from_number {
293 if let Some(back) = self.pull(&repo.id, from).await?.filter(|pull| pull.status.is_active()) {
294 let reply = AgentMessage {
295 id: new_id("msg", now),
296 author: a.actor.username.clone(),
297 body: if a.decline { format!("Declined: {body}") } else { body },
298 created_at: rfc3339(now),
299 delivered_at: None,
300 kind: "answer".to_owned(),
301 from_number: Some(asked.to_number),
302 to_number: from,
303 answer: None,
304 declined: a.decline,
305 hint: None,
306 };
307 self.insert_message(&repo.id, &back.id, &a.actor.id, &reply).await?;
308 }
309 let said = if asked.kind == "handoff" {
310 if a.decline { "declined the handoff from" } else { "took on the handoff from" }
311 } else {
312 "answered the question from"
313 };
314 self.note(
315 &repo.id,
316 asked.to_number,
317 (a.actor.id.as_str(), a.actor.username.as_str()),
318 &format!("{said} the agent on #{from}"),
319 )
320 .await?;
321 }
322 Ok(Outcome::Ok(asked))
323 }
324
325 /// Questions, handoffs and answers between the agents on these pull
326 /// requests, newest first.
327 pub(crate) async fn exchanges(&self, repo_id: &str, numbers: &[u32]) -> Result<Vec<AgentMessage>> {
328 if numbers.is_empty() {
329 return Ok(Vec::new());
330 }
331 let rows = self
332 .db
333 .prepare(
334 "SELECT * FROM agent_messages
335 WHERE repo_id = ? AND kind IN ('question', 'handoff')
336 ORDER BY created_at DESC LIMIT 50",
337 )
338 .bind(&[repo_id.into()])?
339 .all()
340 .await?
341 .results::<MessageRow>()?;
342 Ok(rows
343 .into_iter()
344 .map(AgentMessage::from)
345 .filter(|message| {
346 numbers.contains(&message.to_number)
347 || message.from_number.is_some_and(|from| numbers.contains(&from))
348 })
349 .collect())
350 }
351
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request352 /// The undelivered messages, marked delivered and recorded in the
353 /// session, for the agent's sandbox.
354 pub(crate) async fn take_messages(&self, a: TakeMessagesArgs) -> Result<Outcome<Vec<AgentMessage>>> {
355 if a.actor.kind != PrincipalKind::Agent {
356 return Ok(Outcome::fail(
357 FailureCode::Forbidden,
358 "Only g1t's agents take messages.",
359 ));
360 }
361 let viewer = Some(a.actor.clone());
362 let (_, pull) = match self.pull_at(&a.repo, a.number, &viewer).await? {
363 Outcome::Ok(found) => found,
364 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
365 };
366 let now = rfc3339(now_ms());
367 let taken: Vec<AgentMessage> = self
368 .db
369 .prepare(
370 "UPDATE agent_messages SET delivered_at = ?
371 WHERE pull_id = ? AND delivered_at IS NULL
372 RETURNING *",
373 )
374 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
375 .all()
376 .await?
377 .results::<MessageRow>()?
378 .into_iter()
379 .map(AgentMessage::from)
380 .collect();
381 if !taken.is_empty() {
382 let entries: Vec<NewSessionEntry> = taken
383 .iter()
384 .map(|message| NewSessionEntry {
385 kind: SessionEntryKind::Prompt,
Agents ask each other, hand each other work, and answer386 text: match message.kind.as_str() {
387 "question" => format!("Question from {} ({}): {}", sender(message), message.id, message.body),
388 "handoff" => format!("Work handed over by {} ({}): {}", sender(message), message.id, message.body),
389 "answer" => format!("Answer from {}: {}", sender(message), message.body),
390 _ => format!("Message from {}: {}", message.author, message.body),
391 },
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request392 tool: None,
393 commit: None,
394 })
395 .collect();
396 self.append_entries(&pull, &entries).await?;
397 }
398 Ok(Outcome::Ok(taken))
399 }
400}