feat(media): 微信消息媒体 OSS 归档与下载地址回写
- 新增 MediaArchiveJob、MediaOssArchiveService、WechatMediaArchiveService 与回填命令 - Message/DataProcessing 等支持归档调度与持久化下载 URL - WebSocket 控制器整理;朋友圈与文档/定时任务说明更新 Made-with: Cursor
This commit is contained in:
@@ -4,7 +4,10 @@ namespace app\api\controller;
|
||||
|
||||
use app\api\model\WechatMessageModel;
|
||||
use app\common\service\FriendTransferService;
|
||||
use app\common\service\WechatMediaArchiveService;
|
||||
use app\job\MediaArchiveJob;
|
||||
use think\Db;
|
||||
use think\facade\Log;
|
||||
use think\facade\Request;
|
||||
|
||||
class MessageController extends BaseController
|
||||
@@ -418,6 +421,7 @@ class MessageController extends BaseController
|
||||
'type' => 1,
|
||||
'accountId' => $item['accountId'],
|
||||
'content' => $item['content'],
|
||||
'originalContent' => $item['content'],
|
||||
'createTime' => $createTime,
|
||||
'deleteTime' => $deleteTime,
|
||||
'isDeleted' => $item['isDeleted'] ?? false,
|
||||
@@ -463,6 +467,9 @@ class MessageController extends BaseController
|
||||
}else{
|
||||
$id = $data['id'];
|
||||
unset($data['id']);
|
||||
if (!empty($exists['originalContent'])) {
|
||||
unset($data['originalContent']);
|
||||
}
|
||||
$res = $exists->save($data);
|
||||
}
|
||||
|
||||
@@ -496,6 +503,9 @@ class MessageController extends BaseController
|
||||
}
|
||||
}
|
||||
}
|
||||
if (in_array((int)($item['msgType'] ?? 0), [3, 34, 43, 47, 49], true) && !empty($id)) {
|
||||
MediaArchiveJob::dispatch('message', $id, ['source' => 'saveMessage']);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
@@ -577,6 +587,9 @@ class MessageController extends BaseController
|
||||
throw new \Exception('更新群聊消息记录失败');
|
||||
}
|
||||
}
|
||||
if (in_array((int)($item['msgType'] ?? 0), [3, 34, 43, 47, 49], true)) {
|
||||
MediaArchiveJob::dispatch('message', $item['id'], ['source' => 'saveChatroomMessage']);
|
||||
}
|
||||
return true;
|
||||
} catch (\Exception $e) {
|
||||
// 记录错误日志,便于调试
|
||||
@@ -588,6 +601,29 @@ class MessageController extends BaseController
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 持久化视频/文件下载后的真实地址,并重新触发 OSS 归档
|
||||
* @param array $data
|
||||
* @return bool
|
||||
*/
|
||||
public function updateDownloadedMessageMedia($data)
|
||||
{
|
||||
$messageId = (int)($data['friendMessageId'] ?? $data['chatroomMessageId'] ?? 0);
|
||||
$downloadUrl = trim((string)($data['url'] ?? ''));
|
||||
|
||||
if ($messageId <= 0 || empty($downloadUrl)) {
|
||||
return false;
|
||||
}
|
||||
|
||||
$updated = WechatMediaArchiveService::updateDownloadedMessageMedia($messageId, $downloadUrl);
|
||||
if (!$updated) {
|
||||
return false;
|
||||
}
|
||||
|
||||
MediaArchiveJob::dispatch('message', $messageId, ['source' => $data['type'] ?? 'download_result']);
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理消息内容,提取发送者ID和消息内容
|
||||
* @param string $content 原始消息内容
|
||||
|
||||
@@ -12,7 +12,7 @@ use think\facade\Env;
|
||||
use app\api\model\WechatFriendModel as WechatFriend;
|
||||
use app\api\model\WechatMomentsModel as WechatMoments;
|
||||
use think\facade\Cache;
|
||||
use app\common\util\AliyunOSS;
|
||||
use app\common\service\MediaOssArchiveService;
|
||||
|
||||
|
||||
class WebSocketController extends BaseController
|
||||
@@ -468,7 +468,7 @@ class WebSocketController extends BaseController
|
||||
// 更新数据库:保存原始URL和OSS URL,并标记已上传
|
||||
$updateData = [
|
||||
'resUrls' => $urls,
|
||||
'isOssUploaded' => 1, // 标识已上传到OSS
|
||||
'isOssUploaded' => !empty($ossUrls) ? 1 : 0,
|
||||
'update_time' => time()
|
||||
];
|
||||
|
||||
@@ -513,7 +513,7 @@ class WebSocketController extends BaseController
|
||||
// 更新数据库:保存原始URL和OSS URL,并标记已上传
|
||||
$updateData = [
|
||||
'resUrls' => $urls,
|
||||
'isOssUploaded' => 1, // 标识已上传到OSS
|
||||
'isOssUploaded' => !empty($ossUrls) ? 1 : 0,
|
||||
'update_time' => time()
|
||||
];
|
||||
|
||||
@@ -557,76 +557,23 @@ class WebSocketController extends BaseController
|
||||
if (empty($urls) || !is_array($urls)) {
|
||||
return $ossUrls;
|
||||
}
|
||||
|
||||
try {
|
||||
// 创建临时目录(兼容无 runtime_path() 辅助函数的环境)
|
||||
if (function_exists('runtime_path')) {
|
||||
$baseRuntimePath = rtrim(runtime_path(), DIRECTORY_SEPARATOR) . DIRECTORY_SEPARATOR;
|
||||
} elseif (defined('RUNTIME_PATH')) {
|
||||
$baseRuntimePath = rtrim(RUNTIME_PATH, DIRECTORY_SEPARATOR) . DIRECTORY_SEPARATOR;
|
||||
} else {
|
||||
// 兜底:使用项目根目录下的 runtime 目录
|
||||
$baseRuntimePath = rtrim(ROOT_PATH, DIRECTORY_SEPARATOR) . DIRECTORY_SEPARATOR . 'runtime' . DIRECTORY_SEPARATOR;
|
||||
|
||||
foreach ($urls as $url) {
|
||||
if (!MediaOssArchiveService::isRemoteHttpUrl($url)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
$tempDir = $baseRuntimePath . 'temp' . DIRECTORY_SEPARATOR . 'moments' . DIRECTORY_SEPARATOR . date('Y' . DIRECTORY_SEPARATOR . 'm' . DIRECTORY_SEPARATOR . 'd') . DIRECTORY_SEPARATOR;
|
||||
$resourceType = preg_match('/\.(mp4|mov|avi|webm|mkv)(\?.*)?$/i', $url) ? 'video' : 'image';
|
||||
$result = MediaOssArchiveService::archiveRemoteUrl($url, 'moments', $resourceType, (string)$snsId);
|
||||
if (!empty($result['success']) && !empty($result['url'])) {
|
||||
$ossUrls[] = $result['url'];
|
||||
continue;
|
||||
}
|
||||
|
||||
if (!is_dir($tempDir)) {
|
||||
mkdir($tempDir, 0755, true);
|
||||
}
|
||||
|
||||
foreach ($urls as $index => $url) {
|
||||
if (empty($url)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
try {
|
||||
// 下载图片到临时文件
|
||||
$tempFile = $tempDir . md5($url . $snsId . $index) . '.jpg';
|
||||
|
||||
// 使用curl下载图片
|
||||
$ch = curl_init($url);
|
||||
$fp = fopen($tempFile, 'wb');
|
||||
curl_setopt($ch, CURLOPT_FILE, $fp);
|
||||
curl_setopt($ch, CURLOPT_HEADER, 0);
|
||||
curl_setopt($ch, CURLOPT_FOLLOWLOCATION, true);
|
||||
curl_setopt($ch, CURLOPT_TIMEOUT, 30);
|
||||
curl_setopt($ch, CURLOPT_SSL_VERIFYPEER, false);
|
||||
curl_exec($ch);
|
||||
$httpCode = curl_getinfo($ch, CURLINFO_HTTP_CODE);
|
||||
curl_close($ch);
|
||||
fclose($fp);
|
||||
|
||||
if ($httpCode != 200 || !file_exists($tempFile) || filesize($tempFile) == 0) {
|
||||
Log::warning('下载朋友圈图片失败:' . $url . ', HTTP Code: ' . $httpCode);
|
||||
@unlink($tempFile);
|
||||
continue;
|
||||
}
|
||||
|
||||
// 生成OSS对象名称
|
||||
$objectName = 'moments/' . date('Y/m/d/') . md5($snsId . $index . time()) . '.jpg';
|
||||
|
||||
// 上传到OSS
|
||||
$result = AliyunOSS::uploadFile($tempFile, $objectName);
|
||||
if ($result['success']) {
|
||||
$ossUrls[] = $result['url'];
|
||||
} else {
|
||||
Log::error('朋友圈图片上传OSS失败:' . $url . ', 错误:' . ($result['error'] ?? '未知错误'));
|
||||
}
|
||||
|
||||
// 删除临时文件
|
||||
@unlink($tempFile);
|
||||
|
||||
} catch (\Exception $e) {
|
||||
Log::error('上传朋友圈图片到OSS异常:' . $e->getMessage() . ', URL: ' . $url);
|
||||
if (isset($tempFile) && file_exists($tempFile)) {
|
||||
@unlink($tempFile);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
} catch (\Exception $e) {
|
||||
Log::error('上传朋友圈图片到OSS异常:' . $e->getMessage());
|
||||
Log::error('朋友圈媒体上传OSS失败:' . ($result['error'] ?? '未知错误'), [
|
||||
'snsId' => $snsId,
|
||||
'url' => $url,
|
||||
]);
|
||||
}
|
||||
|
||||
return $ossUrls;
|
||||
@@ -696,8 +643,8 @@ class WebSocketController extends BaseController
|
||||
}
|
||||
|
||||
// 获取资源链接(检查是否已上传到OSS,如果已上传则跳过)
|
||||
if(empty($momentEntity['urls']) || $moment['type'] != 1) {
|
||||
// 如果没有urls或类型不是1,跳过
|
||||
if(empty($momentEntity['urls'])) {
|
||||
// 如果没有urls,跳过
|
||||
} elseif ($isOssUploaded == 1) {
|
||||
// 如果已上传到OSS,跳过采集
|
||||
} else {
|
||||
|
||||
Reference in New Issue
Block a user