using System.Net;
using Amazon.S3;
using Amazon.S3.Model;
namespace GB5Shared.Storage
{
///
/// backed by Amazon S3. Constructed per-tenant by
/// with that tenant's own bucket/credentials
/// (from MDBLEVELSETTING) — never registered as a DI singleton, since bucket/client differ
/// per tenant DB.
///
public sealed class S3StorageProvider : IStorageProvider
{
private readonly IAmazonS3 _client;
private readonly string _bucketName;
public StorageProviderType ProviderType => StorageProviderType.S3;
public S3StorageProvider(IAmazonS3 client, string bucketName)
{
_client = client;
_bucketName = bucketName;
}
public async Task SaveAsync(Stream content, string storagePath, string mimeType, CancellationToken ct = default)
{
var request = new PutObjectRequest
{
BucketName = _bucketName,
Key = storagePath,
InputStream = content,
ContentType = mimeType,
AutoCloseStream = false
};
await _client.PutObjectAsync(request, ct).ConfigureAwait(false);
return storagePath;
}
public async Task GetStreamAsync(string storagePath, CancellationToken ct = default)
{
try
{
var response = await _client.GetObjectAsync(_bucketName, storagePath, ct).ConfigureAwait(false);
return new S3ObjectStream(response);
}
catch (AmazonS3Exception ex) when (ex.StatusCode == HttpStatusCode.NotFound)
{
throw new FileNotFoundException($"File not found at storage path: {storagePath}", storagePath, ex);
}
}
public async Task DeleteAsync(string storagePath, CancellationToken ct = default)
{
if (!await ExistsAsync(storagePath, ct).ConfigureAwait(false))
return false;
await _client.DeleteObjectAsync(_bucketName, storagePath, ct).ConfigureAwait(false);
return true;
}
public async Task ExistsAsync(string storagePath, CancellationToken ct = default)
{
try
{
await _client.GetObjectMetadataAsync(_bucketName, storagePath, ct).ConfigureAwait(false);
return true;
}
catch (AmazonS3Exception ex) when (ex.StatusCode == HttpStatusCode.NotFound)
{
return false;
}
}
public async Task MoveAsync(string oldStoragePath, string newStoragePath, CancellationToken ct = default)
{
if (!await ExistsAsync(oldStoragePath, ct).ConfigureAwait(false))
return false;
await _client.CopyObjectAsync(new CopyObjectRequest
{
SourceBucket = _bucketName,
SourceKey = oldStoragePath,
DestinationBucket = _bucketName,
DestinationKey = newStoragePath
}, ct).ConfigureAwait(false);
await _client.DeleteObjectAsync(_bucketName, oldStoragePath, ct).ConfigureAwait(false);
return true;
}
///
/// Wraps a so its ResponseStream stays valid until the
/// caller disposes the returned stream — the response and its stream share one lifetime
/// and must be disposed together, not separately.
///
private sealed class S3ObjectStream : Stream
{
private readonly GetObjectResponse _response;
private readonly Stream _inner;
public S3ObjectStream(GetObjectResponse response)
{
_response = response;
_inner = response.ResponseStream;
}
public override bool CanRead => true;
public override bool CanSeek => false;
public override bool CanWrite => false;
public override long Length => _response.ContentLength;
public override long Position
{
get => throw new NotSupportedException();
set => throw new NotSupportedException();
}
public override int Read(byte[] buffer, int offset, int count) => _inner.Read(buffer, offset, count);
public override Task ReadAsync(byte[] buffer, int offset, int count, CancellationToken ct) => _inner.ReadAsync(buffer, offset, count, ct);
public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException();
public override void SetLength(long value) => throw new NotSupportedException();
public override void Write(byte[] buffer, int offset, int count) => throw new NotSupportedException();
public override void Flush() { }
protected override void Dispose(bool disposing)
{
if (disposing)
{
_inner.Dispose();
_response.Dispose();
}
base.Dispose(disposing);
}
}
}
}