using System.Text.Json; using Dapper; using Npgsql; namespace Kbx.Shared.ExternalData; public sealed class KbxExternalDataRepository(NpgsqlDataSource dataSource) { public async Task GetAsync(Guid tenantId, string datasetId, string cacheKey, CancellationToken ct) { const string sql = """ select tenant_id as TenantId, dataset_id as DatasetId, provider_id as ProviderId, cache_key as CacheKey, request_descriptor::text as RequestDescriptorJson, normalized_data::text as NormalizedJson, provider_observed_at as ProviderObservedAt, requested_at as RequestedAt, received_at as ReceivedAt, ingested_at as IngestedAt, fresh_until as FreshUntil, usable_until as UsableUntil, payload_sha256 as PayloadSha256, normalizer_version as NormalizerVersion, state, correlation_id as CorrelationId from kbx.external_data_cache where tenant_id=@TenantId and dataset_id=@DatasetId and cache_key=@CacheKey """; await using var connection = await dataSource.OpenConnectionAsync(ct); var row = await connection.QuerySingleOrDefaultAsync(new CommandDefinition(sql, new { TenantId=tenantId, DatasetId=datasetId, CacheKey=cacheKey }, cancellationToken:ct)); return row is null ? null : new KbxExternalDataCacheEntry(row.TenantId,row.DatasetId,row.ProviderId,row.CacheKey,JsonDocument.Parse(row.RequestDescriptorJson),JsonDocument.Parse(row.NormalizedJson),row.ProviderObservedAt,row.RequestedAt,row.ReceivedAt,row.IngestedAt,row.FreshUntil,row.UsableUntil,row.PayloadSha256,row.NormalizerVersion,row.State,row.CorrelationId); } public async Task UpsertAsync(KbxExternalDataCacheEntry entry, CancellationToken ct) { const string sql = """ insert into kbx.external_data_cache(tenant_id,dataset_id,provider_id,cache_key,request_descriptor,normalized_data,provider_observed_at,requested_at,received_at,ingested_at,fresh_until,usable_until,payload_sha256,normalizer_version,state,correlation_id) values(@TenantId,@DatasetId,@ProviderId,@CacheKey,cast(@RequestDescriptorJson as jsonb),cast(@NormalizedJson as jsonb),@ProviderObservedAt,@RequestedAt,@ReceivedAt,@IngestedAt,@FreshUntil,@UsableUntil,@PayloadSha256,@NormalizerVersion,@State,@CorrelationId) on conflict(tenant_id,dataset_id,cache_key) do update set provider_id=excluded.provider_id, request_descriptor=excluded.request_descriptor, normalized_data=excluded.normalized_data, provider_observed_at=excluded.provider_observed_at, requested_at=excluded.requested_at, received_at=excluded.received_at, ingested_at=excluded.ingested_at, fresh_until=excluded.fresh_until, usable_until=excluded.usable_until, payload_sha256=excluded.payload_sha256, normalizer_version=excluded.normalizer_version, state=excluded.state, correlation_id=excluded.correlation_id; insert into kbx.external_data_observations(tenant_id,dataset_id,provider_id,cache_key,payload_sha256,normalizer_version,provider_observed_at,received_at,ingested_at,state,correlation_id) values(@TenantId,@DatasetId,@ProviderId,@CacheKey,@PayloadSha256,@NormalizerVersion,@ProviderObservedAt,@ReceivedAt,@IngestedAt,@State,@CorrelationId); """; await using var connection = await dataSource.OpenConnectionAsync(ct); await using var transaction = await connection.BeginTransactionAsync(ct); await connection.ExecuteAsync(new CommandDefinition(sql,new { entry.TenantId,entry.DatasetId,entry.ProviderId,entry.CacheKey,RequestDescriptorJson=entry.RequestDescriptor.RootElement.GetRawText(),NormalizedJson=entry.NormalizedData.RootElement.GetRawText(),entry.ProviderObservedAt,entry.RequestedAt,entry.ReceivedAt,entry.IngestedAt,entry.FreshUntil,entry.UsableUntil,entry.PayloadSha256,entry.NormalizerVersion,entry.State,entry.CorrelationId },transaction,cancellationToken:ct)); await transaction.CommitAsync(ct); } private sealed record Row(Guid TenantId,string DatasetId,string ProviderId,string CacheKey,string RequestDescriptorJson,string NormalizedJson,DateTimeOffset? ProviderObservedAt,DateTimeOffset RequestedAt,DateTimeOffset ReceivedAt,DateTimeOffset IngestedAt,DateTimeOffset? FreshUntil,DateTimeOffset? UsableUntil,string PayloadSha256,string NormalizerVersion,string State,string CorrelationId); }