Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 13 additions & 6 deletions Directory.Packages.props
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,13 @@
</PropertyGroup>

<ItemGroup Label="Messaging (WolverineFx — Postgres transport default, DR-009)">
<PackageVersion Include="WolverineFx" Version="5.18.0" />
<PackageVersion Include="WolverineFx.Postgresql" Version="5.18.0" />
<PackageVersion Include="WolverineFx" Version="6.22.0" />
<PackageVersion Include="WolverineFx.Postgresql" Version="6.22.0" />
<!-- Wolverine 6 removed the runtime Roslyn compiler from core (GH-2876). Handler code is
compiled at runtime (TypeLoadMode.Dynamic), so the hosts that call UseWolverine need
this package — it auto-registers an IAssemblyGenerator. The alternative is pre-generated
code + TypeLoadMode.Static, which trades a build step for faster cold start. -->
<PackageVersion Include="WolverineFx.RuntimeCompilation" Version="6.22.0" />
</ItemGroup>

<ItemGroup Label="EF Core">
Expand Down Expand Up @@ -53,18 +58,20 @@
<ItemGroup Label="Aspire">
<PackageVersion Include="Aspire.Hosting" Version="13.4.6" />
<PackageVersion Include="Aspire.Hosting.PostgreSQL" Version="13.4.6" />
<PackageVersion Include="Microsoft.Extensions.ServiceDiscovery" Version="9.3.1" />
<PackageVersion Include="Microsoft.Extensions.Http.Resilience" Version="9.4.0" />
<PackageVersion Include="Microsoft.Extensions.ServiceDiscovery" Version="10.8.0" />
<PackageVersion Include="Microsoft.Extensions.Http.Resilience" Version="10.8.0" />
</ItemGroup>

<ItemGroup Label="Testing">
<PackageVersion Include="Microsoft.NET.Test.Sdk" Version="18.8.1" />
<PackageVersion Include="Testcontainers.PostgreSql" Version="4.13.0" />
<PackageVersion Include="xunit" Version="2.9.3" />
<PackageVersion Include="xunit.v3" Version="3.2.2" />
<PackageVersion Include="xunit.runner.visualstudio" Version="3.1.5" />
<PackageVersion Include="Shouldly" Version="4.3.0" />
<PackageVersion Include="NSubstitute" Version="6.0.0" />
<PackageVersion Include="TngTech.ArchUnitNET.xUnit" Version="0.13.3" />
<!-- Core ArchUnitNET only: the .xUnit integration package still depends on xunit.assert 2.x,
which conflicts with xunit v3. ArchRuleAssert in Architecture.Tests replaces it. -->
<PackageVersion Include="TngTech.ArchUnitNET" Version="0.13.3" />
</ItemGroup>

</Project>
1 change: 1 addition & 0 deletions src/Presentation/Agents/Ard/Krautwatch.Agents.Ard.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
<ItemGroup>
<PackageReference Include="WolverineFx" />
<PackageReference Include="WolverineFx.Postgresql" />
<PackageReference Include="WolverineFx.RuntimeCompilation" />
</ItemGroup>

</Project>
7 changes: 7 additions & 0 deletions src/Presentation/Agents/Ard/Program.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
using Krautwatch.Application;
using Krautwatch.Application.Crawling;
using Krautwatch.Infrastructure;
using JasperFx.CodeGeneration.Model;
using Wolverine;
using Wolverine.Postgresql;

Expand Down Expand Up @@ -42,6 +43,12 @@
{
opts.PersistMessagesWithPostgresql(connectionString);
opts.Policies.UseDurableLocalQueues();
// Wolverine 6 changed the default ServiceLocationPolicy to NotAllowed (5.x was AllowedButWarn),
// which refuses to generate a handler needing container resolution. CrawlShowHandler needs it:
// IEnumerable<IBroadcasterCrawler> is an opaque lambda registration, and IEpisodeRepository's
// graph reaches EF's own DbContextOptions factory — not something we control. Restore the 5.x
// behaviour: allowed, but keep Wolverine's warning so the nudge to inline stays visible.
opts.ServiceLocationPolicy = ServiceLocationPolicy.AllowedButWarn;
// Discover the Crawling Action (CrawlShowHandler) in the Application assembly.
opts.Discovery.IncludeAssembly(typeof(CrawlShowCommand).Assembly);
});
Expand Down
1 change: 1 addition & 0 deletions src/Presentation/Agents/Zdf/Krautwatch.Agents.Zdf.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
<ItemGroup>
<PackageReference Include="WolverineFx" />
<PackageReference Include="WolverineFx.Postgresql" />
<PackageReference Include="WolverineFx.RuntimeCompilation" />
</ItemGroup>

