Skip to main content

tdm_server_rust/repository/
gray_release_repo.rs

1//! 灰度发布数据访问层 (Gray Release Repository)
2//!
3//! 封装灰度功能配置表和体验人员表的读写。
4
5use crate::entity::gray_release::{GrayFeatureRecord, GrayFeatureUpdateRequest};
6use chrono::{DateTime, Utc};
7use sqlx::{PgPool, Row};
8use std::collections::{BTreeMap, BTreeSet};
9
10/// 灰度发布仓储。
11pub struct GrayReleaseRepository {
12    /// PostgreSQL 连接池。
13    db: PgPool,
14}
15
16impl GrayReleaseRepository {
17    /// 从 PostgreSQL 连接池构造仓储。
18    pub fn new(db: PgPool) -> Self {
19        Self { db }
20    }
21
22    /// 查询全部灰度功能配置。
23    pub async fn list_features(&self) -> crate::error::ApiResult<Vec<GrayFeatureRecord>> {
24        let feature_rows = sqlx::query(
25            r#"
26            SELECT feature_key, feature_name, description, enabled, rollout_rate, created_at, updated_at
27            FROM gray_release_feature
28            ORDER BY feature_key ASC
29            "#,
30        )
31        .fetch_all(&self.db)
32        .await?;
33        let tester_rows = sqlx::query(
34            r#"
35            SELECT feature_key, member_id
36            FROM gray_release_tester
37            ORDER BY feature_key ASC, member_id ASC
38            "#,
39        )
40        .fetch_all(&self.db)
41        .await?;
42        let mut testers: BTreeMap<String, Vec<i32>> = BTreeMap::new();
43        for row in tester_rows {
44            testers
45                .entry(row.try_get::<String, _>("feature_key")?)
46                .or_default()
47                .push(row.try_get::<i32, _>("member_id")?);
48        }
49        let mut features = Vec::with_capacity(feature_rows.len());
50        for row in feature_rows {
51            let feature_key = row.try_get::<String, _>("feature_key")?;
52            features.push(GrayFeatureRecord {
53                tester_member_ids: testers.remove(&feature_key).unwrap_or_default(),
54                feature_key,
55                feature_name: row.try_get("feature_name")?,
56                description: row.try_get("description")?,
57                enabled: row.try_get("enabled")?,
58                rollout_rate: row.try_get("rollout_rate")?,
59                created_at: row.try_get::<DateTime<Utc>, _>("created_at")?,
60                updated_at: row.try_get::<DateTime<Utc>, _>("updated_at")?,
61            });
62        }
63        Ok(features)
64    }
65
66    /// 按功能标识查询单个灰度功能配置。
67    pub async fn get_feature(
68        &self,
69        feature_key: &str,
70    ) -> crate::error::ApiResult<Option<GrayFeatureRecord>> {
71        let rows = self.list_features().await?;
72        Ok(rows
73            .into_iter()
74            .find(|feature| feature.feature_key == feature_key))
75    }
76
77    /// 新增或更新灰度功能配置和体验人员。
78    pub async fn upsert_feature(
79        &self,
80        feature_key: &str,
81        req: GrayFeatureUpdateRequest,
82    ) -> crate::error::ApiResult<()> {
83        let tester_ids = normalize_member_ids(req.tester_member_ids);
84        let mut tx = self.db.begin().await?;
85        sqlx::query(
86            r#"
87            INSERT INTO gray_release_feature(feature_key, feature_name, description, enabled, rollout_rate)
88            VALUES ($1, $2, $3, $4, $5)
89            ON CONFLICT (feature_key) DO UPDATE
90            SET feature_name = EXCLUDED.feature_name,
91                description = EXCLUDED.description,
92                enabled = EXCLUDED.enabled,
93                rollout_rate = EXCLUDED.rollout_rate
94            "#,
95        )
96        .bind(feature_key)
97        .bind(req.feature_name.trim())
98        .bind(req.description.trim())
99        .bind(req.enabled)
100        .bind(req.rollout_rate)
101        .execute(&mut *tx)
102        .await?;
103        sqlx::query("DELETE FROM gray_release_tester WHERE feature_key = $1")
104            .bind(feature_key)
105            .execute(&mut *tx)
106            .await?;
107        for member_id in tester_ids {
108            sqlx::query(
109                r#"
110                INSERT INTO gray_release_tester(feature_key, member_id)
111                VALUES ($1, $2)
112                ON CONFLICT (feature_key, member_id) DO NOTHING
113                "#,
114            )
115            .bind(feature_key)
116            .bind(member_id)
117            .execute(&mut *tx)
118            .await?;
119        }
120        tx.commit().await?;
121        Ok(())
122    }
123
124    /// 回退单个灰度功能到旧逻辑。
125    pub async fn rollback_feature(&self, feature_key: &str) -> crate::error::ApiResult<()> {
126        sqlx::query(
127            r#"
128            UPDATE gray_release_feature
129            SET enabled = FALSE,
130                rollout_rate = 0
131            WHERE feature_key = $1
132            "#,
133        )
134        .bind(feature_key)
135        .execute(&self.db)
136        .await?;
137        Ok(())
138    }
139
140    /// 回退全部灰度功能到旧逻辑。
141    pub async fn rollback_all(&self) -> crate::error::ApiResult<()> {
142        sqlx::query(
143            r#"
144            UPDATE gray_release_feature
145            SET enabled = FALSE,
146                rollout_rate = 0
147            "#,
148        )
149        .execute(&self.db)
150        .await?;
151        Ok(())
152    }
153}
154
155/// 去重并过滤无效组员 ID。
156fn normalize_member_ids(ids: Vec<i32>) -> Vec<i32> {
157    ids.into_iter()
158        .filter(|id| *id > 0)
159        .collect::<BTreeSet<_>>()
160        .into_iter()
161        .collect()
162}