Skip to main content

ag_store/
review.rs

1//! Session review-request persistence adapters and query helpers.
2
3use ag_session::ReviewRequest;
4use async_trait::async_trait;
5use sqlx::{SqliteConnection, SqlitePool};
6
7use crate::DbError;
8
9/// Durable input for one review-comment resolution operation.
10#[derive(Clone, Debug, Eq, PartialEq)]
11pub struct NewSessionReviewCommentResolution {
12    /// Full commit hash containing the reported fix, once auto-commit succeeds.
13    pub commit_hash: Option<String>,
14    /// Original agent-authored reply retained across retries.
15    pub reply: String,
16    /// Unguessable token embedded in the posted reply.
17    pub reply_token: String,
18    /// Persisted normalized resolution decision.
19    pub resolution: String,
20    /// Forge-native review-request identifier such as `#123` or `!123`.
21    pub review_request_display_id: String,
22    /// Forge-native review-thread identifier.
23    pub thread_id: String,
24}
25
26/// Row returned for one unfinished review-comment resolution operation.
27#[derive(Clone, Debug, Eq, PartialEq)]
28pub struct SessionReviewCommentResolutionRow {
29    /// Full commit hash that must remain reachable from the pushed branch.
30    pub commit_hash: Option<String>,
31    /// Whether a forge reply attempt may already have exposed the token.
32    pub is_posting: bool,
33    /// Original agent-authored reply retained across retries.
34    pub reply: String,
35    /// Unguessable token embedded in the posted reply.
36    pub reply_token: String,
37    /// Persisted normalized resolution decision.
38    pub resolution: String,
39    /// Forge-native review-request identifier such as `#123` or `!123`.
40    pub review_request_display_id: String,
41    /// Forge-native review-thread identifier.
42    pub thread_id: String,
43}
44
45/// Row returned when loading one `session_review_request`.
46#[derive(Clone, Debug, Eq, PartialEq)]
47pub struct SessionReviewRequestRow {
48    /// Forge-native display identifier such as `#123` or `!123`.
49    pub display_id: String,
50    /// Persisted forge-family discriminator.
51    pub forge_kind: String,
52    /// Most recent successful refresh timestamp in Unix seconds.
53    pub last_refreshed_at: i64,
54    /// Review request source branch.
55    pub source_branch: String,
56    /// Persisted normalized lifecycle state.
57    pub state: String,
58    /// Optional normalized checks or merge-status summary.
59    pub status_summary: Option<String>,
60    /// Review request target branch.
61    pub target_branch: String,
62    /// Review request title.
63    pub title: String,
64    /// Browser-openable review request URL.
65    pub web_url: String,
66}
67
68/// Review-request persistence boundary used by app orchestration and tests.
69#[cfg_attr(test, mockall::automock)]
70#[async_trait]
71pub trait ReviewRepository: Send + Sync {
72    /// Binds newly inserted operations to the commit produced for their turn.
73    async fn bind_session_review_comment_resolutions_to_commit(
74        &self,
75        id: &str,
76        resolutions: &[NewSessionReviewCommentResolution],
77        commit_hash: &str,
78    ) -> Result<(), DbError>;
79
80    /// Discards unfinished operations created with the supplied reply tokens.
81    async fn discard_session_review_comment_resolutions(
82        &self,
83        id: &str,
84        resolutions: &[NewSessionReviewCommentResolution],
85    ) -> Result<(), DbError>;
86
87    /// Inserts one atomic set of review-comment resolution operations.
88    ///
89    /// A bound operation for the same thread wins. A fresh agent turn may
90    /// replace an unbound operation left by interrupted commit binding.
91    async fn insert_session_review_comment_resolutions(
92        &self,
93        id: &str,
94        resolutions: &[NewSessionReviewCommentResolution],
95    ) -> Result<(), DbError>;
96
97    /// Loads unfinished review-comment resolution operations for a session.
98    async fn load_session_review_comment_resolutions(
99        &self,
100        id: &str,
101    ) -> Result<Vec<SessionReviewCommentResolutionRow>, DbError>;
102
103    /// Loads the persisted forge review-request linkage for a session.
104    async fn load_session_review_request(
105        &self,
106        id: &str,
107    ) -> Result<Option<SessionReviewRequestRow>, DbError>;
108
109    /// Updates the persisted forge review-request linkage for a session.
110    async fn update_session_review_request(
111        &self,
112        id: &str,
113        review_request: Option<ReviewRequest>,
114    ) -> Result<(), DbError>;
115
116    /// Records that one operation is about to attempt its forge reply.
117    async fn mark_session_review_comment_resolution_posting(
118        &self,
119        id: &str,
120        reply_token: &str,
121    ) -> Result<(), DbError>;
122
123    /// Removes one operation after its requested forge effects finish.
124    async fn remove_session_review_comment_resolution(
125        &self,
126        id: &str,
127        reply_token: &str,
128    ) -> Result<(), DbError>;
129}
130
131/// `SQLite` implementation of [`ReviewRepository`].
132#[derive(Clone)]
133pub(crate) struct SqliteReviewRepository(SqlitePool);
134
135impl SqliteReviewRepository {
136    /// Creates a review repository backed by the provided pool.
137    pub(crate) fn new(pool: SqlitePool) -> Self {
138        Self(pool)
139    }
140}
141
142/// Inserts review-comment resolutions through the caller's transaction.
143pub(crate) async fn insert_review_comment_resolutions(
144    connection: &mut SqliteConnection,
145    session_id: &str,
146    resolutions: &[NewSessionReviewCommentResolution],
147) -> Result<(), DbError> {
148    for resolution in resolutions {
149        sqlx::query!(
150            r"
151INSERT INTO session_review_comment_resolution (
152    session_id,
153    commit_hash,
154    review_request_display_id,
155    thread_id,
156    reply,
157    reply_token,
158    resolution
159)
160VALUES (?, ?, ?, ?, ?, ?, ?)
161ON CONFLICT(session_id, review_request_display_id, thread_id)
162DO UPDATE SET
163    commit_hash = excluded.commit_hash,
164    reply = excluded.reply,
165    reply_token = excluded.reply_token,
166    resolution = excluded.resolution,
167    is_posting = 0
168WHERE session_review_comment_resolution.commit_hash IS NULL
169",
170            session_id,
171            resolution.commit_hash,
172            resolution.review_request_display_id,
173            resolution.thread_id,
174            resolution.reply,
175            resolution.reply_token,
176            resolution.resolution
177        )
178        .execute(&mut *connection)
179        .await?;
180    }
181
182    Ok(())
183}
184
185#[async_trait]
186impl ReviewRepository for SqliteReviewRepository {
187    async fn bind_session_review_comment_resolutions_to_commit(
188        &self,
189        id: &str,
190        resolutions: &[NewSessionReviewCommentResolution],
191        commit_hash: &str,
192    ) -> Result<(), DbError> {
193        let mut transaction = self.0.begin().await?;
194        for resolution in resolutions {
195            sqlx::query!(
196                r"
197UPDATE session_review_comment_resolution
198SET commit_hash = ?
199WHERE session_id = ?
200  AND reply_token = ?
201  AND commit_hash IS NULL
202",
203                commit_hash,
204                id,
205                resolution.reply_token
206            )
207            .execute(&mut *transaction)
208            .await?;
209        }
210        transaction.commit().await?;
211
212        Ok(())
213    }
214
215    async fn discard_session_review_comment_resolutions(
216        &self,
217        id: &str,
218        resolutions: &[NewSessionReviewCommentResolution],
219    ) -> Result<(), DbError> {
220        let mut transaction = self.0.begin().await?;
221        for resolution in resolutions {
222            sqlx::query!(
223                r"
224DELETE FROM session_review_comment_resolution
225WHERE session_id = ?
226  AND reply_token = ?
227",
228                id,
229                resolution.reply_token
230            )
231            .execute(&mut *transaction)
232            .await?;
233        }
234        transaction.commit().await?;
235
236        Ok(())
237    }
238
239    async fn insert_session_review_comment_resolutions(
240        &self,
241        id: &str,
242        resolutions: &[NewSessionReviewCommentResolution],
243    ) -> Result<(), DbError> {
244        let mut transaction = self.0.begin().await?;
245        insert_review_comment_resolutions(&mut transaction, id, resolutions).await?;
246        transaction.commit().await?;
247
248        Ok(())
249    }
250
251    async fn load_session_review_comment_resolutions(
252        &self,
253        id: &str,
254    ) -> Result<Vec<SessionReviewCommentResolutionRow>, DbError> {
255        let resolutions = sqlx::query_as!(
256            SessionReviewCommentResolutionRow,
257            r#"
258SELECT is_posting AS "is_posting: bool",
259       commit_hash,
260       reply,
261       reply_token,
262       resolution,
263       review_request_display_id,
264       thread_id
265FROM session_review_comment_resolution
266WHERE session_id = ?
267ORDER BY rowid
268"#,
269            id
270        )
271        .fetch_all(&self.0)
272        .await?;
273
274        Ok(resolutions)
275    }
276
277    async fn load_session_review_request(
278        &self,
279        id: &str,
280    ) -> Result<Option<SessionReviewRequestRow>, DbError> {
281        let review_request = sqlx::query_as!(
282            SessionReviewRequestRow,
283            r"
284SELECT display_id,
285       forge_kind,
286       last_refreshed_at,
287       source_branch,
288       state,
289       status_summary,
290       target_branch,
291       title,
292       web_url
293FROM session_review_request
294WHERE session_id = ?
295",
296            id
297        )
298        .fetch_optional(&self.0)
299        .await?;
300
301        Ok(review_request)
302    }
303
304    async fn update_session_review_request(
305        &self,
306        id: &str,
307        review_request: Option<ReviewRequest>,
308    ) -> Result<(), DbError> {
309        if let Some(review_request) = review_request.as_ref() {
310            sqlx::query!(
311                r"
312INSERT INTO session_review_request (
313    session_id,
314    display_id,
315    forge_kind,
316    last_refreshed_at,
317    source_branch,
318    state,
319    status_summary,
320    target_branch,
321    title,
322    web_url
323)
324VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
325ON CONFLICT(session_id) DO UPDATE
326SET display_id = excluded.display_id,
327    forge_kind = excluded.forge_kind,
328    last_refreshed_at = excluded.last_refreshed_at,
329    source_branch = excluded.source_branch,
330    state = excluded.state,
331    status_summary = excluded.status_summary,
332    target_branch = excluded.target_branch,
333    title = excluded.title,
334    web_url = excluded.web_url
335",
336                id,
337                review_request.summary.display_id.as_str(),
338                review_request.summary.forge_kind.as_str(),
339                review_request.last_refreshed_at,
340                review_request.summary.source_branch.as_str(),
341                review_request.summary.state.as_str(),
342                review_request.summary.status_summary.as_deref(),
343                review_request.summary.target_branch.as_str(),
344                review_request.summary.title.as_str(),
345                review_request.summary.web_url.as_str()
346            )
347            .execute(&self.0)
348            .await?;
349        } else {
350            sqlx::query!(
351                r"
352DELETE FROM session_review_request
353WHERE session_id = ?
354",
355                id
356            )
357            .execute(&self.0)
358            .await?;
359        }
360
361        Ok(())
362    }
363
364    async fn mark_session_review_comment_resolution_posting(
365        &self,
366        id: &str,
367        reply_token: &str,
368    ) -> Result<(), DbError> {
369        let update_result = sqlx::query!(
370            r"
371UPDATE session_review_comment_resolution
372SET is_posting = 1
373WHERE session_id = ?
374  AND reply_token = ?
375  AND is_posting = 0
376",
377            id,
378            reply_token
379        )
380        .execute(&self.0)
381        .await?;
382        if update_result.rows_affected() != 1 {
383            return Err(DbError::InvalidData {
384                entity: "review-comment operation",
385                reason: format!("posting update matched no pending row for token `{reply_token}`"),
386            });
387        }
388
389        Ok(())
390    }
391
392    async fn remove_session_review_comment_resolution(
393        &self,
394        id: &str,
395        reply_token: &str,
396    ) -> Result<(), DbError> {
397        let delete_result = sqlx::query!(
398            r"
399DELETE FROM session_review_comment_resolution
400WHERE session_id = ?
401  AND reply_token = ?
402",
403            id,
404            reply_token
405        )
406        .execute(&self.0)
407        .await?;
408        if delete_result.rows_affected() != 1 {
409            return Err(DbError::InvalidData {
410                entity: "review-comment operation",
411                reason: format!("delete matched no row for token `{reply_token}`"),
412            });
413        }
414
415        Ok(())
416    }
417}
418
419#[cfg(test)]
420mod tests {
421    use super::*;
422    use crate::AppRepositories;
423
424    #[tokio::test]
425    async fn active_review_comment_operation_preserves_original_reply_across_retries() {
426        // Arrange
427        let repositories = AppRepositories::in_memory()
428            .await
429            .expect("failed to open in-memory repositories");
430        let project_id = repositories
431            .projects()
432            .upsert_project("/tmp/project", Some("main".to_string()))
433            .await
434            .expect("failed to insert project");
435        repositories
436            .sessions()
437            .insert_session("session-id", "codex", "main", "Review", project_id)
438            .await
439            .expect("failed to insert session");
440        let original = review_comment_resolution("Original reply", "token-1");
441        let regenerated = review_comment_resolution("Regenerated reply", "token-2");
442        repositories
443            .reviews()
444            .insert_session_review_comment_resolutions(
445                "session-id",
446                std::slice::from_ref(&original),
447            )
448            .await
449            .expect("failed to insert original operation");
450        repositories
451            .reviews()
452            .bind_session_review_comment_resolutions_to_commit(
453                "session-id",
454                std::slice::from_ref(&original),
455                "commit-original",
456            )
457            .await
458            .expect("failed to bind original operation");
459
460        // Act
461        repositories
462            .reviews()
463            .insert_session_review_comment_resolutions(
464                "session-id",
465                std::slice::from_ref(&regenerated),
466            )
467            .await
468            .expect("failed to ignore regenerated operation");
469        repositories
470            .reviews()
471            .bind_session_review_comment_resolutions_to_commit(
472                "session-id",
473                std::slice::from_ref(&regenerated),
474                "commit-regenerated",
475            )
476            .await
477            .expect("failed to ignore regenerated binding");
478        let active = repositories
479            .reviews()
480            .load_session_review_comment_resolutions("session-id")
481            .await
482            .expect("failed to load active operation");
483        repositories
484            .reviews()
485            .mark_session_review_comment_resolution_posting("session-id", "token-1")
486            .await
487            .expect("failed to mark original operation as posting");
488        let posting = repositories
489            .reviews()
490            .load_session_review_comment_resolutions("session-id")
491            .await
492            .expect("failed to reload posting operation");
493        repositories
494            .reviews()
495            .remove_session_review_comment_resolution("session-id", "token-1")
496            .await
497            .expect("failed to remove original operation");
498        repositories
499            .reviews()
500            .insert_session_review_comment_resolutions("session-id", &[regenerated])
501            .await
502            .expect("failed to insert later operation");
503        let replacement = repositories
504            .reviews()
505            .load_session_review_comment_resolutions("session-id")
506            .await
507            .expect("failed to load replacement operation");
508        let missing_update_error = repositories
509            .reviews()
510            .mark_session_review_comment_resolution_posting("session-id", "missing-token")
511            .await
512            .expect_err("missing operation should reject state update");
513        let missing_delete_error = repositories
514            .reviews()
515            .remove_session_review_comment_resolution("session-id", "missing-token")
516            .await
517            .expect_err("missing operation should reject deletion");
518
519        // Assert
520        assert_eq!(active.len(), 1);
521        assert_eq!(active[0].reply, "Original reply");
522        assert_eq!(active[0].reply_token, "token-1");
523        assert_eq!(active[0].commit_hash.as_deref(), Some("commit-original"));
524        assert!(!active[0].is_posting);
525        assert!(posting[0].is_posting);
526        assert_eq!(replacement.len(), 1);
527        assert_eq!(replacement[0].reply, "Regenerated reply");
528        assert!(matches!(missing_update_error, DbError::InvalidData { .. }));
529        assert!(matches!(missing_delete_error, DbError::InvalidData { .. }));
530    }
531
532    #[tokio::test]
533    async fn discard_failed_retry_preserves_older_conflicting_operation() {
534        // Arrange
535        let repositories = AppRepositories::in_memory()
536            .await
537            .expect("failed to open in-memory repositories");
538        let project_id = repositories
539            .projects()
540            .upsert_project("/tmp/project", Some("main".to_string()))
541            .await
542            .expect("failed to insert project");
543        repositories
544            .sessions()
545            .insert_session("session-id", "codex", "main", "Review", project_id)
546            .await
547            .expect("failed to insert session");
548        let original = review_comment_resolution("Original reply", "token-original");
549        let regenerated = review_comment_resolution("Regenerated reply", "token-regenerated");
550        let mut unrelated = review_comment_resolution("Unrelated reply", "token-unrelated");
551        unrelated.thread_id = "thread-2".to_string();
552        let mut inserted = review_comment_resolution("Inserted reply", "token-inserted");
553        inserted.thread_id = "thread-3".to_string();
554        repositories
555            .reviews()
556            .insert_session_review_comment_resolutions(
557                "session-id",
558                &[original.clone(), unrelated.clone()],
559            )
560            .await
561            .expect("failed to insert review operations");
562        repositories
563            .reviews()
564            .bind_session_review_comment_resolutions_to_commit(
565                "session-id",
566                std::slice::from_ref(&original),
567                "commit-original",
568            )
569            .await
570            .expect("failed to bind original operation");
571        repositories
572            .reviews()
573            .insert_session_review_comment_resolutions(
574                "session-id",
575                &[regenerated.clone(), inserted.clone()],
576            )
577            .await
578            .expect("failed to insert later review operations");
579
580        // Act
581        repositories
582            .reviews()
583            .discard_session_review_comment_resolutions("session-id", &[regenerated, inserted])
584            .await
585            .expect("failed to discard newly inserted review operations");
586        let remaining = repositories
587            .reviews()
588            .load_session_review_comment_resolutions("session-id")
589            .await
590            .expect("failed to load remaining review operations");
591
592        // Assert
593        assert_eq!(remaining.len(), 2);
594        assert_eq!(remaining[0].reply_token, original.reply_token);
595        assert_eq!(remaining[0].reply, original.reply);
596        assert_eq!(remaining[1].reply_token, unrelated.reply_token);
597        assert_eq!(remaining[1].reply, unrelated.reply);
598    }
599
600    #[tokio::test]
601    async fn fresh_retry_replaces_unbound_operation() {
602        // Arrange
603        let repositories = AppRepositories::in_memory()
604            .await
605            .expect("failed to open in-memory repositories");
606        let project_id = repositories
607            .projects()
608            .upsert_project("/tmp/project", Some("main".to_string()))
609            .await
610            .expect("failed to insert project");
611        repositories
612            .sessions()
613            .insert_session("session-id", "codex", "main", "Review", project_id)
614            .await
615            .expect("failed to insert session");
616        let original = review_comment_resolution("Original reply", "token-original");
617        let regenerated = review_comment_resolution("Regenerated reply", "token-regenerated");
618        repositories
619            .reviews()
620            .insert_session_review_comment_resolutions("session-id", &[original])
621            .await
622            .expect("failed to insert unbound operation");
623
624        // Act
625        repositories
626            .reviews()
627            .insert_session_review_comment_resolutions(
628                "session-id",
629                std::slice::from_ref(&regenerated),
630            )
631            .await
632            .expect("failed to replace unbound operation");
633        repositories
634            .reviews()
635            .bind_session_review_comment_resolutions_to_commit(
636                "session-id",
637                std::slice::from_ref(&regenerated),
638                "commit-regenerated",
639            )
640            .await
641            .expect("failed to bind replacement operation");
642        let active = repositories
643            .reviews()
644            .load_session_review_comment_resolutions("session-id")
645            .await
646            .expect("failed to load replacement operation");
647
648        // Assert
649        assert_eq!(active.len(), 1);
650        assert_eq!(active[0].reply, "Regenerated reply");
651        assert_eq!(active[0].reply_token, "token-regenerated");
652        assert_eq!(active[0].commit_hash.as_deref(), Some("commit-regenerated"));
653        assert!(!active[0].is_posting);
654    }
655
656    fn review_comment_resolution(
657        reply: &str,
658        reply_token: &str,
659    ) -> NewSessionReviewCommentResolution {
660        NewSessionReviewCommentResolution {
661            commit_hash: None,
662            reply: reply.to_string(),
663            reply_token: reply_token.to_string(),
664            resolution: "fixed".to_string(),
665            review_request_display_id: "#42".to_string(),
666            thread_id: "thread-1".to_string(),
667        }
668    }
669}