</Project>
7 changes: 7 additions & 0 deletions src/Presentation/Agents/Zdf/Program.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
using Krautwatch.Application;
using Krautwatch.Application.Crawling;
using Krautwatch.Infrastructure;
using JasperFx.CodeGeneration.Model;
using Wolverine;
using Wolverine.Postgresql;

Expand Down Expand Up @@ -38,6 +39,12 @@
{
opts.PersistMessagesWithPostgresql(connectionString);
opts.Policies.UseDurableLocalQueues();
// Wolverine 6 changed the default ServiceLocationPolicy to NotAllowed (5.x was AllowedButWarn),
// which refuses to generate a handler needing container resolution. CrawlShowHandler needs it:
// IEnumerable<IBroadcasterCrawler> is an opaque lambda registration, and IEpisodeRepository's
// graph reaches EF's own DbContextOptions factory — not something we control. Restore the 5.x
// behaviour: allowed, but keep Wolverine's warning so the nudge to inline stays visible.
opts.ServiceLocationPolicy = ServiceLocationPolicy.AllowedButWarn;
// Discover the Crawling Action (CrawlShowHandler) in the Application assembly.
opts.Discovery.IncludeAssembly(typeof(CrawlShowCommand).Assembly);
});
Expand Down
8 changes: 4 additions & 4 deletions tests/Application.Tests/CrawlShowHandlerTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ public async Task Handle_selects_crawler_by_provider_and_upserts_the_crawled_epi
var repo = Substitute.For<IEpisodeRepository>();

var handler = new CrawlShowHandler([otherProvider, zdf], repo, NullLogger<CrawlShowHandler>.Instance);
await handler.HandleAsync(new CrawlShowCommand("zdf", "heute-show"));
await handler.HandleAsync(new CrawlShowCommand("zdf", "heute-show"), TestContext.Current.CancellationToken);

zdf.LastQuery.ShouldBe("heute-show");
otherProvider.LastQuery.ShouldBeNull(); // the ARD crawler must not be invoked
Expand All @@ -56,7 +56,7 @@ public async Task Handle_is_case_insensitive_on_provider_key()
var repo = Substitute.For<IEpisodeRepository>();

var handler = new CrawlShowHandler([zdf], repo, NullLogger<CrawlShowHandler>.Instance);
await handler.HandleAsync(new CrawlShowCommand("ZDF", "heute-show"));
await handler.HandleAsync(new CrawlShowCommand("ZDF", "heute-show"), TestContext.Current.CancellationToken);

zdf.LastQuery.ShouldBe("heute-show");
await repo.Received(1).UpsertManyAsync(Arg.Any<IEnumerable<Episode>>(), Arg.Any<CancellationToken>());
Expand All @@ -68,7 +68,7 @@ public async Task Handle_unknown_provider_is_a_no_op()
var repo = Substitute.For<IEpisodeRepository>();

var handler = new CrawlShowHandler([], repo, NullLogger<CrawlShowHandler>.Instance);
await handler.HandleAsync(new CrawlShowCommand("kika", "Biene Maja"));
await handler.HandleAsync(new CrawlShowCommand("kika", "Biene Maja"), TestContext.Current.CancellationToken);

await repo.DidNotReceive().UpsertManyAsync(Arg.Any<IEnumerable<Episode>>(), Arg.Any<CancellationToken>());
}
Expand All @@ -80,7 +80,7 @@ public async Task Handle_empty_crawl_result_does_not_upsert()
var repo = Substitute.For<IEpisodeRepository>();

var handler = new CrawlShowHandler([ard], repo, NullLogger<CrawlShowHandler>.Instance);
await handler.HandleAsync(new CrawlShowCommand("ard", "Nonexistent Show"));
await handler.HandleAsync(new CrawlShowCommand("ard", "Nonexistent Show"), TestContext.Current.CancellationToken);

