Files

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.
}
}
}