13-douyinInsightManager.js 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641
  1. /**
  2. * 云函数:douyinInsightManager(二阶段:爆款分析、选题池、日报)
  3. * actions:
  4. * analysisCreate | analysisList | analysisGet | analysisUpdate
  5. * topicCreate | topicList | topicUpdate | topicArchive
  6. * dailyReportCreate | dailyReportList | dailyReportGet
  7. * transcriptStart | transcriptGet
  8. *
  9. * 说明:本函数只负责业务资产的持久化和账号隔离。抖音原始数据抓取仍由 12-douyinManager 负责。
  10. */
  11. const PARSE_API_HOST = readEnv('PARSE_API_HOST') || 'https://server.fmode.cn';
  12. const PARSE_APP_ID = readEnv('PARSE_APP_ID') || 'ncloudmaster';
  13. const DOUYIN_API_BASE_URL = (readEnv('DOUYIN_API_BASE_URL') || readEnv('VOC_SOCIAL_BASE_URL') || 'https://server.fmode.cn/api/voc-social').replace(/\/+$/, '');
  14. const IS_TIKHUB_DIRECT = /api\.tikhub\.io/i.test(DOUYIN_API_BASE_URL);
  15. const LOCAL_VOC_TOKEN_FALLBACK = 'r:33c57d404c8fffc9b19199a4da0bb663';
  16. const VOC_SOCIAL_TOKEN = readEnv('DOUYIN_API_TOKEN') || readEnv('VOC_TOKEN') || readEnv('TRANSCRIPTION_VOC_TOKEN') || readEnv('VOICE_TOKEN') || readEnv('OPENCLAW_VOC_TOKEN') || readEnv('VOC_SOCIAL_TOKEN') || LOCAL_VOC_TOKEN_FALLBACK;
  17. const TIKHUB_TOKEN = readEnv('TIKHUB_TOKEN') || 'gqsZHfMWgAiMwV+ITbmZy0qALADWBZVS7QnV7kKJe9CwzgWgJG+7bwK+GQ==';
  18. const DOUYIN_API_TOKEN = IS_TIKHUB_DIRECT ? (readEnv('DOUYIN_API_TOKEN') || TIKHUB_TOKEN) : VOC_SOCIAL_TOKEN;
  19. const TRANSCRIPTION_GATEWAY = (readEnv('IFLYTEK_GATEWAY_BASE_URL') || 'https://server.fmode.cn/api/apig/transcription').replace(/\/+$/, '');
  20. const VOC_TOKEN = readEnv('VOC_TOKEN') || readEnv('TRANSCRIPTION_VOC_TOKEN') || readEnv('VOICE_TOKEN') || readEnv('OPENCLAW_VOC_TOKEN') || readEnv('VOC_SOCIAL_TOKEN') || LOCAL_VOC_TOKEN_FALLBACK;
  21. async function handler(request, response) {
  22. try {
  23. const action = pickParam(request, 'action') || '';
  24. if (action === 'diagnose') {
  25. return diagnose(request, response);
  26. }
  27. await ensureTables();
  28. const session = await requireSession(request);
  29. const requestedUserId = pickParam(request, 'userId') || '';
  30. if (requestedUserId && requestedUserId !== session.userId) {
  31. return response.json({ code: 403, success: false, error: '没有访问该账号数据的权限' });
  32. }
  33. const userId = session.userId;
  34. if (action === 'analysisCreate') return createRow(response, 'VideoflowViralAnalysis', userId, pickParam(request, 'analysis', 'data') || {});
  35. if (action === 'analysisList') return listRows(response, 'VideoflowViralAnalysis', userId, pickParam(request, 'limit') || 100);
  36. if (action === 'analysisGet') return getRow(response, 'VideoflowViralAnalysis', userId, pickParam(request, 'id'));
  37. if (action === 'analysisUpdate') return updateRow(response, 'VideoflowViralAnalysis', userId, pickParam(request, 'id'), pickParam(request, 'patch', 'analysis', 'data') || {});
  38. if (action === 'topicCreate') return createRow(response, 'VideoflowTopicIdea', userId, pickParam(request, 'topic', 'data') || {});
  39. if (action === 'topicList') return listRows(response, 'VideoflowTopicIdea', userId, pickParam(request, 'limit') || 500);
  40. if (action === 'topicUpdate') return updateRow(response, 'VideoflowTopicIdea', userId, pickParam(request, 'id'), pickParam(request, 'patch', 'topic', 'data') || {});
  41. if (action === 'topicArchive') return updateRow(response, 'VideoflowTopicIdea', userId, pickParam(request, 'id'), { status: 'archived' });
  42. if (action === 'dailyReportCreate') return createRow(response, 'VideoflowDailyReport', userId, pickParam(request, 'report', 'data') || {});
  43. if (action === 'dailyReportList') return listRows(response, 'VideoflowDailyReport', userId, pickParam(request, 'limit') || 100);
  44. if (action === 'dailyReportGet') return getRow(response, 'VideoflowDailyReport', userId, pickParam(request, 'id'));
  45. if (action === 'transcriptStart') return startTranscript(request, response, userId);
  46. if (action === 'transcriptGet') return getTranscript(request, response, userId);
  47. response.json({ code: 400, success: false, error: `未知 action: ${action}` });
  48. } catch (error) {
  49. console.error('douyinInsightManager failed:', error.message);
  50. response.json({ code: 500, success: false, error: error.message });
  51. }
  52. }
  53. async function ensureTables() {
  54. for (const table of ['VideoflowViralAnalysis', 'VideoflowTopicIdea', 'VideoflowDailyReport', 'VideoflowTranscriptJob']) {
  55. await Psql.query(`
  56. CREATE TABLE IF NOT EXISTS "${table}" (
  57. "objectId" VARCHAR(50) PRIMARY KEY,
  58. "bizId" VARCHAR(255) NOT NULL,
  59. "userId" VARCHAR(255) NOT NULL,
  60. "data" JSONB NOT NULL DEFAULT '{}',
  61. "status" VARCHAR(50) DEFAULT '',
  62. "createdAt" TIMESTAMPTZ DEFAULT NOW(),
  63. "updatedAt" TIMESTAMPTZ DEFAULT NOW()
  64. )
  65. `);
  66. await Psql.query(`DROP INDEX IF EXISTS idx_${table.toLowerCase()}_biz`);
  67. await Psql.query(`CREATE UNIQUE INDEX IF NOT EXISTS idx_${table.toLowerCase()}_user_biz ON "${table}" ("userId", "bizId")`);
  68. await Psql.query(`CREATE INDEX IF NOT EXISTS idx_${table.toLowerCase()}_user ON "${table}" ("userId")`);
  69. }
  70. await Psql.query(`
  71. CREATE TABLE IF NOT EXISTS "AppSession" (
  72. "token" VARCHAR(120) PRIMARY KEY,
  73. "userId" VARCHAR(50) NOT NULL,
  74. "expiresAt" TIMESTAMPTZ NOT NULL,
  75. "createdAt" TIMESTAMPTZ DEFAULT NOW()
  76. )
  77. `);
  78. }
  79. async function createRow(response, table, userId, data) {
  80. const now = new Date().toISOString();
  81. const bizId = data.id || data.bizId || generateId();
  82. const merged = { ...data, id: bizId, userId, createdAt: data.createdAt || now, updatedAt: now };
  83. const existing = await Psql.query(
  84. `SELECT * FROM "${table}" WHERE "bizId"=$1 AND "userId"=$2 LIMIT 1`,
  85. [bizId, userId]
  86. );
  87. if (existing.length) {
  88. await Psql.query(
  89. `UPDATE "${table}" SET "data"=$1, "status"=$2, "updatedAt"=NOW() WHERE "bizId"=$3 AND "userId"=$4`,
  90. [JSON.stringify(merged), merged.status || '', bizId, userId]
  91. );
  92. return response.json({ code: 200, success: true, data: merged });
  93. }
  94. await Psql.query(
  95. `INSERT INTO "${table}" ("objectId","bizId","userId","data","status")
  96. VALUES ($1,$2,$3,$4,$5)`,
  97. [generateId(), bizId, userId, JSON.stringify(merged), merged.status || '']
  98. );
  99. response.json({ code: 200, success: true, data: merged });
  100. }
  101. async function listRows(response, table, userId, limit) {
  102. const parsedLimit = parseInt(limit || '100', 10);
  103. const safeLimit = Number.isFinite(parsedLimit) ? Math.min(Math.max(parsedLimit, 1), 1000) : 100;
  104. const rows = await Psql.query(
  105. `SELECT * FROM "${table}" WHERE "userId"=$1 ORDER BY "updatedAt" DESC LIMIT $2`,
  106. [userId, safeLimit]
  107. );
  108. response.json({ code: 200, success: true, data: rows.map(rowToObj) });
  109. }
  110. async function getRow(response, table, userId, id) {
  111. if (!id) return response.json({ code: 400, success: false, error: '缺少 id' });
  112. const rows = await Psql.query(
  113. `SELECT * FROM "${table}" WHERE "bizId"=$1 AND "userId"=$2 LIMIT 1`,
  114. [id, userId]
  115. );
  116. if (!rows.length) return response.json({ code: 404, success: false, error: '未找到记录' });
  117. response.json({ code: 200, success: true, data: rowToObj(rows[0]) });
  118. }
  119. async function updateRow(response, table, userId, id, patch) {
  120. if (!id) return response.json({ code: 400, success: false, error: '缺少 id' });
  121. const rows = await Psql.query(
  122. `SELECT * FROM "${table}" WHERE "bizId"=$1 AND "userId"=$2 LIMIT 1`,
  123. [id, userId]
  124. );
  125. if (!rows.length) return response.json({ code: 404, success: false, error: '未找到记录' });
  126. const merged = { ...rowToObj(rows[0]), ...patch, id, userId, updatedAt: new Date().toISOString() };
  127. await Psql.query(
  128. `UPDATE "${table}" SET "data"=$1, "status"=$2, "updatedAt"=NOW() WHERE "bizId"=$3 AND "userId"=$4`,
  129. [JSON.stringify(merged), merged.status || '', id, userId]
  130. );
  131. response.json({ code: 200, success: true, data: merged });
  132. }
  133. async function startTranscript(request, response, userId) {
  134. const awemeId = clean(pickParam(request, 'awemeId'));
  135. const analysisId = clean(pickParam(request, 'analysisId'));
  136. if (!awemeId) return response.json({ code: 400, success: false, error: '缺少 awemeId' });
  137. const now = new Date().toISOString();
  138. const job = {
  139. id: `transcript_${Date.now()}_${Math.random().toString(36).slice(2, 8)}`,
  140. awemeId,
  141. analysisId,
  142. provider: 'iflytek-gateway',
  143. status: 'pending',
  144. warnings: [],
  145. createdAt: now,
  146. updatedAt: now,
  147. };
  148. if (!VOC_TOKEN) {
  149. job.status = 'needs_provider_config';
  150. job.warnings.push('云函数未配置 VOC_TOKEN、TRANSCRIPTION_VOC_TOKEN 或 VOICE_TOKEN,无法调用转写网关。');
  151. await upsertTranscriptJob(userId, job);
  152. return response.json({ code: 200, success: true, data: job });
  153. }
  154. try {
  155. const detail = await fetchDouyinDetail(awemeId);
  156. const media = selectAudioCandidate(detail);
  157. if (!media.url) {
  158. job.status = 'needs_media';
  159. job.warnings.push('视频详情中未找到可直接提交转写的音频地址。');
  160. await upsertTranscriptJob(userId, job);
  161. return response.json({ code: 200, success: true, data: job });
  162. }
  163. const durationMs = media.durationMs || extractDurationMs(detail);
  164. if (!durationMs) {
  165. job.status = 'needs_media';
  166. job.warnings.push('未能确认音频时长,转写网关需要 durationMs。');
  167. await upsertTranscriptJob(userId, job);
  168. return response.json({ code: 200, success: true, data: job });
  169. }
  170. const uploaded = await uploadGatewayAudio(media.url, durationMs);
  171. job.orderId = uploaded.orderId;
  172. job.estimateTime = uploaded.estimateTime || 0;
  173. job.mediaUrl = media.url;
  174. job.durationMs = durationMs;
  175. job.warnings.push(`已提交转写任务:${uploaded.orderId}`);
  176. await upsertTranscriptJob(userId, job);
  177. response.json({ code: 200, success: true, data: job });
  178. } catch (error) {
  179. job.status = 'failed';
  180. job.errorMessage = error.message || '提交逐字稿任务失败';
  181. job.warnings.push(job.errorMessage);
  182. await upsertTranscriptJob(userId, job);
  183. response.json({ code: 200, success: true, data: job });
  184. }
  185. }
  186. async function getTranscript(request, response, userId) {
  187. const id = clean(pickParam(request, 'id'));
  188. if (!id) return response.json({ code: 400, success: false, error: '缺少 id' });
  189. const rows = await Psql.query(
  190. `SELECT * FROM "VideoflowTranscriptJob" WHERE "bizId"=$1 AND "userId"=$2 LIMIT 1`,
  191. [id, userId]
  192. );
  193. if (!rows.length) return response.json({ code: 404, success: false, error: '未找到逐字稿任务' });
  194. const job = rowToObj(rows[0]);
  195. if (!job.orderId || job.status === 'completed' || job.status === 'failed') {
  196. return response.json({ code: 200, success: true, data: job });
  197. }
  198. if (!VOC_TOKEN) {
  199. job.status = 'needs_provider_config';
  200. job.warnings = [...(job.warnings || []), '云函数未配置 VOC_TOKEN、TRANSCRIPTION_VOC_TOKEN 或 VOICE_TOKEN,无法查询转写网关。'];
  201. await upsertTranscriptJob(userId, job);
  202. return response.json({ code: 200, success: true, data: job });
  203. }
  204. try {
  205. const data = await queryGateway(job.orderId);
  206. const status = String(gatewayValue(data, 'status') || '').toLowerCase();
  207. const text = clean(gatewayValue(data, 'text'));
  208. const segments = normalizeGatewaySegments(gatewayValue(data, 'segments'));
  209. if (status === 'completed' || text) {
  210. job.status = 'completed';
  211. job.text = text || segments.map(s => s.text).join('\n');
  212. job.segments = segments;
  213. } else if (status === 'failed' || status === 'error') {
  214. job.status = 'failed';
  215. job.errorMessage = gatewayValue(data, 'error') || gatewayValue(data, 'message') || '转写失败';
  216. job.warnings = [...(job.warnings || []), job.errorMessage];
  217. } else {
  218. job.status = 'pending';
  219. job.warnings = [...new Set([...(job.warnings || []), '转写任务仍在处理中。'])];
  220. }
  221. job.updatedAt = new Date().toISOString();
  222. await upsertTranscriptJob(userId, job);
  223. response.json({ code: 200, success: true, data: job });
  224. } catch (error) {
  225. job.status = 'failed';
  226. job.errorMessage = error.message || '查询逐字稿任务失败';
  227. job.warnings = [...(job.warnings || []), job.errorMessage];
  228. job.updatedAt = new Date().toISOString();
  229. await upsertTranscriptJob(userId, job);
  230. response.json({ code: 200, success: true, data: job });
  231. }
  232. }
  233. async function upsertTranscriptJob(userId, job) {
  234. const existing = await Psql.query(
  235. `SELECT * FROM "VideoflowTranscriptJob" WHERE "bizId"=$1 AND "userId"=$2 LIMIT 1`,
  236. [job.id, userId]
  237. );
  238. if (existing.length) {
  239. await Psql.query(
  240. `UPDATE "VideoflowTranscriptJob" SET "data"=$1, "status"=$2, "updatedAt"=NOW() WHERE "bizId"=$3 AND "userId"=$4`,
  241. [JSON.stringify(job), job.status || '', job.id, userId]
  242. );
  243. return;
  244. }
  245. await Psql.query(
  246. `INSERT INTO "VideoflowTranscriptJob" ("objectId","bizId","userId","data","status")
  247. VALUES ($1,$2,$3,$4,$5)`,
  248. [generateId(), job.id, userId, JSON.stringify(job), job.status || '']
  249. );
  250. }
  251. function rowToObj(row) {
  252. const data = typeof row.data === 'string' ? JSON.parse(row.data) : (row.data || {});
  253. return { ...data, objectId: row.objectId, createdAt: row.createdAt, updatedAt: row.updatedAt };
  254. }
  255. async function fetchDouyinDetail(awemeId) {
  256. if (!DOUYIN_API_TOKEN) {
  257. throw new Error('抖音数据网关未配置有效 token:FMode voc-social 请配置 DOUYIN_API_TOKEN、VOC_TOKEN 或 VOC_SOCIAL_TOKEN;TikHub 直连才使用 TIKHUB_TOKEN。');
  258. }
  259. const url = new URL(`${DOUYIN_API_BASE_URL}${normalizeDouyinPath('/douyin/app/v3/fetch_one_video_v3')}`);
  260. url.searchParams.set('aweme_id', awemeId);
  261. let resp;
  262. try {
  263. resp = await fetch(url.toString(), {
  264. method: 'GET',
  265. headers: {
  266. 'Accept': 'application/json',
  267. 'Authorization': bearerAuth(DOUYIN_API_TOKEN),
  268. },
  269. });
  270. } catch (error) {
  271. throw new Error(`抖音数据网关网络请求失败:${formatFetchError(error)}。base=${maskBaseUrl(DOUYIN_API_BASE_URL)} route=/douyin/app/v3/fetch_one_video_v3`);
  272. }
  273. const data = await resp.json().catch(() => ({}));
  274. if (!resp.ok || data.success === false) throw new Error(readGatewayError(data) || `抖音详情获取失败 HTTP ${resp.status}`);
  275. return data;
  276. }
  277. async function diagnose(request, response) {
  278. const awemeId = String(pickParam(request, 'awemeId') || '7592116912205630761').trim();
  279. const result = {
  280. code: 200,
  281. success: true,
  282. data: {
  283. parseApiHost: PARSE_API_HOST,
  284. douyinBaseUrl: maskBaseUrl(DOUYIN_API_BASE_URL),
  285. isTikhubDirect: IS_TIKHUB_DIRECT,
  286. douyinTokenConfigured: !!DOUYIN_API_TOKEN,
  287. douyinTokenSource: douyinTokenSource(),
  288. transcriptTokenConfigured: !!VOC_TOKEN,
  289. transcriptTokenSource: transcriptTokenSource(),
  290. transcriptionGateway: maskBaseUrl(TRANSCRIPTION_GATEWAY),
  291. probe: null,
  292. }
  293. };
  294. if (String(pickParam(request, 'probe') || '') === '1') {
  295. try {
  296. const detail = await fetchDouyinDetail(awemeId);
  297. result.data.probe = {
  298. code: detail?.code,
  299. success: detail?.success !== false,
  300. hasData: !!detail?.data,
  301. keys: detail && typeof detail === 'object' ? Object.keys(detail).slice(0, 10) : [],
  302. };
  303. } catch (error) {
  304. result.data.probe = {
  305. code: 502,
  306. success: false,
  307. error: error && error.message ? error.message : String(error || 'probe failed'),
  308. };
  309. }
  310. }
  311. if (String(pickParam(request, 'probeDb') || '') === '1') {
  312. try {
  313. await ensureTables();
  314. result.data.database = { success: true };
  315. } catch (error) {
  316. result.data.database = {
  317. success: false,
  318. error: error && error.message ? error.message : String(error || 'database probe failed'),
  319. };
  320. }
  321. }
  322. const sessionToken = clean(pickParam(request, 'sessionToken'));
  323. if (sessionToken) {
  324. try {
  325. const user = await verifyParseSessionFromDb(sessionToken) || await verifyParseSession(sessionToken);
  326. result.data.session = {
  327. success: !!user?.objectId,
  328. userId: user?.objectId || '',
  329. source: user?.source || '',
  330. };
  331. } catch (error) {
  332. result.data.session = {
  333. success: false,
  334. error: error && error.message ? error.message : String(error || 'session probe failed'),
  335. };
  336. }
  337. }
  338. return response.json(result);
  339. }
  340. function normalizeDouyinPath(path) {
  341. if (IS_TIKHUB_DIRECT && !path.startsWith('/api/v1/')) {
  342. return `/api/v1${path}`;
  343. }
  344. return path;
  345. }
  346. function selectAudioCandidate(detail) {
  347. const candidates = [];
  348. collectMediaUrls(detail, candidates, []);
  349. const audio = candidates.find(item => item.kind === 'audio') || candidates.find(item => /audio|mp4a|music|sound/i.test(item.url));
  350. return audio || { url: '', durationMs: extractDurationMs(detail) };
  351. }
  352. function collectMediaUrls(value, out, path) {
  353. if (!value) return;
  354. if (Array.isArray(value)) {
  355. value.forEach((item, index) => collectMediaUrls(item, out, [...path, String(index)]));
  356. return;
  357. }
  358. if (typeof value !== 'object') return;
  359. for (const [key, raw] of Object.entries(value)) {
  360. const keyPath = [...path, key].join('.');
  361. if (typeof raw === 'string' && /^https?:\/\//i.test(raw)) {
  362. const joined = keyPath.toLowerCase();
  363. const kind = /audio|music|sound|mp4a/.test(joined) ? 'audio' : /video|play|download|media/.test(joined) ? 'video' : '';
  364. if (kind) out.push({ url: raw, kind, keyPath, durationMs: extractDurationMs(value) });
  365. } else if (raw && typeof raw === 'object') {
  366. collectMediaUrls(raw, out, [...path, key]);
  367. }
  368. }
  369. }
  370. function extractDurationMs(value) {
  371. const found = findFirstNumber(value, ['duration', 'duration_ms', 'durationMs', 'video_duration']);
  372. if (!found) return 0;
  373. return found > 1000 ? Math.round(found) : Math.round(found * 1000);
  374. }
  375. function findFirstNumber(value, keys, depth = 0) {
  376. if (!value || depth > 5) return 0;
  377. if (Array.isArray(value)) {
  378. for (const item of value) {
  379. const hit = findFirstNumber(item, keys, depth + 1);
  380. if (hit) return hit;
  381. }
  382. return 0;
  383. }
  384. if (typeof value !== 'object') return 0;
  385. for (const [key, raw] of Object.entries(value)) {
  386. if (keys.includes(key)) {
  387. const number = Number(raw);
  388. if (Number.isFinite(number) && number > 0) return number;
  389. }
  390. const hit = findFirstNumber(raw, keys, depth + 1);
  391. if (hit) return hit;
  392. }
  393. return 0;
  394. }
  395. async function uploadGatewayAudio(mediaUrl, durationMs) {
  396. let mediaResp;
  397. try {
  398. mediaResp = await fetch(mediaUrl);
  399. } catch (error) {
  400. throw new Error(`音频下载网络失败:${formatFetchError(error)}`);
  401. }
  402. if (!mediaResp.ok) throw new Error(`音频下载失败 HTTP ${mediaResp.status}`);
  403. const blob = await mediaResp.blob();
  404. const form = new FormData();
  405. form.append('audio', blob, `douyin-audio-${Date.now()}.m4a`);
  406. form.append('durationMs', String(durationMs));
  407. form.append('roleType', '1');
  408. form.append('roleNum', '0');
  409. let resp;
  410. try {
  411. resp = await fetch(`${TRANSCRIPTION_GATEWAY}/upload`, {
  412. method: 'POST',
  413. headers: {
  414. Authorization: bearerAuth(VOC_TOKEN),
  415. Accept: 'application/json',
  416. },
  417. body: form,
  418. });
  419. } catch (error) {
  420. throw new Error(`转写网关上传网络失败:${formatFetchError(error)}。gateway=${maskBaseUrl(TRANSCRIPTION_GATEWAY)}`);
  421. }
  422. const data = await resp.json().catch(() => ({}));
  423. if (!resp.ok || data.success === false) throw new Error(data.error?.message || data.error || data.message || `转写上传失败 HTTP ${resp.status}`);
  424. const orderId = data.orderId || data.content?.orderId || data.data?.orderId || data.result?.orderId;
  425. if (!orderId) throw new Error('转写网关未返回 orderId');
  426. return { orderId, estimateTime: Number(data.estimateTime || data.content?.estimateTime || data.data?.estimateTime || 0) };
  427. }
  428. async function queryGateway(orderId) {
  429. let resp;
  430. try {
  431. resp = await fetch(`${TRANSCRIPTION_GATEWAY}/result`, {
  432. method: 'POST',
  433. headers: {
  434. Authorization: bearerAuth(VOC_TOKEN),
  435. Accept: 'application/json',
  436. 'Content-Type': 'application/json',
  437. },
  438. body: JSON.stringify({ orderId }),
  439. });
  440. } catch (error) {
  441. throw new Error(`转写网关查询网络失败:${formatFetchError(error)}。gateway=${maskBaseUrl(TRANSCRIPTION_GATEWAY)}`);
  442. }
  443. const data = await resp.json().catch(() => ({}));
  444. if (!resp.ok) throw new Error(data.error?.message || data.error || data.message || `转写查询失败 HTTP ${resp.status}`);
  445. return data;
  446. }
  447. function gatewayValue(data, key) {
  448. return data?.[key] ?? data?.data?.[key] ?? data?.result?.[key] ?? data?.content?.[key];
  449. }
  450. function readGatewayError(data) {
  451. const detail = data?.detail;
  452. if (detail === 'Not Found') return '抖音数据接口地址未找到,请检查 DOUYIN_API_BASE_URL 是否配置为 https://server.fmode.cn/api/voc-social';
  453. const message = data?.mess || data?.message || data?.msg || data?.error?.message || data?.error || detail || '';
  454. if (/company或用户信息不存在/.test(String(message))) {
  455. return '当前抖音数据网关 token 未绑定有效用户或公司,请在云函数配置 DOUYIN_API_TOKEN/VOC_TOKEN/VOC_SOCIAL_TOKEN,不能使用 TikHub token。';
  456. }
  457. return message;
  458. }
  459. function normalizeGatewaySegments(segments) {
  460. const arr = Array.isArray(segments) ? segments : [];
  461. return arr.map(segment => ({
  462. start: normalizeTime(segment.start ?? segment.begin ?? segment.bg),
  463. end: normalizeTime(segment.end ?? segment.ed),
  464. text: clean(segment.text || segment.onebest || segment.content),
  465. })).filter(segment => segment.text);
  466. }
  467. function normalizeTime(value) {
  468. const number = Number(value);
  469. if (!Number.isFinite(number)) return null;
  470. return number > 1000 ? number / 1000 : number;
  471. }
  472. function bearerAuth(token) {
  473. const value = clean(token);
  474. return /^Bearer\s+/i.test(value) ? value : `Bearer ${value}`;
  475. }
  476. function douyinTokenSource() {
  477. if (readEnv('DOUYIN_API_TOKEN')) return 'DOUYIN_API_TOKEN';
  478. if (readEnv('VOC_TOKEN')) return 'VOC_TOKEN';
  479. if (readEnv('TRANSCRIPTION_VOC_TOKEN')) return 'TRANSCRIPTION_VOC_TOKEN';
  480. if (readEnv('VOICE_TOKEN')) return 'VOICE_TOKEN';
  481. if (readEnv('OPENCLAW_VOC_TOKEN')) return 'OPENCLAW_VOC_TOKEN';
  482. if (readEnv('VOC_SOCIAL_TOKEN')) return 'VOC_SOCIAL_TOKEN';
  483. if (IS_TIKHUB_DIRECT && readEnv('TIKHUB_TOKEN')) return 'TIKHUB_TOKEN';
  484. if (!IS_TIKHUB_DIRECT && LOCAL_VOC_TOKEN_FALLBACK) return 'LOCAL_VOC_TOKEN_FALLBACK';
  485. return IS_TIKHUB_DIRECT ? 'TIKHUB_TOKEN_FALLBACK' : '';
  486. }
  487. function transcriptTokenSource() {
  488. if (readEnv('VOC_TOKEN')) return 'VOC_TOKEN';
  489. if (readEnv('TRANSCRIPTION_VOC_TOKEN')) return 'TRANSCRIPTION_VOC_TOKEN';
  490. if (readEnv('VOICE_TOKEN')) return 'VOICE_TOKEN';
  491. if (readEnv('OPENCLAW_VOC_TOKEN')) return 'OPENCLAW_VOC_TOKEN';
  492. if (readEnv('VOC_SOCIAL_TOKEN')) return 'VOC_SOCIAL_TOKEN';
  493. if (LOCAL_VOC_TOKEN_FALLBACK) return 'LOCAL_VOC_TOKEN_FALLBACK';
  494. return '';
  495. }
  496. function formatFetchError(error) {
  497. const message = error && error.message ? error.message : String(error || 'fetch failed');
  498. const cause = error && error.cause ? `;cause=${error.cause.code || error.cause.message || error.cause}` : '';
  499. return `${message}${cause}`;
  500. }
  501. function maskBaseUrl(value) {
  502. return String(value || '').replace(/(token=)[^&]+/ig, '$1***');
  503. }
  504. function pickParam(request, ...names) {
  505. const sources = [request.params, request.body, request];
  506. for (const src of sources) {
  507. if (!src || typeof src !== 'object') continue;
  508. for (const name of names) {
  509. const value = src[name];
  510. if (value !== undefined && value !== null && value !== '') return value;
  511. }
  512. }
  513. return null;
  514. }
  515. async function requireSession(request) {
  516. const token = clean(pickParam(request, 'sessionToken'));
  517. if (!token) throw new Error('请先登录');
  518. const rows = await Psql.query(
  519. `SELECT * FROM "AppSession" WHERE "token"=$1 AND "expiresAt" > NOW() LIMIT 1`,
  520. [token]
  521. );
  522. if (rows.length) return rows[0];
  523. const parseUser = await verifyParseSessionFromDb(token) || await verifyParseSession(token);
  524. if (parseUser?.objectId) return { token, userId: parseUser.objectId, source: 'parse' };
  525. throw new Error('登录已过期,请重新登录');
  526. }
  527. async function verifyParseSessionFromDb(sessionToken) {
  528. try {
  529. const rows = await Psql.query(
  530. `SELECT s."objectId" AS "sessionObjectId",
  531. s."sessionToken" AS "sessionToken",
  532. s."expiresAt" AS "expiresAt",
  533. s."_p_user" AS "userPointer",
  534. u."objectId" AS "objectId",
  535. u."username" AS "username"
  536. FROM "_Session" s
  537. LEFT JOIN "_User" u ON s."_p_user" = CONCAT('_User$', u."objectId")
  538. WHERE s."sessionToken"=$1
  539. AND (s."expiresAt" IS NULL OR s."expiresAt" > NOW())
  540. LIMIT 1`,
  541. [sessionToken]
  542. );
  543. if (!rows.length) return null;
  544. const row = rows[0];
  545. const pointerUserId = String(row.userPointer || '').replace(/^_User\$/, '');
  546. const objectId = row.objectId || pointerUserId;
  547. return objectId ? { objectId, username: row.username || '', source: 'parse-db' } : null;
  548. } catch (error) {
  549. console.warn('verifyParseSessionFromDb skipped:', error.message);
  550. return null;
  551. }
  552. }
  553. async function verifyParseSession(sessionToken) {
  554. if (typeof fetch !== 'function') return null;
  555. try {
  556. const resp = await fetch(`${PARSE_API_HOST}/parse/users/me?include=company`, {
  557. method: 'GET',
  558. headers: {
  559. 'X-Parse-Application-Id': PARSE_APP_ID,
  560. 'X-Parse-Session-Token': sessionToken,
  561. },
  562. });
  563. const data = await resp.json().catch(() => ({}));
  564. if (!resp.ok) throw new Error(data.error || data.message || `HTTP ${resp.status}`);
  565. return data.objectId ? { ...data, source: 'parse-rest' } : null;
  566. } catch (error) {
  567. throw new Error(`Parse 会话网络校验失败:${formatFetchError(error)}`);
  568. }
  569. }
  570. function clean(value) {
  571. return String(value || '').trim();
  572. }
  573. function generateId() {
  574. const chars = 'ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789';
  575. let s = '';
  576. for (let i = 0; i < 10; i++) s += chars.charAt(Math.floor(Math.random() * chars.length));
  577. return s;
  578. }
  579. function readEnv(name) {
  580. if (typeof process !== 'undefined' && process.env && process.env[name]) {
  581. return process.env[name];
  582. }
  583. return '';
  584. }