211 lines
19 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 Ping();
Task<DateTime?> WorkerLastSeen();
Task RecordWorkerHeartbeat();
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<StaffUser?> FindAccountLink(string hash);
Task<bool> TryInsertStaff(StaffUser user);
Task<bool> ConsumeAccountLink(StaffUser user,long version,string hash);
Task<OAuthRequest?> ConsumeOAuth(string id, string hotel, string user);
Task<List<Mailbox>> Mailboxes();
Task SaveMailbox(Mailbox mailbox);
Task<bool> TryInsertMailbox(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();
Task<bool> TryInsertPmsChange(PmsChange change);
Task<bool> TryInsertPayment(PaymentRequest payment);
Task<bool> TryAutoReplyClaim(AutoReplyClaim claim);
Task<List<Conversation>> AutoReplyCandidates(string hotel,string mailbox,DateTime since);
}
public sealed class MongoStore : IStore
{
public async Task Ping()=>await db.RunCommandAsync<MongoDB.Bson.BsonDocument>(new MongoDB.Bson.BsonDocument("ping",1));
public async Task<DateTime?> WorkerLastSeen()=>(await db.GetCollection<WorkerHeartbeat>("workerheartbeat").Find(x=>x.Id=="worker").FirstOrDefaultAsync())?.At;
public async Task RecordWorkerHeartbeat()=>await db.GetCollection<WorkerHeartbeat>("workerheartbeat").ReplaceOneAsync(x=>x.Id=="worker",new WorkerHeartbeat(),new ReplaceOptions{IsUpsert=true});
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.AccountLinkHash)));
foreach(var field in new[]{"ThreadKey","RecipientDay","DaySlot"})await Collection<AutoReplyClaim>().Indexes.CreateOneAsync(new CreateIndexModel<AutoReplyClaim>(Builders<AutoReplyClaim>.IndexKeys.Ascending(x=>x.HotelId).Ascending(field),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.AutoReplyCheckedAt).Ascending(x=>x.ReceivedAt)));
await Collection<PaymentRequest>().Indexes.CreateOneAsync(new CreateIndexModel<PaymentRequest>(Builders<PaymentRequest>.IndexKeys.Ascending(x=>x.HotelId).Ascending(x=>x.Reference),new(){Unique=true}));
await Collection<PaymentRequest>().Indexes.CreateOneAsync(new CreateIndexModel<PaymentRequest>(Builders<PaymentRequest>.IndexKeys.Ascending(x=>x.HotelId).Descending(x=>x.UpdatedAt)));
await Collection<PmsChange>().Indexes.CreateOneAsync(new CreateIndexModel<PmsChange>(Builders<PmsChange>.IndexKeys.Ascending(x=>x.HotelId).Ascending(x=>x.ReservationId),new CreateIndexOptions<PmsChange>{Unique=true,PartialFilterExpression=Builders<PmsChange>.Filter.In(x=>x.State,new[]{"Review","Applying","NeedsReview"})}));
await Collection<PmsChange>().Indexes.CreateOneAsync(new CreateIndexModel<PmsChange>(Builders<PmsChange>.IndexKeys.Ascending(x=>x.HotelId).Descending(x=>x.UpdatedAt)));
await Collection<PmsSnapshot>().Indexes.CreateOneAsync(new CreateIndexModel<PmsSnapshot>(Builders<PmsSnapshot>.IndexKeys.Ascending(x=>x.FetchedAt),new(){ExpireAfter=TimeSpan.FromDays(1)}));
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"));
if (typeof(T) == typeof(PmsChange) || typeof(T) == typeof(PaymentRequest)) query = query.Sort(Builders<T>.Sort.Descending("UpdatedAt"));
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 versionFilter=Builders<T>.Filter.Eq("Version",version);
if((typeof(T)==typeof(StaffUser)||typeof(T)==typeof(Mailbox))&&version==0)versionFilter|=Builders<T>.Filter.Exists("Version",false);
var result = await Collection<T>().ReplaceOneAsync(Scope<T>(hotel) & Builders<T>.Filter.Eq(x => x.Id, id) & versionFilter, 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?> FindAccountLink(string hash)=>await Collection<StaffUser>().Find(x=>x.AccountLinkHash==hash&&x.AccountLinkExpiresAt>DateTime.UtcNow).FirstOrDefaultAsync();
public async Task<bool> TryInsertStaff(StaffUser user){try{await Insert(user);return true;}catch(MongoWriteException ex) when(ex.WriteError.Category==ServerErrorCategory.DuplicateKey){return false;}}
public async Task<bool> ConsumeAccountLink(StaffUser user,long version,string hash)
{
var filter=Scope<StaffUser>(user.HotelId)&Builders<StaffUser>.Filter.Eq(x=>x.Id,user.Id)&Builders<StaffUser>.Filter.Eq(x=>x.Version,version)&Builders<StaffUser>.Filter.Eq(x=>x.AccountLinkHash,hash)&Builders<StaffUser>.Filter.Gt(x=>x.AccountLinkExpiresAt,DateTime.UtcNow);
return (await Collection<StaffUser>().ReplaceOneAsync(filter,user)).ModifiedCount==1;
}
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<bool> TryInsertMailbox(Mailbox mailbox){try{await Insert(mailbox);return true;}catch(MongoWriteException ex) when(ex.WriteError.Category==ServerErrorCategory.DuplicateKey){return false;}}
public async Task SaveSync(Mailbox mailbox) => await Collection<Mailbox>().UpdateOneAsync(
Builders<Mailbox>.Filter.Eq(x=>x.HotelId,mailbox.HotelId)&Builders<Mailbox>.Filter.Eq(x=>x.Id,mailbox.Id)&Builders<Mailbox>.Filter.Eq(x=>x.ProtectedRefreshToken,mailbox.ProtectedRefreshToken)&Builders<Mailbox>.Filter.Eq(x=>x.Status,"Connected")&(mailbox.Version==0?(Builders<Mailbox>.Filter.Eq(x=>x.Version,0)|Builders<Mailbox>.Filter.Exists("Version",false)):Builders<Mailbox>.Filter.Eq(x=>x.Version,mailbox.Version)),
Builders<Mailbox>.Update.Set(x => x.LastSyncAt, mailbox.LastSyncAt).Set(x => x.SyncError, mailbox.SyncError)
.Set(x=>x.Status,mailbox.Status).Set(x=>x.LastAttemptAt,mailbox.LastAttemptAt).Set(x=>x.NextAttemptAt,mailbox.NextAttemptAt).Set(x=>x.FailureCount,mailbox.FailureCount).Set(x=>x.SyncErrorCode,mailbox.SyncErrorCode)
.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 */ }
}
public Task<List<Conversation>> AutoReplyCandidates(string hotel,string mailbox,DateTime since)=>Collection<Conversation>().Find(Scope<Conversation>(hotel)&Builders<Conversation>.Filter.Eq(x=>x.MailboxId,mailbox)&Builders<Conversation>.Filter.Eq(x=>x.AutoReplyCheckedAt,null)&Builders<Conversation>.Filter.Gte(x=>x.ReceivedAt,since)).SortBy(x=>x.ReceivedAt).Limit(100).ToListAsync();
public async Task<bool> TryAutoReplyClaim(AutoReplyClaim claim)
{
try{await Insert(claim);return true;}catch(MongoWriteException ex) when(ex.WriteError.Category==ServerErrorCategory.DuplicateKey){return false;}
}
public async Task<bool> TryInsertPayment(PaymentRequest payment)
{
try { await Insert(payment);return true; }
catch(MongoWriteException ex) when(ex.WriteError.Category==ServerErrorCategory.DuplicateKey){return false;}
}
public async Task<bool> TryInsertPmsChange(PmsChange change)
{
try { await Insert(change);return true; }
catch(MongoWriteException ex) when(ex.WriteError.Category==ServerErrorCategory.DuplicateKey){return false;}
}
}
// Explicit Development-only preview store. Production never falls back to this.
public sealed class PreviewStore : IStore
{
public Task Ping()=>Task.CompletedTask;
public Task<DateTime?> WorkerLastSeen()=>Task.FromResult<DateTime?>(null);
public Task RecordWorkerHeartbeat()=>Task.CompletedTask;
public async Task<List<Conversation>> AutoReplyCandidates(string hotel,string mailbox,DateTime since)=>(await List<Conversation>(hotel)).Where(x=>x.MailboxId==mailbox&&x.AutoReplyCheckedAt==null&&x.ReceivedAt>=since).OrderBy(x=>x.ReceivedAt).Take(100).ToList();
public Task<bool> TryAutoReplyClaim(AutoReplyClaim claim)
{
lock(gate){if(rows.Where(x=>x.Key.StartsWith("AutoReplyClaim:")).Select(x=>Clone<AutoReplyClaim>(x.Value)).Any(x=>x.HotelId==claim.HotelId&&(x.ThreadKey==claim.ThreadKey||x.RecipientDay==claim.RecipientDay||x.DaySlot==claim.DaySlot)))return Task.FromResult(false);return Task.FromResult(rows.TryAdd(Key<AutoReplyClaim>(claim.Id),Json(claim)));}
}
public Task<bool> TryInsertPayment(PaymentRequest payment)
{
lock(gate)
{
if(rows.Where(x=>x.Key.StartsWith("PaymentRequest:")).Select(x=>Clone<PaymentRequest>(x.Value)).Any(x=>x.HotelId==payment.HotelId&&x.Reference==payment.Reference))return Task.FromResult(false);
return Task.FromResult(rows.TryAdd(Key<PaymentRequest>(payment.Id),Json(payment)));
}
}
public Task<bool> TryInsertPmsChange(PmsChange change)
{
lock(gate)
{
if(rows.Where(x=>x.Key.StartsWith("PmsChange:")).Select(x=>Clone<PmsChange>(x.Value)).Any(x=>x.HotelId==change.HotelId&&x.ReservationId==change.ReservationId&&PmsChange.Active(x.State)))return Task.FromResult(false);
return Task.FromResult(rows.TryAdd(Key<PmsChange>(change.Id),Json(change)));
}
}
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?> FindAccountLink(string hash)=>Task.FromResult(rows.Where(x=>x.Key.StartsWith("StaffUser:")).Select(x=>Clone<StaffUser>(x.Value)).SingleOrDefault(x=>x.AccountLinkHash==hash&&x.AccountLinkExpiresAt>DateTime.UtcNow));
public Task<bool> TryInsertStaff(StaffUser user){lock(gate){if(rows.Where(x=>x.Key.StartsWith("StaffUser:")).Select(x=>Clone<StaffUser>(x.Value)).Any(x=>x.Email==user.Email))return Task.FromResult(false);return Task.FromResult(rows.TryAdd(Key<StaffUser>(user.Id),Json(user)));}}
public Task<bool> ConsumeAccountLink(StaffUser user,long version,string hash)
{
lock(gate){if(!rows.TryGetValue(Key<StaffUser>(user.Id),out var raw))return Task.FromResult(false);var old=Clone<StaffUser>(raw);if(old.HotelId!=user.HotelId||old.Version!=version||old.AccountLinkHash!=hash||old.AccountLinkExpiresAt<=DateTime.UtcNow||old.AccountLinkExpiresAt==null)return Task.FromResult(false);rows[Key<StaffUser>(user.Id)]=Json(user);return Task.FromResult(true);}
}
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<bool> TryInsertMailbox(Mailbox mailbox){lock(gate){if(rows.Where(x=>x.Key.StartsWith("Mailbox:")).Select(x=>Clone<Mailbox>(x.Value)).Any(x=>x.Email==mailbox.Email))return Task.FromResult(false);return Task.FromResult(rows.TryAdd(Key<Mailbox>(mailbox.Id),Json(mailbox)));}}
public Task SaveSync(Mailbox mailbox){lock(gate){if(rows.TryGetValue(Key<Mailbox>(mailbox.Id),out var raw)){var old=Clone<Mailbox>(raw);if(old.HotelId==mailbox.HotelId&&old.Version==mailbox.Version&&old.Status=="Connected"&&old.ProtectedRefreshToken==mailbox.ProtectedRefreshToken)rows[Key<Mailbox>(mailbox.Id)]=Json(mailbox);}return Task.CompletedTask;}}
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); }
}