await repo.DidNotReceive().UpsertManyAsync(Arg.Any<IEnumerable<Episode>>(), Arg.Any<CancellationToken>());
}
Expand Down
63 changes: 31 additions & 32 deletions tests/Application.Tests/DownloadHandlerTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -67,43 +67,42 @@ public class StartDownloadHandlerTests
public async Task ValidRequest_CreatesJobAndEnqueues()
{
var episode = Fixtures.MakeEpisode();
_episodes.GetByIdAsync("ep-1", default).Returns(episode);
_episodes.GetByIdAsync("ep-1", Arg.Any<CancellationToken>()).Returns(episode);

var result = await Handler().HandleAsync(new StartDownloadRequest("ep-1", "stream-1"));
var result = await Handler().HandleAsync(new StartDownloadRequest("ep-1", "stream-1"), TestContext.Current.CancellationToken);

result.ShouldNotBeNull();
result!.EpisodeId.ShouldBe("ep-1");
result.Status.ShouldBe(nameof(DownloadStatus.Queued));

await _jobs.Received(1).AddAsync(Arg.Any<DownloadJob>(), default);
await _jobs.Received(1).AddAsync(Arg.Any<DownloadJob>(), TestContext.Current.CancellationToken);
await _queue.Received(1).EnqueueAsync(
Arg.Any<Guid>(),
"https://example.com/ep.mp4",
default);
"https://example.com/ep.mp4", TestContext.Current.CancellationToken);
}

[Fact]
public async Task EpisodeNotFound_ReturnsNull_NoJobCreated()
{
_episodes.GetByIdAsync("bad-id", default).Returns((Episode?)null);
_episodes.GetByIdAsync("bad-id", Arg.Any<CancellationToken>()).Returns((Episode?)null);

var result = await Handler().HandleAsync(new StartDownloadRequest("bad-id", "stream-1"));
var result = await Handler().HandleAsync(new StartDownloadRequest("bad-id", "stream-1"), TestContext.Current.CancellationToken);

result.ShouldBeNull();
await _jobs.DidNotReceive().AddAsync(Arg.Any<DownloadJob>(), default);
await _queue.DidNotReceive().EnqueueAsync(Arg.Any<Guid>(), Arg.Any<string>(), default);
await _jobs.DidNotReceive().AddAsync(Arg.Any<DownloadJob>(), TestContext.Current.CancellationToken);
await _queue.DidNotReceive().EnqueueAsync(Arg.Any<Guid>(), Arg.Any<string>(), TestContext.Current.CancellationToken);
}

[Fact]
public async Task StreamNotFound_ReturnsNull_NoJobCreated()
{
var episode = Fixtures.MakeEpisode("real-stream");
_episodes.GetByIdAsync("ep-1", default).Returns(episode);
_episodes.GetByIdAsync("ep-1", Arg.Any<CancellationToken>()).Returns(episode);

var result = await Handler().HandleAsync(new StartDownloadRequest("ep-1", "wrong-stream"));
var result = await Handler().HandleAsync(new StartDownloadRequest("ep-1", "wrong-stream"), TestContext.Current.CancellationToken);

result.ShouldBeNull();
await _jobs.DidNotReceive().AddAsync(Arg.Any<DownloadJob>(), default);
await _jobs.DidNotReceive().AddAsync(Arg.Any<DownloadJob>(), TestContext.Current.CancellationToken);
}
}

Expand All @@ -125,13 +124,13 @@ public async Task ActiveStatus_CancelsAndReturnsTrue(DownloadStatus status)
{
var job = Fixtures.MakeJob();
ApplyStatus(job, status);
_jobs.GetByIdAsync(job.Id, default).Returns(job);
_jobs.GetByIdAsync(job.Id, Arg.Any<CancellationToken>()).Returns(job);

var result = await Handler().HandleAsync(job.Id);
var result = await Handler().HandleAsync(job.Id, TestContext.Current.CancellationToken);

result.ShouldBeTrue();
job.Status.ShouldBe(DownloadStatus.Cancelled);
await _jobs.Received(1).UpdateAsync(job, default);
await _jobs.Received(1).UpdateAsync(job, TestContext.Current.CancellationToken);
}

[Theory]
Expand All @@ -144,20 +143,20 @@ public async Task TerminalStatus_ReturnsFalse_NoUpdate(DownloadStatus status)
{
var job = Fixtures.MakeJob();
ApplyStatus(job, status);
_jobs.GetByIdAsync(job.Id, default).Returns(job);
_jobs.GetByIdAsync(job.Id, Arg.Any<CancellationToken>()).Returns(job);

var result = await Handler().HandleAsync(job.Id);
var result = await Handler().HandleAsync(job.Id, TestContext.Current.CancellationToken);

result.ShouldBeFalse();
await _jobs.DidNotReceive().UpdateAsync(Arg.Any<DownloadJob>(), default);
await _jobs.DidNotReceive().UpdateAsync(Arg.Any<DownloadJob>(), TestContext.Current.CancellationToken);
}

