-
Notifications
You must be signed in to change notification settings - Fork 2
/
Copy pathStreamOperatorService.cs
186 lines (174 loc) · 9.03 KB
/
StreamOperatorService.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
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
using System;
using System.Threading;
using System.Threading.Tasks;
using Akka;
using Akka.Streams;
using Akka.Streams.Dsl;
using Akka.Streams.Supervision;
using Akka.Util;
using Akka.Util.Extensions;
using Arcane.Operator.Extensions;
using Arcane.Operator.Models;
using Arcane.Operator.Models.StreamClass.Base;
using Arcane.Operator.Models.StreamDefinitions.Base;
using Arcane.Operator.Services.Base;
using Arcane.Operator.Services.Models;
using k8s;
using k8s.Autorest;
using k8s.Models;
using Microsoft.Extensions.Logging;
using Snd.Sdk.ActorProviders;
using Snd.Sdk.Tasks;
namespace Arcane.Operator.Services.Operator;
public class StreamOperatorService : IStreamOperatorService
{
private const int parallelism = 1;
private readonly ILogger<StreamOperatorService> logger;
private readonly IStreamingJobOperatorService operatorService;
private readonly IStreamDefinitionRepository streamDefinitionRepository;
private readonly IStreamClass streamClass;
private readonly IMetricsReporter metricsReporter;
public StreamOperatorService(IStreamClass streamClass,
IStreamingJobOperatorService operatorService,
IStreamDefinitionRepository streamDefinitionRepository,
IMetricsReporter metricsReporter,
ILogger<StreamOperatorService> logger)
{
this.streamClass = streamClass;
this.streamDefinitionRepository = streamDefinitionRepository;
this.operatorService = operatorService;
this.logger = logger;
this.metricsReporter = metricsReporter;
}
public IRunnableGraph<Task> GetStreamDefinitionEventsGraph(CancellationToken cancellationToken)
{
var request = new CustomResourceApiRequest(
this.operatorService.StreamJobNamespace,
this.streamClass.ApiGroupRef,
this.streamClass.VersionRef,
this.streamClass.PluralNameRef
);
this.logger.LogInformation("Start listening to event stream for {@streamClass}", request);
var restartSettings = RestartSettings.Create(
TimeSpan.FromSeconds(10),
TimeSpan.FromMinutes(3),
0.2);
var eventsSource = this.streamDefinitionRepository.GetEvents(request, this.streamClass.MaxBufferCapacity)
.RecoverWithRetries(exception =>
{
if (exception is HttpOperationException { Response.StatusCode: System.Net.HttpStatusCode.NotFound })
{
this.logger.LogWarning("The resource definition {@streamClass} not found", request);
}
throw exception;
}, 1);
return RestartSource.OnFailuresWithBackoff(() => eventsSource, restartSettings)
.Via(cancellationToken.AsFlow<ResourceEvent<IStreamDefinition>>(true))
.Select(this.metricsReporter.ReportTrafficMetrics)
.SelectAsync(parallelism, this.OnEvent)
.WithAttributes(ActorAttributes.CreateSupervisionStrategy(this.HandleError))
.CollectOption()
.SelectAsync(parallelism,
response => this.streamDefinitionRepository.SetStreamStatus(response.Namespace,
response.Kind,
response.Id,
response.ToStatus()))
.WithAttributes(ActorAttributes.CreateSupervisionStrategy(this.HandleError))
.ToMaterialized(Sink.Ignore<Option<IStreamDefinition>>(), Keep.Right);
}
private Task<Option<StreamOperatorResponse>> OnEvent(ResourceEvent<IStreamDefinition> resourceEvent)
{
return resourceEvent switch
{
(WatchEventType.Added, var streamDefinition) => this.OnAdded(streamDefinition),
(WatchEventType.Modified, var streamDefinition) => this.OnModified(streamDefinition),
_ => Task.FromResult(Option<StreamOperatorResponse>.None)
};
}
private Directive HandleError(Exception exception)
{
this.logger.LogError(exception, "Failed to handle stream definition event");
return exception switch
{
BufferOverflowException => Directive.Stop,
_ => Directive.Resume
};
}
private Task<Option<StreamOperatorResponse>> OnModified(IStreamDefinition streamDefinition)
{
this.logger.LogInformation("Modified a stream definition with id {streamId}", streamDefinition.StreamId);
return this.operatorService.GetStreamingJob(streamDefinition.StreamId)
.Map(maybeJob =>
{
return maybeJob switch
{
{ HasValue: false } when streamDefinition.CrashLoopDetected
=> Task.FromResult(StreamOperatorResponse.CrashLoopDetected(streamDefinition.Namespace(),
streamDefinition.Kind,
streamDefinition.StreamId)
.AsOption()),
{ HasValue: true } when streamDefinition.CrashLoopDetected
=> this.operatorService.DeleteJob(streamDefinition.Kind, streamDefinition.StreamId),
{ HasValue: false } when streamDefinition.ReloadRequested
=> this.streamDefinitionRepository
.RemoveReloadingAnnotation(streamDefinition.Namespace(), streamDefinition.Kind,
streamDefinition.StreamId)
.Map(sd => sd.HasValue
? this.operatorService.StartRegisteredStream(sd.Value, true, this.streamClass)
: Task.FromResult(Option<StreamOperatorResponse>.None))
.Flatten(),
{ HasValue: true } when streamDefinition.ReloadRequested
=> this.streamDefinitionRepository
.RemoveReloadingAnnotation(streamDefinition.Namespace(), streamDefinition.Kind,
streamDefinition.StreamId)
.Map(sd => sd.HasValue
? this.operatorService.RequestStreamingJobReload(streamDefinition.StreamId)
: Task.FromResult(Option<StreamOperatorResponse>.None))
.Flatten(),
{ HasValue: true } when streamDefinition.Suspended
=> this.operatorService.DeleteJob(streamDefinition.Kind, streamDefinition.StreamId),
{ HasValue: false } when streamDefinition.Suspended
=> Task.FromResult(StreamOperatorResponse.Suspended(streamDefinition.Namespace(),
streamDefinition.Kind,
streamDefinition.StreamId)
.AsOption()),
{ Value: var job } when job.GetConfigurationChecksum() ==
streamDefinition.GetConfigurationChecksum()
=> Task.FromResult(Option<StreamOperatorResponse>.None),
{ Value: var job } when !string.IsNullOrEmpty(job.GetConfigurationChecksum()) &&
job.GetConfigurationChecksum() !=
streamDefinition.GetConfigurationChecksum()
=> this.operatorService.RequestStreamingJobRestart(streamDefinition.StreamId),
{ HasValue: false }
=> this.operatorService.StartRegisteredStream(streamDefinition, false, this.streamClass),
_ => Task.FromResult(Option<StreamOperatorResponse>.None)
};
}).Flatten();
}
private Task<Option<StreamOperatorResponse>> OnAdded(IStreamDefinition streamDefinition)
{
this.logger.LogInformation("Added a stream definition with id {streamId}", streamDefinition.StreamId);
return streamDefinition.Suspended
? Task.FromResult(StreamOperatorResponse.Suspended(
streamDefinition.Namespace(),
streamDefinition.Kind,
streamDefinition.StreamId).AsOption())
: this.operatorService.GetStreamingJob(streamDefinition.StreamId)
.Map(maybeJob => maybeJob switch
{
{ HasValue: true, Value: var job } when job.IsReloading()
=> Task.FromResult(StreamOperatorResponse.Reloading(
streamDefinition.Metadata.Namespace(),
streamDefinition.Kind,
streamDefinition.StreamId)
.AsOption()),
{ HasValue: true, Value: var job } when !job.IsReloading()
=> Task.FromResult(StreamOperatorResponse.Running(
streamDefinition.Metadata.Namespace(),
streamDefinition.Kind,
streamDefinition.StreamId)
.AsOption()),
{ HasValue: false } => this.operatorService.StartRegisteredStream(streamDefinition, true, this.streamClass)
}).Flatten();
}
}