Skip to main content

tdm_server_rust/repository/
oss_repo.rs

1//! OSS 对象存储数据访问层 (OSS Repository)
2//!
3//! 封装 `oss` 表记录管理及 COS STS 凭证生成(SeaORM)。
4
5use crate::config::AppConfig;
6use crate::db::DbConn;
7use crate::entity::oss::{OssCredential, OssDto};
8use crate::sea_entity::{episode_role_taker, mangaepisodetb, oss};
9use crate::utils::{agent_debug, cos_cdn, cos_presign, cos_sts};
10use sea_orm::{ActiveModelTrait, ColumnTrait, DatabaseConnection, EntityTrait, QueryFilter, Set};
11use sqlx::postgres::PgPool;
12use std::time::Duration;
13use uuid::Uuid;
14
15/// OSS 仓储
16pub struct OssRepository {
17    /// SeaORM 数据库连接
18    db: DbConn,
19    /// 应用配置
20    config: AppConfig,
21}
22
23impl OssRepository {
24    /// 从 `PgPool` 构造
25    pub fn new(pool: PgPool, config: AppConfig) -> Self {
26        Self {
27            db: crate::db::from_sqlx_pool(pool),
28            config,
29        }
30    }
31
32    /// 从 `DatabaseConnection` 构造
33    pub fn from_db(db: DatabaseConnection, config: AppConfig) -> Self {
34        Self { db, config }
35    }
36
37    /// 新增或更新 OSS 记录,绑定到话数岗位列
38    #[tracing::instrument(skip_all, level = "debug")]
39    pub async fn upsert_oss(&self, dto: &OssDto, member_id: i32) -> crate::error::ApiResult<i32> {
40        let bucket = &self.config.tencent.bucket;
41        let region = &self.config.tencent.region;
42        let manga_id = self.get_episode_manga_id(dto.episode_id).await?;
43        let object_key = dto.object_key.clone().unwrap_or_else(|| {
44            self.build_object_key(manga_id, dto.episode_id, &dto.post_name, &dto.filename)
45        });
46        self.verify_object_exists_if_configured(bucket, region, &object_key)
47            .await?;
48        let url = format!("https://{bucket}.cos.{region}.myqcloud.com/{object_key}");
49        let now = chrono::Utc::now().naive_utc();
50        let ext = file_ext(&dto.filename).unwrap_or_default();
51
52        let existing_oss_id = self
53            .get_episode_oss_id(dto.episode_id, &dto.post_name, Some(member_id))
54            .await?;
55        let oss_id = if let Some(id) = existing_oss_id {
56            let model = oss::ActiveModel {
57                id: Set(id),
58                filename: Set(dto.filename.clone()),
59                object_key: Set(object_key.clone()),
60                url: Set(url.clone()),
61                bucket_name: Set(Some(bucket.clone())),
62                region: Set(Some(region.clone())),
63                update_time: Set(Some(now)),
64                update_by: Set(Some(member_id)),
65                ..Default::default()
66            };
67            model.update(&self.db).await?;
68            id
69        } else {
70            let model = oss::ActiveModel {
71                filename: Set(dto.filename.clone()),
72                file_ext: Set(ext),
73                object_key: Set(object_key.clone()),
74                url: Set(url.clone()),
75                bucket_name: Set(Some(bucket.clone())),
76                region: Set(Some(region.clone())),
77                create_time: Set(Some(now)),
78                create_by: Set(Some(member_id)),
79                update_time: Set(Some(now)),
80                update_by: Set(Some(member_id)),
81                ..Default::default()
82            };
83            let inserted = model.insert(&self.db).await?;
84            self.bind_episode_oss_id(dto.episode_id, &dto.post_name, member_id, inserted.id)
85                .await?;
86            inserted.id
87        };
88        mangaepisodetb::Entity::update_many()
89            .col_expr(
90                mangaepisodetb::Column::UpdateTime,
91                sea_orm::sea_query::Expr::current_timestamp().into(),
92            )
93            .filter(mangaepisodetb::Column::Id.eq(dto.episode_id))
94            .exec(&self.db)
95            .await?;
96        Ok(oss_id)
97    }
98
99    /// 删除 OSS 记录
100    #[tracing::instrument(skip_all, level = "debug")]
101    pub async fn delete_oss(&self, oss_id: i32) -> crate::error::ApiResult<Option<String>> {
102        let row = oss::Entity::find_by_id(oss_id).one(&self.db).await?;
103        let key = row.map(|r| r.object_key);
104        oss::Entity::delete_by_id(oss_id).exec(&self.db).await?;
105        Ok(key)
106    }
107
108    /// 获取上传凭证
109    #[tracing::instrument(skip_all, level = "debug")]
110    pub async fn get_upload_credential(
111        &self,
112        episode_id: i32,
113        post_name: &str,
114        filename: &str,
115    ) -> crate::error::ApiResult<OssCredential> {
116        self.ensure_episode_exists(episode_id).await?;
117        self.ensure_ext_allowed(filename, &self.config.tencent.ext_whitelist)?;
118        let manga_id = self.get_episode_manga_id(episode_id).await?;
119        let object_key = self.build_object_key(manga_id, episode_id, post_name, filename);
120        let max_bytes = self.config.tencent.max_file_size * 1024 * 1024;
121        agent_debug::log(
122            "H5",
123            "oss_repo.rs:get_upload_credential",
124            "upload_credential_request",
125            serde_json::json!({
126                "episodeId": episode_id,
127                "postName": post_name,
128                "objectKeyLen": object_key.len(),
129                "hasNonAscii": !object_key.is_ascii(),
130                "maxBytes": max_bytes
131            }),
132        );
133        let presigned_url = self
134            .presigned_put_url(&self.config.tencent.bucket, &object_key, max_bytes)
135            .await?;
136        Ok(OssCredential {
137            presigned_url,
138            object_key: Some(object_key),
139            original_filename: Some(filename.to_string()),
140        })
141    }
142
143    /// 获取下载凭证
144    #[tracing::instrument(skip_all, level = "debug")]
145    pub async fn get_download_credential(
146        &self,
147        episode_id: i32,
148        post_name: &str,
149    ) -> crate::error::ApiResult<OssCredential> {
150        let oss_id: i32 = self
151            .get_episode_oss_id(episode_id, post_name, None)
152            .await?
153            .ok_or_else(|| {
154                agent_debug::log(
155                    "H2",
156                    "oss_repo.rs:get_download_credential",
157                    "oss_id_missing",
158                    serde_json::json!({
159                        "episodeId": episode_id,
160                        "postName": post_name
161                    }),
162                );
163                crate::error::AppError::Oss {
164                    code: None,
165                    msg: format!("当前单话尚未上传 [{post_name}] 文件"),
166                }
167            })?;
168        let row = oss::Entity::find_by_id(oss_id)
169            .one(&self.db)
170            .await?
171            .ok_or_else(|| crate::error::AppError::Oss {
172                code: None,
173                msg: format!("OSS 记录 {oss_id} 不存在"),
174            })?;
175        let object_key = row.object_key.clone();
176        let original_filename = Some(row.filename.clone());
177        let bucket = resolve_download_bucket(&row, &self.config.tencent.bucket);
178        let region = row
179            .region
180            .as_deref()
181            .filter(|s| !s.trim().is_empty())
182            .unwrap_or(&self.config.tencent.region);
183        let cdn_domain = self.resolve_download_cdn_domain(&row);
184        let presigned_url = if let Some(cdn_url) = cdn_domain.as_deref().and_then(|domain| {
185            cos_cdn::generate_cdn_url(domain, &self.config.tencent.cdn_key, &object_key)
186        }) {
187            agent_debug::log(
188                "H1",
189                "oss_repo.rs:get_download_credential",
190                "cdn_signed_url",
191                serde_json::json!({
192                    "hasSign": cdn_url.contains("sign="),
193                    "hasT": cdn_url.contains("&t="),
194                    "objectKeyLen": object_key.len(),
195                    "bucket": bucket,
196                    "cdnDomain": cdn_domain
197                }),
198            );
199            cdn_url
200        } else {
201            agent_debug::log(
202                "H1",
203                "oss_repo.rs:get_download_credential",
204                "fallback_cos_presign",
205                serde_json::json!({
206                    "cdnDomainEmpty": self.config.tencent.cdn_domain.is_empty(),
207                    "cdnKeyEmpty": self.config.tencent.cdn_key.is_empty(),
208                    "bucket": bucket,
209                    "region": region
210                }),
211            );
212            self.presigned_get_url(&bucket, region, &object_key).await?
213        };
214        Ok(OssCredential {
215            presigned_url,
216            object_key: Some(object_key),
217            original_filename,
218        })
219    }
220
221    /// 获取图片上传凭证
222    #[tracing::instrument(skip_all, level = "debug")]
223    pub async fn get_image_upload_credential(
224        &self,
225        image_type: &str,
226        filename: &str,
227    ) -> crate::error::ApiResult<OssCredential> {
228        let image_bucket = self.config.tencent.image_bucket.trim();
229        if image_bucket.is_empty() {
230            return Err(crate::error::AppError::Oss {
231                code: None,
232                msg: "图片存储桶未配置".into(),
233            });
234        }
235        self.ensure_ext_allowed(filename, &self.config.tencent.image_ext_whitelist)?;
236        let object_key = self.build_image_object_key(image_type, filename);
237        let max_bytes = self.config.tencent.image_max_file_size * 1024 * 1024;
238        let presigned_url = self
239            .presigned_put_url(image_bucket, &object_key, max_bytes)
240            .await?;
241        Ok(OssCredential {
242            presigned_url,
243            object_key: Some(object_key),
244            original_filename: Some(filename.to_string()),
245        })
246    }
247
248    /// 生成图片桶的预签名 PUT URL(服务端导入图源时上传用)
249    ///
250    /// 与 `get_image_upload_credential` 不同:此处不校验图片后缀白名单,
251    /// 且使用更大的大小上限,用于服务端批量上传解压后的页面图片。
252    #[tracing::instrument(skip_all, level = "debug")]
253    pub async fn presigned_image_put_url(
254        &self,
255        object_key: &str,
256        max_bytes: u64,
257    ) -> crate::error::ApiResult<String> {
258        let image_bucket = self.config.tencent.image_bucket.trim();
259        if image_bucket.is_empty() {
260            return Err(crate::error::AppError::Oss {
261                code: None,
262                msg: "图片存储桶未配置".into(),
263            });
264        }
265        self.presigned_put_url(image_bucket, object_key, max_bytes)
266            .await
267    }
268
269    /// 批量生成图片桶预签名 PUT URL(前端图源解压后逐张直传用)
270    ///
271    /// 以 `{prefix}*` 前缀策略**仅申请一次** STS 联合身份临时凭证,
272    /// 再用同一份临时密钥对每个对象键分别签名,避免逐张调用 STS 拖慢上传。
273    ///
274    /// ## 参数
275    /// - `prefix`: 对象键公共前缀(如 `editor/episode_123/`),策略按 `{prefix}*` 授权
276    /// - `object_keys`: 全部待签名对象键(必须都落在 `prefix` 下)
277    /// - `max_bytes`: 单个对象允许的最大字节数
278    #[tracing::instrument(skip_all, level = "debug")]
279    pub async fn presigned_image_put_urls(
280        &self,
281        prefix: &str,
282        object_keys: &[String],
283        max_bytes: u64,
284    ) -> crate::error::ApiResult<Vec<String>> {
285        let image_bucket = self.config.tencent.image_bucket.trim();
286        if image_bucket.is_empty() {
287            return Err(crate::error::AppError::Oss {
288                code: None,
289                msg: "图片存储桶未配置".into(),
290            });
291        }
292        if object_keys.is_empty() {
293            return Ok(Vec::new());
294        }
295        let region = &self.config.tencent.region;
296        // 通配前缀策略,单次 STS 覆盖本话数全部页面对象键
297        let policy = cos_sts::build_put_policy(
298            image_bucket,
299            region,
300            &format!("{}*", prefix.trim_start_matches('/')),
301            max_bytes,
302        )?;
303        let creds = cos_sts::get_federation_token(
304            &self.config.tencent.secret_id,
305            &self.config.tencent.secret_key,
306            region,
307            policy,
308            self.config.tencent.duration_seconds,
309        )
310        .await?;
311        let mut urls = Vec::with_capacity(object_keys.len());
312        for key in object_keys {
313            let url = cos_presign::presigned_url(
314                &creds.tmp_secret_id,
315                &creds.tmp_secret_key,
316                image_bucket,
317                region,
318                key,
319                "PUT",
320                self.config.tencent.duration_seconds,
321                &creds.session_token,
322            )?;
323            urls.push(url);
324        }
325        Ok(urls)
326    }
327
328    /// 构造图片访问 URL(优先图片 CDN 域名,否则 COS 直链)
329    pub fn image_display_url(&self, object_key: &str) -> String {
330        let cdn = self.config.tencent.image_cdn_domain.trim();
331        if !cdn.is_empty() {
332            let domain = cdn.trim_end_matches('/');
333            let proto = if domain.starts_with("http://") || domain.starts_with("https://") {
334                ""
335            } else {
336                "https://"
337            };
338            return format!("{proto}{domain}/{object_key}");
339        }
340        format!(
341            "https://{}.cos.{}.myqcloud.com/{}",
342            self.config.tencent.image_bucket, self.config.tencent.region, object_key
343        )
344    }
345
346    /// 图片桶名
347    pub fn image_bucket_name(&self) -> String {
348        self.config.tencent.image_bucket.clone()
349    }
350
351    /// STS + 预签名 PUT URL
352    #[tracing::instrument(skip_all, level = "debug")]
353    async fn presigned_put_url(
354        &self,
355        bucket: &str,
356        object_key: &str,
357        max_bytes: u64,
358    ) -> crate::error::ApiResult<String> {
359        let policy =
360            cos_sts::build_put_policy(bucket, &self.config.tencent.region, object_key, max_bytes)?;
361        let creds = cos_sts::get_federation_token(
362            &self.config.tencent.secret_id,
363            &self.config.tencent.secret_key,
364            &self.config.tencent.region,
365            policy,
366            self.config.tencent.duration_seconds,
367        )
368        .await?;
369        cos_presign::presigned_url(
370            &creds.tmp_secret_id,
371            &creds.tmp_secret_key,
372            bucket,
373            &self.config.tencent.region,
374            object_key,
375            "PUT",
376            self.config.tencent.duration_seconds,
377            &creds.session_token,
378        )
379    }
380
381    /// STS + 预签名 GET URL
382    #[tracing::instrument(skip_all, level = "debug")]
383    async fn presigned_get_url(
384        &self,
385        bucket: &str,
386        region: &str,
387        object_key: &str,
388    ) -> crate::error::ApiResult<String> {
389        if self.config.tencent.secret_id.trim().is_empty()
390            || self.config.tencent.secret_key.trim().is_empty()
391        {
392            return Err(crate::error::AppError::Oss {
393                code: None,
394                msg: "当前环境缺少 COS 下载密钥,无法生成对象存储直链".into(),
395            });
396        }
397
398        let policy = cos_sts::build_get_policy(bucket, region, object_key)?;
399        let creds = cos_sts::get_federation_token(
400            &self.config.tencent.secret_id,
401            &self.config.tencent.secret_key,
402            region,
403            policy,
404            self.config.tencent.duration_seconds,
405        )
406        .await?;
407        cos_presign::presigned_url(
408            &creds.tmp_secret_id,
409            &creds.tmp_secret_key,
410            bucket,
411            region,
412            object_key,
413            "GET",
414            self.config.tencent.duration_seconds,
415            &creds.session_token,
416        )
417    }
418
419    /// 按 OSS 行内信息选择下载 CDN 域名。
420    ///
421    /// 只在 OSS 记录属于当前 profile 桶时使用当前 profile 的 CDN key。
422    /// 如果记录来自其他 bucket,不能拿当前环境 key 去签别的 CDN 域名,
423    /// 应降级为对应 bucket/region 的 COS 预签名直链。
424    fn resolve_download_cdn_domain(&self, row: &oss::Model) -> Option<String> {
425        if self.config.tencent.cdn_key.trim().is_empty() {
426            return None;
427        }
428
429        if should_use_config_cdn(
430            row.bucket_name.as_deref(),
431            &row.url,
432            &self.config.tencent.bucket,
433            &self.config.tencent.cdn_domain,
434        ) {
435            return non_empty_trimmed(&self.config.tencent.cdn_domain);
436        }
437
438        None
439    }
440
441    /// 在密钥可用时校验对象已经真实存在于 COS。
442    ///
443    /// 对齐 Java 旧服务在 upsert 前调用 `getObjectMetadata` 的保护思路。
444    /// 本地开发/测试经常没有腾讯云 SecretKey,此时跳过网络校验,避免阻断
445    /// 不依赖 OSS 的开发流程。
446    async fn verify_object_exists_if_configured(
447        &self,
448        bucket: &str,
449        region: &str,
450        object_key: &str,
451    ) -> crate::error::ApiResult<()> {
452        if self.config.tencent.secret_id.trim().is_empty()
453            || self.config.tencent.secret_key.trim().is_empty()
454        {
455            agent_debug::log(
456                "H2",
457                "oss_repo.rs:verify_object_exists",
458                "skip_missing_secret",
459                serde_json::json!({
460                    "bucket": bucket,
461                    "objectKeyLen": object_key.len()
462                }),
463            );
464            return Ok(());
465        }
466
467        let url = self.presigned_get_url(bucket, region, object_key).await?;
468        let client = reqwest::Client::builder()
469            .timeout(Duration::from_secs(10))
470            .build()
471            .map_err(|e| {
472                crate::error::AppError::Internal(format!("构建 COS 校验客户端失败: {e}"))
473            })?;
474        let resp = client
475            .get(&url)
476            .header(reqwest::header::RANGE, "bytes=0-0")
477            .send()
478            .await
479            .map_err(|e| crate::error::AppError::Oss {
480                code: None,
481                msg: format!("校验 COS 对象存在失败: {e}"),
482            })?;
483        if resp.status().is_success() || resp.status() == reqwest::StatusCode::PARTIAL_CONTENT {
484            return Ok(());
485        }
486        Err(crate::error::AppError::Oss {
487            code: None,
488            msg: format!(
489                "COS 对象不存在或无法访问: {object_key} (HTTP {})",
490                resp.status()
491            ),
492        })
493    }
494
495    /// 按 ID 查询 objectKey
496    #[tracing::instrument(skip_all, level = "debug")]
497    pub async fn get_object_key(&self, oss_id: i32) -> crate::error::ApiResult<Option<String>> {
498        let row = oss::Entity::find_by_id(oss_id).one(&self.db).await?;
499        Ok(row.map(|r| r.object_key))
500    }
501
502    /// 查询话数某岗位是否已绑定 OSS 记录(对外只读)
503    ///
504    /// 用于「下载图源」判断:provider 未绑定压缩包时改走编辑器页面现打包。
505    #[tracing::instrument(skip_all, level = "debug")]
506    pub async fn episode_post_oss_id(
507        &self,
508        episode_id: i32,
509        post_name: &str,
510    ) -> crate::error::ApiResult<Option<i32>> {
511        self.get_episode_oss_id(episode_id, post_name, None).await
512    }
513
514    /// 检测话数是否存在
515    #[tracing::instrument(skip_all, level = "debug")]
516    async fn ensure_episode_exists(&self, episode_id: i32) -> crate::error::ApiResult<()> {
517        let exists = mangaepisodetb::Entity::find_by_id(episode_id)
518            .one(&self.db)
519            .await?
520            .is_some();
521        if !exists {
522            return Err(crate::error::AppError::business(format!(
523                "漫画单话 {episode_id} 不存在"
524            )));
525        }
526        Ok(())
527    }
528
529    /// 查询话数岗位绑定的 OSS ID
530    #[tracing::instrument(skip_all, level = "debug")]
531    async fn get_episode_oss_id(
532        &self,
533        episode_id: i32,
534        post_name: &str,
535        member_id: Option<i32>,
536    ) -> crate::error::ApiResult<Option<i32>> {
537        validate_post_name(post_name)?;
538        if let Some(member_id) = member_id {
539            if let Some(taker) =
540                find_episode_role_taker(&self.db, episode_id, post_name, member_id).await?
541            {
542                return Ok(taker.file_oss_id);
543            }
544        }
545        let row = mangaepisodetb::Entity::find_by_id(episode_id)
546            .one(&self.db)
547            .await?;
548        Ok(row.and_then(|ep| read_post_oss_id(&ep, post_name)))
549    }
550
551    /// 绑定话数岗位 OSS ID
552    #[tracing::instrument(skip_all, level = "debug")]
553    async fn bind_episode_oss_id(
554        &self,
555        episode_id: i32,
556        post_name: &str,
557        member_id: i32,
558        oss_id: i32,
559    ) -> crate::error::ApiResult<()> {
560        validate_post_name(post_name)?;
561        if let Some(taker) =
562            find_episode_role_taker(&self.db, episode_id, post_name, member_id).await?
563        {
564            let mut am: episode_role_taker::ActiveModel = taker.into();
565            am.file_oss_id = Set(Some(oss_id));
566            am.update(&self.db).await?;
567            return Ok(());
568        }
569        let Some(ep) = mangaepisodetb::Entity::find_by_id(episode_id)
570            .one(&self.db)
571            .await?
572        else {
573            return Err(crate::error::AppError::business("话数不存在喵"));
574        };
575        let mut am: mangaepisodetb::ActiveModel = ep.into();
576        write_post_oss_id(&mut am, post_name, oss_id)?;
577        am.update_time = Set(chrono::Utc::now().naive_utc());
578        am.update(&self.db).await?;
579        Ok(())
580    }
581
582    /// 查询话数所属漫画 ID
583    #[tracing::instrument(skip_all, level = "debug")]
584    async fn get_episode_manga_id(&self, episode_id: i32) -> crate::error::ApiResult<i32> {
585        let row = mangaepisodetb::Entity::find_by_id(episode_id)
586            .one(&self.db)
587            .await?;
588        row.map(|r| r.manga_id).ok_or_else(|| {
589            crate::error::AppError::business(format!("漫画单话 {episode_id} 不存在"))
590        })
591    }
592
593    /// 生成话数稿件 objectKey
594    fn build_object_key(
595        &self,
596        manga_id: i32,
597        episode_id: i32,
598        post_name: &str,
599        filename: &str,
600    ) -> String {
601        let file_type_dir = match post_name.to_lowercase().as_str() {
602            "provider" => "manga-raw",
603            "translator" => "translation",
604            "proofreader" => "proofread",
605            "letterer" => "manga-cooked",
606            "timer" => "timing",
607            _ => "other",
608        };
609        format!("manga_{manga_id}/episode_{episode_id}/{file_type_dir}/{filename}")
610    }
611
612    /// 生成图片 objectKey
613    fn build_image_object_key(&self, image_type: &str, filename: &str) -> String {
614        let mut normalized = image_type.trim().trim_matches('/').to_string();
615        if normalized.is_empty() {
616            normalized = "misc".into();
617        }
618        let ext = file_ext(filename).unwrap_or_else(|| "jpg".into());
619        format!("{normalized}/{}.{}", Uuid::new_v4(), ext)
620    }
621
622    /// 校验后缀白名单
623    fn ensure_ext_allowed(
624        &self,
625        filename: &str,
626        whitelist: &[String],
627    ) -> crate::error::ApiResult<()> {
628        let ext = file_ext(filename).ok_or_else(|| crate::error::AppError::Oss {
629            code: None,
630            msg: format!("无法识别文件后缀: {filename}"),
631        })?;
632        if !whitelist.iter().any(|w| w.eq_ignore_ascii_case(&ext)) {
633            return Err(crate::error::AppError::Oss {
634                code: None,
635                msg: format!("非法文件后缀: {ext}"),
636            });
637        }
638        Ok(())
639    }
640}
641
642/// 校验岗位名
643fn validate_post_name(post_name: &str) -> crate::error::ApiResult<()> {
644    match post_name.to_lowercase().as_str() {
645        "provider" | "translator" | "proofreader" | "letterer" | "timer" => Ok(()),
646        _ => Err(crate::error::AppError::Oss {
647            code: None,
648            msg: format!("不支持的岗位类型:{post_name}"),
649        }),
650    }
651}
652
653fn non_empty_trimmed(value: &str) -> Option<String> {
654    let trimmed = value.trim();
655    (!trimmed.is_empty()).then(|| trimmed.to_string())
656}
657
658fn resolve_download_bucket(row: &oss::Model, config_bucket: &str) -> String {
659    row.bucket_name
660        .as_deref()
661        .and_then(non_empty_trimmed)
662        .or_else(|| cos_bucket_from_url(&row.url))
663        .unwrap_or_else(|| config_bucket.to_string())
664}
665
666fn should_use_config_cdn(
667    row_bucket: Option<&str>,
668    row_url: &str,
669    config_bucket: &str,
670    config_cdn_domain: &str,
671) -> bool {
672    let row_bucket_matches = row_bucket.is_none_or(|bucket| {
673        let bucket = bucket.trim();
674        bucket.is_empty() || bucket == config_bucket
675    });
676    row_bucket_matches
677        && stored_url_matches_config_resource(row_url, config_bucket, config_cdn_domain)
678}
679
680fn stored_url_matches_config_resource(
681    row_url: &str,
682    config_bucket: &str,
683    config_cdn_domain: &str,
684) -> bool {
685    let Some(host) = url_host(row_url) else {
686        return true;
687    };
688
689    if let Some(bucket) = cos_bucket_from_host(&host) {
690        return bucket == config_bucket;
691    }
692
693    let config_cdn_host = normalize_domain_host(config_cdn_domain);
694    config_cdn_host.as_deref() == Some(host.as_str())
695}
696
697fn cos_bucket_from_url(url: &str) -> Option<String> {
698    url_host(url).and_then(|host| cos_bucket_from_host(&host))
699}
700
701fn cos_bucket_from_host(host: &str) -> Option<String> {
702    host.split_once(".cos.").and_then(|(bucket, _)| {
703        let bucket = bucket.trim();
704        (!bucket.is_empty()).then(|| bucket.to_string())
705    })
706}
707
708fn normalize_domain_host(domain: &str) -> Option<String> {
709    let trimmed = domain.trim().trim_end_matches('/');
710    if trimmed.is_empty() {
711        return None;
712    }
713    let without_scheme = trimmed
714        .strip_prefix("https://")
715        .or_else(|| trimmed.strip_prefix("http://"))
716        .unwrap_or(trimmed);
717    without_scheme
718        .split('/')
719        .next()
720        .map(str::trim)
721        .filter(|host| !host.is_empty())
722        .map(ToString::to_string)
723}
724
725fn url_host(url: &str) -> Option<String> {
726    let trimmed = url.trim();
727    if trimmed.is_empty() {
728        return None;
729    }
730    let without_scheme = trimmed
731        .strip_prefix("https://")
732        .or_else(|| trimmed.strip_prefix("http://"))?;
733    without_scheme
734        .split('/')
735        .next()
736        .map(str::trim)
737        .filter(|host| !host.is_empty())
738        .map(ToString::to_string)
739}
740
741/// 读取岗位 OSS ID 列
742fn read_post_oss_id(ep: &mangaepisodetb::Model, post_name: &str) -> Option<i32> {
743    match post_name.to_lowercase().as_str() {
744        "provider" => ep.provider_file_oss_id,
745        "translator" => ep.translator_file_oss_id,
746        "proofreader" => ep.proofreader_file_oss_id,
747        "letterer" => ep.letterer_file_oss_id,
748        "timer" => ep.timer_file_oss_id,
749        _ => None,
750    }
751}
752
753/// 查询当前上传人是否为翻译 / 校对共同接稿人。
754async fn find_episode_role_taker(
755    db: &DatabaseConnection,
756    episode_id: i32,
757    post_name: &str,
758    member_id: i32,
759) -> crate::error::ApiResult<Option<episode_role_taker::Model>> {
760    let role = post_name.to_lowercase();
761    if !matches!(role.as_str(), "translator" | "proofreader") {
762        return Ok(None);
763    }
764    let row = episode_role_taker::Entity::find()
765        .filter(episode_role_taker::Column::EpisodeId.eq(episode_id))
766        .filter(episode_role_taker::Column::Role.eq(role))
767        .filter(episode_role_taker::Column::MemberId.eq(member_id))
768        .one(db)
769        .await?;
770    Ok(row)
771}
772
773/// 写入岗位 OSS ID 列
774fn write_post_oss_id(
775    am: &mut mangaepisodetb::ActiveModel,
776    post_name: &str,
777    oss_id: i32,
778) -> crate::error::ApiResult<()> {
779    match post_name.to_lowercase().as_str() {
780        "provider" => am.provider_file_oss_id = Set(Some(oss_id)),
781        "translator" => am.translator_file_oss_id = Set(Some(oss_id)),
782        "proofreader" => am.proofreader_file_oss_id = Set(Some(oss_id)),
783        "letterer" => am.letterer_file_oss_id = Set(Some(oss_id)),
784        "timer" => am.timer_file_oss_id = Set(Some(oss_id)),
785        _ => {
786            return Err(crate::error::AppError::Oss {
787                code: None,
788                msg: format!("不支持的岗位类型:{post_name}"),
789            });
790        }
791    }
792    Ok(())
793}
794
795/// 取小写扩展名
796fn file_ext(filename: &str) -> Option<String> {
797    filename.rsplit('.').next().map(|s| s.to_lowercase())
798}
799
800#[cfg(test)]
801mod tests {
802    use super::*;
803
804    /// 岗位列读取与写入映射
805    #[test]
806    fn post_oss_column_mapping() {
807        let ep = mangaepisodetb::Model {
808            id: 1,
809            manga_id: 1,
810            manga_episode: "1".into(),
811            episode_type: "MAIN".into(),
812            manga_episode_name: None,
813            provider_id: None,
814            translator_id: None,
815            proofreader_id: None,
816            letterer_id: None,
817            timer_id: None,
818            reviewer_id: None,
819            setup_time: chrono::Utc::now().naive_utc(),
820            update_time: chrono::Utc::now().naive_utc(),
821            translator_file: None,
822            proofreader_file: None,
823            timer_file: None,
824            publish_link: None,
825            provider_file_oss_id: Some(10),
826            translator_file_oss_id: None,
827            proofreader_file_oss_id: None,
828            letterer_file_oss_id: None,
829            timer_file_oss_id: None,
830        };
831        assert_eq!(read_post_oss_id(&ep, "provider"), Some(10));
832        assert_eq!(read_post_oss_id(&ep, "unknown"), None);
833    }
834
835    #[test]
836    fn uses_config_cdn_only_for_current_or_missing_bucket() {
837        assert!(should_use_config_cdn(
838            None,
839            "",
840            "manga-trans-1317356496",
841            "ossdev.yuriful.top"
842        ));
843        assert!(should_use_config_cdn(
844            Some(""),
845            "https://manga-trans-1317356496.cos.ap-guangzhou.myqcloud.com/manga_1/file.txt",
846            "manga-trans-1317356496",
847            "ossdev.yuriful.top"
848        ));
849        assert!(should_use_config_cdn(
850            Some("manga-trans-1317356496"),
851            "https://ossdev.yuriful.top/manga_1/file.txt",
852            "manga-trans-1317356496",
853            "ossdev.yuriful.top"
854        ));
855    }
856
857    #[test]
858    fn rejects_config_cdn_for_different_bucket_or_domain() {
859        assert!(!should_use_config_cdn(
860            Some("manga-trans-prod-1317356496"),
861            "https://oss.yuriful.top/manga_1/file.txt",
862            "manga-trans-1317356496",
863            "ossdev.yuriful.top"
864        ));
865        assert!(!should_use_config_cdn(
866            None,
867            "https://manga-trans-prod-1317356496.cos.ap-guangzhou.myqcloud.com/manga_1/file.txt",
868            "manga-trans-1317356496",
869            "ossdev.yuriful.top"
870        ));
871        assert!(!should_use_config_cdn(
872            None,
873            "https://oss.yuriful.top/manga_1/file.txt",
874            "manga-trans-1317356496",
875            "ossdev.yuriful.top"
876        ));
877        assert!(!should_use_config_cdn(
878            None,
879            "https://unknown.example.com/manga_1/file.txt",
880            "manga-trans-1317356496",
881            "ossdev.yuriful.top"
882        ));
883    }
884
885    #[test]
886    fn extracts_bucket_from_cos_url() {
887        assert_eq!(
888            cos_bucket_from_url(
889                "https://manga-trans-prod-1317356496.cos.ap-guangzhou.myqcloud.com/manga_1/file.txt"
890            )
891            .as_deref(),
892            Some("manga-trans-prod-1317356496")
893        );
894    }
895}