[Fact]
public async Task JobNotFound_ReturnsFalse()
{
_jobs.GetByIdAsync(Arg.Any<Guid>(), default).Returns((DownloadJob?)null);
_jobs.GetByIdAsync(Arg.Any<Guid>(), Arg.Any<CancellationToken>()).Returns((DownloadJob?)null);

var result = await Handler().HandleAsync(Guid.NewGuid());
var result = await Handler().HandleAsync(Guid.NewGuid(), TestContext.Current.CancellationToken);

result.ShouldBeFalse();
}
Expand Down Expand Up @@ -197,50 +196,50 @@ public async Task FailedOrCancelledJob_CreatesNewJobAndEnqueues(DownloadStatus s
{
var original = Fixtures.MakeJob();
ApplyTerminal(original, status);
_jobs.GetByIdAsync(original.Id, default).Returns(original);
_jobs.GetByIdAsync(original.Id, Arg.Any<CancellationToken>()).Returns(original);

var result = await Handler().HandleAsync(original.Id);
var result = await Handler().HandleAsync(original.Id, TestContext.Current.CancellationToken);

result.ShouldNotBeNull();
result!.Status.ShouldBe(nameof(DownloadStatus.Queued));
// New job created — not the original id
result.JobId.ShouldNotBe(original.Id);

await _jobs.Received(1).AddAsync(Arg.Any<DownloadJob>(), default);
await _queue.Received(1).RequeueAsync(Arg.Any<Guid>(), original.StreamUrl, default);
await _jobs.Received(1).AddAsync(Arg.Any<DownloadJob>(), TestContext.Current.CancellationToken);
await _queue.Received(1).RequeueAsync(Arg.Any<Guid>(), original.StreamUrl, TestContext.Current.CancellationToken);
}

[Fact]
public async Task CompletedJob_ReturnsNull_NothingEnqueued()
{
var job = Fixtures.MakeJob();
job.MarkCompleted("/out/file.mp4", 1024);
_jobs.GetByIdAsync(job.Id, default).Returns(job);
_jobs.GetByIdAsync(job.Id, Arg.Any<CancellationToken>()).Returns(job);

var result = await Handler().HandleAsync(job.Id);
var result = await Handler().HandleAsync(job.Id, TestContext.Current.CancellationToken);

result.ShouldBeNull();
await _queue.DidNotReceive().RequeueAsync(Arg.Any<Guid>(), Arg.Any<string>(), default);
await _queue.DidNotReceive().RequeueAsync(Arg.Any<Guid>(), Arg.Any<string>(), TestContext.Current.CancellationToken);
}

[Fact]
public async Task ActiveJob_ReturnsNull_NothingEnqueued()
{
var job = Fixtures.MakeJob(); // Queued — IsTerminal = false
_jobs.GetByIdAsync(job.Id, default).Returns(job);
_jobs.GetByIdAsync(job.Id, Arg.Any<CancellationToken>()).Returns(job);

var result = await Handler().HandleAsync(job.Id);
var result = await Handler().HandleAsync(job.Id, TestContext.Current.CancellationToken);

result.ShouldBeNull();
await _queue.DidNotReceive().RequeueAsync(Arg.Any<Guid>(), Arg.Any<string>(), default);
await _queue.DidNotReceive().RequeueAsync(Arg.Any<Guid>(), Arg.Any<string>(), TestContext.Current.CancellationToken);
}

[Fact]
public async Task JobNotFound_ReturnsNull()
{
_jobs.GetByIdAsync(Arg.Any<Guid>(), default).Returns((DownloadJob?)null);
_jobs.GetByIdAsync(Arg.Any<Guid>(), Arg.Any<CancellationToken>()).Returns((DownloadJob?)null);

var result = await Handler().HandleAsync(Guid.NewGuid());
var result = await Handler().HandleAsync(Guid.NewGuid(), TestContext.Current.CancellationToken);

result.ShouldBeNull();
}
Expand Down
6 changes: 3 additions & 3 deletions tests/Application.Tests/IndexingTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ public async Task Query_maps_episodes_to_releases_with_stable_guid_and_token()
repo.SearchAsync("heute-show", Arg.Any<CancellationToken>())
.Returns(new[] { Ep("zdf:doc-1", SeriesType.Daily, null, null) });

var releases = await new SearchReleasesHandler(repo).HandleAsync(new SearchReleasesQuery("heute-show"));
var releases = await new SearchReleasesHandler(repo).HandleAsync(new SearchReleasesQuery("heute-show"), TestContext.Current.CancellationToken);

var release = releases.ShouldHaveSingleItem();
release.Guid.ShouldBe("zdf:doc-1");
Expand All @@ -76,7 +76,7 @@ public async Task Empty_query_reads_the_recent_feed()
repo.GetRecentAsync(Arg.Any<int>(), Arg.Any<CancellationToken>())
.Returns(new[] { Ep("zdf:doc-9", SeriesType.Daily, null, null) });

var releases = await new SearchReleasesHandler(repo).HandleAsync(new SearchReleasesQuery(Q: null));
var releases = await new SearchReleasesHandler(repo).HandleAsync(new SearchReleasesQuery(Q: null), TestContext.Current.CancellationToken);

releases.ShouldHaveSingleItem().Guid.ShouldBe("zdf:doc-9");
await repo.DidNotReceive().SearchAsync(Arg.Any<string>(), Arg.Any<CancellationToken>());
Expand All @@ -93,7 +93,7 @@ public async Task Season_and_episode_filter_a_standard_series()
});

