Skip to content

Commit 9390940

Browse files
committed
Removed TimeoutWatchdog
1 parent 392671a commit 9390940

67 files changed

Lines changed: 610 additions & 1860 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

Core/Cleipnir.ResilientFunctions.Tests/InMemoryTests/LeaseUpdaterTests/LeaseUpdaterTestFunctionStore.cs

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@ public class LeaseUpdaterTestFunctionStore : IFunctionStore
2020
public ITypeStore TypeStore => _inner.TypeStore;
2121
public IMessageStore MessageStore => _inner.MessageStore;
2222
public IEffectsStore EffectsStore => _inner.EffectsStore;
23-
public ITimeoutStore TimeoutStore => _inner.TimeoutStore;
2423
public ICorrelationStore CorrelationStore => _inner.CorrelationStore;
2524
public Utilities Utilities => _inner.Utilities;
2625
public IMigrator Migrator => _inner.Migrator;

Core/Cleipnir.ResilientFunctions.Tests/InMemoryTests/RFunctionTests/SuspensionTests.cs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -95,4 +95,8 @@ public override Task InterruptSuspendedFlows()
9595
[TestMethod]
9696
public override Task AwaitMessageAfterAppendShouldNotCauseSuspension()
9797
=> AwaitMessageAfterAppendShouldNotCauseSuspension(FunctionStoreFactory.Create());
98+
99+
[TestMethod]
100+
public override Task DelayedFlowIsRestartedOnce()
101+
=> DelayedFlowIsRestartedOnce(FunctionStoreFactory.Create());
98102
}

Core/Cleipnir.ResilientFunctions.Tests/InMemoryTests/TimeoutStoreTests.cs

Lines changed: 0 additions & 46 deletions
This file was deleted.

