Skip to content

Commit 054f859

Browse files
refactor: simplies code
1 parent 0dc0a55 commit 054f859

16 files changed

Lines changed: 160 additions & 115 deletions

File tree

Lines changed: 7 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,10 @@
11
using FluentValidation;
22
using FluentValidation.Results;
33

4-
using Iris.WebApi.Modules.Indicators.Features.Ingestion.Models;
5-
using Iris.WebApi.Modules.Indicators.Mappers;
64
using Iris.WebApi.Modules.Indicators.Models;
5+
using Iris.WebApi.Modules.Indicators.Repositories;
76
using Iris.WebApi.Shared.Validation;
87

9-
using StackExchange.Redis;
10-
118
namespace Iris.WebApi.Modules.Indicators.Features.GetByRange;
129

1310
public static class GetIndicatorsByRangeEndpoint
@@ -25,7 +22,7 @@ private static async Task<IResult> HandleAsync(
2522
HttpContext context,
2623
[AsParameters] GetIndicatorsByRangeRequest request,
2724
IValidator<GetIndicatorsByRangeRequest> validator,
28-
IDatabase redis)
25+
IIndicatorTimeSeriesRepository repository)
2926
{
3027
ValidationResult validationResult = await validator.ValidateAsync(request);
3128

@@ -34,19 +31,11 @@ private static async Task<IResult> HandleAsync(
3431
return Results.BadRequest(validationResult.ToApiResponse());
3532
}
3633

37-
RedisResult timeSeries = await redis.ExecuteAsync(
38-
"TS.RANGE",
39-
IndicatorConfigs.GetByCode(request.Code).RedisKey,
40-
request.From.ToUnixMilliseconds(),
41-
request.To.ToUnixMilliseconds());
42-
43-
if (timeSeries.IsNull || timeSeries.Length == 0)
44-
{
45-
return Results.NoContent();
46-
}
47-
48-
IEnumerable<Indicator> data = IndicatorMapper.Map((RedisResult[])timeSeries!);
34+
IEnumerable<Indicator> indicators = await repository
35+
.GetIndicatorsAsync(request.Code, request.From, request.To);
4936

50-
return Results.Ok(new GetIndicatorsByRangeResponse(request.Code, data));
37+
return indicators is null || !indicators.Any()
38+
? Results.NoContent()
39+
: Results.Ok(new GetIndicatorsByRangeResponse(request.Code, indicators));
5140
}
5241
}

src/Iris.WebApi/Modules/Indicators/Features/GetByRange/GetIndicatorsByRangeRequestValidator.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
using FluentValidation;
22

3-
using Iris.WebApi.Modules.Indicators.Features.Ingestion.Models;
3+
using Iris.WebApi.Modules.Indicators.Features.Ingestion;
44

55
namespace Iris.WebApi.Modules.Indicators.Features.GetByRange;
66

Lines changed: 8 additions & 75 deletions
Original file line numberDiff line numberDiff line change
@@ -1,55 +1,41 @@
1-
using System.Globalization;
2-
3-
using Iris.WebApi.Modules.Indicators.Features.Ingestion.Models;
1+
using Iris.WebApi.Modules.Indicators.Mappers;
2+
using Iris.WebApi.Modules.Indicators.Repositories;
43
using Iris.WebApi.Shared.Infra.Http.Clients.BCB;
54
using Iris.WebApi.Shared.Infra.Http.Clients.BCB.Loggers;
65
using Iris.WebApi.Shared.Infra.Http.Clients.BCB.Models;
76

87
using Refit;
98

10-
using StackExchange.Redis;
11-
129
namespace Iris.WebApi.Modules.Indicators.Features.Ingestion;
1310

