huz1xuan

完善高光时刻入口

@@ -44,5 +44,15 @@ @@ -44,5 +44,15 @@
44 "outputNamespace": "", 44 "outputNamespace": "",
45 "outputBaseUrl": "https://xdymp4.xuedianyun.com" 45 "outputBaseUrl": "https://xdymp4.xuedianyun.com"
46 }, 46 },
  47 + "RECORDINGV2CONFIG": {
  48 + "taskListUrl": "https://saas.xuedianyun.com/3m/api/recording/getRecordingTasksPrivate.do",
  49 + "pageSize": 100,
  50 + "maxPages": 1000,
  51 + "maxTasksPerRun": 4,
  52 + "maxFullConcurrent": 2,
  53 + "apiTimeoutMs": 10000,
  54 + "apiRetryCount": 2,
  55 + "apiRetryBaseDelayMs": 500
  56 + },
47 "classLastNumber":["0","1","2","3","4","5","6","7","8","9"] 57 "classLastNumber":["0","1","2","3","4","5","6","7","8","9"]
48 } 58 }
@@ -4,6 +4,7 @@ @@ -4,6 +4,7 @@
4 4
5 - 原 `/recordingTask`、`/fileExists` 保持不变。 5 - 原 `/recordingTask`、`/fileExists` 保持不变。
6 - `/recordingTaskV2` 兼容原整课请求并支持 `onlyHighlight=1`。 6 - `/recordingTaskV2` 兼容原整课请求并支持 `onlyHighlight=1`。
  7 +- `/recordingTaskV2/scheduled` 供 cron 触发,由 WebScreen 自行拉取 SaaS 任务。
7 - `/fileExistsV2` 是新功能的客户查询接口。 8 - `/fileExistsV2` 是新功能的客户查询接口。
8 - `/highlight/*` 仅允许运维内网访问。 9 - `/highlight/*` 仅允许运维内网访问。
9 10
@@ -23,39 +24,37 @@ @@ -23,39 +24,37 @@
23 24
24 还需配置 OSS AccessKey、`web_capture_c`、显示服务及现有 OSS 搬运任务。 25 还需配置 OSS AccessKey、`web_capture_c`、显示服务及现有 OSS 搬运任务。
25 26
26 -不录制站点已由 xdySDK 上游流程筛除,WebScreen V2 不再维护第二份站点策略。  
27 -  
28 -## 外部 cron 27 +录制任务拉取配置:
29 28
30 -正式任务链路由现有外部 cron 分发,不使用 WebScreen 本机的  
31 -`/highlight/recording/scheduled` 作为正式入口。 29 +```json
  30 +"RECORDINGV2CONFIG": {
  31 + "taskListUrl": "https://saas.xuedianyun.com/3m/api/recording/getRecordingTasksPrivate.do",
  32 + "pageSize": 100,
  33 + "maxPages": 1000,
  34 + "maxTasksPerRun": 4,
  35 + "maxFullConcurrent": 2,
  36 + "apiTimeoutMs": 10000,
  37 + "apiRetryCount": 2,
  38 + "apiRetryBaseDelayMs": 500
  39 +}
  40 +```
32 41
33 -cron 继续调用 SaaS 的 `getRecordingTasksPrivate.do`,但需要将目标地址从: 42 +不录制站点已由 xdySDK 上游流程筛除,WebScreen V2 不再维护第二份站点策略。
34 43
35 -```text  
36 -POST /recordingTask  
37 -``` 44 +## cron
38 45
39 -改为: 46 +cron 不查询 SaaS,也不转换 `taskList`,只触发 WebScreen:
40 47
41 -```text  
42 -POST /recordingTaskV2 48 +```cron
  49 +31 7-20 * * * curl -fsS -X POST http://127.0.0.1:3001/recordingTaskV2/scheduled >> /var/log/webscreen_recording_v2.log 2>&1
43 ``` 50 ```
44 51
45 -原来的字段映射可以继续使用,只需增加 `onlyHighlight`: 52 +此入口只接受本机回环地址请求。
46 53
47 -```json  
48 -{  
49 - "classId": "task.meetingNumber",  
50 - "siteId": "task.siteId",  
51 - "yymmdd": "原有课堂日期",  
52 - "onlyHighlight": "task.onlyHighlight"  
53 -}  
54 -``` 54 +WebScreen 随后按照 `RECORDINGV2CONFIG` 分页调用
  55 +`getRecordingTasksPrivate.do`,并将任务交给 V2 录制服务。
55 56
56 -也可以不做映射,直接转发 SaaS 的完整 `taskList`;V2 同时兼容  
57 -`classId/meetingNumber`、`taskId/id` 两种字段名。`id` 和  
58 -`beginTime/endTime` 不是状态回写的必填字段。 57 +`/highlight/recording/scheduled` 仅用于按站点批量补录,不作为正式任务入口。
59 58
60 ## 灰度与回滚 59 ## 灰度与回滚
61 60
@@ -7,6 +7,8 @@ @@ -7,6 +7,8 @@
7 - [部署说明](./HIGHLIGHT_DEPLOYMENT.md) 7 - [部署说明](./HIGHLIGHT_DEPLOYMENT.md)
8 - [部署与验收测试](./HIGHLIGHT_DEPLOYMENT_USAGE_TEST.md) 8 - [部署与验收测试](./HIGHLIGHT_DEPLOYMENT_USAGE_TEST.md)
9 9
  10 +正式 cron 入口:`POST /recordingTaskV2/scheduled`。cron 只触发该入口,SaaS 任务拉取和分流由 WebScreen 完成。
  11 +
