File size: 17,481 Bytes
90647f1
 
 
ee71e80
9391d4e
90647f1
 
7a3a29b
90647f1
 
 
e450726
 
9391d4e
 
7a3a29b
0cf9806
46665f5
90647f1
bd24c9a
 
1cadb77
 
 
226deb2
 
 
 
1cadb77
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
90647f1
 
 
 
 
 
 
96b2e68
fe06e2a
90647f1
79a3810
 
517e1b0
 
90647f1
 
 
 
 
 
ee71e80
 
90647f1
ee71e80
 
 
 
 
 
e450726
ee71e80
 
 
 
 
 
6fc89d1
ee71e80
 
 
 
 
 
 
 
 
e450726
 
 
 
6fc89d1
e450726
 
 
ee71e80
90647f1
 
4f45ad2
 
 
 
 
 
 
 
 
 
 
 
 
 
90647f1
 
 
 
 
 
 
 
 
 
4f45ad2
 
 
 
 
 
 
 
90647f1
 
 
 
 
 
 
 
 
 
 
 
 
 
 
96b2e68
90647f1
ddd8a50
 
 
ee71e80
 
90647f1
ee71e80
 
 
 
90647f1
ddd8a50
90647f1
ddd8a50
 
 
 
 
 
90647f1
 
 
 
 
 
 
 
96b2e68
 
 
 
 
 
 
 
7a3a29b
 
 
 
0cf9806
 
 
90647f1
96b2e68
0cf9806
5221da3
7a3a29b
517e1b0
 
90647f1
 
 
 
 
9391d4e
7ad8ce5
 
 
10d10f3
7ad8ce5
 
4f45ad2
9391d4e
f267d9e
 
c9393b9
f267d9e
9391d4e
 
5221da3
9391d4e
7a3a29b
 
0cf9806
 
7a3a29b
 
 
 
 
 
 
a19488c
 
 
 
 
 
 
 
7a3a29b
 
 
 
 
517e1b0
7a3a29b
 
 
 
 
 
 
 
 
 
 
 
a19488c
 
 
 
 
 
 
 
7a3a29b
 
 
9391d4e
188b01b
 
 
f267d9e
9391d4e
 
 
58ac1ce
517e1b0
 
 
 
90647f1
 
 
 
 
9391d4e
c9393b9
9391d4e
 
90647f1
58ac1ce
517e1b0
58ac1ce
 
 
 
90647f1
 
9391d4e
90647f1
1cadb77
 
 
 
 
 
 
 
 
 
 
8fc98b6
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
ee71e80
90647f1
ee71e80
90647f1
 
0e98142
46665f5
0e98142
 
 
 
46665f5
0e98142
 
 
 
46665f5
0e98142
 
6fc89d1
46665f5
 
 
6fc89d1
 
0e98142
 
6fc89d1
58ac1ce
 
 
 
 
 
6fc89d1
8fc98b6
1cadb77
 
 
 
 
 
 
 
 
 
 
8fc98b6
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
226deb2
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
6fc89d1
 
 
ee71e80
 
 
 
 
90647f1
 
 
 
 
e450726
 
 
90647f1
e450726
 
90647f1
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
/**
 * 任务处理器
 * 使用 TaskExecutor 执行慢任务
 * ★★★ 并发模式:每个任务独立 executor 实例 ★★★
 * ★★★ 修复:上传到腾讯云 COS,设置真实 downloadUrl ★★★
 */
import { getRedis } from '../lib/redis';
import { TaskParams, TaskResult, ProgressReport, DesignTaskType } from './task-executor';
import { ModelDispatcher } from '../lib/model-dispatcher';
import { NvidiaKeyRotator } from '../lib/nvidia-rotator';
import { LatencyTracker } from '../lib/latency-tracker';
import { isTaskAborted, clearAbortFlag, setGlobalExecutor, getGlobalExecutor } from './index';
import { TaskExecutor } from './task-executor';
import { uploadDocument, completeTask } from '../lib/storage-provider';
import { markdownToDocx } from '../lib/docx-engine';
import { generateOpenDesignAssetAsync } from '../lib/open-design-engine';
import { selectStyle } from '../lib/style-profiles';
import { sanitizeErrorMessage } from '../lib/error-sanitizer';

