diff --git a/Directory.Packages.props b/Directory.Packages.props
index fa05691..f59273e 100644
--- a/Directory.Packages.props
+++ b/Directory.Packages.props
@@ -54,6 +54,7 @@
+
diff --git a/README.md b/README.md
index 0baf5bd..245d456 100644
--- a/README.md
+++ b/README.md
@@ -127,6 +127,32 @@ The token is required — `/setup` is closed without it, so nobody on your netwo
before you do. It lives in memory only and rotates if the process restarts. Once an administrator exists,
`/setup` never reopens.
+### What gets searched (query-driven, DR-011)
+
+A Newznab search for a show nothing has crawled yet **resolves it live** against the broadcasters, so
+Krautwatch works with no `*arr` configuration at all.
+
+**How the first search behaves is your choice** (a setting, editable in the UI — not a rebuild):
+
+| Mode | Behaviour |
+|---|---|
+| **Return results fast** *(default)* | Answer after a short wait with whatever has resolved so far, and let the crawl finish in the background. The first search may under-report; the next one is complete and instant. Advanced: set the wait in seconds (1–300, default 8). |
+| **Wait for complete result on first query** | Wait for the resolution to finish so the first search is already complete. Slower — and if it exceeds Sonarr's own indexer timeout, Sonarr may treat the indexer as failing. Still bounded by `CrawlTimeout`; no wait is ever unbounded. |
+
+Operational knobs stay in config:
+
+```
+Indexing:OnDemandResolution:Enabled # default true — kill switch
+Indexing:OnDemandResolution:CrawlTimeout # default 00:02:00 — background crawl budget,
+ # and the ceiling on "wait for complete"
+Indexing:OnDemandResolution:PositiveTtl # default 06:00:00 — trust a hit this long
+Indexing:OnDemandResolution:NegativeTtl # default 00:45:00 — trust a miss this long
+Indexing:OnDemandResolution:MaxConcurrentResolutions # default 2 — politeness cap toward ARD/ZDF
+```
+
+The RSS feed (no query) is never resolved — it serves the standing crawl list, since RSS-Sync polls
+constantly with no particular target.
+
### Downloads
`Download:Directory` sets the output path (the dev fleet points it at a temp dir; in production it's
diff --git a/docs/plans/2026-07-28 - query-driven-search.md b/docs/plans/2026-07-28 - query-driven-search.md
new file mode 100644
index 0000000..de9644c
--- /dev/null
+++ b/docs/plans/2026-07-28 - query-driven-search.md
@@ -0,0 +1,111 @@
+# 2026-07-28 — Query-driven search: resolve `t=tvsearch` on demand (#58)
+
+**Status:** planned · **Milestone:** ARD and KiKA indexer · **Implements:** [DR-011](../architecture/DR-011-search-driven-indexing.md)
+
+`SearchReleasesHandler` reads only `IEpisodeRepository`, and the crawl list is three hardcoded seed shows
+(`extra 3`, `Die Biene Maja`, `heute-show`). So Sonarr searching for anything else gets an empty feed —
+indistinguishable from "not available in the Mediathek". DR-011 makes on-demand resolution the target,
+because that is what lets Krautwatch work with **zero** `*arr` configuration.
+
+## Where resolution runs, and why that matters
+
+The Newznab host currently registers **no broadcaster crawlers** — only the agents do
+(`AddArdCrawler`/`AddZdfCrawler` are absent from `Api/NewznabIndexerApi/Program.cs`). Live resolution
+therefore needs one of:
+
+1. **In-process in the API host** — register the crawlers there and call the port directly.
+2. **Dispatch to an agent over the durable bus** and poll for completion.
+
+**Decision: (1).** Option 2 puts a message round-trip plus polling inside Sonarr's HTTP request, for a
+result that is needed synchronously — far more latency and machinery than the problem warrants. This makes
+the API host run an IO-driven Action, which is the same narrow DR-009 deviation already recorded for
+`TestArrConnection`: Actions are *supposed* to live on agents, but a synchronous request/response cannot
+wait on a bus. Recorded here rather than discovered later.
+
+## Slice isolation forces a small duplication
+
+The obvious move — call `Crawling`'s `CrawlShowHandler` from `Indexing` — is **forbidden** by the
+`Slice_does_not_depend_on_sibling_slices` architecture test, and rightly so. `Indexing` gets its own Action
+using the `IBroadcasterCrawler` port and `IEpisodeRepository` directly. The "crawl then upsert" shape
+repeats; that is the intended cost of slice isolation, not an accident to refactor away.
+
+## Shape
+
+### 1. Resolution cache (`ResolvedQuery`)
+
+Without a marker of "we already looked", every repeat search re-crawls ARD. Sonarr retries the *same*
+failing query on a schedule, so **negative caching matters more than positive** here.
+
+- `ResolvedQuery` entity: normalised query text (key), `LastAttemptedAt`, `ResultCount`, `ProvidersTried`.
+- TTLs, configurable: **positive 6h**, **negative 45m**. Rationale: public-TV episodes appear on
+ broadcast schedules, not continuously, so 6h is ample; 45m keeps a mistyped or genuinely-absent show from
+ hammering ARD every RSS-Sync cycle while still recovering within an hour of the show appearing.
+- Migration `AddResolvedQueries`.
+
+### 2. On-demand resolution Action (`Application/Indexing`)
+
+- Fan out to **all** registered crawlers concurrently (ard, kika, zdf) under one shared deadline, since we
+ cannot know which broadcaster carries a title.
+- **Bound the wait, not the crawl.** The request waits a configurable deadline (default 8s) and then
+ serves whatever has landed — but the crawl **keeps running in the background** to completion, so the next
+ call (search or RSS) gets the full set. This is the key decision: abandoning a half-finished crawl would
+ throw away work already paid for in ARD round-trips, and would leave the cache permanently partial.
+- **Therefore the crawl must not use the request's CancellationToken.** It runs on a queue drained by a
+ hosted service, tied to the *host* lifetime (`ApplicationStopping`) with its own longer budget. The
+ request merely awaits a completion signal up to its deadline. Getting this wrong — passing the request
+ token through — would silently cancel every crawl the instant the response is written.
+- **Coalesce** concurrent identical queries: the in-flight table means a Sonarr library refresh issuing the
+ same query twice crawls once, and a second caller simply waits on the first one's signal. Per-process
+ only — several API replicas would each crawl once, which is acceptable and not worth distributed locking.
+- **Outbound politeness:** cap concurrent background resolutions so a library refresh cannot become a crawl
+ storm against ARD.
+
+### 3. Wire into search
+
+`SearchReleasesHandler`: on a DB miss for a non-empty `q`, and if the resolution cache is stale, resolve →
+upsert → re-read → serve in the same response. The RSS path (`q` empty) is untouched — per DR-011 it keeps
+serving the standing crawl list.
+
+## Decisions
+
+**Everything is configurable** under `Indexing:OnDemandResolution` — the whole point is that an operator on
+a slow link or a fast LAN can tune it without a rebuild:
+
+> **Revised during implementation:** how long to wait is not ours to decide — it is a genuine trade-off
+> between a complete first answer and a responsive indexer, so it became a **user setting** on `AppSettings`
+> (`SearchWaitMode` = `ReturnFast` | `WaitForComplete`, plus `SearchWaitSeconds` as an advanced value for the
+> fast mode). It lives in the database rather than config precisely so the settings page can own it.
+> `WaitForComplete` is still capped by `CrawlTimeout` — an unbounded wait would hang the request on a stuck
+> crawl.
+
+| Setting | Default | Purpose |
+|---|---|---|
+| `Enabled` | `true` | Kill switch |
+| `CrawlTimeout` | `00:02:00` | Budget for the background crawl, independent of the request |
+| `PositiveTtl` | `06:00:00` | How long a successful resolution is trusted |
+| `NegativeTtl` | `00:45:00` | How long an empty result is trusted |
+| `MaxConcurrentResolutions` | `2` | Politeness cap on outbound crawling |
+
+**Only `q` searches resolve, never the RSS feed.** RSS-Sync polls constantly with no query; resolving there
+would mean crawling on a timer for no defined target.
+
+**A miss still returns 200 with an empty feed**, never an error. Sonarr treats indexer errors as an
+availability problem and will disable the indexer after repeated failures — "no results" and "broken" must
+stay distinguishable.
+
+**Deliberately out of scope: adding resolved shows to the standing crawl list.** It is tempting (a show you
+searched for should keep producing episodes in RSS), but it grows the crawl list without bound and needs its
+own retention policy. Filed separately rather than smuggled in.
+
+## Tests
+
+- Resolution cache: fresh positive suppresses a crawl; stale positive re-crawls; fresh **negative**
+ suppresses (the Sonarr-retry case); stale negative re-crawls.
+- Coalescing: two concurrent identical queries produce **one** crawl.
+- Deadline: a crawler that never returns does not hang the search; whatever is in the DB is served.
+- **The background crawl survives the response.** A crawl that finishes *after* the request deadline still
+ persists its episodes, so the next call sees the full set — this is the behaviour that would break if the
+ request token were threaded into the crawl, and it is the single most important test here.
+- A crawler throwing does not fail the whole search — the others still contribute.
+- RSS path (`q` empty) never triggers resolution.
+- A miss returns an empty result rather than throwing.
diff --git a/src/Application/ApplicationServiceExtensions.cs b/src/Application/ApplicationServiceExtensions.cs
index f8742d1..a794bd1 100644
--- a/src/Application/ApplicationServiceExtensions.cs
+++ b/src/Application/ApplicationServiceExtensions.cs
@@ -55,4 +55,20 @@ public static IServiceCollection AddApplication(this IServiceCollection services
return services;
}
+
+ ///
+ /// Enables query-driven search (#58 / DR-011): a Newznab search for a show no crawler has visited yet
+ /// resolves it on demand. Call only from a host that also registers broadcaster crawlers — without them
+ /// there is nothing to resolve against.
+ ///
+ public static IServiceCollection AddOnDemandResolution(
+ this IServiceCollection services, OnDemandResolutionOptions options)
+ {
+ // Singletons: the coalescing state and the queue must be shared process-wide, which is also why
+ // OnDemandResolver takes IServiceScopeFactory rather than a scoped repository.
+ services.AddSingleton(options);
+ services.AddSingleton();
+ services.AddHostedService();
+ return services;
+ }
}
diff --git a/src/Application/Indexing/OnDemandResolution.cs b/src/Application/Indexing/OnDemandResolution.cs
new file mode 100644
index 0000000..83ccf02
--- /dev/null
+++ b/src/Application/Indexing/OnDemandResolution.cs
@@ -0,0 +1,299 @@
+using System.Collections.Concurrent;
+using System.Threading.Channels;
+using Krautwatch.Domain.Entities;
+using Krautwatch.Domain.Enums;
+using Krautwatch.Domain.Interfaces;
+using Microsoft.Extensions.DependencyInjection;
+using Microsoft.Extensions.Hosting;
+using Microsoft.Extensions.Logging;
+
+namespace Krautwatch.Application.Indexing;
+
+// ══════════════════════════════════════════════════════════════
+// Options
+// ══════════════════════════════════════════════════════════════
+
+///
+/// Tuning for query-driven search (#58 / DR-011). Bound from Indexing:OnDemandResolution.
+///
+public class OnDemandResolutionOptions
+{
+ public const string SectionName = "Indexing:OnDemandResolution";
+
+ /// Kill switch — when false, search only ever reads what a crawler already persisted.
+ public bool Enabled { get; set; } = true;
+
+ ///
+ /// Budget for the background crawl, independent of any request. The ARD path is multi-hop
+ /// (A-Z widget → show page → episode list → item page), so this is generous by design. Also the ceiling
+ /// on how long a search can wait in mode — no wait is ever
+ /// unbounded.
+ ///
+ public TimeSpan CrawlTimeout { get; set; } = TimeSpan.FromMinutes(2);
+
+ /// How long a successful resolution is trusted before re-crawling.
+ public TimeSpan PositiveTtl { get; set; } = TimeSpan.FromHours(6);
+
+ ///
+ /// How long an empty result is trusted. Shorter than , but the more important
+ /// of the two: Sonarr re-issues the same failing query every RSS-Sync cycle.
+ ///
+ public TimeSpan NegativeTtl { get; set; } = TimeSpan.FromMinutes(45);
+
+ /// Politeness cap, so a Sonarr library refresh cannot become a crawl storm against ARD.
+ public int MaxConcurrentResolutions { get; set; } = 2;
+}
+
+// ══════════════════════════════════════════════════════════════
+// Action (IO-driven, DR-009) — runs in the API host; see the plan
+// ══════════════════════════════════════════════════════════════
+
+///
+/// Resolves a search term against the broadcasters on demand, so Sonarr can find a show no crawler has
+/// visited yet (#58).
+///
+///
+///
+/// The wait is bounded; the crawl is not. A caller waits at most
+/// and is then released to serve whatever landed,
+/// while the crawl runs to completion in the background so the next call gets the full set. Abandoning a
+/// half-finished crawl would discard ARD round-trips already paid for and leave the cache permanently
+/// partial.
+///
+///
+/// Consequently the crawl must not observe the request's CancellationToken — it is queued here and
+/// drained by under the host lifetime. Threading the request token
+/// through would cancel every crawl the instant the HTTP response was written, which presents as
+/// "resolution mysteriously never works".
+///
+///
+/// Identical concurrent queries are coalesced: a second caller waits on the first one's completion instead
+/// of starting a duplicate crawl. Per-process only — several API replicas would each crawl once, which is
+/// acceptable and not worth distributed locking.
+///
+///
+/// Singleton, so it deliberately takes rather than a repository: holding
+/// a scoped repository (and its DbContext) for the process lifetime would be a captive dependency.
+///
+///
+public sealed class OnDemandResolver(
+ IServiceScopeFactory scopeFactory,
+ OnDemandResolutionOptions options,
+ ILogger logger)
+{
+ private readonly ConcurrentDictionary _inFlight = new(StringComparer.Ordinal);
+ private readonly ConcurrentDictionary _pending = new(StringComparer.Ordinal);
+ // Fully qualified: Domain has its own Channel entity (a broadcaster channel).
+ private readonly System.Threading.Channels.Channel _queue =
+ System.Threading.Channels.Channel.CreateUnbounded();
+
+ internal ChannelReader Queue => _queue.Reader;
+
+ ///
+ /// Ensures the term has been (or is being) resolved, waiting as long as the operator has asked for.
+ /// Returns true when a resolution completed inside that window, so the caller knows a re-read is
+ /// worthwhile.
+ ///
+ public async Task EnsureResolvedAsync(string query, CancellationToken ct = default)
+ {
+ if (!options.Enabled || string.IsNullOrWhiteSpace(query))
+ return false;
+
+ var normalised = ResolvedQuery.Normalise(query);
+ if (normalised.Length == 0)
+ return false;
+
+ if (await IsFreshAsync(normalised, ct))
+ return false; // looked recently — nothing to wait for
+
+ // Join an in-flight resolution, or start one. GetOrAdd keeps it to a single crawl per term.
+ var resolution = _inFlight.GetOrAdd(normalised, key =>
+ {
+ var completion = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ _pending[key] = completion;
+ _queue.Writer.TryWrite(key); // unbounded, so this always succeeds
+ return completion.Task;
+ });
+
+ // How long to wait is the operator's choice (AppSettings), not a compiled-in constant.
+ var wait = await ResolveWaitAsync(ct);
+
+ var finished = await Task.WhenAny(resolution, Task.Delay(wait, ct));
+ if (finished == resolution)
+ return true;
+
+ logger.LogDebug(
+ "Resolution of '{Query}' still running after {Wait}; serving what is available and letting the "
+ + "crawl finish in the background.", normalised, wait);
+ return false;
+ }
+
+ ///
+ /// Releases waiters and clears the in-flight entry so the term can be resolved again once its TTL
+ /// lapses. One method rather than two, so a resolution can never signal without releasing (leaking the
+ /// term forever) or release without signalling (stranding a caller until its deadline).
+ ///
+ internal void ReleaseAndSignal(string key)
+ {
+ _inFlight.TryRemove(key, out _);
+ if (_pending.TryRemove(key, out var completion))
+ completion.TrySetResult();
+ }
+
+ ///
+ /// How long to wait for a resolution, per the operator's preference.
+ /// waits up to the crawl's own ceiling — never
+ /// indefinitely, since a stuck crawl would otherwise hang the request forever.
+ ///
+ private async Task ResolveWaitAsync(CancellationToken ct)
+ {
+ try
+ {
+ using var scope = scopeFactory.CreateScope();
+ var settings = await scope.ServiceProvider
+ .GetRequiredService().GetAsync(ct);
+
+ return settings.SearchWaitMode == SearchWaitMode.WaitForComplete
+ ? options.CrawlTimeout
+ : TimeSpan.FromSeconds(Math.Clamp(settings.SearchWaitSeconds, 1, 300));
+ }
+ catch (Exception ex) when (ex is not OperationCanceledException)
+ {
+ // Never fail a search because the preference could not be read — fall back to the safe default.
+ logger.LogWarning(ex, "Could not read the search wait preference; defaulting to 8s.");
+ return TimeSpan.FromSeconds(8);
+ }
+ }
+
+ private async Task IsFreshAsync(string normalised, CancellationToken ct)
+ {
+ using var scope = scopeFactory.CreateScope();
+ var previous = await scope.ServiceProvider
+ .GetRequiredService()
+ .GetAsync(normalised, ct);
+
+ if (previous is null)
+ return false;
+
+ var ttl = previous.ResultCount > 0 ? options.PositiveTtl : options.NegativeTtl;
+ return DateTimeOffset.UtcNow - previous.LastAttemptedAt < ttl;
+ }
+}
+
+///
+/// Drains the resolution queue, crawling each term against every registered broadcaster and persisting what
+/// comes back. Runs under the host lifetime, so a crawl survives the HTTP response that triggered it.
+///
+public sealed class OnDemandResolutionService(
+ OnDemandResolver resolver,
+ IServiceScopeFactory scopeFactory,
+ OnDemandResolutionOptions options,
+ ILogger logger) : BackgroundService
+{
+ protected override async Task ExecuteAsync(CancellationToken stoppingToken)
+ {
+ using var gate = new SemaphoreSlim(Math.Max(1, options.MaxConcurrentResolutions));
+
+ try
+ {
+ await foreach (var query in resolver.Queue.ReadAllAsync(stoppingToken))
+ {
+ await gate.WaitAsync(stoppingToken);
+
+ // Deliberately not awaited: one slow crawl must not block the queue. The semaphore bounds
+ // concurrency; stoppingToken — never a request token — bounds lifetime.
+ _ = ResolveAsync(query, stoppingToken)
+ .ContinueWith(_ => gate.Release(), TaskScheduler.Default);
+ }
+ }
+ catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
+ {
+ // normal shutdown
+ }
+ }
+
+ private async Task ResolveAsync(string normalised, CancellationToken stoppingToken)
+ {
+ using var timeout = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken);
+ timeout.CancelAfter(options.CrawlTimeout);
+
+ var persisted = 0;
+ var providers = Array.Empty();
+
+ try
+ {
+ using var scope = scopeFactory.CreateScope();
+ var crawlers = scope.ServiceProvider.GetServices().ToList();
+ var episodes = scope.ServiceProvider.GetRequiredService();
+
+ // Captured before the fan-out: building this inside the concurrent lambdas would be a
+ // data race on a non-thread-safe collection.
+ providers = crawlers.Select(c => c.ProviderKey).ToArray();
+
+ // We cannot know which broadcaster carries a title, so ask all of them.
+ var results = await Task.WhenAll(crawlers.Select(async crawler =>
+ {
+ try
+ {
+ return await crawler.CrawlShowAsync(normalised, timeout.Token);
+ }
+ catch (Exception ex) when (ex is not OperationCanceledException)
+ {
+ // One broadcaster failing must not sink the others.
+ logger.LogWarning(ex, "Crawler {Provider} failed resolving '{Query}'.",
+ crawler.ProviderKey, normalised);
+ return [];
+ }
+ }));
+
+ var found = results.SelectMany(r => r).ToList();
+ if (found.Count > 0)
+ {
+ await episodes.UpsertManyAsync(found, timeout.Token);
+ persisted = found.Count;
+ }
+
+ logger.LogInformation("Resolved '{Query}' on demand: {Count} episode(s) from {Providers}.",
+ normalised, persisted, string.Join(", ", providers));
+ }
+ catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
+ {
+ // Don't record an attempt that never ran, or shutdown would poison the cache with a false miss.
+ logger.LogDebug("Resolution of '{Query}' abandoned — host is shutting down.", normalised);
+ resolver.ReleaseAndSignal(normalised);
+ return;
+ }
+ catch (Exception ex)
+ {
+ logger.LogWarning(ex, "Resolution of '{Query}' failed.", normalised);
+ }
+
+ await RecordAttemptAsync(normalised, persisted, providers, stoppingToken);
+
+ // Signalled last, so a caller still inside its deadline sees the persisted episodes.
+ resolver.ReleaseAndSignal(normalised);
+ }
+
+ private async Task RecordAttemptAsync(
+ string normalised, int persisted, string[] providers, CancellationToken stoppingToken)
+ {
+ try
+ {
+ using var scope = scopeFactory.CreateScope();
+ await scope.ServiceProvider.GetRequiredService().RecordAsync(
+ new ResolvedQuery
+ {
+ Query = normalised,
+ LastAttemptedAt = DateTimeOffset.UtcNow,
+ ResultCount = persisted,
+ ProvidersTried = providers.Length > 0 ? string.Join(",", providers) : null,
+ },
+ stoppingToken);
+ }
+ catch (Exception ex)
+ {
+ // Failing to record only costs a redundant crawl later — never fail the resolution over it.
+ logger.LogWarning(ex, "Could not record the resolution attempt for '{Query}'.", normalised);
+ }
+ }
+}
diff --git a/src/Application/Indexing/SearchReleases.cs b/src/Application/Indexing/SearchReleases.cs
index e0720db..552da22 100644
--- a/src/Application/Indexing/SearchReleases.cs
+++ b/src/Application/Indexing/SearchReleases.cs
@@ -9,16 +9,43 @@ namespace Krautwatch.Application.Indexing;
///
public record SearchReleasesQuery(string? Q = null, int? Season = null, int? Episode = null, int Limit = 100);
-public class SearchReleasesHandler(IEpisodeRepository episodes)
+///
+/// Serves Newznab results from the catalog, resolving against the broadcasters on demand when a search term
+/// has not been crawled yet (#58 / DR-011).
+///
+///
+/// resolver is optional, so a host that only reads the catalog needs no broadcaster clients wired in.
+/// When it is absent the behaviour is exactly the pre-#58 read-only search.
+///
+public class SearchReleasesHandler(IEpisodeRepository episodes, OnDemandResolver? resolver = null)
{
public async Task> HandleAsync(SearchReleasesQuery query, CancellationToken ct = default)
{
var limit = Math.Clamp(query.Limit, 1, 500);
- var found = string.IsNullOrWhiteSpace(query.Q)
- ? await episodes.GetRecentAsync(limit, ct)
- : await episodes.SearchAsync(query.Q!, ct);
+ // RSS (no query) is never resolved: RSS-Sync polls constantly with no particular target, so
+ // resolving here would mean crawling on a timer for nothing specific. Per DR-011 it serves the
+ // standing crawl list.
+ if (string.IsNullOrWhiteSpace(query.Q))
+ return Project(await episodes.GetRecentAsync(limit, ct), query, limit);
+
+ var matches = Project(await episodes.SearchAsync(query.Q!, ct), query, limit);
+ if (matches.Count > 0 || resolver is null)
+ return matches;
+
+ // Nothing in the catalog. Ask the broadcasters, waiting only for the configured deadline — the
+ // crawl continues in the background either way, so a later call gets the full set.
+ if (await resolver.EnsureResolvedAsync(query.Q!, ct))
+ return Project(await episodes.SearchAsync(query.Q!, ct), query, limit);
+ // The deadline passed with the crawl still running. Return empty rather than an error: Sonarr
+ // treats indexer errors as an availability problem and will disable an indexer that keeps failing,
+ // so "no results yet" has to stay distinguishable from "broken".
+ return matches;
+ }
+
+ private static List Project(IReadOnlyList found, SearchReleasesQuery query, int limit)
+ {
IEnumerable filtered = found;
if (query.Season is not null)
filtered = filtered.Where(e => e.SeasonNumber == query.Season);
diff --git a/src/Application/Settings/SettingsHandlers.cs b/src/Application/Settings/SettingsHandlers.cs
index 003876e..fcc6be3 100644
--- a/src/Application/Settings/SettingsHandlers.cs
+++ b/src/Application/Settings/SettingsHandlers.cs
@@ -1,5 +1,6 @@
using FluentValidation;
using Krautwatch.Domain.Entities;
+using Krautwatch.Domain.Enums;
using Krautwatch.Domain.Interfaces;
namespace Krautwatch.Application.Settings;
@@ -12,12 +13,16 @@ public record SettingsResponse(
string DownloadDirectory,
int MaxConcurrentDownloads,
int CatalogRefreshIntervalHours,
- string CatalogProviderKey);
+ string CatalogProviderKey,
+ SearchWaitMode SearchWaitMode,
+ int SearchWaitSeconds);
public record SaveSettingsRequest(
string DownloadDirectory,
int MaxConcurrentDownloads,
- int CatalogRefreshIntervalHours);
+ int CatalogRefreshIntervalHours,
+ SearchWaitMode SearchWaitMode = SearchWaitMode.ReturnFast,
+ int SearchWaitSeconds = 8);
// ──────────────────────────────────────────────────────────────
// Validator
@@ -38,6 +43,12 @@ public SaveSettingsRequestValidator()
RuleFor(x => x.CatalogRefreshIntervalHours)
.InclusiveBetween(1, 168) // 1 hour to 1 week
.WithMessage("Refresh interval must be between 1 and 168 hours.");
+
+ // Only meaningful in ReturnFast mode, but validated regardless so a stale value cannot become
+ // active later by flipping the mode back.
+ RuleFor(x => x.SearchWaitSeconds)
+ .InclusiveBetween(1, 300)
+ .WithMessage("Search wait must be between 1 and 300 seconds.");
}
}
@@ -65,6 +76,8 @@ public async Task HandleAsync(
settings.DownloadDirectory = request.DownloadDirectory;
settings.MaxConcurrentDownloads = request.MaxConcurrentDownloads;
settings.CatalogRefreshIntervalHours = request.CatalogRefreshIntervalHours;
+ settings.SearchWaitMode = request.SearchWaitMode;
+ settings.SearchWaitSeconds = request.SearchWaitSeconds;
await repository.SaveAsync(settings, ct);
return SettingsMapper.ToResponse(settings);
@@ -77,5 +90,7 @@ file static class SettingsMapper
DownloadDirectory: s.DownloadDirectory,
MaxConcurrentDownloads: s.MaxConcurrentDownloads,
CatalogRefreshIntervalHours: s.CatalogRefreshIntervalHours,
- CatalogProviderKey: s.CatalogProviderKey);
+ CatalogProviderKey: s.CatalogProviderKey,
+ SearchWaitMode: s.SearchWaitMode,
+ SearchWaitSeconds: s.SearchWaitSeconds);
}
diff --git a/src/Domain/Entities/AppSettings.cs b/src/Domain/Entities/AppSettings.cs
index 9fd7675..ea42389 100644
--- a/src/Domain/Entities/AppSettings.cs
+++ b/src/Domain/Entities/AppSettings.cs
@@ -1,3 +1,5 @@
+using Krautwatch.Domain.Enums;
+
namespace Krautwatch.Domain.Entities;
///
@@ -10,4 +12,17 @@ public class AppSettings
public int MaxConcurrentDownloads { get; set; } = 2;
public int CatalogRefreshIntervalHours { get; set; } = 6;
public string CatalogProviderKey { get; set; } = "mediathekview";
+
+ ///
+ /// What a search should do when the show has not been crawled yet (#58). Defaults to
+ /// , because Sonarr treats a slow indexer as a broken one.
+ ///
+ public SearchWaitMode SearchWaitMode { get; set; } = SearchWaitMode.ReturnFast;
+
+ ///
+ /// How many seconds a search waits before answering, when
+ /// is . Advanced setting —
+ /// ignored entirely in mode.
+ ///
+ public int SearchWaitSeconds { get; set; } = 8;
}
diff --git a/src/Domain/Entities/ResolvedQuery.cs b/src/Domain/Entities/ResolvedQuery.cs
new file mode 100644
index 0000000..8242858
--- /dev/null
+++ b/src/Domain/Entities/ResolvedQuery.cs
@@ -0,0 +1,30 @@
+namespace Krautwatch.Domain.Entities;
+
+///
+/// A record that we have already tried to resolve a search term against the broadcasters, so repeat
+/// searches don't re-crawl (#58).
+///
+///
+/// The **negative** case is the important one. Sonarr re-issues the same query on a schedule, so without a
+/// marker for "we looked and found nothing", every RSS-Sync cycle would trigger a fresh multi-hop crawl of
+/// ARD for a show that isn't there. Hence is recorded rather than just a
+/// timestamp: successes and misses are trusted for different lengths of time.
+///
+public class ResolvedQuery
+{
+ /// Normalised search term — lower-cased and whitespace-collapsed. The natural key.
+ public string Query { get; set; } = string.Empty;
+
+ public DateTimeOffset LastAttemptedAt { get; set; } = DateTimeOffset.UtcNow;
+
+ /// Episodes persisted by the last attempt. Zero means a genuine miss, not a failure to run.
+ public int ResultCount { get; set; }
+
+ /// Which provider keys were tried, for diagnostics when a show is expected but absent.
+ public string? ProvidersTried { get; set; }
+
+ /// Collapses whitespace and case so trivially different spellings share one cache entry.
+ public static string Normalise(string query) =>
+ string.Join(' ', query.Trim().ToLowerInvariant()
+ .Split(' ', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries));
+}
diff --git a/src/Domain/Enums/SearchWaitMode.cs b/src/Domain/Enums/SearchWaitMode.cs
new file mode 100644
index 0000000..0a69e37
--- /dev/null
+++ b/src/Domain/Enums/SearchWaitMode.cs
@@ -0,0 +1,26 @@
+namespace Krautwatch.Domain.Enums;
+
+///
+/// What a Newznab search should do when the requested show has not been crawled yet and has to be resolved
+/// against the broadcasters (#58 / DR-011). A genuine trade-off, so it is the operator's choice.
+///
+public enum SearchWaitMode
+{
+ ///
+ /// Answer quickly with whatever has been resolved so far, letting the crawl finish in the background —
+ /// the first search may under-report, the next one is complete and instant. The safe default: Sonarr
+ /// treats a slow indexer as a broken one.
+ ///
+ ReturnFast = 0,
+
+ ///
+ /// Wait for the resolution to finish so the very first search is complete. Still bounded by
+ /// CrawlTimeout — an unbounded wait would hang the request forever on a stuck crawl.
+ ///
+ /// Be aware this puts a full multi-hop broadcaster crawl inside Sonarr's HTTP request. If it exceeds
+ /// Sonarr's own indexer timeout, Sonarr gives up and may mark the indexer as failing — so a generous
+ /// wait here can look like an outage from the other side.
+ ///
+ ///
+ WaitForComplete = 1,
+}
diff --git a/src/Domain/Interfaces/IRepositories.cs b/src/Domain/Interfaces/IRepositories.cs
index 32b4f9c..0566553 100644
--- a/src/Domain/Interfaces/IRepositories.cs
+++ b/src/Domain/Interfaces/IRepositories.cs
@@ -47,6 +47,15 @@ public interface IEpisodeRepository
Task UpsertManyAsync(IEnumerable episodes, CancellationToken ct = default);
}
+/// Tracks which search terms have already been resolved against the broadcasters (#58).
+public interface IResolvedQueryRepository
+{
+ Task GetAsync(string normalisedQuery, CancellationToken ct = default);
+
+ /// Records an attempt, inserting or replacing the existing entry for this query.
+ Task RecordAsync(ResolvedQuery attempt, CancellationToken ct = default);
+}
+
public interface ISettingsRepository
{
Task GetAsync(CancellationToken ct = default);
diff --git a/src/Infrastructure/InfrastructureServiceExtensions.cs b/src/Infrastructure/InfrastructureServiceExtensions.cs
index f37925b..ff25dff 100644
--- a/src/Infrastructure/InfrastructureServiceExtensions.cs
+++ b/src/Infrastructure/InfrastructureServiceExtensions.cs
@@ -62,6 +62,7 @@ public static IServiceCollection AddInfrastructure(
// Auth — local credential store + password hashing behind Domain ports (#48).
services.AddScoped();
+ services.AddScoped();
services.AddScoped();
services.AddSingleton();
diff --git a/src/Infrastructure/Persistence/AppDbContext.cs b/src/Infrastructure/Persistence/AppDbContext.cs
index e6309f6..b5cc974 100644
--- a/src/Infrastructure/Persistence/AppDbContext.cs
+++ b/src/Infrastructure/Persistence/AppDbContext.cs
@@ -16,6 +16,7 @@ public class AppDbContext(DbContextOptions options) : DbContext(op
public DbSet Proxies => Set();
public DbSet AdminAccounts => Set();
public DbSet ArrInstances => Set();
+ public DbSet ResolvedQueries => Set();
protected override void OnModelCreating(ModelBuilder modelBuilder)
{
@@ -194,6 +195,11 @@ protected override void OnModelCreating(ModelBuilder modelBuilder)
{
e.HasKey(x => x.Id);
e.Property(x => x.DownloadDirectory).IsRequired().HasMaxLength(500);
+
+ // Text like the other enums, so the column stays readable in the database.
+ e.Property(x => x.SearchWaitMode)
+ .HasConversion(v => v.ToString(), v => Enum.Parse(v))
+ .HasMaxLength(20);
e.HasData(new AppSettings
{
Id = 1,
@@ -239,6 +245,17 @@ protected override void OnModelCreating(ModelBuilder modelBuilder)
e.HasIndex(x => x.BaseUrl).IsUnique();
});
+ // --------------------------------------------------------
+ // ResolvedQuery — the on-demand search resolution cache (#58).
+ // Keyed on the normalised query itself; no surrogate id, because the query IS the identity.
+ // --------------------------------------------------------
+ modelBuilder.Entity(e =>
+ {
+ e.HasKey(x => x.Query);
+ e.Property(x => x.Query).HasMaxLength(300);
+ e.Property(x => x.ProvidersTried).HasMaxLength(200);
+ });
+
// --------------------------------------------------------
// TickerQ — job scheduler tables (TimeTickers, CronTickers, etc.)
// UseModelCustomizerForMigrations() is the alternative but we use
diff --git a/src/Infrastructure/Persistence/Migrations/20260727224258_AddResolvedQueries.Designer.cs b/src/Infrastructure/Persistence/Migrations/20260727224258_AddResolvedQueries.Designer.cs
new file mode 100644
index 0000000..23713bc
--- /dev/null
+++ b/src/Infrastructure/Persistence/Migrations/20260727224258_AddResolvedQueries.Designer.cs
@@ -0,0 +1,500 @@
+//
+using System;
+using Krautwatch.Infrastructure.Persistence;
+using Microsoft.EntityFrameworkCore;
+using Microsoft.EntityFrameworkCore.Infrastructure;
+using Microsoft.EntityFrameworkCore.Migrations;
+using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
+using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata;
+
+#nullable disable
+
+namespace Krautwatch.Infrastructure.Persistence.Migrations
+{
+ [DbContext(typeof(AppDbContext))]
+ [Migration("20260727224258_AddResolvedQueries")]
+ partial class AddResolvedQueries
+ {
+ ///
+ protected override void BuildTargetModel(ModelBuilder modelBuilder)
+ {
+#pragma warning disable 612, 618
+ modelBuilder
+ .HasAnnotation("ProductVersion", "10.0.10")
+ .HasAnnotation("Relational:MaxIdentifierLength", 63);
+
+ NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.AdminAccount", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("integer");
+
+ b.Property("CreatedAt")
+ .HasColumnType("timestamp with time zone");
+
+ b.Property("LastLoginAt")
+ .HasColumnType("timestamp with time zone");
+
+ b.Property("PasswordHash")
+ .IsRequired()
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.Property("Username")
+ .IsRequired()
+ .HasMaxLength(100)
+ .HasColumnType("character varying(100)");
+
+ b.HasKey("Id");
+
+ b.ToTable("AdminAccounts");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.AppSettings", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("CatalogProviderKey")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("CatalogRefreshIntervalHours")
+ .HasColumnType("integer");
+
+ b.Property("DownloadDirectory")
+ .IsRequired()
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.Property("MaxConcurrentDownloads")
+ .HasColumnType("integer");
+
+ b.HasKey("Id");
+
+ b.ToTable("Settings");
+
+ b.HasData(
+ new
+ {
+ Id = 1,
+ CatalogProviderKey = "mediathekview",
+ CatalogRefreshIntervalHours = 6,
+ DownloadDirectory = "/downloads",
+ MaxConcurrentDownloads = 2
+ });
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.ArrInstance", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("uuid");
+
+ b.Property("ApiKey")
+ .IsRequired()
+ .HasMaxLength(200)
+ .HasColumnType("character varying(200)");
+
+ b.Property("BaseUrl")
+ .IsRequired()
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.Property("CreatedAt")
+ .HasColumnType("timestamp with time zone");
+
+ b.Property("Enabled")
+ .HasColumnType("boolean");
+
+ b.Property("Kind")
+ .IsRequired()
+ .HasMaxLength(20)
+ .HasColumnType("character varying(20)");
+
+ b.Property("LastTestMessage")
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.Property("LastTestOk")
+ .HasColumnType("boolean");
+
+ b.Property("LastTestedAt")
+ .HasColumnType("timestamp with time zone");
+
+ b.Property("Name")
+ .IsRequired()
+ .HasMaxLength(100)
+ .HasColumnType("character varying(100)");
+
+ b.HasKey("Id");
+
+ b.HasIndex("BaseUrl")
+ .IsUnique();
+
+ b.ToTable("ArrInstances");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Channel", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("text");
+
+ b.Property("Name")
+ .IsRequired()
+ .HasMaxLength(100)
+ .HasColumnType("character varying(100)");
+
+ b.Property("ProviderKey")
+ .IsRequired()
+ .HasMaxLength(50)
+ .HasColumnType("character varying(50)");
+
+ b.HasKey("Id");
+
+ b.ToTable("Channels");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.DownloadJob", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("uuid");
+
+ b.Property("CompletedAt")
+ .HasColumnType("text");
+
+ b.Property("ContentLengthBytes")
+ .HasColumnType("bigint");
+
+ b.Property("CreatedAt")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("EpisodeId")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("ErrorMessage")
+ .HasMaxLength(2000)
+ .HasColumnType("character varying(2000)");
+
+ b.Property("FileSizeBytes")
+ .HasColumnType("bigint");
+
+ b.Property("GeoRestricted")
+ .HasColumnType("boolean");
+
+ b.Property("OutputPath")
+ .HasMaxLength(1000)
+ .HasColumnType("character varying(1000)");
+
+ b.Property("ProgressPercent")
+ .HasColumnType("double precision");
+
+ b.Property("Quality")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("StartedAt")
+ .HasColumnType("text");
+
+ b.Property("Status")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("StreamType")
+ .HasMaxLength(10)
+ .HasColumnType("character varying(10)");
+
+ b.Property("StreamUrl")
+ .IsRequired()
+ .HasMaxLength(2000)
+ .HasColumnType("character varying(2000)");
+
+ b.Property("TempPath")
+ .HasMaxLength(1000)
+ .HasColumnType("character varying(1000)");
+
+ b.Property("WorkerId")
+ .HasColumnType("text");
+
+ b.HasKey("Id");
+
+ b.HasIndex("CreatedAt");
+
+ b.HasIndex("EpisodeId");
+
+ b.HasIndex("Status");
+
+ b.ToTable("DownloadJobs");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Episode", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("text");
+
+ b.Property("AbsoluteEpisodeNumber")
+ .HasColumnType("integer");
+
+ b.Property("AvailableUntil")
+ .HasColumnType("text");
+
+ b.Property("BroadcastDate")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("ContentType")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("Description")
+ .HasMaxLength(5000)
+ .HasColumnType("character varying(5000)");
+
+ b.Property("Duration")
+ .HasColumnType("double precision");
+
+ b.Property("EpisodeNumber")
+ .HasColumnType("integer");
+
+ b.Property("GeoRestricted")
+ .HasColumnType("boolean");
+
+ b.Property("SeasonNumber")
+ .HasColumnType("integer");
+
+ b.Property("ShowId")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("Title")
+ .IsRequired()
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.HasKey("Id");
+
+ b.HasIndex("BroadcastDate");
+
+ b.HasIndex("ContentType");
+
+ b.HasIndex("ShowId");
+
+ b.HasIndex("ShowId", "SeasonNumber", "EpisodeNumber");
+
+ b.ToTable("Episodes");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.EpisodeStream", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("text");
+
+ b.Property("EpisodeId")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("Format")
+ .IsRequired()
+ .HasMaxLength(10)
+ .HasColumnType("character varying(10)");
+
+ b.Property("Quality")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("Url")
+ .IsRequired()
+ .HasMaxLength(2000)
+ .HasColumnType("character varying(2000)");
+
+ b.HasKey("Id");
+
+ b.HasIndex("EpisodeId");
+
+ b.ToTable("EpisodeStreams");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Proxy", b =>
+ {
+ b.Property("Id")
+ .HasMaxLength(100)
+ .HasColumnType("character varying(100)");
+
+ b.Property("AnonymityLevel")
+ .HasMaxLength(20)
+ .HasColumnType("character varying(20)");
+
+ b.Property("Country")
+ .HasMaxLength(10)
+ .HasColumnType("character varying(10)");
+
+ b.Property("CreatedAt")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("Host")
+ .IsRequired()
+ .HasMaxLength(100)
+ .HasColumnType("character varying(100)");
+
+ b.Property("LastProbeOk")
+ .HasColumnType("boolean");
+
+ b.Property("LastProbedAt")
+ .HasColumnType("text");
+
+ b.Property("Latency")
+ .HasColumnType("double precision");
+
+ b.Property("Port")
+ .HasColumnType("integer");
+
+ b.Property("Protocol")
+ .IsRequired()
+ .HasMaxLength(10)
+ .HasColumnType("character varying(10)");
+
+ b.Property("ResponseTime")
+ .HasColumnType("integer");
+
+ b.Property("Source")
+ .IsRequired()
+ .HasMaxLength(50)
+ .HasColumnType("character varying(50)");
+
+ b.Property("SourceLastChecked")
+ .HasColumnType("text");
+
+ b.Property("Speed")
+ .HasColumnType("integer");
+
+ b.Property("UpTime")
+ .HasColumnType("double precision");
+
+ b.Property("UpdatedAt")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("VerifiedEgressCountry")
+ .HasMaxLength(10)
+ .HasColumnType("character varying(10)");
+
+ b.HasKey("Id");
+
+ b.HasIndex("Country");
+
+ b.ToTable("Proxies");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.ResolvedQuery", b =>
+ {
+ b.Property("Query")
+ .HasMaxLength(300)
+ .HasColumnType("character varying(300)");
+
+ b.Property("LastAttemptedAt")
+ .HasColumnType("timestamp with time zone");
+
+ b.Property("ProvidersTried")
+ .HasMaxLength(200)
+ .HasColumnType("character varying(200)");
+
+ b.Property("ResultCount")
+ .HasColumnType("integer");
+
+ b.HasKey("Query");
+
+ b.ToTable("ResolvedQueries");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Show", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("text");
+
+ b.Property("ChannelId")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("SeriesType")
+ .IsRequired()
+ .HasMaxLength(20)
+ .HasColumnType("character varying(20)");
+
+ b.Property("Title")
+ .IsRequired()
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.Property("TvdbId")
+ .HasColumnType("integer");
+
+ b.HasKey("Id");
+
+ b.HasIndex("ChannelId");
+
+ b.ToTable("Shows");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.DownloadJob", b =>
+ {
+ b.HasOne("Krautwatch.Domain.Entities.Episode", "Episode")
+ .WithMany()
+ .HasForeignKey("EpisodeId")
+ .OnDelete(DeleteBehavior.Restrict)
+ .IsRequired();
+
+ b.Navigation("Episode");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Episode", b =>
+ {
+ b.HasOne("Krautwatch.Domain.Entities.Show", "Show")
+ .WithMany("Episodes")
+ .HasForeignKey("ShowId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired();
+
+ b.Navigation("Show");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.EpisodeStream", b =>
+ {
+ b.HasOne("Krautwatch.Domain.Entities.Episode", null)
+ .WithMany("Streams")
+ .HasForeignKey("EpisodeId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired();
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Show", b =>
+ {
+ b.HasOne("Krautwatch.Domain.Entities.Channel", "Channel")
+ .WithMany()
+ .HasForeignKey("ChannelId")
+ .OnDelete(DeleteBehavior.Restrict)
+ .IsRequired();
+
+ b.Navigation("Channel");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Episode", b =>
+ {
+ b.Navigation("Streams");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Show", b =>
+ {
+ b.Navigation("Episodes");
+ });
+#pragma warning restore 612, 618
+ }
+ }
+}
diff --git a/src/Infrastructure/Persistence/Migrations/20260727224258_AddResolvedQueries.cs b/src/Infrastructure/Persistence/Migrations/20260727224258_AddResolvedQueries.cs
new file mode 100644
index 0000000..a1ca42c
--- /dev/null
+++ b/src/Infrastructure/Persistence/Migrations/20260727224258_AddResolvedQueries.cs
@@ -0,0 +1,36 @@
+using System;
+using Microsoft.EntityFrameworkCore.Migrations;
+
+#nullable disable
+
+namespace Krautwatch.Infrastructure.Persistence.Migrations
+{
+ ///
+ public partial class AddResolvedQueries : Migration
+ {
+ ///
+ protected override void Up(MigrationBuilder migrationBuilder)
+ {
+ migrationBuilder.CreateTable(
+ name: "ResolvedQueries",
+ columns: table => new
+ {
+ Query = table.Column(type: "character varying(300)", maxLength: 300, nullable: false),
+ LastAttemptedAt = table.Column(type: "timestamp with time zone", nullable: false),
+ ResultCount = table.Column(type: "integer", nullable: false),
+ ProvidersTried = table.Column(type: "character varying(200)", maxLength: 200, nullable: true)
+ },
+ constraints: table =>
+ {
+ table.PrimaryKey("PK_ResolvedQueries", x => x.Query);
+ });
+ }
+
+ ///
+ protected override void Down(MigrationBuilder migrationBuilder)
+ {
+ migrationBuilder.DropTable(
+ name: "ResolvedQueries");
+ }
+ }
+}
diff --git a/src/Infrastructure/Persistence/Migrations/20260727230423_AddSearchWaitPreference.Designer.cs b/src/Infrastructure/Persistence/Migrations/20260727230423_AddSearchWaitPreference.Designer.cs
new file mode 100644
index 0000000..720b408
--- /dev/null
+++ b/src/Infrastructure/Persistence/Migrations/20260727230423_AddSearchWaitPreference.Designer.cs
@@ -0,0 +1,510 @@
+//
+using System;
+using Krautwatch.Infrastructure.Persistence;
+using Microsoft.EntityFrameworkCore;
+using Microsoft.EntityFrameworkCore.Infrastructure;
+using Microsoft.EntityFrameworkCore.Migrations;
+using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
+using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata;
+
+#nullable disable
+
+namespace Krautwatch.Infrastructure.Persistence.Migrations
+{
+ [DbContext(typeof(AppDbContext))]
+ [Migration("20260727230423_AddSearchWaitPreference")]
+ partial class AddSearchWaitPreference
+ {
+ ///
+ protected override void BuildTargetModel(ModelBuilder modelBuilder)
+ {
+#pragma warning disable 612, 618
+ modelBuilder
+ .HasAnnotation("ProductVersion", "10.0.10")
+ .HasAnnotation("Relational:MaxIdentifierLength", 63);
+
+ NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.AdminAccount", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("integer");
+
+ b.Property("CreatedAt")
+ .HasColumnType("timestamp with time zone");
+
+ b.Property("LastLoginAt")
+ .HasColumnType("timestamp with time zone");
+
+ b.Property("PasswordHash")
+ .IsRequired()
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.Property("Username")
+ .IsRequired()
+ .HasMaxLength(100)
+ .HasColumnType("character varying(100)");
+
+ b.HasKey("Id");
+
+ b.ToTable("AdminAccounts");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.AppSettings", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("CatalogProviderKey")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("CatalogRefreshIntervalHours")
+ .HasColumnType("integer");
+
+ b.Property("DownloadDirectory")
+ .IsRequired()
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.Property("MaxConcurrentDownloads")
+ .HasColumnType("integer");
+
+ b.Property("SearchWaitMode")
+ .IsRequired()
+ .HasMaxLength(20)
+ .HasColumnType("character varying(20)");
+
+ b.Property("SearchWaitSeconds")
+ .HasColumnType("integer");
+
+ b.HasKey("Id");
+
+ b.ToTable("Settings");
+
+ b.HasData(
+ new
+ {
+ Id = 1,
+ CatalogProviderKey = "mediathekview",
+ CatalogRefreshIntervalHours = 6,
+ DownloadDirectory = "/downloads",
+ MaxConcurrentDownloads = 2,
+ SearchWaitMode = "ReturnFast",
+ SearchWaitSeconds = 8
+ });
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.ArrInstance", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("uuid");
+
+ b.Property("ApiKey")
+ .IsRequired()
+ .HasMaxLength(200)
+ .HasColumnType("character varying(200)");
+
+ b.Property("BaseUrl")
+ .IsRequired()
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.Property("CreatedAt")
+ .HasColumnType("timestamp with time zone");
+
+ b.Property("Enabled")
+ .HasColumnType("boolean");
+
+ b.Property("Kind")
+ .IsRequired()
+ .HasMaxLength(20)
+ .HasColumnType("character varying(20)");
+
+ b.Property("LastTestMessage")
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.Property("LastTestOk")
+ .HasColumnType("boolean");
+
+ b.Property("LastTestedAt")
+ .HasColumnType("timestamp with time zone");
+
+ b.Property("Name")
+ .IsRequired()
+ .HasMaxLength(100)
+ .HasColumnType("character varying(100)");
+
+ b.HasKey("Id");
+
+ b.HasIndex("BaseUrl")
+ .IsUnique();
+
+ b.ToTable("ArrInstances");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Channel", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("text");
+
+ b.Property("Name")
+ .IsRequired()
+ .HasMaxLength(100)
+ .HasColumnType("character varying(100)");
+
+ b.Property("ProviderKey")
+ .IsRequired()
+ .HasMaxLength(50)
+ .HasColumnType("character varying(50)");
+
+ b.HasKey("Id");
+
+ b.ToTable("Channels");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.DownloadJob", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("uuid");
+
+ b.Property("CompletedAt")
+ .HasColumnType("text");
+
+ b.Property("ContentLengthBytes")
+ .HasColumnType("bigint");
+
+ b.Property("CreatedAt")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("EpisodeId")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("ErrorMessage")
+ .HasMaxLength(2000)
+ .HasColumnType("character varying(2000)");
+
+ b.Property("FileSizeBytes")
+ .HasColumnType("bigint");
+
+ b.Property("GeoRestricted")
+ .HasColumnType("boolean");
+
+ b.Property("OutputPath")
+ .HasMaxLength(1000)
+ .HasColumnType("character varying(1000)");
+
+ b.Property("ProgressPercent")
+ .HasColumnType("double precision");
+
+ b.Property("Quality")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("StartedAt")
+ .HasColumnType("text");
+
+ b.Property("Status")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("StreamType")
+ .HasMaxLength(10)
+ .HasColumnType("character varying(10)");
+
+ b.Property("StreamUrl")
+ .IsRequired()
+ .HasMaxLength(2000)
+ .HasColumnType("character varying(2000)");
+
+ b.Property("TempPath")
+ .HasMaxLength(1000)
+ .HasColumnType("character varying(1000)");
+
+ b.Property("WorkerId")
+ .HasColumnType("text");
+
+ b.HasKey("Id");
+
+ b.HasIndex("CreatedAt");
+
+ b.HasIndex("EpisodeId");
+
+ b.HasIndex("Status");
+
+ b.ToTable("DownloadJobs");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Episode", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("text");
+
+ b.Property("AbsoluteEpisodeNumber")
+ .HasColumnType("integer");
+
+ b.Property("AvailableUntil")
+ .HasColumnType("text");
+
+ b.Property("BroadcastDate")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("ContentType")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("Description")
+ .HasMaxLength(5000)
+ .HasColumnType("character varying(5000)");
+
+ b.Property("Duration")
+ .HasColumnType("double precision");
+
+ b.Property("EpisodeNumber")
+ .HasColumnType("integer");
+
+ b.Property("GeoRestricted")
+ .HasColumnType("boolean");
+
+ b.Property("SeasonNumber")
+ .HasColumnType("integer");
+
+ b.Property("ShowId")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("Title")
+ .IsRequired()
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.HasKey("Id");
+
+ b.HasIndex("BroadcastDate");
+
+ b.HasIndex("ContentType");
+
+ b.HasIndex("ShowId");
+
+ b.HasIndex("ShowId", "SeasonNumber", "EpisodeNumber");
+
+ b.ToTable("Episodes");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.EpisodeStream", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("text");
+
+ b.Property("EpisodeId")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("Format")
+ .IsRequired()
+ .HasMaxLength(10)
+ .HasColumnType("character varying(10)");
+
+ b.Property("Quality")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("Url")
+ .IsRequired()
+ .HasMaxLength(2000)
+ .HasColumnType("character varying(2000)");
+
+ b.HasKey("Id");
+
+ b.HasIndex("EpisodeId");
+
+ b.ToTable("EpisodeStreams");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Proxy", b =>
+ {
+ b.Property("Id")
+ .HasMaxLength(100)
+ .HasColumnType("character varying(100)");
+
+ b.Property("AnonymityLevel")
+ .HasMaxLength(20)
+ .HasColumnType("character varying(20)");
+
+ b.Property("Country")
+ .HasMaxLength(10)
+ .HasColumnType("character varying(10)");
+
+ b.Property("CreatedAt")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("Host")
+ .IsRequired()
+ .HasMaxLength(100)
+ .HasColumnType("character varying(100)");
+
+ b.Property("LastProbeOk")
+ .HasColumnType("boolean");
+
+ b.Property("LastProbedAt")
+ .HasColumnType("text");
+
+ b.Property("Latency")
+ .HasColumnType("double precision");
+
+ b.Property("Port")
+ .HasColumnType("integer");
+
+ b.Property("Protocol")
+ .IsRequired()
+ .HasMaxLength(10)
+ .HasColumnType("character varying(10)");
+
+ b.Property("ResponseTime")
+ .HasColumnType("integer");
+
+ b.Property("Source")
+ .IsRequired()
+ .HasMaxLength(50)
+ .HasColumnType("character varying(50)");
+
+ b.Property("SourceLastChecked")
+ .HasColumnType("text");
+
+ b.Property("Speed")
+ .HasColumnType("integer");
+
+ b.Property("UpTime")
+ .HasColumnType("double precision");
+
+ b.Property("UpdatedAt")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("VerifiedEgressCountry")
+ .HasMaxLength(10)
+ .HasColumnType("character varying(10)");
+
+ b.HasKey("Id");
+
+ b.HasIndex("Country");
+
+ b.ToTable("Proxies");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.ResolvedQuery", b =>
+ {
+ b.Property("Query")
+ .HasMaxLength(300)
+ .HasColumnType("character varying(300)");
+
+ b.Property("LastAttemptedAt")
+ .HasColumnType("timestamp with time zone");
+
+ b.Property("ProvidersTried")
+ .HasMaxLength(200)
+ .HasColumnType("character varying(200)");
+
+ b.Property("ResultCount")
+ .HasColumnType("integer");
+
+ b.HasKey("Query");
+
+ b.ToTable("ResolvedQueries");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Show", b =>
+ {
+ b.Property("Id")
+ .HasColumnType("text");
+
+ b.Property("ChannelId")
+ .IsRequired()
+ .HasColumnType("text");
+
+ b.Property("SeriesType")
+ .IsRequired()
+ .HasMaxLength(20)
+ .HasColumnType("character varying(20)");
+
+ b.Property("Title")
+ .IsRequired()
+ .HasMaxLength(500)
+ .HasColumnType("character varying(500)");
+
+ b.Property("TvdbId")
+ .HasColumnType("integer");
+
+ b.HasKey("Id");
+
+ b.HasIndex("ChannelId");
+
+ b.ToTable("Shows");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.DownloadJob", b =>
+ {
+ b.HasOne("Krautwatch.Domain.Entities.Episode", "Episode")
+ .WithMany()
+ .HasForeignKey("EpisodeId")
+ .OnDelete(DeleteBehavior.Restrict)
+ .IsRequired();
+
+ b.Navigation("Episode");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Episode", b =>
+ {
+ b.HasOne("Krautwatch.Domain.Entities.Show", "Show")
+ .WithMany("Episodes")
+ .HasForeignKey("ShowId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired();
+
+ b.Navigation("Show");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.EpisodeStream", b =>
+ {
+ b.HasOne("Krautwatch.Domain.Entities.Episode", null)
+ .WithMany("Streams")
+ .HasForeignKey("EpisodeId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired();
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Show", b =>
+ {
+ b.HasOne("Krautwatch.Domain.Entities.Channel", "Channel")
+ .WithMany()
+ .HasForeignKey("ChannelId")
+ .OnDelete(DeleteBehavior.Restrict)
+ .IsRequired();
+
+ b.Navigation("Channel");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Episode", b =>
+ {
+ b.Navigation("Streams");
+ });
+
+ modelBuilder.Entity("Krautwatch.Domain.Entities.Show", b =>
+ {
+ b.Navigation("Episodes");
+ });
+#pragma warning restore 612, 618
+ }
+ }
+}
diff --git a/src/Infrastructure/Persistence/Migrations/20260727230423_AddSearchWaitPreference.cs b/src/Infrastructure/Persistence/Migrations/20260727230423_AddSearchWaitPreference.cs
new file mode 100644
index 0000000..51c9902
--- /dev/null
+++ b/src/Infrastructure/Persistence/Migrations/20260727230423_AddSearchWaitPreference.cs
@@ -0,0 +1,50 @@
+using Microsoft.EntityFrameworkCore.Migrations;
+
+#nullable disable
+
+namespace Krautwatch.Infrastructure.Persistence.Migrations
+{
+ ///
+ public partial class AddSearchWaitPreference : Migration
+ {
+ ///
+ protected override void Up(MigrationBuilder migrationBuilder)
+ {
+ migrationBuilder.AddColumn(
+ name: "SearchWaitMode",
+ table: "Settings",
+ type: "character varying(20)",
+ maxLength: 20,
+ nullable: false,
+ // "ReturnFast", not "": the column is parsed as an enum, so an empty string would throw on
+ // read for any row inserted without an explicit value.
+ defaultValue: "ReturnFast");
+
+ migrationBuilder.AddColumn(
+ name: "SearchWaitSeconds",
+ table: "Settings",
+ type: "integer",
+ nullable: false,
+ defaultValue: 8); // matches the entity default; 0 would be an invalid wait
+
+ migrationBuilder.UpdateData(
+ table: "Settings",
+ keyColumn: "Id",
+ keyValue: 1,
+ columns: new[] { "SearchWaitMode", "SearchWaitSeconds" },
+ values: new object[] { "ReturnFast", 8 });
+ }
+
+ ///
+ protected override void Down(MigrationBuilder migrationBuilder)
+ {
+ migrationBuilder.DropColumn(
+ name: "SearchWaitMode",
+ table: "Settings");
+
+ migrationBuilder.DropColumn(
+ name: "SearchWaitSeconds",
+ table: "Settings");
+ }
+ }
+}
diff --git a/src/Infrastructure/Persistence/Migrations/AppDbContextModelSnapshot.cs b/src/Infrastructure/Persistence/Migrations/AppDbContextModelSnapshot.cs
index 81942ab..296bc80 100644
--- a/src/Infrastructure/Persistence/Migrations/AppDbContextModelSnapshot.cs
+++ b/src/Infrastructure/Persistence/Migrations/AppDbContextModelSnapshot.cs
@@ -71,6 +71,14 @@ protected override void BuildModel(ModelBuilder modelBuilder)
b.Property("MaxConcurrentDownloads")
.HasColumnType("integer");
+ b.Property("SearchWaitMode")
+ .IsRequired()
+ .HasMaxLength(20)
+ .HasColumnType("character varying(20)");
+
+ b.Property("SearchWaitSeconds")
+ .HasColumnType("integer");
+
b.HasKey("Id");
b.ToTable("Settings");
@@ -82,7 +90,9 @@ protected override void BuildModel(ModelBuilder modelBuilder)
CatalogProviderKey = "mediathekview",
CatalogRefreshIntervalHours = 6,
DownloadDirectory = "/downloads",
- MaxConcurrentDownloads = 2
+ MaxConcurrentDownloads = 2,
+ SearchWaitMode = "ReturnFast",
+ SearchWaitSeconds = 8
});
});
@@ -390,6 +400,27 @@ protected override void BuildModel(ModelBuilder modelBuilder)
b.ToTable("Proxies");
});
+ modelBuilder.Entity("Krautwatch.Domain.Entities.ResolvedQuery", b =>
+ {
+ b.Property("Query")
+ .HasMaxLength(300)
+ .HasColumnType("character varying(300)");
+
+ b.Property("LastAttemptedAt")
+ .HasColumnType("timestamp with time zone");
+
+ b.Property("ProvidersTried")
+ .HasMaxLength(200)
+ .HasColumnType("character varying(200)");
+
+ b.Property("ResultCount")
+ .HasColumnType("integer");
+
+ b.HasKey("Query");
+
+ b.ToTable("ResolvedQueries");
+ });
+
modelBuilder.Entity("Krautwatch.Domain.Entities.Show", b =>
{
b.Property("Id")
diff --git a/src/Infrastructure/Persistence/ResolvedQueryRepository.cs b/src/Infrastructure/Persistence/ResolvedQueryRepository.cs
new file mode 100644
index 0000000..7f66950
--- /dev/null
+++ b/src/Infrastructure/Persistence/ResolvedQueryRepository.cs
@@ -0,0 +1,30 @@
+using Krautwatch.Domain.Entities;
+using Krautwatch.Domain.Interfaces;
+using Microsoft.EntityFrameworkCore;
+
+namespace Krautwatch.Infrastructure.Persistence;
+
+/// EF-backed resolution cache for query-driven search (#58).
+public class ResolvedQueryRepository(AppDbContext db) : IResolvedQueryRepository
+{
+ public Task GetAsync(string normalisedQuery, CancellationToken ct = default) =>
+ db.ResolvedQueries.AsNoTracking().FirstOrDefaultAsync(q => q.Query == normalisedQuery, ct);
+
+ public async Task RecordAsync(ResolvedQuery attempt, CancellationToken ct = default)
+ {
+ var existing = await db.ResolvedQueries.FirstOrDefaultAsync(q => q.Query == attempt.Query, ct);
+
+ if (existing is null)
+ {
+ db.ResolvedQueries.Add(attempt);
+ }
+ else
+ {
+ existing.LastAttemptedAt = attempt.LastAttemptedAt;
+ existing.ResultCount = attempt.ResultCount;
+ existing.ProvidersTried = attempt.ProvidersTried;
+ }
+
+ await db.SaveChangesAsync(ct);
+ }
+}
diff --git a/src/Presentation/Api/NewznabIndexerApi/Program.cs b/src/Presentation/Api/NewznabIndexerApi/Program.cs
index 63dca40..a2ca604 100644
--- a/src/Presentation/Api/NewznabIndexerApi/Program.cs
+++ b/src/Presentation/Api/NewznabIndexerApi/Program.cs
@@ -1,5 +1,6 @@
using Krautwatch.Api.NewznabIndexerApi.Endpoints;
using Krautwatch.Application;
+using Krautwatch.Application.Indexing;
using Krautwatch.Infrastructure;
// Krautwatch Newznab indexer (DR-010) — the public *arr-facing surface. Reads the catalog and
@@ -20,6 +21,20 @@
});
builder.Services.AddApplication();
+// Query-driven search (#58 / DR-011). Sonarr searching for a show no crawler has visited must not get an
+// empty feed, so this host resolves against the broadcasters on demand — which means it needs the crawler
+// clients that until now only the agents had. That makes it run an IO-driven Action, the same narrow DR-009
+// deviation recorded for TestArrConnection: a synchronous request cannot wait on the durable bus.
+var resolutionOptions = new OnDemandResolutionOptions();
+builder.Configuration.GetSection(OnDemandResolutionOptions.SectionName).Bind(resolutionOptions);
+
+if (resolutionOptions.Enabled)
+{
+ builder.Services.AddArdCrawlers(); // ARD + KiKA
+ builder.Services.AddZdfCrawler();
+ builder.Services.AddOnDemandResolution(resolutionOptions);
+}
+
var app = builder.Build();
app.MapDefaultEndpoints(); // /health, /alive from ServiceDefaults
diff --git a/tests/Application.Tests/Krautwatch.Application.Tests.csproj b/tests/Application.Tests/Krautwatch.Application.Tests.csproj
index 615c20e..ec7b6bf 100644
--- a/tests/Application.Tests/Krautwatch.Application.Tests.csproj
+++ b/tests/Application.Tests/Krautwatch.Application.Tests.csproj
@@ -16,6 +16,9 @@
+
+
diff --git a/tests/Application.Tests/OnDemandResolutionTests.cs b/tests/Application.Tests/OnDemandResolutionTests.cs
new file mode 100644
index 0000000..e8e7627
--- /dev/null
+++ b/tests/Application.Tests/OnDemandResolutionTests.cs
@@ -0,0 +1,370 @@
+using Krautwatch.Application.Indexing;
+using Krautwatch.Domain.Entities;
+using Krautwatch.Domain.Enums;
+using Krautwatch.Domain.Interfaces;
+using Microsoft.Extensions.DependencyInjection;
+using Microsoft.Extensions.Logging.Abstractions;
+using NSubstitute;
+using Shouldly;
+using Xunit;
+
+namespace Krautwatch.Application.Tests;
+
+///
+/// Covers query-driven search (#58). The behaviour that matters most is that a crawl **outlives the request
+/// that triggered it** — if the request's CancellationToken ever reached the crawl, every resolution would
+/// be cancelled the instant the response was written and resolution would silently never work.
+///
+public class OnDemandResolutionTests
+{
+ // ── harness ───────────────────────────────────────────────
+
+ private sealed class Harness : IDisposable
+ {
+ public readonly FakeCrawler Crawler = new();
+ public readonly IResolvedQueryRepository ResolvedQueries = Substitute.For();
+ public readonly IEpisodeRepository Episodes = Substitute.For();
+ public readonly ISettingsRepository Settings = Substitute.For();
+ public readonly OnDemandResolutionOptions Options;
+ public readonly OnDemandResolver Resolver;
+ private readonly OnDemandResolutionService _service;
+ private readonly IServiceProvider _provider;
+ private readonly CancellationTokenSource _hostLifetime = new();
+
+ public Harness(
+ OnDemandResolutionOptions? options = null,
+ SearchWaitMode waitMode = SearchWaitMode.ReturnFast,
+ int waitSeconds = 1)
+ {
+ Options = options ?? new OnDemandResolutionOptions { CrawlTimeout = TimeSpan.FromSeconds(5) };
+
+ Settings.GetAsync(Arg.Any()).Returns(new AppSettings
+ {
+ SearchWaitMode = waitMode,
+ SearchWaitSeconds = waitSeconds,
+ });
+
+ var services = new ServiceCollection();
+ services.AddSingleton(ResolvedQueries);
+ services.AddSingleton(Episodes);
+ services.AddSingleton(Settings);
+ services.AddSingleton(Crawler);
+ _provider = services.BuildServiceProvider();
+
+ var scopeFactory = _provider.GetRequiredService();
+ Resolver = new OnDemandResolver(scopeFactory, Options, NullLogger.Instance);
+ _service = new OnDemandResolutionService(
+ Resolver, scopeFactory, Options, NullLogger.Instance);
+ }
+
+ public Task StartAsync() => _service.StartAsync(_hostLifetime.Token);
+
+ public void Dispose()
+ {
+ _hostLifetime.Cancel();
+ (_provider as IDisposable)?.Dispose();
+ _hostLifetime.Dispose();
+ }
+ }
+
+ private sealed class FakeCrawler : IBroadcasterCrawler
+ {
+ public string ProviderKey => "ard";
+ public int Calls;
+ public TaskCompletionSource? Gate;
+ public Exception? Throw;
+ public IReadOnlyList Result = [];
+
+ public async Task> CrawlShowAsync(string showQuery, CancellationToken ct = default)
+ {
+ Interlocked.Increment(ref Calls);
+ if (Gate is not null)
+ await Gate.Task.WaitAsync(ct);
+ if (Throw is not null)
+ throw Throw;
+ return Result;
+ }
+ }
+
+ private static Episode Ep(string id = "ard:1") => new()
+ {
+ Id = id, Title = "ep", ShowId = "ard:show",
+ BroadcastDate = DateTimeOffset.UtcNow, Duration = TimeSpan.FromMinutes(30),
+ };
+
+ private static ResolvedQuery Previous(int resultCount, TimeSpan age) => new()
+ {
+ Query = "tatort", ResultCount = resultCount, LastAttemptedAt = DateTimeOffset.UtcNow - age,
+ };
+
+ // ── the cache ─────────────────────────────────────────────
+
+ [Fact]
+ public async Task A_fresh_successful_resolution_is_not_repeated()
+ {
+ using var h = new Harness();
+ await h.StartAsync();
+ h.ResolvedQueries.GetAsync("tatort", Arg.Any())
+ .Returns(Previous(resultCount: 5, age: TimeSpan.FromMinutes(10)));
+
+ var resolved = await h.Resolver.EnsureResolvedAsync("Tatort", TestContext.Current.CancellationToken);
+
+ resolved.ShouldBeFalse();
+ h.Crawler.Calls.ShouldBe(0);
+ }
+
+ [Fact]
+ public async Task A_fresh_empty_resolution_is_not_repeated()
+ {
+ // The important half: Sonarr re-issues the same failing query every RSS-Sync cycle, so without
+ // negative caching each cycle would trigger a fresh multi-hop crawl of ARD.
+ using var h = new Harness();
+ await h.StartAsync();
+ h.ResolvedQueries.GetAsync("tatort", Arg.Any())
+ .Returns(Previous(resultCount: 0, age: TimeSpan.FromMinutes(5)));
+
+ await h.Resolver.EnsureResolvedAsync("Tatort", TestContext.Current.CancellationToken);
+
+ h.Crawler.Calls.ShouldBe(0);
+ }
+
+ [Fact]
+ public async Task A_stale_empty_resolution_is_retried()
+ {
+ using var h = new Harness();
+ await h.StartAsync();
+ h.ResolvedQueries.GetAsync("tatort", Arg.Any())
+ .Returns(Previous(resultCount: 0, age: TimeSpan.FromHours(2))); // past the 45m negative TTL
+ h.Crawler.Result = [Ep()];
+
+ var resolved = await h.Resolver.EnsureResolvedAsync("Tatort", TestContext.Current.CancellationToken);
+
+ resolved.ShouldBeTrue();
+ h.Crawler.Calls.ShouldBe(1);
+ }
+
+ [Fact]
+ public async Task A_stale_successful_resolution_is_retried()
+ {
+ using var h = new Harness();
+ await h.StartAsync();
+ h.ResolvedQueries.GetAsync("tatort", Arg.Any())
+ .Returns(Previous(resultCount: 3, age: TimeSpan.FromHours(12))); // past the 6h positive TTL
+
+ await h.Resolver.EnsureResolvedAsync("Tatort", TestContext.Current.CancellationToken);
+
+ h.Crawler.Calls.ShouldBe(1);
+ }
+
+ [Fact]
+ public void Query_normalisation_collapses_case_and_whitespace()
+ {
+ ResolvedQuery.Normalise(" Die Biene MAJA ").ShouldBe("die biene maja");
+ }
+
+ // ── resolution behaviour ──────────────────────────────────
+
+ [Fact]
+ public async Task Persists_what_the_crawlers_return_and_records_the_attempt()
+ {
+ using var h = new Harness();
+ await h.StartAsync();
+ h.Crawler.Result = [Ep("ard:1"), Ep("ard:2")];
+
+ var resolved = await h.Resolver.EnsureResolvedAsync("Tatort", TestContext.Current.CancellationToken);
+
+ resolved.ShouldBeTrue();
+ await h.Episodes.Received(1).UpsertManyAsync(
+ Arg.Is>(e => e != null && e.Count() == 2), Arg.Any());
+ await h.ResolvedQueries.Received(1).RecordAsync(
+ Arg.Is(q => q != null && q.Query == "tatort" && q.ResultCount == 2),
+ Arg.Any());
+ }
+
+ [Fact]
+ public async Task An_empty_crawl_is_recorded_as_a_miss_not_skipped()
+ {
+ using var h = new Harness();
+ await h.StartAsync();
+ h.Crawler.Result = [];
+
+ await h.Resolver.EnsureResolvedAsync("Nonexistent", TestContext.Current.CancellationToken);
+
+ await h.ResolvedQueries.Received(1).RecordAsync(
+ Arg.Is(q => q != null && q.ResultCount == 0), Arg.Any());
+ await h.Episodes.DidNotReceive().UpsertManyAsync(
+ Arg.Any>(), Arg.Any());
+ }
+
+ [Fact]
+ public async Task A_throwing_crawler_does_not_break_the_resolution()
+ {
+ using var h = new Harness();
+ await h.StartAsync();
+ h.Crawler.Throw = new HttpRequestException("ard is down");
+
+ var resolved = await h.Resolver.EnsureResolvedAsync("Tatort", TestContext.Current.CancellationToken);
+
+ // Still completes and still records the attempt, so one broken broadcaster cannot wedge search.
+ resolved.ShouldBeTrue();
+ await h.ResolvedQueries.Received(1).RecordAsync(
+ Arg.Any(), Arg.Any());
+ }
+
+ [Fact]
+ public async Task Concurrent_identical_queries_crawl_once()
+ {
+ using var h = new Harness();
+ await h.StartAsync();
+ h.Crawler.Gate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ h.Crawler.Result = [Ep()];
+
+ var first = h.Resolver.EnsureResolvedAsync("Tatort", TestContext.Current.CancellationToken);
+ var second = h.Resolver.EnsureResolvedAsync("tatort", TestContext.Current.CancellationToken);
+
+ h.Crawler.Gate.SetResult();
+ await Task.WhenAll(first, second);
+
+ h.Crawler.Calls.ShouldBe(1); // one crawl serving both callers
+ }
+
+ // ── the deadline, and the crawl outliving it ──────────────
+
+ [Fact]
+ public async Task The_request_deadline_releases_the_caller_while_the_crawl_continues()
+ {
+ using var h = new Harness(
+ new OnDemandResolutionOptions { CrawlTimeout = TimeSpan.FromSeconds(10) },
+ SearchWaitMode.ReturnFast, waitSeconds: 1);
+ await h.StartAsync();
+ h.Crawler.Gate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ h.Crawler.Result = [Ep()];
+
+ // The crawl is held open, so the wait must expire first.
+ var resolved = await h.Resolver.EnsureResolvedAsync("Tatort", TestContext.Current.CancellationToken);
+ resolved.ShouldBeFalse(); // released without a completed resolution
+ h.Crawler.Calls.ShouldBe(1); // ...but the crawl did start
+
+ // THE KEY ASSERTION: the request is long gone, yet letting the crawl finish still persists its
+ // episodes. If the request token reached the crawl, this upsert would never happen.
+ h.Crawler.Gate.SetResult();
+ await WaitUntilAsync(() => h.Episodes.ReceivedCalls().Any(
+ c => c.GetMethodInfo().Name == nameof(IEpisodeRepository.UpsertManyAsync)));
+
+ await h.Episodes.Received(1).UpsertManyAsync(
+ Arg.Any>(), Arg.Any());
+ }
+
+ [Fact]
+ public async Task Cancelling_the_request_does_not_cancel_the_crawl()
+ {
+ // Stronger than the deadline case: here the caller's token is actively cancelled, as it would be if
+ // Sonarr dropped the connection. The crawl must still finish and persist, because it runs under the
+ // host lifetime. If the request token were ever threaded into the crawl, this would upsert nothing.
+ using var h = new Harness(
+ new OnDemandResolutionOptions { CrawlTimeout = TimeSpan.FromSeconds(10) },
+ SearchWaitMode.ReturnFast, waitSeconds: 1);
+ await h.StartAsync();
+ h.Crawler.Gate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ h.Crawler.Result = [Ep()];
+
+ using var request = new CancellationTokenSource();
+ await h.Resolver.EnsureResolvedAsync("Tatort", request.Token);
+ await request.CancelAsync(); // the caller is gone
+
+ h.Crawler.Gate.SetResult();
+ await WaitUntilAsync(() => h.Episodes.ReceivedCalls().Any(
+ c => c.GetMethodInfo().Name == nameof(IEpisodeRepository.UpsertManyAsync)));
+
+ await h.Episodes.Received(1).UpsertManyAsync(
+ Arg.Any>(), Arg.Any