131 lines
10 KiB
C#

using System.Collections.Concurrent;
using MongoDB.Driver;
namespace GuestOps.Web;
// All application reads/writes require a server-derived hotel ID. Global access is
// confined to authentication and the mailbox worker, not exposed as query parameters.
public interface IStore
{
Task Initialize();
Task<List<T>> List<T>(string hotel) where T : TenantDocument;
Task<T?> Get<T>(string hotel, string id) where T : TenantDocument;
Task Insert<T>(T document) where T : TenantDocument;
Task<bool> Replace<T>(string hotel, string id, long version, T document) where T : TenantDocument;
Task Delete<T>(string hotel, string id) where T : TenantDocument;
Task<StaffUser?> FindLogin(string email);
Task<OAuthRequest?> ConsumeOAuth(string id, string hotel, string user);
Task<List<Mailbox>> Mailboxes();
Task SaveMailbox(Mailbox mailbox);
Task SaveSync(Mailbox mailbox);
Task<bool> TryLease(string id, string owner);
Task ReleaseLease(string id, string owner);
Task Import(Conversation message);
Task<List<Conversation>> Deliveries();
}
public sealed class MongoStore : IStore
{
public Task<List<Conversation>> Deliveries() => Collection<Conversation>().Find(x => x.Delivery != null && (x.Delivery.State == "Pending" || x.Delivery.State == "Sending")).SortBy(x => x.Delivery!.UpdatedAt).Limit(100).ToListAsync();
private readonly IMongoDatabase db;
public MongoStore(IConfiguration config)
{
var connection = config["Mongo:ConnectionString"] ?? throw new InvalidOperationException("Mongo:ConnectionString is required.");
var settings = MongoClientSettings.FromConnectionString(connection);
settings.ServerSelectionTimeout = TimeSpan.FromSeconds(8);
db = new MongoClient(settings).GetDatabase(config["Mongo:Database"] ?? "guestops");
}
IMongoCollection<T> Collection<T>() => db.GetCollection<T>(typeof(T).Name.ToLowerInvariant());
FilterDefinition<T> Scope<T>(string hotel) where T : TenantDocument
{
if (string.IsNullOrWhiteSpace(hotel)) throw new InvalidOperationException("Hotel scope required.");
return Builders<T>.Filter.Eq(x => x.HotelId, hotel);
}
public async Task Initialize()
{
await Collection<StaffUser>().Indexes.CreateOneAsync(new CreateIndexModel<StaffUser>(Builders<StaffUser>.IndexKeys.Ascending(x => x.Email), new() { Unique = true }));
await Collection<Conversation>().Indexes.CreateOneAsync(new CreateIndexModel<Conversation>(Builders<Conversation>.IndexKeys.Ascending(x => x.HotelId).Ascending(x => x.MailboxId).Ascending(x => x.ProviderMessageId), new() { Unique = true }));
await Collection<Conversation>().Indexes.CreateOneAsync(new CreateIndexModel<Conversation>(Builders<Conversation>.IndexKeys.Ascending(x => x.HotelId).Descending(x => x.ReceivedAt)));
await Collection<Conversation>().Indexes.CreateOneAsync(new CreateIndexModel<Conversation>(Builders<Conversation>.IndexKeys.Ascending("Delivery.State").Ascending("Delivery.UpdatedAt")));
await Collection<Mailbox>().Indexes.CreateOneAsync(new CreateIndexModel<Mailbox>(Builders<Mailbox>.IndexKeys.Ascending(x => x.Email), new() { Unique = true }));
await Collection<OAuthRequest>().Indexes.CreateOneAsync(new CreateIndexModel<OAuthRequest>(Builders<OAuthRequest>.IndexKeys.Ascending(x => x.ExpiresAt), new() { ExpireAfter = TimeSpan.Zero }));
}
public Task<List<T>> List<T>(string hotel) where T : TenantDocument
{
var query = Collection<T>().Find(Scope<T>(hotel));
if (typeof(T) == typeof(Conversation)) query = query.Sort(Builders<T>.Sort.Descending("ReceivedAt"));
if (typeof(T) == typeof(Activity)) query = query.Sort(Builders<T>.Sort.Descending("At"));
return query.Limit(500).ToListAsync();
}
public async Task<T?> Get<T>(string hotel, string id) where T : TenantDocument => await Collection<T>().Find(Scope<T>(hotel) & Builders<T>.Filter.Eq(x => x.Id, id)).FirstOrDefaultAsync();
public Task Insert<T>(T document) where T : TenantDocument
{
_ = Scope<T>(document.HotelId);
return Collection<T>().InsertOneAsync(document);
}
public async Task<bool> Replace<T>(string hotel, string id, long version, T document) where T : TenantDocument
{
if (document.HotelId != hotel || document.Id != id) throw new InvalidOperationException("Invalid document scope.");
var result = await Collection<T>().ReplaceOneAsync(Scope<T>(hotel) & Builders<T>.Filter.Eq(x => x.Id, id) & Builders<T>.Filter.Eq("Version", version), document);
return result.ModifiedCount == 1;
}
public async Task Delete<T>(string hotel, string id) where T : TenantDocument => await Collection<T>().DeleteOneAsync(Scope<T>(hotel) & Builders<T>.Filter.Eq(x => x.Id, id));
public async Task<StaffUser?> FindLogin(string email) => await Collection<StaffUser>().Find(x => x.Email == email).FirstOrDefaultAsync();
public async Task<OAuthRequest?> ConsumeOAuth(string id, string hotel, string user) => await Collection<OAuthRequest>().FindOneAndDeleteAsync(x => x.Id == id && x.HotelId == hotel && x.UserId == user && x.ExpiresAt > DateTime.UtcNow);
public Task<List<Mailbox>> Mailboxes() => Collection<Mailbox>().Find(x => x.Status == "Connected").ToListAsync();
public async Task SaveMailbox(Mailbox mailbox) => await Collection<Mailbox>().ReplaceOneAsync(x => x.HotelId == mailbox.HotelId && x.Id == mailbox.Id, mailbox, new ReplaceOptions { IsUpsert = true });
public async Task SaveSync(Mailbox mailbox) => await Collection<Mailbox>().UpdateOneAsync(
x => x.HotelId == mailbox.HotelId && x.Id == mailbox.Id && x.ProtectedRefreshToken == mailbox.ProtectedRefreshToken,
Builders<Mailbox>.Update.Set(x => x.LastSyncAt, mailbox.LastSyncAt).Set(x => x.SyncError, mailbox.SyncError)
.Set(x => x.PageToken, mailbox.PageToken).Set(x => x.WindowStart, mailbox.WindowStart).Set(x => x.WindowEnd, mailbox.WindowEnd));
public async Task<bool> TryLease(string id, string owner)
{
try
{
var result = await Collection<WorkerLease>().FindOneAndUpdateAsync(x => x.Id == id && x.Until < DateTime.UtcNow,
Builders<WorkerLease>.Update.Set(x => x.Owner, owner).Set(x => x.Until, DateTime.UtcNow.AddMinutes(5)), new() { IsUpsert = true, ReturnDocument = ReturnDocument.After });
return result.Owner == owner;
}
catch (MongoCommandException ex) when (ex.Code == 11000) { return false; }
catch (MongoWriteException ex) when (ex.WriteError.Category == ServerErrorCategory.DuplicateKey) { return false; }
}
public async Task ReleaseLease(string id, string owner) => await Collection<WorkerLease>().DeleteOneAsync(x => x.Id == id && x.Owner == owner);
public async Task Import(Conversation message)
{
try { await Insert(message); }
catch (MongoWriteException ex) when (ex.WriteError.Category == ServerErrorCategory.DuplicateKey) { /* already durable */ }
}
}
// Explicit Development-only preview store. Production never falls back to this.
public sealed class PreviewStore : IStore
{
public Task<List<Conversation>> Deliveries() => Task.FromResult(rows.Where(x => x.Key.StartsWith("Conversation:")).Select(x => Clone<Conversation>(x.Value)).Where(x => x.Delivery?.State is "Pending" or "Sending").ToList());
private readonly ConcurrentDictionary<string, string> rows = new();
private readonly object gate = new();
static string Key<T>(string id) => typeof(T).Name + ":" + id;
static T Clone<T>(string text) => System.Text.Json.JsonSerializer.Deserialize<T>(text)!;
static string Json<T>(T value) => System.Text.Json.JsonSerializer.Serialize(value);
public Task Initialize() => Task.CompletedTask;
public Task<List<T>> List<T>(string hotel) where T : TenantDocument => Task.FromResult(rows.Where(x => x.Key.StartsWith(typeof(T).Name + ":")).Select(x => Clone<T>(x.Value)).Where(x => x.HotelId == hotel).ToList());
public async Task<T?> Get<T>(string hotel, string id) where T : TenantDocument => (await List<T>(hotel)).SingleOrDefault(x => x.Id == id);
public Task Insert<T>(T document) where T : TenantDocument { if (!rows.TryAdd(Key<T>(document.Id), Json(document))) throw new InvalidOperationException("Duplicate document"); return Task.CompletedTask; }
public Task<bool> Replace<T>(string hotel, string id, long version, T document) where T : TenantDocument
{
lock (gate)
{
if (!rows.TryGetValue(Key<T>(id), out var raw)) return Task.FromResult(false);
var old = Clone<T>(raw);
if (old.HotelId != hotel || document.HotelId != hotel || document.Id != id || (long)typeof(T).GetProperty("Version")!.GetValue(old)! != version) return Task.FromResult(false);
rows[Key<T>(id)] = Json(document); return Task.FromResult(true);
}
}
public async Task Delete<T>(string hotel, string id) where T : TenantDocument { if (await Get<T>(hotel, id) != null) rows.TryRemove(Key<T>(id), out _); }
public Task<StaffUser?> FindLogin(string email) => Task.FromResult(rows.Where(x => x.Key.StartsWith("StaffUser:")).Select(x => Clone<StaffUser>(x.Value)).SingleOrDefault(x => x.Email == email));
public Task<OAuthRequest?> ConsumeOAuth(string id, string hotel, string user) { lock(gate) { var x = rows.TryGetValue(Key<OAuthRequest>(id), out var raw) ? Clone<OAuthRequest>(raw) : null; if(x?.HotelId != hotel || x.UserId != user || x.ExpiresAt <= DateTime.UtcNow) return Task.FromResult<OAuthRequest?>(null); rows.TryRemove(Key<OAuthRequest>(id),out _); return Task.FromResult<OAuthRequest?>(x); } }
public Task<List<Mailbox>> Mailboxes() => Task.FromResult(new List<Mailbox>());
public Task SaveMailbox(Mailbox mailbox) => throw new InvalidOperationException("Real mailbox connections are unavailable in preview mode.");
public Task SaveSync(Mailbox mailbox) => throw new InvalidOperationException("Real mailbox connections are unavailable in preview mode.");
public Task<bool> TryLease(string id, string owner) => Task.FromResult(false);
public Task ReleaseLease(string id, string owner) => Task.CompletedTask;
public async Task Import(Conversation message) { if (!(await List<Conversation>(message.HotelId)).Any(x => x.MailboxId == message.MailboxId && x.ProviderMessageId == message.ProviderMessageId)) await Insert(message); }
}