610 lines
26 KiB
C#
610 lines
26 KiB
C#
using adas_core.Application.Exceptions;
|
|
using adas_core.Application.Repositories.Interfaces;
|
|
using adas_core.Application.Services.Interfaces;
|
|
using adas_core.Application.Subscriptions;
|
|
using adas_core.Domain.Enums;
|
|
using adas_core.Domain.Models;
|
|
using adas_core.Domain.Models.AppSettings;
|
|
using adas_core.Domain.Models.Filter;
|
|
using adas_core.Domain.Models.MongoModels;
|
|
using adas_core.Domain.Models.Pumps;
|
|
using adas_core.Domain.Models.Responses;
|
|
using adas_core.Domain.Utils;
|
|
using Microsoft.AspNetCore.Http;
|
|
using Microsoft.Extensions.Logging;
|
|
using Microsoft.Extensions.Options;
|
|
using MongoDB.Bson;
|
|
using Patient = adas_core.Domain.Models.MongoModels.Patient;
|
|
|
|
namespace adas_core.Application.Services
|
|
{
|
|
/// <summary>
|
|
/// Servicio maestro de gestión de bombas
|
|
/// - Procesa ApiRequest (HL7/Alaris transformado)
|
|
/// - Histórico clínico (pump_observations)
|
|
/// - Histórico de alarmas (pump_alarm_events)
|
|
/// - Alarmas activas (pump_alarm_state)
|
|
/// - Snapshot (pump_state)
|
|
/// - Broadcast de snapshots (PumpState + PumpAlarmState)
|
|
/// </summary>
|
|
public class PumpService(
|
|
IPumpObservationRepository pumpObservationRepository,
|
|
IPumpStateRepository pumpStateRepo,
|
|
IPumpAlarmEventRepository alarmEventRepo,
|
|
IPumpAlarmStateRepository alarmStateRepo,
|
|
IPumpArchiveRepository pumpArchiveRepo,
|
|
IPatientService patientService,
|
|
IConfigPumpsService configPumpsService,
|
|
IOptions<ApiSettings> apiSettings,
|
|
ILogger<PumpService> logger,
|
|
ISubscribersService subscribersService,
|
|
IClientMessageService clientMessageService,
|
|
Lazy<ICalculatedObservationsService> calculatedObservationsService,
|
|
IHttpContextAccessor httpContextAccessor,
|
|
ILocalAuditService auditService,
|
|
IConfigUnitsService configUnitsService)
|
|
: IPumpService
|
|
{
|
|
|
|
// Settings
|
|
private readonly int _pumpExpiresSeconds = apiSettings.Value.PumpExpiresSeconds;
|
|
private readonly bool _sendPumpsZero = apiSettings.Value.SendPumpsZero;
|
|
|
|
// ======================================================================
|
|
// ENTRYPOINT
|
|
// ======================================================================
|
|
public async Task SaveRequest(ApiRequest req)
|
|
{
|
|
// Normalizar single vs list (Alaris puede mandar 1 sola)
|
|
if (req.PumpObservation != null && (req.PumpObservations == null || req.PumpObservations.Count == 0))
|
|
req.PumpObservations = [req.PumpObservation];
|
|
|
|
if (req.PumpObservations == null || req.PumpObservations.Count == 0)
|
|
{
|
|
logger.LogWarning("ApiRequest contains 0 PumpObservations");
|
|
return;
|
|
}
|
|
|
|
// Busca paciente (lookup/create) a partir de PatientNumber / Patient / PatientId string
|
|
var resolvedPatient = await patientService.FindPatientByApiRequest(req);
|
|
var foundPatientId = resolvedPatient?.Id;
|
|
|
|
foreach (var pobs in req.PumpObservations)
|
|
{
|
|
UpdatePatientFromRequest(req, pobs, foundPatientId);
|
|
|
|
// Regla Alaris: si el origen es AlarisPump y no hay PatientId -> descartar
|
|
if (string.Equals(req.Type, "AlarisPump", StringComparison.OrdinalIgnoreCase) &&
|
|
pobs.PatientId == null)
|
|
{
|
|
logger.LogWarning("AlarisPump: Observation descartada por ausencia de PatientId. DeviceId={device}", pobs.DeviceId);
|
|
continue;
|
|
}
|
|
|
|
if (pobs.Time == DateTime.MinValue)
|
|
pobs.Time = DateTime.UtcNow;
|
|
|
|
try
|
|
{
|
|
// 1) Procesar por tipo HL7 o según ObservationType (Alaris)
|
|
switch (req.Type)
|
|
{
|
|
case "ORU_R01": // PCD-01
|
|
case "ORU_R42": // PCD-10
|
|
pobs.MessageType = PumpEnum.PumpMessageType.Observation;
|
|
await ProcessObservation(pobs);
|
|
break;
|
|
|
|
case "ORU_R40": // PCD-04
|
|
pobs.MessageType = PumpEnum.PumpMessageType.Alarm;
|
|
await ProcessAlarm(pobs);
|
|
break;
|
|
|
|
case "AlarisPump":
|
|
if (pobs.MessageType == PumpEnum.PumpMessageType.Alarm)
|
|
await ProcessAlarm(pobs);
|
|
else
|
|
await ProcessObservation(pobs);
|
|
break;
|
|
|
|
default:
|
|
logger.LogWarning("Tipo de request desconocido para PumpService: {type}", req.Type);
|
|
break;
|
|
}
|
|
|
|
// 2) Snapshot de bomba
|
|
var state = await UpdatePumpState(pobs);
|
|
|
|
// 3) Alarmas activas del dispositivo
|
|
var activeAlarms = await alarmStateRepo.FindAllActiveByDeviceAsync(pobs.DeviceId!);
|
|
|
|
// 4) Broadcast de snapshots (PumpState + todas las PumpAlarmState)
|
|
await SendSnapshotsBroadcast(state, activeAlarms, pobs.PatientId ?? foundPatientId, req);
|
|
|
|
// 5) Retención (sobre histórico de observaciones) - reglas
|
|
await DoRetentionActions(pobs);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
logger.LogError(ex, "Error procesando observation/alarm DeviceId={deviceId}", pobs.DeviceId);
|
|
}
|
|
}
|
|
}
|
|
|
|
// ======================================================================
|
|
// OBSERVACIONES (PCD-01 / PCD-10)
|
|
// ======================================================================
|
|
private async Task ProcessObservation(PumpObservation obs)
|
|
{
|
|
logger.LogDebug("Insertando OBSERVATION DeviceId={dev} Time={time}", obs.DeviceId, obs.Time);
|
|
|
|
obs.Id = ObjectId.GenerateNewId();
|
|
obs.Expires = _pumpExpiresSeconds;
|
|
|
|
var mapped = await MapPumpObservation(obs);
|
|
if (mapped != null)
|
|
{
|
|
await pumpObservationRepository.InsertAsync(mapped);
|
|
await auditService.CreateAuditLogAsync(httpContextAccessor.HttpContext?.User!, null, mapped);
|
|
|
|
}
|
|
}
|
|
|
|
// ======================================================================
|
|
// ALARMAS (PCD-04)
|
|
// ======================================================================
|
|
private async Task ProcessAlarm(PumpObservation obs)
|
|
{
|
|
logger.LogDebug("Insertando ALARM DeviceId={dev}, Phase={phase}, Type={type}",
|
|
obs.DeviceId, obs.EventPhase, obs.AlarmType);
|
|
|
|
var alarmEvent = new PumpAlarmEvent
|
|
{
|
|
Id = ObjectId.GenerateNewId(),
|
|
DeviceId = obs.DeviceId,
|
|
RackId = obs.RackId,
|
|
DeviceTypeMdc = obs.DeviceTypeMdc,
|
|
DeviceIp = obs.DeviceIp,
|
|
PillarAssembly = obs.PillarAssembly,
|
|
PillarRackSlot = obs.PillarRackSlot,
|
|
|
|
Time = obs.Time,
|
|
InfusionId = obs.InfusionId,
|
|
|
|
AlarmType = obs.AlarmType,
|
|
AlarmTypeMdc = obs.AlarmTypeMdc,
|
|
AlarmDescription = obs.AlarmDescription,
|
|
AlarmPriority = obs.AlarmPriority,
|
|
AlarmState = obs.AlarmState,
|
|
AlarmInactivationState = obs.AlarmInactivationState,
|
|
EventPhase = obs.EventPhase,
|
|
|
|
AlertSourceMdc = obs.AlertSourceMdc,
|
|
|
|
PatientId = obs.PatientId
|
|
};
|
|
|
|
await alarmEventRepo.InsertAsync(alarmEvent);
|
|
|
|
await UpdateAlarmState(obs);
|
|
}
|
|
|
|
private async Task UpdateAlarmState(PumpObservation obs)
|
|
{
|
|
if (obs.EventPhase == null) return;
|
|
|
|
var phaseLower = obs.EventPhase.Value.ToString().ToLowerInvariant();
|
|
|
|
// Cierre
|
|
if (phaseLower == "end")
|
|
{
|
|
await alarmStateRepo.RemoveAsync(obs.DeviceId!, obs.AlarmType, obs.AlarmTypeMdc);
|
|
return;
|
|
}
|
|
|
|
// Start/continue → upsert
|
|
var state = new PumpAlarmState
|
|
{
|
|
Id = ObjectId.GenerateNewId(),
|
|
DeviceId = obs.DeviceId,
|
|
|
|
AlarmType = obs.AlarmType,
|
|
AlarmCodeMdc = obs.AlarmTypeMdc,
|
|
AlarmDescription = obs.AlarmDescription,
|
|
AlarmPriority = obs.AlarmPriority,
|
|
AlarmState = obs.AlarmState,
|
|
LastPhase = obs.EventPhase,
|
|
|
|
FirstSeen = DateTime.UtcNow,
|
|
LastUpdated = DateTime.UtcNow,
|
|
|
|
AlertSourceMdc = obs.AlertSourceMdc,
|
|
InfusionId = obs.InfusionId,
|
|
|
|
PatientId = obs.PatientId
|
|
};
|
|
|
|
await alarmStateRepo.UpsertActiveAsync(state);
|
|
}
|
|
|
|
// ======================================================================
|
|
// SNAPSHOT (devuelve el PumpState actualizado)
|
|
// ======================================================================
|
|
private async Task<PumpState> UpdatePumpState(PumpObservation obs)
|
|
{
|
|
var current = await pumpStateRepo.FindByDeviceIdAsync(obs.DeviceId!)
|
|
?? new PumpState
|
|
{
|
|
Id = ObjectId.GenerateNewId(),
|
|
DeviceId = obs.DeviceId
|
|
};
|
|
|
|
MergePumpState(current, obs);
|
|
current.LastUpdated = DateTime.UtcNow;
|
|
|
|
await pumpStateRepo.UpsertAsync(current);
|
|
return current;
|
|
}
|
|
|
|
private static void MergePumpState(PumpState state, PumpObservation obs)
|
|
{
|
|
// Identidad / físico
|
|
if (!string.IsNullOrWhiteSpace(obs.RackId)) state.RackId = obs.RackId;
|
|
if (!string.IsNullOrWhiteSpace(obs.DeviceTypeMdc)) state.DeviceTypeMdc = obs.DeviceTypeMdc;
|
|
if (!string.IsNullOrWhiteSpace(obs.DeviceIp)) state.DeviceIp = obs.DeviceIp;
|
|
if (!string.IsNullOrWhiteSpace(obs.PillarAssembly)) state.PillarAssembly = obs.PillarAssembly;
|
|
if (!string.IsNullOrWhiteSpace(obs.PillarRackSlot)) state.PillarRackSlot = obs.PillarRackSlot;
|
|
|
|
// Infusión
|
|
if (!string.IsNullOrWhiteSpace(obs.InfusionId)) state.InfusionId = obs.InfusionId;
|
|
|
|
// Estado
|
|
if (obs.InfusingStatus != null)
|
|
{
|
|
state.InfusingStatus = obs.InfusingStatus;
|
|
state.IsInfusing = obs.IsInfusing ?? false;
|
|
}
|
|
if (obs.Status != null) state.Status = obs.Status;
|
|
if (obs.PumpMode != null) state.PumpMode = obs.PumpMode;
|
|
|
|
if (!string.IsNullOrWhiteSpace(obs.ActiveSourceInfo)) state.ActiveSourceInfo = obs.ActiveSourceInfo;
|
|
if (!string.IsNullOrWhiteSpace(obs.InfusionModeDetail)) state.InfusionModeDetail = obs.InfusionModeDetail;
|
|
if (!string.IsNullOrWhiteSpace(obs.NotDeliveringReason)) state.NotDeliveringReason = obs.NotDeliveringReason;
|
|
if (!string.IsNullOrWhiteSpace(obs.Source)) state.Source = obs.Source;
|
|
|
|
// Métricas
|
|
if (HasValue(obs.FlowFluid)) state.FlowFluid = obs.FlowFluid;
|
|
if (HasValue(obs.Rate)) state.Rate = obs.Rate;
|
|
if (HasValue(obs.VolumeInfused)) state.VolumeInfused = obs.VolumeInfused;
|
|
if (HasValue(obs.FluidDelivTotal)) state.FluidDelivTotal = obs.FluidDelivTotal;
|
|
if (HasValue(obs.FluidDelivTotalSet)) state.FluidDelivTotalSet = obs.FluidDelivTotalSet;
|
|
if (HasValue(obs.VolumeRemaining)) state.VolumeRemaining = obs.VolumeRemaining;
|
|
if (HasValue(obs.Vtbi)) state.Vtbi = obs.Vtbi;
|
|
if (HasValue(obs.TimeRemaining)) state.TimeRemaining = obs.TimeRemaining;
|
|
if (HasValue(obs.TimeProgrammed)) state.TimeProgrammed = obs.TimeProgrammed;
|
|
|
|
// Medicación
|
|
if (!string.IsNullOrWhiteSpace(obs.DrugName)) state.DrugName = obs.DrugName;
|
|
if (!string.IsNullOrWhiteSpace(obs.DrugId)) state.DrugId = obs.DrugId;
|
|
|
|
if (HasValue(obs.Concentration)) state.Concentration = obs.Concentration;
|
|
if (HasValue(obs.DoseRate)) state.DoseRate = obs.DoseRate;
|
|
if (HasValue(obs.DrugAmount)) state.DrugAmount = obs.DrugAmount;
|
|
if (HasValue(obs.DrugDoseDelivered)) state.DrugDoseDelivered = obs.DrugDoseDelivered;
|
|
|
|
if (HasValue(obs.PatientWeight)) state.PatientWeight = obs.PatientWeight;
|
|
if (obs.Syringe != null) state.Syringe = obs.Syringe;
|
|
|
|
// Eventos / Alarmas resumen
|
|
if (obs.Event != null) state.Event = obs.Event;
|
|
if (obs.EventPhase != null) state.EventPhase = obs.EventPhase;
|
|
|
|
if (obs.AlarmType != null) state.AlarmType = obs.AlarmType;
|
|
if (!string.IsNullOrWhiteSpace(obs.AlarmDescription)) state.AlarmDescription = obs.AlarmDescription;
|
|
if (!string.IsNullOrWhiteSpace(obs.AlarmState)) state.AlarmState = obs.AlarmState;
|
|
if (!string.IsNullOrWhiteSpace(obs.AlarmInactivationState)) state.AlarmInactivationState = obs.AlarmInactivationState;
|
|
if (!string.IsNullOrWhiteSpace(obs.AlarmPriority)) state.AlarmPriority = obs.AlarmPriority;
|
|
if (!string.IsNullOrWhiteSpace(obs.AlarmTypeMdc)) state.AlarmCodeMdc = obs.AlarmTypeMdc;
|
|
|
|
// Paciente
|
|
state.PatientId = obs.PatientId;
|
|
}
|
|
|
|
private static bool HasValue(CommonPumpTypes.PumpValue? v) => v is { Value: not null };
|
|
|
|
// MAP
|
|
|
|
public async Task<PumpObservation?> MapPumpObservation(PumpObservation obs)
|
|
{
|
|
var obs2 = await configPumpsService.Map(obs);
|
|
var obs3 = await configUnitsService.Map(obs2);
|
|
var obs4 = await calculatedObservationsService.Value.Map(obs3);
|
|
|
|
if (obs4 == null)
|
|
logger.LogDebug("Mapping ignorado para obs");
|
|
|
|
return obs3;
|
|
}
|
|
|
|
//Métodos
|
|
|
|
public Task SaveRequestAsync(ApiRequest req) => SaveRequest(req);
|
|
|
|
|
|
// Últimas N observaciones por paciente
|
|
public async Task<List<PumpObservation>> FindLastPumpObservations(ObjectId patientId, int num = 1)
|
|
{
|
|
var list = await pumpObservationRepository.FindByPatientId(patientId);
|
|
if (list is List<PumpObservation> pumpObservations)
|
|
return pumpObservations is { Count: 0 }
|
|
? []
|
|
: pumpObservations.OrderByDescending(x => x.Time).Take(num).ToList();
|
|
return [];
|
|
}
|
|
|
|
// Última fecha de observación por paciente (para todos los pacientes)
|
|
public async Task<Dictionary<ObjectId, DateTime>> FindAllLastPatientObservationTime()
|
|
{
|
|
return await pumpObservationRepository.FindAllLastPatientObservationTimeAsync();
|
|
}
|
|
|
|
// Borrado completo por paciente (observaciones + alarmas + alarmState)
|
|
public async Task DeleteByPatientId(ObjectId id)
|
|
{
|
|
logger.LogDebug("Delete Pump data by Patient Id {id}", id);
|
|
await pumpObservationRepository.DeleteByPatientId(id);
|
|
await alarmEventRepo.DeleteByPatientId(id);
|
|
await alarmStateRepo.DeleteByPatientId(id);
|
|
}
|
|
|
|
// Archivo → por entidad paciente
|
|
public async Task Archive(Patient patient) => await ArchiveByPatientId(patient.Id);
|
|
|
|
// Archivo → por PatientId (mueve a archive_pumpobservations y elimina del activo)
|
|
public async Task ArchiveByPatientId(ObjectId id)
|
|
{
|
|
var list = await pumpObservationRepository.FindByPatientId(id);
|
|
var pumpObservations = list.ToList();
|
|
if (pumpObservations.Count != 0)
|
|
await pumpArchiveRepo.InsertManyAsync(pumpObservations);
|
|
|
|
await DeleteByPatientId(id);
|
|
logger.LogDebug("Archived Pump observations & deleted active data by Patient Id {id}", id);
|
|
}
|
|
|
|
// Actualización masiva
|
|
public async Task UpdateManyObjectId(string nameId, ObjectId id, ObjectId oldId)
|
|
{
|
|
|
|
if (string.IsNullOrWhiteSpace(nameId))
|
|
throw new ArgumentException("nameId no puede ser nulo o vacío.", nameof(nameId));
|
|
|
|
var updatedObs = await pumpObservationRepository.UpdateManyObjectIdByFieldAsync(nameId, id, oldId);
|
|
|
|
// actualizar también alarmas activas e históricas
|
|
_ = await alarmEventRepo.UpdateManyObjectIdByFiledNameAsync(nameId, id, oldId);
|
|
_ = await alarmStateRepo.UpdateManyObjectIdByFieldNameAsync(nameId, id, oldId);
|
|
|
|
logger.LogInformation(
|
|
"UpdateManyObjectId completado. Campo={field}, oldId={oldId}, newId={newId}. Obs actualizadas={obsUpdated}",
|
|
nameId, oldId, id, updatedObs);
|
|
|
|
// Si quieres auditar el cambio:
|
|
await auditService.CreateAuditLogAsync(
|
|
httpContextAccessor.HttpContext?.User!,
|
|
new { Field = nameId, OldId = oldId, NewId = id, Scope = "PumpObservation" },
|
|
null);
|
|
|
|
}
|
|
|
|
// ConfigPumpsService
|
|
public async Task<List<ConfigPumpItem>?> GetItemsById(string id)
|
|
=> await configPumpsService.GetConfigItems(id);
|
|
|
|
public async Task<List<ConfigPumps>?> GetAllPumpConfig()
|
|
=> await configPumpsService.GetAllPumpConfigs()
|
|
?? throw new NotFoundException(HttpEnum.ErrorMessage.NotFoundResourceMissing);
|
|
|
|
public async Task<ConfigPumps?> GetPumpConfigsById(string id)
|
|
=> await configPumpsService.GetPumpConfigById(id)
|
|
?? throw new NotFoundException(HttpEnum.ErrorMessage.NotFoundResourceMissing);
|
|
|
|
public async Task<ConfigPumps?> UpdatePumpConfig(ConfigPumps config)
|
|
{
|
|
var oldConfig = await configPumpsService.GetPumpConfigById(config.Id);
|
|
await auditService.CreateAuditLogAsync(httpContextAccessor.HttpContext?.User!, oldConfig, config);
|
|
return await configPumpsService.UpdatePumpConfig(config)
|
|
?? throw new NotFoundException(HttpEnum.ErrorMessage.NotFoundResourceMissing);
|
|
}
|
|
|
|
public async Task<ConfigPumps?> InsertPumpConfig(ConfigPumps config)
|
|
{
|
|
await auditService.CreateAuditLogAsync(httpContextAccessor.HttpContext?.User!, null, config);
|
|
return await configPumpsService.InsertPumpConfig(config)
|
|
?? throw new ConflictException(HttpEnum.ErrorMessage.ConflictCreationFailed);
|
|
}
|
|
|
|
public async Task<bool> DeletePumpConfig(ConfigPumps config)
|
|
{
|
|
var result = await configPumpsService.DeletePumpConfig(config);
|
|
if (!result) throw new ConflictException(HttpEnum.ErrorMessage.ConflictDeleteFailed);
|
|
await auditService.CreateAuditLogAsync(httpContextAccessor.HttpContext?.User!, config, null);
|
|
return result;
|
|
}
|
|
|
|
// Paginación de observaciones
|
|
public async Task<PaginationResponse<PumpObservation>?> GetPaginatedPump(PaginationFilter filter)
|
|
{
|
|
// Implementación compatible sin nuevos métodos en los repos:
|
|
// 1) Si llega PatientId, paginamos en memoria desde FindByPatientId.
|
|
// 2) Si llega DeviceId, usamos FindByDeviceIdAsync y paginamos en memoria.
|
|
// 3) Si no hay filtro, devolvemos vacío para evitar lecturas completas.
|
|
|
|
var page = filter.PageNumber <= 0 ? 1 : filter.PageNumber;
|
|
var size = filter.PageSize <= 0 ? 20 : filter.PageSize;
|
|
|
|
var fr = filter.FilteredRequest;
|
|
if (fr == null) return new PaginationResponse<PumpObservation>([], page, size, 0);
|
|
|
|
List<PumpObservation> all;
|
|
|
|
if (!string.IsNullOrWhiteSpace(fr.PatientId))
|
|
{
|
|
if (!ObjectId.TryParse(fr.PatientId, out var patientId))
|
|
return new PaginationResponse<PumpObservation>([], page, size, 0);
|
|
|
|
var list = await pumpObservationRepository.FindByPatientId(patientId);
|
|
all = list.ToList();
|
|
|
|
if (fr.StartDate.HasValue)
|
|
all = all.Where(o => o.Time >= fr.StartDate.Value).ToList();
|
|
if (fr.EndDate.HasValue)
|
|
all = all.Where(o => o.Time <= fr.EndDate.Value).ToList();
|
|
}
|
|
else if (!string.IsNullOrWhiteSpace(fr.DeviceId))
|
|
{
|
|
var found = await pumpObservationRepository.FindByDeviceIdAsync(fr.DeviceId, fr.StartDate, fr.EndDate);
|
|
all = found.ToList();
|
|
}
|
|
else
|
|
{
|
|
all = [];
|
|
}
|
|
|
|
var count = all.Count;
|
|
var pageData = all
|
|
.OrderByDescending(o => o.Time)
|
|
.Skip((page - 1) * size)
|
|
.Take(size)
|
|
.ToList();
|
|
|
|
return new PaginationResponse<PumpObservation>(pageData, page, size, count);
|
|
}
|
|
|
|
// Inserción manual
|
|
public async Task InsertPumpObservation(PumpObservation obs)
|
|
{
|
|
logger.LogDebug("Insert {obs}", obs);
|
|
obs.Id = ObjectId.GenerateNewId();
|
|
if (obs.Time == DateTime.MinValue) obs.Time = DateTime.UtcNow;
|
|
var obsMapped = await MapPumpObservation(obs);
|
|
|
|
if (obsMapped != null)
|
|
{
|
|
await pumpObservationRepository.InsertAsync(obsMapped);
|
|
await auditService.CreateAuditLogAsync(httpContextAccessor.HttpContext?.User!, null, obsMapped);
|
|
|
|
if (!_sendPumpsZero && obsMapped.Number == 0) return;
|
|
|
|
// Tras inserción manual, actualizar y emitir snapshots
|
|
var state = await UpdatePumpState(obsMapped);
|
|
var activeAlarms = await alarmStateRepo.FindAllActiveByDeviceAsync(obsMapped.DeviceId!);
|
|
await SendSnapshotsBroadcast(state, activeAlarms, obsMapped.PatientId, null);
|
|
|
|
await DoRetentionActions(obsMapped);
|
|
}
|
|
}
|
|
|
|
private async Task DoRetentionActions(PumpObservation obs)
|
|
{
|
|
var result = await configPumpsService.RetentionActions(obs);
|
|
if (result is not { RetentionPolicyValue: not null })
|
|
return;
|
|
|
|
switch (result.RetentionPolicy)
|
|
{
|
|
case RetentionPolicy.DeleteOlderDays:
|
|
{
|
|
var removed = await pumpObservationRepository.DeleteOlderThanDaysAsync(
|
|
result.RetentionPolicyValue.Value);
|
|
|
|
logger.LogInformation(
|
|
"Retention DeleteOlderDays: {removed} deleted (>{days} days)",
|
|
removed, result.RetentionPolicyValue);
|
|
|
|
break;
|
|
}
|
|
|
|
case RetentionPolicy.DeleteOlderNumber:
|
|
{
|
|
var removed = await pumpObservationRepository.DeleteKeepLastNAsync(
|
|
result.RetentionPolicyValue.Value);
|
|
|
|
logger.LogInformation(
|
|
"Retention DeleteOlderNumber: {removed} deleted (keeping {max})",
|
|
removed, result.RetentionPolicyValue);
|
|
|
|
break;
|
|
}
|
|
|
|
case RetentionPolicy.NoDelete:
|
|
case RetentionPolicy.DeleteOlderSeconds:
|
|
default:
|
|
logger.LogWarning("Unknown retention policy: {policy}", result.RetentionPolicy);
|
|
break;
|
|
}
|
|
}
|
|
|
|
// BROADCAST SNAPSHOTS: PumpState + PumpAlarmState (activos)
|
|
private async Task SendSnapshotsBroadcast(
|
|
PumpState state,
|
|
IEnumerable<PumpAlarmState> activeAlarms,
|
|
ObjectId? patientId,
|
|
ApiRequest? req)
|
|
{
|
|
var subscribers = new List<WsSubscriber>();
|
|
|
|
if (patientId != null)
|
|
{
|
|
var patient = await patientService.FindById(patientId.Value);
|
|
if (patient != null)
|
|
{
|
|
subscribers = subscribersService.GetSubscribers()
|
|
.Where(s => !s.Locations.IsNullOrEmpty()
|
|
&& s.Locations.Any(c =>
|
|
c.UnitName == patient.Location.UnitName &&
|
|
c.Bed == patient.Location.Bed &&
|
|
c.Room == patient.Location.Room))
|
|
.ToList();
|
|
}
|
|
}
|
|
else if (req?.Location != null)
|
|
{
|
|
var loc = req.Location;
|
|
subscribers = subscribersService.GetSubscribers()
|
|
.Where(s => !s.Locations.IsNullOrEmpty()
|
|
&& s.Locations.Any(c =>
|
|
c.UnitName == loc.UnitName &&
|
|
c.Bed == loc.Bed &&
|
|
c.Room == loc.Room))
|
|
.ToList();
|
|
}
|
|
|
|
if (subscribers.Count == 0) return;
|
|
|
|
// Enviar PumpState
|
|
foreach (var sub in subscribers)
|
|
await clientMessageService.SendAsync(sub.Id, OperationType.Pump, state);
|
|
|
|
// Enviar todas las PumpAlarmState activas
|
|
foreach (var alarm in activeAlarms)
|
|
{
|
|
foreach (var sub in subscribers)
|
|
await clientMessageService.SendAsync(sub.Id, OperationType.PumpAlarm, alarm);
|
|
}
|
|
}
|
|
|
|
|
|
private static void UpdatePatientFromRequest(ApiRequest req, PumpObservation obs, ObjectId? foundPatientId)
|
|
{
|
|
if (obs.PatientId == null && foundPatientId != null)
|
|
obs.PatientId = foundPatientId;
|
|
|
|
|
|
if (obs.PatientId != null || string.IsNullOrWhiteSpace(req.PatientId)) return;
|
|
|
|
if (ObjectId.TryParse(req.PatientId, out var parsed))
|
|
obs.PatientId = parsed;
|
|
}
|
|
}
|
|
} |