Skip to content

Akka.Streams: build island-boundary subscribers without runtime generics (Native AOT) - #8732

Merged
Aaronontheweb merged 3 commits into
akkadotnet:devfrom
Aaronontheweb:fix/streams-aot-island-boundaries
Oct 3, 2026
Merged

Aaronontheweb merged 3 commits into
akkadotnet:devfrom
Aaronontheweb:fix/streams-aot-island-boundaries

Conversation

@Aaronontheweb

@Aaronontheweb Aaronontheweb commented Oct 2, 2026 •

Copy link
Copy Markdown
Member

Akka.Streams built the types on every island boundary (BoundarySubscriber<T>, BoundaryPublisher<T>, ActorOutputBoundary<T>, the untyped wrappers, and a few more) with MakeGenericType + Activator. Native AOT trims those constructors, so any stream that crosses an island boundary failed with MissingMethodException. This PR builds each of them from code that already knows T, so AOT apps need no trimmer roots.

Changes

  • Inlet<T>/Outlet<T> get internal factories for BoundarySubscriber<T>, BoundaryPublisher<T>, ActorOutputBoundary<T> and MaterializedValueSource<T>; the materializer, GraphInterpreterShell.Init and Fusing call them instead of reflecting on the port type.
  • UntypedSubscriber/UntypedPublisher drop their reflective FromTyped(object) and gain typed operations for sink subscribers, the materialization-panic path (ErrorPublisher<T>/CancellingSubscriber<T>) and ActorSubscription.Create; ProcessorModule wraps its processor with its own TIn/TOut. Panic clean-up is now best effort, so the original materialization failure always surfaces.
  • The old reflective path stays only for port/publisher/subscriber types implemented outside Akka.Streams, gated on Akka.DynamicTypeLoading (RuntimeGenerics); with the switch off it throws NotSupportedException naming the switch.
  • SubscriberManagement's TStreamBuffer gets [DynamicallyAccessedMembers(PublicConstructors)], so Sink.AsPublisher(fanout: true) keeps its buffer constructor.
  • Hosting AOT canary: new StreamsScenarios.cs runs .Async(), AsPublisher→FromPublisher, fan-out, Source.ActorPublisher and a materialized-value graph, each with int and string. Docs: new "Akka.Streams Under Native AOT" section in native-aot.md; the stream-ref error now points to Stream refs under Native AOT: source-generated ISourceRef<T>/ISinkRef<T> fields (no element-type reflection) #8673.

Testing

  • Hosting AOT canary (linux-x64, no TrimmerRootDescriptor): on dev it fails with MissingMethodException on BoundarySubscriber<Int32> and emits 8 Akka.Streams IL warnings; with this PR it prints [canary-hosting] OK with 0 warnings (baseline stays empty).
  • Full Akka.Streams.Tests and Akka.API.Tests; new MaterializerSessionSpec covers the sink-subscriber and panic paths and the switch-off NotSupportedException.
  • Once this merges, the StreamsRoots.xml trimmer-root workaround in Persistence: register plugins in code through Akka.Persistence.Hosting (AOT) #8730's persistence canary (and the Embedded plugin canary) can be deleted.

Breaking changes

  • Obsoleted (not removed): Construct.Instantiate, TypeExtensions.GetSubscribedType/GetPublishedType.
  • Otherwise none expected; all other new members are internal.

Closes #8731

…ics (Native AOT)

The materializer built BoundarySubscriber<T>, BoundaryPublisher<T>,
ActorOutputBoundary<T>, MaterializedValueSource<T>, ActorSubscription<T>
and the UntypedSubscriberImpl<T>/UntypedPublisherImpl<T> wrappers with
MakeGenericType + Activator. Native AOT trims those constructors, so any
stream crossing an island boundary failed with MissingMethodException.

Every one of these types is now built by code that already knows T:
Inlet<T>/Outlet<T> (internal factories), the typed untyped-wrappers
(sink subscribers, panic-path ErrorPublisher/CancellingSubscriber,
ActorSubscription) and ProcessorModule. The old reflective path only
remains for port/publisher/subscriber types implemented outside
Akka.Streams, behind Akka.DynamicTypeLoading. SubscriberManagement's
stream-buffer type parameter gets [DynamicallyAccessedMembers] so the
fan-out publisher's buffer constructor survives trimming.

The Hosting AOT canary now runs streams across island boundaries
(value and reference element types) with no trimmer roots.

