flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/work/src/messages.rs

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