console.log('[STARTUP] processor loaded, BUILD: 20260522-v8-MIMO-FIX');

// ─── Server-side analytics ───
const ANALYTICS_ENDPOINT = 'https://www.5e1.com/api/analytics/event';

// ─── 退款接口(game.5e1.com)───
const GAME_API_BASE = 'https://game.5e1.com';
const INTERNAL_SECRET = process.env.INTERNAL_API_SECRET || '';

function recordEvent(options: {
  eventType: string;
  userId?: string | null;
  page?: string;
  meta?: Record<string, unknown>;
}): void {
  const { eventType, userId, page, meta = {} } = options;
  fetch(ANALYTICS_ENDPOINT, {
    method: 'POST',
    headers: { 'Content-Type': 'application/json' },
    body: JSON.stringify({
      eventType,
      product: 'game',
      page: page || null,
      userId: userId || null,
      meta: JSON.stringify(meta),
    }),
    signal: AbortSignal.timeout(5000),
  }).catch((err) => {
    console.warn('[Analytics] Failed to record event:', eventType, err);
  });
}

interface SlowTask {
  taskId: string;
  userId: string;
  tier: string;
  taskType: string;
  status: string;
  chatContext?: Array<{ role: string; content: string }>;
  prompt?: string;  // ★ 新增:直接 prompt 字段兼容
  title?: string;   // ★ 任务标题(由 submit API 生成)
  createdAt: number;
  callbackUrl?: string;  // ★ 回调 URL
  metadata?: Record<string, unknown>;  // ★ 附加元数据
  draftContent?: string;   // ★ 草案内容(HTML artifact)
  skipPhase1?: boolean;    // ★ 是否跳过 Phase 1
}

const QUEUE_KEYS = {
  results: 'slow_task:results',
};

// ★★★ 并发模式:每个任务独立 executor,注册到 Map 供看门狗使用 ★★★
const taskExecutors = new Map<string, TaskExecutor>();

export function abortTask(taskId: string): boolean {
  const executor = taskExecutors.get(taskId);
  if (executor) {
    executor.abort();
    console.log(`[Processor] 🛑 Aborted executor for task: ${taskId}`);
    return true;
  }
  return false;
}

function registerExecutor(taskId: string, executor: TaskExecutor): void {
  taskExecutors.set(taskId, executor);
}

function unregisterExecutor(taskId: string): void {
  taskExecutors.delete(taskId);
}

/**
 * 为每个任务创建独立的 executor 实例
 * 共享同一个 NvidiaKeyRotator(轮询分发 key)和 Redis 连接
 */
async function createTaskExecutor(): Promise<TaskExecutor> {
  const redis = getRedis();
  if (!redis) {
    throw new Error('Redis not available');
  }

  const keyRotator = new NvidiaKeyRotator();
  const latencyTracker = new LatencyTracker(redis);
  const dispatcher = new ModelDispatcher(keyRotator, latencyTracker, redis);
  return new TaskExecutor(dispatcher, redis);
}

// ─── 更新任务状态到 Redis ───

// ★★★ 安全 JSON 解析(处理 Upstash SDK 自动解析)★★★
function safeParseTaskState(raw: unknown): Record<string, any> | null {
  if (!raw) return null;
  if (typeof raw === 'object') return raw as Record<string, any>;
  try {
    return JSON.parse(raw as string);
  } catch {
    console.warn('[Processor] JSON parse failed for task state');
    return null;
  }
}

async function updateTaskStatus(
  taskId: string,
  status: string,
  extra: Record<string, any> = {}
): Promise<void> {
  const redis = getRedis();
  if (!redis) return;

  const existing = await redis.hget(QUEUE_KEYS.results, taskId);

  // ★★★ 如果没有现有数据,创建初始状态 ★★★
  const state = safeParseTaskState(existing) || {
    taskId,
    status: 'queued',
    progress: 0,
    createdAt: Date.now(),
  };

  const updated = {
    ...state,
    status,
    updatedAt: Date.now(),
    ...extra,
  };

  await redis.hset(QUEUE_KEYS.results, {
    [taskId]: JSON.stringify(updated),
  });
}