10 SaaS 已有接口: 12 SaaS 已有接口:
11 13
12 - `/3m/api/recording/getRecordingTasksPrivate.do` 14 - `/3m/api/recording/getRecordingTasksPrivate.do`
@@ -13,6 +13,7 @@ POST /fileExists @@ -13,6 +13,7 @@ POST /fileExists
13 13
14 ```http 14 ```http
15 POST /recordingTaskV2 15 POST /recordingTaskV2
  16 +POST /recordingTaskV2/scheduled
16 POST /fileExistsV2 17 POST /fileExistsV2
17 ``` 18 ```
18 19
@@ -23,14 +24,76 @@ V2 兼容原整课录制格式,并通过 SaaS 字段 `onlyHighlight=1` 支持 @@ -23,14 +24,76 @@ V2 兼容原整课录制格式,并通过 SaaS 字段 `onlyHighlight=1` 支持
23 | 不传或 `0` | 复用原整课录制实现 | 24 | 不传或 `0` | 复用原整课录制实现 |
24 | `1` | 仅录制该课堂的高光时刻 | 25 | `1` | 仅录制该课堂的高光时刻 |
25 26
26 -## 2. 创建 V2 录制任务 27 +## 2. cron 统一入口
  28 +
  29 +正式定时任务只调用 WebScreen 内部入口,不负责查询或转换 SaaS 数据:
  30 +
  31 +```http
  32 +POST /recordingTaskV2/scheduled
  33 +```
  34 +
  35 +请求体为空。WebScreen 会自行完成:
  36 +
  37 +1. 分页调用 SaaS `/3m/api/recording/getRecordingTasksPrivate.do`。
  38 +2. 生成接口要求的 `timestamp` 和 `authId`。
  39 +3. 根据 `onlyHighlight` 将任务交给现有 V2 整课或高光流程。
  40 +4. 使用本地任务快照防止 cron 重复投递。
  41 +
  42 +该入口只允许从 WebScreen 所在服务器的回环地址调用,不作为公网接口开放。
  43 +
  44 +配置:
  45 +
  46 +```json
  47 +"RECORDINGV2CONFIG": {
  48 + "taskListUrl": "https://saas.xuedianyun.com/3m/api/recording/getRecordingTasksPrivate.do",
  49 + "pageSize": 100,
  50 + "maxPages": 1000,
  51 + "maxTasksPerRun": 4,
  52 + "maxFullConcurrent": 2,
  53 + "apiTimeoutMs": 10000,
  54 + "apiRetryCount": 2,
  55 + "apiRetryBaseDelayMs": 500
  56 +}
  57 +```
  58 +
  59 +`maxTasksPerRun` 限制一次 cron 最多新接收多少个任务;重复任务不占用该额度。
  60 +`maxFullConcurrent` 限制同时运行的整课录制进程数。高光录制仍由
  61 +`HIGHLIGHTCONFIG.maxConcurrent` 单独限制。
  62 +
  63 +cron 示例:
  64 +
  65 +```cron
  66 +31 7-20 * * * curl -fsS -X POST http://127.0.0.1:3001/recordingTaskV2/scheduled >> /var/log/webscreen_recording_v2.log 2>&1
  67 +```
  68 +
  69 +成功响应:
  70 +
  71 +```json
  72 +{
  73 + "code": "0",
  74 + "message": "success",
  75 + "data": {
  76 + "received": 12,
  77 + "taskCount": 12,
  78 + "pages": 1,
  79 + "inspected": 5,
  80 + "accepted": 4,
  81 + "duplicates": 1,
  82 + "noMedia": 0,
  83 + "invalid": 0,
  84 + "limited": true
  85 + }
  86 +}
  87 +```
  88 +
  89 +## 3. 创建 V2 录制任务
