feat: add Phase 3 TTS task pipeline
This commit is contained in:
parent
b0abd8b08a
commit
f5ddf761e4
@ -1,26 +1,29 @@
|
|||||||
# TTS 接口契约与适配边界
|
# TTS 接口契约与适配边界
|
||||||
|
|
||||||
状态:Phase 0 审计基线
|
状态:Phase 0 审计基线,已补充旧 TTS 工程只读核查事实
|
||||||
审计日期:2026-09-08
|
审计日期:2026-09-08
|
||||||
重要说明:本文件将“已核实事实”和“目标契约”分开。当前没有在仓库或测试机中发现可调用的现有 TTS 接口;目标样例不能作为现有能力证明。
|
重要说明:本文件将“已核实事实”和“目标契约”分开。已授权只读检查外部工程 `E:\RD\Kaotings\kts_site\KT26-0903_Big-TTS`,未复制或修改该工程。
|
||||||
|
|
||||||
## 1. 现有接口事实
|
## 1. 现有接口事实
|
||||||
|
|
||||||
| 项目 | 结论 | 证据/处理 |
|
| 项目 | 结论 | 证据/处理 |
|
||||||
| --- | --- | --- |
|
| --- | --- | --- |
|
||||||
| 上游地址 | 待确认 | 仓库无配置;测试机仅监听 SSH/DNS |
|
| 上游地址 | 代码默认:`http://43.248.188.28:44480`,路径 `/v1/audio/speech` | `ReadMe.md` 和 `server.js` 初始 SQLite 配置;不是实际运行地址证明 |
|
||||||
| 协议和方法 | 待确认 | 未发现 OpenAPI、客户端或请求样例 |
|
| 协议和方法 | `POST`,OpenAI 兼容 JSON 请求 | 旧工程 `/api/tts` 使用 `fetch` |
|
||||||
| 鉴权 | 待确认 | 未发现密钥名、Token 方案或服务账号 |
|
| 鉴权 | 配置有 API Key 时发送 `Authorization: Bearer <key>`;默认 Key 为空 | Key 只存旧工程服务端 SQLite,配置接口返回掩码 |
|
||||||
| 输入文本字段 | 待确认 | 规划只规定业务层语义,不代表上游字段 |
|
| 输入文本字段 | `input` | 旧工程将业务文本映射为 `input` |
|
||||||
| 音色 ID/语言 | 待确认 | 未发现音色清单 |
|
| 音色 ID/语言 | `voice`;默认 `default`,逗号分隔配置;无真实语言清单 | 旧工程 `tts_config` 和前端 `/api/config` |
|
||||||
| 参数范围 | 待确认 | 未发现 speed/pitch/format 等上游定义 |
|
| 参数范围 | `response_format`: `wav`/`mp3`;`speed`: `0.5` 至 `2.0`;文本 `1` 至 `5000` 字符 | `server.js` 服务端校验 |
|
||||||
| 同步/异步 | 待确认 | 未发现任务 ID 或结果轮询协议 |
|
| 模型 | 默认 `qwen3-tts`,可由旧工程管理员配置 | `ReadMe.md`、`tts_config` |
|
||||||
| 结果格式 | 待确认 | 未发现音频 MIME、字节流、URL 或 JSON 约定 |
|
| 同步/异步 | 同步 | 单次 `fetch` 等待完整响应,无 task ID |
|
||||||
| 错误语义 | 待确认 | 未发现状态码/错误码映射 |
|
| 结果格式 | 音频字节流;WAV `audio/wav`,MP3 `audio/mpeg` | 旧工程返回完整 `arrayBuffer` |
|
||||||
| 超时和大小限制 | 待确认 | 由联调和容量测试确定,不能凭规划虚构 |
|
| 错误语义 | 输入错误 `400`;未登录 `401`;上游非 2xx/超时统一 `502` | 旧工程业务层映射 |
|
||||||
| 幂等支持 | 待确认 | 不假设上游支持;业务层必须先自行防重复结算 |
|
| 超时和大小限制 | 上游超时 `120000ms`;请求体最多 `1 MiB`;文本最多 `5000` 字符;未见上游响应大小限制 | `body()` 与 `AbortSignal.timeout(120000)` |
|
||||||
|
| 幂等支持 | 未实现 | 无幂等键、provider task ID 或重复请求去重 |
|
||||||
|
|
||||||
真实请求测试因缺少授权上游地址、测试凭据和测试文本而未执行。不得将此状态写成“接口不可用”或“接口已通”,正确表述是“接口未核实”。
|
代码契约已核实;对默认地址的真实最小文本联调返回 `401 Unauthorized`,说明地址可达但当前请求缺少有效鉴权。不能把代码默认值写成已确认的实际部署地址。当前测试服务器没有 TTS 监听端口。
|
||||||
|
|
||||||
|
旧工程自身的业务接口为:登录后 `GET /api/config` 获取模型/音色元数据,`POST /api/tts` 同步返回音频,`GET /api/history` 返回当前用户最多 100 条历史;管理员配置接口为 `/api/admin/config`。Phase 3 只复用上游边界,不复制旧工程用户体系或整套工程。
|
||||||
|
|
||||||
## 2. 业务层目标契约
|
## 2. 业务层目标契约
|
||||||
|
|
||||||
@ -54,7 +57,7 @@
|
|||||||
}
|
}
|
||||||
```
|
```
|
||||||
|
|
||||||
服务端必须重新计算 Unicode 码点计量、校验文本长度/音色/参数/权益/并发和频率,不能信任客户端计数或隐藏字段。原始文本不进入普通业务日志。
|
服务端必须重新计算 Unicode 码点计量、校验文本长度/音色/参数/权益/并发和频率,不能信任客户端计数或隐藏字段。旧工程会把完整原始文本写入 SQLite 历史,Phase 3 不直接复制该隐私行为。
|
||||||
|
|
||||||
### 2.3 创建成功响应
|
### 2.3 创建成功响应
|
||||||
|
|
||||||
@ -135,10 +138,11 @@
|
|||||||
|
|
||||||
## 5. 待确认清单和验证顺序
|
## 5. 待确认清单和验证顺序
|
||||||
|
|
||||||
1. 由服务负责人提供已授权的上游主机、端点、测试凭据注入方式和网络允许范围。
|
1. 确认 `43.248.188.28:44480` 是否仍是获授权的测试上游;代码默认值不作为部署事实。
|
||||||
2. 获取真实 OpenAPI/接口文档或脱敏请求响应,确认同步/异步、任务查询、结果下载、音色、参数、错误和计量。
|
2. 提供该上游测试环境的 API Key 或其他鉴权凭据注入方式;本次无鉴权最小请求返回 `401`,未保存或输出响应内容。
|
||||||
3. 确认上游是否支持幂等、取消、重试和 provider task ID 查询;不支持时由业务 Worker 承担任务恢复与防重复结算。
|
3. 获取当前音色 ID、语言、参数和计量补充资料;旧工程只有 `default` 示例,没有真实清单。
|
||||||
4. 在测试环境执行最小成功、参数拒绝、超时、上游失败、重复请求和响应不明测试;不得用 Mock 结果替代。
|
4. 确认上游是否支持幂等、取消、重试和 provider task ID 查询;当前旧工程不提供这些能力。
|
||||||
5. 依据真实结果固定 `TTSVoice`、`TTSTask`、错误映射、超时、重试、音频校验和额度计量配置。
|
5. 在测试环境执行最小成功、参数拒绝、超时、上游失败、重复请求和响应不明测试;不得用 Mock 结果替代。
|
||||||
|
6. 依据真实结果固定 `TTSVoice`、`TTSTask`、错误映射、超时、重试、音频校验和额度计量配置。
|
||||||
|
|
||||||
在上述步骤完成前,Phase 3 不得宣称“真实 TTS 已接入”。
|
在上述步骤完成前,Phase 3 不得宣称“真实 TTS 已接入”。
|
||||||
|
|||||||
@ -1,14 +1,18 @@
|
|||||||
# TTS 上游资料请求清单
|
# TTS 上游资料请求清单
|
||||||
|
|
||||||
状态:Phase 2 核查仍未取得真实上游资料。
|
状态:已取得旧工程代码契约;实际部署状态、完整音色资料和联调授权仍待确认。
|
||||||
日期:2026-09-08
|
日期:2026-09-08
|
||||||
|
|
||||||
## 已核查范围
|
## 已核查范围
|
||||||
|
|
||||||
- 当前仓库:未发现 TTS 客户端、适配器、接口文档、音色清单、上游地址或凭据引用。
|
- 当前仓库:未发现新的 TTS 客户端或适配器;目标业务边界仍在 `API_CONTRACT.md`。
|
||||||
|
- 已授权旧工程:`E:\RD\Kaotings\kts_site\KT26-0903_Big-TTS`,只读核查 `ReadMe.md`、`server.js`、`package.json`、Docker 配置和前端调用。
|
||||||
|
- 旧工程代码默认上游:`http://43.248.188.28:44480/v1/audio/speech`,这是默认配置,不是实际部署证明。
|
||||||
|
- 旧工程实际调用是同步 `POST`,JSON 字段为 `model`、`input`、`voice`、`response_format`、`speed`,返回 WAV/MP3 字节流。
|
||||||
- 测试服务器:当前仅有 PostgreSQL、API、Web、Caddy 和 SSH 监听;未发现 TTS 监听端口或 TTS 服务配置。
|
- 测试服务器:当前仅有 PostgreSQL、API、Web、Caddy 和 SSH 监听;未发现 TTS 监听端口或 TTS 服务配置。
|
||||||
- `services/api`:仅实现认证、会话、权益和额度基础;没有 Phase 3 生成路由或上游调用。
|
- `services/api`:仅实现认证、会话、权益和额度基础;没有 Phase 3 生成路由或上游调用。
|
||||||
- `docs/API_CONTRACT.md`:只有脱敏目标契约,明确标记为待确认,不能替代真实上游协议。
|
- 旧工程没有真实音色语言元数据、异步任务查询、provider task ID 或幂等键实现。
|
||||||
|
- 代码默认地址的最小文本联调已执行:返回 `401 Unauthorized`;地址可达,但缺少有效测试鉴权凭据。
|
||||||
|
|
||||||
## 需要内部服务负责人提供
|
## 需要内部服务负责人提供
|
||||||
|
|
||||||
@ -24,6 +28,6 @@
|
|||||||
8. 连接/读取/任务超时、响应大小、计量口径和错误码映射。
|
8. 连接/读取/任务超时、响应大小、计量口径和错误码映射。
|
||||||
9. 最小联调用例、失败用例和允许保存的测试音频/文本保留期。
|
9. 最小联调用例、失败用例和允许保存的测试音频/文本保留期。
|
||||||
|
|
||||||
## 当前阻断
|
## 当前仍缺少
|
||||||
|
|
||||||
缺少以上资料前,Phase 3 不启动真实生成接入,不创建伪造音色,不使用 Mock 成功替代真实联调,不尝试扫描或猜测内部 TTS 地址。
|
当前最具体阻断是:缺少上游测试 API Key/鉴权凭据,无法完成成功音频联调。资料补齐前,Phase 3 不固定生产音色和计量策略,不创建伪造音色,不使用 Mock 成功替代真实联调,不扫描或猜测其他服务地址。
|
||||||
|
|||||||
@ -35,6 +35,10 @@ EMAIL_VERIFICATION_ENABLED=false
|
|||||||
PHONE_VERIFICATION_ENABLED=false
|
PHONE_VERIFICATION_ENABLED=false
|
||||||
TEST_FREE_PERIOD_LIMIT=10000
|
TEST_FREE_PERIOD_LIMIT=10000
|
||||||
TEST_VIP_PERIOD_LIMIT=100000
|
TEST_VIP_PERIOD_LIMIT=100000
|
||||||
|
TTS_UPSTREAM_URL=
|
||||||
|
TTS_API_KEY=
|
||||||
|
TTS_TIMEOUT_SECONDS=120
|
||||||
|
AUDIO_STORAGE_DIR=/home/flym/kaotings-audio
|
||||||
EOF
|
EOF
|
||||||
chmod 600 /etc/kaotings/api.env
|
chmod 600 /etc/kaotings/api.env
|
||||||
echo "database and API environment provisioned"
|
echo "database and API environment provisioned"
|
||||||
|
|||||||
@ -9,3 +9,7 @@ EMAIL_VERIFICATION_ENABLED=false
|
|||||||
PHONE_VERIFICATION_ENABLED=false
|
PHONE_VERIFICATION_ENABLED=false
|
||||||
TEST_FREE_PERIOD_LIMIT=10000
|
TEST_FREE_PERIOD_LIMIT=10000
|
||||||
TEST_VIP_PERIOD_LIMIT=100000
|
TEST_VIP_PERIOD_LIMIT=100000
|
||||||
|
TTS_UPSTREAM_URL=
|
||||||
|
TTS_API_KEY=
|
||||||
|
TTS_TIMEOUT_SECONDS=120
|
||||||
|
AUDIO_STORAGE_DIR=/home/flym/kaotings-audio
|
||||||
|
|||||||
@ -21,6 +21,10 @@ class Settings:
|
|||||||
phone_verification_enabled: bool
|
phone_verification_enabled: bool
|
||||||
test_free_period_limit: int
|
test_free_period_limit: int
|
||||||
test_vip_period_limit: int
|
test_vip_period_limit: int
|
||||||
|
tts_upstream_url: str
|
||||||
|
tts_api_key: str
|
||||||
|
tts_timeout_seconds: int
|
||||||
|
audio_storage_dir: str
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def from_env(cls) -> "Settings":
|
def from_env(cls) -> "Settings":
|
||||||
@ -40,6 +44,10 @@ class Settings:
|
|||||||
phone_verification_enabled=as_bool(os.getenv("PHONE_VERIFICATION_ENABLED")),
|
phone_verification_enabled=as_bool(os.getenv("PHONE_VERIFICATION_ENABLED")),
|
||||||
test_free_period_limit=int(os.getenv("TEST_FREE_PERIOD_LIMIT", "10000")),
|
test_free_period_limit=int(os.getenv("TEST_FREE_PERIOD_LIMIT", "10000")),
|
||||||
test_vip_period_limit=int(os.getenv("TEST_VIP_PERIOD_LIMIT", "100000")),
|
test_vip_period_limit=int(os.getenv("TEST_VIP_PERIOD_LIMIT", "100000")),
|
||||||
|
tts_upstream_url=os.getenv("TTS_UPSTREAM_URL", "").rstrip("/"),
|
||||||
|
tts_api_key=os.getenv("TTS_API_KEY", ""),
|
||||||
|
tts_timeout_seconds=int(os.getenv("TTS_TIMEOUT_SECONDS", "120")),
|
||||||
|
audio_storage_dir=os.getenv("AUDIO_STORAGE_DIR", "./data/audio"),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@ -1,10 +1,15 @@
|
|||||||
|
import asyncio
|
||||||
|
import hashlib
|
||||||
|
import json
|
||||||
from datetime import datetime, timedelta, timezone
|
from datetime import datetime, timedelta, timezone
|
||||||
|
from pathlib import Path
|
||||||
from typing import Any
|
from typing import Any
|
||||||
from uuid import UUID
|
from uuid import UUID
|
||||||
|
|
||||||
import psycopg
|
import psycopg
|
||||||
from fastapi import Depends, FastAPI, HTTPException, Request, Response, status
|
from fastapi import Depends, FastAPI, HTTPException, Request, Response, status
|
||||||
from fastapi.responses import JSONResponse
|
import httpx
|
||||||
|
from fastapi.responses import FileResponse, JSONResponse
|
||||||
from psycopg import Connection
|
from psycopg import Connection
|
||||||
from psycopg.types.json import Json
|
from psycopg.types.json import Json
|
||||||
|
|
||||||
@ -21,6 +26,9 @@ from .schemas import (
|
|||||||
UserPublic,
|
UserPublic,
|
||||||
VerificationConfirmRequest,
|
VerificationConfirmRequest,
|
||||||
VerificationSendRequest,
|
VerificationSendRequest,
|
||||||
|
TtsTaskRequest,
|
||||||
|
TtsTaskPublic,
|
||||||
|
TtsVoicePublic,
|
||||||
normalize_email,
|
normalize_email,
|
||||||
)
|
)
|
||||||
from .security import (
|
from .security import (
|
||||||
@ -110,6 +118,104 @@ def usage_view(quota: dict) -> dict[str, Any]:
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def task_view(connection: Connection, task: dict) -> dict[str, Any]:
|
||||||
|
audio = connection.execute("SELECT expires_at, status FROM audio_files WHERE task_id = %s", (task["id"],)).fetchone()
|
||||||
|
return {
|
||||||
|
"id": task["id"], "status": task["status"], "text_length": task["text_length"],
|
||||||
|
"voice_id": task["provider_voice_id"], "parameters": task["parameters"],
|
||||||
|
"error_code": task["error_code"], "audio_available": bool(audio and audio["status"] == "available"),
|
||||||
|
"audio_expires_at": audio["expires_at"] if audio else None, "created_at": task["created_at"],
|
||||||
|
"started_at": task["started_at"], "finished_at": task["finished_at"],
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def claim_next_task() -> dict | None:
|
||||||
|
with psycopg.connect(settings.database_url, row_factory=psycopg.rows.dict_row) as connection:
|
||||||
|
task = connection.execute("SELECT * FROM tts_tasks WHERE status = 'queued' ORDER BY created_at FOR UPDATE SKIP LOCKED LIMIT 1").fetchone()
|
||||||
|
if not task:
|
||||||
|
return None
|
||||||
|
updated = connection.execute("UPDATE tts_tasks SET status = 'running', started_at = now(), updated_at = now(), attempt_count = attempt_count + 1 WHERE id = %s RETURNING *", (task["id"],)).fetchone()
|
||||||
|
connection.commit()
|
||||||
|
return updated
|
||||||
|
|
||||||
|
|
||||||
|
def finish_task_failure(task: dict, code: str) -> None:
|
||||||
|
with psycopg.connect(settings.database_url, row_factory=psycopg.rows.dict_row) as connection:
|
||||||
|
connection.execute("UPDATE tts_tasks SET status = 'failed', error_code = %s, finished_at = now(), updated_at = now() WHERE id = %s AND status = 'running'", (code, task["id"]))
|
||||||
|
connection.execute("UPDATE quota_accounts SET reserved = GREATEST(0, reserved - %s), version = version + 1 WHERE id = %s", (task["reserved_amount"], task["quota_account_id"]))
|
||||||
|
connection.execute("INSERT INTO usage_records(user_id, quota_account_id, task_id, type, amount) VALUES (%s, %s, %s, 'release', %s)", (task["user_id"], task["quota_account_id"], task["id"], task["reserved_amount"]))
|
||||||
|
connection.commit()
|
||||||
|
|
||||||
|
|
||||||
|
def finish_task_success(task: dict, audio: bytes, mime_type: str) -> None:
|
||||||
|
storage_dir = Path(settings.audio_storage_dir)
|
||||||
|
storage_dir.mkdir(parents=True, exist_ok=True)
|
||||||
|
extension = "mp3" if mime_type == "audio/mpeg" else "wav"
|
||||||
|
storage_key = f"{task['user_id']}/{task['id']}.{extension}"
|
||||||
|
target = storage_dir / storage_key
|
||||||
|
target.parent.mkdir(parents=True, exist_ok=True)
|
||||||
|
target.write_bytes(audio)
|
||||||
|
checksum = hashlib.sha256(audio).hexdigest()
|
||||||
|
with psycopg.connect(settings.database_url, row_factory=psycopg.rows.dict_row) as connection:
|
||||||
|
connection.execute("INSERT INTO audio_files(task_id, owner_id, storage_key, mime_type, size_bytes, checksum) VALUES (%s, %s, %s, %s, %s, %s)", (task["id"], task["user_id"], storage_key, mime_type, len(audio), checksum))
|
||||||
|
connection.execute("UPDATE tts_tasks SET status = 'succeeded', finished_at = now(), updated_at = now() WHERE id = %s AND status = 'running'", (task["id"],))
|
||||||
|
connection.execute("UPDATE quota_accounts SET reserved = GREATEST(0, reserved - %s), used = used + %s, version = version + 1 WHERE id = %s", (task["reserved_amount"], task["reserved_amount"], task["quota_account_id"]))
|
||||||
|
connection.execute("INSERT INTO usage_records(user_id, quota_account_id, task_id, type, amount) VALUES (%s, %s, %s, 'consume', %s)", (task["user_id"], task["quota_account_id"], task["id"], task["reserved_amount"]))
|
||||||
|
connection.commit()
|
||||||
|
|
||||||
|
|
||||||
|
async def process_task(task: dict) -> None:
|
||||||
|
if not settings.tts_upstream_url:
|
||||||
|
await asyncio.to_thread(finish_task_failure, task, "UPSTREAM_NOT_CONFIGURED")
|
||||||
|
return
|
||||||
|
parameters = task["parameters"]
|
||||||
|
response_format = parameters.get("format", "wav")
|
||||||
|
payload = {"model": parameters.get("model", "qwen3-tts"), "input": task["text"], "voice": task["provider_voice_id"], "response_format": response_format, "speed": parameters.get("speed", 1.0)}
|
||||||
|
headers = {"Content-Type": "application/json"}
|
||||||
|
if settings.tts_api_key:
|
||||||
|
headers["Authorization"] = f"Bearer {settings.tts_api_key}"
|
||||||
|
try:
|
||||||
|
async with httpx.AsyncClient(timeout=settings.tts_timeout_seconds) as client:
|
||||||
|
response = await client.post(f"{settings.tts_upstream_url}/v1/audio/speech", headers=headers, json=payload)
|
||||||
|
if response.status_code == 401:
|
||||||
|
raise RuntimeError("UPSTREAM_UNAUTHORIZED")
|
||||||
|
if response.status_code >= 400:
|
||||||
|
raise RuntimeError(f"UPSTREAM_HTTP_{response.status_code}")
|
||||||
|
audio = response.content
|
||||||
|
if not audio or len(audio) > 50 * 1024 * 1024:
|
||||||
|
raise RuntimeError("UPSTREAM_AUDIO_INVALID")
|
||||||
|
mime = "audio/mpeg" if response_format == "mp3" else "audio/wav"
|
||||||
|
await asyncio.to_thread(finish_task_success, task, audio, mime)
|
||||||
|
except httpx.TimeoutException:
|
||||||
|
await asyncio.to_thread(finish_task_failure, task, "UPSTREAM_TIMEOUT")
|
||||||
|
except Exception as exc:
|
||||||
|
await asyncio.to_thread(finish_task_failure, task, str(exc)[:80])
|
||||||
|
|
||||||
|
|
||||||
|
async def task_worker() -> None:
|
||||||
|
while True:
|
||||||
|
task = await asyncio.to_thread(claim_next_task)
|
||||||
|
if task:
|
||||||
|
await process_task(task)
|
||||||
|
else:
|
||||||
|
await asyncio.sleep(1)
|
||||||
|
|
||||||
|
|
||||||
|
worker_task: asyncio.Task | None = None
|
||||||
|
|
||||||
|
|
||||||
|
@app.on_event("startup")
|
||||||
|
async def start_worker():
|
||||||
|
global worker_task
|
||||||
|
worker_task = asyncio.create_task(task_worker())
|
||||||
|
|
||||||
|
|
||||||
|
@app.on_event("shutdown")
|
||||||
|
async def stop_worker():
|
||||||
|
if worker_task:
|
||||||
|
worker_task.cancel()
|
||||||
|
|
||||||
|
|
||||||
@app.exception_handler(HTTPException)
|
@app.exception_handler(HTTPException)
|
||||||
async def http_error_handler(_: Request, exc: HTTPException):
|
async def http_error_handler(_: Request, exc: HTTPException):
|
||||||
detail = exc.detail if isinstance(exc.detail, dict) else {"code": "REQUEST_FAILED", "message": str(exc.detail)}
|
detail = exc.detail if isinstance(exc.detail, dict) else {"code": "REQUEST_FAILED", "message": str(exc.detail)}
|
||||||
@ -233,6 +339,93 @@ def verification_confirm(payload: VerificationConfirmRequest, request: Request,
|
|||||||
return {"status": "ok"}
|
return {"status": "ok"}
|
||||||
|
|
||||||
|
|
||||||
|
@app.get("/api/v1/tts/voices", response_model=list[TtsVoicePublic])
|
||||||
|
def tts_voices(connection: Connection = Depends(get_connection)):
|
||||||
|
return connection.execute("SELECT id, provider_voice_id, name, language, description, supported_parameters FROM tts_voices WHERE enabled = true ORDER BY name").fetchall()
|
||||||
|
|
||||||
|
|
||||||
|
@app.post("/api/v1/tts/tasks", status_code=202)
|
||||||
|
def create_tts_task(payload: TtsTaskRequest, request: Request, connection: Connection = Depends(get_connection), user: dict = Depends(current_user), _: None = Depends(require_csrf)):
|
||||||
|
idempotency_key = request.headers.get("idempotency-key", "").strip()
|
||||||
|
if len(idempotency_key) < 8 or len(idempotency_key) > 120:
|
||||||
|
raise error("IDEMPOTENCY_REQUIRED", "需要有效的 Idempotency-Key", 400)
|
||||||
|
plan = effective_plan(connection, user)
|
||||||
|
policy = connection.execute("SELECT * FROM plan_policies WHERE code = %s", (plan,)).fetchone()
|
||||||
|
text_length = len(payload.text)
|
||||||
|
if not policy or text_length > policy["max_text_length"]:
|
||||||
|
raise error("TEXT_TOO_LONG", "文本超过当前计划限制", 422)
|
||||||
|
parameters = dict(payload.parameters)
|
||||||
|
response_format = str(parameters.get("format", "wav")).lower()
|
||||||
|
speed = float(parameters.get("speed", 1.0))
|
||||||
|
if response_format not in {"wav", "mp3"} or not 0.5 <= speed <= 2:
|
||||||
|
raise error("INVALID_PARAMETERS", "音频格式或语速不可用", 422)
|
||||||
|
voice = connection.execute("SELECT * FROM tts_voices WHERE provider_voice_id = %s AND enabled = true", (payload.voice_id,)).fetchone()
|
||||||
|
if not voice or plan not in (voice["allowed_plans"] or ["free", "vip"]):
|
||||||
|
raise error("VOICE_NOT_ALLOWED", "音色不可用", 422)
|
||||||
|
request_hash = hashlib.sha256(json.dumps({"text": payload.text, "voice": payload.voice_id, "parameters": parameters}, ensure_ascii=False, sort_keys=True).encode()).hexdigest()
|
||||||
|
existing = connection.execute("SELECT * FROM tts_tasks WHERE user_id = %s AND idempotency_key = %s", (user["id"], idempotency_key)).fetchone()
|
||||||
|
if existing:
|
||||||
|
if existing["request_hash"] != request_hash:
|
||||||
|
raise error("IDEMPOTENCY_CONFLICT", "幂等键已用于其他请求", 409)
|
||||||
|
return task_view(connection, existing)
|
||||||
|
quota = ensure_quota(connection, user["id"], plan)
|
||||||
|
quota = connection.execute("SELECT * FROM quota_accounts WHERE id = %s FOR UPDATE", (quota["id"],)).fetchone()
|
||||||
|
available = quota["limit_snapshot"] + quota["adjustment"] - quota["used"] - quota["reserved"]
|
||||||
|
if available < text_length:
|
||||||
|
raise error("QUOTA_EXCEEDED", "额度不足", 409)
|
||||||
|
task = connection.execute(
|
||||||
|
"""
|
||||||
|
INSERT INTO tts_tasks(user_id, text, text_length, voice_id, provider_voice_id, parameters, idempotency_key, request_hash, policy_version, quota_account_id, reserved_amount)
|
||||||
|
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) RETURNING *
|
||||||
|
""",
|
||||||
|
(user["id"], payload.text, text_length, voice["id"], payload.voice_id, Json(parameters), idempotency_key, request_hash, policy["version"], quota["id"], text_length),
|
||||||
|
).fetchone()
|
||||||
|
connection.execute("UPDATE quota_accounts SET reserved = reserved + %s, version = version + 1 WHERE id = %s", (text_length, quota["id"]))
|
||||||
|
connection.execute("INSERT INTO usage_records(user_id, quota_account_id, task_id, type, amount, idempotency_key) VALUES (%s, %s, %s, 'reserve', %s, %s)", (user["id"], quota["id"], task["id"], text_length, f"reserve:{task['id']}"))
|
||||||
|
connection.commit()
|
||||||
|
return task_view(connection, task)
|
||||||
|
|
||||||
|
|
||||||
|
@app.get("/api/v1/tts/tasks", response_model=list[TtsTaskPublic])
|
||||||
|
def list_tts_tasks(user: dict = Depends(current_user), connection: Connection = Depends(get_connection)):
|
||||||
|
tasks = connection.execute("SELECT * FROM tts_tasks WHERE user_id = %s ORDER BY created_at DESC LIMIT 100", (user["id"],)).fetchall()
|
||||||
|
return [task_view(connection, task) for task in tasks]
|
||||||
|
|
||||||
|
|
||||||
|
def owned_task(task_id: UUID, user: dict, connection: Connection) -> dict:
|
||||||
|
task = connection.execute("SELECT * FROM tts_tasks WHERE id = %s", (task_id,)).fetchone()
|
||||||
|
if not task or (task["user_id"] != user["id"] and user["role"] != "admin"):
|
||||||
|
raise error("NOT_FOUND", "任务不存在", 404)
|
||||||
|
return task
|
||||||
|
|
||||||
|
|
||||||
|
@app.get("/api/v1/tts/tasks/{task_id}", response_model=TtsTaskPublic)
|
||||||
|
def get_tts_task(task_id: UUID, user: dict = Depends(current_user), connection: Connection = Depends(get_connection)):
|
||||||
|
return task_view(connection, owned_task(task_id, user, connection))
|
||||||
|
|
||||||
|
|
||||||
|
def audio_response(task_id: UUID, user: dict, connection: Connection, download: bool):
|
||||||
|
task = owned_task(task_id, user, connection)
|
||||||
|
audio = connection.execute("SELECT * FROM audio_files WHERE task_id = %s AND status = 'available'", (task_id,)).fetchone()
|
||||||
|
if task["status"] != "succeeded" or not audio:
|
||||||
|
raise error("AUDIO_NOT_AVAILABLE", "音频尚不可用", 404)
|
||||||
|
path = Path(settings.audio_storage_dir) / audio["storage_key"]
|
||||||
|
if not path.is_file():
|
||||||
|
raise error("AUDIO_NOT_AVAILABLE", "音频文件不可用", 404)
|
||||||
|
filename = f"kaotings-{task_id}.{'mp3' if audio['mime_type'] == 'audio/mpeg' else 'wav'}" if download else None
|
||||||
|
return FileResponse(path, media_type=audio["mime_type"], filename=filename)
|
||||||
|
|
||||||
|
|
||||||
|
@app.get("/api/v1/tts/tasks/{task_id}/audio")
|
||||||
|
def play_tts_audio(task_id: UUID, user: dict = Depends(current_user), connection: Connection = Depends(get_connection)):
|
||||||
|
return audio_response(task_id, user, connection, False)
|
||||||
|
|
||||||
|
|
||||||
|
@app.get("/api/v1/tts/tasks/{task_id}/download")
|
||||||
|
def download_tts_audio(task_id: UUID, user: dict = Depends(current_user), connection: Connection = Depends(get_connection)):
|
||||||
|
return audio_response(task_id, user, connection, True)
|
||||||
|
|
||||||
|
|
||||||
@app.get("/api/v1/admin/users")
|
@app.get("/api/v1/admin/users")
|
||||||
def admin_users(user: dict = Depends(require_admin), connection: Connection = Depends(get_connection)):
|
def admin_users(user: dict = Depends(require_admin), connection: Connection = Depends(get_connection)):
|
||||||
rows = connection.execute("SELECT id, email, phone, role, plan, status, email_verified, phone_verified, created_at, last_login_at FROM users ORDER BY created_at DESC LIMIT 100").fetchall()
|
rows = connection.execute("SELECT id, email, phone, role, plan, status, email_verified, phone_verified, created_at, last_login_at FROM users ORDER BY created_at DESC LIMIT 100").fetchall()
|
||||||
|
|||||||
@ -65,5 +65,34 @@ class VerificationSendRequest(BaseModel):
|
|||||||
purpose: Literal["registration", "contact_binding", "contact_change", "password_reset"]
|
purpose: Literal["registration", "contact_binding", "contact_change", "password_reset"]
|
||||||
|
|
||||||
|
|
||||||
|
class TtsTaskRequest(BaseModel):
|
||||||
|
text: str = Field(min_length=1, max_length=10000)
|
||||||
|
voice_id: str = Field(min_length=1, max_length=120)
|
||||||
|
parameters: dict[str, Any] = Field(default_factory=dict)
|
||||||
|
|
||||||
|
|
||||||
|
class TtsVoicePublic(BaseModel):
|
||||||
|
id: UUID
|
||||||
|
provider_voice_id: str
|
||||||
|
name: str
|
||||||
|
language: str | None
|
||||||
|
description: str | None
|
||||||
|
supported_parameters: dict[str, Any]
|
||||||
|
|
||||||
|
|
||||||
|
class TtsTaskPublic(BaseModel):
|
||||||
|
id: UUID
|
||||||
|
status: str
|
||||||
|
text_length: int
|
||||||
|
voice_id: str
|
||||||
|
parameters: dict[str, Any]
|
||||||
|
error_code: str | None = None
|
||||||
|
audio_available: bool = False
|
||||||
|
audio_expires_at: datetime | None = None
|
||||||
|
created_at: datetime
|
||||||
|
started_at: datetime | None = None
|
||||||
|
finished_at: datetime | None = None
|
||||||
|
|
||||||
|
|
||||||
def normalize_email(value: str) -> str:
|
def normalize_email(value: str) -> str:
|
||||||
return value.strip().lower()
|
return value.strip().lower()
|
||||||
|
|||||||
56
services/api/migrations/002_tts.sql
Normal file
56
services/api/migrations/002_tts.sql
Normal file
@ -0,0 +1,56 @@
|
|||||||
|
CREATE TABLE IF NOT EXISTS tts_voices (
|
||||||
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||||
|
provider_voice_id TEXT NOT NULL UNIQUE,
|
||||||
|
name TEXT NOT NULL,
|
||||||
|
language TEXT,
|
||||||
|
description TEXT,
|
||||||
|
enabled BOOLEAN NOT NULL DEFAULT true,
|
||||||
|
supported_parameters JSONB NOT NULL DEFAULT '{}'::jsonb,
|
||||||
|
allowed_plans JSONB NOT NULL DEFAULT '["free", "vip"]'::jsonb,
|
||||||
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||||
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||||
|
);
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS tts_tasks (
|
||||||
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||||
|
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
||||||
|
text TEXT NOT NULL,
|
||||||
|
text_length INTEGER NOT NULL CHECK (text_length > 0),
|
||||||
|
voice_id UUID REFERENCES tts_voices(id),
|
||||||
|
provider_voice_id TEXT NOT NULL,
|
||||||
|
parameters JSONB NOT NULL DEFAULT '{}'::jsonb,
|
||||||
|
status TEXT NOT NULL DEFAULT 'queued' CHECK (status IN ('queued', 'running', 'succeeded', 'failed')),
|
||||||
|
provider_task_id TEXT,
|
||||||
|
idempotency_key TEXT NOT NULL,
|
||||||
|
request_hash TEXT NOT NULL,
|
||||||
|
policy_version INTEGER NOT NULL,
|
||||||
|
quota_account_id UUID NOT NULL REFERENCES quota_accounts(id),
|
||||||
|
reserved_amount BIGINT NOT NULL CHECK (reserved_amount >= 0),
|
||||||
|
error_code TEXT,
|
||||||
|
attempt_count INTEGER NOT NULL DEFAULT 0,
|
||||||
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||||
|
started_at TIMESTAMPTZ,
|
||||||
|
finished_at TIMESTAMPTZ,
|
||||||
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||||
|
UNIQUE(user_id, idempotency_key)
|
||||||
|
);
|
||||||
|
CREATE INDEX IF NOT EXISTS tts_tasks_user_created_idx ON tts_tasks(user_id, created_at DESC);
|
||||||
|
CREATE INDEX IF NOT EXISTS tts_tasks_queue_idx ON tts_tasks(status, updated_at);
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS audio_files (
|
||||||
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||||
|
task_id UUID NOT NULL UNIQUE REFERENCES tts_tasks(id) ON DELETE CASCADE,
|
||||||
|
owner_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
||||||
|
storage_key TEXT NOT NULL UNIQUE,
|
||||||
|
mime_type TEXT NOT NULL,
|
||||||
|
size_bytes BIGINT NOT NULL CHECK (size_bytes >= 0),
|
||||||
|
checksum TEXT NOT NULL,
|
||||||
|
duration_ms INTEGER,
|
||||||
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||||
|
expires_at TIMESTAMPTZ,
|
||||||
|
status TEXT NOT NULL DEFAULT 'available' CHECK (status IN ('available', 'expired', 'deleted'))
|
||||||
|
);
|
||||||
|
|
||||||
|
INSERT INTO tts_voices(provider_voice_id, name, description, supported_parameters, allowed_plans)
|
||||||
|
VALUES ('default', 'default', '旧 TTS 工程代码中的默认音色标识', '{"format":["wav","mp3"],"speed":{"min":0.5,"max":2}}'::jsonb, '["free","vip"]'::jsonb)
|
||||||
|
ON CONFLICT (provider_voice_id) DO NOTHING;
|
||||||
@ -4,3 +4,4 @@ psycopg[binary]==3.2.9
|
|||||||
argon2-cffi==25.1.0
|
argon2-cffi==25.1.0
|
||||||
pydantic==2.11.7
|
pydantic==2.11.7
|
||||||
email-validator==2.2.0
|
email-validator==2.2.0
|
||||||
|
httpx==0.28.1
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user