TaskCenterService.php 9.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301
  1. <?php
  2. namespace App\Services;
  3. use App\Facade\Site;
  4. use App\Models\MpTaskCenter;
  5. use Illuminate\Support\Facades\DB;
  6. class TaskCenterService
  7. {
  8. /**
  9. * 创建任务中心记录
  10. *
  11. * @param string $type 任务类型:text/image/video
  12. * @param array $data 任务数据(title/ref_task_id/prompt/params等)
  13. * @return MpTaskCenter
  14. */
  15. public function createTask(string $type, array $data = [])
  16. {
  17. $uid = 0;
  18. $cpid = 0;
  19. try {
  20. $uid = (int)Site::getUid();
  21. $cpid = (int)Site::getCpid();
  22. } catch (\Throwable $e) {
  23. // 非请求上下文(如命令行)下可能无法获取用户ID,置为0
  24. }
  25. return MpTaskCenter::create([
  26. 'uid' => $data['uid'] ?? $uid,
  27. 'cpid' => $data['cpid'] ?? $cpid,
  28. 'task_type' => $type,
  29. 'title' => $data['title'] ?? '',
  30. 'ref_task_id' => $data['ref_task_id'] ?? 0,
  31. 'status' => $data['status'] ?? MpTaskCenter::STATUS_PROCESSING,
  32. 'result' => $data['result'] ?? null,
  33. 'error_message'=> $data['error_message'] ?? null,
  34. 'prompt' => $data['prompt'] ?? null,
  35. 'params' => $data['params'] ?? null,
  36. 'params_md5' => $data['params_md5'] ?? null,
  37. ]);
  38. }
  39. /**
  40. * 更新任务中心记录
  41. *
  42. * @param int $taskId
  43. * @param array $data
  44. * @return bool
  45. */
  46. public function updateTask(int $taskId, array $data = []): bool
  47. {
  48. return (bool)MpTaskCenter::where('id', $taskId)->update($data);
  49. }
  50. /**
  51. * 查询任务详情
  52. *
  53. * @param int $taskId
  54. * @return MpTaskCenter|null
  55. */
  56. public function getTaskDetail(int $taskId)
  57. {
  58. $query = MpTaskCenter::where('id', $taskId);
  59. // 请求上下文下按当前用户过滤
  60. $uid = 0;
  61. try {
  62. $uid = (int)Site::getUid();
  63. } catch (\Throwable $e) {
  64. }
  65. if ($uid > 0) {
  66. $query->where('uid', $uid);
  67. }
  68. return $query->first();
  69. }
  70. /**
  71. * 分页查询任务列表
  72. *
  73. * @param array $params
  74. * @return \Illuminate\Contracts\Pagination\LengthAwarePaginator
  75. */
  76. public function getTaskList(array $params = [])
  77. {
  78. $query = MpTaskCenter::query();
  79. // 请求上下文下按当前用户过滤
  80. $uid = 0;
  81. try {
  82. $uid = (int)Site::getUid();
  83. } catch (\Throwable $e) {
  84. }
  85. if ($uid > 0) {
  86. $query->where('uid', $uid);
  87. }
  88. // 按任务ID筛选
  89. if (!empty($params['task_id'])) {
  90. $query->where('id', (int)$params['task_id']);
  91. }
  92. // 按任务状态筛选(支持逗号分隔多状态)
  93. if (!empty($params['status'])) {
  94. $statuses = is_array($params['status'])
  95. ? $params['status']
  96. : array_filter(array_map('trim', explode(',', (string)$params['status'])));
  97. if (!empty($statuses)) {
  98. $query->whereIn('status', $statuses);
  99. }
  100. }
  101. // 按任务类型筛选
  102. if (!empty($params['task_type'])) {
  103. $query->where('task_type', $params['task_type']);
  104. }
  105. $pageSize = (int)($params['page_size'] ?? 20);
  106. if ($pageSize <= 0 || $pageSize > 100) {
  107. $pageSize = 20;
  108. }
  109. return $query->orderBy('created_at', 'desc')->orderBy('id', 'desc')->paginate($pageSize);
  110. }
  111. /**
  112. * 定时同步任务中心状态和结果
  113. *
  114. * 图片任务关联 mp_generate_pic_tasks,视频任务关联 mp_generate_video_tasks,
  115. * 将底层任务的最新状态、结果和错误信息同步到任务中心。
  116. *
  117. * @return int 同步更新的记录数
  118. */
  119. public function syncTaskStatus(): int
  120. {
  121. $updated = 0;
  122. // 只同步尚未结束的任务(避免重复扫描已完成记录)
  123. $tasks = MpTaskCenter::whereIn('status', [
  124. MpTaskCenter::STATUS_PENDING,
  125. MpTaskCenter::STATUS_PROCESSING,
  126. ])
  127. ->where('ref_task_id', '>', 0)
  128. ->orderBy('id', 'desc')
  129. ->limit(500)
  130. ->get();
  131. foreach ($tasks as $task) {
  132. try {
  133. if ($task->task_type === MpTaskCenter::TYPE_IMAGE) {
  134. $updated += $this->syncImageTask($task);
  135. } elseif ($task->task_type === MpTaskCenter::TYPE_VIDEO) {
  136. $updated += $this->syncVideoTask($task);
  137. }
  138. } catch (\Exception $e) {
  139. dLog('command')->error('任务中心同步失败: ' . $e->getMessage(), [
  140. 'task_id' => $task->id,
  141. 'ref_task_id' => $task->ref_task_id,
  142. 'task_type' => $task->task_type,
  143. ]);
  144. }
  145. }
  146. return $updated;
  147. }
  148. /**
  149. * 同步图片任务状态到任务中心
  150. *
  151. * @param MpTaskCenter $task
  152. * @return int
  153. */
  154. private function syncImageTask(MpTaskCenter $task): int
  155. {
  156. $ref = DB::table('mp_generate_pic_tasks')->where('id', $task->ref_task_id)->first();
  157. if (!$ref) {
  158. return 0;
  159. }
  160. $result = null;
  161. if (!empty($ref->result_url)) {
  162. $urls = $this->normalizeResultUrls($ref->result_url);
  163. $urls = array_values(array_filter($urls));
  164. // 与文生图 completed 返回格式保持一致
  165. $resultData = [
  166. 'msg' => '',
  167. 'code' => 0,
  168. 'data' => $urls[0] ?? '',
  169. 'task_center_id' => $task->id,
  170. ];
  171. if (count($urls) > 1) {
  172. $resultData['image_urls'] = $urls;
  173. }
  174. $result = json_encode($resultData, JSON_UNESCAPED_UNICODE);
  175. }
  176. return $this->applyRefStatus($task, $ref->status, $result, $ref->error_message);
  177. }
  178. /**
  179. * 规范化图片结果URL(兼容字符串/数组/JSON字符串/双重编码)
  180. *
  181. * @param mixed $resultUrl
  182. * @return array
  183. */
  184. private function normalizeResultUrls($resultUrl): array
  185. {
  186. if (is_array($resultUrl)) {
  187. return array_values(array_filter($resultUrl));
  188. }
  189. $value = (string)$resultUrl;
  190. // 最多解析两层 JSON(防御双重编码)
  191. for ($i = 0; $i < 2; $i++) {
  192. if (!is_string($value) || !is_json($value)) {
  193. break;
  194. }
  195. $decoded = json_decode($value, true);
  196. if (!is_array($decoded)) {
  197. // 解码结果是字符串且仍是JSON,继续解析下一层
  198. if (is_string($decoded) && is_json($decoded)) {
  199. $value = $decoded;
  200. continue;
  201. }
  202. $value = $decoded;
  203. break;
  204. }
  205. $value = $decoded;
  206. }
  207. if (is_array($value)) {
  208. return array_values(array_filter($value));
  209. }
  210. if (is_string($value) && $value !== '') {
  211. return [$value];
  212. }
  213. return [];
  214. }
  215. /**
  216. * 同步视频任务状态到任务中心
  217. *
  218. * @param MpTaskCenter $task
  219. * @return int
  220. */
  221. private function syncVideoTask(MpTaskCenter $task): int
  222. {
  223. $ref = DB::table('mp_generate_video_tasks')->where('id', $task->ref_task_id)->first();
  224. if (!$ref) {
  225. return 0;
  226. }
  227. $result = null;
  228. if (!empty($ref->result_url) || !empty($ref->compressed_url) || !empty($ref->last_frame_url)) {
  229. // 与文生视频 completed 返回格式保持一致
  230. $result = json_encode([
  231. 'task_id' => $ref->id,
  232. 'status' => $ref->status,
  233. 'video_url' => $ref->compressed_url ?: $ref->result_url,
  234. 'origin_video_url' => $ref->result_url,
  235. 'last_frame_url' => $ref->last_frame_url,
  236. 'error_message' => $ref->error_message ? mapErrorMessage($ref->error_message) : '',
  237. ], JSON_UNESCAPED_UNICODE);
  238. }
  239. return $this->applyRefStatus($task, $ref->status, $result, $ref->error_message);
  240. }
  241. /**
  242. * 将底层任务状态应用到任务中心记录
  243. *
  244. * @param MpTaskCenter $task
  245. * @param string $refStatus
  246. * @param string|null $result
  247. * @param string|null $errorMessage
  248. * @return int
  249. */
  250. private function applyRefStatus(MpTaskCenter $task, string $refStatus, $result, $errorMessage): int
  251. {
  252. $statusMap = [
  253. 'pending' => MpTaskCenter::STATUS_PENDING,
  254. 'processing' => MpTaskCenter::STATUS_PROCESSING,
  255. 'success' => MpTaskCenter::STATUS_SUCCESS,
  256. 'failed' => MpTaskCenter::STATUS_FAILED,
  257. ];
  258. $newStatus = $statusMap[$refStatus] ?? $task->status;
  259. $updateData = ['status' => $newStatus];
  260. if ($result !== null) {
  261. $updateData['result'] = $result;
  262. }
  263. if ($errorMessage !== null) {
  264. $updateData['error_message'] = $errorMessage;
  265. }
  266. return $this->updateTask((int)$task->id, $updateData) ? 1 : 0;
  267. }
  268. }