27 90
28 ```http 91 ```http
29 POST /recordingTaskV2 92 POST /recordingTaskV2
30 Content-Type: application/json 93 Content-Type: application/json
31 ``` 94 ```
32 95
33 -### 2.1 兼容原整课请求 96 +### 3.1 兼容原整课请求
34 97
35 ```json 98 ```json
36 { 99 {
@@ -47,7 +110,7 @@ Content-Type: application/json @@ -47,7 +110,7 @@ Content-Type: application/json
47 110
48 未传 `onlyHighlight` 时,V2 调用现有整课录制代码,文件仍为 `{classId}.mp4`。 111 未传 `onlyHighlight` 时,V2 调用现有整课录制代码,文件仍为 `{classId}.mp4`。
49 112
50 -### 2.2 仅高光请求 113 +### 3.2 仅高光请求
51 114
52 V2 可直接接受 SaaS `getRecordingTasksPrivate.do` 返回项的字段形式: 115 V2 可直接接受 SaaS `getRecordingTasksPrivate.do` 返回项的字段形式:
53 116
@@ -90,7 +153,7 @@ V2 可直接接受 SaaS `getRecordingTasksPrivate.do` 返回项的字段形式 @@ -90,7 +153,7 @@ V2 可直接接受 SaaS `getRecordingTasksPrivate.do` 返回项的字段形式
90 153
91 `getBySitePrivate.do` 不属于单课堂任务主流程,只用于批量排查或补录。 154 `getBySitePrivate.do` 不属于单课堂任务主流程,只用于批量排查或补录。
92 155
93 -### 2.3 响应 156 +### 3.3 响应
94 157
95 ```json 158 ```json
96 { 159 {
@@ -109,7 +172,7 @@ V2 可直接接受 SaaS `getRecordingTasksPrivate.do` 返回项的字段形式 @@ -109,7 +172,7 @@ V2 可直接接受 SaaS `getRecordingTasksPrivate.do` 返回项的字段形式
109 | `duplicates` | 已处理或正在执行的重复任务数 | 172 | `duplicates` | 已处理或正在执行的重复任务数 |
110 | `noMedia` | `onlyHighlight=1` 但 SaaS 正常返回零条高光的任务数 | 173 | `noMedia` | `onlyHighlight=1` 但 SaaS 正常返回零条高光的任务数 |
111 174
112 -## 3. 查询 V2 录制文件 175 +## 4. 查询 V2 录制文件
113 176
114 ```http 177 ```http
115 POST /fileExistsV2 178 POST /fileExistsV2
@@ -186,7 +249,7 @@ V2 根据 `/recordingTaskV2` 保存的任务快照判断文件类型,客户不 @@ -186,7 +249,7 @@ V2 根据 `/recordingTaskV2` 保存的任务快照判断文件类型,客户不
186 249
187 查询接口只检查明确的 OSS Key,不触发录制,也不重新调用 SaaS 高光接口。 250 查询接口只检查明确的 OSS Key,不触发录制,也不重新调用 SaaS 高光接口。
188 251
189 -## 4. 文件规则 252 +## 5. 文件规则
190 253
191 ```text 254 ```text
192 整课:oss/{siteId}/{yyyyMMdd}/{classId}.mp4 255 整课:oss/{siteId}/{yyyyMMdd}/{classId}.mp4
@@ -103,10 +103,19 @@ class MediaCreat { @@ -103,10 +103,19 @@ class MediaCreat {
103 fs.appendFileSync(logFile, new Date().toLocaleString() + " " + text + '\r\n'); 103 fs.appendFileSync(logFile, new Date().toLocaleString() + " " + text + '\r\n');
104 } 104 }
105 105
106 - recordingCreat(id, siteId, type,yymmdd="") { 106 + recordingCreat(id, siteId, type,yymmdd="", done) {
  107 + let completed = false
  108 + const finish = error => {
  109 + if (completed || typeof done !== 'function') return
  110 + completed = true
  111 + done(error || null)
  112 + }
107 this.wrieLog(" 课堂录制开始:------>" + id) 113 this.wrieLog(" 课堂录制开始:------>" + id)
108 let fileConfig = this.getConfigFileJson() 114 let fileConfig = this.getConfigFileJson()
109 - if (!fileConfig) return false 115 + if (!fileConfig) {
  116 + finish(new Error('无效的fileConfig'))
  117 + return false
  118 + }
110 const { BACKMEDIACONFIG, PROJECTWINCATALOG, PROJECTCATALOG } = JSON.parse(fileConfig) 119 const { BACKMEDIACONFIG, PROJECTWINCATALOG, PROJECTCATALOG } = JSON.parse(fileConfig)
111 let mediaDir = PROJECTCATALOG + "/media/" 120 let mediaDir = PROJECTCATALOG + "/media/"
112 let classDir = PROJECTCATALOG + "/media/" + siteId 121 let classDir = PROJECTCATALOG + "/media/" + siteId
@@ -137,6 +146,7 @@ class MediaCreat { @@ -137,6 +146,7 @@ class MediaCreat {
137 this.wrieLog("files:" + files) 146 this.wrieLog("files:" + files)
138 if (files.indexOf(id + ".mp4") != -1) { 147 if (files.indexOf(id + ".mp4") != -1) {
139 this.wrieLog("已存在:" + id + "课堂号,停止继续录制") 148 this.wrieLog("已存在:" + id + "课堂号,停止继续录制")
  149 + finish()
140 if (type == 'post') { 150 if (type == 'post') {
141 if (classidPost.length) { 151 if (classidPost.length) {
142 let shiftData = classidPost.shift() 152 let shiftData = classidPost.shift()
@@ -163,13 +173,16 @@ class MediaCreat { @@ -163,13 +173,16 @@ class MediaCreat {
163 this.wrieLog(" 错误" + id + ":" + err) 173 this.wrieLog(" 错误" + id + ":" + err)
164 this.wrieLog(" 错误 stdout" + id + ":" + stdout) 174 this.wrieLog(" 错误 stdout" + id + ":" + stdout)
165 this.wrieLog(" 错误 stderr" + id + ":" + stderr) 175 this.wrieLog(" 错误 stderr" + id + ":" + stderr)
  176 + finish(err)
166 return 177 return
167 } 178 }
168 let files = fs.readdirSync(ymdDir); 179 let files = fs.readdirSync(ymdDir);
169 // let interValGetFile = setInterval(()=>{ 180 // let interValGetFile = setInterval(()=>{
170 if (files.indexOf(id + ".mp4") == -1) { 181 if (files.indexOf(id + ".mp4") == -1) {
171 this.wrieLog(" 课堂录制未发现该" + id + "课堂号") 182 this.wrieLog(" 课堂录制未发现该" + id + "课堂号")
  183 + finish(new Error(`课堂录制未生成文件: ${id}`))
172 } else { 184 } else {
  185 + finish()
173 if (type == 'get') { 186 if (type == 'get') {
174 if (classid.length) { 187 if (classid.length) {
175 let shiftData = classid.shift() 188 let shiftData = classid.shift()
@@ -410,18 +423,28 @@ router.post('/recordingTask', async function (req, res, next) { @@ -410,18 +423,28 @@ router.post('/recordingTask', async function (req, res, next) {
410 res.send({ code: "0",message:"success",v:version }); 423 res.send({ code: "0",message:"success",v:version });
411 }) 424 })
412 425
413 -router.post('/recordingTaskV2', async function (req, res, next) {  
414 - try {  
415 - const result = await recordingTaskService.acceptTasks(req.body || {}, {  
416 - recordFullClass: task => { 426 +function getRecordingTaskV2Handlers() {
  427 + return {
  428 + recordFullClass: task => new Promise((resolve, reject) => {
417 new MediaCreat().recordingCreat( 429 new MediaCreat().recordingCreat(
418 task.classId, 430 task.classId,
419 task.siteId, 431 task.siteId,
420 'post', 432 'post',
421 - task.classStartTime || task.classDate 433 + task.classStartTime || task.classDate,
  434 + error => error ? reject(error) : resolve()
422 ); 435 );
423 - }  
424 - }); 436 + })
  437 + };
  438 +}
  439 +
  440 +function isLoopbackRequest(req) {
  441 + const address = String(req.socket && req.socket.remoteAddress || '');
  442 + return address === '127.0.0.1' || address === '::1' || address === '::ffff:127.0.0.1';
  443 +}
  444 +
  445 +router.post('/recordingTaskV2', async function (req, res, next) {
  446 + try {
  447 + const result = await recordingTaskService.acceptTasks(req.body || {}, getRecordingTaskV2Handlers());
425 return res.send({ 448 return res.send({
426 code: "0", 449 code: "0",
427 message: "success", 450 message: "success",
@@ -441,4 +464,28 @@ router.post('/recordingTaskV2', async function (req, res, next) { @@ -441,4 +464,28 @@ router.post('/recordingTaskV2', async function (req, res, next) {
441 } 464 }
442 }) 465 })
443 466
  467 +// cron 只触发本入口;SaaS 任务查询、分页和 V2 分流均由 WebScreen 完成。
  468 +router.post('/recordingTaskV2/scheduled', async function (req, res, next) {
  469 + if (!isLoopbackRequest(req)) {
  470 + return res.status(403).send({ code: "1", message: "仅允许本机调用", data: {}, v: version });
  471 + }
  472 + try {
  473 + const result = await recordingTaskService.runScheduledTaskPull(getRecordingTaskV2Handlers());
  474 + return res.send({
  475 + code: "0",
  476 + message: result.received ? "success" : "无录制任务",
  477 + data: result,
  478 + v: version
  479 + });
  480 + } catch (error) {
  481 + if (error instanceof RecordingTaskValidationError || error instanceof HighlightValidationError) {
  482 + return res.status(400).send({ code: "1", message: error.message, data: {}, v: version });
  483 + }
  484 + if (error instanceof HighlightUpstreamError) {
  485 + return res.status(502).send({ code: String(error.upstreamCode), message: error.message, data: {}, v: version });
  486 + }
  487 + return next(error);
  488 + }
  489 +})
  490 +
444 module.exports = router 491 module.exports = router
@@ -103,6 +103,10 @@ function buildManifestKey(task) { @@ -103,6 +103,10 @@ function buildManifestKey(task) {
103 return `${task.siteId}:${task.classId}`; 103 return `${task.siteId}:${task.classId}`;
104 } 104 }
105 105
  106 +function delay(ms) {
  107 + return new Promise(resolve => setTimeout(resolve, ms));
  108 +}
  109 +
106 class RecordingTaskService { 110 class RecordingTaskService {
107 constructor(options) { 111 constructor(options) {
108 const opts = options || {}; 112 const opts = options || {};
@@ -110,9 +114,13 @@ class RecordingTaskService { @@ -110,9 +114,13 @@ class RecordingTaskService {
110 this.highlightService = opts.highlightService || new HighlightRecordingService({ configPath: this.configPath }); 114 this.highlightService = opts.highlightService || new HighlightRecordingService({ configPath: this.configPath });
111 this.inspectObjectOverride = opts.inspectObject; 115 this.inspectObjectOverride = opts.inspectObject;
112 this.updateTaskStatusOverride = opts.updateTaskStatus; 116 this.updateTaskStatusOverride = opts.updateTaskStatus;
  117 + this.fetchTaskPageOverride = opts.fetchTaskPage;
113 this.ossClient = opts.ossClient || null; 118 this.ossClient = opts.ossClient || null;
114 this.inFlight = new Set(); 119 this.inFlight = new Set();
115 this.objectStatusCache = new Map(); 120 this.objectStatusCache = new Map();
  121 + this.scheduledPullActive = false;
  122 + this.fullClassQueue = [];
  123 + this.fullClassActiveCount = 0;
116 } 124 }
117 125
118 readConfig() { 126 readConfig() {
@@ -154,6 +162,217 @@ class RecordingTaskService { @@ -154,6 +162,217 @@ class RecordingTaskService {
154 'https://xdymp4.xuedianyun.com').replace(/\/$/, ''); 162 'https://xdymp4.xuedianyun.com').replace(/\/$/, '');
155 } 163 }
156 164
  165 + getTaskPullConfig() {
  166 + const raw = this.readConfig();
  167 + const config = Object.assign({
  168 + taskListUrl: 'https://saas.xuedianyun.com/3m/api/recording/getRecordingTasksPrivate.do',
  169 + pageSize: 100,
  170 + maxPages: 1000,
  171 + maxTasksPerRun: 4,
  172 + maxFullConcurrent: 2,
  173 + apiTimeoutMs: 10000,
  174 + apiRetryCount: 2,
  175 + apiRetryBaseDelayMs: 500
  176 + }, raw.RECORDINGV2CONFIG || {});
  177 + try {
  178 + const url = new URL(String(config.taskListUrl || ''));
  179 + if (!['http:', 'https:'].includes(url.protocol)) throw new Error('protocol');
  180 + config.taskListUrl = url.toString();
  181 + } catch (error) {
  182 + throw new RecordingTaskValidationError('RECORDINGV2CONFIG.taskListUrl 无效');
  183 + }
  184 + const integerFields = [
  185 + ['pageSize', 1, 1000],
  186 + ['maxPages', 1, 1000],
  187 + ['maxTasksPerRun', 1, 1000],
  188 + ['maxFullConcurrent', 1, 100],
  189 + ['apiTimeoutMs', 1, 300000],
  190 + ['apiRetryCount', 0, 10],
  191 + ['apiRetryBaseDelayMs', 0, 60000]
  192 + ];
  193 + for (const [field, minimum, maximum] of integerFields) {
  194 + const value = Number(config[field]);
  195 + if (!Number.isSafeInteger(value) || value < minimum || value > maximum) {
  196 + throw new RecordingTaskValidationError(`RECORDINGV2CONFIG.${field} 无效`);
  197 + }
  198 + config[field] = value;
  199 + }
  200 + return config;
  201 + }
  202 +
  203 + async fetchTaskPage(pageNo, config) {
  204 + if (this.fetchTaskPageOverride) {
  205 + return this.fetchTaskPageOverride(pageNo, config.pageSize);
  206 + }
  207 + const timestamp = String(Date.now());
  208 + const body = {
  209 + pageNo,
  210 + pageSize: config.pageSize,
  211 + timestamp,
  212 + authId: crypto.createHash('md5')
  213 + .update(`${pageNo}${config.pageSize}${timestamp}`, 'utf8')
  214 + .digest('hex')
  215 + };
  216 + let response;
  217 + let lastError;
  218 + for (let attempt = 0; attempt <= config.apiRetryCount; attempt += 1) {
  219 + try {
  220 + response = await axios.post(config.taskListUrl, querystring.stringify(body), {
  221 + headers: { 'Content-Type': 'application/x-www-form-urlencoded' },
  222 + timeout: config.apiTimeoutMs
  223 + });
  224 + break;
  225 + } catch (error) {
  226 + lastError = error;
  227 + const status = error && error.response && error.response.status;
  228 + const retryable = !status || status >= 500;
  229 + if (!retryable || attempt >= config.apiRetryCount) {
  230 + const upstreamError = new HighlightUpstreamError(-1,
  231 + `录制任务列表请求失败${status ? ` HTTP ${status}` : ''}`);
  232 + upstreamError.cause = error;
  233 + throw upstreamError;
  234 + }
  235 + await delay(config.apiRetryBaseDelayMs * Math.pow(2, attempt));
  236 + }
  237 + }
  238 + if (!response) {
  239 + const upstreamError = new HighlightUpstreamError(-1, '录制任务列表请求失败');
  240 + upstreamError.cause = lastError;
  241 + throw upstreamError;
  242 + }
  243 + return response.data || {};
  244 + }
  245 +
  246 + async fetchReadyTasks() {
  247 + const config = this.getTaskPullConfig();
  248 + const tasks = [];
  249 + const seen = new Set();
  250 + let taskCount = 0;
  251 + let pages = 0;
  252 + for (let pageNo = 1; pageNo <= config.maxPages; pageNo += 1) {
  253 + const result = await this.fetchTaskPage(pageNo, config);
  254 + const code = Number(result.code);
  255 + if (code !== 0) {
  256 + throw new HighlightUpstreamError(code, `录制任务列表返回错误 code=${result.code}`);
  257 + }
  258 + if (!Array.isArray(result.taskList)) {
  259 + throw new HighlightUpstreamError(10, '录制任务列表 taskList 不是数组');
  260 + }
  261 + pages = pageNo;
  262 + const count = Number(result.taskCount);
  263 + if (Number.isFinite(count) && count >= 0) taskCount = count;
  264 + for (const task of result.taskList) {
  265 + const identity = task && task.id != null
  266 + ? `id:${task.id}`
  267 + : `task:${task && task.siteId}:${task && (task.meetingNumber || task.classId)}:${task && task.onlyHighlight}`;
  268 + if (!seen.has(identity)) {
  269 + seen.add(identity);
  270 + tasks.push(task);
  271 + }
  272 + }
  273 + const responsePageSize = Number(result.pageSize) || config.pageSize;
  274 + const responsePageNo = Number(result.pageNo) || pageNo;
  275 + const complete = result.taskList.length === 0 ||
  276 + (Number.isFinite(count) && responsePageNo * responsePageSize >= count) ||
  277 + (!Number.isFinite(count) && result.taskList.length < config.pageSize);
  278 + if (complete) break;
  279 + if (pageNo === config.maxPages) {
  280 + throw new HighlightUpstreamError(2, '录制任务列表分页超过安全上限');
  281 + }
  282 + }
  283 + return { config, tasks, taskCount, pages };
  284 + }
  285 +
  286 + async runScheduledTaskPull(handlers) {
  287 + if (this.scheduledPullActive) {
  288 + throw new RecordingTaskValidationError('V2录制任务正在读取中');
  289 + }
  290 + this.scheduledPullActive = true;
  291 + try {
  292 + const fetched = await this.fetchReadyTasks();
  293 + const summary = {
  294 + received: fetched.tasks.length,
  295 + taskCount: fetched.taskCount,
  296 + pages: fetched.pages,
  297 + inspected: 0,
  298 + accepted: 0,
  299 + duplicates: 0,
  300 + noMedia: 0,
  301 + invalid: 0,
  302 + limited: false
  303 + };
  304 + for (const task of fetched.tasks) {
  305 + if (summary.accepted >= fetched.config.maxTasksPerRun) {
  306 + summary.limited = true;
  307 + break;
  308 + }
  309 + summary.inspected += 1;
  310 + try {
  311 + const result = await this.acceptTasks({ list: [task] }, handlers);
  312 + summary.accepted += result.accepted;
  313 + summary.duplicates += result.duplicates;
  314 + summary.noMedia += result.noMedia;
  315 + } catch (error) {
  316 + if (error instanceof RecordingTaskValidationError || error instanceof HighlightValidationError) {
  317 + summary.invalid += 1;
  318 + continue;
  319 + }
  320 + throw error;
  321 + }
  322 + }
  323 + return summary;
  324 + } finally {
  325 + this.scheduledPullActive = false;
  326 + }
  327 + }
  328 +
  329 + enqueueFullClass(task, recordFullClass) {
  330 + if (typeof recordFullClass !== 'function') {
  331 + throw new RecordingTaskValidationError('整课录制处理器未配置');
  332 + }
  333 + const key = buildManifestKey(task);
  334 + this.inFlight.add(key);
  335 + this.fullClassQueue.push({ key, task, recordFullClass });
  336 + this.drainFullClassQueue();
  337 + }
  338 +
  339 + drainFullClassQueue() {
  340 + const maxConcurrent = this.getTaskPullConfig().maxFullConcurrent;
  341 + while (this.fullClassActiveCount < maxConcurrent && this.fullClassQueue.length > 0) {
  342 + const queued = this.fullClassQueue.shift();
  343 + this.fullClassActiveCount += 1;
  344 + this.processFullClassTask(queued).catch(error => {
  345 + console.error(`V2整课录制任务执行失败 ${queued.key}:`, error);
  346 + }).finally(() => {
  347 + this.fullClassActiveCount -= 1;
  348 + this.inFlight.delete(queued.key);
  349 + this.drainFullClassQueue();
  350 + });
  351 + }
  352 + }
  353 +
  354 + async processFullClassTask(queued) {
  355 + const task = queued.task;
  356 + try {
  357 + await queued.recordFullClass(task);
  358 + await this.updateTaskStatus(task, 2);
  359 + const manifest = this.loadManifest(queued.key);
  360 + if (manifest) {
  361 + manifest.status = 'completed';
  362 + manifest.updatedAt = Date.now();
  363 + this.saveManifest(manifest);
  364 + }
  365 + } catch (error) {
  366 + const manifest = this.loadManifest(queued.key);
  367 + if (manifest) {
  368 + manifest.status = 'failed';
  369 + manifest.updatedAt = Date.now();
  370 + this.saveManifest(manifest);
  371 + }
  372 + throw error;
  373 + }
  374 + }
  375 +
157 buildFullFile(task) { 376 buildFullFile(task) {
158 const ossKey = `oss/${task.siteId}/${task.classDate}/${task.classId}.mp4`; 377 const ossKey = `oss/${task.siteId}/${task.classDate}/${task.classId}.mp4`;
159 return { type: 'full', ossKey, url: `${this.getOutputBaseUrl()}/${ossKey}` }; 378 return { type: 'full', ossKey, url: `${this.getOutputBaseUrl()}/${ossKey}` };
@@ -338,12 +557,9 @@ class RecordingTaskService { @@ -338,12 +557,9 @@ class RecordingTaskService {
338 const processTask = async task => { 557 const processTask = async task => {
339 const key = buildManifestKey(task); 558 const key = buildManifestKey(task);
340 const existing = this.loadManifest(key); 559 const existing = this.loadManifest(key);
341 - const sameTaskId = task.taskId && existing && Array.isArray(existing.taskIds) &&  
342 - existing.taskIds.includes(task.taskId);  
343 const sameMode = existing && existing.onlyHighlight === task.onlyHighlight; 560 const sameMode = existing && existing.onlyHighlight === task.onlyHighlight;
344 const completed = existing && (existing.status === 'completed' || existing.status === 'no_media'); 561 const completed = existing && (existing.status === 'completed' || existing.status === 'no_media');
345 - const acceptedFullTask = sameTaskId && task.onlyHighlight === 0 && existing.status === 'accepted';  
346 - if (this.inFlight.has(key) || (sameMode && (completed || acceptedFullTask))) { 562 + if (this.inFlight.has(key) || (sameMode && completed)) {
347 duplicates += 1; 563 duplicates += 1;
348 return; 564 return;
349 } 565 }
@@ -351,7 +567,7 @@ class RecordingTaskService { @@ -351,7 +567,7 @@ class RecordingTaskService {
351 if (task.onlyHighlight === 0) { 567 if (task.onlyHighlight === 0) {
352 const manifest = this.buildManifest(task, [], 'accepted'); 568 const manifest = this.buildManifest(task, [], 'accepted');
353 this.saveManifest(manifest); 569 this.saveManifest(manifest);
354 - if (typeof recordFullClass === 'function') await recordFullClass(task); 570 + this.enqueueFullClass(task, recordFullClass);
355 accepted += 1; 571 accepted += 1;
356 return; 572 return;
357 } 573 }
@@ -33,6 +33,7 @@ async function run() { @@ -33,6 +33,7 @@ async function run() {
33 await new Promise(resolve => server.listen(0, '127.0.0.1', resolve)); 33 await new Promise(resolve => server.listen(0, '127.0.0.1', resolve));
34 const originalAcceptTasks = recordingTaskService.acceptTasks; 34 const originalAcceptTasks = recordingTaskService.acceptTasks;
35 const originalGetFileResult = recordingTaskService.getFileResult; 35 const originalGetFileResult = recordingTaskService.getFileResult;
  36 + const originalRunScheduledTaskPull = recordingTaskService.runScheduledTaskPull;
36 let acceptedCalls = 0; 37 let acceptedCalls = 0;
37 let fileCalls = 0; 38 let fileCalls = 0;
38 const acceptedBodies = []; 39 const acceptedBodies = [];
@@ -54,6 +55,20 @@ async function run() { @@ -54,6 +55,20 @@ async function run() {
54 files: [{ type: 'highlight', highlightId: 5, url: 'https://example/highlight.mp4' }] 55 files: [{ type: 'highlight', highlightId: 5, url: 'https://example/highlight.mp4' }]
55 }; 56 };
56 }; 57 };
  58 + recordingTaskService.runScheduledTaskPull = async handlers => {
  59 + assert.strictEqual(typeof handlers.recordFullClass, 'function');
  60 + return {
  61 + received: 3,
  62 + taskCount: 3,
  63 + pages: 1,
  64 + inspected: 3,
  65 + accepted: 2,
  66 + duplicates: 1,
  67 + noMedia: 0,
  68 + invalid: 0,
  69 + limited: false
  70 + };
  71 + };
57 72
58 const oldTask = await request(server, '/recordingTask', {}); 73 const oldTask = await request(server, '/recordingTask', {});
59 assert.strictEqual(oldTask.status, 200); 74 assert.strictEqual(oldTask.status, 200);
@@ -89,6 +104,12 @@ async function run() { @@ -89,6 +104,12 @@ async function run() {
89 assert.strictEqual(acceptedCalls, 2); 104 assert.strictEqual(acceptedCalls, 2);
90 assert.strictEqual(acceptedBodies[1].list[0].onlyHighlight, undefined); 105 assert.strictEqual(acceptedBodies[1].list[0].onlyHighlight, undefined);
91 106
  107 + const scheduledTask = await request(server, '/recordingTaskV2/scheduled', {});
  108 + assert.strictEqual(scheduledTask.status, 200);
  109 + assert.strictEqual(scheduledTask.body.code, '0');
  110 + assert.strictEqual(scheduledTask.body.data.accepted, 2);
  111 + assert.strictEqual(scheduledTask.body.data.duplicates, 1);
  112 +
92 const oldFile = await request(server, '/fileExists', { 113 const oldFile = await request(server, '/fileExists', {
93 siteId: 'doctest', classId: '1001' 114 siteId: 'doctest', classId: '1001'
94 }); 115 });
@@ -108,6 +129,7 @@ async function run() { @@ -108,6 +129,7 @@ async function run() {
108 } finally { 129 } finally {
109 recordingTaskService.acceptTasks = originalAcceptTasks; 130 recordingTaskService.acceptTasks = originalAcceptTasks;
110 recordingTaskService.getFileResult = originalGetFileResult; 131 recordingTaskService.getFileResult = originalGetFileResult;
  132 + recordingTaskService.runScheduledTaskPull = originalRunScheduledTaskPull;
111 await new Promise(resolve => server.close(resolve)); 133 await new Promise(resolve => server.close(resolve));
112 } 134 }
113 console.log('recording V2 routes tests passed'); 135 console.log('recording V2 routes tests passed');
@@ -63,6 +63,15 @@ async function run() { @@ -63,6 +63,15 @@ async function run() {
63 const configPath = path.join(tempRoot, 'config.json'); 63 const configPath = path.join(tempRoot, 'config.json');
64 fs.writeFileSync(configPath, JSON.stringify({ 64 fs.writeFileSync(configPath, JSON.stringify({
65 PROJECTCATALOG: tempRoot, 65 PROJECTCATALOG: tempRoot,
  66 + RECORDINGV2CONFIG: {
  67 + taskListUrl: 'https://saas.xuedianyun.com/3m/api/recording/getRecordingTasksPrivate.do',
  68 + pageSize: 2,
  69 + maxPages: 10,
  70 + maxTasksPerRun: 2,
  71 + apiTimeoutMs: 10000,
  72 + apiRetryCount: 0,
  73 + apiRetryBaseDelayMs: 0
  74 + },
66 HIGHLIGHTCONFIG: { 75 HIGHLIGHTCONFIG: {
67 enabled: true, 76 enabled: true,
68 siteIds: ['doctest'], 77 siteIds: ['doctest'],
@@ -155,7 +164,7 @@ async function run() { @@ -155,7 +164,7 @@ async function run() {
155 assert.strictEqual(calls.fetch, 1); 164 assert.strictEqual(calls.fetch, 1);
156 assert.strictEqual(calls.enqueue, 1); 165 assert.strictEqual(calls.enqueue, 1);
157 assert.strictEqual(calls.wait, 1); 166 assert.strictEqual(calls.wait, 1);
158 - assert.deepStrictEqual(statusUpdates, ['1002:2']); 167 + assert.deepStrictEqual(statusUpdates, ['1001:2', '1002:2']);
159 168
160 const highManifest = service.loadManifest(buildManifestKey({ siteId: 'doctest', classId: '1002' })); 169 const highManifest = service.loadManifest(buildManifestKey({ siteId: 'doctest', classId: '1002' }));
161 assert.strictEqual(highManifest.onlyHighlight, 1); 170 assert.strictEqual(highManifest.onlyHighlight, 1);
@@ -190,7 +199,7 @@ async function run() { @@ -190,7 +199,7 @@ async function run() {
190 onlyHighlight: 1 199 onlyHighlight: 1
191 }] }); 200 }] });
192 assert.deepStrictEqual(noMedia, { accepted: 1, duplicates: 0, noMedia: 1 }); 201 assert.deepStrictEqual(noMedia, { accepted: 1, duplicates: 0, noMedia: 1 });
193 - assert.deepStrictEqual(statusUpdates, ['1002:2', '1003:3']); 202 + assert.deepStrictEqual(statusUpdates, ['1001:2', '1002:2', '1003:3']);
194 const emptyResult = await service.getFileResult({ siteId: 'doctest', classId: '1003' }); 203 const emptyResult = await service.getFileResult({ siteId: 'doctest', classId: '1003' });
195 assert.strictEqual(emptyResult.fileExists, false); 204 assert.strictEqual(emptyResult.fileExists, false);
196 205
@@ -259,6 +268,82 @@ async function run() { @@ -259,6 +268,82 @@ async function run() {
259 releaseConcurrentFetch([]); 268 releaseConcurrentFetch([]);
260 assert.deepStrictEqual(await firstConcurrent, { accepted: 1, duplicates: 0, noMedia: 1 }); 269 assert.deepStrictEqual(await firstConcurrent, { accepted: 1, duplicates: 0, noMedia: 1 });
261 270
  271 + const scheduledFullCalls = [];
  272 + const scheduledService = new RecordingTaskService({
  273 + configPath,
  274 + highlightService,
  275 + updateTaskStatus: async () => {},
  276 + fetchTaskPage: async pageNo => ({
  277 + code: 0,
  278 + taskCount: 3,
  279 + pageNo,
  280 + pageSize: 2,
  281 + taskList: pageNo === 1 ? [{
  282 + id: 'scheduled-1', siteId: 'doctest', meetingNumber: '1010',
  283 + beginTime: '2026-08-05 10:00:00', endTime: '2026-08-05 11:00:00', onlyHighlight: 0
  284 + }, {
  285 + id: 'scheduled-2', siteId: 'doctest', meetingNumber: '1011',
  286 + beginTime: '2026-08-05 10:00:00', endTime: '2026-08-05 11:00:00', onlyHighlight: 0
  287 + }] : [{
  288 + id: 'scheduled-3', siteId: 'doctest', meetingNumber: '1012',
  289 + beginTime: '2026-08-05 10:00:00', endTime: '2026-08-05 11:00:00', onlyHighlight: 0
  290 + }]
  291 + })
  292 + });
  293 + const scheduledHandlers = {
  294 + recordFullClass: task => { scheduledFullCalls.push(task.classId); }
  295 + };
  296 + const firstScheduled = await scheduledService.runScheduledTaskPull(scheduledHandlers);
  297 + assert.deepStrictEqual(firstScheduled, {
  298 + received: 3,
  299 + taskCount: 3,
  300 + pages: 2,
  301 + inspected: 2,
  302 + accepted: 2,
  303 + duplicates: 0,
  304 + noMedia: 0,
  305 + invalid: 0,
  306 + limited: true
  307 + });
  308 + const secondScheduled = await scheduledService.runScheduledTaskPull(scheduledHandlers);
  309 + assert.strictEqual(secondScheduled.inspected, 3);
  310 + assert.strictEqual(secondScheduled.accepted, 1);
  311 + assert.strictEqual(secondScheduled.duplicates, 2);
  312 + assert.strictEqual(secondScheduled.limited, false);
  313 + assert.deepStrictEqual(scheduledFullCalls, ['1010', '1011', '1012']);
  314 + await waitForBackground(scheduledService);
  315 +
  316 + const fullQueueStarts = [];
  317 + const fullQueueReleases = [];
  318 + const fullQueueService = new RecordingTaskService({
  319 + configPath,
  320 + highlightService,
  321 + updateTaskStatus: async () => {}
  322 + });
  323 + const fullQueueResult = await fullQueueService.acceptTasks({
  324 + list: ['1020', '1021', '1022'].map(classId => ({
  325 + id: `queue-${classId}`,
  326 + siteId: 'doctest',
  327 + classId,
  328 + yymmdd: '20260805'
  329 + }))
  330 + }, {
  331 + recordFullClass: task => new Promise(resolve => {
  332 + fullQueueStarts.push(task.classId);
  333 + fullQueueReleases.push(resolve);
  334 + })
  335 + });
  336 + assert.strictEqual(fullQueueResult.accepted, 3);
  337 + assert.deepStrictEqual(fullQueueStarts, ['1020', '1021']);
  338 + assert.strictEqual(fullQueueService.fullClassActiveCount, 2);
  339 + assert.strictEqual(fullQueueService.fullClassQueue.length, 1);
  340 + fullQueueReleases[0]();
  341 + await new Promise(resolve => setImmediate(resolve));
  342 + assert.deepStrictEqual(fullQueueStarts, ['1020', '1021', '1022']);
  343 + fullQueueReleases[1]();
  344 + fullQueueReleases[2]();
  345 + await waitForBackground(fullQueueService);
  346 +
262 fs.rmSync(tempRoot, { recursive: true, force: true }); 347 fs.rmSync(tempRoot, { recursive: true, force: true });
263 console.log('recording task V2 service tests passed'); 348 console.log('recording task V2 service tests passed');
264 } 349 }