pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/capture.rs

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