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 WorkerLastSeen(); Task RecordWorkerHeartbeat(); Task> List(string hotel) where T : TenantDocument; Task ConversationPage(string hotel,DateTime? before,string beforeId,int limit); Task Get(string hotel, string id) where T : TenantDocument; Task Insert(T document) where T : TenantDocument; Task Replace(string hotel, string id, long version, T document) where T : TenantDocument; Task Delete(string hotel, string id) where T : TenantDocument; Task FindLogin(string email); Task FindAccountLink(string hash); Task TryInsertStaff(StaffUser user); Task ConsumeAccountLink(StaffUser user,long version,string hash); Task ConsumeOAuth(string id, string hotel, string user); Task> Mailboxes(); Task SaveMailbox(Mailbox mailbox); Task TryInsertMailbox(Mailbox mailbox); Task SaveSync(Mailbox mailbox); Task TryLease(string id, string owner); Task ReleaseLease(string id, string owner); Task Import(Conversation message); Task> Deliveries(); Task> AccountMails(); Task ClaimAccountMail(string id, string owner); Task TryInsertPmsChange(PmsChange change); Task TryInsertPayment(PaymentRequest payment); Task TryAutoReplyClaim(AutoReplyClaim claim); Task> AutoReplyCandidates(string hotel,string mailbox,DateTime since); } public sealed class MongoStore : IStore { public async Task Ping()=>await db.RunCommandAsync(new MongoDB.Bson.BsonDocument("ping",1)); public async Task WorkerLastSeen()=>(await db.GetCollection("workerheartbeat").Find(x=>x.Id=="worker").FirstOrDefaultAsync())?.At; public async Task RecordWorkerHeartbeat()=>await db.GetCollection("workerheartbeat").ReplaceOneAsync(x=>x.Id=="worker",new WorkerHeartbeat(),new ReplaceOptions{IsUpsert=true}); public Task> Deliveries() => Collection().Find(x => x.Delivery != null && (x.Delivery.State == "Pending" || x.Delivery.State == "Sending")).SortBy(x => x.Delivery!.UpdatedAt).Limit(100).ToListAsync(); public Task> AccountMails() => Collection().Find(x => (x.State == "Pending" || x.State == "Sending") && x.NextAttemptAt <= DateTime.UtcNow).SortBy(x => x.NextAttemptAt).Limit(100).ToListAsync(); public async Task ClaimAccountMail(string id,string owner)=>await Collection().FindOneAndUpdateAsync( x=>x.Id==id&&(x.State=="Pending"||(x.State=="Sending"&&x.LeaseUntil.Update.Set(x=>x.State,"Sending").Set(x=>x.LeaseOwner,owner).Set(x=>x.LeaseUntil,DateTime.UtcNow.AddMinutes(2)).Inc(x=>x.Version,1), new(){ReturnDocument=ReturnDocument.After}); 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 Collection() => db.GetCollection(typeof(T).Name.ToLowerInvariant()); FilterDefinition Scope(string hotel) where T : TenantDocument { if (string.IsNullOrWhiteSpace(hotel)) throw new InvalidOperationException("Hotel scope required."); return Builders.Filter.Eq(x => x.HotelId, hotel); } public async Task Initialize() { await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x=>x.AccountLinkHash))); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x=>x.State).Ascending(x=>x.NextAttemptAt))); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x=>x.MessageId),new(){Unique=true})); foreach(var field in new[]{"ThreadKey","RecipientDay","DaySlot"})await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x=>x.HotelId).Ascending(field),new(){Unique=true})); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x=>x.HotelId).Ascending(x=>x.MailboxId).Ascending(x=>x.AutoReplyCheckedAt).Ascending(x=>x.ReceivedAt))); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x=>x.HotelId).Ascending(x=>x.Reference),new(){Unique=true})); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x=>x.HotelId).Descending(x=>x.UpdatedAt))); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x=>x.HotelId).Ascending(x=>x.ReservationId),new CreateIndexOptions{Unique=true,PartialFilterExpression=Builders.Filter.In(x=>x.State,new[]{"Review","Applying","NeedsReview"})})); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x=>x.HotelId).Descending(x=>x.UpdatedAt))); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x=>x.FetchedAt),new(){ExpireAfter=TimeSpan.FromDays(1)})); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x => x.Email), new() { Unique = true })); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x => x.HotelId).Ascending(x => x.MailboxId).Ascending(x => x.ProviderMessageId), new() { Unique = true })); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x => x.HotelId).Descending(x => x.ReceivedAt))); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending("Delivery.State").Ascending("Delivery.UpdatedAt"))); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x => x.Email), new() { Unique = true })); await Collection().Indexes.CreateOneAsync(new CreateIndexModel(Builders.IndexKeys.Ascending(x => x.ExpiresAt), new() { ExpireAfter = TimeSpan.Zero })); } public Task> List(string hotel) where T : TenantDocument { var query = Collection().Find(Scope(hotel)); if (typeof(T) == typeof(Conversation)) query = query.Sort(Builders.Sort.Descending("ReceivedAt")); if (typeof(T) == typeof(Activity)) query = query.Sort(Builders.Sort.Descending("At")); if (typeof(T) == typeof(PmsChange) || typeof(T) == typeof(PaymentRequest)) query = query.Sort(Builders.Sort.Descending("UpdatedAt")); return query.Limit(500).ToListAsync(); } public async Task ConversationPage(string hotel,DateTime? before,string beforeId,int limit) { var filter=Scope(hotel);if(before!=null)filter&=Builders.Filter.Lt(x=>x.ReceivedAt,before.Value)|(Builders.Filter.Eq(x=>x.ReceivedAt,before.Value)&Builders.Filter.Lt(x=>x.Id,beforeId)); var rows=await Collection().Find(filter).SortByDescending(x=>x.ReceivedAt).ThenByDescending(x=>x.Id).Limit(limit+1).ToListAsync(); var more=rows.Count>limit;if(more)rows.RemoveAt(rows.Count-1);return new(rows,more?ConversationPaging.Encode(rows[^1]):null); } public async Task Get(string hotel, string id) where T : TenantDocument => await Collection().Find(Scope(hotel) & Builders.Filter.Eq(x => x.Id, id)).FirstOrDefaultAsync(); public Task Insert(T document) where T : TenantDocument { _ = Scope(document.HotelId); return Collection().InsertOneAsync(document); } public async Task Replace(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.Filter.Eq("Version",version); if((typeof(T)==typeof(StaffUser)||typeof(T)==typeof(Mailbox))&&version==0)versionFilter|=Builders.Filter.Exists("Version",false); var result = await Collection().ReplaceOneAsync(Scope(hotel) & Builders.Filter.Eq(x => x.Id, id) & versionFilter, document); return result.ModifiedCount == 1; } public async Task Delete(string hotel, string id) where T : TenantDocument => await Collection().DeleteOneAsync(Scope(hotel) & Builders.Filter.Eq(x => x.Id, id)); public async Task FindAccountLink(string hash)=>await Collection().Find(x=>x.AccountLinkHash==hash&&x.AccountLinkExpiresAt>DateTime.UtcNow).FirstOrDefaultAsync(); public async Task TryInsertStaff(StaffUser user){try{await Insert(user);return true;}catch(MongoWriteException ex) when(ex.WriteError.Category==ServerErrorCategory.DuplicateKey){return false;}} public async Task ConsumeAccountLink(StaffUser user,long version,string hash) { var filter=Scope(user.HotelId)&Builders.Filter.Eq(x=>x.Id,user.Id)&Builders.Filter.Eq(x=>x.Version,version)&Builders.Filter.Eq(x=>x.AccountLinkHash,hash)&Builders.Filter.Gt(x=>x.AccountLinkExpiresAt,DateTime.UtcNow); return (await Collection().ReplaceOneAsync(filter,user)).ModifiedCount==1; } public async Task FindLogin(string email) => await Collection().Find(x => x.Email == email).FirstOrDefaultAsync(); public async Task ConsumeOAuth(string id, string hotel, string user) => await Collection().FindOneAndDeleteAsync(x => x.Id == id && x.HotelId == hotel && x.UserId == user && x.ExpiresAt > DateTime.UtcNow); public Task> Mailboxes() => Collection().Find(x => x.Status == "Connected").ToListAsync(); public async Task SaveMailbox(Mailbox mailbox) => await Collection().ReplaceOneAsync(x => x.HotelId == mailbox.HotelId && x.Id == mailbox.Id, mailbox, new ReplaceOptions { IsUpsert = true }); public async Task 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().UpdateOneAsync( Builders.Filter.Eq(x=>x.HotelId,mailbox.HotelId)&Builders.Filter.Eq(x=>x.Id,mailbox.Id)&Builders.Filter.Eq(x=>x.ProtectedRefreshToken,mailbox.ProtectedRefreshToken)&Builders.Filter.Eq(x=>x.Status,"Connected")&(mailbox.Version==0?(Builders.Filter.Eq(x=>x.Version,0)|Builders.Filter.Exists("Version",false)):Builders.Filter.Eq(x=>x.Version,mailbox.Version)), Builders.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 TryLease(string id, string owner) { try { var result = await Collection().FindOneAndUpdateAsync(x => x.Id == id && x.Until < DateTime.UtcNow, Builders.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().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> AutoReplyCandidates(string hotel,string mailbox,DateTime since)=>Collection().Find(Scope(hotel)&Builders.Filter.Eq(x=>x.MailboxId,mailbox)&Builders.Filter.Eq(x=>x.AutoReplyCheckedAt,null)&Builders.Filter.Gte(x=>x.ReceivedAt,since)).SortBy(x=>x.ReceivedAt).Limit(100).ToListAsync(); public async Task TryAutoReplyClaim(AutoReplyClaim claim) { try{await Insert(claim);return true;}catch(MongoWriteException ex) when(ex.WriteError.Category==ServerErrorCategory.DuplicateKey){return false;} } public async Task TryInsertPayment(PaymentRequest payment) { try { await Insert(payment);return true; } catch(MongoWriteException ex) when(ex.WriteError.Category==ServerErrorCategory.DuplicateKey){return false;} } public async Task 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 WorkerLastSeen()=>Task.FromResult(null); public Task RecordWorkerHeartbeat()=>Task.CompletedTask; public async Task> AutoReplyCandidates(string hotel,string mailbox,DateTime since)=>(await List(hotel)).Where(x=>x.MailboxId==mailbox&&x.AutoReplyCheckedAt==null&&x.ReceivedAt>=since).OrderBy(x=>x.ReceivedAt).Take(100).ToList(); public Task TryAutoReplyClaim(AutoReplyClaim claim) { lock(gate){if(rows.Where(x=>x.Key.StartsWith("AutoReplyClaim:")).Select(x=>Clone(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(claim.Id),Json(claim)));} } public Task TryInsertPayment(PaymentRequest payment) { lock(gate) { if(rows.Where(x=>x.Key.StartsWith("PaymentRequest:")).Select(x=>Clone(x.Value)).Any(x=>x.HotelId==payment.HotelId&&x.Reference==payment.Reference))return Task.FromResult(false); return Task.FromResult(rows.TryAdd(Key(payment.Id),Json(payment))); } } public Task TryInsertPmsChange(PmsChange change) { lock(gate) { if(rows.Where(x=>x.Key.StartsWith("PmsChange:")).Select(x=>Clone(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(change.Id),Json(change))); } } public Task> Deliveries() => Task.FromResult(rows.Where(x => x.Key.StartsWith("Conversation:")).Select(x => Clone(x.Value)).Where(x => x.Delivery?.State is "Pending" or "Sending").ToList()); public Task> AccountMails()=>Task.FromResult(rows.Where(x=>x.Key.StartsWith("AccountMail:")).Select(x=>Clone(x.Value)).Where(x=>x.State is "Pending" or "Sending"&&x.NextAttemptAt<=DateTime.UtcNow).ToList()); public Task ClaimAccountMail(string id,string owner) { lock(gate) { if(!rows.TryGetValue(Key(id),out var raw))return Task.FromResult(null); var mail=Clone(raw);if(mail.NextAttemptAt>DateTime.UtcNow||mail.State!="Pending"&&!(mail.State=="Sending"&&mail.LeaseUntil(null); mail.State="Sending";mail.LeaseOwner=owner;mail.LeaseUntil=DateTime.UtcNow.AddMinutes(2);mail.Version++;rows[Key(id)]=Json(mail);return Task.FromResult(mail); } } private readonly ConcurrentDictionary rows = new(); private readonly object gate = new(); static string Key(string id) => typeof(T).Name + ":" + id; static T Clone(string text) => System.Text.Json.JsonSerializer.Deserialize(text)!; static string Json(T value) => System.Text.Json.JsonSerializer.Serialize(value); public Task Initialize() => Task.CompletedTask; public Task> List(string hotel) where T : TenantDocument => Task.FromResult(rows.Where(x => x.Key.StartsWith(typeof(T).Name + ":")).Select(x => Clone(x.Value)).Where(x => x.HotelId == hotel).ToList()); public async Task ConversationPage(string hotel,DateTime? before,string beforeId,int limit) { var query=(await List(hotel)).OrderByDescending(x=>x.ReceivedAt).ThenByDescending(x=>x.Id).Where(x=>before==null||x.ReceivedAtlimit;if(more)query.RemoveAt(query.Count-1);return new(query,more?ConversationPaging.Encode(query[^1]):null); } public async Task Get(string hotel, string id) where T : TenantDocument => (await List(hotel)).SingleOrDefault(x => x.Id == id); public Task Insert(T document) where T : TenantDocument { if (!rows.TryAdd(Key(document.Id), Json(document))) throw new InvalidOperationException("Duplicate document"); return Task.CompletedTask; } public Task Replace(string hotel, string id, long version, T document) where T : TenantDocument { lock (gate) { if (!rows.TryGetValue(Key(id), out var raw)) return Task.FromResult(false); var old = Clone(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(id)] = Json(document); return Task.FromResult(true); } } public async Task Delete(string hotel, string id) where T : TenantDocument { if (await Get(hotel, id) != null) rows.TryRemove(Key(id), out _); } public Task FindAccountLink(string hash)=>Task.FromResult(rows.Where(x=>x.Key.StartsWith("StaffUser:")).Select(x=>Clone(x.Value)).SingleOrDefault(x=>x.AccountLinkHash==hash&&x.AccountLinkExpiresAt>DateTime.UtcNow)); public Task TryInsertStaff(StaffUser user){lock(gate){if(rows.Where(x=>x.Key.StartsWith("StaffUser:")).Select(x=>Clone(x.Value)).Any(x=>x.Email==user.Email))return Task.FromResult(false);return Task.FromResult(rows.TryAdd(Key(user.Id),Json(user)));}} public Task ConsumeAccountLink(StaffUser user,long version,string hash) { lock(gate){if(!rows.TryGetValue(Key(user.Id),out var raw))return Task.FromResult(false);var old=Clone(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(user.Id)]=Json(user);return Task.FromResult(true);} } public Task FindLogin(string email) => Task.FromResult(rows.Where(x => x.Key.StartsWith("StaffUser:")).Select(x => Clone(x.Value)).SingleOrDefault(x => x.Email == email)); public Task ConsumeOAuth(string id, string hotel, string user) { lock(gate) { var x = rows.TryGetValue(Key(id), out var raw) ? Clone(raw) : null; if(x?.HotelId != hotel || x.UserId != user || x.ExpiresAt <= DateTime.UtcNow) return Task.FromResult(null); rows.TryRemove(Key(id),out _); return Task.FromResult(x); } } public Task> Mailboxes() => Task.FromResult(new List()); public Task SaveMailbox(Mailbox mailbox) => throw new InvalidOperationException("Real mailbox connections are unavailable in preview mode."); public Task TryInsertMailbox(Mailbox mailbox){lock(gate){if(rows.Where(x=>x.Key.StartsWith("Mailbox:")).Select(x=>Clone(x.Value)).Any(x=>x.Email==mailbox.Email))return Task.FromResult(false);return Task.FromResult(rows.TryAdd(Key(mailbox.Id),Json(mailbox)));}} public Task SaveSync(Mailbox mailbox){lock(gate){if(rows.TryGetValue(Key(mailbox.Id),out var raw)){var old=Clone(raw);if(old.HotelId==mailbox.HotelId&&old.Version==mailbox.Version&&old.Status=="Connected"&&old.ProtectedRefreshToken==mailbox.ProtectedRefreshToken)rows[Key(mailbox.Id)]=Json(mailbox);}return Task.CompletedTask;}} public Task 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(message.HotelId)).Any(x => x.MailboxId == message.MailboxId && x.ProviderMessageId == message.ProviderMessageId)) await Insert(message); } }