var releases = await new SearchReleasesHandler(repo)
.HandleAsync(new SearchReleasesQuery("Die Biene Maja", Season: 2, Episode: 52));
.HandleAsync(new SearchReleasesQuery("Die Biene Maja", Season: 2, Episode: 52), TestContext.Current.CancellationToken);

var release = releases.ShouldHaveSingleItem();
release.Guid.ShouldBe("kika:2");
Expand Down
7 changes: 6 additions & 1 deletion tests/Application.Tests/Krautwatch.Application.Tests.csproj
Original file line number Diff line number Diff line change
@@ -1,13 +1,18 @@
<Project Sdk="Microsoft.NET.Sdk">

<PropertyGroup>
<!-- xunit v3 test projects are executables. -->
<OutputType>Exe</OutputType>
</PropertyGroup>

<ItemGroup>
<ProjectReference Include="..\..\src\Application\Krautwatch.Application.csproj" />
<ProjectReference Include="..\..\src\Domain\Krautwatch.Domain.csproj" />
</ItemGroup>

<ItemGroup>
<PackageReference Include="Microsoft.NET.Test.Sdk" />
<PackageReference Include="xunit" />
<PackageReference Include="xunit.v3" />
<PackageReference Include="xunit.runner.visualstudio" />
<PackageReference Include="Shouldly" />
<PackageReference Include="NSubstitute" />
Expand Down
4 changes: 2 additions & 2 deletions tests/Application.Tests/RefreshProxyListHandlerTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ public async Task Upserts_the_fetched_candidates()
source.FetchAsync(Arg.Any<CancellationToken>()).Returns([P("1.1.1.1"), P("2.2.2.2")]);
var repo = Substitute.For<IProxyRepository>();

await new RefreshProxyListHandler(source, repo, NullLogger<RefreshProxyListHandler>.Instance).HandleAsync();
await new RefreshProxyListHandler(source, repo, NullLogger<RefreshProxyListHandler>.Instance).HandleAsync(TestContext.Current.CancellationToken);

await repo.Received(1).UpsertBatchAsync(
Arg.Is<IEnumerable<Proxy>>(ps => ps != null && ps.Count() == 2), Arg.Any<CancellationToken>());
Expand All @@ -35,7 +35,7 @@ public async Task An_empty_fetch_keeps_the_cached_rows_untouched()
source.FetchAsync(Arg.Any<CancellationToken>()).Returns([]);
var repo = Substitute.For<IProxyRepository>();

await new RefreshProxyListHandler(source, repo, NullLogger<RefreshProxyListHandler>.Instance).HandleAsync();
await new RefreshProxyListHandler(source, repo, NullLogger<RefreshProxyListHandler>.Instance).HandleAsync(TestContext.Current.CancellationToken);

await repo.DidNotReceive().UpsertBatchAsync(Arg.Any<IEnumerable<Proxy>>(), Arg.Any<CancellationToken>());
}
Expand Down
Loading