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
33 changes: 33 additions & 0 deletions NeoReports.sln
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,12 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Destinations", "Destination
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "06-multi-sheet-xlsx", "samples\06-multi-sheet-xlsx\06-multi-sheet-xlsx.csproj", "{4481B01D-0C8B-4343-89ED-B55EA3F116C5}"
EndProject
Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Sources", "Sources", "{371C8358-20A2-CD07-4715-AFA0A31FB8ED}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "NeoReports.Sources.Join", "src\Sources\NeoReports.Sources.Join\NeoReports.Sources.Join.csproj", "{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "NeoReports.Sources.Join.UnitTests", "tests\NeoReports.Sources.Join.UnitTests\NeoReports.Sources.Join.UnitTests.csproj", "{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
Expand Down Expand Up @@ -403,6 +409,30 @@ Global
{4481B01D-0C8B-4343-89ED-B55EA3F116C5}.Release|x64.Build.0 = Release|Any CPU
{4481B01D-0C8B-4343-89ED-B55EA3F116C5}.Release|x86.ActiveCfg = Release|Any CPU
{4481B01D-0C8B-4343-89ED-B55EA3F116C5}.Release|x86.Build.0 = Release|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Debug|Any CPU.Build.0 = Debug|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Debug|x64.ActiveCfg = Debug|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Debug|x64.Build.0 = Debug|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Debug|x86.ActiveCfg = Debug|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Debug|x86.Build.0 = Debug|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Release|Any CPU.ActiveCfg = Release|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Release|Any CPU.Build.0 = Release|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Release|x64.ActiveCfg = Release|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Release|x64.Build.0 = Release|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Release|x86.ActiveCfg = Release|Any CPU
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B}.Release|x86.Build.0 = Release|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Debug|Any CPU.Build.0 = Debug|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Debug|x64.ActiveCfg = Debug|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Debug|x64.Build.0 = Debug|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Debug|x86.ActiveCfg = Debug|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Debug|x86.Build.0 = Debug|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Release|Any CPU.ActiveCfg = Release|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Release|Any CPU.Build.0 = Release|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Release|x64.ActiveCfg = Release|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Release|x64.Build.0 = Release|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Release|x86.ActiveCfg = Release|Any CPU
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F}.Release|x86.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
Expand Down Expand Up @@ -437,5 +467,8 @@ Global
{0C2934E0-ADD4-4DF2-9CD7-A3C5714F28DB} = {22222222-2222-2222-2222-222222222222}
{5785BC42-0EC0-98FB-3B2B-12E9D9C18C4E} = {11111111-1111-1111-1111-111111111111}
{4481B01D-0C8B-4343-89ED-B55EA3F116C5} = {44444444-4444-4444-4444-444444444444}
{371C8358-20A2-CD07-4715-AFA0A31FB8ED} = {11111111-1111-1111-1111-111111111111}
{D0B727FF-C7C3-4F1E-AFFE-7AF0C843BD6B} = {371C8358-20A2-CD07-4715-AFA0A31FB8ED}
{C4A8B5EF-482B-49B5-A0A7-1BE6A592C37F} = {22222222-2222-2222-2222-222222222222}
EndGlobalSection
EndGlobal
7 changes: 4 additions & 3 deletions PLAN.md
Original file line number Diff line number Diff line change
Expand Up @@ -163,9 +163,10 @@ Blueprint: [`docs/epic-b2-multisource.md`](docs/epic-b2-multisource.md); **D28**
source the existing pipeline consumes unchanged. Open sub-decisions (Pro vs free, package name, join
types, dynamic config, validation gate) in the doc must be settled before B2.3.

