use uuid::Uuid; use crate::error::AppError; use crate::models::issues::IssueSubscriber; use crate::service::IssueService; use crate::session::Session; use super::util::{clamp_limit_offset, ensure_affected}; impl IssueService { pub async fn issue_subscribers( &self, ctx: &Session, wk_name: &str, number: i64, limit: i64, offset: i64, ) -> Result, AppError> { let user_uid = ctx.user().ok_or(AppError::Unauthorized)?; let issue = self.resolve_issue(wk_name, number).await?; let issue_id = issue.id; self.ensure_issue_readable(user_uid, &issue).await?; let (limit, offset) = clamp_limit_offset(limit, offset); sqlx::query_as::<_, IssueSubscriber>( "SELECT id, issue_id, user_id, reason, muted, created_at, updated_at \ FROM issue_subscriber WHERE issue_id = $1 ORDER BY created_at ASC LIMIT $2 OFFSET $3", ) .bind(issue_id) .bind(limit) .bind(offset) .fetch_all(self.ctx.db.reader()) .await .map_err(AppError::Database) } pub async fn issue_subscribe( &self, ctx: &Session, wk_name: &str, number: i64, ) -> Result { let user_uid = ctx.user().ok_or(AppError::Unauthorized)?; let issue = self.resolve_issue(wk_name, number).await?; let issue_id = issue.id; self.ensure_issue_readable(user_uid, &issue).await?; let now = chrono::Utc::now(); let mut txn = self .ctx .db .writer() .begin() .await .map_err(|_| AppError::TxnError)?; sqlx::query("SET LOCAL app.current_user_id = $1") .bind(user_uid) .execute(&mut *txn) .await .map_err(AppError::Database)?; let sub = sqlx::query_as::<_, IssueSubscriber>( "INSERT INTO issue_subscriber (id, issue_id, user_id, reason, muted, created_at, updated_at) \ VALUES ($1, $2, $3, 'manual', false, $4, $4) ON CONFLICT (issue_id, user_id) DO NOTHING \ RETURNING id, issue_id, user_id, reason, muted, created_at, updated_at", ) .bind(Uuid::now_v7()).bind(issue_id).bind(user_uid).bind(now) .fetch_optional(&mut *txn).await.map_err(AppError::Database)? .ok_or(AppError::Conflict("already subscribed".into()))?; sqlx::query("UPDATE issue_stats SET subscribers_count = subscribers_count + 1, updated_at = $1 WHERE issue_id = $2") .bind(now).bind(issue_id).execute(&mut *txn).await.map_err(AppError::Database)?; txn.commit().await.map_err(|_| AppError::TxnError)?; Ok(sub) } pub async fn issue_unsubscribe( &self, ctx: &Session, wk_name: &str, number: i64, ) -> Result<(), AppError> { let user_uid = ctx.user().ok_or(AppError::Unauthorized)?; let issue = self.resolve_issue(wk_name, number).await?; let issue_id = issue.id; let now = chrono::Utc::now(); let mut txn = self .ctx .db .writer() .begin() .await .map_err(|_| AppError::TxnError)?; sqlx::query("SET LOCAL app.current_user_id = $1") .bind(user_uid) .execute(&mut *txn) .await .map_err(AppError::Database)?; let result = sqlx::query("DELETE FROM issue_subscriber WHERE issue_id = $1 AND user_id = $2") .bind(issue_id) .bind(user_uid) .execute(&mut *txn) .await .map_err(AppError::Database)?; ensure_affected(result.rows_affected(), "not subscribed")?; sqlx::query("UPDATE issue_stats SET subscribers_count = GREATEST(subscribers_count - 1, 0), updated_at = $1 WHERE issue_id = $2") .bind(now).bind(issue_id).execute(&mut *txn).await.map_err(AppError::Database)?; txn.commit().await.map_err(|_| AppError::TxnError)?; Ok(()) } pub async fn issue_mute( &self, ctx: &Session, wk_name: &str, number: i64, muted: bool, ) -> Result<(), AppError> { let user_uid = ctx.user().ok_or(AppError::Unauthorized)?; let issue = self.resolve_issue(wk_name, number).await?; let issue_id = issue.id; let result = sqlx::query( "UPDATE issue_subscriber SET muted = $1, updated_at = $2 WHERE issue_id = $3 AND user_id = $4", ) .bind(muted).bind(chrono::Utc::now()).bind(issue_id).bind(user_uid) .execute(self.ctx.db.writer()).await.map_err(AppError::Database)?; ensure_affected(result.rows_affected(), "not subscribed") } }