fix: align backend APIs and upload flow
This commit is contained in:
@@ -1,173 +1,70 @@
|
||||
using FileService.Application.Ports;
|
||||
using FileService.Application.Ports;
|
||||
using FileService.Application.StorageContracts;
|
||||
using FileService.Domain.ValueObjects;
|
||||
using IM.Commons;
|
||||
using IM.InitCommon;
|
||||
using Microsoft.Extensions.Options;
|
||||
using StackExchange.Redis;
|
||||
|
||||
namespace FileService.Infrastructure.Storage
|
||||
namespace FileService.Infrastructure.Storage;
|
||||
public class LocalStorageAdapter(IStorageRedisCache redis, IOptionsSnapshot<StorageOptions> options) : IObjectStoragePort, ILocalChunkStorage
|
||||
{
|
||||
public class LocalStorageAdapter(IStorageRedisCache redis, IOptions<StorageOptions> options) : IObjectStoragePort, ILocalChunkStorage
|
||||
private StorageProviderOptions Provider => options.Value.Providers["Local"];
|
||||
public string ProviderCode => "Local";
|
||||
public static string SafePath(string root, params string[] segments)
|
||||
{
|
||||
private readonly IStorageRedisCache redis = redis;
|
||||
private readonly IOptions<StorageOptions> options = options;
|
||||
private readonly StorageProviderOptions providerOptions = options.Value.Providers[options.Value.DefaultProviderCode];
|
||||
|
||||
public string ProviderCode => "Local";
|
||||
|
||||
/// <summary>
|
||||
/// 单次直传:直接写入 LocalRootPath/{bucket}/{objectKey}。
|
||||
/// bucket 传公开桶名即落到公开目录,可被静态托管直链访问。
|
||||
/// </summary>
|
||||
public async Task<StorageLocation> PutObjectAsync(PutObjectCommand command, CancellationToken token)
|
||||
{
|
||||
var fullPath = Path.Combine(providerOptions.LocalRootPath!, command.Bucket, command.ObjectKey);
|
||||
Directory.CreateDirectory(Path.GetDirectoryName(fullPath)!);
|
||||
|
||||
await using (var fs = new FileStream(fullPath, FileMode.Create))
|
||||
{
|
||||
await command.Content.CopyToAsync(fs, token);
|
||||
await fs.FlushAsync(token);
|
||||
var fullRoot = Path.GetFullPath(root).TrimEnd(Path.DirectorySeparatorChar) + Path.DirectorySeparatorChar;
|
||||
var path = Path.GetFullPath(Path.Combine(new[] { fullRoot }.Concat(segments).ToArray()));
|
||||
if (!path.StartsWith(fullRoot, OperatingSystem.IsWindows() ? StringComparison.OrdinalIgnoreCase : StringComparison.Ordinal)) throw new InvalidOperationException("存储路径超出允许根目录");
|
||||
for (var dir = Path.GetDirectoryName(path); dir != null && dir.Length >= fullRoot.Length; dir = Path.GetDirectoryName(dir))
|
||||
if (Directory.Exists(dir) && File.GetAttributes(dir).HasFlag(FileAttributes.ReparsePoint)) throw new InvalidOperationException("不允许使用存储目录链接");
|
||||
if (File.Exists(path) && File.GetAttributes(path).HasFlag(FileAttributes.ReparsePoint)) throw new InvalidOperationException("不允许使用文件链接");
|
||||
return path;
|
||||
}
|
||||
public async Task<StorageLocation> PutObjectAsync(PutObjectCommand command, CancellationToken token)
|
||||
{
|
||||
var path = SafePath(Provider.LocalRootPath!, command.Bucket, command.ObjectKey);
|
||||
Directory.CreateDirectory(Path.GetDirectoryName(path)!);
|
||||
await using var stream = new FileStream(path, FileMode.CreateNew);
|
||||
await command.Content.CopyToAsync(stream, token);
|
||||
if (stream.Length != command.ContentLength) throw new InvalidOperationException("文件实际大小与申报大小不符");
|
||||
return new(ProviderCode, command.Bucket, command.ObjectKey, Provider.Region);
|
||||
}
|
||||
public string? GetPublicUrl(StorageLocation location) => !string.IsNullOrEmpty(Provider.PublicBucket) && location.Bucket == Provider.PublicBucket
|
||||
? $"{(Provider.PublicBaseUrl ?? Provider.LocalUploadApiBaseUrl ?? "").TrimEnd('/')}/static/{string.Join('/', location.ObjectKey.Replace('\\', '/').Split('/').Select(Uri.EscapeDataString))}" : null;
|
||||
public Task<Stream> OpenReadAsync(StorageLocation location, CancellationToken token) => Task.FromResult<Stream>(File.OpenRead(SafePath(Provider.LocalRootPath!, location.Bucket, location.ObjectKey)));
|
||||
public Task<InitiateUploadResult> InitUploadAsync(InitiateUploadCommand command, CancellationToken token) => Task.FromResult(new InitiateUploadResult(Guid.NewGuid().ToString(), new StorageLocation(ProviderCode, command.Bucket, command.ObjectKey, Provider.Region)));
|
||||
public Task<PresignedUrl> GenerateUploadUrlAsync(GenerateUploadUrlCommand command, CancellationToken token) => Task.FromResult(new PresignedUrl(
|
||||
Provider.LocalUploadApiBaseUrl!.TrimEnd('/') + $"/local/parts/upload?sessionId={Uri.EscapeDataString(command.UploadSessionId)}&partNumber={command.PartNumber}", "POST", new Dictionary<string, string>(), DateTimeOffset.UtcNow.Add(command.ExpiresIn)));
|
||||
public async Task SavePartAsync(SaveLocalPartCommand command)
|
||||
{
|
||||
if (!Guid.TryParse(command.UploadSessionId, out _) || command.PartNumber < 1) throw new InvalidOperationException("分片参数无效");
|
||||
var path = SafePath(Provider.LocalRootPath!, "staging", command.UploadSessionId, $"{command.PartNumber}.part");
|
||||
Directory.CreateDirectory(Path.GetDirectoryName(path)!);
|
||||
await using var stream = File.Create(path);
|
||||
await command.Stream.CopyToAsync(stream);
|
||||
if (stream.Length != command.ContentLength) throw new InvalidOperationException("分片实际大小不符");
|
||||
}
|
||||
public async Task<CompleteUploadResult> CompleteUploadAsync(CompleteUploadCommand command, CancellationToken token)
|
||||
{
|
||||
var cache = await redis.GetAsync(command.UploadSessionId) ?? throw new InvalidOperationException("上传任务已过期");
|
||||
var final = SafePath(Provider.LocalRootPath!, command.Bucket, command.ObjectKey);
|
||||
Directory.CreateDirectory(Path.GetDirectoryName(final)!);
|
||||
var temporary = final + ".merging";
|
||||
await using (var output = File.Create(temporary)) {
|
||||
foreach (var part in command.Parts.OrderBy(x => x.PartNumber)) {
|
||||
await using var input = File.OpenRead(SafePath(Provider.LocalRootPath!, "staging", command.UploadSessionId, $"{part.PartNumber}.part"));
|
||||
await input.CopyToAsync(output, token);
|
||||
}
|
||||
|
||||
return new StorageLocation(
|
||||
storageProvider: ProviderCode,
|
||||
bucket: command.Bucket,
|
||||
objectKey: command.ObjectKey,
|
||||
region: providerOptions.Region);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// 公开桶文件返回静态托管直链;私有文件返回 null。
|
||||
/// </summary>
|
||||
public string? GetPublicUrl(StorageLocation location)
|
||||
{
|
||||
if (string.IsNullOrEmpty(providerOptions.PublicBucket) ||
|
||||
!string.Equals(location.Bucket, providerOptions.PublicBucket, StringComparison.OrdinalIgnoreCase))
|
||||
{
|
||||
return null;
|
||||
}
|
||||
|
||||
var baseUrl = (providerOptions.PublicBaseUrl ?? providerOptions.LocalUploadApiBaseUrl ?? string.Empty)
|
||||
.TrimEnd('/');
|
||||
var key = location.ObjectKey.Replace('\\', '/').TrimStart('/');
|
||||
return $"{baseUrl}/static/{key}";
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// 打开本地文件读取流:LocalRootPath/{bucket}/{objectKey}。
|
||||
/// </summary>
|
||||
public Task<Stream> OpenReadAsync(StorageLocation location, CancellationToken token)
|
||||
{
|
||||
var fullPath = Path.Combine(providerOptions.LocalRootPath!, location.Bucket, location.ObjectKey);
|
||||
if (!File.Exists(fullPath))
|
||||
{
|
||||
throw new FileNotFoundException(fullPath);
|
||||
}
|
||||
|
||||
Stream stream = new FileStream(fullPath, FileMode.Open, FileAccess.Read, FileShare.Read);
|
||||
return Task.FromResult(stream);
|
||||
}
|
||||
|
||||
public async Task<CompleteUploadResult> CompleteUploadAsync(CompleteUploadCommand command, CancellationToken token)
|
||||
{
|
||||
var res = await MergeAsync(command.UploadSessionId, command.ObjectKey, command.Parts);
|
||||
return new CompleteUploadResult(new StorageLocation(
|
||||
storageProvider: command.ProviderCode,
|
||||
bucket: command.Bucket,
|
||||
objectKey: command.ObjectKey,
|
||||
region: command.Region
|
||||
), null, command.Parts.Sum(x => x.Size).Value);
|
||||
}
|
||||
public async Task<Result<object>> MergeAsync(string sessionId, string objectKey, IReadOnlyList<UploadPart> parts)
|
||||
{
|
||||
var rootPath = options.Value.Providers[options.Value.DefaultProviderCode].LocalRootPath;
|
||||
var tempPath = Path.Combine(rootPath, sessionId, "parts"); // 项目根目录下 uploads // 最终文件存储路径(这里可以用你之前 ObjectNameGenerator 生成的名字)
|
||||
var finalPath = Path.Combine(rootPath, objectKey);
|
||||
var finalDir = Path.GetDirectoryName(finalPath);
|
||||
Directory.CreateDirectory(finalDir);
|
||||
|
||||
var storageCache = await redis.GetAsync(sessionId);
|
||||
var totalChunks = storageCache.TotalPartCount;
|
||||
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");
|
||||
if (!File.Exists(chunkPath))
|
||||
return Result.Fail(ResultCode.CHUNK_NOT_FOUND);
|
||||
using (var chunkStream = new FileStream(chunkPath, FileMode.Open))
|
||||
{
|
||||
await chunkStream.CopyToAsync(finalStream);
|
||||
}
|
||||
}
|
||||
Directory.Delete(tempPath, true);
|
||||
await redis.DeleteAsync(sessionId);
|
||||
}
|
||||
|
||||
return Result.Success();
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
//_logger.LogError(e, e.Message);
|
||||
throw;
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
public async Task<PresignedUrl> GenerateUploadUrlAsync(GenerateUploadUrlCommand command, CancellationToken token)
|
||||
{
|
||||
var baseUrl = options.Value.Providers[options.Value.DefaultProviderCode].LocalUploadApiBaseUrl;
|
||||
return new PresignedUrl(
|
||||
baseUrl + $"local/parts/upload?sessionId={command.UploadSessionId}&partNumber={command.PartNumber}",
|
||||
new Dictionary<string, string>(),
|
||||
ExpiresAt: DateTimeOffset.Now.Add(options.Value.Providers[options.Value.DefaultProviderCode].UploadUrlExpiresIn)
|
||||
);
|
||||
}
|
||||
|
||||
public async Task<InitiateUploadResult> InitUploadAsync(InitiateUploadCommand command, CancellationToken token)
|
||||
{
|
||||
var sessionId = Guid.NewGuid();
|
||||
var location = new StorageLocation();
|
||||
return new InitiateUploadResult(sessionId.ToString(),location);
|
||||
}
|
||||
|
||||
public async Task SavePartAsync(SaveLocalPartCommand command)
|
||||
{
|
||||
var path = BuildPartPath(
|
||||
command.UploadSessionId,
|
||||
command.PartNumber);
|
||||
|
||||
Directory.CreateDirectory(
|
||||
Path.GetDirectoryName(path)!);
|
||||
|
||||
await using var fs = File.Create(path);
|
||||
|
||||
await command.Stream.CopyToAsync(fs);
|
||||
|
||||
await fs.FlushAsync();
|
||||
}
|
||||
private string BuildPartPath(
|
||||
string uploadSessionId,
|
||||
int partNumber)
|
||||
{
|
||||
return Path.Combine(
|
||||
providerOptions.LocalRootPath,
|
||||
uploadSessionId,
|
||||
"parts",
|
||||
$"{partNumber}.part");
|
||||
if (output.Length != cache.FileSize) throw new InvalidOperationException("合并文件大小不符");
|
||||
}
|
||||
File.Move(temporary, final, true);
|
||||
// Keep parts available for idempotent retry after response loss.
|
||||
return new(new StorageLocation(ProviderCode, command.Bucket, command.ObjectKey, command.Region), null, cache.FileSize);
|
||||
}
|
||||
public async Task<Result<object>> MergeAsync(string sessionId, string objectKey, IReadOnlyList<UploadPart> parts)
|
||||
{
|
||||
var cache = await redis.GetAsync(sessionId) ?? throw new InvalidOperationException("上传任务已过期");
|
||||
await CompleteUploadAsync(new CompleteUploadCommand(ProviderCode: cache.ProviderCode, Bucket: cache.Bucket, Region: cache.Region, ObjectKey: objectKey, UploadSessionId: sessionId, Parts: parts), CancellationToken.None);
|
||||
return Result.Success();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,22 +1,26 @@
|
||||
using FileService.Application.Ports;
|
||||
using FileService.Application.Ports;
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using System.Text;
|
||||
using System.Threading.Tasks;
|
||||
using IM.InitCommon;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
namespace FileService.Infrastructure.Storage
|
||||
{
|
||||
public class ObjectStorageRouter : IObjectStorageRouter
|
||||
public class ObjectStorageRouter : IObjectStorageRouter, IDisposable
|
||||
{
|
||||
private IReadOnlyDictionary<string, IObjectStoragePort> adpters;
|
||||
|
||||
public ObjectStorageRouter(IEnumerable<IObjectStoragePort> storages)
|
||||
public ObjectStorageRouter(IEnumerable<IObjectStoragePort> storages, IOptionsSnapshot<StorageOptions> options)
|
||||
{
|
||||
this.adpters = storages.ToDictionary(x => x.ProviderCode, StringComparer.OrdinalIgnoreCase);
|
||||
var adapters = storages.ToDictionary(x => x.ProviderCode, StringComparer.OrdinalIgnoreCase);
|
||||
foreach (var (code, provider) in options.Value.Providers.Where(x => x.Value.ProviderType is StorageProviderType.AwsS3 or StorageProviderType.Minio)) adapters[code] = new S3StorageAdapter(code, provider);
|
||||
this.adpters = adapters;
|
||||
}
|
||||
|
||||
public IObjectStoragePort Route(string providerCode)
|
||||
public void Dispose() { foreach (var adapter in adpters.Values.OfType<S3StorageAdapter>()) adapter.Dispose(); } public IObjectStoragePort Route(string providerCode)
|
||||
{
|
||||
return this.adpters[providerCode];
|
||||
}
|
||||
|
||||
@@ -0,0 +1,53 @@
|
||||
using Amazon.S3;
|
||||
using Amazon.S3.Model;
|
||||
using FileService.Application.Ports;
|
||||
using FileService.Application.StorageContracts;
|
||||
using FileService.Domain.ValueObjects;
|
||||
using IM.InitCommon;
|
||||
|
||||
namespace FileService.Infrastructure.Storage;
|
||||
public sealed class S3StorageAdapter(string code, StorageProviderOptions options) : IObjectStoragePort, IDisposable
|
||||
{
|
||||
private readonly AmazonS3Client client = new(options.AccessKeyId, options.AccessKeySecret, new AmazonS3Config { ServiceURL = options.Endpoint, AuthenticationRegion = string.IsNullOrWhiteSpace(options.Region) ? "us-east-1" : options.Region, ForcePathStyle = true });
|
||||
public string ProviderCode => code;
|
||||
public async Task<UploadPart> WritePartAsync(UploadRuntimeCache task, int partNumber, Stream content, long size, CancellationToken ct) {
|
||||
var part = await client.UploadPartAsync(new UploadPartRequest { BucketName = task.Bucket, Key = task.ObjectKey, UploadId = task.UploadSessionId, PartNumber = partNumber, InputStream = content, PartSize = size }, ct);
|
||||
return new(partNumber, part.ETag, size);
|
||||
}
|
||||
public async Task<InitiateUploadResult> InitUploadAsync(InitiateUploadCommand command, CancellationToken token) {
|
||||
var r = await client.InitiateMultipartUploadAsync(new InitiateMultipartUploadRequest { BucketName = command.Bucket, Key = command.ObjectKey, ContentType = command.ContentType }, token);
|
||||
return new(r.UploadId, new StorageLocation(code, command.Bucket, command.ObjectKey, options.Region));
|
||||
}
|
||||
public Task<PresignedUrl> GenerateUploadUrlAsync(GenerateUploadUrlCommand command, CancellationToken token) {
|
||||
var expires = DateTimeOffset.UtcNow.Add(command.ExpiresIn);
|
||||
var request = new GetPreSignedUrlRequest { BucketName = command.Bucket, Key = command.ObjectKey, Verb = HttpVerb.PUT, Expires = expires.UtcDateTime, UploadId = command.UploadSessionId, PartNumber = command.PartNumber ?? 1 };
|
||||
return Task.FromResult(new PresignedUrl(client.GetPreSignedURL(request), "PUT", new Dictionary<string, string>(), expires));
|
||||
}
|
||||
public async Task<CompleteUploadResult> CompleteUploadAsync(CompleteUploadCommand command, CancellationToken token) {
|
||||
var response = await client.CompleteMultipartUploadAsync(new CompleteMultipartUploadRequest { BucketName = command.Bucket, Key = command.ObjectKey, UploadId = command.UploadSessionId, PartETags = command.Parts.OrderBy(x => x.PartNumber).Select(x => new PartETag(x.PartNumber, x.ETag)).ToList() }, token);
|
||||
var meta = await client.GetObjectMetadataAsync(command.Bucket, command.ObjectKey, token);
|
||||
return new(new StorageLocation(code, command.Bucket, command.ObjectKey, options.Region), response.ETag, meta.ContentLength, VersionId: response.VersionId);
|
||||
}
|
||||
public async Task<StorageLocation> PutObjectAsync(PutObjectCommand command, CancellationToken token) {
|
||||
await client.PutObjectAsync(new PutObjectRequest { BucketName = command.Bucket, Key = command.ObjectKey, ContentType = command.ContentType, InputStream = command.Content, AutoCloseStream = false }, token);
|
||||
return new(code, command.Bucket, command.ObjectKey, options.Region);
|
||||
}
|
||||
// Public rendering also passes through FileService. Private access is always authorized there.
|
||||
public string? GetPublicUrl(StorageLocation location) => null;
|
||||
public async Task<Stream> OpenReadAsync(StorageLocation location, CancellationToken token) { var response = await client.GetObjectAsync(location.Bucket, location.ObjectKey, token); return new ResponseStream(response); }
|
||||
public async Task Test(CancellationToken ct) {
|
||||
var key = "im-admin-connectivity/" + Guid.NewGuid().ToString("N");
|
||||
try { await client.PutObjectAsync(new PutObjectRequest { BucketName = options.Bucket, Key = key, ContentBody = "IM connectivity test" }, ct); using var read = await client.GetObjectAsync(options.Bucket, key, ct); }
|
||||
finally { await client.DeleteObjectAsync(options.Bucket, key, CancellationToken.None); }
|
||||
}
|
||||
public void Dispose() => client.Dispose();
|
||||
sealed class ResponseStream(GetObjectResponse response) : Stream {
|
||||
readonly Stream inner = response.ResponseStream;
|
||||
public override bool CanRead => inner.CanRead; public override bool CanSeek => inner.CanSeek; public override bool CanWrite => false;
|
||||
public override long Length => response.ContentLength; public override long Position { get => inner.Position; set => inner.Position = value; }
|
||||
public override void Flush() => inner.Flush(); public override int Read(byte[] b, int o, int c) => inner.Read(b, o, c);
|
||||
public override ValueTask<int> ReadAsync(Memory<byte> b, CancellationToken ct = default) => inner.ReadAsync(b, ct);
|
||||
public override long Seek(long o, SeekOrigin origin) => inner.Seek(o, origin); public override void SetLength(long v) => throw new NotSupportedException(); public override void Write(byte[] b, int o, int c) => throw new NotSupportedException();
|
||||
protected override void Dispose(bool disposing) { if (disposing) response.Dispose(); base.Dispose(disposing); }
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user