diff --git a/Notesnook.API/Controllers/S3Controller.cs b/Notesnook.API/Controllers/S3Controller.cs
index f72e45d..6c4a9e0 100644
--- a/Notesnook.API/Controllers/S3Controller.cs
+++ b/Notesnook.API/Controllers/S3Controller.cs
@@ -20,10 +20,12 @@ along with this program. If not, see .
using System;
using System.Net.Http;
using System.Security.Claims;
+using System.Threading;
using System.Threading.Tasks;
using Microsoft.AspNetCore.Authorization;
using Microsoft.AspNetCore.Http.Extensions;
using Microsoft.AspNetCore.Mvc;
+using Microsoft.AspNetCore.RateLimiting;
using Microsoft.Extensions.Logging;
using MongoDB.Driver;
using Notesnook.API.Helpers;
@@ -32,7 +34,6 @@ using Notesnook.API.Models;
using Streetwriters.Common;
using Streetwriters.Common.Accessors;
using Streetwriters.Common.Extensions;
-using Streetwriters.Common.Models;
namespace Notesnook.API.Controllers
{
@@ -40,39 +41,45 @@ namespace Notesnook.API.Controllers
[Route("s3")]
[ResponseCache(NoStore = true, Location = ResponseCacheLocation.None)]
[Authorize("Sync")]
- public class S3Controller(IS3Service s3Service, ISyncItemsRepositoryAccessor repositories, WampServiceAccessor serviceAccessor, ILogger logger) : ControllerBase
+ public class S3Controller(IS3Service s3Service, ISyncItemsRepositoryAccessor repositories, WampServiceAccessor serviceAccessor, IHttpClientFactory httpClientFactory, ILogger logger) : ControllerBase
{
[HttpPut]
+ [EnableRateLimiting("s3-direct")]
public async Task Upload([FromQuery] string name)
{
- return BadRequest(new { error = "Attachment storage is temporarily unavailable. Please try again later." });
- // try
- // {
- // var userId = this.User.GetUserId();
+ try
+ {
+ var userId = this.User.GetUserId();
- // var fileSize = HttpContext.Request.ContentLength ?? 0;
- // bool hasBody = fileSize > 0;
+ var fileSize = HttpContext.Request.ContentLength ?? 0;
+ bool hasBody = fileSize > 0;
- // if (!hasBody)
- // {
- // return Ok(Request.GetEncodedUrl() + "&access_token=" + Request.Headers.Authorization.ToString().Replace("Bearer ", ""));
- // }
+ if (!hasBody)
+ {
+ return Ok(Request.GetEncodedUrl() + "&access_token=" + Request.Headers.Authorization.ToString().Replace("Bearer ", ""));
+ }
- // if (Constants.IS_SELF_HOSTED) await UploadFileAsync(userId, name, fileSize);
- // else await UploadFileWithChecksAsync(userId, name, fileSize);
+ if (Constants.IS_SELF_HOSTED) await UploadFileAsync(userId, name, fileSize);
+ else await UploadFileWithChecksAsync(userId, name, fileSize);
- // return Ok();
- // }
- // catch (Exception ex)
- // {
- // logger.LogError(ex, "Error uploading attachment for user.");
- // return BadRequest(new { error = "Failed to upload attachment." });
- // }
+ return Ok();
+ }
+ catch (Exception ex) when (S3TransientErrorClassifier.IsTransient(ex))
+ {
+ return StorageUnavailable(ex);
+ }
+ catch (Exception ex)
+ {
+ logger.LogError(ex, "Error uploading attachment for user.");
+ return BadRequest(new { error = "Failed to upload attachment." });
+ }
}
private async Task UploadFileWithChecksAsync(string userId, string name, long fileSize)
{
- var userSettings = await repositories.UsersSettings.FindOneAsync((u) => u.UserId == userId);
+ var userSettings = await repositories.UsersSettings.FindOneAsync((u) => u.UserId == userId)
+ ?? throw new Exception("User settings not found.");
+ var observedStorageLimit = userSettings.StorageLimit;
var subscription = await serviceAccessor.UserSubscriptionService.GetUserSubscriptionAsync(Clients.Notesnook.Id, userId) ?? throw new Exception("User subscription not found.");
@@ -85,49 +92,68 @@ namespace Notesnook.API.Controllers
var uploadedFileSize = await UploadFileAsync(userId, name, fileSize);
- userSettings.StorageLimit.Value += uploadedFileSize;
- await repositories.UsersSettings.Collection.UpdateOneAsync(
- Builders.Filter.Eq(u => u.UserId, userId),
- Builders.Update.Set(u => u.StorageLimit, userSettings.StorageLimit)
- );
-
- // extra check in case user sets wrong ContentLength in the HTTP header
- if (uploadedFileSize != fileSize && StorageHelper.IsStorageLimitReached(subscription, userSettings.StorageLimit.Value))
+ try
{
- await s3Service.DeleteObjectAsync(userId, name);
- throw new Exception("Storage limit exceeded.");
+ await s3Service.IncrementStorageUsageAsync(userId, uploadedFileSize, observedStorageLimit);
}
+ catch (Exception ex)
+ {
+ logger.LogError(ex, "Failed to account for attachment usage for user {UserId}.", userId);
+ }
+
}
private async Task UploadFileAsync(string userId, string name, long fileSize)
{
var url = await s3Service.GetInternalUploadObjectUrlAsync(userId, name) ?? throw new Exception("Could not create signed url.");
- var httpClient = new HttpClient();
+ var httpClient = httpClientFactory.CreateClient("S3Upload");
var content = new StreamContent(HttpContext.Request.BodyReader.AsStream());
content.Headers.ContentLength = fileSize;
- var response = await httpClient.SendRequestAsync(url, null, HttpMethod.Put, content);
- if (!response.Success) throw new Exception(response.Content != null ? await response.Content.ReadAsStringAsync() : "Could not upload file.");
+ using var request = new HttpRequestMessage(HttpMethod.Put, url)
+ {
+ Content = content
+ };
+ using var uploadTimeout = CancellationTokenSource.CreateLinkedTokenSource(HttpContext.RequestAborted);
+ uploadTimeout.CancelAfter(TimeSpan.FromMinutes(15 + (15 * fileSize / (1024d * 1024d * 1024d))));
+ using var response = await httpClient.SendAsync(
+ request,
+ HttpCompletionOption.ResponseHeadersRead,
+ uploadTimeout.Token);
+ if (!response.IsSuccessStatusCode)
+ {
+ var statusCode = (int)response.StatusCode;
+ if (S3TransientErrorClassifier.IsTransient(statusCode))
+ throw new S3StorageUnavailableException("Attachment storage is temporarily unavailable.");
- return await s3Service.GetObjectSizeAsync(userId, name);
+ throw new Exception(await response.Content.ReadAsStringAsync(uploadTimeout.Token));
+ }
+
+ // The PUT response confirms this request. A follow-up HEAD can lag or
+ // be throttled, so direct accounting uses the accepted request bytes.
+ return fileSize;
}
[HttpGet("multipart")]
+ [EnableRateLimiting("s3-multipart-control")]
public async Task MultipartUpload([FromQuery] string name, [FromQuery] int parts, [FromQuery] string? uploadId)
{
- return BadRequest(new { error = "Attachment storage is temporarily unavailable. Please try again later." });
- // var userId = this.User.GetUserId();
- // try
- // {
- // var meta = await s3Service.StartMultipartUploadAsync(userId, name, parts, uploadId);
- // return Ok(meta);
- // }
- // catch (Exception ex)
- // {
- // logger.LogError(ex, "Error starting multipart upload for user.");
- // return BadRequest(new { error = "Failed to start multipart upload." });
- // }
+ var userId = this.User.GetUserId();
+ try
+ {
+ var meta = await s3Service.StartMultipartUploadAsync(userId, name, parts, uploadId);
+ return Ok(meta);
+ }
+ catch (Exception ex) when (S3TransientErrorClassifier.IsTransient(ex))
+ {
+ return StorageUnavailable(ex);
+ }
+ catch (Exception ex)
+ {
+ logger.LogError(ex, "Error starting multipart upload for user.");
+ return BadRequest(new { error = "Failed to start multipart upload." });
+ }
}
[HttpDelete("multipart")]
@@ -139,6 +165,10 @@ namespace Notesnook.API.Controllers
await s3Service.AbortMultipartUploadAsync(userId, name, uploadId);
return Ok();
}
+ catch (Exception ex) when (S3TransientErrorClassifier.IsTransient(ex))
+ {
+ return StorageUnavailable(ex);
+ }
catch (Exception ex)
{
logger.LogError(ex, "Error aborting multipart upload for user.");
@@ -147,6 +177,7 @@ namespace Notesnook.API.Controllers
}
[HttpPost("multipart")]
+ [EnableRateLimiting("s3-multipart-control")]
public async Task CompleteMultipartUpload([FromBody] CompleteMultipartUploadRequestWrapper uploadRequestWrapper)
{
var userId = this.User.GetUserId();
@@ -155,6 +186,10 @@ namespace Notesnook.API.Controllers
await s3Service.CompleteMultipartUploadAsync(userId, uploadRequestWrapper.ToRequest());
return Ok();
}
+ catch (Exception ex) when (S3TransientErrorClassifier.IsTransient(ex))
+ {
+ return StorageUnavailable(ex);
+ }
catch (Exception ex)
{
logger.LogError(ex, "Error completing multipart upload for user.");
@@ -162,6 +197,12 @@ namespace Notesnook.API.Controllers
}
}
+ private IActionResult StorageUnavailable(Exception exception)
+ {
+ logger.LogWarning(exception, "Attachment storage is temporarily unavailable.");
+ return StatusCode(503, new { error = "Attachment storage is temporarily unavailable. Please try again later." });
+ }
+
[HttpGet]
public async Task Download([FromQuery] string name)
{
diff --git a/Notesnook.API/Helpers/S3FailoverHelper.cs b/Notesnook.API/Helpers/S3FailoverHelper.cs
index 9ba4b5c..4eaea39 100644
--- a/Notesnook.API/Helpers/S3FailoverHelper.cs
+++ b/Notesnook.API/Helpers/S3FailoverHelper.cs
@@ -20,6 +20,8 @@ along with this program. If not, see .
using System;
using System.Collections.Generic;
using System.Linq;
+using System.Net.Http;
+using System.Net.Sockets;
using System.Threading.Tasks;
using Amazon.S3;
using Amazon.S3.Model;
@@ -47,12 +49,6 @@ namespace Notesnook.API.Helpers
///
public bool UseExponentialBackoff { get; set; } = true;
- ///
- /// Whether to allow failover for write operations (PUT, POST, DELETE).
- /// Default is false to prevent data consistency issues.
- ///
- public bool AllowWriteFailover { get; set; } = false;
-
///
/// List of exception types that should trigger failover
///
@@ -77,6 +73,39 @@ namespace Notesnook.API.Helpers
};
}
+ public sealed class S3StorageUnavailableException : Exception
+ {
+ public S3StorageUnavailableException(string message, Exception? innerException = null)
+ : base(message, innerException)
+ {
+ }
+ }
+
+ public static class S3TransientErrorClassifier
+ {
+ public static bool IsTransient(int statusCode)
+ {
+ return statusCode == 408 || statusCode == 429 || statusCode >= 500;
+ }
+
+ public static bool IsTransient(Exception exception)
+ {
+ if (exception is S3StorageUnavailableException)
+ return true;
+
+ if (exception is AmazonS3Exception s3Exception)
+ {
+ return IsTransient((int)s3Exception.StatusCode)
+ || s3Exception.ErrorCode is "SlowDown" or "ServiceUnavailable" or "InternalError" or "RequestTimeout";
+ }
+
+ if (exception is HttpRequestException or SocketException or TimeoutException or TaskCanceledException)
+ return true;
+
+ return exception.InnerException != null && IsTransient(exception.InnerException);
+ }
+ }
+
///
/// Result of a failover operation
///
@@ -161,10 +190,11 @@ namespace Notesnook.API.Helpers
var result = new S3FailoverResult();
Exception? lastException = null;
- // Determine max clients to try based on write operation flag
- var maxClientsToTry = (isWriteOperation && !config.AllowWriteFailover) ? 1 : clients.Count;
+ // Writes stay on the provider that owns the object or multipart upload.
+ // They get one application attempt and never fail over.
+ var maxClientsToTry = isWriteOperation ? 1 : clients.Count;
- if (isWriteOperation && !config.AllowWriteFailover && clients.Count > 1)
+ if (isWriteOperation && clients.Count > 1)
{
logger?.LogDebug(
"Write operation {Operation} will only use primary endpoint. Failover is disabled for write operations.",
@@ -185,7 +215,12 @@ namespace Notesnook.API.Helpers
operationName, clientName, i + 1, maxClientsToTry);
}
- var (success, value, exception, attempts) = await TryExecuteAsync(client, operation, operationName, clientName);
+ var (success, value, exception, attempts) = await TryExecuteAsync(
+ client,
+ operation,
+ operationName,
+ clientName,
+ isWriteOperation ? 0 : config.MaxRetries);
result.AttemptsUsed += attempts;
if (success && value != null)
@@ -222,19 +257,22 @@ namespace Notesnook.API.Helpers
operationName, maxClientsToTry, result.AttemptsUsed);
return result;
- } ///
- /// Try to execute an operation with retries
- ///
+ }
+
+ ///
+ /// Try to execute an operation with retries
+ ///
private async Task<(bool success, T? value, Exception? exception, int attempts)> TryExecuteAsync(
AmazonS3Client client,
Func> operation,
string operationName,
- string endpointName)
+ string endpointName,
+ int maxRetries)
{
Exception? lastException = null;
int attempts = 0;
- for (int retry = 0; retry <= config.MaxRetries; retry++)
+ for (int retry = 0; retry <= maxRetries; retry++)
{
attempts++;
try
@@ -246,12 +284,12 @@ namespace Notesnook.API.Helpers
{
lastException = ex;
- if (retry < config.MaxRetries && ShouldRetry(ex))
+ if (retry < maxRetries && ShouldRetry(ex))
{
var delay = CalculateRetryDelay(retry);
logger?.LogWarning(ex,
"Attempt {Attempt}/{MaxAttempts} failed for {Operation} on {Endpoint}. Retrying in {Delay}ms",
- retry + 1, config.MaxRetries + 1, operationName, endpointName, delay);
+ retry + 1, maxRetries + 1, operationName, endpointName, delay);
await Task.Delay(delay);
}
@@ -322,10 +360,10 @@ namespace Notesnook.API.Helpers
string operationName = "S3Operation",
bool isWriteOperation = false)
{
- await ExecuteWithFailoverAsync