// 任务处理入口
export async function processTask(task: SlowTask): Promise<void> {
  console.log('[Processor] Processing task:', task.taskId, 'type:', task.taskType);
  console.log('[Processor] Task prompt:', task.prompt?.slice(0, 50) || 'from chatContext');

  // ★★★ 清除中止信号 ★★★
  clearAbortFlag();

  let executor: TaskExecutor | null = null;

  try {
    // ★★★ 并发模式:为每个任务创建独立 executor ★★★
    executor = await createTaskExecutor();
    registerExecutor(task.taskId, executor);
    console.log('[Processor] Executor created for task:', task.taskId);

    // ★★★ 注册进度回调 - 同时检查中止信号 ★★★
    executor.onProgress(async (report: ProgressReport) => {
      // ★★★ 如果任务已被看门狗熔断,立即退出 ★★★
      if (isTaskAborted()) {
        console.warn('[Processor] 任务已被看门狗熔断,停止执行');
        throw new Error('Task aborted by watchdog');
      }

      console.log('[Processor] Progress:', report.progress, '%', report.phase);
      await updateTaskStatus(task.taskId, report.status, {
        progress: report.progress,
        phase: report.phase,
        detail: report.detail,
      });
    });

    // 构建任务参数(兼容 chatContext 和 prompt 字段)
    const promptContent = task.chatContext?.map(m => m.content).join('\n') || task.prompt || '';
    console.log('[Processor] Prompt content length:', promptContent.length);

    if (!promptContent) {
      throw new Error('No prompt content provided');
    }

    // ★★★ 设计任务类型检测 ★★★
    const DESIGN_TASK_TYPES: DesignTaskType[] = ['html_prototype', 'pitch_deck', 'ui_design_system'];
    const isDesign = DESIGN_TASK_TYPES.includes(task.taskType as DesignTaskType);

    // ★★★ 设计任务:随机选择视觉风格 ★★★
    const designStyle = isDesign ? selectStyle() : undefined;

    const params: TaskParams = {
      prompt: promptContent,
      style: designStyle,
      tier: (task.tier as 'free' | 'starter' | 'pro' | 'enterprise') || 'free',
      taskType: task.taskType as TaskParams['taskType'],
      draftContent: task.draftContent,   // ★ 透传草案内容
      skipPhase1: task.skipPhase1,       // ★ 透传跳过标记
    };

    // 执行任务
    const result: TaskResult = await executor.execute(task.taskId, params);

    // ★★★ 保存生成内容到 Redis(备份)★★★
    const redis = getRedis();
    if (redis) {
      const contentKey = `slow_task:content:${task.taskId}`;
      await redis.set(contentKey, result.content, 'EX', 86400 * 7); // 7天过期
      console.log('[Processor] Content saved to Redis:', contentKey, 'length:', result.content.length);
    }

    // ★★★ 上传到腾讯云 COS ★★★
    // ★★★ 修复:初始化变量,确保始终有值 ★★★
    let downloadUrl: string = `https://game.5e1.com/api/slow-task/content-raw/${task.taskId}`;
    let downloadSource: 'cos' | 'local' = 'local';
    let expiresAt: number = Date.now() + 86400 * 7;

    try {
      const tier = (task.tier as 'free' | 'starter' | 'pro' | 'enterprise') || 'free';

      // ★★★ 设计任务:使用 OpenDesignEngine 生成 ZIP/PPTX ★★★
      if (isDesign && result.outputType) {
        console.log('[Processor] 设计任务,使用 OpenDesignEngine:', task.taskType, 'style:', designStyle?.name);
        const asset = await generateOpenDesignAssetAsync(task.taskType!, result.content, designStyle);
        console.log('[Processor] 设计资产生成完成, size:', asset.buffer.length, 'bytes, type:', asset.contentType);

        // 上传到 COS
        const uploadResult = await uploadDocument(task.taskId, asset.buffer, tier, undefined, {
          fileExtension: asset.fileExtension,
          contentType: asset.contentType,
        });
        if (uploadResult.publicUrl.includes('/mock/')) {
          downloadSource = 'local';
          console.warn('[Processor] ⚠ COS returned mock URL, using local fallback');
        } else {
          downloadUrl = uploadResult.publicUrl;
          downloadSource = 'cos';
          expiresAt = uploadResult.expiresAt;
        }

        console.log('[Processor] 设计资产已上传:', downloadUrl, 'expiresAt:', expiresAt);
      } else {
        // ★★★ 文档任务:生成 DOCX 文件 ★★★
        const isFull = tier !== 'free';
        const docTitle = result.title || task.title || (isFull ? '大师全案' : '灵感草案');
        console.log('[Processor] 开始生成 DOCX...', docTitle, 'tier:', tier);
        const docxBuffer = await markdownToDocx(result.content, {
          title: docTitle,
          version: 'v1.0',
          confidentiality: isFull ? '机密文档' : '内部文档',
          tier,
        });
        console.log('[Processor] DOCX generated, size:', docxBuffer.length, 'bytes');

        // 上传到 COS
        console.log('[Processor] 开始上传到 COS, tier:', tier);
        const uploadResult = await uploadDocument(task.taskId, docxBuffer, tier);
        if (uploadResult.publicUrl.includes('/mock/')) {
          downloadSource = 'local';
          console.warn('[Processor] ⚠ COS returned mock URL, using local fallback');
        } else {
          downloadUrl = uploadResult.publicUrl;
          downloadSource = 'cos';
          expiresAt = uploadResult.expiresAt;
        }

        console.log('[Processor] Uploaded to COS:', downloadUrl, 'expiresAt:', expiresAt);
      }
    } catch (uploadErr: any) {
      console.warn('[Processor] Upload/Design generation failed:', uploadErr.message);
      console.warn('[Processor] Error type:', uploadErr.constructor?.name);
      console.warn('[Processor] Error stack:', uploadErr.stack?.slice(0, 300));
      // 失败时使用备用 URL(已初始化,无需重新赋值)
    }

    // ★★★ 更新最终结果 ★★★
    // ★ 保留初始字段(title/taskType/tier/userId),防止 updateTaskStatus 丢失
    // ★★★ 优先使用 LLM 生成的标题(result.title),其次用 submit 传入的标题 ★★★
    const finalTitle = result.title || task.title || '';
    console.log('[Processor] 最终标题:', finalTitle, '(LLM:', result.title, '| submit:', task.title, ')');

    await updateTaskStatus(task.taskId, 'completed', {
      progress: 100,
      completedAt: Date.now(),
      sectionsGenerated: result.sectionsGenerated,
      duration: result.duration,
      downloadUrl,  // ★ 使用真实 COS URL 或备用 URL
      downloadSource,  // ★ OPT-05: 标记下载来源
      expiresAt,
      qualityGrade: result.phases?.review?.model ? 'A' : 'B',
      phases: result.phases,
      // ★ 补回初始字段
      title: finalTitle,
      taskType: task.taskType,
      tier: task.tier,
      userId: task.userId,
      createdAt: task.createdAt,
    });

    console.log('[Processor] Task completed:', task.taskId, 'duration:', result.duration, 'ms', 'downloadUrl:', downloadUrl);

    // ★ 追踪:任务完成事件
    recordEvent({
      eventType: 'slow_task_complete',
      userId: task.userId,
      page: '/api/slow-task/callback',
      meta: {
        taskType: task.taskType,
        totalTimeMs: result.duration,
      },
    });

    // ★★★ 回调通知:任务完成后 POST 结果到 callbackUrl ★★★
    try {
      const existingState = await redis?.hget(QUEUE_KEYS.results, task.taskId);
      const state = existingState ? (typeof existingState === 'string' ? JSON.parse(existingState) : existingState) : {};
      const callbackUrl = state.callbackUrl || task.callbackUrl;
      if (callbackUrl) {
        console.log('[Processor] Sending callback to:', callbackUrl);
        const callbackBody = {
          taskId: task.taskId,
          taskType: task.taskType,
          status: 'completed',
          content: result.content,
          downloadUrl,
          qualityGrade: result.phases?.review?.model ? 'A' : 'B',
          duration: result.duration,
          sectionsGenerated: result.sectionsGenerated,
          metadata: state.metadata || {},
        };
        fetch(callbackUrl, {
          method: 'POST',
          headers: { 'Content-Type': 'application/json' },
          body: JSON.stringify(callbackBody),
          signal: AbortSignal.timeout(15000),
        }).then(r => {
          console.log('[Processor] Callback response:', r.status);
        }).catch(cbErr => {
          console.warn('[Processor] Callback failed:', cbErr.message);
        });
      }
    } catch (cbErr) {
      console.warn('[Processor] Callback error:', cbErr);
    }

    // 重置执行器状态并注销
    executor.reset();
    unregisterExecutor(task.taskId);

  } catch (err: any) {
    // ★★★ 增强错误捕获 - 确保错误信息正确序列化 ★★★
    const rawErrorMessage = typeof err === 'string' ? err :
                         (err?.message ? err.message :
                         (typeof err === 'object' ? JSON.stringify(err) : 'Unknown error'));
    const errorStack = err?.stack?.split('\n').slice(0, 5).join('\n') || 'No stack trace';

    console.error('[Processor] Task processing failed:', rawErrorMessage);
    console.error('[Processor] Error type:', typeof err);
    console.error('[Processor] Error stack:', errorStack);

    // 如果是JSON解析错误,尝试记录原始响应
    if (rawErrorMessage.includes('JSON') && err?.rawResponse) {
      console.error('[Processor] Raw response that caused JSON error:', err.rawResponse);
    }

    // ★★★ 脱敏清洗:确保不暴露模型名、API Key 等技术细节 ★★★
    const errorMessage = sanitizeErrorMessage(rawErrorMessage);

    try {
      await updateTaskStatus(task.taskId, 'failed', {
        errorMessage: errorMessage,
        errorStack: errorStack,
        completedAt: Date.now(),
        // ★ 补回初始字段
        title: task.title,
        taskType: task.taskType,
        tier: task.tier,
        userId: task.userId,
        createdAt: task.createdAt,
      });

      // ★ 追踪:任务失败事件
      recordEvent({
        eventType: 'slow_task_fail',
        userId: task.userId,
        page: '/api/slow-task/callback',
        meta: {
          taskType: task.taskType,
          error: errorMessage,
        },
      });

      // ★★★ 回调通知:任务失败时也通知 callbackUrl ★★★
      try {
        const redis = getRedis();
        const existingState = redis ? await redis.hget(QUEUE_KEYS.results, task.taskId) : null;
        const state = existingState ? (typeof existingState === 'string' ? JSON.parse(existingState) : existingState) : {};
        const callbackUrl = state.callbackUrl || task.callbackUrl;
        if (callbackUrl) {
          fetch(callbackUrl, {
            method: 'POST',
            headers: { 'Content-Type': 'application/json' },
            body: JSON.stringify({
              taskId: task.taskId,
              taskType: task.taskType,
              status: 'failed',
              errorMessage,
            }),
            signal: AbortSignal.timeout(10000),
          }).catch(() => {});
        }
      } catch {}

      // ★★★ 退款:任务失败时退还配额和 Credits ★★★
      try {
        const isLite = task.metadata?.planType === 'lite';
        const refundEndpoint = `${GAME_API_BASE}/api/slow-task/refund`;
        fetch(refundEndpoint, {
          method: 'POST',
          headers: {
            'Authorization': `Bearer ${INTERNAL_SECRET}`,
            'Content-Type': 'application/json',
          },
          body: JSON.stringify({
            taskId: task.taskId,
            userId: task.userId,
            taskType: task.taskType,
            isLite,
            preDeductedCredits: task.metadata?.preDeductedCredits || 0,
          }),
          signal: AbortSignal.timeout(10000),
        }).then(r => {
          console.log('[Processor] Refund response:', r.status);
        }).catch(refundErr => {
          console.warn('[Processor] Refund request failed:', refundErr.message);
        });
      } catch {}
    } catch (updateErr) {
      console.error('[Processor] Failed to update task status:', updateErr);
    }

    // ★★★ 注销 executor(无论成功失败)★★★
    if (executor) {
      unregisterExecutor(task.taskId);
    }
  }
}

// 获取调度器状态摘要
export async function getDispatcherSummary(): Promise<any> {
  const executor = getGlobalExecutor();
  if (!executor) {
    return { error: 'Executor not initialized' };
  }
  // 从全局执行器获取状态(如果需要)
  return { status: 'active', message: 'Executor available' };
}