Core/Cleipnir.ResilientFunctions.Tests/Messaging/TestTemplates/CustomMessageSerializerTests.cs

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -34,8 +34,15 @@ await functionStore.CreateFunction(
3434
var messagesWriter = new MessageWriter(storedId, functionStore, eventSerializer, scheduleReInvocation: (_, _) => Task.CompletedTask);
3535
var lazyExistingEffects = new Lazy<Task<IReadOnlyList<StoredEffect>>>(() => Task.FromResult((IReadOnlyList<StoredEffect>) new List<StoredEffect>()));
3636
var effectResults = new EffectResults(flowId, storedId, lazyExistingEffects, functionStore.EffectsStore, DefaultSerializer.Instance);
37-
var effect = new Effect(effectResults, utcNow: () => DateTime.UtcNow);
38-
var registeredTimeouts = new RegisteredTimeouts(storedId, functionStore.TimeoutStore, effect, () => DateTime.UtcNow);
37+
var minimumTimeout = new FlowMinimumTimeout();
38+
var effect = new Effect(effectResults, utcNow: () => DateTime.UtcNow, minimumTimeout);
39+
var registeredTimeouts = new FlowRegisteredTimeouts(
40+
effect,
41+
utcNow: () => DateTime.UtcNow,
42+
minimumTimeout,
43+
publishTimeoutEvent: t => messagesWriter.AppendMessage(t),
44+
unhandledExceptionHandler: new UnhandledExceptionHandler(_ => {}),
45+
flowId);
3946
var messagesPullerAndEmitter = new MessagesPullerAndEmitter(
4047
storedId,
4148
defaultDelay: TimeSpan.FromSeconds(1),

Core/Cleipnir.ResilientFunctions.Tests/Messaging/TestTemplates/MessagesTests.cs

Lines changed: 83 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -35,7 +35,15 @@ await functionStore.CreateFunction(
3535
owner: null
3636
);
3737
var messagesWriter = new MessageWriter(storedId, functionStore, DefaultSerializer.Instance, scheduleReInvocation: (_, _) => Task.CompletedTask);
38-
var registeredTimeouts = new RegisteredTimeouts(storedId, functionStore.TimeoutStore, CreateEffect(storedId, flowId, functionStore), () => DateTime.UtcNow);
38+
var minimumTimeout = new FlowMinimumTimeout();
39+
using var registeredTimeouts = new FlowRegisteredTimeouts(
40+
CreateEffect(storedId, flowId, functionStore, minimumTimeout),
41+
() => DateTime.UtcNow,
42+
minimumTimeout,
43+
t => messagesWriter.AppendMessage(t),
44+
new UnhandledExceptionHandler(_ => {}),
45+
flowId
46+
);
3947
var messagesPullerAndEmitter = new MessagesPullerAndEmitter(
4048
storedId,
4149
defaultDelay: TimeSpan.FromMilliseconds(250),
@@ -103,7 +111,15 @@ await functionStore.CreateFunction(
103111
owner: null
104112
);
105113
var messagesWriter = new MessageWriter(storedId, functionStore, DefaultSerializer.Instance, scheduleReInvocation: (_, _) => Task.CompletedTask);
106-
var registeredTimeouts = new RegisteredTimeouts(storedId, functionStore.TimeoutStore, CreateEffect(storedId, flowId, functionStore), () => DateTime.UtcNow);
114+
var minimumTimeout = new FlowMinimumTimeout();
115+
using var registeredTimeouts = new FlowRegisteredTimeouts(
116+
CreateEffect(storedId, flowId, functionStore, minimumTimeout),
117+
() => DateTime.UtcNow,
118+
minimumTimeout,
119+
t => messagesWriter.AppendMessage(t),
120+
new UnhandledExceptionHandler(_ => {}),
121+
flowId
122+
);
107123
var messagesPullerAndEmitter = new MessagesPullerAndEmitter(
108124
storedId,
109125
defaultDelay: TimeSpan.FromMilliseconds(250),
@@ -145,7 +161,15 @@ await functionStore.CreateFunction(
145161
owner: null
146162
);
147163
var messagesWriter = new MessageWriter(storedId, functionStore, DefaultSerializer.Instance, scheduleReInvocation: (_, _) => Task.CompletedTask);
148-
var registeredTimeouts = new RegisteredTimeouts(storedId, functionStore.TimeoutStore, CreateEffect(storedId, flowId, functionStore), () => DateTime.UtcNow);
164+
var minimumTimeout = new FlowMinimumTimeout();
165+
using var registeredTimeouts = new FlowRegisteredTimeouts(
166+
CreateEffect(storedId, flowId, functionStore, minimumTimeout),
167+
() => DateTime.UtcNow,
168+
minimumTimeout,
169+
t => messagesWriter.AppendMessage(t),
170+
new UnhandledExceptionHandler(_ => {}),
171+
flowId
172+
);
149173
var messagesPullerAndEmitter = new MessagesPullerAndEmitter(
150174
storedId,
151175
defaultDelay: TimeSpan.FromMilliseconds(250),
@@ -187,7 +211,15 @@ await functionStore.CreateFunction(
187211
owner: null
188212
);
189213
var messagesWriter = new MessageWriter(storedId, functionStore, DefaultSerializer.Instance, scheduleReInvocation: (_, _) => Task.CompletedTask);
190-
var registeredTimeouts = new RegisteredTimeouts(storedId, functionStore.TimeoutStore, CreateEffect(storedId, flowId, functionStore), () => DateTime.UtcNow);
214+
var minimumTimeout = new FlowMinimumTimeout();
215+
using var registeredTimeouts = new FlowRegisteredTimeouts(
216+
CreateEffect(storedId, flowId, functionStore, minimumTimeout),
217+
() => DateTime.UtcNow,
218+
minimumTimeout,
219+
t => messagesWriter.AppendMessage(t),
220+
new UnhandledExceptionHandler(_ => {}),
221+
flowId
222+
);
191223
var messagesPullerAndEmitter = new MessagesPullerAndEmitter(
192224
storedId,
193225
defaultDelay: TimeSpan.FromMilliseconds(250),
@@ -231,7 +263,15 @@ await functionStore.CreateFunction(
231263
owner: null
232264
);
233265
var messagesWriter = new MessageWriter(storedId, functionStore, DefaultSerializer.Instance, scheduleReInvocation: (_, _) => Task.CompletedTask);
234-
var registeredTimeouts = new RegisteredTimeouts(storedId, functionStore.TimeoutStore, CreateEffect(storedId, flowId, functionStore), () => DateTime.UtcNow);
266+
var minimumTimeout = new FlowMinimumTimeout();
267+
using var registeredTimeouts = new FlowRegisteredTimeouts(
268+
CreateEffect(storedId, flowId, functionStore, minimumTimeout),
269+
() => DateTime.UtcNow,
270+
minimumTimeout,
271+
t => messagesWriter.AppendMessage(t),
272+
new UnhandledExceptionHandler(_ => {}),
273+
flowId
274+
);
235275
var messagesPullerAndEmitter = new MessagesPullerAndEmitter(
236276
storedId,
237277
defaultDelay: TimeSpan.FromMilliseconds(250),
@@ -280,7 +320,15 @@ await functionStore.CreateFunction(
280320
owner: null
281321
);
282322
var messagesWriter = new MessageWriter(storedId, functionStore, DefaultSerializer.Instance, scheduleReInvocation: (_, _) => Task.CompletedTask);
283-
var registeredTimeouts = new RegisteredTimeouts(storedId, functionStore.TimeoutStore, CreateEffect(storedId, flowId, functionStore), () => DateTime.UtcNow);
323+
var minimumTimeout = new FlowMinimumTimeout();
324+
using var registeredTimeouts = new FlowRegisteredTimeouts(
325+
CreateEffect(storedId, flowId, functionStore, minimumTimeout),
326+
() => DateTime.UtcNow,
327+
minimumTimeout,
328+
t => messagesWriter.AppendMessage(t),
329+
new UnhandledExceptionHandler(_ => {}),
330+
flowId
331+
);
284332
var messagesPullerAndEmitter = new MessagesPullerAndEmitter(
285333
storedId,
286334
defaultDelay: TimeSpan.FromMilliseconds(250),
@@ -328,7 +376,15 @@ await functionStore.CreateFunction(
328376
owner: null
329377
);
330378
var messagesWriter = new MessageWriter(storedId, functionStore, DefaultSerializer.Instance, scheduleReInvocation: (_, _) => Task.CompletedTask);
331-
var registeredTimeouts = new RegisteredTimeouts(storedId, functionStore.TimeoutStore, CreateEffect(storedId, flowId, functionStore), () => DateTime.UtcNow);
379+
var minimumTimeout = new FlowMinimumTimeout();
380+
using var registeredTimeouts = new FlowRegisteredTimeouts(
381+
CreateEffect(storedId, flowId, functionStore, minimumTimeout),
382+
() => DateTime.UtcNow,
383+
minimumTimeout,
384+
t => messagesWriter.AppendMessage(t),
385+
new UnhandledExceptionHandler(_ => {}),
386+
flowId
387+
);
332388
var messagesPullerAndEmitter = new MessagesPullerAndEmitter(
333389
storedId,
334390
defaultDelay: TimeSpan.FromMilliseconds(250),
@@ -372,7 +428,15 @@ await functionStore.CreateFunction(
372428
owner: null
373429
);
374430
var messagesWriter = new MessageWriter(storedId, functionStore, DefaultSerializer.Instance, scheduleReInvocation: (_, _) => Task.CompletedTask);
375-
var registeredTimeouts = new RegisteredTimeouts(storedId, functionStore.TimeoutStore, CreateEffect(storedId, flowId, functionStore), () => DateTime.UtcNow);
431+
var minimumTimeout = new FlowMinimumTimeout();
432+
using var registeredTimeouts = new FlowRegisteredTimeouts(
433+
CreateEffect(storedId, flowId, functionStore, minimumTimeout),
434+
() => DateTime.UtcNow,
435+
minimumTimeout,
436+
t => messagesWriter.AppendMessage(t),
437+
new UnhandledExceptionHandler(_ => {}),
438+
flowId
439+
);
376440
var messagesPullerAndEmitter = new MessagesPullerAndEmitter(
377441
storedId,
378442
defaultDelay: TimeSpan.FromMilliseconds(250),
@@ -428,7 +492,15 @@ await functionStore.CreateFunction(
428492
owner: null
429493
);
430494
var messagesWriter = new MessageWriter(storedId, functionStore, DefaultSerializer.Instance, scheduleReInvocation: (_, _) => Task.CompletedTask);
431-
var registeredTimeouts = new RegisteredTimeouts(storedId, functionStore.TimeoutStore, CreateEffect(storedId, flowId, functionStore), () => DateTime.UtcNow);
495+
var minimumTimeout = new FlowMinimumTimeout();
496+
using var registeredTimeouts = new FlowRegisteredTimeouts(
497+
CreateEffect(storedId, flowId, functionStore, minimumTimeout),
498+
() => DateTime.UtcNow,
499+
minimumTimeout,
500+
t => messagesWriter.AppendMessage(t),
501+
new UnhandledExceptionHandler(_ => {}),
502+
flowId
503+
);
432504
var messagesPullerAndEmitter = new MessagesPullerAndEmitter(
433505
storedId,
434506
defaultDelay: TimeSpan.FromMilliseconds(250),
@@ -587,11 +659,11 @@ async Task<string> (_, workflow) => (await workflow.Messages.First(maxWait: Time
587659
result.ShouldBe("Hallo World!");
588660
}
589661

590-
private Effect CreateEffect(StoredId storedId, FlowId flowId, IFunctionStore functionStore)
662+
private Effect CreateEffect(StoredId storedId, FlowId flowId, IFunctionStore functionStore, FlowMinimumTimeout flowMinimumTimeout)
591663
{
592664
var lazyExistingEffects = new Lazy<Task<IReadOnlyList<StoredEffect>>>(() => Task.FromResult((IReadOnlyList<StoredEffect>) new List<StoredEffect>()));
593665
var effectResults = new EffectResults(flowId, storedId, lazyExistingEffects, functionStore.EffectsStore, DefaultSerializer.Instance);
594-
var effect = new Effect(effectResults, utcNow: () => DateTime.UtcNow);
666+
var effect = new Effect(effectResults, utcNow: () => DateTime.UtcNow, flowMinimumTimeout);
595667
return effect;
596668
}
597669

Core/Cleipnir.ResilientFunctions.Tests/ReactiveTests/NoOpTimeoutProvider.cs

Lines changed: 10 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -10,16 +10,22 @@ namespace Cleipnir.ResilientFunctions.Tests.ReactiveTests;
1010
public class NoOpRegisteredTimeouts : IRegisteredTimeouts
1111
{
1212
public static NoOpRegisteredTimeouts Instance { get; } = new();
13-
public Task RegisterTimeout(EffectId timeoutId, DateTime expiresAt)
14-
=> Task.CompletedTask;
1513

16-
public Task RegisterTimeout(EffectId timeoutId, TimeSpan expiresIn)
17-
=> Task.CompletedTask;
14+
public Task<Tuple<TimeoutStatus, DateTime>> RegisterTimeout(EffectId timeoutId, DateTime expiresAt, bool publishMessage)
15+
=> Tuple.Create(TimeoutStatus.Registered, expiresAt).ToTask();
1816

17+
public Task<Tuple<TimeoutStatus, DateTime>> RegisterTimeout(EffectId timeoutId, TimeSpan expiresIn, bool publishMessage)
18+
=> Tuple.Create(TimeoutStatus.Registered, DateTime.UtcNow.Add(expiresIn)).ToTask();
19+
1920
public Task CancelTimeout(EffectId timeoutId)
2021
=> Task.CompletedTask;
2122

23+
public Task CompleteTimeout(EffectId timeoutId)
24+
=> Task.CompletedTask;
25+
2226
public Task<IReadOnlyList<RegisteredTimeout>> PendingTimeouts() => new List<RegisteredTimeout>()
2327
.CastTo<IReadOnlyList<RegisteredTimeout>>()
2428
.ToTask();
29+
30+
public void Dispose() { }
2531
}

Core/Cleipnir.ResilientFunctions.Tests/ReactiveTests/ReactiveIntegrationTests.cs

Lines changed: 1 addition & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,5 @@
1-
using System;
2-
using System.Threading.Tasks;
3-
using Cleipnir.ResilientFunctions.Domain;
4-
using Cleipnir.ResilientFunctions.Domain.Exceptions;
5-
using Cleipnir.ResilientFunctions.Helpers;
1+
using System.Threading.Tasks;
62
using Cleipnir.ResilientFunctions.Reactive.Extensions;
7-
using Cleipnir.ResilientFunctions.Storage;
83
using Cleipnir.ResilientFunctions.Tests.Utils;
94
using Microsoft.VisualStudio.TestTools.UnitTesting;
105
using Shouldly;
@@ -14,28 +9,6 @@ namespace Cleipnir.ResilientFunctions.Tests.ReactiveTests;
149
[TestClass]
1510
public class ReactiveIntegrationTests
1611
{
17-
[TestMethod]
18-
public async Task FunctionCanBeSuspendedForASecondSuccessfully()
19-
{
20-
var store = new InMemoryFunctionStore();
21-
using var functionsRegistry = new FunctionsRegistry(store);
22-
var functionId = TestFlowId.Create();
23-
var (flowType, flowInstance) = functionId;
24-
var rAction = functionsRegistry.RegisterAction<string>(
25-
flowType,
26-
inner: async (_, workflow) =>
27-
{
28-
var messages = workflow.Messages;
29-
await messages.SuspendFor(timeoutEventId: "timeout", resumeAfter: TimeSpan.FromSeconds(1));
30-
});
31-
32-
await Should.ThrowAsync<InvocationSuspendedException>(rAction.Invoke(flowInstance.Value, "param"));
33-
34-
await BusyWait.Until(() =>
35-
store.GetFunction(rAction.MapToStoredId(functionId.Instance)).SelectAsync(sf => sf?.Status == Status.Succeeded)
36-
);
37-
}
38-
3912
[TestMethod]
4013
public async Task SyncingStopsAfterReactiveChainCompletion()
4114
{

Core/Cleipnir.ResilientFunctions.Tests/ReactiveTests/TimeoutTests.cs

Lines changed: 9 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -120,8 +120,6 @@ public async Task ExistingTimeoutEventInMessagesAvoidRegisteredTimeoutsCancellat
120120

121121
var task = await source.TakeUntilTimeout(timeoutId, expiresAt).FirstOrNone();
122122
task.HasValue.ShouldBeFalse();
123-
124-
registeredTimeoutsStub.Cancelled.ShouldBeNull();
125123
}
126124

127125
private class RegisteredTimeoutsStub : IRegisteredTimeouts
@@ -142,24 +140,28 @@ public List<Tuple<EffectId, DateTime>> Registrations
142140
private readonly Lock _sync = new();
143141
private readonly Dictionary<EffectId, DateTime> _registrations = new();
144142

145-
public Task RegisterTimeout(EffectId timeoutId, DateTime expiresAt)
143+
public Task<Tuple<TimeoutStatus, DateTime>> RegisterTimeout(EffectId timeoutId, DateTime expiresAt, bool publishMessage)
146144
{
147145
lock (_sync)
148146
_registrations[timeoutId] = expiresAt;
149147

150-
return Task.CompletedTask;
148+
return Tuple.Create(TimeoutStatus.Registered, expiresAt).ToTask();
151149
}
152150

153-
public Task RegisterTimeout(EffectId timeoutId, TimeSpan expiresIn)
154-
=> RegisterTimeout(timeoutId, DateTime.UtcNow.Add(expiresIn));
151+
public Task<Tuple<TimeoutStatus, DateTime>> RegisterTimeout(EffectId timeoutId, TimeSpan expiresIn, bool publishMessage)
152+
=> RegisterTimeout(timeoutId, DateTime.UtcNow.Add(expiresIn), publishMessage);
155153

156154
public Task CancelTimeout(EffectId timeoutId)
157155
{
158156
Cancelled = timeoutId;
159157
return Task.CompletedTask;
160158
}
161-
159+
160+
public Task CompleteTimeout(EffectId timeoutId) => Task.CompletedTask;
161+
162162
public Task<IReadOnlyList<RegisteredTimeout>> PendingTimeouts()
163163
=> Task.FromException<IReadOnlyList<RegisteredTimeout>>(new Exception("Stub-method invocation"));
164+
165+
public void Dispose() { }
164166
}
165167
}

0 commit comments

Comments
 (0)