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/capture.rs

863 lines35,732 bytesCodeBlame
1//! Memory that fills itself (see `g1t_contracts::capture`).
2//!
3//! Candidates arrive from four places:
4//!
5//! - **Runs.** At the end of a run the harness asks the agent what it
6//! learned and reports it with the run's token (`report_learned`).
7//! - **Reviews.** A person's request for changes, or a comment that
8//! corrects the agent ("use the shared client instead"), on a pull
9//! request, becomes a convention candidate quoting the comment.
10//! - **Merges.** A merged pull request's title, why and files become a
11//! decision candidate. No model is asked: the pull request's own words
12//! are the decision, and a candidate waits for review anyway.
13//! - **Docs.** The context service reads a project's README, AGENTS.md,
14//! docs and manifests and sends what they say (`capture_memories`).
15//!
16//! The same thing said again, by an independent source, is one memory seen
17//! twice, which keeps it (`promotes`). A dismissed memory stays dismissed,
18//! so the same wording is never suggested again. Nothing that looks like a
19//! secret is stored, in the text or in its evidence.
20
21use std::collections::HashMap;
22
23use g1t_contracts::agents::*;
24use g1t_contracts::capture::*;
25use g1t_contracts::events::{Event, NewEvent, Publish};
26use g1t_contracts::repos::{PathByIdArgs, RepoPath};
27use g1t_contracts::time::rfc3339;
28use g1t_contracts::work::{Comment, CommentKind, Pull, PullStatus, Verdict};
29use g1t_contracts::{FailureCode, Outcome, PrincipalKind, new_id};
30use g1t_kit::now_ms;
31use serde::Deserialize;
32use worker::Result;
33use worker::wasm_bindgen::JsValue;
34
35use crate::Work;
36use crate::rows::PULL_COLUMNS;
37use crate::checks::hash;
38use crate::rows::{CommentRow, PullRow};
39
40/// The most memories one project, or one workspace, keeps (as memory.rs).
41const MAX_PER_SCOPE: u32 = 500;
42/// How sure g1t is of what an agent says it learned.
43const RUN_CONFIDENCE: f64 = 0.6;
44/// Of a person's correction in a review.
45const REVIEW_CONFIDENCE: f64 = 0.7;
46/// Of a merged pull request's decision, summed up from its own words.
47const MERGE_CONFIDENCE: f64 = 0.5;
48/// The longest correction or decision kept, in characters.
49const MAX_SUMMARY_CHARS: usize = 400;
50/// The most files a decision names.
51const MAX_FILES_NAMED: usize = 6;
52/// Who people's and agents' names are not.
53const NOT_PEOPLE: [&str; 3] = ["g1t-agent", "g1t", "agent"];
54
55/// Words that mark a sentence as telling the agent how things are done.
56const CORRECTIVE: [&str; 20] = [
57 "instead",
58 "don't",
59 "do not",
60 "never",
61 "always",
62 "please use",
63 "we use",
64 "use the",
65 "should",
66 "shouldn't",
67 "prefer",
68 "avoid",
69 "must",
70 "rather than",
71 "convention",
72 "not how",
73 "we don't",
74 "make sure",
75 "remember to",
76 "the right way",
77];
78
79#[derive(Deserialize)]
80struct Existing {
81 id: String,
82 status: String,
83 sources: String,
84 confidence: Option<f64>,
85}
86
87#[derive(Deserialize)]
88struct RunTicket {
89 id: String,
90 workspace: String,
91 repo_id: String,
92 number: Option<u32>,
93 token_hash: String,
94}
95
96/// `text` cut to `max` characters at a word, with an ellipsis if cut.
97fn clip(text: &str, max: usize) -> String {
98 let text = text.split_whitespace().collect::<Vec<_>>().join(" ");
99 if text.chars().count() <= max {
100 return text;
101 }
102 let cut: String = text.chars().take(max - 1).collect();
103 let cut = cut.rsplit_once(' ').map_or(cut.as_str(), |(head, _)| head).to_owned();
104 format!("{cut}…")
105}
106
107/// Markdown reduced to plain prose: no code blocks, quotes, headings or
108/// list markers.
109fn prose(markdown: &str) -> Vec<String> {
110 let mut lines = Vec::new();
111 let mut fenced = false;
112 for line in markdown.lines() {
113 let trimmed = line.trim();
114 if trimmed.starts_with("```") || trimmed.starts_with("~~~") {
115 fenced = !fenced;
116 continue;
117 }
118 if fenced || trimmed.starts_with('>') || trimmed.starts_with('#') || trimmed.starts_with("<!--") {
119 lines.push(String::new());
120 continue;
121 }
122 let item = trimmed.trim_start_matches(['-', '*', '+']).trim_start();
123 lines.push(item.to_owned());
124 }
125 lines
126}
127
128/// The sentences of a comment, as plain prose.
129fn sentences(body: &str) -> Vec<String> {
130 let mut out = Vec::new();
131 for line in prose(body) {
132 let mut current = String::new();
133 for c in line.chars() {
134 current.push(c);
135 if matches!(c, '.' | '!' | '?') {
136 let sentence = current.trim().to_owned();
137 if !sentence.is_empty() {
138 out.push(sentence);
139 }
140 current.clear();
141 }
142 }
143 let rest = current.trim();
144 if !rest.is_empty() {
145 out.push(rest.to_owned());
146 }
147 }
148 out
149}
150
151/// Whether a sentence tells the agent how something is done here.
152fn corrective(sentence: &str) -> bool {
153 let lower = format!(" {} ", sentence.to_lowercase());
154 !sentence.trim_end().ends_with('?') && CORRECTIVE.iter().any(|word| lower.contains(word))
155}
156
157/// What a person's review comment teaches, as one convention: its first
158/// sentence that says how things are done, or, for a request for changes,
159/// its first real sentence. None when it teaches nothing reusable: a
160/// question, a thank-you, a nit too short to mean anything.
161pub(crate) fn correction(body: &str, requested_changes: bool) -> Option<String> {
162 let all = sentences(body);
163 let found = all
164 .iter()
165 .find(|sentence| sentence.chars().count() >= 15 && corrective(sentence))
166 .or_else(|| {
167 requested_changes.then(|| {
168 all.iter()
169 .find(|sentence| sentence.chars().count() >= 20 && !sentence.ends_with('?'))
170 })?
171 })?;
172 Some(clip(found, MAX_SUMMARY_CHARS))
173}
174
175/// A merged pull request's decision, from its own words: its title, the
176/// first paragraph of its description (why), and the files it changed.
177/// None when it says no more than its title.
178pub(crate) fn decision(number: u32, title: &str, body: Option<&str>, files: &[String]) -> Option<String> {
179 let body = body.unwrap_or_default();
180 let mut paragraph = String::new();
181 for line in prose(body) {
182 if line.trim().is_empty() {
183 if !paragraph.is_empty() {
184 break;
185 }
186 continue;
187 }
188 // "Summary:" and the like say nothing by themselves.
189 if paragraph.is_empty() && line.trim_end_matches(':').split_whitespace().count() <= 1 {
190 continue;
191 }
192 if !paragraph.is_empty() {
193 paragraph.push(' ');
194 }
195 paragraph.push_str(line.trim());
196 }
197 if paragraph.chars().count() < 30 {
198 return None;
199 }
200 let why = clip(&paragraph, MAX_SUMMARY_CHARS);
201 let mut text = format!("Decided in #{number}, \"{}\": {why}", clip(title, 120));
202 if !files.is_empty() {
203 let named: Vec<&str> = files.iter().take(MAX_FILES_NAMED).map(String::as_str).collect();
204 let more = files.len().saturating_sub(MAX_FILES_NAMED);
205 text.push_str(&format!(
206 " Changed {}{}.",
207 named.join(", "),
208 if more > 0 { format!(" and {more} more") } else { String::new() }
209 ));
210 }
211 Some(clip(&text, MAX_MEMORY_CHARS))
212}
213
214/// The JSON list of sources, with `reference` added once.
215fn with_source(sources: &str, reference: &str) -> Vec<String> {
216 let mut list: Vec<String> = serde_json::from_str(sources).unwrap_or_default();
217 if !list.iter().any(|seen| seen == reference) {
218 list.push(reference.to_owned());
219 }
220 list
221}
222
223/// An item tidied, or why it cannot be kept.
224fn checked(item: &CaptureItem) -> std::result::Result<(String, Option<String>), &'static str> {
225 let text = item.text.split_whitespace().collect::<Vec<_>>().join(" ");
226 if text.chars().count() < 8 {
227 return Err("too short");
228 }
229 if text.chars().count() > MAX_MEMORY_CHARS {
230 return Err("too long");
231 }
232 if secret_in(&text).is_some() {
233 return Err("looks like a secret");
234 }
235 let evidence = item
236 .evidence
237 .as_deref()
238 .map(|evidence| clip(evidence, MAX_EVIDENCE_CHARS))
239 .filter(|evidence| !evidence.is_empty());
240 // Evidence quoting a key would put the key in memory all the same.
241 if evidence.as_deref().is_some_and(|evidence| secret_in(evidence).is_some()) {
242 return Err("looks like a secret");
243 }
244 Ok((text, evidence))
245}
246
247impl Work {
248 /// The `owner/name` of a repository, by id, remembered for one capture.
249 async fn path_of(&self, cache: &mut HashMap<String, Option<RepoPath>>, repo_id: &str) -> Result<Option<RepoPath>> {
250 if let Some(found) = cache.get(repo_id) {
251 return Ok(found.clone());
252 }
253 let found: Option<RepoPath> =
254 g1t_kit::call(&self.repos, "path_by_id", &PathByIdArgs { id: repo_id.to_owned() }).await?;
255 cache.insert(repo_id.to_owned(), found.clone());
256 Ok(found)
257 }
258
259 /// Tells subscribers (the context service's search) that a memory
260 /// changed. Never fails what changed it.
261 pub(crate) async fn memory_changed(&self, id: &str, workspace: &str, status: &str, repo_id: Option<&str>) {
262 let event = NewEvent {
263 kind: "memory.changed",
264 source: "work",
265 repo_id: repo_id.map(str::to_owned),
266 actor: None,
267 data: MemoryChanged {
268 memory_id: id.to_owned(),
269 workspace: workspace.to_owned(),
270 status: status.to_owned(),
271 },
272 };
273 let sent: Result<()> = g1t_kit::call(&self.events, "publish", &Publish { events: vec![event] }).await;
274 if let Err(error) = sent {
275 worker::console_error!("memory.changed for {id} was not published: {error}");
276 }
277 }
278
279 /// Adds candidates, or counts them as another sighting of what is there.
280 pub(crate) async fn capture(&self, workspace: &str, items: &[CaptureItem], by: &str) -> Result<Captured> {
281 let workspace = workspace.to_lowercase();
282 let mut done = Captured::default();
283 let mut paths = HashMap::new();
284 for item in items.iter().take(MAX_CAPTURE) {
285 let Ok((text, evidence)) = checked(item) else {
286 done.refused += 1;
287 continue;
288 };
289 let (key, repo) = match (item.scope, &item.repo_id) {
290 (MemoryScope::Workspace, _) => (workspace.clone(), item.repo_id.clone()),
291 (MemoryScope::Project, Some(repo_id)) => (repo_id.clone(), Some(repo_id.clone())),
292 (MemoryScope::Project, None) => {
293 done.refused += 1;
294 continue;
295 }
296 };
297 let path = match &repo {
298 Some(repo_id) => self.path_of(&mut paths, repo_id).await?,
299 None => None,
300 };
301 // A project's memory belongs to a repository of this workspace.
302 if item.scope == MemoryScope::Project
303 && !path.as_ref().is_some_and(|path| path.namespace.to_lowercase() == workspace)
304 {
305 done.refused += 1;
306 continue;
307 }
308 let print = fingerprint(&text);
309 let now = rfc3339(now_ms());
310 let existing = self
311 .db
312 .prepare(
313 "SELECT id, status, sources, confidence FROM memories
314 WHERE scope = ? AND scope_key = ? AND (fingerprint = ? OR text = ?) LIMIT 1",
315 )
316 .bind(&[item.scope.as_str().into(), key.as_str().into(), print.as_str().into(), text.as_str().into()])?
317 .first::<Existing>(None)
318 .await?;
319 if let Some(existing) = existing {
320 done.merged += 1;
321 if existing.status == MemoryStatus::Dismissed.as_str() {
322 continue;
323 }
324 let sources = with_source(&existing.sources, &item.reference);
325 let confidence = match (existing.confidence, item.confidence) {
326 (Some(a), Some(b)) => Some(a.max(b)),
327 (a, b) => a.or(b),
328 };
329 let keep = existing.status == MemoryStatus::Candidate.as_str()
330 && promotes(sources.len(), item.source, confidence);
331 if keep {
332 done.kept += 1;
333 }
334 let status = if keep { MemoryStatus::Kept.as_str() } else { existing.status.as_str() };
335 self.db
336 .prepare(
337 "UPDATE memories SET sources = ?, confidence = ?, status = ?, fingerprint = ?, updated_at = ?
338 WHERE id = ?",
339 )
340 .bind(&[
341 serde_json::to_string(&sources)?.into(),
342 confidence.map_or(JsValue::NULL, JsValue::from),
343 status.into(),
344 print.as_str().into(),
345 now.as_str().into(),
346 existing.id.as_str().into(),
347 ])?
348 .run()
349 .await?;
350 if keep {
351 self.memory_changed(&existing.id, &workspace, status, repo.as_deref()).await;
352 }
353 continue;
354 }
355 let count = self
356 .db
357 .prepare("SELECT count(*) AS value FROM memories WHERE scope = ? AND scope_key = ? AND status != 'dismissed'")
358 .bind(&[item.scope.as_str().into(), key.as_str().into()])?
359 .first::<u32>(Some("value"))
360 .await?
361 .unwrap_or_default();
362 if count >= MAX_PER_SCOPE {
363 done.refused += 1;
364 continue;
365 }
366 let keep = promotes(1, item.source, item.confidence);
367 let status = if keep { MemoryStatus::Kept } else { MemoryStatus::Candidate };
368 let id = new_id("mem", now_ms());
369 let repo_text = path.as_ref().map(|path| format!("{}/{}", path.namespace, path.name));
370 self.db
371 .prepare(
372 "INSERT INTO memories
373 (id, scope, scope_key, workspace, repo, text, kind, source_kind, source_run,
374 source_repo, source_number, created_by, pinned, created_at, updated_at,
375 status, confidence, source_ref, evidence, fingerprint, sources)
376 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 0, ?, ?, ?, ?, ?, ?, ?, ?)",
377 )
378 .bind(&[
379 id.as_str().into(),
380 item.scope.as_str().into(),
381 key.as_str().into(),
382 workspace.as_str().into(),
383 if item.scope == MemoryScope::Project { repo_text.clone().map_or(JsValue::NULL, JsValue::from) } else { JsValue::NULL },
384 text.as_str().into(),
385 item.kind.as_str().into(),
386 item.source.as_str().into(),
387 item.run_id.as_deref().map_or(JsValue::NULL, JsValue::from),
388 repo_text.map_or(JsValue::NULL, JsValue::from),
389 item.number.map_or(JsValue::NULL, JsValue::from),
390 by.into(),
391 now.as_str().into(),
392 now.as_str().into(),
393 status.as_str().into(),
394 item.confidence.map_or(JsValue::NULL, JsValue::from),
395 item.reference.as_str().into(),
396 evidence.map_or(JsValue::NULL, JsValue::from),
397 print.as_str().into(),
398 serde_json::to_string(&[&item.reference])?.into(),
399 ])?
400 .run()
401 .await?;
402 done.added += 1;
403 if keep {
404 done.kept += 1;
405 }
406 self.memory_changed(&id, &workspace, status.as_str(), repo.as_deref()).await;
407 }
408 Ok(done)
409 }
410
411 pub(crate) async fn capture_memories(&self, a: CaptureMemoriesArgs) -> Result<Captured> {
412 self.capture(&a.workspace, &a.items, a.by.as_deref().unwrap_or("g1t")).await
413 }
414
415 pub(crate) async fn report_learned(&self, a: ReportLearnedArgs) -> Result<Outcome<Captured>> {
416 let run = self
417 .db
418 .prepare("SELECT id, workspace, repo_id, number, token_hash FROM agent_runs WHERE id = ?")
419 .bind(&[a.run_id.as_str().into()])?
420 .first::<RunTicket>(None)
421 .await?
422 .filter(|run| !a.token.is_empty() && run.token_hash == hash(&a.token));
423 let Some(run) = run else {
424 return Ok(Outcome::fail(FailureCode::NotFound, "Run not found."));
425 };
426 let items: Vec<CaptureItem> = a
427 .items
428 .iter()
429 .take(MAX_LEARNED)
430 .map(|learned| CaptureItem {
431 scope: if learned.scope.as_deref() == Some("workspace") {
432 MemoryScope::Workspace
433 } else {
434 MemoryScope::Project
435 },
436 repo_id: Some(run.repo_id.clone()),
437 kind: learned.kind.as_deref().and_then(MemoryKind::parse).unwrap_or_default(),
438 text: learned.text.clone(),
439 confidence: Some(RUN_CONFIDENCE),
440 source: CaptureSource::Run,
441 reference: format!("run:{}", run.id),
442 evidence: learned.evidence.clone(),
443 number: run.number,
444 run_id: Some(run.id.clone()),
445 })
446 .collect();
447 Ok(Outcome::Ok(self.capture(&run.workspace, &items, "g1t-agent").await?))
448 }
449
450 pub(crate) async fn list_candidates(&self, a: ListCandidatesArgs) -> Result<Outcome<Vec<Memory>>> {
451 let workspace = a.workspace.to_lowercase();
452 if !a.viewer.as_ref().is_some_and(|viewer| viewer.is_member(&workspace)) {
453 return Ok(Outcome::fail(FailureCode::Forbidden, "Memory is for members of the workspace and its agents."));
454 }
455 let statement = match &a.repo {
456 Some(path) => {
457 let repo = match self.repo(path, &a.viewer).await? {
458 Outcome::Ok(repo) => repo,
459 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
460 };
461 // The project's own, and the workspace's learned there.
462 self.db
463 .prepare(
464 "SELECT * FROM memories WHERE workspace = ? AND status = 'candidate'
465 AND (scope_key = ? OR (scope = 'workspace' AND source_repo = ?))
466 ORDER BY updated_at DESC LIMIT 200",
467 )
468 .bind(&[
469 workspace.as_str().into(),
470 repo.id.as_str().into(),
471 format!("{}/{}", repo.namespace, repo.name).into(),
472 ])?
473 }
474 None => self
475 .db
476 .prepare("SELECT * FROM memories WHERE workspace = ? AND status = 'candidate' ORDER BY updated_at DESC LIMIT 200")
477 .bind(&[workspace.as_str().into()])?,
478 };
479 let rows = statement.all().await?.results::<crate::memory::MemoryRow>()?;
480 Ok(Outcome::Ok(rows.into_iter().map(Memory::from).collect()))
481 }
482
483 pub(crate) async fn review_memory(&self, a: ReviewMemoryArgs) -> Result<Outcome<Memory>> {
484 let workspace = a.workspace.to_lowercase();
485 if !a.actor.is_member(&workspace) || a.actor.kind != PrincipalKind::User || !a.actor.verified {
486 return Ok(Outcome::fail(FailureCode::Forbidden, "Only members of the workspace can review its memory."));
487 }
488 let Some(row) = self
489 .db
490 .prepare("SELECT * FROM memories WHERE id = ? AND workspace = ?")
491 .bind(&[a.id.as_str().into(), workspace.as_str().into()])?
492 .first::<crate::memory::MemoryRow>(None)
493 .await?
494 else {
495 return Ok(Outcome::fail(FailureCode::NotFound, "Memory not found."));
496 };
497 let text = match a.text.as_deref().map(str::trim).filter(|text| !text.is_empty()) {
498 Some(text) => {
499 if text.chars().count() > MAX_MEMORY_CHARS {
500 return Ok(Outcome::fail(FailureCode::Invalid, format!("Keep a memory to {MAX_MEMORY_CHARS} characters.")));
501 }
502 if let Some(what) = secret_in(text) {
503 return Ok(Outcome::fail(
504 FailureCode::Invalid,
505 format!("That looks like {what}. Memory never holds secrets; say where the secret lives instead."),
506 ));
507 }
508 Some(text.to_owned())
509 }
510 None => None,
511 };
512 let status = match a.decision {
513 ReviewDecision::Keep => MemoryStatus::Kept,
514 ReviewDecision::Dismiss => MemoryStatus::Dismissed,
515 };
516 let now = rfc3339(now_ms());
517 let print = text.as_deref().map(fingerprint);
518 self.db
519 .prepare(
520 "UPDATE memories SET status = ?, text = COALESCE(?, text), fingerprint = COALESCE(?, fingerprint),
521 kind = COALESCE(?, kind), reviewed_by = ?, reviewed_at = ?, updated_at = ?
522 WHERE id = ? AND workspace = ?",
523 )
524 .bind(&[
525 status.as_str().into(),
526 text.map_or(JsValue::NULL, JsValue::from),
527 print.map_or(JsValue::NULL, JsValue::from),
528 a.kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
529 a.actor.username.as_str().into(),
530 now.as_str().into(),
531 now.as_str().into(),
532 a.id.as_str().into(),
533 workspace.as_str().into(),
534 ])?
535 .run()
536 .await?;
537 let repo_id = (row.scope == "project").then(|| row.scope_key.clone());
538 self.memory_changed(&a.id, &workspace, status.as_str(), repo_id.as_deref()).await;
539 let reviewed = self
540 .db
541 .prepare("SELECT * FROM memories WHERE id = ?")
542 .bind(&[a.id.as_str().into()])?
543 .first::<crate::memory::MemoryRow>(None)
544 .await?;
545 Ok(match reviewed {
546 Some(row) => Outcome::Ok(Memory::from(row)),
547 None => Outcome::fail(FailureCode::NotFound, "Memory not found."),
548 })
549 }
550
551 pub(crate) async fn memories_by_id(&self, a: MemoriesByIdArgs) -> Result<Vec<Memory>> {
552 let ids: Vec<&String> = a.ids.iter().take(100).collect();
553 let rows = self
554 .db
555 .prepare("SELECT * FROM memories WHERE workspace = ? AND id IN (SELECT value FROM json_each(?))")
556 .bind(&[a.workspace.to_lowercase().into(), serde_json::to_string(&ids)?.into()])?
557 .all()
558 .await?
559 .results::<crate::memory::MemoryRow>()?;
560 Ok(rows.into_iter().map(Memory::from).collect())
561 }
562
563 pub(crate) async fn search_memories(&self, a: SearchMemoriesArgs) -> Result<Vec<Memory>> {
564 let status = a.status.unwrap_or(MemoryStatus::Kept);
565 let limit = a.limit.unwrap_or(20).clamp(1, 100) as usize;
566 let statement = match &a.repo_ids {
567 Some(ids) => self
568 .db
569 .prepare(
570 "SELECT * FROM memories WHERE workspace = ? AND status = ?
571 AND (scope = 'workspace' OR scope_key IN (SELECT value FROM json_each(?)))
572 ORDER BY pinned DESC, COALESCE(last_used_at, updated_at) DESC LIMIT 1000",
573 )
574 .bind(&[a.workspace.to_lowercase().into(), status.as_str().into(), serde_json::to_string(ids)?.into()])?,
575 None => self
576 .db
577 .prepare(
578 "SELECT * FROM memories WHERE workspace = ? AND status = ?
579 ORDER BY pinned DESC, COALESCE(last_used_at, updated_at) DESC LIMIT 1000",
580 )
581 .bind(&[a.workspace.to_lowercase().into(), status.as_str().into()])?,
582 };
583 let words: Vec<String> = a
584 .query
585 .as_deref()
586 .unwrap_or_default()
587 .split_whitespace()
588 .map(str::to_lowercase)
589 .collect();
590 Ok(statement
591 .all()
592 .await?
593 .results::<crate::memory::MemoryRow>()?
594 .into_iter()
595 .map(Memory::from)
596 .filter(|memory| {
597 let text = memory.text.to_lowercase();
598 words.iter().all(|word| text.contains(word.as_str()))
599 })
600 .take(limit)
601 .collect())
602 }
603
604 /// What a merged pull request decided, as a candidate.
605 fn decision_item(pull: &Pull) -> Option<CaptureItem> {
606 let files: Vec<String> = pull.files.iter().map(|file| file.path.clone()).collect();
607 let text = decision(pull.number, &pull.title, pull.body.as_deref(), &files)?;
608 Some(CaptureItem {
609 scope: MemoryScope::Project,
610 repo_id: Some(pull.repo_id.clone()),
611 kind: MemoryKind::Decision,
612 text,
613 confidence: Some(MERGE_CONFIDENCE),
614 source: CaptureSource::Pr,
615 reference: format!("pull:{}#{}", pull.repo_id, pull.number),
616 evidence: Some(clip(&pull.title, MAX_EVIDENCE_CHARS)),
617 number: Some(pull.number),
618 run_id: None,
619 })
620 }
621
622 /// What a person's comment on a pull request corrected, as a candidate.
623 fn correction_item(pull: &Pull, comment: &Comment) -> Option<CaptureItem> {
624 if comment.kind != CommentKind::Comment
625 || comment.author.kind == PrincipalKind::Agent
626 || NOT_PEOPLE.contains(&comment.author.username.as_str())
627 {
628 return None;
629 }
630 let text = correction(&comment.body, comment.verdict == Some(Verdict::RequestChanges))?;
631 Some(CaptureItem {
632 scope: MemoryScope::Project,
633 repo_id: Some(pull.repo_id.clone()),
634 kind: MemoryKind::Convention,
635 text,
636 confidence: Some(REVIEW_CONFIDENCE),
637 source: CaptureSource::Review,
638 reference: format!("comment:{}", comment.id),
639 evidence: Some(format!(
640 "{} on #{}{}: {}",
641 comment.author.username,
642 pull.number,
643 comment.path.as_deref().map(|path| format!(" ({path})")).unwrap_or_default(),
644 clip(&comment.body, MAX_EVIDENCE_CHARS - 60)
645 )),
646 number: Some(pull.number),
647 run_id: None,
648 })
649 }
650
651 async fn workspace_of(&self, repo_id: &str) -> Result<Option<String>> {
652 let mut cache = HashMap::new();
653 Ok(self.path_of(&mut cache, repo_id).await?.map(|path| path.namespace.to_lowercase()))
654 }
655
656 pub(crate) async fn seed_from_pulls(&self, a: SeedFromPullsArgs) -> Result<Captured> {
657 let Some(workspace) = self.workspace_of(&a.repo_id).await? else {
658 return Ok(Captured::default());
659 };
660 let limit = a.limit.unwrap_or(20).clamp(1, 50);
661 let pulls: Vec<Pull> = self
662 .db
663 .prepare(format!("SELECT {PULL_COLUMNS} FROM pulls WHERE repo_id = ? AND status = 'merged' ORDER BY merged_at DESC LIMIT ?"))
664 .bind(&[a.repo_id.as_str().into(), limit.into()])?
665 .all()
666 .await?
667 .results::<PullRow>()?
668 .into_iter()
669 .map(Pull::from)
670 .collect();
671 let mut items = Vec::new();
672 for pull in &pulls {
673 items.extend(Self::decision_item(pull));
674 let comments = self
675 .db
676 .prepare("SELECT * FROM comments WHERE repo_id = ? AND number = ? ORDER BY id LIMIT 100")
677 .bind(&[pull.repo_id.as_str().into(), pull.number.into()])?
678 .all()
679 .await?
680 .results::<CommentRow>()?;
681 items.extend(
682 comments
683 .into_iter()
684 .map(Comment::from)
685 .filter_map(|comment| Self::correction_item(pull, &comment)),
686 );
687 }
688 let mut done = Captured::default();
689 for chunk in items.chunks(MAX_CAPTURE) {
690 let part = self.capture(&workspace, chunk, "g1t").await?;
691 done.added += part.added;
692 done.merged += part.merged;
693 done.kept += part.kept;
694 done.refused += part.refused;
695 }
696 Ok(done)
697 }
698
699 /// Captures from the bus: corrections in review, decisions on merge.
700 async fn capture_event(&self, event: &Event) -> Result<()> {
701 let field = |name: &str| event.data[name].as_str().map(str::to_owned);
702 match event.kind.as_str() {
703 "pull.merged" => {
704 let Some(pull) = (match field("pullId") {
705 Some(id) => self.pull_by_id(&id).await?,
706 None => None,
707 }) else {
708 return Ok(());
709 };
710 if pull.status != PullStatus::Merged {
711 return Ok(());
712 }
713 if let (Some(item), Some(workspace)) = (Self::decision_item(&pull), self.workspace_of(&pull.repo_id).await?) {
714 self.capture(&workspace, &[item], "g1t").await?;
715 }
716 }
717 "comment.created" => {
718 let (Some(comment_id), Some(pull_id)) = (field("commentId"), field("pullId")) else {
719 return Ok(());
720 };
721 let Some(pull) = self.pull_by_id(&pull_id).await? else {
722 return Ok(());
723 };
724 let Some(comment) = self
725 .db
726 .prepare("SELECT * FROM comments WHERE id = ?")
727 .bind(&[comment_id.as_str().into()])?
728 .first::<CommentRow>(None)
729 .await?
730 .map(Comment::from)
731 else {
732 return Ok(());
733 };
734 if let (Some(item), Some(workspace)) =
735 (Self::correction_item(&pull, &comment), self.workspace_of(&pull.repo_id).await?)
736 {
737 self.capture(&workspace, &[item], "g1t").await?;
738 }
739 }
740 _ => {}
741 }
742 Ok(())
743 }
744}
745
746/// The methods this module serves at `POST /rpc/<method>`.
747pub(crate) const METHODS: [&str; 7] = [
748 "capture_memories",
749 "report_learned",
750 "list_candidates",
751 "review_memory",
752 "memories_by_id",
753 "search_memories",
754 "seed_from_pulls",
755];
756
757pub(crate) async fn dispatch(work: &Work, method: &str, body: serde_json::Value) -> Result<worker::Response> {
758 use g1t_kit::{args, reply};
759 match method {
760 "capture_memories" => reply(&work.capture_memories(args(body)?).await?),
761 "report_learned" => reply(&work.report_learned(args(body)?).await?),
762 "list_candidates" => reply(&work.list_candidates(args(body)?).await?),
763 "review_memory" => reply(&work.review_memory(args(body)?).await?),
764 "memories_by_id" => reply(&work.memories_by_id(args(body)?).await?),
765 "search_memories" => reply(&work.search_memories(args(body)?).await?),
766 "seed_from_pulls" => reply(&work.seed_from_pulls(args(body)?).await?),
767 _ => worker::Response::error("Unknown method", 404),
768 }
769}
770
771/// Captures what an event teaches. Never holds up the queue: a failure is
772/// logged and the event goes on to its other handlers.
773pub(crate) async fn on_event(work: &Work, event: &Event) {
774 if let Err(error) = work.capture_event(event).await {
775 worker::console_error!("capturing memory from {} {} failed: {error}", event.kind, event.id);
776 }
777}
778
779#[cfg(test)]
780mod tests {
781 use super::*;
782
783 #[test]
784 fn a_correction_keeps_the_sentence_that_says_how() {
785 let body = "Thanks for this! We use the shared `api` client instead of calling fetch directly. Otherwise looks fine.";
786 assert_eq!(
787 correction(body, false).as_deref(),
788 Some("We use the shared `api` client instead of calling fetch directly.")
789 );
790 }
791
792 #[test]
793 fn questions_and_praise_teach_nothing() {
794 assert_eq!(correction("Why not use the cache here?", false), None);
795 assert_eq!(correction("LGTM, nice work.", false), None);
796 assert_eq!(correction("> never do this\n\nok", false), None);
797 assert_eq!(correction("```\nnever run this\n```", false), None);
798 }
799
800 #[test]
801 fn a_request_for_changes_without_a_marker_keeps_its_first_sentence() {
802 assert_eq!(
803 correction("The migration has to be idempotent for D1 replays.", true).as_deref(),
804 Some("The migration has to be idempotent for D1 replays.")
805 );
806 assert_eq!(correction("The migration has to be idempotent for D1 replays.", false), None);
807 }
808
809 #[test]
810 fn a_decision_names_why_and_the_files() {
811 let files: Vec<String> = (1..=8).map(|n| format!("src/f{n}.rs")).collect();
812 let text = decision(
813 12,
814 "Keep the v1 webhook payload",
815 Some("## Summary\n\nTwo customers still parse the v1 payload, so it stays alongside v2.\n\nMore detail."),
816 &files,
817 )
818 .unwrap();
819 assert!(text.starts_with("Decided in #12, \"Keep the v1 webhook payload\": Two customers still parse"));
820 assert!(text.contains("src/f1.rs") && text.contains("and 2 more"));
821 assert!(!text.contains("More detail"));
822 }
823
824 #[test]
825 fn a_pull_request_that_says_only_its_title_decides_nothing() {
826 assert_eq!(decision(3, "Fix typo", Some("Fix typo."), &[]), None);
827 assert_eq!(decision(3, "Fix typo", None, &[]), None);
828 }
829
830 #[test]
831 fn sources_are_counted_once_each() {
832 assert_eq!(with_source("[\"run:a\"]", "run:a"), ["run:a"]);
833 assert_eq!(with_source("[\"run:a\"]", "comment:b"), ["run:a", "comment:b"]);
834 assert_eq!(with_source("not json", "doc:x"), ["doc:x"]);
835 }
836
837 fn item(text: &str, evidence: Option<&str>) -> CaptureItem {
838 CaptureItem {
839 scope: MemoryScope::Project,
840 repo_id: Some("rep_1".into()),
841 kind: MemoryKind::Gotcha,
842 text: text.into(),
843 confidence: Some(0.6),
844 source: CaptureSource::Run,
845 reference: "run:1".into(),
846 evidence: evidence.map(str::to_owned),
847 number: None,
848 run_id: None,
849 }
850 }
851
852 #[test]
853 fn secrets_are_refused_on_capture_in_text_and_evidence() {
854 // Key-shaped values are joined at run time, as in the contracts' tests.
855 let key = format!("{}{}", "gh", "p_abcdefghijklmnopqrstuvwxyz0123456789");
856 assert!(checked(&item(&format!("Clone with {key} for the mirror."), None)).is_err());
857 assert!(checked(&item("The mirror needs a token from the MIRROR_TOKEN secret.", Some(&format!("export T={key}")))).is_err());
858 let (text, evidence) = checked(&item(" Tests need TZ=UTC. ", Some("cargo test failed on dates"))).unwrap();
859 assert_eq!(text, "Tests need TZ=UTC.");
860 assert_eq!(evidence.as_deref(), Some("cargo test failed on dates"));
861 assert!(checked(&item("short", None)).is_err());
862 }
863}