Skip to content

Commit 0dcf9c8

Browse files
Enhance Migration Engine and Unit Tests
- Improved error logging in MigrationEngine to provide clearer context on exceptions. - Added functionality to InMemoryMigrationSource to simulate shutdown behavior during body reads. - Enhanced InMemoryMigrationTarget to allow custom exceptions for failure scenarios. - Introduced TimerRecordingTimeProvider to facilitate testing of timer-related functionality. - Expanded HaltThresholdTests to cover edge cases involving mixed benign and fault skips. - Updated MigrationEngineBodyRetryTests to ensure proper handling of shutdowns during body reads. - Added tests to verify behavior when categories are empty or already completed. - Implemented checks for checkpoint conflicts and cancellation scenarios in MigrationEngineFailurePathTests. - Enhanced MigrationEngineHaltTests to evaluate mixed skip scenarios and their impact on halting. - Improved MigrationEngineOrderingTests to ensure blocked categories behave correctly. - Added tests to MigrationEngineResumeTests to verify that restart behavior maintains original start time. - Updated MigrationEngineRunCategoriesTests to handle cases where no categories are selected. - Enhanced MigrationEngineSkipReasonTests to ensure accurate recording of multiple skip reasons. - Added tests to MigrationEngineThrottleTests to verify behavior during pauses and shutdowns.
1 parent 4bc951e commit 0dcf9c8

14 files changed

Lines changed: 608 additions & 35 deletions

‎src/ServiceControl.Persistence/DataMigration/MigrationEngine.cs‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -159,7 +159,8 @@ public async Task<MigrationCheckpoint> RunCategoryAsync(MigrationCategory catego
159159
}
160160
catch (Exception ex)
161161
{
162-
var reason = $"{ex.GetType().Name} at cursor {checkpoint.Cursor ?? "the start"}: {ex.Message}";
162+
var position = checkpoint.Cursor is null ? "at the start" : $"at cursor {checkpoint.Cursor}";
163+
var reason = $"{ex.GetType().Name} {position}: {ex.Message}";
163164
logger.LogError(ex, "Category {CategoryId} halted at cursor {Cursor}", category.Id, checkpoint.Cursor);
164165
return await Settle(checkpoint with { State = MigrationCategoryState.Halted, LastError = reason }, cancellationToken);
165166
}

‎src/ServiceControl.UnitTests/Migration/Fakes/InMemoryMigrationSource.cs‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,9 @@ public sealed class InMemoryMigrationSource : IMigrationSource
2424

2525
public void FailBodyReads(string sourceId, int times, Exception failure) => bodyFailures[sourceId] = (times, failure);
2626

27+
/// <summary>Makes the body read for this row behave like a host shutting down: the token is cancelled and the read throws.</summary>
28+
public (string SourceId, CancellationTokenSource Source)? StopOnBodyRead { get; set; }
29+
2730
public int BodyReadAttempts(string sourceId) => bodyReadAttempts.GetValueOrDefault(sourceId);
2831

