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