Closes akkadotnet#8731
@Aaronontheweb Aaronontheweb added this to the 1.6.0 milestone Oct 2, 2026
@Aaronontheweb Aaronontheweb added akka-streams AOT Ahead-of-Time (AOT) Compilation labels Oct 2, 2026
Aaronontheweb added a commit to Aaronontheweb/akka.net that referenced this pull request Oct 2, 2026
The persistence canary now registers a journal, a snapshot store and a stash
overflow configurator that Akka.Persistence does not ship and runs them with
Akka.DynamicTypeLoading off, next to the built-in inmem plugins. The comment
that called two sends a burst now says what the actor does.

Tests: add the unregistered snapshot store guard, switch-on tests for a
registered snapshot store, event adapter, binding and stash configurator, and
a run of PersistencePluginProxy with the switch off. Fold the exact-message
event adapter test into the existing one, keep Merge in its own test, and
make the read journal switch-on test assert the default config injection.

Docs: show how to add a read journal's default HOCON with the switch off, and
note in the baseline that Streams joins the gate once akkadotnet#8732 merges.

@Aaronontheweb Aaronontheweb left a comment •

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review: what each change and test is for. Inline comments have the detail.

What this PR does

  • Builds every island-boundary type from code that already knows T: Inlet<T>/Outlet<T> factories, typed methods on the untyped wrappers, and ProcessorModule. Native AOT never compiles a closed generic type that only MakeGenericType asks for.
  • Keeps the reflective path only for port, publisher and subscriber types implemented outside Akka.Streams, gated on Akka.DynamicTypeLoading.
  • Reproduced on 87fd339: the Hosting AOT canary publishes with 0 IL warnings and prints OK.

Overlap

  • The two MaterializerSessionSpec panic tests differ only by element type (matters under AOT, not on the JIT).
  • The Source.Queue canary scenario covers nothing the first .Async() scenario doesn't.

Not covered

  • The NotSupportedException path with the switch off (custom Inlet/Outlet subclass, foreign IUntypedPublisher/IUntypedSubscriber).
  • ActorSubscription.Create's typed path (only the legacy FanIn<T>/FanOut<T> actors reach it).
  • The IUntypedVirtualPublisher branches (nothing constructs VirtualPublisher<T>).

A cleanup pass addressing these is in progress.

Comment thread src/core/Akka.Streams/Shape.cs
/// <returns>TBD</returns>
public override Inlet CarbonCopy() => new Inlet<T>(Name);

internal override IUntypedSubscriber CreateBoundarySubscriber(IActorRef parent, GraphInterpreterShell shell, int id)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We build BoundarySubscriber<T> here because Inlet<T> knows T at compile time, so ILC compiles it for every element type the app uses. This replaces inlet.GetType().GetGenericArguments().First() + Instantiate in the materializer, the MissingMethodException from #8731.

public abstract Outlet CarbonCopy();


// INTERNAL API. The three factories below build the types that need this port's element type.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We give Outlet the same pattern for the three other types the materializer and fuser built by reflection: BoundaryPublisher<T>, ActorOutputBoundary<T> and MaterializedValueSource<T>. The base versions are the gated fallback; Outlet<T> overrides all three.

/// <returns>TBD</returns>
public override Outlet CarbonCopy() => new Outlet<T>(Name);

internal override IUntypedPublisher CreateBoundaryPublisher(IActorRef parent, GraphInterpreterShell shell, int id, out IActorPublisher actorPublisher)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We return the BoundaryPublisher<T> inside an UntypedPublisherImpl<T> and hand the raw publisher back through out for ExposedPublisher. Behavior change: the materializer's port-to-publisher map used to hold the BoundaryPublisher<T> itself (it implements IUntypedPublisher through ActorPublisher<T>); it now holds the wrapper, so DoSubscribe and the panic path can use the typed UntypedPublisher methods. Cost: one extra allocation per island outlet per materialization.

var t = module.CreateProcessor();
var processor = t.Item1;
var materialized = t.Item2;
var (subscriber, publisher, materialized) = module.CreateUntypedProcessor();

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We let ProcessorModule wrap its processor with its own TIn/TOut, because the old FromTyped(object) read the element types off the runtime type with GetInterfaces + MakeGenericType. The unused ProcessorForMethod lookup is also removed. Public ProcessorModule.CreateProcessor() stays for API compatibility, but nothing in the repo calls it now.

Comment thread src/aot/Akka.Hosting.AOT.App/StreamsScenarios.cs Outdated
.RunWith(Sink.Seq<int>(), mat),
Enumerable.Range(1, 5), "fan-out AsPublisher (int)");

