using System; using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; using adas_core.Application.Repositories.Interfaces; using adas_core.Domain.Models.AppSettings; using adas_core.Domain.Models.Pumps; using Microsoft.Extensions.Options; using MongoDB.Bson; using MongoDB.Driver; namespace adas_core.Infrastructure.Repositories { /// /// Repositorio de archivo para observaciones de bombas. /// Colección: archive_pumpobservations (configurable por ApiSettings.ArchivePumpObservations). /// public class PumpArchiveRepository : MongoRepository, IPumpArchiveRepository { private readonly ApiSettings _apiSettings; public PumpArchiveRepository(IOptions apiSettings, IMongoDatabase database) : base(database) { _apiSettings = apiSettings.Value; } public override string GetCollectionName() { // Nombre de colección pactado: "archive_pumpobservations" return _apiSettings.ArchivePatientsPumpobservations ?? "archive_pumpobservations"; } public override async Task CreateIndexes() { var indexModels = new List> { // Búsquedas por paciente (audit / restauraciones) new CreateIndexModel( Builders.IndexKeys.Ascending(x => x.PatientId), new CreateIndexOptions { Name = "ix_patientId" }), // Timeline por dispositivo (útil para auditorías por equipo) new CreateIndexModel( Builders.IndexKeys .Ascending(x => x.DeviceId) .Descending(x => x.Time), new CreateIndexOptions { Name = "ix_deviceId_time" }), // Orden temporal simple new CreateIndexModel( Builders.IndexKeys.Descending(x => x.Time), new CreateIndexOptions { Name = "ix_time" }) }; await Collection.Indexes.CreateManyAsync(indexModels); } /// /// Asynchronously inserts a into the underlying MongoDB collection. /// /// The pump observation to persist. public async Task InsertAsync(PumpObservation obs) { await Collection.InsertOneAsync(obs); } /// /// Asynchronously inserts a batch of pump observations into the underlying data store. If the collection is empty, the method completes without performing any insertion. /// /// The pump observations to insert into the collection. public async Task InsertManyAsync(IEnumerable observations) { var list = observations as IList ?? observations.ToList(); if (list.Count == 0) return; await Collection.InsertManyAsync(list); } /// /// Retrieves a collection of records for a specific patient, optionally filtered by a time range and limited in count, sorted by time in descending order. /// /// The identifier of the patient whose pump observations are being queried. /// Optional inclusive lower bound for the observation time. When provided, only observations on or after this time are returned. /// Optional inclusive upper bound for the observation time. When provided, only observations on or before this time are returned. /// Optional maximum number of observations to return. When null, all matching observations are returned. /// A task that represents the asynchronous operation. The task result contains an of matching observations ordered from newest to oldest. public async Task> FindByPatientIdAsync( ObjectId patientId, DateTime? from = null, DateTime? to = null, int? limit = null) { var filter = Builders.Filter.Eq(x => x.PatientId, patientId); if (from.HasValue) filter &= Builders.Filter.Gte(x => x.Time, from.Value); if (to.HasValue) filter &= Builders.Filter.Lte(x => x.Time, to.Value); var query = Collection.Find(filter).SortByDescending(x => x.Time); if (limit.HasValue) query = query.Limit(limit.Value) as IOrderedFindFluent; return await query.ToListAsync(); } /// /// Deletes all pump observations whose recording time is earlier than the specified cutoff date. /// /// The cutoff date; observations with a timestamp before this value are removed. public async Task DeleteBeforeDate(DateTime addDays) { var filter = Builders.Filter.Lt(x => x.Time, addDays); await Collection.DeleteManyAsync(filter); } } }