260 lines
9.2 KiB
C#
260 lines
9.2 KiB
C#
using System.Net.Http.Json;
|
|
using System.Text.Json;
|
|
using Microsoft.EntityFrameworkCore;
|
|
using Microsoft.Extensions.DependencyInjection;
|
|
using Microsoft.Extensions.Options;
|
|
using Minio.DataModel.Args;
|
|
using Parquet;
|
|
using StackExchange.Redis;
|
|
|
|
[Collection("Integration")]
|
|
public class DataLakePhase9Tests
|
|
{
|
|
private readonly ApiFixture _fixture;
|
|
private readonly HttpClient _http;
|
|
private readonly MinioOptions _minioOpts;
|
|
|
|
public DataLakePhase9Tests(ApiFixture fixture)
|
|
{
|
|
_fixture = fixture;
|
|
_http = fixture.CreateClient();
|
|
_minioOpts = fixture.Services
|
|
.GetRequiredService<IOptions<MinioOptions>>().Value;
|
|
}
|
|
|
|
[Fact]
|
|
public async Task ObservationEvents_FlushesToParquetInMinIO()
|
|
{
|
|
await EnsureThresholdsAsync();
|
|
var patientId = await CreatePatientAsync();
|
|
var encounterId = await CreateActiveEncounterAsync(patientId);
|
|
|
|
for (var i = 0; i < 3; i++)
|
|
{
|
|
var resp = await _http.PostAsJsonAsync(
|
|
$"/api/v1/encounters/{encounterId}/observations",
|
|
new BatchIngestRequest(new List<IngestObservationRequest>
|
|
{
|
|
new("HEART_RATE", 72 + i, "bpm", ObservationSource.Device, DateTimeOffset.UtcNow, null)
|
|
}));
|
|
resp.EnsureSuccessStatusCode();
|
|
}
|
|
|
|
await Task.Delay(TimeSpan.FromSeconds(20));
|
|
|
|
var keys = await ListMinioObjectsAsync("observations/");
|
|
Assert.True(keys.Count > 0,
|
|
$"No Parquet files found under observations/ in MinIO bucket '{_minioOpts.BucketName}'.");
|
|
Assert.All(keys, k => Assert.EndsWith(".parquet", k));
|
|
}
|
|
|
|
[Fact]
|
|
public async Task AlertEvents_FlushesToParquetInMinIO()
|
|
{
|
|
await EnsureThresholdsAsync();
|
|
for (var i = 0; i < 3; i++)
|
|
{
|
|
var pid = await CreatePatientAsync();
|
|
var eid = await CreateActiveEncounterAsync(pid);
|
|
|
|
var resp = await _http.PostAsJsonAsync(
|
|
$"/api/v1/encounters/{eid}/observations",
|
|
new BatchIngestRequest(new List<IngestObservationRequest>
|
|
{
|
|
new("POTASSIUM_MEQ_L", 2.1m, "mEq/L", ObservationSource.Device, DateTimeOffset.UtcNow, null)
|
|
}));
|
|
resp.EnsureSuccessStatusCode();
|
|
}
|
|
|
|
await Task.Delay(TimeSpan.FromSeconds(20));
|
|
|
|
var keys = await ListMinioObjectsAsync("alerts/");
|
|
Assert.True(keys.Count > 0,
|
|
$"No Parquet files found under alerts/ in MinIO bucket '{_minioOpts.BucketName}'.");
|
|
}
|
|
|
|
[Fact]
|
|
public async Task EncounterStatusEvents_FlushesToParquetInMinIO()
|
|
{
|
|
await EnsureThresholdsAsync();
|
|
for (var i = 0; i < 3; i++)
|
|
{
|
|
var pid = await CreatePatientAsync();
|
|
_ = await CreateActiveEncounterAsync(pid);
|
|
}
|
|
|
|
await Task.Delay(TimeSpan.FromSeconds(20));
|
|
|
|
var keys = await ListMinioObjectsAsync("encounters/");
|
|
Assert.True(keys.Count > 0,
|
|
$"No Parquet files found under encounters/ in MinIO bucket '{_minioOpts.BucketName}'.");
|
|
}
|
|
|
|
[Fact]
|
|
public async Task ObservationParquetFile_ContainsCorrectColumns()
|
|
{
|
|
await EnsureThresholdsAsync();
|
|
var patientId = await CreatePatientAsync();
|
|
var encounterId = await CreateActiveEncounterAsync(patientId);
|
|
|
|
for (var i = 0; i < 3; i++)
|
|
{
|
|
var resp = await _http.PostAsJsonAsync(
|
|
$"/api/v1/encounters/{encounterId}/observations",
|
|
new BatchIngestRequest(new List<IngestObservationRequest>
|
|
{
|
|
new("TEMP_C", 37.5m, "°C", ObservationSource.Device, DateTimeOffset.UtcNow, null)
|
|
}));
|
|
resp.EnsureSuccessStatusCode();
|
|
}
|
|
|
|
await Task.Delay(TimeSpan.FromSeconds(20));
|
|
|
|
var keys = await ListMinioObjectsAsync("observations/");
|
|
Assert.True(keys.Count > 0);
|
|
|
|
var bytes = await DownloadMinioObjectAsync(keys[0]);
|
|
Assert.True(bytes.Length > 0);
|
|
|
|
using var ms = new MemoryStream(bytes);
|
|
using var reader = await ParquetReader.CreateAsync(ms);
|
|
|
|
var columnNames = reader.Schema.DataFields.Select(f => f.Name).ToHashSet();
|
|
Assert.Contains("observation_id", columnNames);
|
|
Assert.Contains("encounter_id", columnNames);
|
|
Assert.Contains("observation_code", columnNames);
|
|
Assert.Contains("value", columnNames);
|
|
Assert.Contains("kafka_partition", columnNames);
|
|
Assert.Contains("kafka_offset", columnNames);
|
|
}
|
|
|
|
private async Task<List<string>> ListMinioObjectsAsync(string prefix)
|
|
{
|
|
var client = MinioClientFactory.Build(_minioOpts);
|
|
var keys = new List<string>();
|
|
|
|
var listArgs = new ListObjectsArgs()
|
|
.WithBucket(_minioOpts.BucketName)
|
|
.WithPrefix(prefix)
|
|
.WithRecursive(true);
|
|
|
|
await foreach (var item in client.ListObjectsEnumAsync(listArgs))
|
|
{
|
|
keys.Add(item.Key);
|
|
}
|
|
|
|
return keys;
|
|
}
|
|
|
|
private async Task<byte[]> DownloadMinioObjectAsync(string objectKey)
|
|
{
|
|
var client = MinioClientFactory.Build(_minioOpts);
|
|
using var ms = new MemoryStream();
|
|
|
|
await client.GetObjectAsync(new GetObjectArgs()
|
|
.WithBucket(_minioOpts.BucketName)
|
|
.WithObject(objectKey)
|
|
.WithCallbackStream(stream => stream.CopyTo(ms)));
|
|
|
|
return ms.ToArray();
|
|
}
|
|
|
|
private async Task<Guid> CreatePatientAsync()
|
|
{
|
|
var resp = await _http.PostAsJsonAsync("/api/v1/patients", new
|
|
{
|
|
firstName = "DataLake",
|
|
lastName = "Test",
|
|
dateOfBirth = "1980-11-12",
|
|
gender = "Female",
|
|
});
|
|
resp.EnsureSuccessStatusCode();
|
|
|
|
var body = await resp.Content.ReadFromJsonAsync<JsonDocument>();
|
|
return body!.RootElement.GetProperty("data").GetProperty("id").GetGuid();
|
|
}
|
|
|
|
private async Task<Guid> CreateActiveEncounterAsync(Guid patientId)
|
|
{
|
|
var resp = await _http.PostAsJsonAsync(
|
|
$"/api/v1/patients/{patientId}/encounters",
|
|
new
|
|
{
|
|
encounterType = "INPATIENT",
|
|
department = "ICU",
|
|
attendingPhysician = "Dr. Osei",
|
|
admittedAt = DateTimeOffset.UtcNow,
|
|
});
|
|
resp.EnsureSuccessStatusCode();
|
|
|
|
var body = await resp.Content.ReadFromJsonAsync<JsonDocument>();
|
|
var encounterId = body!.RootElement.GetProperty("data").GetProperty("id").GetGuid();
|
|
|
|
// Encounter may be created as Scheduled depending on fixture seed path; move it
|
|
// to Active if required, but tolerate Conflict when already Active.
|
|
var activateResp = await _http.PatchAsJsonAsync(
|
|
$"/api/v1/encounters/{encounterId}/status",
|
|
new { status = "Active" });
|
|
if (activateResp.StatusCode != System.Net.HttpStatusCode.Conflict)
|
|
activateResp.EnsureSuccessStatusCode();
|
|
|
|
return encounterId;
|
|
}
|
|
|
|
private async Task EnsureThresholdsAsync()
|
|
{
|
|
using var scope = _fixture.Services.CreateScope();
|
|
var db = scope.ServiceProvider.GetRequiredService<AppDbContext>();
|
|
|
|
await EnsureThresholdAsync(db, "HEART_RATE", "Heart Rate", "bpm", 30, 50, 100, 150);
|
|
await EnsureThresholdAsync(db, "POTASSIUM_MEQ_L", "Serum Potassium", "mEq/L", 2.5m, 3.5m, 5.0m, 6.5m);
|
|
await EnsureThresholdAsync(db, "TEMP_C", "Temperature", "°C", 34m, 36m, 37.8m, 40m);
|
|
|
|
var redis = scope.ServiceProvider.GetRequiredService<IConnectionMultiplexer>();
|
|
var cache = redis.GetDatabase(1);
|
|
await cache.StringSetAsync("threshold:HEART_RATE",
|
|
"""{"ObservationCode":"HEART_RATE","CriticalLow":30,"WarningLow":50,"WarningHigh":100,"CriticalHigh":150}""");
|
|
await cache.StringSetAsync("threshold:POTASSIUM_MEQ_L",
|
|
"""{"ObservationCode":"POTASSIUM_MEQ_L","CriticalLow":2.5,"WarningLow":3.5,"WarningHigh":5.0,"CriticalHigh":6.5}""");
|
|
await cache.StringSetAsync("threshold:TEMP_C",
|
|
"""{"ObservationCode":"TEMP_C","CriticalLow":34,"WarningLow":36,"WarningHigh":37.8,"CriticalHigh":40}""");
|
|
}
|
|
|
|
private static async Task EnsureThresholdAsync(
|
|
AppDbContext db,
|
|
string code,
|
|
string displayName,
|
|
string unit,
|
|
decimal criticalLow,
|
|
decimal warningLow,
|
|
decimal warningHigh,
|
|
decimal criticalHigh)
|
|
{
|
|
var existing = await db.AlertThresholds.FirstOrDefaultAsync(t => t.ObservationCode == code);
|
|
if (existing is not null) return;
|
|
|
|
db.AlertThresholds.Add(new AlertThreshold
|
|
{
|
|
Id = Guid.NewGuid(),
|
|
ObservationCode = code,
|
|
DisplayName = displayName,
|
|
Unit = unit,
|
|
CriticalLow = criticalLow,
|
|
WarningLow = warningLow,
|
|
WarningHigh = warningHigh,
|
|
CriticalHigh = criticalHigh,
|
|
CreatedAt = DateTimeOffset.UtcNow
|
|
});
|
|
|
|
try
|
|
{
|
|
await db.SaveChangesAsync();
|
|
}
|
|
catch (DbUpdateException ex) when (
|
|
ex.InnerException is Npgsql.PostgresException { SqlState: "23505" })
|
|
{
|
|
// Another test inserted the same observation code concurrently.
|
|
}
|
|
}
|
|
}
|