1use ag_session::ReviewRequest;
4use async_trait::async_trait;
5use sqlx::{SqliteConnection, SqlitePool};
6
7use crate::DbError;
8
9#[derive(Clone, Debug, Eq, PartialEq)]
11pub struct NewSessionReviewCommentResolution {
12 pub commit_hash: Option<String>,
14 pub reply: String,
16 pub reply_token: String,
18 pub resolution: String,
20 pub review_request_display_id: String,
22 pub thread_id: String,
24}
25
26#[derive(Clone, Debug, Eq, PartialEq)]
28pub struct SessionReviewCommentResolutionRow {
29 pub commit_hash: Option<String>,
31 pub is_posting: bool,
33 pub reply: String,
35 pub reply_token: String,
37 pub resolution: String,
39 pub review_request_display_id: String,
41 pub thread_id: String,
43}
44
45#[derive(Clone, Debug, Eq, PartialEq)]
47pub struct SessionReviewRequestRow {
48 pub display_id: String,
50 pub forge_kind: String,
52 pub last_refreshed_at: i64,
54 pub source_branch: String,
56 pub state: String,
58 pub status_summary: Option<String>,
60 pub target_branch: String,
62 pub title: String,
64 pub web_url: String,
66}
67
68#[cfg_attr(test, mockall::automock)]
70#[async_trait]
71pub trait ReviewRepository: Send + Sync {
72 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 async fn discard_session_review_comment_resolutions(
82 &self,
83 id: &str,
84 resolutions: &[NewSessionReviewCommentResolution],
85 ) -> Result<(), DbError>;
86
87 async fn insert_session_review_comment_resolutions(
92 &self,
93 id: &str,
94 resolutions: &[NewSessionReviewCommentResolution],
95 ) -> Result<(), DbError>;
96
97 async fn load_session_review_comment_resolutions(
99 &self,
100 id: &str,
101 ) -> Result<Vec<SessionReviewCommentResolutionRow>, DbError>;
102
103 async fn load_session_review_request(
105 &self,
106 id: &str,
107 ) -> Result<Option<SessionReviewRequestRow>, DbError>;
108
109 async fn update_session_review_request(
111 &self,
112 id: &str,
113 review_request: Option<ReviewRequest>,
114 ) -> Result<(), DbError>;
115
116 async fn mark_session_review_comment_resolution_posting(
118 &self,
119 id: &str,
120 reply_token: &str,
121 ) -> Result<(), DbError>;
122
123 async fn remove_session_review_comment_resolution(
125 &self,
126 id: &str,
127 reply_token: &str,
128 ) -> Result<(), DbError>;
129}
130
131#[derive(Clone)]
133pub(crate) struct SqliteReviewRepository(SqlitePool);
134
135impl SqliteReviewRepository {
136 pub(crate) fn new(pool: SqlitePool) -> Self {
138 Self(pool)
139 }
140}
141
142pub(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 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 repositories
462 .reviews()
463 .insert_session_review_comment_resolutions(
464 "session-id",
465 std::slice::from_ref(®enerated),
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(®enerated),
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_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 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 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_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 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 repositories
626 .reviews()
627 .insert_session_review_comment_resolutions(
628 "session-id",
629 std::slice::from_ref(®enerated),
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(®enerated),
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_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}