1411
[AutomaticRetry(Attempts = 1)]
1512
public abstract class BaseIndicatorIngestionJob<TJob>(
1613
ILogger<TJob> logger,
1714
IBCBHttpClient httpClient,
18-
IConnectionMultiplexer redis) where TJob : BaseIndicatorIngestionJob<TJob>
15+
IIndicatorTimeSeriesRepository repository) where TJob : BaseIndicatorIngestionJob<TJob>
1916
{
2017
protected abstract IndicatorConfig Config { get; }
2118
protected readonly IBCBHttpClient _httpClient = httpClient;
2219

2320
public async Task ExecuteAsync()
2421
{
25-
IDatabase db = redis.GetDatabase();
26-
27-
if (!await db.KeyExistsAsync(Config.RedisKey))
28-
{
29-
await db.ExecuteAsync(
30-
"TS.CREATE",
31-
Config.RedisKey,
32-
"RETENTION", 0,
33-
"CHUNK_SIZE", 128,
34-
"DUPLICATE_POLICY", "LAST",
35-
"LABELS", "code", Config.Code.ToLower());
36-
}
22+
await repository.EnsureTimeSeriesExistsAsync(Config);
3723

3824
DateOnly today = DateOnly.FromDateTime(DateTime.Now);
39-
DateOnly from = await GetStartDateAsync(db, today);
25+
DateOnly from = await repository.GetNextIngestionDateAsync(Config, today);
4026
DateOnly to = today;
4127

4228
try
4329
{
44-
IEnumerable<Indicator> indicators = await GetIndicatorDataAsync(new IndicatorQueryParams(from, to));
30+
IEnumerable<RawIndicator> indicators = await GetIndicatorDataAsync(new IndicatorQueryParams(from, to));
4531

4632
if (indicators is null || !indicators.Any())
4733
{
4834
logger.LogNoDataFound(from, to);
4935
return;
5036
}
5137

52-
await AddIndicatorsAsync(db, indicators);
38+
await repository.AppendAsync(Config, IndicatorMapper.Map(indicators));
5339
}
5440
catch (ApiException ex) when (ex.StatusCode == System.Net.HttpStatusCode.NotFound)
5541
{
@@ -62,58 +48,5 @@ await db.ExecuteAsync(
6248
}
6349
}
6450

65-
protected abstract Task<IEnumerable<Indicator>> GetIndicatorDataAsync(IndicatorQueryParams queryParams);
66-
67-
private async Task<DateOnly> GetStartDateAsync(IDatabase db, DateOnly today)
68-
{
69-
try
70-
{
71-
RedisResult result = await db.ExecuteAsync("TS.GET", Config.RedisKey);
72-
73-
if (result.IsNull)
74-
{
75-
return today.AddYears(-10);
76-
}
77-
78-
RedisResult[] parts = (RedisResult[])result!;
79-
80-
if (parts == null || parts.Length == 0)
81-
{
82-
return today.AddYears(-10);
83-
}
84-
85-
long lastTimestamp = (long)parts[0];
86-
87-
DateOnly lastDate = DateOnly.FromDateTime(
88-
DateTimeOffset.FromUnixTimeMilliseconds(lastTimestamp).DateTime
89-
);
90-
91-
return lastDate.AddDays(1);
92-
}
93-
catch (Exception)
94-
{
95-
return today.AddYears(-10);
96-
}
97-
}
98-
99-
private async Task AddIndicatorsAsync(IDatabase db, IEnumerable<Indicator> indicators)
100-
{
101-
var indicatorsArray = indicators.ToArray();
102-
if (indicatorsArray.Length == 0) return;
103-
104-
var args = new object[indicatorsArray.Length * 3];
105-
int index = 0;
106-
107-
foreach (var indicator in indicatorsArray)
108-
{
109-
long timestamp = new DateTimeOffset(indicator.Date.ToDateTime(TimeOnly.MinValue))
110-
.ToUnixTimeMilliseconds();
111-
112-
args[index++] = Config.RedisKey;
113-
args[index++] = timestamp;
114-
args[index++] = indicator.Value.ToString(CultureInfo.InvariantCulture);
115-
}
116-
117-
await db.ExecuteAsync("TS.MADD", args);
118-
}
51+
protected abstract Task<IEnumerable<RawIndicator>> GetIndicatorDataAsync(IndicatorQueryParams queryParams);
11952
}

src/Iris.WebApi/Modules/Indicators/Features/Ingestion/IndicatorConfig.cs

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,4 @@
1-
using Iris.WebApi.Modules.Indicators.Features.GetByRange;
2-
3-
namespace Iris.WebApi.Modules.Indicators.Features.Ingestion.Models;
1+
namespace Iris.WebApi.Modules.Indicators.Features.Ingestion;
42

53
public record struct IndicatorConfig(
64
string Code,

src/Iris.WebApi/Modules/Indicators/Features/Ingestion/IpcaIngestionJob.cs

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,17 @@
1-
using Iris.WebApi.Modules.Indicators.Features.Ingestion.Models;
1+
using Iris.WebApi.Modules.Indicators.Repositories;
22
using Iris.WebApi.Shared.Infra.Http.Clients.BCB;
33
using Iris.WebApi.Shared.Infra.Http.Clients.BCB.Models;
44

5-
using StackExchange.Redis;
6-
75
namespace Iris.WebApi.Modules.Indicators.Features.Ingestion;
86

97
public class IpcaIngestionJob(
108
ILogger<IpcaIngestionJob> logger,
119
IBCBHttpClient httpClient,
12-
IConnectionMultiplexer redis) : BaseIndicatorIngestionJob<IpcaIngestionJob>(logger, httpClient, redis)
10+
IIndicatorTimeSeriesRepository repository) : BaseIndicatorIngestionJob<IpcaIngestionJob>(logger, httpClient, repository)
1311
{
1412
protected override IndicatorConfig Config => IndicatorConfigs.Ipca;
1513

16-
protected override async Task<IEnumerable<Indicator>> GetIndicatorDataAsync(IndicatorQueryParams queryParams)
14+
protected override async Task<IEnumerable<RawIndicator>> GetIndicatorDataAsync(IndicatorQueryParams queryParams)
1715
{
1816
return await _httpClient.GetIpcaAsync(queryParams);
1917
}

src/Iris.WebApi/Modules/Indicators/Features/Ingestion/SelicIngestionJob.cs

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,17 @@
1-
using Iris.WebApi.Modules.Indicators.Features.Ingestion.Models;
1+
using Iris.WebApi.Modules.Indicators.Repositories;
22
using Iris.WebApi.Shared.Infra.Http.Clients.BCB;
33
using Iris.WebApi.Shared.Infra.Http.Clients.BCB.Models;
44

5-
using StackExchange.Redis;
6-
75
namespace Iris.WebApi.Modules.Indicators.Features.Ingestion;
86

97
public class SelicIngestionJob(
108
ILogger<SelicIngestionJob> logger,
119
IBCBHttpClient httpClient,
12-
IConnectionMultiplexer redis) : BaseIndicatorIngestionJob<SelicIngestionJob>(logger, httpClient, redis)
10+
IIndicatorTimeSeriesRepository repository) : BaseIndicatorIngestionJob<SelicIngestionJob>(logger, httpClient, repository)
1311
{
1412
protected override IndicatorConfig Config => IndicatorConfigs.Selic;
1513

16-
protected override async Task<IEnumerable<Indicator>> GetIndicatorDataAsync(IndicatorQueryParams queryParams)
14+
protected override async Task<IEnumerable<RawIndicator>> GetIndicatorDataAsync(IndicatorQueryParams queryParams)
1715
{
1816
return await _httpClient.GetSelicAsync(queryParams);
1917
}

src/Iris.WebApi/Modules/Indicators/IoC/DependencyInjector.cs

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,12 +2,18 @@
22

33
using Iris.WebApi.Modules.Indicators.Features.GetByRange;
44
using Iris.WebApi.Modules.Indicators.Features.Ingestion;
5-
using Iris.WebApi.Modules.Indicators.Features.Ingestion.Models;
5+
using Iris.WebApi.Modules.Indicators.Repositories;
66

77
namespace Iris.WebApi.Modules.Indicators.IoC;
88

99
public static class DependencyInjector
1010
{
11+
public static IServiceCollection AddIndicatorsRepositories(this IServiceCollection services)
12+
{
13+
services.AddScoped<IIndicatorTimeSeriesRepository, RedisTimeSeriesRepository>();
14+
return services;
15+
}
16+
1117
public static IServiceCollection AddIndicatorsValidators(this IServiceCollection services)
1218
{
1319
services.AddScoped<IValidator<GetIndicatorsByRangeRequest>, GetIndicatorsByRangeRequestValidator>();

src/Iris.WebApi/Modules/Indicators/Mappers/IndicatorMapper.cs

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
using System.Globalization;
22

33
using Iris.WebApi.Modules.Indicators.Models;
4+
using Iris.WebApi.Shared.Infra.Http.Clients.BCB.Models;
45

56
using StackExchange.Redis;
67

@@ -23,4 +24,9 @@ public static IEnumerable<Indicator> Map(RedisResult[] results)
2324
return new Indicator(parsedDate, parsedValue);
2425
});
2526
}
27+
28+
public static Indicator Map(RawIndicator raw) => new(raw.Date, raw.Value);
29+
30+
public static IEnumerable<Indicator> Map(IEnumerable<RawIndicator> raw) =>
31+
raw.Select(Map);
2632
}
Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,15 @@
1+
using Iris.WebApi.Modules.Indicators.Features.Ingestion;
2+
using Iris.WebApi.Modules.Indicators.Models;
3+
4+
namespace Iris.WebApi.Modules.Indicators.Repositories;
5+
6+
public interface IIndicatorTimeSeriesRepository
7+
{
8+
Task EnsureTimeSeriesExistsAsync(IndicatorConfig config);
9+
10+
Task<DateOnly> GetNextIngestionDateAsync(IndicatorConfig config, DateOnly today);
11+
12+
Task AppendAsync(IndicatorConfig config, IEnumerable<Indicator> indicators);
13+
14+
Task<IEnumerable<Indicator>> GetIndicatorsAsync(string code, DateOnly from, DateOnly to);
15+
}
Lines changed: 101 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,101 @@
1+
using System.Globalization;
2+
3+
using Iris.WebApi.Modules.Indicators.Features.Ingestion;
4+
using Iris.WebApi.Modules.Indicators.Mappers;
5+
using Iris.WebApi.Modules.Indicators.Models;
6+
7+
using StackExchange.Redis;
8+
9+
namespace Iris.WebApi.Modules.Indicators.Repositories;
10+
11+
public sealed class RedisTimeSeriesRepository(IDatabase redis) : IIndicatorTimeSeriesRepository
12+
{
13+
public async Task EnsureTimeSeriesExistsAsync(IndicatorConfig config)
14+
{
15+
if (await redis.KeyExistsAsync(config.RedisKey))
16+
return;
17+
18+
await redis.ExecuteAsync(
19+
"TS.CREATE",
20+
config.RedisKey,
21+
"RETENTION", 0,
22+
"CHUNK_SIZE", 128,
23+
"DUPLICATE_POLICY", "LAST",
24+
"LABELS", "code", config.Code.ToLowerInvariant());
25+
}
26+
27+
public async Task<DateOnly> GetNextIngestionDateAsync(IndicatorConfig config, DateOnly today)
28+
{
29+
try
30+
{
31+
RedisResult result = await redis.ExecuteAsync("TS.GET", config.RedisKey);
32+
33+
if (result.IsNull)
34+
{
35+
return today.AddYears(-10);
36+
}
37+
38+
RedisResult[] parts = (RedisResult[])result!;
39+
40+
if (parts is null || parts.Length == 0)
41+
{
42+
return today.AddYears(-10);
43+
}
44+
45+
long lastTimestampMs = (long)parts[0];
46+
47+
DateOnly lastDate = DateOnly.FromDateTime(
48+
DateTimeOffset.FromUnixTimeMilliseconds(lastTimestampMs).DateTime);
49+
50+
return lastDate.AddDays(1);
51+
}
52+
catch
53+
{
54+
return today.AddYears(-10);
55+
}
56+
}
57+
58+
public async Task<IEnumerable<Indicator>> GetIndicatorsAsync(string code, DateOnly from, DateOnly to)
59+
{
60+
RedisResult? timeSeries = await redis.ExecuteAsync(
61+
"TS.RANGE",
62+
IndicatorConfigs.GetByCode(code).RedisKey,
63+
from.ToUnixMilliseconds(),
64+
to.ToUnixMilliseconds());
65+
66+
if (timeSeries.IsNull || timeSeries.Length == 0)
67+
{
68+
return [];
69+
}
70+
71+
IEnumerable<Indicator> data = IndicatorMapper.Map((RedisResult[])timeSeries!);
72+
73+
return data;
74+
}
75+
76+
public async Task AppendAsync(IndicatorConfig config, IEnumerable<Indicator> indicators)
77+
{
78+
Indicator[] arr = indicators as Indicator[] ?? [.. indicators];
79+
80+
if (arr.Length == 0)
81+
{
82+
return;
83+
}
84+
85+
object[] args = new object[arr.Length * 3];
86+
87+
int index = 0;
88+
89+
foreach (var indicator in arr)
90+
{
91+
long timestampMs = new DateTimeOffset(indicator.Date.ToDateTime(TimeOnly.MinValue))
92+
.ToUnixTimeMilliseconds();
93+
94+
args[index++] = config.RedisKey;
95+
args[index++] = timestampMs;
96+
args[index++] = indicator.Value.ToString(CultureInfo.InvariantCulture);
97+
}
98+
99+
await redis.ExecuteAsync("TS.MADD", args);
100+
}
101+
}

0 commit comments

Comments
 (0)