- [ ] **B2.1 — Enrichment** (`.Enrich(key, lookup, map)`): an `IBatchSource<TResult>` wrapper that
batch-looks-up related data once per page and maps it in (O(pageSize), no N+1). Tests: batched
per page, correct mapping, missing-key handling.
- [x] **B2.1 — Enrichment** (`.Enrich(key, lookup, map)`): `EnrichingBatchSource<...>` in a new
`NeoReports.Sources.Join` package (IsPackable=false; license/distribution deferred to B2.3) —
one batched lookup per page, O(pageSize), no N+1; a standard `IBatchSource<TResult>` the pipeline
consumes unchanged. ✅ 2 green tests (batched-per-page with distinct keys; missing-key → default).
- [ ] **B2.2 — Keyset merge-join** (`Source.MergeJoin(left, right, on, map)`): a streaming merge of
two same-key-ordered sources; inner + left-outer; constant memory (bounded key group). Tests:
ordered-merge correctness, memory, Testcontainers E2E across two SQL sources.
Expand Down
42 changes: 42 additions & 0 deletions src/Sources/NeoReports.Sources.Join/EnrichingBatchSource.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
using NeoReports.Abstractions;

namespace NeoReports.Sources.Join;

/// <summary>
/// Wraps a primary <see cref="IBatchSource{T}"/> and transforms each page into enriched result rows
/// via a page-level delegate (built by <see cref="Enrichment.Enrich{TPrimary, TKey, TLookup, TResult}"/>,
/// which does the batched key lookup). Keeps only one page in flight (O(pageSize)); cursor/paging come
/// straight from the primary.
/// </summary>
/// <typeparam name="TPrimary">The primary row type.</typeparam>
/// <typeparam name="TResult">The enriched result row type.</typeparam>
public sealed class EnrichingBatchSource<TPrimary, TResult> : IBatchSource<TResult>
{
private readonly IBatchSource<TPrimary> _primary;
private readonly Func<IReadOnlyList<TPrimary>, CancellationToken, Task<IReadOnlyList<TResult>>> _enrichPage;

/// <summary>Creates an enriching source.</summary>
/// <param name="primary">The primary source, read page by page.</param>
/// <param name="enrichPage">Transforms a page of primary rows into enriched result rows.</param>
public EnrichingBatchSource(
IBatchSource<TPrimary> primary,
Func<IReadOnlyList<TPrimary>, CancellationToken, Task<IReadOnlyList<TResult>>> enrichPage)
{
_primary = primary ?? throw new ArgumentNullException(nameof(primary));
_enrichPage = enrichPage ?? throw new ArgumentNullException(nameof(enrichPage));
}

/// <inheritdoc />
public ReportSchema Schema => _primary.Schema;

/// <inheritdoc />
public async Task<BatchResult<TResult>> ReadBatchAsync(BatchContext context, CancellationToken cancellationToken)
{
BatchResult<TPrimary> page = await _primary.ReadBatchAsync(context, cancellationToken).ConfigureAwait(false);
if (page.Records.Count == 0)
return new BatchResult<TResult>(Array.Empty<TResult>(), page.NextCursor, page.HasMore);

IReadOnlyList<TResult> results = await _enrichPage(page.Records, cancellationToken).ConfigureAwait(false);
return new BatchResult<TResult>(results, page.NextCursor, page.HasMore);
}
}
52 changes: 52 additions & 0 deletions src/Sources/NeoReports.Sources.Join/Enrichment.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
using System.Diagnostics.CodeAnalysis;
using NeoReports.Abstractions;

namespace NeoReports.Sources.Join;