// Source.ActorPublisher is what the persistence query read journals are built on.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We run Source.ActorPublisher because it proves the persistence-query shape from #8730 (a source module feeding a graph-stage island) via the same BoundarySubscriber<T> that failed there.

Comment thread src/aot/Akka.Hosting.AOT.App/StreamsScenarios.cs Outdated
queue.Complete();
await ExpectAsync(done, new[] { 1, 2, 3 }, "Source.Queue (int)");

// A graph that reads its own materialized value (an int) back as a stream element:

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We read the materialized value back as an element because it proves auto-fusing rebuilds MaterializedValueSource<T> via Outlet<T>.CreateMaterializedValueSource; int only.

#
# Empty since #8671 put the Akka.Streams stream-ref reflection behind Akka.DynamicTypeLoading. Any
# Empty since #8671 put the Akka.Streams stream-ref reflection behind Akka.DynamicTypeLoading, and
# #8731 stopped the materializer building island-boundary types with MakeGenericType. Any

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We keep the baseline empty because the publish now emits no Akka.Streams IL warnings. Reproduced on 87fd339 (dotnet publish -r linux-x64 -c Release): 0 IL warnings in total, canary prints [canary-hosting] OK.

- Materialization panic clean-up is best effort: a failing step (e.g. a
  foreign IUntypedSubscriber with Akka.DynamicTypeLoading off) no longer
  replaces the original exception or skips the remaining ports.
- Obsolete Construct.Instantiate (now [RequiresDynamicCode] and
  [RequiresUnreferencedCode]) and TypeExtensions.GetSubscribedType /
  GetPublishedType (now [DynamicallyAccessedMembers(Interfaces)]); API
  approvals updated.
- MaterializerSessionSpec: switch-off NotSupportedException specs for a
  non-generic Inlet/Outlet and a foreign IUntypedPublisher, a panic spec
  proving the original failure surfaces; int/string duplicate collapsed.
- Hosting AOT canary: every streams scenario runs with int and string,
  Source.Queue dropped, fan-out comment corrected.
- Stream-ref error message and docs point to akkadotnet#8673.
- docs: "Akka.Streams Under Native AOT" section in native-aot.md.
Aaronontheweb added a commit to Aaronontheweb/akka.net that referenced this pull request Oct 2, 2026
PersistenceSetup.Create().WithEmbeddedPersistence() registers the journal, snapshot
store and read journal by plugin id, each with its reference configuration, so a user
writes no class names and no WithFallback. WithEmbeddedJournal, WithEmbeddedSnapshotStore
and WithEmbeddedReadJournal take a plugin id; the journal takes event adapters. The
read journal provider gets its plugin id from the registration.

The AOT canary uses the new API and now also runs an event adapter, a Both-mode
database with a delete, and keeps StreamsRoots.xml until akkadotnet#8732 has merged.

@Aaronontheweb Aaronontheweb left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM - enables support for custom ISubscriber<T> and IPublisher<T> types to the extent that it's possible. Needed to make .Async() work in a number of cases.

// .Async() splits the graph into islands: BoundarySubscriber, BoundaryPublisher and
// ActorOutputBoundary on every edge between them.
await ExpectAsync(
Source.From(Enumerable.Range(1, 10)).Async().Select(x => x * 2).Async().RunWith(Sink.Seq<int>(), mat),

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Exercises the .Async() path inside the AOT system.

@Aaronontheweb
Aaronontheweb merged commit 0ad59d3 into akkadotnet:dev Oct 3, 2026
16 checks passed
Aaronontheweb added a commit that referenced this pull request Oct 3, 2026
…ings

#8732 builds island-boundary subscribers without runtime generics, so the
canary's queries run without StreamsRoots.xml. Delete it with its csproj item
and README section, and add src/core/Akka.Streams/ to the canary's gated
warning scope. The publish shows no Akka.Streams warnings.
Aaronontheweb added a commit that referenced this pull request Oct 5, 2026
…ings

#8732 builds island-boundary subscribers without runtime generics, so the
canary's queries run without StreamsRoots.xml. Delete it with its csproj item
and README section, and add src/core/Akka.Streams/ to the canary's gated
warning scope. The publish shows no Akka.Streams warnings.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

akka-streams AOT Ahead-of-Time (AOT) Compilation

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Akka.Streams: island boundary subscribers use runtime generics and fail under Native AOT

1 participant