更新
This commit is contained in:
2026-02-23 18:52:32 +08:00
parent d429560511
commit 123fa6a7aa
81 changed files with 5952 additions and 407 deletions
@@ -0,0 +1,236 @@
using AutoMapper;
using IM_API.Configs.Options;
using IM_API.Domain.Events;
using IM_API.Dtos;
using IM_API.Exceptions;
using IM_API.Interface.Services;
using IM_API.Models.Upload;
using IM_API.Tools;
using IM_API.VOs;
using MassTransit;
using Microsoft.EntityFrameworkCore.Storage;
using StackExchange.Redis;
using System.Security.AccessControl;
using System.Security.Claims;
using System.Threading.Tasks;
using IDatabase = StackExchange.Redis.IDatabase;
namespace IM_API.Services
{
public class LocalStorageService : IStorageService
{
private readonly IMapper _mapper;
private readonly IHttpContextAccessor _httpContext;
private FileUploadOptions _options;
private readonly IUploadTaskService _uploadTaskService;
private readonly IDatabase _redis;
private readonly IHostEnvironment _env;
private readonly ILogger<LocalStorageService> _logger;
private readonly IPublishEndpoint _endpoint;
public LocalStorageService(IMapper mapper, IHttpContextAccessor httpContextAccessor,
IConfiguration configuration, IUploadTaskService uploadTaskService,
IConnectionMultiplexer connectionMultiplexer, IHostEnvironment hostEnvironment
, ILogger<LocalStorageService> logger, IPublishEndpoint publishEndpoint)
{
_mapper = mapper;
_httpContext = httpContextAccessor;
_options = configuration.GetSection("FileUploadOptions").Get<FileUploadOptions>()!;
_uploadTaskService = uploadTaskService;
_redis = connectionMultiplexer.GetDatabase();
_env = hostEnvironment;
_logger = logger;
_endpoint = publishEndpoint;
}
public UploadMode Mode => UploadMode.Proxy;
public string ProviderName => "Local";
public async Task<Guid> CompleteAsync(Guid taskId, List<UploadPartDto> parts)
{
var task = await _uploadTaskService.GetTaskAsync(taskId);
if(task is null)
throw new BaseException(CodeDefine.CHUNKE_NOT_FOUND);
var partsToCheck = Enumerable.Range(1, task.TotalChunks)
.Select(i => (RedisValue)i).ToArray();
var results = await _redis.SetContainsAsync(RedisKeys.GetUploadPartKey(taskId), partsToCheck);
// 3. 快速判断是否全部存在
bool isAllUploaded = results.All(exists => exists);
if (!isAllUploaded) throw new BaseException(CodeDefine.CHUNKE_NOT_FOUND);
await _endpoint.Publish(new UploadMergeEvent
{
AggregateId = taskId.ToString(),
OccurredAt = DateTime.UtcNow,
EventId = Guid.NewGuid(),
OperatorId = 0,
Parts = parts,
TaskId = taskId,
ChunckCount = task.TotalChunks,
ObjectName = task.ObjectName
});
return taskId;
}
public async Task MergeAsync(Guid taskId, string objectName, int totalChunks, List<UploadPartDto> parts)
{
var baseDir = Path.Combine(_env.ContentRootPath, "uploads");
var tempPath = Path.Combine(baseDir, "temp", taskId.ToString()); // 项目根目录下 uploads // 最终文件存储路径(这里可以用你之前 ObjectNameGenerator 生成的名字)
var finalPath = Path.Combine(baseDir, "files", objectName);
var finalDir = Path.GetDirectoryName(finalPath);
Directory.CreateDirectory(finalDir);
try
{
using (var finalStream = new FileStream(finalPath, FileMode.Create))
{
for (var i = 1; i <= totalChunks; i++)
{
var progress = (i * 100.0 / totalChunks);
if (i % 5 == 0 || i == totalChunks)
{
await _redis.HashSetAsync(RedisKeys.MergeStatus(taskId), new HashEntry[]
{
new("status", "processing"),
new("progress", progress.ToString("F2"))
});
}
var chunkPath = Path.Combine(tempPath, $"{i}.part.tmp");
if (!File.Exists(chunkPath))
throw new BaseException(CodeDefine.CHUNKE_NOT_FOUND);
using (var chunkStream = new FileStream(chunkPath, FileMode.Open))
{
await chunkStream.CopyToAsync(finalStream);
}
}
Directory.Delete(tempPath, true);
await _redis.KeyDeleteAsync(RedisKeys.GetUploadPartKey(taskId));
await _uploadTaskService.UpdateStatusAsync(taskId, UploadStatus.Completed);
await _redis.HashSetAsync(RedisKeys.MergeStatus(taskId), new HashEntry[]
{
new("status", "Completed"),
new("progress", "100"),
new("url", objectName)
});
}
}
catch (Exception e) when (e is not BaseException)
{
_logger.LogError(e, e.Message);
throw new BaseException(CodeDefine.CHUNKE_COMBINE_FAIL);
}
}
public async Task<UploadPartInstructionVo> CreatePartInstructionAsync(Guid taskId, int partNumer)
{
if (await _redis.SetContainsAsync(RedisKeys.GetUploadPartKey(taskId), partNumer)){
return new UploadPartInstructionVo
{
PartNumber = partNumer,
Skip = true,
Headers = new Dictionary<string, string>()
};
}
var request = _httpContext.HttpContext!.Request;
var scheme = request.Scheme; // http 或 https
var host = request.Host.Value; // localhost:5000 或域名
var baseUrl = $"{scheme}://{host}/api/upload/local/{taskId}/parts/{partNumer}";
var headers = new Dictionary<string, string>();
headers.Add("Content-Type", "multipart/form-data");
return new UploadPartInstructionVo
{
Method = "POST",
PartNumber = partNumer,
Skip = false,
Url = baseUrl,
Headers = headers
};
}
public async Task<UploadTask> UploadSmallFileAsync(Stream stream, string fileName, string fileType, long size, string hash)
{
var taskOld = await _uploadTaskService.GetTaskAsync(hash);
if (taskOld is not null) return taskOld;
var userId = _httpContext.HttpContext?.User.FindFirstValue(ClaimTypes.NameIdentifier);
var objectname = ObjectNameGenerator.Generate(new ObjectNameContext
{
ContentType = fileType,
FileName = fileName,
UserId = int.Parse(userId)
});
var path = GetDownloadUrl(objectname);
// 4. 将 Stream 写入本地文件
using (var fileStream = new FileStream(path, FileMode.Create, FileAccess.Write, FileShare.None))
{
await stream.CopyToAsync(fileStream);
}
var task = new UploadTask
{
CreatedAt = DateTime.UtcNow,
ChunkSize = (int)size,
ContentType = fileType,
FileHash = hash,
FileName = fileName,
FileSize = size,
Id = Guid.NewGuid(),
ObjectName = objectname,
ProviderUploadId = Guid.NewGuid().ToString(),
Status = UploadStatus.Completed,
StorageProvider = ProviderName,
TotalChunks = 1
};
await _uploadTaskService.AddAsync(task);
return task;
}
public string GetDownloadUrl(string objectname)
{
var baseDir = Path.Combine(_env.ContentRootPath, "uploads"); // 最终文件存储路径(这里可以用你之前 ObjectNameGenerator 生成的名字)
var finalPath = Path.Combine(baseDir, "files", objectname);
var finalDir = Path.GetDirectoryName(finalPath);
Directory.CreateDirectory(finalDir);
return finalPath;
}
public async Task<CreateUploadTaskVo> InitTaskAsync(CreateUploadTaskDto dto)
{
var userId = _httpContext.HttpContext.User.FindFirstValue(ClaimTypes.NameIdentifier);
UploadTask task = _mapper.Map<UploadTask>(dto);
var taskOld = await _uploadTaskService.GetTaskAsync(dto.FileHash);
if(taskOld != null)
{
var t = _mapper.Map<CreateUploadTaskVo>(taskOld);
t.Skip = false;
if (taskOld.Status == UploadStatus.Completed)
{
t.Skip = true;
}
return t;
task = taskOld;
}
task.ObjectName = ObjectNameGenerator.Generate(new ObjectNameContext
{
ContentType = task.ContentType,
FileName = task.FileName,
UserId = int.Parse(userId)
});
task.StorageProvider = ProviderName;
task.ProviderUploadId = Guid.NewGuid().ToString();
task.ChunkSize = _options.ChunkSize;
task.TotalChunks = (int)Math.Ceiling((double)task.FileSize / _options.ChunkSize);
await _uploadTaskService.AddAsync(task);
return _mapper.Map<CreateUploadTaskVo>(task);
}
}
}
+50 -3
View File
@@ -10,9 +10,11 @@ using IM_API.Tools;
using IM_API.VOs.Message;
using MassTransit;
using Microsoft.EntityFrameworkCore;
using Newtonsoft.Json;
using StackExchange.Redis;
using System.Text.Json;
using System.Text.RegularExpressions;
using System.Threading.Tasks;
using static MassTransit.Monitoring.Performance.BuiltInCounters;
using static Microsoft.EntityFrameworkCore.DbLoggerCategory;
@@ -28,10 +30,13 @@ namespace IM_API.Services
private readonly IPublishEndpoint _endpoint;
private readonly ISequenceIdService _sequenceIdService;
private readonly IUserService _userService;
private readonly IUploadTaskService _uploadService;
private readonly IHttpContextAccessor _httpContextAccessor;
public MessageService(
ImContext context, ILogger<MessageService> logger, IMapper mapper,
IPublishEndpoint publishEndpoint, ISequenceIdService sequenceIdService,
IUserService userService
IUserService userService, IUploadTaskService uploadTaskService,
IHttpContextAccessor httpContextAccessor
)
{
_context = context;
@@ -41,6 +46,8 @@ namespace IM_API.Services
_endpoint = publishEndpoint;
_sequenceIdService = sequenceIdService;
_userService = userService;
_uploadService = uploadTaskService;
_httpContextAccessor = httpContextAccessor;
}
public async Task<List<MessageBaseVo>> GetMessagesAsync(int userId,MessageQueryDto dto)
@@ -91,6 +98,12 @@ namespace IM_API.Services
foreach (var item in messages)
{
if(item.Type != MessageMsgType.Text)
{
var request = _httpContextAccessor.HttpContext?.Request;
var baseUrl = $"{request.Scheme}://{request.Host}";
item.Content = UrlTools.ProcessMessageUrl(item.Content, baseUrl);
}
if(userDict.TryGetValue(item.SenderId, out var user))
{
item.SenderName = user.NickName;
@@ -143,7 +156,10 @@ namespace IM_API.Services
var message = _mapper.Map<Message>(dto);
message.StreamKey = StreamKeyBuilder.Group(groupId);
message.SequenceId = await _sequenceIdService.GetNextSquenceIdAsync(message.StreamKey);
await _endpoint.Publish(_mapper.Map<MessageCreatedEvent>(message));
var publishData = _mapper.Map<MessageCreatedEvent>(message);
var request = _httpContextAccessor.HttpContext?.Request;
publishData.BaseUrl = $"{request.Scheme}://{request.Host}";
await _endpoint.Publish(publishData);
return _mapper.Map<MessageBaseVo>(message);
}
@@ -156,9 +172,40 @@ namespace IM_API.Services
var message = _mapper.Map<Message>(dto);
message.StreamKey = StreamKeyBuilder.Private(senderId, receiverId);
message.SequenceId = await _sequenceIdService.GetNextSquenceIdAsync(message.StreamKey);
await _endpoint.Publish(_mapper.Map<MessageCreatedEvent>(message));
var publishData = _mapper.Map<MessageCreatedEvent>(message);
var request = _httpContextAccessor.HttpContext?.Request;
publishData.BaseUrl = $"{request.Scheme}://{request.Host}";
await _endpoint.Publish(publishData);
return _mapper.Map<MessageBaseVo>(message);
}
#endregion
public async Task<MessageBaseDto> HandleFileMessageContentAsync(MessageBaseDto dto)
{
if(dto.Type == MessageMsgType.Text)
{
return dto;
}
var dic = JsonConvert.DeserializeObject<Dictionary<string, object>>(dto.Content);
if (dic == null || !dic.TryGetValue("fileId", out var fileIdObj))
throw new BaseException(CodeDefine.PARAMETER_ERROR);
var fileInfo = await _uploadService.GetTaskAsync(new Guid(fileIdObj.ToString()));
if (fileInfo is null)
throw new BaseException(CodeDefine.FILE_NOT_FOUND);
dic["url"] = fileInfo.ObjectName;
dic["provider"] = fileInfo.StorageProvider;
dic["size"] = fileInfo.FileSize;
dto.Content = JsonConvert.SerializeObject(dic);
return dto;
}
}
}
@@ -0,0 +1,43 @@
using IM_API.Interface.Services;
using IM_API.Models;
using IM_API.Models.Upload;
using Microsoft.EntityFrameworkCore;
namespace IM_API.Services
{
public class UploadTaskService : IUploadTaskService
{
private readonly ImContext _context;
public UploadTaskService(ImContext context)
{
_context = context;
}
public async Task AddAsync(UploadTask task)
{
_context.UploadTasks.Add(task);
await _context.SaveChangesAsync();
}
public async Task<UploadTask?> GetTaskAsync(Guid taskId)
{
return await _context.UploadTasks.FirstOrDefaultAsync(x => x.Id == taskId);
}
public async Task<UploadTask?> GetTaskAsync(string hash)
{
return await _context.UploadTasks.FirstOrDefaultAsync(x => x.FileHash == hash);
}
public async Task UpdateStatusAsync(Guid taskId, UploadStatus status)
{
var task = await _context.UploadTasks.FirstOrDefaultAsync(x => x.Id == taskId);
if (task != null)
{
task.Status = status;
_context.UploadTasks.Update(task);
await _context.SaveChangesAsync();
}
}
}
}