From e0d278e6eb8dc0a1c11a873385710e5b0ba01920 Mon Sep 17 00:00:00 2001 From: kjh2064 Date: Mon, 13 Jul 2026 01:00:36 +0900 Subject: [PATCH] refactor(dotnet): split collection read and write contracts --- .../Services/CollectionReadModelService.cs | 4 +- .../Services/KisDataCollectionOrchestrator.cs | 20 +++--- .../KisDataCollectionOrchestratorTests.cs | 61 +++++++++++-------- .../Interfaces/IDataCollectionStore.cs | 19 ++++++ .../Repositories/CollectionRepository.cs | 2 +- src/dotnet/QuantEngine.Web/Program.cs | 5 +- 6 files changed, 72 insertions(+), 39 deletions(-) diff --git a/src/dotnet/QuantEngine.Application/Services/CollectionReadModelService.cs b/src/dotnet/QuantEngine.Application/Services/CollectionReadModelService.cs index 89760d3b..beb28ab4 100644 --- a/src/dotnet/QuantEngine.Application/Services/CollectionReadModelService.cs +++ b/src/dotnet/QuantEngine.Application/Services/CollectionReadModelService.cs @@ -5,9 +5,9 @@ namespace QuantEngine.Application.Services; public sealed class CollectionReadModelService : ICollectionReadModelService { - private readonly ICollectionRepository _repository; + private readonly ICollectionReadRepository _repository; - public CollectionReadModelService(ICollectionRepository repository) + public CollectionReadModelService(ICollectionReadRepository repository) { _repository = repository; } diff --git a/src/dotnet/QuantEngine.Application/Services/KisDataCollectionOrchestrator.cs b/src/dotnet/QuantEngine.Application/Services/KisDataCollectionOrchestrator.cs index 8c4c1d01..cde1a169 100644 --- a/src/dotnet/QuantEngine.Application/Services/KisDataCollectionOrchestrator.cs +++ b/src/dotnet/QuantEngine.Application/Services/KisDataCollectionOrchestrator.cs @@ -13,7 +13,8 @@ namespace QuantEngine.Application.Services; public class KisDataCollectionOrchestrator : ICollectionOrchestrator { private readonly IKisApiClient _kisApiClient; - private readonly ICollectionRepository _repository; + private readonly ICollectionWriteRepository _writeRepository; + private readonly ICollectionReadRepository _readRepository; private readonly PriceDataNormalizer _normalizer; private readonly SourcePriorityResolver _priorityResolver; private readonly ILogger _logger; @@ -21,14 +22,16 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator public KisDataCollectionOrchestrator( IKisApiClient kisApiClient, - ICollectionRepository repository, + ICollectionWriteRepository repository, + ICollectionReadRepository readRepository, PriceDataNormalizer normalizer, SourcePriorityResolver priorityResolver, ILogger logger, IRuntimeAuditTrailService auditTrail) { _kisApiClient = kisApiClient; - _repository = repository; + _writeRepository = repository; + _readRepository = readRepository; _normalizer = normalizer; _priorityResolver = priorityResolver; _logger = logger; @@ -66,7 +69,7 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator CollectionSnapshotRecord? cachedSnapshot = null; if (IsMarketClosed()) { - var latest = await _repository.GetLatestSnapshotsForTickerAsync(ticker, 1); + var latest = await _readRepository.GetLatestSnapshotsForTickerAsync(ticker, 1); var todayPrefix = DateTime.UtcNow.AddHours(9).ToString("yyyy-MM-dd"); if (latest.Count > 0 && latest[0].CapturedAt.StartsWith(todayPrefix)) { @@ -96,7 +99,7 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator } // Save to DB - await _repository.SaveSnapshotAsync(new CollectionSnapshotRecord( + await _writeRepository.SaveSnapshotAsync(new CollectionSnapshotRecord( RunId: runId, DatasetName: "data_feed", Ticker: ticker, @@ -119,7 +122,7 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator _logger.LogWarning("Skipped invalid OHLCV bar for {Ticker}: constraints not satisfied", ticker); continue; } - await _repository.SavePriceHistoryDailyAsync(priceRecord); + await _writeRepository.SavePriceHistoryDailyAsync(priceRecord); } } } @@ -147,7 +150,7 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator { "error_kind", ex.GetType().Name } }); - await _repository.SaveErrorAsync(new CollectionErrorRecord( + await _writeRepository.SaveErrorAsync(new CollectionErrorRecord( RunId: runId, SourceName: "kis_collector", ErrorKind: ex.GetType().Name, @@ -166,7 +169,7 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator _auditTrail.Append("collection_audit", runId, new CollectionExecutionAudit(runId, result.Status, DateTimeOffset.Parse(startedAt), DateTimeOffset.Parse(finishedAt), result.SuccessCount, result.ErrorCount, "finished")); // Save run record - await _repository.SaveRunAsync(new CollectionRunRecord( + await _writeRepository.SaveRunAsync(new CollectionRunRecord( RunId: runId, Status: result.Status, StartedAt: startedAt, @@ -364,4 +367,3 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator } } } - diff --git a/src/dotnet/QuantEngine.Core.Tests/KisDataCollectionOrchestratorTests.cs b/src/dotnet/QuantEngine.Core.Tests/KisDataCollectionOrchestratorTests.cs index e919bbb5..622a175c 100644 --- a/src/dotnet/QuantEngine.Core.Tests/KisDataCollectionOrchestratorTests.cs +++ b/src/dotnet/QuantEngine.Core.Tests/KisDataCollectionOrchestratorTests.cs @@ -13,7 +13,8 @@ namespace QuantEngine.Core.Tests; public class KisDataCollectionOrchestratorTests { private readonly Mock _kisApiClientMock; - private readonly Mock _repositoryMock; + private readonly Mock _writeRepositoryMock; + private readonly Mock _readRepositoryMock; private readonly Mock> _loggerMock; private readonly Mock _auditTrailMock; private readonly PriceDataNormalizer _normalizer; @@ -23,15 +24,21 @@ public class KisDataCollectionOrchestratorTests public KisDataCollectionOrchestratorTests() { _kisApiClientMock = new Mock(); - _repositoryMock = new Mock(); + var repositoryMock = new Mock(); + _writeRepositoryMock = repositoryMock.As(); + _readRepositoryMock = repositoryMock.As(); _loggerMock = new Mock>(); _auditTrailMock = new Mock(); _priorityResolver = new SourcePriorityResolver(); _normalizer = new PriceDataNormalizer(_priorityResolver); + _kisApiClientMock + .Setup(k => k.GetDailyItemChartPriceAsync(It.IsAny(), It.IsAny(), It.IsAny(), "D", It.IsAny())) + .ReturnsAsync(new Dictionary()); _orchestrator = new KisDataCollectionOrchestrator( _kisApiClientMock.Object, - _repositoryMock.Object, + _writeRepositoryMock.Object, + _readRepositoryMock.Object, _normalizer, _priorityResolver, _loggerMock.Object, @@ -56,19 +63,19 @@ public class KisDataCollectionOrchestratorTests CapturedAt: $"{todayPrefix}T14:30:00" ); - _repositoryMock + _readRepositoryMock .Setup(r => r.GetLatestSnapshotsForTickerAsync(ticker, It.IsAny())) .ReturnsAsync(new List { cachedSnapshot }); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SaveSnapshotAsync(It.IsAny())) .Returns(Task.CompletedTask); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SavePriceHistoryDailyAsync(It.IsAny())) .Returns(Task.CompletedTask); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SaveRunAsync(It.IsAny())) .Returns(Task.CompletedTask); @@ -84,7 +91,7 @@ public class KisDataCollectionOrchestratorTests "IsMarketClosed should return true and cached snapshot should be used, so KIS API should not be called" ); - _repositoryMock.Verify( + _writeRepositoryMock.Verify( r => r.SaveSnapshotAsync(It.Is(s => s.SourceName.Contains("(Cached)"))), Times.Once, @@ -98,7 +105,7 @@ public class KisDataCollectionOrchestratorTests var runId = "test-run-002"; var ticker = "005930"; var account = "mock"; - _repositoryMock + _readRepositoryMock .Setup(r => r.GetLatestSnapshotsForTickerAsync(ticker, It.IsAny())) .ReturnsAsync(new List()); @@ -115,15 +122,15 @@ public class KisDataCollectionOrchestratorTests .Setup(k => k.GetDailyItemChartPriceAsync(ticker, It.IsAny(), It.IsAny(), "D", account)) .ReturnsAsync(new Dictionary()); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SaveSnapshotAsync(It.IsAny())) .Returns(Task.CompletedTask); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SavePriceHistoryDailyAsync(It.IsAny())) .Returns(Task.CompletedTask); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SaveRunAsync(It.IsAny())) .Returns(Task.CompletedTask); @@ -158,7 +165,7 @@ public class KisDataCollectionOrchestratorTests CapturedAt: $"{priorDay}T14:30:00" ); - _repositoryMock + _readRepositoryMock .Setup(r => r.GetLatestSnapshotsForTickerAsync(ticker, It.IsAny())) .ReturnsAsync(new List { priorDaySnapshot }); @@ -174,15 +181,15 @@ public class KisDataCollectionOrchestratorTests .Setup(k => k.GetDailyItemChartPriceAsync(ticker, It.IsAny(), It.IsAny(), "D", account)) .ReturnsAsync(new Dictionary()); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SaveSnapshotAsync(It.IsAny())) .Returns(Task.CompletedTask); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SavePriceHistoryDailyAsync(It.IsAny())) .Returns(Task.CompletedTask); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SaveRunAsync(It.IsAny())) .Returns(Task.CompletedTask); @@ -206,7 +213,7 @@ public class KisDataCollectionOrchestratorTests var ticker = "005930"; var account = "mock"; - _repositoryMock + _readRepositoryMock .Setup(r => r.GetLatestSnapshotsForTickerAsync(ticker, It.IsAny())) .ReturnsAsync(new List()); @@ -222,15 +229,15 @@ public class KisDataCollectionOrchestratorTests .Setup(k => k.GetDailyItemChartPriceAsync(ticker, It.IsAny(), It.IsAny(), "D", account)) .ReturnsAsync(new Dictionary()); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SaveSnapshotAsync(It.IsAny())) .Returns(Task.CompletedTask); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SavePriceHistoryDailyAsync(It.IsAny())) .Returns(Task.CompletedTask); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SaveRunAsync(It.IsAny())) .Returns(Task.CompletedTask); @@ -261,7 +268,7 @@ public class KisDataCollectionOrchestratorTests var account = "mock"; var tickers = new List { "005930", "000660" }; - _repositoryMock + _readRepositoryMock .Setup(r => r.GetLatestSnapshotsForTickerAsync(It.IsAny(), It.IsAny())) .ReturnsAsync(new List()); @@ -286,7 +293,7 @@ public class KisDataCollectionOrchestratorTests .ReturnsAsync(new Dictionary()); var callCount = 0; - _repositoryMock + _writeRepositoryMock .Setup(r => r.SaveSnapshotAsync(It.IsAny())) .Returns((CollectionSnapshotRecord snapshot) => { @@ -296,15 +303,15 @@ public class KisDataCollectionOrchestratorTests return Task.CompletedTask; }); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SavePriceHistoryDailyAsync(It.IsAny())) .Returns(Task.CompletedTask); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SaveErrorAsync(It.IsAny())) .Returns(Task.CompletedTask); - _repositoryMock + _writeRepositoryMock .Setup(r => r.SaveRunAsync(It.IsAny())) .Returns(Task.CompletedTask); @@ -315,7 +322,7 @@ public class KisDataCollectionOrchestratorTests Assert.Equal(1, result.SuccessCount); Assert.Equal(1, result.ErrorCount); - _repositoryMock.Verify( + _writeRepositoryMock.Verify( r => r.SaveErrorAsync(It.Is(e => e.Ticker == "000660" && e.ErrorMessage == "Storage Error")), Times.Once @@ -411,3 +418,5 @@ public class KisDataCollectionOrchestratorTests throw new InvalidOperationException("Repository root not found."); } } + + diff --git a/src/dotnet/QuantEngine.Core/Interfaces/IDataCollectionStore.cs b/src/dotnet/QuantEngine.Core/Interfaces/IDataCollectionStore.cs index fb11d998..7d3888b3 100644 --- a/src/dotnet/QuantEngine.Core/Interfaces/IDataCollectionStore.cs +++ b/src/dotnet/QuantEngine.Core/Interfaces/IDataCollectionStore.cs @@ -44,6 +44,25 @@ public interface IDataCollectionStore Task GetDashboardStateAsync(); } +public interface ICollectionWriteRepository +{ + Task SaveRunAsync(CollectionRunRecord run); + Task UpdateRunStatusAsync(string runId, string status, string? finishedAt = null, int? totalSnapshots = null, int? totalErrors = null); + Task SaveSnapshotAsync(CollectionSnapshotRecord snapshot); + Task SaveErrorAsync(CollectionErrorRecord error); + Task SavePriceHistoryDailyAsync(PriceHistoryDailyRecord record); +} + +public interface ICollectionReadRepository +{ + Task> GetRecentRunsAsync(int limit = 20); + Task> GetRunSnapshotsAsync(string runId); + Task> GetRunErrorsAsync(string runId, int limit = 50); + Task GetDashboardStateAsync(); + Task> GetLatestSnapshotsForTickerAsync(string ticker, int limit = 10); + Task> GetPriceHistorySummaryAsync(); +} + /// /// Collection run record (maps Python CollectionRun). /// diff --git a/src/dotnet/QuantEngine.Infrastructure/Repositories/CollectionRepository.cs b/src/dotnet/QuantEngine.Infrastructure/Repositories/CollectionRepository.cs index 7239fcb5..3c5254d5 100644 --- a/src/dotnet/QuantEngine.Infrastructure/Repositories/CollectionRepository.cs +++ b/src/dotnet/QuantEngine.Infrastructure/Repositories/CollectionRepository.cs @@ -8,7 +8,7 @@ using QuantEngine.Infrastructure.Data; namespace QuantEngine.Infrastructure.Repositories { - public class CollectionRepository : ICollectionRepository + public class CollectionRepository : ICollectionRepository, ICollectionReadRepository, ICollectionWriteRepository { private readonly IDbConnectionFactory _connectionFactory; diff --git a/src/dotnet/QuantEngine.Web/Program.cs b/src/dotnet/QuantEngine.Web/Program.cs index e81d782e..d475b413 100644 --- a/src/dotnet/QuantEngine.Web/Program.cs +++ b/src/dotnet/QuantEngine.Web/Program.cs @@ -107,7 +107,10 @@ try builder.Services.AddScoped(); builder.Services.AddScoped(); builder.Services.AddScoped(); - builder.Services.AddScoped(); + builder.Services.AddScoped(); + builder.Services.AddScoped(sp => sp.GetRequiredService()); + builder.Services.AddScoped(sp => sp.GetRequiredService()); + builder.Services.AddScoped(sp => sp.GetRequiredService()); builder.Services.AddScoped(); builder.Services.AddSingleton(); builder.Services.AddScoped();