using AutoMapper; using FileService.Application.Ports; using FileService.Application.StorageContracts; using FileService.Domain.IReposities; using IM.Commons; using IM.Commons.IntegrationEvents; using IM.InitCommon; using MassTransit; using Microsoft.Extensions.Options; namespace FileService.Application.UploadFileTask { public class UploadFileTaskService(IUploadTaskReposity reposity, IMapper mapper, IObjectStorageRouter router, IOptions options, IStorageRedisCache redis, IPublishEndpoint endpoint, ILocalChunkStorage localChunkStorage ) { private readonly IUploadTaskReposity reposity = reposity; private readonly IMapper mapper = mapper; private readonly IObjectStorageRouter router = router; private readonly IOptions options = options; private readonly IStorageRedisCache redis = redis; private readonly IPublishEndpoint endpoint = endpoint; private readonly ILocalChunkStorage localChunkStorage = localChunkStorage; private readonly IObjectStoragePort storage = router.Route(options.Value.DefaultProviderCode); public async Task> InitTaskAsync(UploadTaskInitCommand command) { CancellationToken cancellationToken = CancellationToken.None; var task = command.ToUploadTask(); var date = DateTime.Now; var storageOption = options.Value.Providers[options.Value.DefaultProviderCode]; var storage = router.Route(storageOption.ProviderCode); var initUpdateCommand = new StorageContracts.InitiateUploadCommand( ProviderCode: storageOption.ProviderCode, Bucket: storageOption.Bucket, ObjectKey: $"{storageOption.LocalRootPath}\\{date.Year}\\{date.Month}\\{date.Day}\\{command.FileName}", ContentType: task.ContentType.Value, ContentLength: command.FileSize, null); var initRes = await storage.InitUploadAsync(initUpdateCommand, cancellationToken); var res = mapper.Map(initRes); res.TaskId = task.Id; task.StartUpload(); reposity.Create(task); await redis.SetAsync(new StorageContracts.UploadRuntimeCache( taskId: task.Id.ToString(), providerCode: storage.ProviderCode, uploadSessionId: res.UploadSessionId, bucket: storageOption.Bucket, region: storageOption.Region, objectKey: initUpdateCommand.ObjectKey, fileSize: task.FileSize, totalPartCount: (int)(task.FileSize % storageOption.DefaultPartSizeBytes > 0 ? (task.FileSize / storageOption.DefaultPartSizeBytes) + 1 : task.FileSize / storageOption.DefaultPartSizeBytes) )); return Result.Success(res); } public async Task> GenerateUrlAsync(string sessionId, int partNum, Guid userId, CancellationToken token = default) { var taskCache = await redis.GetAsync(sessionId); if(taskCache is null) { return Result.Fail(ResultCode.CHUNK_NOT_FOUND); } if(taskCache.TotalPartCount < partNum) { return Result.Fail(ResultCode.CHUNK_NOT_FOUND); } var presignUrl = await storage.GenerateUploadUrlAsync(new GenerateUploadUrlCommand( ProviderCode: taskCache.ProviderCode, Bucket: taskCache.Bucket, ObjectKey: taskCache.ObjectKey, UploadSessionId: taskCache.UploadSessionId, PartNumber: partNum, ExpiresIn: options.Value.Providers[options.Value.DefaultProviderCode].UploadUrlExpiresIn ), token); return Result.Success(presignUrl); } public async Task> CompleteTaskAsync(UploadTaskCompleteCommand command, CancellationToken cancellationToken = default) { var taskCache = await redis.GetAsync(command.UploadSessionId); if(taskCache is null) { return Result.Fail(ResultCode.CHUNK_NOT_FOUND); } if(taskCache.Parts.Count < taskCache.TotalPartCount) { return Result.Fail(ResultCode.CHUNK_COMBINE_FAIL); } var task = await reposity.FindByIdAsync(Guid.Parse(taskCache.TaskId)); //var res = await storage.CompleteUploadAsync(new CompleteUploadCommand( // ProviderCode: taskCache.ProviderCode, // Bucket: taskCache.Bucket, // Region: taskCache.Region, // ObjectKey: taskCache.ObjectKey, // UploadSessionId: taskCache.UploadSessionId, // Parts: command.Parts // ), cancellationToken); task.CompleteUpload(new Domain.ValueObjects.StorageLocation( taskCache.ProviderCode, taskCache.Bucket, taskCache.ObjectKey, taskCache.Region )); await endpoint.Publish(new UploadTaskCompleteEvent() { OperatorId = command.userId, Bucket = taskCache.Bucket, FileName = task.FileName.ToString(), ObjectKey = taskCache.ObjectKey, Parts = command.Parts.Select(s => new IM.Commons.IntegrationEvents.UploadPart( s.PartNumber, s.ETag, s.Size, s.Checksum) ).ToList(), ProviderCode = taskCache.ProviderCode, Region = taskCache.Region, SessionId = command.UploadSessionId, TaskId = task.Id, FileSize = task.FileSize, ContentType = task.ContentType.ToString(), CheckSun = task.CheckSum.Value }, cancellationToken); return Result.Success(mapper.Map(task)); } public async Task> UploadPartAsync(UploadPartCommand command) { var taskCache = await redis.GetAsync(command.SessionId); if(taskCache is null) { return Result.Fail(ResultCode.CHUNK_NOT_FOUND); } await localChunkStorage.SavePartAsync(new SaveLocalPartCommand( UploadSessionId: command.SessionId, PartNumber: command.PartNum, Stream: command.Stream, ContentLength: command.ContentLength )); taskCache.AddOrUpdatePart(new StorageContracts.UploadPart(command.PartNum, command.PartNum.ToString(), command.ContentLength)); await redis.SetAsync(taskCache); var location = new Domain.ValueObjects.StorageLocation( storageProvider: taskCache.ProviderCode, bucket: taskCache.Bucket, objectKey: taskCache.ObjectKey, region: taskCache.Region ); return Result.Success(new CompleteUploadResult(location, command.PartNum.ToString(), command.ContentLength)); } } }