2932
public Task Open(CancellationToken cancellationToken = default) => Task.CompletedTask;
@@ -70,6 +73,13 @@ public async IAsyncEnumerable<MigrationBatch> Read(
7073
await Task.Yield();
7174

7275
var attempt = bodyReadAttempts[sourceId] = BodyReadAttempts(sourceId) + 1;
76+
77+
if (StopOnBodyRead is { } stop && stop.SourceId == sourceId)
78+
{
79+
await stop.Source.CancelAsync();
80+
throw new OperationCanceledException(stop.Source.Token);
81+
}
82+
7383
if (bodyFailures.TryGetValue(sourceId, out var failures) && attempt <= failures.Times)
7484
{
7585
throw failures.Failure;

‎src/ServiceControl.UnitTests/Migration/Fakes/InMemoryMigrationTarget.cs‎

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,12 @@ public sealed class InMemoryMigrationTarget(IMigrationCheckpointStore checkpoint
1616
public int DefaultBatchSize { get; set; } = 3;
1717
public string NoBatchSizeFor { get; set; }
1818
public int? FailOnCallNumber { get; set; }
19+
20+
/// <summary>What FailOnCallNumber throws, when the default simulated failure is the wrong shape for the test.</summary>
21+
public Exception FailWith { get; set; }
22+
23+
/// <summary>Cancels the token on this call and then writes normally, so the stop surfaces from the source's next batch.</summary>
24+
public (int CallNumber, CancellationTokenSource Source)? CancelOnCall { get; set; }
1925
public (int CallNumber, CancellationTokenSource Source)? StopOnCall { get; set; }
2026
int callCount;
2127

@@ -45,7 +51,12 @@ public async Task<MigrationWriteResult> Write(
4551

4652
if (FailOnCallNumber == callCount)
4753
{
48-
throw new InvalidOperationException($"Simulated failure on write {callCount}");
54+
throw FailWith ?? new InvalidOperationException($"Simulated failure on write {callCount}");
55+
}
56+
57+
if (CancelOnCall is { } cancel && cancel.CallNumber == callCount)
58+
{
59+
await cancel.Source.CancelAsync();
4960
}
5061

5162
var keys = writtenKeysByCategory.TryGetValue(category.Id, out var existingKeys) ? existingKeys : writtenKeysByCategory[category.Id] = [];
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
#nullable enable
2+
namespace ServiceControl.UnitTests.Migration.Fakes;
3+
4+
using System;
5+
using System.Collections.Generic;
6+
using System.Threading;
7+
using Microsoft.Extensions.Time.Testing;
8+
9+
// Signals each timer the engine creates, so a test only advances the clock once a pause is waiting on it.
10+
public sealed class TimerRecordingTimeProvider : TimeProvider
11+
{
12+
readonly FakeTimeProvider clock = new();
13+
14+
public SemaphoreSlim TimerCreated { get; } = new(0);
15+
16+
public List<TimeSpan> DueTimes { get; } = [];
17+
18+
public void Advance(TimeSpan delta) => clock.Advance(delta);
19+
20+
public override DateTimeOffset GetUtcNow() => clock.GetUtcNow();
21+
22+
public override long GetTimestamp() => clock.GetTimestamp();
23+
24+
public override long TimestampFrequency => clock.TimestampFrequency;
25+
26+
public override TimeZoneInfo LocalTimeZone => clock.LocalTimeZone;
27+
28+
public override ITimer CreateTimer(TimerCallback callback, object? state, TimeSpan dueTime, TimeSpan period)
29+
{
30+
var timer = clock.CreateTimer(callback, state, dueTime, period);
31+
DueTimes.Add(dueTime);
32+
TimerCreated.Release();
33+
return timer;
34+
}
35+
}

‎src/ServiceControl.UnitTests/Migration/HaltThresholdTests.cs‎

Lines changed: 29 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,12 +36,40 @@ public void Does_not_halt_a_large_category_with_only_a_small_proportion_skipped(
3636
[Test]
3737
public void Exactly_at_both_boundaries_does_not_halt_because_both_must_be_exceeded()
3838
{
39-
// 100 of 2,000 is exactly 5% and exactly the floor. "Exceed" means strictly past, not "at".
39+
// Exactly the floor, which settles it before the proportion is ever worked out. "Exceed" means strictly past, not "at".
4040
var exceeded = HaltThreshold.Exceeded(skippedCount: 100, totalCount: 2_000, percentThreshold: 5, minimumFloor: 100);
4141

4242
Assert.That(exceeded, Is.False);
4343
}
4444

45+
[Test]
46+
public void Exactly_on_the_proportion_does_not_halt_once_the_floor_is_behind_it()
47+
{
48+
// 101 of 2,020 is exactly 5% with the floor already passed, so this is the only shape that
49+
// reaches the proportion comparison and depends on it being strictly greater.
50+
var exceeded = HaltThreshold.Exceeded(skippedCount: 101, totalCount: 2_020, percentThreshold: 5, minimumFloor: 100);
51+
52+
Assert.That(exceeded, Is.False);
53+
}
54+
55+
[Test]
56+
public void A_hair_past_the_proportion_halts_once_the_floor_is_behind_it()
57+
{
58+
// 101 of 2,000 is 5.05%: the same skip count as above, one row's worth over the line.
59+
var exceeded = HaltThreshold.Exceeded(skippedCount: 101, totalCount: 2_000, percentThreshold: 5, minimumFloor: 100);
60+
61+
Assert.That(exceeded, Is.True);
62+
}
63+
64+
[Test]
65+
public void A_skip_count_with_nothing_processed_never_halts_and_never_divides_by_zero()
66+
{
67+
// Past the floor with a zero total, which is the only input that reaches the division guard.
68+
var exceeded = HaltThreshold.Exceeded(skippedCount: 101, totalCount: 0, percentThreshold: 5, minimumFloor: 100);
69+
70+
Assert.That(exceeded, Is.False);
71+
}
72+
4573
[Test]
4674
public void No_rows_processed_yet_never_halts()
4775
{

‎src/ServiceControl.UnitTests/Migration/MigrationEngineBodyRetryTests.cs‎

Lines changed: 82 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ namespace ServiceControl.UnitTests.Migration;
44
using System;
55
using System.Collections.Generic;
66
using System.Linq;
7+
using System.Threading;
78
using System.Threading.Tasks;
89
using Microsoft.Extensions.Logging.Abstractions;
910
using Microsoft.Extensions.Time.Testing;
@@ -43,13 +44,25 @@ public async Task A_body_that_fails_until_the_last_attempt_then_succeeds_is_writ
4344
}
4445
}
4546

46-
[Test]
47-
public async Task A_body_read_that_fails_as_a_defect_halts_the_category_on_the_first_attempt_without_skipping_the_message()
47+
// Every type the engine treats as a defect. Each fails the same way on every attempt, so retrying one
48+
// would only turn a code fault into skipped messages.
49+
static readonly Exception[] Defects =
50+
[
51+
new NotSupportedException("this source cannot read bodies"),
52+
new NotImplementedException("not written yet"),
53+
new InvalidOperationException("the session is closed"),
54+
new ArgumentException("the id is not a document id"),
55+
new NullReferenceException("no attachment"),
56+
new InvalidCastException("not an attachment")
57+
];
58+
59+
[TestCaseSource(nameof(Defects))]
60+
public async Task A_body_read_that_fails_as_a_defect_halts_the_category_on_the_first_attempt_without_skipping_the_message(Exception defect)
4861
{
4962
var category = MigrationCategoryRegistry.Find("UnresolvedAndRetryIssuedFailedMessages")!;
5063
var source = new InMemoryMigrationSource();
5164
source.Seed(category.Id, Row("msg-1"));
52-
source.FailBodyReads("msg-1", MigrationEngine.MaxBodyReadAttempts, new NotSupportedException("this source cannot read bodies"));
65+
source.FailBodyReads("msg-1", MigrationEngine.MaxBodyReadAttempts, defect);
5366
var checkpointStore = new InMemoryMigrationCheckpointStore();
5467
var target = new InMemoryMigrationTarget(checkpointStore);
5568
var options = new MigrationEngineOptions(TimeSpan.Zero, 5, 100, []) { BodyRetryBackoff = TimeSpan.Zero };
@@ -60,13 +73,78 @@ public async Task A_body_read_that_fails_as_a_defect_halts_the_category_on_the_f
6073
using (Assert.EnterMultipleScope())
6174
{
6275
Assert.That(checkpoint.State, Is.EqualTo(MigrationCategoryState.Halted));
63-
Assert.That(checkpoint.LastError, Does.Contain(nameof(NotSupportedException)));
76+
Assert.That(checkpoint.LastError, Does.Contain(defect.GetType().Name));
6477
Assert.That(source.BodyReadAttempts("msg-1"), Is.EqualTo(1));
6578
Assert.That(checkpoint.SkippedCount, Is.Zero);
6679
Assert.That(checkpoint.SkipReasons, Is.Null);
6780
}
6881
}
6982

83+
[Test]
84+
public async Task A_shutdown_during_a_body_read_stops_the_run_instead_of_skipping_the_message()
85+
{
86+
// Retrying a shutdown to the attempt limit and then recording the message as permanently
87+
// unreadable is the one path here that silently loses a customer's failed message.
88+
var category = MigrationCategoryRegistry.Find("UnresolvedAndRetryIssuedFailedMessages")!;
89+
var source = new InMemoryMigrationSource();
90+
source.Seed(category.Id, Row("msg-1"));
91+
using var stopping = new CancellationTokenSource();
92+
source.StopOnBodyRead = ("msg-1", stopping);
93+
var checkpointStore = new InMemoryMigrationCheckpointStore();
94+
var target = new InMemoryMigrationTarget(checkpointStore);
95+
var options = new MigrationEngineOptions(TimeSpan.Zero, 5, 100, []) { BodyRetryBackoff = TimeSpan.Zero };
96+
var engine = new MigrationEngine(source, target, checkpointStore, new FakeTimeProvider(), options, NullLogger<MigrationEngine>.Instance);
97+
98+
Assert.ThrowsAsync<OperationCanceledException>(() => engine.RunCategoryAsync(category, stopping.Token));
99+
100+
var persisted = await checkpointStore.Read(category.Id);
101+
using (Assert.EnterMultipleScope())
102+
{
103+
Assert.That(source.BodyReadAttempts("msg-1"), Is.EqualTo(1), "a shutdown is not a transient body failure, so it is not retried");
104+
Assert.That(persisted!.State, Is.EqualTo(MigrationCategoryState.InProgress));
105+
Assert.That(persisted.SkippedCount, Is.Zero, "the message is still there to copy on the next run");
106+
Assert.That(target.WrittenRows(category.Id), Is.Empty);
107+
}
108+
}
109+
110+
[Test]
111+
public async Task The_configured_backoff_is_waited_between_body_read_attempts()
112+
{
113+
// Without the wait, three attempts against a body store that is briefly down all fail inside a
114+
// millisecond and the message is skipped for an outage it would have survived.
115+
var category = MigrationCategoryRegistry.Find("UnresolvedAndRetryIssuedFailedMessages")!;
116+
var source = new InMemoryMigrationSource();
117+
source.Seed(category.Id, Row("msg-1"));
118+
var body = new MigrationBody(new byte[] { 1 }, "text/plain");
119+
source.SetBody("msg-1", body);
120+
source.FailBodyReads("msg-1", MigrationEngine.MaxBodyReadAttempts - 1, new TimeoutException("body store unreachable"));
121+
var checkpointStore = new InMemoryMigrationCheckpointStore();
122+
var target = new InMemoryMigrationTarget(checkpointStore);
123+
var clock = new TimerRecordingTimeProvider();
124+
var backoff = TimeSpan.FromMilliseconds(200);
125+
var options = new MigrationEngineOptions(TimeSpan.Zero, 5, 100, []) { BodyRetryBackoff = backoff };
126+
var engine = new MigrationEngine(source, target, checkpointStore, clock, options, NullLogger<MigrationEngine>.Instance);
127+
128+
var runTask = engine.RunCategoryAsync(category);
129+
130+
// Two failures, so a wait after each before the attempt that succeeds.
131+
for (var waitNumber = 1; waitNumber <= MigrationEngine.MaxBodyReadAttempts - 1; waitNumber++)
132+
{
133+
Assert.That(await clock.TimerCreated.WaitAsync(TimeSpan.FromSeconds(5)), Is.True, $"backoff {waitNumber} never started");
134+
Assert.That(runTask.IsCompleted, Is.False, $"backoff {waitNumber} should still be pending");
135+
clock.Advance(backoff);
136+
}
137+
138+
var checkpoint = await runTask.WaitAsync(TimeSpan.FromSeconds(5));
139+
140+
using (Assert.EnterMultipleScope())
141+
{
142+
Assert.That(checkpoint.State, Is.EqualTo(MigrationCategoryState.Complete));
143+
Assert.That(clock.DueTimes, Is.EqualTo(new[] { backoff, backoff }), "no wait after the final attempt, which has nothing left to retry");
144+
Assert.That(target.WrittenRows(category.Id).Single().Body, Is.EqualTo(body));
145+
}
146+
}
147+
70148
[Test]
71149
public async Task The_skip_warning_for_an_unreadable_body_carries_the_last_attempts_exception()
72150
{

‎src/ServiceControl.UnitTests/Migration/MigrationEngineCopyTests.cs‎

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,72 @@ public async Task Copies_every_row_and_finishes_Complete_when_nothing_was_skippe
4040
}
4141
}
4242

43+
[Test]
44+
public async Task A_category_with_no_rows_finishes_Complete_without_a_cursor()
45+
{
46+
// The ordinary state of several required categories on a small instance: nothing to copy is a
47+
// finished category, not a category that never ran.
48+
var category = MigrationCategoryRegistry.Find("MessageRedirects")!;
49+
var source = new InMemoryMigrationSource();
50+
var checkpointStore = new InMemoryMigrationCheckpointStore();
51+
var target = new InMemoryMigrationTarget(checkpointStore);
52+
var options = new MigrationEngineOptions(TimeSpan.Zero, 5, 100, []);
53+
var engine = new MigrationEngine(source, target, checkpointStore, new FakeTimeProvider(), options, NullLogger<MigrationEngine>.Instance);
54+
55+
var checkpoint = await engine.RunCategoryAsync(category);
56+
57+
using (Assert.EnterMultipleScope())
58+
{
59+
Assert.That(checkpoint.State, Is.EqualTo(MigrationCategoryState.Complete));
60+
Assert.That((checkpoint.CopiedCount, checkpoint.SkippedCount), Is.EqualTo((0L, 0L)));
61+
Assert.That(checkpoint.Cursor, Is.Null);
62+
Assert.That(checkpoint.SettledAt, Is.Not.Null);
63+
}
64+
}
65+
66+
[Test]
67+
public async Task The_moment_a_category_settles_comes_from_the_injected_clock()
68+
{
69+
var category = MigrationCategoryRegistry.Find("KnownEndpoints")!;
70+
var source = new InMemoryMigrationSource();
71+
source.Seed(category.Id, Row("a"));
72+
var checkpointStore = new InMemoryMigrationCheckpointStore();
73+
var target = new InMemoryMigrationTarget(checkpointStore);
74+
var settledAt = new DateTimeOffset(2026, 3, 4, 5, 6, 7, TimeSpan.Zero);
75+
var clock = new FakeTimeProvider(settledAt);
76+
var options = new MigrationEngineOptions(TimeSpan.Zero, 5, 100, []);
77+
var engine = new MigrationEngine(source, target, checkpointStore, clock, options, NullLogger<MigrationEngine>.Instance);
78+
79+
var checkpoint = await engine.RunCategoryAsync(category);
80+
81+
Assert.That(checkpoint.SettledAt, Is.EqualTo(settledAt.UtcDateTime), "a wall-clock read here would drift from every other time the migration reports");
82+
}
83+
84+
[Test]
85+
public async Task A_category_already_CompleteWithErrors_is_left_alone_on_a_second_run()
86+
{
87+
// Finished with a few skips is finished. Re-reading it would copy the whole category again and
88+
// count its skips a second time.
89+
var category = MigrationCategoryRegistry.Find("KnownEndpoints")!;
90+
var source = new InMemoryMigrationSource();
91+
source.Seed(category.Id, Row("a"), Row("b"));
92+
var checkpointStore = new InMemoryMigrationCheckpointStore();
93+
var finishedWithSkips = new MigrationCheckpoint(category.Id, MigrationCategoryState.CompleteWithErrors, "b", 1, 1, 2,
94+
new Dictionary<MigrationSkipReason, long> { [MigrationSkipReason.BodyUnreadable] = 1 }, DateTime.UtcNow, DateTime.UtcNow, DateTime.UtcNow, null);
95+
await checkpointStore.Upsert(finishedWithSkips);
96+
var target = new InMemoryMigrationTarget(checkpointStore);
97+
var options = new MigrationEngineOptions(TimeSpan.Zero, 5, 100, []);
98+
var engine = new MigrationEngine(source, target, checkpointStore, new FakeTimeProvider(), options, NullLogger<MigrationEngine>.Instance);
99+
100+
var checkpoint = await engine.RunCategoryAsync(category);
101+
102+
using (Assert.EnterMultipleScope())
103+
{
104+
Assert.That(checkpoint, Is.EqualTo(finishedWithSkips with { Version = 1 }), "the row is read back untouched, at the version the seeding save left it");
105+
Assert.That(target.WrittenRows(category.Id), Is.Empty);
106+
}
107+
}
108+
43109
[Test]
44110
public async Task A_category_already_Complete_is_left_alone_on_a_second_run()
45111
{

0 commit comments

Comments
 (0)