/
AltinnServiceUpdateConsumer.cs
57 lines (48 loc) · 2.11 KB
/
AltinnServiceUpdateConsumer.cs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
using Altinn.Notifications.Core.Models.AltinnServiceUpdate;
using Altinn.Notifications.Core.Services.Interfaces;
using Altinn.Notifications.Integrations.Configuration;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
namespace Altinn.Notifications.Integrations.Kafka.Consumers
{
/// <summary>
/// Kafka consumer class for Altinn service updates
/// </summary>
public class AltinnServiceUpdateConsumer : KafkaConsumerBase<AltinnServiceUpdateConsumer>
{
private readonly IAltinnServiceUpdateService _serviceUpdate;
private readonly ILogger<AltinnServiceUpdateConsumer> _logger;
/// <summary>
/// Initializes a new instance of the <see cref="AltinnServiceUpdateConsumer"/> class.
/// </summary>
public AltinnServiceUpdateConsumer(
IAltinnServiceUpdateService serviceUpdate,
IOptions<KafkaSettings> settings,
ILogger<AltinnServiceUpdateConsumer> logger)
: base(settings, logger, settings.Value.AltinnServiceUpdateTopicName)
{
_serviceUpdate = serviceUpdate;
_logger = logger;
}
/// <inheritdoc/>
protected override Task ExecuteAsync(CancellationToken stoppingToken)
{
return Task.Run(() => ConsumeMessage(ProcessServiceUpdate, RetryServiceUpdate, stoppingToken), stoppingToken);
}
private async Task ProcessServiceUpdate(string message)
{
bool succeeded = GenericServiceUpdate.TryParse(message, out GenericServiceUpdate update);
if (!succeeded)
{
_logger.LogError("// AltinnServiceUpdateConsumer // ProcessServiceUpdate // Deserialization of message failed. {Message}", message);
return;
}
await _serviceUpdate.HandleServiceUpdate(update.Source.ToLower().Trim(), update.Schema, update.Data);
}
private async Task RetryServiceUpdate(string message)
{
// Making a second attempt, but no further action if it fails again.
await ProcessServiceUpdate(message);
}
}
}