/// <summary>Fluent entry point for enriching a source with a batched per-page lookup.</summary>
public static class Enrichment
{
/// <summary>
/// Enriches a primary source: for each page, one batched <paramref name="lookup"/> call resolves
/// the page's distinct keys, and <paramref name="map"/> combines each primary row with its
/// looked-up value (or <c>null</c> when absent). O(pageSize); no N+1.
/// </summary>
/// <typeparam name="TPrimary">The primary row type.</typeparam>
/// <typeparam name="TKey">The join key type.</typeparam>
/// <typeparam name="TLookup">The looked-up value type.</typeparam>
/// <typeparam name="TResult">The enriched result row type.</typeparam>
/// <param name="primary">The primary source.</param>
/// <param name="key">Extracts the join key from a primary row.</param>
/// <param name="lookup">Batched lookup: given a page's distinct keys, returns their values.</param>
/// <param name="map">Maps a primary row plus its looked-up value (or <c>null</c>) to the result.</param>
/// <returns>A source that yields enriched rows through the standard pipeline.</returns>
[SuppressMessage(
"Major Code Smell", "S2436:Types and methods should not have too many generic parameters",
Justification = "A batched join/enrichment API needs the primary, key, lookup-value and result types — " +
"the same four-type shape as the BCL's Enumerable.GroupJoin<TOuter,TInner,TKey,TResult>.")]
public static IBatchSource<TResult> Enrich<TPrimary, TKey, TLookup, TResult>(
this IBatchSource<TPrimary> primary,
Func<TPrimary, TKey> key,
Func<IReadOnlyList<TKey>, CancellationToken, Task<IReadOnlyDictionary<TKey, TLookup>>> lookup,
Func<TPrimary, TLookup?, TResult> map)
where TKey : notnull
{
ArgumentNullException.ThrowIfNull(primary);
ArgumentNullException.ThrowIfNull(key);
ArgumentNullException.ThrowIfNull(lookup);
ArgumentNullException.ThrowIfNull(map);

return new EnrichingBatchSource<TPrimary, TResult>(primary, async (records, cancellationToken) =>
{
var keys = records.Select(key).Distinct().ToList();
IReadOnlyDictionary<TKey, TLookup> values =
await lookup(keys, cancellationToken).ConfigureAwait(false) ?? new Dictionary<TKey, TLookup>();

return records.Select(record =>
{
values.TryGetValue(key(record), out TLookup? value);
return map(record, value);
}).ToArray();
});
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
<Project Sdk="Microsoft.NET.Sdk">

<PropertyGroup>
<TargetFrameworks>net8.0;net9.0</TargetFrameworks>
<Description>Multi-source composition for NeoReports: enrichment (batched per-page lookup) and keyset merge-join.</Description>
<PackageTags>reports;reporting;multi-source;join;enrichment</PackageTags>

<!-- Packaging & license are an open Epic B2 decision (Pro vs free — D29), settled in B2.3.
Until then this package is not auto-published. -->
<IsPackable>false</IsPackable>
</PropertyGroup>

<ItemGroup>
<ProjectReference Include="..\..\NeoReports.Abstractions\NeoReports.Abstractions.csproj" />
</ItemGroup>

</Project>
23 changes: 23 additions & 0 deletions src/Sources/NeoReports.Sources.Join/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
# NeoReports.Sources.Join

Multi-source composition for NeoReports. **B2.1 — enrichment** is here; the keyset **merge-join**
(B2.2) follows.

> Packaging and license (Pro vs free) are an open Epic B2 decision (**D29**), settled in B2.3. This
> package is not auto-published yet.

## Enrichment

For each page of a primary source, one **batched** lookup resolves the page's distinct keys (never one
call per row), then each row is mapped with its looked-up value:

```csharp
.From(Source.Sql(conn, sqlCustomers).Keyset<Customer, long>(c => c.Id)
.Enrich(
key: c => c.Id,
lookup: (keys, ct) => LoadOrderCountsAsync(keys, ct), // ONE call per page
map: (c, orderCount) => new CustomerSummary(c, orderCount)))
```

O(pageSize) memory; the batched-per-page shape structurally prevents the N+1 trap. The result is an
`IBatchSource<TResult>` the standard pipeline consumes unchanged.
95 changes: 95 additions & 0 deletions tests/NeoReports.Sources.Join.UnitTests/EnrichmentTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
using Microsoft.Extensions.Logging.Abstractions;
using NeoReports.Abstractions;
using NeoReports.Sources.Join;
using Shouldly;
using Xunit;

namespace NeoReports.Sources.Join.UnitTests;

/// <summary>
/// B2.1: enrichment reads the primary page by page and, per page, makes exactly one batched lookup
/// (never one per row), mapping each row with its looked-up value. Missing keys map to the default.
/// </summary>
public class EnrichmentTests
{
private static readonly int[] Page1Orders = { 10, 20 };
private static readonly int[] Page2Orders = { 30 };
private static readonly long[] Page1Keys = { 1, 2 };
private static readonly long[] Page2Keys = { 3 };

private sealed record Customer(long Id, string Name);

private sealed record CustomerSummary(long Id, string Name, int Orders);

private static BatchContext Ctx(int pageNumber, string? cursor) =>
new(new ReportExecutionContext("job", "r", null, NullLogger.Instance, CancellationToken.None), 10, cursor, pageNumber);

[Fact]
public async Task Enriches_each_page_with_a_single_batched_lookup()
{
var primary = new PagedCustomers(
new[] { new Customer(1, "A"), new Customer(2, "B") },
new[] { new Customer(3, "C") });

var lookupCalls = new List<IReadOnlyList<long>>();
IBatchSource<CustomerSummary> enriched = primary.Enrich(
key: c => c.Id,
lookup: (keys, _) =>
{
lookupCalls.Add(keys);
IReadOnlyDictionary<long, int> map = keys.ToDictionary(k => k, k => (int)(k * 10));
return Task.FromResult(map);
},
map: (c, orders) => new CustomerSummary(c.Id, c.Name, orders));

BatchResult<CustomerSummary> page1 = await enriched.ReadBatchAsync(Ctx(1, null), CancellationToken.None);
page1.Records.Select(r => r.Orders).ShouldBe(Page1Orders);
page1.HasMore.ShouldBeTrue();

BatchResult<CustomerSummary> page2 = await enriched.ReadBatchAsync(Ctx(2, page1.NextCursor), CancellationToken.None);
page2.Records.Select(r => r.Orders).ShouldBe(Page2Orders);
page2.HasMore.ShouldBeFalse();

// One batched call per page (not per row), with that page's distinct keys.
lookupCalls.Count.ShouldBe(2);
lookupCalls[0].ShouldBe(Page1Keys);
lookupCalls[1].ShouldBe(Page2Keys);
}

[Fact]
public async Task Missing_keys_map_to_the_default_lookup_value()
{
var primary = new PagedCustomers(new[] { new Customer(1, "A"), new Customer(2, "B") });

IBatchSource<CustomerSummary> enriched = primary.Enrich<Customer, long, int, CustomerSummary>(
key: c => c.Id,
lookup: (_, _) => Task.FromResult((IReadOnlyDictionary<long, int>)new Dictionary<long, int> { [1] = 99 }),
map: (c, orders) => new CustomerSummary(c.Id, c.Name, orders));

BatchResult<CustomerSummary> page = await enriched.ReadBatchAsync(Ctx(1, null), CancellationToken.None);

page.Records[0].Orders.ShouldBe(99); // found
page.Records[1].Orders.ShouldBe(0); // missing -> default(int)
}

/// <summary>In-memory primary that returns one supplied page per page number.</summary>
private sealed class PagedCustomers : IBatchSource<Customer>
{
private readonly IReadOnlyList<Customer>[] _pages;

public PagedCustomers(params IReadOnlyList<Customer>[] pages) => _pages = pages;

public ReportSchema Schema { get; } = new(new[] { new ReportColumn("Id", ColumnType.Integer) });

public Task<BatchResult<Customer>> ReadBatchAsync(BatchContext context, CancellationToken cancellationToken)
{
var index = context.PageNumber - 1;
if (index >= _pages.Length)
return Task.FromResult(BatchResult<Customer>.Empty);

var hasMore = index + 1 < _pages.Length;
var next = hasMore ? (context.PageNumber + 1).ToString(System.Globalization.CultureInfo.InvariantCulture) : null;
return Task.FromResult(new BatchResult<Customer>(_pages[index], next, hasMore));
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
<Project Sdk="Microsoft.NET.Sdk">

<PropertyGroup>
<TargetFramework>net8.0</TargetFramework>
<IsPackable>false</IsPackable>
<GenerateDocumentationFile>false</GenerateDocumentationFile>
</PropertyGroup>

<ItemGroup>
<ProjectReference Include="..\..\src\Sources\NeoReports.Sources.Join\NeoReports.Sources.Join.csproj" />
</ItemGroup>

<ItemGroup>
<PackageReference Include="Microsoft.NET.Test.Sdk" />
<PackageReference Include="xunit" />
<PackageReference Include="xunit.runner.visualstudio" />
<PackageReference Include="Shouldly" />
<PackageReference Include="Microsoft.Extensions.Logging.Abstractions" />
</ItemGroup>

</Project>