This is an automated email from the ASF dual-hosted git repository.
He-Pin pushed a change to branch feature/bidi-stream-closing-semantics
in repository https://gitbox.apache.org/repos/asf/pekko.git
omit 6a5992bf1c feat: Add BidirectionalGracefulShutdown CancellationStrategy
add 4d17e08c5b test: widen cluster shutdownAll await for aeron-udp drain
on JDK 25 nightly (#3017)
add 61c3d960bf optimize: avoid Holder allocation in ordered mapAsync for
already-completed futures (#3018)
add 59a2cc81af fix: await in-flight completion on partner termination in
SinkRef (#3016)
add 22e38951a9 Update jackson-core to 3.1.4 (#3021)
add 575a79a3f7 docs: add MiMa fix requirement to CLAUDE.md and AGENTS.md
(#3025)
add 42a4596780 fix: guard against int overflow in supervisor strategy and
frequency sketch (#3022)
add 7c1d0a623d Fix sbt load warnings by excluding unused lint keys (#3029)
add bbb1e8efc9 chore(stream): remove dead ActorProcessor,
ActorProcessorImpl, and ExposedPublisherReceive (#3027)
add 4e1508c38d chore(deps): bump sbt/setup-sbt from 1.1.24 to 1.2.1 (#3034)
add fbc05e26e8 chore(deps): bump scalacenter/sbt-dependency-submission
(#3033)
add e04e721ab6 Replace Scala string interpolation in log calls with lazy
{} placeholder format (#3032)
add a30df08865 2.0.0-M3 release notes (#3014)
add 27f6960dda ClusterShardingSettings: remove backward-compat passivation
cruft (#3026)
add 3e7b5e9257 Update netty-handler, netty-transport to 4.2.15.Final
(#3039)
add e012dd39ab Update logback-classic to 1.5.34 (#3038)
add 1dc09320aa fix: add missing space in StageActor PoisonPill/Kill
warning message (#3040)
add 153c151950 fix: widen RemoteSendConsistencySpec await window to 30s
(#3041)
add f6f988235b regenerate protobuf classes with protoc 4.35.0 (#3042)
add 70f6bbda56 Update biz.aQute.bndlib to 7.3.0 (#3044)
add 6f6fc8b7b5 Update jackson-core to 2.22.0 (#3045)
add 9a56d49ceb Update typesafe:config to 1.4.9 (#3046)
add 29dc5e8df3 chore: Rewrite to scala3 syntax (#3048)
add ca4b3aec02 test: retry temporary port bind in RemotingSpec
lazy-connect test (#3015)
add 6371e5b4be chore(deps): bump actions/checkout from 6.0.2 to 6.0.3
(#3053)
add bf96f5afa6 fix: correct typos in comments, docs, and string literals
across codebase (#3055)
add 0fc6120ddd Update Scala 3 version to 3.3.8 (#3056)
add 718c3f5675 Update protobuf-java to 4.35.1 (#3057)
add 52234c0199 Fix unresolvable scaladoc links to the ShardRegion actor
#353 (#3064)
add bf54798054 Add -Yfuture-lazy-vals for Scala 3.3.x builds (#3059)
add 28d61c5368 perf: replace ArrayList consumer wheel with LongMap for
O(1) keyed removal (#3063)
add 03ebaf59bb perf: optimize stream materializer wiring with HashMap and
ArrayList replacements (#3062)
add 645b76fbc9 Runtime plugin configuration for DurableState (#3058)
add b3219211ee Update sbt, scripted-plugin to 1.12.12 (#3069)
add 83e45001d2 Update sbt-api-mappings to 3.0.3 (#3068)
add f253c25de6 Update sbt-scalafix to 0.14.7 (#3067)
add 3e578d8a75 chore: upgrade sbt-java-formatter to 0.12.0 (#3066)
add f42ab033ab Update jackson-core to 3.2.0 (#3054)
add 400354ee8f aeron 1.51.0 (#3070)
add ba4e950edb Optimize lazy stage actor dispatch (#3035)
add ee6965cd58 fix: Scala 3.8 syntax compat fixes and import cleanup
(#3074)
add 6482aa1d9f refactor: rename Jdk9 sbt plugin to Jdk21 to reflect actual
purpose (#3075)
add 8a09821cc5 Add TCK tests for replay bounds across a deletion gap
(#3076)
add 51931817f6 Update sbt-mima-plugin to 1.1.6 (#3078)
add 4749eda97f refactor: replace private[this] with private and = _ with
explicit defaults for Scala 3.8 (#3072)
add 3cdbc7c8bb fix: ClassTag context bound compat for Scala 3.8
cross-build (#3073)
add 7b6b7ed6c1 make one of the new persistence-tck journal tests optional
(#3083)
add 005e242489 fix(actor-typed): discard stale ReceiveTimeout after
cancelReceiveTimeout to avoid NPE (#3086)
add 1d75d4e852 docs: Licensing Rules 顶部增加新文件 header 强制流程提示 (#3087)
add 4d942a0ca7 use assertInstanceof (#3088)
add 0830f58d28 fix: correct various typos across pekko modules (#3084)
add baa1229037 refactor: modernize BlogPostEntity Java test code with Java
17 records (#3085)
add 41b142dd23 perf(stream): optimize Source.actorRef with direct-push
fast path and message-only FunctionRef (#3091)
add 048c5f4a5a test: select the active FileSource actor in dispatcher
checks (#3002)
add 3424402ec1 pin sbt version entry in .scala-steward.conf (#3115)
add d75bf95017 fix: keep lazy dispatch state on VarHandle (#3116)
add b61d9d0b5e chore(deps): bump actions/checkout from 6.0.3 to 7.0.0
(#3117)
add 3eb757c061 chore(deps): bump sbt/setup-sbt from 1.2.1 to 1.4.0 (#3118)
add 85bc2aa1d3 perf(stream): extract reportStageError out of
GraphInterpreter.execute hot loop (#3119)
add e6d1533a72 Fix a few broken Scaladoc links reported by unidoc #353
(#3120)
add b09ee8b2f4 Update sbt-welcome to 0.6.0 (#3121)
add d1d2c1ef4c Update sbt, scripted-plugin to 1.12.13 (#3123)
add c0ed0aff7a Update sbt-sbom to 0.6.0 (#3122)
add f57cc814c7 fix(stream): log supervision Resume/Restart in MapAsync
operators (#3124)
add f571b64f78 refactor: use RandomGenerator interface instead of
java.util.Random (#3153)
add e1e2f2a4b1 refactor: use ProcessHandle.current().pid() instead of
RuntimeMXBean hack (#3152)
add f60702a4de refactor: replace manual stream copy loops with
transferTo/readAllBytes (#3146)
add fe5771080c refactor: use String.formatted() instead of String.format()
(#3144)
add 6628442f16 fix: remove redundant Float cast in CountMinSketch pattern
match (#3142)
add 71588172e8 refactor: use StackWalker for stack trace inspection in
TestKitUtils (#3141)
add 52755c47c0 refactor: use ByteBuffer fluent API chains with JDK 9
covariant returns (#3139)
add 3ac83fbd31 fix(stream): prevent double materialization of SourceRef
and SinkRef with clear error (#3114)
add 58a9bf0077 perf: use Arrays.copyOf/copyOfRange instead of manual new
Array + arraycopy (#3150)
add 6ad445271c refactor: use StandardCharsets.UTF_8 with
URLEncoder/URLDecoder (#3147)
add 0bbdbf78ec refactor: simplify Collections.unmodifiableSet(emptySet())
to Set.of() (#3138)
add acd4926ac1 refactor: use CompletableFuture.orTimeout in
FutureTimeoutSupport (#3137)
add b7ba750296 refactor: replace sun.reflect.Reflection with StackWalker
in Reflect (#3151)
add 22d44aa429 refactor: use Files.readString() in PemManagersProvider
(#3143)
add 88320c6978 refactor: use HexFormat for hex encoding instead of manual
formatting (#3145)
add 713ad01e16 perf: add Thread.onSpinWait() to CAS spin loops (#3149)
add d957aa9887 docs: update Scaladoc examples from Arrays.asList to
List.of (#3140)
add cd980eac30 use decodeString(StandardCharsets.UTF_8) (#3155)
add 640e790be5 Fix comparing release against pre-release versions (#3154)
add 72c92cc678 fix: use Timers for finishConnect retry instead of raw
scheduler callback (#3156)
add 1ba7efbe86 fix: initialize timer stage deadlines in preStart instead
of at construction (#3111)
add f3654268e7 Preserve throttler tokens on failed writes (#3097)
add 61e165c23d use Locale.ROOT for toLowerCase/toUpperCase (#3159)
add 5b2f4b468b Ignore stale classic remoting ACKs #3126 (#3158)
add c532ea043b fix: delay unfoldResourceAsync close until pending read
completes (#3157)
add 744dbd634f feat: replace CountMinSketch with FastFrequencySketch in
Artery compression (#3023)
add 63d22b60fe Clean up stopped FSM transition listeners (#3096)
add 9afe0e9109 fix(artery): distribute ActorSelection messages across
outbound lanes by target path (#3092)
add 14ed9a18b7 chore: add sbt-salad-days to reduce Scaladoc archive size
(#3100)
add 44e2aaecae test: make InputStreamSource TCK publisher finite (#2999)
add 74b9ab2dec fix: harden EndpointReader against NonFatal dispatch errors
and unwrap WrappedMessage in writer logs (#3166)
add 0a5529fda0 Update sbt-salad-days to 0.2.0 (#3171)
add 3757df0128 Update logback-classic to 1.5.35 (#3170)
add 07a346cc4b fix: validate full UniqueAddress in Artery system message
ACK/NACK handling (#3175)
add 69e5216f42 fix: avoid double close in unfoldResource (#3172)
add 08ca126ca9 fix: improve TLS observability - handshake failure logging
and fail-fast SSL init (#3165)
add 51d64ff1dd refactor: use pattern matching instanceof (Java 16+) (#3192)
add 9be70a423c refactor: apply misc Java 11-17 improvements (#3193)
add 96b716b6fb refactor: migrate switch statements to modern switch syntax
(Java 14+) (#3191)
add b50f8a6991 refactor: migrate String.format to String.formatted (Java
15+) (#3190)
add 4c4d0b7665 refactor: replace .collect(Collectors.toList()) with
.toList() (#3189)
add 78cd528313 refactor: migrate Collections.unmodifiable* to
List.copyOf/Set.copyOf/Map.copyOf (Java 10+) (#3194)
add 2e944ef1bf Update logback-classic to 1.5.37 (#3196)
add 526cb8c6d6 deprecation (#3197)
add 612851ec4e fix: add supervision strategy support for batch costFn
(#3184)
add a673f55e59 refactor: migrate Collections/Arrays.asList to Java 9+
collection factories (#3188)
add bf8adfb50e feat: add optional TLS hostname verification for classic
remoting (#3164)
add 7becc119d3 fix(remote): eager SSLContext init for artery
ConfigSSLEngineProvider fail-fast (#3200)
add 965b6b85c6 fix: ignore stale classic remoting reader ACKs (#3173)
add 0231702fd5 Remove fanout publisher sink wrapper (#3176)
add 22084da85c feat: add supervision strategy support for
groupedWeightedWithin costFn (#3180)
add 036327de8b fix(remote): restore eager SSLContext init for classic
remoting fail-fast (#3201)
add a626d6212a feat(stream): add withContext filtering and truncating
operators (#3199)
add b1148ae256 fix code to use correct
pekko.remote.classic.use-passive-connections name (#3207)
add 7328e559e0 test: stabilize mixed protocol netty ssl cluster join
(#3203)
add fdbf6a4317 feat: add supervision strategy support for delayWith (#3182)
add 07cda18251 feat: add supervision strategy support for doOnFirst (#3179)
add e077d55cef feat: add supervision strategy support for throttle
costCalculation (#3183)
add 99fe5f78a6 fix(remote): configure TLS resources for bind canonical
tests (#3210)
add 0b9755a079 fix: add supervision strategy support for zipWith and
zipWithN zipper functions (#3186)
add 1f8479bcfa feat: add supervision strategy support for
groupedAdjacentBy/Weighted (#3181)
add a993c16083 fix: add supervision strategy support for
aggregateWithBoundary callbacks (#3187)
add 9e5b807897 Update junit-jupiter-engine to 6.1.1 (#3260)
add 0874a04892 Update junit-platform-launcher to 6.1.1 (#3261)
add c0e9406bfd fix: expand Array arguments in MarkerLoggingAdapter log
templates (#3263)
add 64e26c8e51 docs: Clarify retry attempts parameter (#3266)
add cfaec70e81 fix: replace broken scalafmt-native-action with
coursier/setup-action (#3269)
add 84aa2b20e4 fix: add supervision strategy support for expand and
extrapolate (#3185)
add cf6784b952 fix: return proper thread-safe Cancellable from
StubbedActorContext.scheduleOnce (#3268)
add 26d8b60f19 feat(stream): add alsoTo overload with configurable
cancellation propagation (#3127)
add 6287b9de07 Include exception type in StatusReply toString when not
text (#3271)
add 5d1b67dedf docs: document preMaterialize error-propagation caveat
(#3272)
add b9bd7d0f5a feat: add FilteredPayload placeholder for filtered journal
events (#3275)
add 05b300dafb Port scalable fan-out persistence query abstraction from
Akka (#3277)
add 26c5086c7c docs: how to write a custom Durable State store plugin
(#3274)
add 5120ec6e48 docs: fix import order in state store plugin examples
(#3278)
add ba462fb3af Update aeron-client, aeron-driver to 1.51.1 (#3280)
add 2e9734bee9 Fix ambiguous scaladoc overload links in
persistence-typed #353 (#3281)
add 1a251f8510 Fix ambiguous scaladoc overload links in persistence module
family #353 (#3282)
add 0516cd6f40 Update sbt-reproducible-builds to 0.33 (#3287)
add c4a28ba211 Update aeron-client, aeron-driver to 1.52.0 (#3286)
add 59dadc172c fix: suppress Java compilation warnings in actor, remote
and docs modules (#3283)
add 8a3774add3 refactor: remove deprecated
RepointableActorRef.isTerminated (since Akka 2.2) (#3289)
add 40ad2491c1 fix: fix javadoc warnings in actor module Java sources
(#3284)
add b2ae2fafa7 fix: deliver stashed timer messages for UnboundedStash and
UnrestrictedStash (#3264)
add 3b21abca49 fix(remote): resume classic remoting reader after failed
handoff (#3205)
add bafe8a944d Fix timer current message for supervision #3265 (#3294)
add 535a5facc3 refactor: remove deprecated materializer settings APIs
(#3293)
add 2610cda2af perf: batch async boundary elements (#3288)
add e520a7c642 test: cover async boundary failure draining (#3296)
add ceca95fe0d Close unregistered TCP channels on setup failure (#3273)
add 7ad07e692c refactor: remove deprecated SubstreamCancelStrategy and
related splitWhen/splitAfter overloads (#3292)
add bbbf7b8bbe fix: fix scaladoc warnings (#3285)
add 79779612ab chore(deps): bump sbt/setup-sbt from 1.4.0 to 1.5.0 (#3298)
add 7f1a1b9140 refactor: replace deprecated `extends App` with explicit
main method (#3299)
add ce8e081a0b fix: fix compilation warnings across Scala 2.13 and 3.3
cross-builds (#3297)
add 4ae53394ac feat: add Scala 3.8.4 cross-compilation target (#3295)
add d7369dbd1e Update lz4-java to 1.11.1 (#3302)
add 8353a4f16c chore: use StandardCharsets.UTF_8 (#3301)
add 339537eec3 fix: avoid NPE in outbound TCP stage when upstream finishes
before connect (#3259)
add 17fb00e51e consistent-hashing shard allocation strategy (#3270)
add 44316a57f5 fix: release reference to initial interpreter shell after
it is shut down (#3031)
add c522cfa626 Update netty-handler, netty-transport to 4.2.16.Final
(#3307)
add c65e56f822 Update sbt-boilerplate to 0.8.1 (#3306)
add e157f92b14 feat: Add BidirectionalGracefulShutdown CancellationStrategy
This update added new revisions after undoing existing revisions.
That is to say, some revisions that were in the old version of the
branch are not in the new version. This situation occurs
when a user --force pushes a change and generates a repository
containing something like this:
* -- * -- B -- O -- O -- O (6a5992bf1c)
\
N -- N -- N refs/heads/feature/bidi-stream-closing-semantics
(e157f92b14)
You should already have received notification emails for all of the O
revisions, and so the following emails describe only the N revisions
from the common base, B.
Any revisions marked "omit" are not gone; other references still
refer to them. Any revisions marked "discard" are gone forever.
No new revisions were added by this update.
Summary of changes:
.github/scripts/resolve-scala-version.sh | 49 ++
.github/workflows/binary-compatibility-checks.yml | 11 +-
.github/workflows/build-test-prValidation.yml | 32 +-
.github/workflows/dependency-graph.yml | 8 +-
.github/workflows/format.yml | 11 +-
.github/workflows/generate-doc-check.yml | 6 +-
.github/workflows/headers.yml | 4 +-
.github/workflows/link-validator.yml | 4 +-
.github/workflows/nightly-1.0-builds.yml | 12 +-
.github/workflows/nightly-1.1-builds.yml | 12 +-
.github/workflows/nightly-1.2-builds.yml | 12 +-
.github/workflows/nightly-1.3-builds.yml | 12 +-
.github/workflows/nightly-1.4-builds.yml | 12 +-
.github/workflows/nightly-1.5-builds.yml | 12 +-
.github/workflows/nightly-1.6-builds.yml | 12 +-
.github/workflows/nightly-1.7-builds.yml | 12 +-
.github/workflows/nightly-builds-aeron.yml | 4 +-
.github/workflows/nightly-builds.yml | 12 +-
.github/workflows/publish-1.0-docs.yml | 4 +-
.github/workflows/publish-1.0-nightly.yml | 4 +-
.github/workflows/publish-1.1-docs.yml | 4 +-
.github/workflows/publish-1.1-nightly.yml | 4 +-
.github/workflows/publish-1.2-docs.yml | 4 +-
.github/workflows/publish-1.2-nightly.yml | 4 +-
.github/workflows/publish-1.3-docs.yml | 4 +-
.github/workflows/publish-1.3-nightly.yml | 4 +-
.github/workflows/publish-1.4-docs.yml | 4 +-
.github/workflows/publish-1.4-nightly.yml | 4 +-
.github/workflows/publish-1.5-docs.yml | 4 +-
.github/workflows/publish-1.5-nightly.yml | 4 +-
.github/workflows/publish-1.6-docs.yml | 4 +-
.github/workflows/publish-1.6-nightly.yml | 4 +-
.github/workflows/publish-1.7-nightly.yml | 4 +-
.github/workflows/publish-2.0-docs.yml | 16 +-
.github/workflows/publish-nightly.yml | 4 +-
.github/workflows/stage-release-candidate.yml | 8 +-
.github/workflows/timing-tests.yml | 4 +-
.scala-steward.conf | 4 +-
.scalafmt.conf | 12 +-
AGENTS.md | 5 +-
CLAUDE.md | 2 +-
CONTRIBUTING.md | 2 +-
.../apache/pekko/actor/testkit/typed/Effect.scala | 22 +-
.../actor/testkit/typed/TestKitSettings.scala | 10 +-
.../testkit/typed/internal/ActorSystemStub.scala | 2 +-
.../typed/internal/LoggingTestKitImpl.scala | 10 +-
.../typed/internal/StubbedActorContext.scala | 11 +-
.../testkit/typed/internal/TestKitUtils.scala | 6 +-
.../testkit/typed/internal/TestProbeImpl.scala | 2 +-
.../actor/testkit/typed/javadsl/ActorTestKit.scala | 8 +-
.../testkit/typed/javadsl/BehaviorTestKit.scala | 8 +-
.../typed/javadsl/JUnit5TestKitBuilder.scala | 4 +-
.../typed/javadsl/JUnitJupiterTestKitBuilder.scala | 4 +-
.../testkit/typed/javadsl/LoggingTestKit.scala | 6 +-
.../actor/testkit/typed/javadsl/ManualTime.scala | 2 +-
.../typed/javadsl/SerializationTestKit.scala | 2 +-
.../actor/testkit/typed/javadsl/TestInbox.scala | 2 +-
.../typed/javadsl/TestKitJUnit5Extension.scala | 2 +-
.../javadsl/TestKitJUnitJupiterExtension.scala | 2 +-
.../typed/javadsl/TestKitJunitResource.scala | 2 +-
.../actor/testkit/typed/javadsl/TestProbe.scala | 12 +-
.../testkit/typed/scaladsl/ActorTestKit.scala | 8 +-
.../testkit/typed/scaladsl/LoggingTestKit.scala | 2 +-
.../actor/testkit/typed/scaladsl/ManualTime.scala | 4 +-
.../typed/scaladsl/ScalaTestWithActorTestKit.scala | 2 +-
.../typed/scaladsl/SerializationTestKit.scala | 2 +-
.../actor/testkit/typed/scaladsl/TestInbox.scala | 2 +-
.../actor/testkit/typed/scaladsl/TestProbe.scala | 4 +-
.../typed/javadsl/AsyncTestingExampleTest.java | 5 +-
.../typed/javadsl/SyncTestingExampleTest.java | 105 +--
.../testkit/typed/javadsl/TestConfigExample.java | 12 +-
.../testkit/typed/javadsl/BehaviorTestKitTest.java | 205 ++---
.../actor/testkit/typed/javadsl/TestProbeTest.java | 3 +-
.../typed/scaladsl/SyncTestingExampleSpec.scala | 4 +-
.../testkit/typed/scaladsl/ActorTestKitSpec.scala | 4 +-
.../typed/scaladsl/BehaviorTestKitSpec.scala | 42 +-
.../apache/pekko/actor/AbstractFSMActorTest.java | 3 +-
.../org/apache/pekko/actor/ActorCreationTest.java | 2 +-
.../apache/pekko/actor/StashJavaAPITestActors.java | 1 +
.../pekko/dispatch/CompletionStagesTests.java | 7 +-
.../java/org/apache/pekko/event/ActorWithMDC.java | 16 +-
.../org/apache/pekko/event/LoggingAdapterTest.java | 4 +-
.../apache/pekko/japi/pf/ReceiveBuilderTest.java | 31 +-
.../org/apache/pekko/pattern/PatternsTest.java | 40 +-
.../org/apache/pekko/util/LineNumberSpec.scala | 4 +-
.../org/apache/pekko/PekkoExceptionSpec.scala | 2 +-
.../actor/AbstractActorPreRestartFinalSpec.scala | 6 +-
.../org/apache/pekko/actor/ActorMailboxSpec.scala | 2 +-
.../org/apache/pekko/actor/DynamicAccessSpec.scala | 2 +-
.../org/apache/pekko/actor/ExtensionSpec.scala | 2 +-
.../org/apache/pekko/actor/FSMTransitionSpec.scala | 28 +-
.../apache/pekko/actor/RelativeActorPathSpec.scala | 3 +-
.../pekko/actor/SupervisorHierarchySpec.scala | 4 +-
.../org/apache/pekko/actor/SupervisorSpec.scala | 26 +
.../scala/org/apache/pekko/actor/TimerSpec.scala | 96 +++
.../apache/pekko/dispatch/MailboxConfigSpec.scala | 2 +-
.../org/apache/pekko/event/MarkerLoggingSpec.scala | 18 +
.../org/apache/pekko/io/TcpConnectionSpec.scala | 42 ++
.../org/apache/pekko/io/TcpIntegrationSpec.scala | 9 +-
.../pekko/io/UdpConnectedIntegrationSpec.scala | 2 +-
.../apache/pekko/io/dns/DockerBindDnsService.scala | 4 +-
.../scala/org/apache/pekko/pattern/RetrySpec.scala | 37 +-
.../org/apache/pekko/pattern/StatusReplySpec.scala | 6 +-
.../org/apache/pekko/routing/BalancingSpec.scala | 3 +-
.../serialization/SerializationSetupSpec.scala | 2 +-
.../apache/pekko/serialization/SerializeSpec.scala | 12 +-
.../pekko/util/BoundedBlockingQueueSpec.scala | 2 +-
.../pekko/util/ByteStringInitializationSpec.scala | 2 +-
.../org/apache/pekko/util/ByteStringSpec.scala | 4 +-
.../apache/pekko/util/FrequencySketchSpec.scala | 42 ++
.../scala/org/apache/pekko/util/ReflectSpec.scala | 28 +
.../scala/org/apache/pekko/util/SWARUtilSpec.scala | 2 +-
.../scala/org/apache/pekko/util/VersionSpec.scala | 1 +
.../org/apache/pekko/util/WildcardIndexSpec.scala | 6 +-
.../jdocs/org/apache/pekko/typed/Aggregator.java | 1 +
.../org/apache/pekko/typed/AggregatorTest.java | 82 +-
.../org/apache/pekko/typed/BubblingSample.java | 22 +-
.../jdocs/org/apache/pekko/typed/FSMDocTest.java | 2 +-
.../apache/pekko/typed/GracefulStopDocTest.java | 2 +
.../InteractionPatternsAskWithStatusTest.java | 2 +
.../pekko/typed/InteractionPatternsTest.java | 256 ++-----
.../jdocs/org/apache/pekko/typed/IntroTest.java | 10 +
.../org/apache/pekko/typed/LoggingDocExamples.java | 3 +
.../jdocs/org/apache/pekko/typed/OOIntroTest.java | 5 +
.../apache/pekko/typed/SpawnProtocolDocTest.java | 1 +
.../org/apache/pekko/typed/StashDocSample.java | 44 +-
.../apache/pekko/typed/StyleGuideDocExamples.java | 211 +-----
.../jdocs/org/apache/pekko/typed/TailChopping.java | 1 +
.../coexistence/ClassicWatchingTypedTest.java | 11 +-
.../pekko/typed/fromclassic/TypedSample.java | 28 +-
.../eventstream/EventStreamSuperClassDocTest.java | 1 +
.../typed/javadsl/ActorContextPipeToSelfTest.java | 4 +-
.../pekko/typed/InteractionPatterns3Spec.scala | 22 +-
.../pekko/typed/InteractionPatternsSpec.scala | 4 +-
.../apache/pekko/typed/LoggingDocExamples.scala | 4 +-
.../apache/pekko/typed/StyleGuideDocExamples.scala | 2 +-
.../pekko/typed/extensions/ExtensionDocSpec.scala | 6 +-
.../pekko/actor/typed/ActorContextSpec.scala | 5 +-
.../org/apache/pekko/actor/typed/AskSpec.scala | 2 +-
.../actor/typed/CancelReceiveTimeoutSpec.scala | 106 +++
.../apache/pekko/actor/typed/ExtensionsSpec.scala | 20 +-
.../LocalActorRefProviderLogMessagesSpec.scala | 8 +-
.../pekko/actor/typed/MailboxSelectorSpec.scala | 2 +-
.../apache/pekko/actor/typed/SupervisionSpec.scala | 6 +-
.../pekko/actor/typed/TransformMessagesSpec.scala | 3 +-
.../delivery/DurableProducerControllerSpec.scala | 4 +-
.../typed/delivery/DurableWorkPullingSpec.scala | 4 +-
.../delivery/ReliableDeliveryRandomSpec.scala | 4 +-
.../actor/typed/internal/ActorSystemSpec.scala | 2 +-
.../{adpater => adapter}/PropsAdapterSpec.scala | 2 +-
.../receptionist/LocalReceptionistSpec.scala | 2 +-
.../actor/typed/scaladsl/MailboxSelectorSpec.scala | 4 +-
.../typed/internal/receptionist/Platform.scala | 4 +-
.../org/apache/pekko/actor/typed/ActorRef.scala | 2 +-
.../pekko/actor/typed/ActorRefResolver.scala | 10 +-
.../org/apache/pekko/actor/typed/ActorSystem.scala | 2 +-
.../org/apache/pekko/actor/typed/Behavior.scala | 6 +-
.../pekko/actor/typed/BehaviorInterceptor.scala | 8 +-
.../org/apache/pekko/actor/typed/Extensions.scala | 10 +-
.../apache/pekko/actor/typed/SpawnProtocol.scala | 2 +-
.../pekko/actor/typed/SupervisorStrategy.scala | 4 +-
.../actor/typed/delivery/ConsumerController.scala | 4 +-
.../typed/delivery/DurableProducerQueue.scala | 4 +-
.../actor/typed/delivery/ProducerController.scala | 10 +-
.../delivery/WorkPullingProducerController.scala | 10 +-
.../delivery/internal/ConsumerControllerImpl.scala | 18 +-
.../delivery/internal/ProducerControllerImpl.scala | 4 +-
.../WorkPullingProducerControllerImpl.scala | 2 +-
.../actor/typed/eventstream/EventStream.scala | 5 +-
.../actor/typed/internal/ActorContextImpl.scala | 16 +-
.../actor/typed/internal/ActorFlightRecorder.scala | 2 +-
.../pekko/actor/typed/internal/ActorRefImpl.scala | 4 +-
.../typed/internal/EventStreamExtension.scala | 4 +-
.../actor/typed/internal/ExtensionsImpl.scala | 14 +-
.../actor/typed/internal/InterceptorImpl.scala | 10 +-
.../pekko/actor/typed/internal/LoggerClass.scala | 10 +-
.../typed/internal/MiscMessageSerializer.scala | 4 +-
.../pekko/actor/typed/internal/PoisonPill.scala | 2 +-
.../pekko/actor/typed/internal/Supervision.scala | 21 +-
.../internal/WithMdcBehaviorInterceptor.scala | 4 +-
.../typed/internal/adapter/ActorAdapter.scala | 14 +-
.../internal/adapter/ActorContextAdapter.scala | 8 +-
.../typed/internal/adapter/ActorRefAdapter.scala | 4 +-
.../internal/adapter/ActorSystemAdapter.scala | 10 +-
.../typed/internal/adapter/PropsAdapter.scala | 2 +-
.../pekko/actor/typed/internal/jfr/Events.scala | 3 +-
.../internal/receptionist/LocalReceptionist.scala | 8 +-
.../internal/receptionist/ReceptionistImpl.scala | 2 +-
.../receptionist/ReceptionistMessages.scala | 6 +-
.../receptionist/ServiceKeySerializer.scala | 4 +-
.../typed/internal/routing/GroupRouterImpl.scala | 2 +-
.../typed/internal/routing/PoolRouterImpl.scala | 2 +-
.../pekko/actor/typed/javadsl/ActorContext.scala | 2 +-
.../apache/pekko/actor/typed/javadsl/Adapter.scala | 18 +-
.../pekko/actor/typed/javadsl/AskPattern.scala | 2 +-
.../actor/typed/javadsl/BehaviorBuilder.scala | 4 +-
.../pekko/actor/typed/javadsl/Behaviors.scala | 19 +-
.../pekko/actor/typed/javadsl/ReceiveBuilder.scala | 4 +-
.../apache/pekko/actor/typed/pubsub/Topic.scala | 6 +-
.../actor/typed/receptionist/Receptionist.scala | 26 +-
.../pekko/actor/typed/scaladsl/ActorContext.scala | 2 +-
.../pekko/actor/typed/scaladsl/AskPattern.scala | 10 +-
.../actor/typed/scaladsl/adapter/package.scala | 8 +-
.../pekko/japi/function/Functions.scala.template | 4 +-
.../io/UnsynchronizedByteArrayInputStream.java | 4 +-
.../org/apache/pekko/japi/pf/AbstractMatch.java | 4 +-
.../apache/pekko/japi/pf/AbstractPFBuilder.java | 4 +-
.../pekko/japi/pf/FSMStateFunctionBuilder.java | 2 +-
.../org/apache/pekko/japi/pf/FSMStopBuilder.java | 2 +-
.../pekko/japi/pf/FSMTransitionHandlerBuilder.java | 2 +-
.../main/java/org/apache/pekko/japi/pf/Match.java | 2 +-
.../java/org/apache/pekko/japi/pf/PFBuilder.java | 2 +-
.../org/apache/pekko/japi/pf/ReceiveBuilder.java | 2 +-
.../java/org/apache/pekko/japi/pf/UnitMatch.java | 4 +-
.../org/apache/pekko/japi/pf/UnitPFBuilder.java | 4 +-
.../java/org/apache/pekko/util/OptionalUtil.java | 5 +-
.../fsm-transition-listeners.excludes} | 4 +-
actor/src/main/resources/reference.conf | 2 +-
.../org/apache/pekko/util/ByteIterator.scala | 159 ++--
.../org/apache/pekko/actor/AbstractProps.scala | 10 +-
.../main/scala/org/apache/pekko/actor/Actor.scala | 2 +-
.../scala/org/apache/pekko/actor/ActorCell.scala | 8 +-
.../scala/org/apache/pekko/actor/ActorPath.scala | 2 +-
.../scala/org/apache/pekko/actor/ActorRef.scala | 4 +-
.../org/apache/pekko/actor/ActorRefProvider.scala | 10 +-
.../org/apache/pekko/actor/ActorSelection.scala | 6 +-
.../scala/org/apache/pekko/actor/ActorSystem.scala | 16 +-
.../apache/pekko/actor/CoordinatedShutdown.scala | 4 +-
.../scala/org/apache/pekko/actor/Deployer.scala | 8 +-
.../org/apache/pekko/actor/DynamicAccess.scala | 6 +-
.../scala/org/apache/pekko/actor/Extension.scala | 2 +-
.../main/scala/org/apache/pekko/actor/FSM.scala | 52 +-
.../org/apache/pekko/actor/FaultHandling.scala | 22 +-
.../apache/pekko/actor/IndirectActorProducer.scala | 18 +-
.../pekko/actor/LightArrayRevolverScheduler.scala | 6 +-
.../main/scala/org/apache/pekko/actor/Props.scala | 14 +-
.../pekko/actor/ReflectiveDynamicAccess.scala | 10 +-
.../apache/pekko/actor/RepointableActorRef.scala | 15 +-
.../main/scala/org/apache/pekko/actor/Stash.scala | 2 +-
.../main/scala/org/apache/pekko/actor/Timers.scala | 7 +-
.../pekko/actor/dungeon/ChildrenContainer.scala | 4 +-
.../org/apache/pekko/actor/dungeon/Dispatch.scala | 2 +-
.../pekko/actor/dungeon/TimerSchedulerImpl.scala | 2 +-
.../pekko/actor/setup/ActorSystemSetup.scala | 2 +-
.../apache/pekko/dispatch/AbstractDispatcher.scala | 34 +-
.../apache/pekko/dispatch/BatchingExecutor.scala | 10 +-
.../apache/pekko/dispatch/CompletionStages.scala | 18 +-
.../org/apache/pekko/dispatch/Dispatchers.scala | 2 +-
.../dispatch/ForkJoinExecutorConfigurator.scala | 6 +-
.../scala/org/apache/pekko/dispatch/Mailbox.scala | 86 ++-
.../org/apache/pekko/dispatch/Mailboxes.scala | 26 +-
.../apache/pekko/dispatch/ThreadPoolBuilder.scala | 8 +-
.../dispatch/VirtualizedExecutorService.scala | 10 +-
.../pekko/dispatch/affinity/AffinityPool.scala | 28 +-
.../pekko/dispatch/sysmsg/SystemMessage.scala | 2 +-
.../scala/org/apache/pekko/event/EventStream.scala | 14 +-
.../scala/org/apache/pekko/event/Logging.scala | 161 ++--
.../org/apache/pekko/io/DirectByteBufferPool.scala | 4 +-
.../scala/org/apache/pekko/io/DnsProvider.scala | 4 +-
.../org/apache/pekko/io/SelectionHandler.scala | 18 +-
.../scala/org/apache/pekko/io/SimpleDnsCache.scala | 4 +-
actor/src/main/scala/org/apache/pekko/io/Tcp.scala | 8 +-
.../scala/org/apache/pekko/io/TcpConnection.scala | 27 +-
.../apache/pekko/io/TcpOutgoingConnection.scala | 15 +-
actor/src/main/scala/org/apache/pekko/io/Udp.scala | 2 +-
.../scala/org/apache/pekko/io/UdpConnection.scala | 5 +-
.../scala/org/apache/pekko/io/UdpListener.scala | 5 +-
.../org/apache/pekko/io/dns/DnsProtocol.scala | 2 +-
.../org/apache/pekko/io/dns/IdGenerator.scala | 6 +-
.../main/scala/org/apache/pekko/japi/JavaAPI.scala | 8 +-
.../scala/org/apache/pekko/japi/Throwables.scala | 2 +-
.../org/apache/pekko/japi/function/Function.scala | 14 +-
.../org/apache/pekko/pattern/AskSupport.scala | 27 +-
.../org/apache/pekko/pattern/BackoffOptions.scala | 6 +-
.../apache/pekko/pattern/BackoffSupervisor.scala | 4 +-
.../org/apache/pekko/pattern/CircuitBreaker.scala | 38 +-
.../pekko/pattern/CircuitBreakersRegistry.scala | 2 +-
.../pekko/pattern/FutureTimeoutSupport.scala | 32 +-
.../scala/org/apache/pekko/pattern/Patterns.scala | 67 +-
.../org/apache/pekko/pattern/RetrySupport.scala | 24 +-
.../org/apache/pekko/pattern/StatusReply.scala | 12 +-
.../internal/BackoffOnRestartSupervisor.scala | 2 +-
.../pattern/internal/CircuitBreakerTelemetry.scala | 4 +-
.../scala/org/apache/pekko/routing/Balancing.scala | 2 +-
.../scala/org/apache/pekko/routing/Broadcast.scala | 2 +-
.../org/apache/pekko/routing/ConsistentHash.scala | 3 +-
.../scala/org/apache/pekko/routing/Random.scala | 2 +-
.../org/apache/pekko/routing/RoundRobin.scala | 2 +-
.../org/apache/pekko/routing/RoutedActorRef.scala | 2 +-
.../org/apache/pekko/routing/RouterConfig.scala | 6 +-
.../org/apache/pekko/routing/SmallestMailbox.scala | 6 +-
.../pekko/serialization/PrimitiveSerializers.scala | 10 +-
.../apache/pekko/serialization/Serialization.scala | 30 +-
.../pekko/serialization/SerializationSetup.scala | 6 +-
.../apache/pekko/serialization/Serializer.scala | 24 +-
.../apache/pekko/util/BoundedBlockingQueue.scala | 14 +-
.../scala/org/apache/pekko/util/BoxedType.scala | 4 +-
.../scala/org/apache/pekko/util/ByteString.scala | 33 +-
.../pekko/util/ClassLoaderObjectInputStream.scala | 2 +-
.../scala/org/apache/pekko/util/Collections.scala | 4 +-
.../scala/org/apache/pekko/util/ConstantFun.scala | 2 +-
.../org/apache/pekko/util/DoubleLinkedList.scala | 6 +-
.../org/apache/pekko/util/FrequencySketch.scala | 50 +-
.../main/scala/org/apache/pekko/util/Helpers.scala | 2 +-
.../org/apache/pekko/util/ImmutableIntMap.scala | 6 +-
.../scala/org/apache/pekko/util/LineNumbers.scala | 4 +-
.../main/scala/org/apache/pekko/util/Reflect.scala | 29 +-
.../apache/pekko/util/StablePriorityQueue.scala | 2 +-
.../org/apache/pekko/util/SubclassifiedIndex.scala | 2 +-
.../scala/org/apache/pekko/util/TokenBucket.scala | 6 +-
.../org/apache/pekko/util/TypedMultiMap.scala | 2 +-
.../main/scala/org/apache/pekko/util/Version.scala | 2 +-
.../pekko/actor/typed/TypedBenchmarkActors.scala | 2 +-
.../pekko/remote/artery/CodecBenchmark.scala | 2 +-
.../ActorGraphInterpreterBoundaryBenchmark.scala | 2 +-
.../pekko/stream/ActorRefSourceBenchmark.scala | 107 +++
.../stream/AsyncBoundaryThroughputBenchmark.scala | 123 +++
.../pekko/stream/BroadcastHubBenchRunner.scala | 113 +++
.../pekko/stream/BroadcastHubBenchmark.scala | 50 +-
.../org/apache/pekko/stream/FlowMapBenchmark.scala | 2 +-
.../apache/pekko/stream/FusedGraphsBenchmark.scala | 4 +-
.../stream/GraphStageConstructionBenchmark.scala | 15 +-
.../pekko/stream/MaterializerWiringBenchmark.scala | 132 ++++
.../apache/pekko/stream/RangeSourceBenchmark.scala | 2 +-
.../pekko/stream/StageActorRefBenchmark.scala | 136 ++++
.../org/apache/pekko/stream/io/TlsBenchmark.scala | 3 +-
.../util/ByteStringParser_readNum_Benchmark.scala | 4 +-
.../pekko/util/FastFrequencySketchBenchmark.scala | 10 +-
.../pekko/util/FrequencySketchBenchmark.scala | 10 +-
.../apache/pekko/util/ImmutableIntMapBench.scala | 12 +-
build.sbt | 17 +-
.../protobuf/msg/ClusterMetricsMessages.java | 230 +++---
.../cluster/metrics/ClusterMetricsCollector.scala | 10 +-
.../cluster/metrics/ClusterMetricsExtension.scala | 4 +-
.../cluster/metrics/ClusterMetricsStrategy.scala | 2 +-
.../pekko/cluster/metrics/MetricsCollector.scala | 2 +-
.../metrics/protobuf/MessageSerializer.scala | 14 +-
.../metrics/protobuf/NumberInputStream.scala | 2 +-
.../typed/internal/protobuf/ShardingMessages.java | 110 +--
.../remove-old-passivation-strategy.excludes | 6 +-
.../sharding/typed/ClusterShardingQuery.scala | 8 +-
.../sharding/typed/ClusterShardingSettings.scala | 69 +-
.../typed/ReplicatedShardingExtension.scala | 4 +-
.../typed/ShardedDaemonProcessSettings.scala | 6 +-
.../sharding/typed/ShardingMessageExtractor.scala | 2 +-
.../delivery/ShardingConsumerController.scala | 4 +-
.../delivery/ShardingProducerController.scala | 10 +-
.../internal/ShardingConsumerControllerImpl.scala | 2 +-
.../internal/ShardingProducerControllerImpl.scala | 4 +-
.../typed/internal/ClusterShardingImpl.scala | 21 +-
.../internal/ReplicatedShardingExtensionImpl.scala | 2 +-
.../internal/ShardedDaemonProcessCoordinator.scala | 2 +-
.../typed/internal/ShardedDaemonProcessImpl.scala | 39 +-
.../typed/internal/ShardingSerializer.scala | 6 +-
.../typed/internal/testkit/TestEntityRefImpl.scala | 4 +-
.../sharding/typed/javadsl/ClusterSharding.scala | 6 +-
.../typed/javadsl/ShardedDaemonProcess.scala | 4 +-
.../sharding/typed/scaladsl/ClusterSharding.scala | 8 +-
.../typed/scaladsl/ShardedDaemonProcess.scala | 2 +-
.../typed/delivery/DeliveryThroughputSpec.scala | 2 +-
.../jdocs/delivery/PointToPointDocExample.java | 38 +-
.../java/jdocs/delivery/ShardingDocExample.java | 128 +---
.../java/jdocs/delivery/WorkPullingDocExample.java | 132 +---
.../sharding/typed/AccountExampleDocTest.java | 15 +-
.../AccountExamplePersistenceProbeDocTest.java | 219 +++---
.../cluster/sharding/typed/AccountExampleTest.java | 8 +-
.../AccountExampleWithEventHandlersInState.java | 106 +--
.../typed/AccountExampleWithMutableState.java | 106 +--
.../typed/AccountExampleWithNullDurableState.java | 81 +-
.../typed/AccountExampleWithNullState.java | 106 +--
...stentHashingShardAllocationCompileOnlyTest.java | 141 ++++
.../typed/HelloWorldPersistentEntityExample.java | 3 +-
.../typed/ReplicatedShardingCompileOnlySpec.java | 4 +-
.../sharding/typed/ShardingCompileOnlyTest.java | 22 +-
.../typed/ShardingReplyCompileOnlyTest.java | 18 +-
.../sharding/typed/ReplicatedShardingTest.java | 7 +-
.../ShardedDaemonProcessCompileOnlyTest.java | 3 +-
...edEntityWithEnforcedRepliesCompileOnlyTest.java | 2 +-
...tentHashingShardAllocationCompileOnlySpec.scala | 94 +++
...urableStateStoreQueryUsageCompileOnlySpec.scala | 2 +-
.../ExternalShardAllocationCompileOnlySpec.scala | 2 +-
.../typed/HelloWorldPersistentEntityExample.scala | 2 +-
.../typed/ReplicatedShardingCompileOnlySpec.scala | 2 +-
.../pekko/cluster/sharding/FlightRecording.scala | 2 +-
...oinConfigCompatCheckerClusterShardingSpec.scala | 4 +-
.../sharding/typed/ReplicatedShardingSpec.scala | 2 +-
.../typed/delivery/DurableShardingSpec.scala | 8 +-
.../delivery/ReliableDeliveryShardingSpec.scala | 2 +-
.../internal/ShardedDaemonProcessIdSpec.scala | 1 +
.../scaladsl/ClusterShardingPersistenceSpec.scala | 2 +-
.../protobuf/msg/ClusterShardingMessages.java | 588 +++++++++------
.../consistent-hashing-allocation.excludes} | 12 +-
.../remove-old-passivation-strategy.excludes | 8 +-
cluster-sharding/src/main/resources/reference.conf | 6 -
.../pekko/cluster/sharding/ClusterSharding.scala | 86 +--
.../cluster/sharding/ClusterShardingSettings.scala | 80 +-
.../ConsistentHashingShardAllocationStrategy.scala | 163 ++++
.../org/apache/pekko/cluster/sharding/Shard.scala | 3 +-
.../pekko/cluster/sharding/ShardCoordinator.scala | 15 +-
.../pekko/cluster/sharding/ShardRegion.scala | 31 +-
.../cluster/sharding/ShardingFlightRecorder.scala | 2 +-
.../AbstractLeastShardAllocationStrategy.scala | 87 +--
...egy.scala => ClusterShardAllocationMixin.scala} | 71 +-
.../internal/LeastShardAllocationStrategy.scala | 4 +-
.../cluster/sharding/internal/jfr/Events.scala | 3 +-
.../ClusterShardingMessageSerializer.scala | 12 +-
.../ClusterShardingCoordinatorRoleSpec.scala | 6 +-
.../sharding/MultiNodeClusterShardingConfig.scala | 2 +-
.../sharding/ClusterShardingSettingsSpec.scala | 11 -
.../sharding/ConcurrentStartupShardingSpec.scala | 2 +-
...sistentHashingShardAllocationStrategySpec.scala | 257 +++++++
...eprecatedLeastShardAllocationStrategySpec.scala | 24 +-
.../LeastShardAllocationStrategySpec.scala | 23 +-
.../sharding/PersistentShardingMigrationSpec.scala | 2 +-
.../RememberEntitiesAndStartEntitySpec.scala | 3 +-
.../sharding/ShardRegionDataTypesSpec.scala | 116 +++
.../cluster/sharding/ShardingQueriesSpec.scala | 6 +-
.../passivation/simulator/AccessPattern.scala | 3 +-
.../passivation/simulator/SimulatorSettings.scala | 23 +-
.../client/protobuf/msg/ClusterClientMessages.java | 28 +-
.../protobuf/msg/DistributedPubSubMessages.java | 244 +++---
.../pekko/cluster/client/ClusterClient.scala | 5 +-
.../cluster/pubsub/DistributedPubSubMediator.scala | 5 +-
.../DistributedPubSubMessageSerializer.scala | 12 +-
.../singleton/ClusterSingletonManager.scala | 10 +-
.../cluster/singleton/ClusterSingletonProxy.scala | 8 +-
.../pekko/cluster/client/ClusterClientTest.java | 8 +-
.../typed/internal/protobuf/ClusterMessages.java | 46 +-
.../typed/internal/protobuf/ReliableDelivery.java | 190 +++--
.../receptionist/ClusterReceptionistProtocol.scala | 2 +-
.../ddata/typed/internal/ReplicatorBehavior.scala | 16 +-
.../ddata/typed/javadsl/DistributedData.scala | 10 +-
.../typed/javadsl/ReplicatorMessageAdapter.scala | 8 +-
.../ddata/typed/javadsl/ReplicatorSettings.scala | 2 +-
.../ddata/typed/scaladsl/DistributedData.scala | 6 +-
.../typed/scaladsl/ReplicatorMessageAdapter.scala | 8 +-
.../ddata/typed/scaladsl/ReplicatorSettings.scala | 4 +-
.../org/apache/pekko/cluster/typed/Cluster.scala | 8 +-
.../pekko/cluster/typed/ClusterSingleton.scala | 16 +-
.../typed/internal/AdaptedClusterImpl.scala | 4 +-
.../internal/AdaptedClusterSingletonImpl.scala | 10 +-
.../internal/PekkoClusterTypedSerializer.scala | 8 +-
.../delivery/ReliableDeliverySerializer.scala | 26 +-
.../receptionist/ClusterReceptionist.scala | 20 +-
.../receptionist/ClusterReceptionistSettings.scala | 2 +-
.../typed/internal/receptionist/Registry.scala | 12 +-
.../ddata/typed/javadsl/ReplicatorDocSample.java | 56 +-
.../cluster/typed/PingSerializerExampleTest.java | 4 +-
.../cluster/typed/SingletonCompileOnlyTest.java | 11 +-
.../cluster/typed/BasicClusterExampleSpec.scala | 10 +-
.../cluster/typed/DistributedPubSubExample.scala | 6 +-
.../pekko/cluster/typed/GroupRouterSpec.scala | 8 +-
.../cluster/protobuf/msg/ClusterMessages.java | 501 +++++++-----
.../scala/org/apache/pekko/cluster/Cluster.scala | 15 +-
.../org/apache/pekko/cluster/ClusterDaemon.scala | 4 +-
.../org/apache/pekko/cluster/ClusterEvent.scala | 4 +-
.../pekko/cluster/CrossDcClusterHeartbeat.scala | 8 +-
.../org/apache/pekko/cluster/MembershipState.scala | 3 +-
.../org/apache/pekko/cluster/Reachability.scala | 4 +-
.../org/apache/pekko/cluster/VectorClock.scala | 6 +-
.../protobuf/ClusterMessageSerializer.scala | 12 +-
.../apache/pekko/cluster/sbr/DowningStrategy.scala | 7 +-
.../cluster/sbr/SplitBrainResolverSettings.scala | 8 +-
.../cluster/MultiDcHeartbeatTakingOverSpec.scala | 6 +-
.../pekko/cluster/MultiDcSunnyWeatherSpec.scala | 6 +-
.../org/apache/pekko/cluster/RestartNodeSpec.scala | 2 +-
.../cluster/SurviveNetworkInstabilitySpec.scala | 2 +-
.../pekko/cluster/ClusterJavaCompileTest.java | 3 +-
.../org/apache/pekko/cluster/ClusterTestKit.scala | 11 +-
.../cluster/JoinConfigCompatCheckerSpec.scala | 4 +-
.../pekko/cluster/MixedProtocolClusterSpec.scala | 169 +----
.../pekko/cluster/ReachabilityPerfSpec.scala | 2 +-
.../lease/scaladsl/LeaseProvider.scala | 2 +-
.../apache/pekko/discovery/ServiceDiscovery.scala | 3 +-
.../pekko/discovery/dns/DnsDiscoverySpec.scala | 2 +-
.../ddata/protobuf/msg/ReplicatedDataMessages.java | 450 ++++++-----
.../ddata/protobuf/msg/ReplicatorMessages.java | 557 ++++++++------
.../org/apache/pekko/cluster/ddata/GSet.scala | 8 +-
.../apache/pekko/cluster/ddata/DurableStore.scala | 3 +-
.../scala/org/apache/pekko/cluster/ddata/Key.scala | 4 +-
.../org/apache/pekko/cluster/ddata/LWWMap.scala | 2 +-
.../apache/pekko/cluster/ddata/LWWRegister.scala | 2 +-
.../org/apache/pekko/cluster/ddata/ORMap.scala | 38 +-
.../apache/pekko/cluster/ddata/ORMultiMap.scala | 2 +-
.../org/apache/pekko/cluster/ddata/ORSet.scala | 28 +-
.../apache/pekko/cluster/ddata/PNCounterMap.scala | 2 +-
.../apache/pekko/cluster/ddata/Replicator.scala | 20 +-
.../ddata/protobuf/ReplicatedDataSerializer.scala | 134 ++--
.../protobuf/ReplicatorMessageSerializer.scala | 56 +-
.../ddata/protobuf/SerializationSupport.scala | 12 +-
.../cluster/ddata/JepsenInspiredInsertSpec.scala | 16 +-
.../pekko/cluster/ddata/ReplicatorChaosSpec.scala | 6 +-
.../pekko/cluster/ddata/ReplicatorDeltaSpec.scala | 24 +-
.../pekko/cluster/ddata/ReplicatorGossipSpec.scala | 4 +-
.../cluster/ddata/ReplicatorORSetDeltaSpec.scala | 2 +-
.../pekko/cluster/ddata/ReplicatorSpec.scala | 10 +-
.../cluster/ddata/WildcardSubscribeSpec.scala | 26 +-
.../pekko/cluster/ddata/LocalConcurrencySpec.scala | 2 +-
.../apache/pekko/cluster/ddata/LotsOfDataBot.scala | 2 +-
.../ddata/ReplicatorWildcardSubscriptionSpec.scala | 12 +-
.../ddata/protobuf/msg/TwoPhaseSetMessages.java | 50 +-
.../docs/persistence/proto/FlightAppModels.java | 28 +-
.../docs/persistence/state/MyJavaStateStore.java | 78 ++
.../state/MyJavaStateStoreProvider.java | 45 ++
docs/src/main/paradox/discovery/index.md | 2 +-
.../paradox/durable-state/state-store-plugin.md | 45 ++
docs/src/main/paradox/fsm.md | 5 +-
docs/src/main/paradox/futures.md | 2 +-
docs/src/main/paradox/persistence-query.md | 9 +-
.../src/main/paradox/release-notes/releases-2.0.md | 13 +
docs/src/main/paradox/serialization-jackson.md | 2 +-
.../stream/operators/Sink/preMaterialize.md | 3 +
.../Source-or-Flow/aggregateWithBoundary.md | 4 +
.../stream/operators/Source-or-Flow/alsoTo.md | 11 +-
.../stream/operators/Source-or-Flow/alsoToAll.md | 2 +-
.../stream/operators/Source-or-Flow/batch.md | 5 +
.../operators/Source-or-Flow/batchWeighted.md | 5 +
.../stream/operators/Source-or-Flow/delayWith.md | 2 +
.../stream/operators/Source-or-Flow/doOnFirst.md | 2 +
.../stream/operators/Source-or-Flow/expand.md | 8 +-
.../stream/operators/Source-or-Flow/extrapolate.md | 6 +
.../operators/Source-or-Flow/groupedAdjacentBy.md | 2 +
.../Source-or-Flow/groupedAdjacentByWeighted.md | 4 +-
.../Source-or-Flow/groupedWeightedWithin.md | 1 +
.../operators/Source-or-Flow/preMaterialize.md | 4 +
.../stream/operators/Source-or-Flow/throttle.md | 2 +
.../stream/operators/Source-or-Flow/zipWith.md | 4 +
.../paradox/stream/operators/Source/zipWithN.md | 5 +-
docs/src/main/paradox/stream/operators/index.md | 2 +-
docs/src/main/paradox/stream/stream-context.md | 11 +
docs/src/main/paradox/stream/stream-refs.md | 2 +-
.../typed/cluster-sharded-daemon-process.md | 38 +-
docs/src/main/paradox/typed/cluster-sharding.md | 40 +-
docs/src/main/paradox/typed/cluster-singleton.md | 6 +-
docs/src/main/paradox/typed/distributed-data.md | 2 +-
.../typed/index-persistence-durable-state.md | 1 +
.../src/main/paradox/typed/persistence-snapshot.md | 2 +-
docs/src/main/paradox/typed/persistence-testing.md | 2 +-
docs/src/main/paradox/typed/style-guide.md | 2 +-
.../docs/persistence/state/MyStateStore.scala | 71 ++
docs/src/test/java/jdocs/actor/ActorDocTest.java | 24 +-
.../jdocs/actor/ByteBufferSerializerDocTest.java | 1 +
.../jdocs/actor/DependencyInjectionDocTest.java | 2 +
.../java/jdocs/actor/FaultHandlingDocSample.java | 14 +-
.../test/java/jdocs/actor/ImmutableMessage.java | 20 +-
docs/src/test/java/jdocs/actor/Messages.java | 20 +-
.../test/java/jdocs/actor/SchedulerDocTest.java | 2 -
docs/src/test/java/jdocs/actor/fsm/Buncher.java | 4 +-
.../actor/typed/CoordinatedActorShutdownTest.java | 1 +
.../test/java/jdocs/cluster/FactorialFrontend.java | 9 +-
.../src/test/java/jdocs/cluster/StatsMessages.java | 36 +-
docs/src/test/java/jdocs/cluster/StatsService.java | 14 +-
.../java/jdocs/cluster/TransformationBackend.java | 2 +-
.../java/jdocs/cluster/TransformationMessages.java | 42 +-
.../singleton/ClusterSingletonSupervision.java | 1 +
.../java/jdocs/ddata/DistributedDataDocTest.java | 4 +-
docs/src/test/java/jdocs/ddata/ShoppingCart.java | 81 +-
.../java/jdocs/dispatcher/DispatcherDocTest.java | 3 +
.../src/test/java/jdocs/event/EventBusDocTest.java | 4 +
docs/src/test/java/jdocs/event/LoggingDocTest.java | 3 +
.../jdocs/extension/SettingsExtensionDocTest.java | 1 +
docs/src/test/java/jdocs/future/FutureDocTest.java | 94 +--
docs/src/test/java/jdocs/io/JavaUdpMulticast.java | 3 +-
.../test/java/jdocs/io/UdpConnectedDocTest.java | 1 -
docs/src/test/java/jdocs/io/japi/IODocTest.java | 3 +
docs/src/test/java/jdocs/io/japi/Message.java | 18 +-
.../persistence/LambdaPersistenceDocTest.java | 32 +-
.../LambdaPersistencePluginDocTest.java | 1 +
.../persistence/testkit/PersistenceInitTest.java | 23 +-
.../routing/ConsistentHashingRouterDocTest.java | 2 +
.../java/jdocs/routing/CustomRouterDocTest.java | 1 -
.../src/test/java/jdocs/routing/RouterDocTest.java | 5 +-
.../jdocs/serialization/SerializationDocTest.java | 2 +
.../java/jdocs/sharding/ClusterShardingTest.java | 12 +-
.../test/java/jdocs/stream/BidiFlowDocTest.java | 18 +-
.../test/java/jdocs/stream/CompositionDocTest.java | 4 +-
docs/src/test/java/jdocs/stream/FlowDocTest.java | 19 +-
.../test/java/jdocs/stream/FlowErrorDocTest.java | 22 +-
.../test/java/jdocs/stream/GraphCyclesDocTest.java | 4 +-
.../test/java/jdocs/stream/GraphStageDocTest.java | 20 +-
.../test/java/jdocs/stream/IntegrationDocTest.java | 13 +-
.../test/java/jdocs/stream/KillSwitchDocTest.java | 10 +-
.../test/java/jdocs/stream/QuickStartDocTest.java | 2 +-
.../jdocs/stream/RateTransformationDocTest.java | 11 +-
.../jdocs/stream/StreamBuffersRateDocTest.java | 4 +-
.../jdocs/stream/StreamPartialGraphDSLDocTest.java | 6 +-
.../java/jdocs/stream/StreamTestKitDocTest.java | 11 +-
.../test/java/jdocs/stream/SubstreamDocTest.java | 32 +-
.../stream/TwitterStreamQuickstartDocTest.java | 8 +-
.../stream/javadsl/cookbook/RecipeByteStrings.java | 9 +-
.../stream/javadsl/cookbook/RecipeFlattenList.java | 5 +-
.../javadsl/cookbook/RecipeLoggingElements.java | 8 +-
.../javadsl/cookbook/RecipeManualTrigger.java | 6 +-
.../javadsl/cookbook/RecipeMultiGroupByTest.java | 11 +-
.../stream/javadsl/cookbook/RecipeParseLines.java | 4 +-
.../javadsl/cookbook/RecipeReduceByKeyTest.java | 5 +-
.../jdocs/stream/javadsl/cookbook/RecipeSeq.java | 7 +-
.../stream/javadsl/cookbook/RecipeSplitter.java | 13 +-
.../stream/javadsl/cookbook/RecipeWorkerPool.java | 3 +-
.../stream/operators/BroadcastDocExample.java | 1 +
.../stream/operators/JavaCollectorDocExamples.java | 3 +-
.../stream/operators/MergeSequenceDocExample.java | 1 +
.../stream/operators/PartitionDocExample.java | 1 +
.../jdocs/stream/operators/SinkDocExamples.java | 10 +-
.../jdocs/stream/operators/SourceDocExamples.java | 6 +-
.../java/jdocs/stream/operators/SourceOrFlow.java | 133 ++--
.../jdocs/stream/operators/WithContextTest.java | 32 +-
.../java/jdocs/stream/operators/flow/DiMap.java | 4 +-
.../jdocs/stream/operators/flow/StatefulMap.java | 12 +-
.../stream/operators/flow/StatefulMapConcat.java | 20 +-
.../stream/operators/source/AsSubscriber.java | 2 +-
.../jdocs/stream/operators/source/Combine.java | 6 +-
.../java/jdocs/stream/operators/source/From.java | 5 +-
.../stream/operators/source/FromPublisher.java | 2 +-
.../java/jdocs/stream/operators/source/Zip.java | 21 +-
.../operators/sourceorflow/FlatMapConcat.java | 6 +-
.../operators/sourceorflow/FlatMapMerge.java | 6 +-
.../stream/operators/sourceorflow/Intersperse.java | 4 +-
.../stream/operators/sourceorflow/MapConcat.java | 6 +-
.../stream/operators/sourceorflow/MapError.java | 4 +-
.../operators/sourceorflow/MapWithResource.java | 3 +-
.../stream/operators/sourceorflow/MergeLatest.java | 6 +-
.../stream/operators/sourceorflow/Monitor.java | 4 +-
.../stream/operators/sourceorflow/Sliding.java | 4 +-
.../jdocs/stream/operators/sourceorflow/Split.java | 4 +-
.../tutorial_1/ActorHierarchyExperiments.java | 4 +
.../jdocs/typed/tutorial_5/DeviceGroupQuery.java | 9 +-
docs/src/test/scala/docs/actor/ActorDocSpec.scala | 20 +-
.../scala/docs/actor/FaultHandlingDocSample.scala | 38 +-
.../docs/actor/SharedMutableStateDocSpec.scala | 2 +-
.../scala/docs/ddata/DistributedDataDocSpec.scala | 10 +-
docs/src/test/scala/docs/ddata/ShoppingCart.scala | 4 +-
.../ddata/protobuf/TwoPhaseSetSerializer.scala | 4 +-
.../ddata/protobuf/TwoPhaseSetSerializer2.scala | 2 +-
.../src/test/scala/docs/event/LoggingDocSpec.scala | 2 +-
.../src/test/scala/docs/future/FutureDocSpec.scala | 2 +-
docs/src/test/scala/docs/io/EchoServer.scala | 20 +-
.../src/test/scala/docs/io/ScalaUdpMulticast.scala | 3 +-
.../persistence/PersistencePluginDocSpec.scala | 30 +-
.../persistence/PersistenceSerializerDocSpec.scala | 4 +-
.../docs/persistence/PersistentActorExample.scala | 22 +-
.../state/PersistenceStatePluginDocSpec.scala | 83 ++
.../docs/serialization/SerializationDocSpec.scala | 2 +-
.../test/scala/docs/stream/GraphDSLDocSpec.scala | 4 +-
.../test/scala/docs/stream/QuickStartDocSpec.scala | 8 +-
.../test/scala/docs/stream/SubstreamDocSpec.scala | 5 +-
.../docs/stream/cookbook/RecipeAdhocSource.scala | 2 +-
.../stream/operators/JavaCollectorDocExample.scala | 2 +-
.../docs/stream/operators/WithContextSpec.scala | 25 +
.../docs/stream/operators/source/Restart.scala | 112 +--
.../scala/docs/stream/operators/source/Zip.scala | 2 +-
.../sourceorflow/ExtrapolateAndExpand.scala | 8 +-
.../docs/stream/operators/sourceorflow/Fold.scala | 34 +-
.../stream/operators/sourceorflow/FoldAsync.scala | 24 +-
.../stream/operators/sourceorflow/FoldWhile.scala | 20 +-
.../operators/sourceorflow/Intersperse.scala | 20 +-
.../docs/stream/operators/sourceorflow/Limit.scala | 2 +-
.../operators/sourceorflow/LimitWeighted.scala | 2 +-
.../stream/operators/sourceorflow/MapAsyncs.scala | 78 +-
.../stream/operators/sourceorflow/MapError.scala | 44 +-
.../operators/sourceorflow/MergeLatest.scala | 42 +-
.../stream/operators/sourceorflow/Throttle.scala | 56 +-
.../tutorial_1/ActorHierarchyExperiments.scala | 8 +-
.../testconductor/TestConductorProtocol.java | 150 ++--
.../pekko/remote/testconductor/Conductor.scala | 8 +-
.../apache/pekko/remote/testconductor/Player.scala | 11 +-
.../pekko/remote/testkit/MultiNodeSpec.scala | 6 +-
.../pekko/remote/testkit/PerfFlamesSupport.scala | 4 +-
.../apache/pekko/osgi/ActorSystemActivator.scala | 4 +-
.../pekko/osgi/BundleDelegatingClassLoader.scala | 4 +-
.../org/apache/pekko/osgi/PojoSRTestSupport.scala | 2 +-
.../query/internal/protobuf/QueryMessages.java | 63 +-
.../persistence/query/javadsl/package-info.java | 4 +-
.../src/main/resources/reference.conf | 53 ++
.../persistence/query/ReadJournalProvider.scala | 4 +-
.../query/internal/QuerySerializer.scala | 4 +-
.../journal/leveldb/AllPersistenceIdsStage.scala | 2 +-
.../leveldb/EventsByPersistenceIdStage.scala | 2 +-
.../query/journal/leveldb/EventsByTagStage.scala | 2 +-
.../leveldb/javadsl/LeveldbReadJournal.scala | 2 +-
.../leveldb/scaladsl/LeveldbReadJournal.scala | 8 +-
.../persistence/query/typed/EventEnvelope.scala | 22 +-
.../EventsBySliceFirehoseReadJournalProvider.scala | 33 +
.../typed/internal/EventsBySliceFirehose.scala | 839 +++++++++++++++++++++
...ByPersistenceIdStartingFromSnapshotQuery.scala} | 16 +-
...tEventsBySliceStartingFromSnapshotsQuery.scala} | 20 +-
...ByPersistenceIdStartingFromSnapshotQuery.scala} | 16 +-
.../typed/javadsl/EventsBySliceFirehoseQuery.scala | 95 +++
.../query/typed/javadsl/EventsBySliceQuery.scala | 2 +
... EventsBySliceStartingFromSnapshotsQuery.scala} | 22 +-
...ByPersistenceIdStartingFromSnapshotQuery.scala} | 16 +-
...tEventsBySliceStartingFromSnapshotsQuery.scala} | 20 +-
...ByPersistenceIdStartingFromSnapshotQuery.scala} | 16 +-
.../scaladsl/EventsBySliceFirehoseQuery.scala | 119 +++
.../query/typed/scaladsl/EventsBySliceQuery.scala | 2 +
... EventsBySliceStartingFromSnapshotsQuery.scala} | 22 +-
.../apache/pekko/persistence/query/TestClock.scala | 59 ++
.../query/typed/EventEnvelopeSpec.scala | 53 ++
.../typed/internal/EventsBySliceFirehoseSpec.scala | 440 +++++++++++
.../persistence/serialization/SerializerSpec.scala | 4 +-
.../apache/pekko/persistence/CapabilityFlags.scala | 10 +
.../japi/journal/JavaJournalPerfSpec.scala | 2 +
.../japi/state/JavaDurableStateStoreSpec.scala | 2 +-
.../pekko/persistence/journal/JournalSpec.scala | 54 ++
.../persistence/state/DurableStateStoreSpec.scala | 2 +-
.../journal/inmem/InmemJournalSpec.scala | 1 +
.../journal/leveldb/LeveldbJournalJavaSpec.scala | 3 +
.../journal/leveldb/LeveldbJournalNativeSpec.scala | 2 +
...bJournalNoAtomicPersistMultipleEventsSpec.scala | 2 +
.../persistence/testkit/javadsl/package-info.java | 4 +-
.../persistence/testkit/SnapshotStorage.scala | 3 -
.../internal/EventSourcedBehaviorTestKitImpl.scala | 6 +-
.../testkit/internal/PersistenceProbeImpl.scala | 10 +-
.../SnapshotStorageEmulatorExtension.scala | 2 +-
.../javadsl/EventSourcedBehaviorTestKit.scala | 22 +-
.../testkit/javadsl/PersistenceProbeBehavior.scala | 3 +-
.../testkit/javadsl/PersistenceTestKit.scala | 2 +-
.../testkit/javadsl/SnapshotTestKit.scala | 2 +-
.../scaladsl/PersistenceTestKitReadJournal.scala | 1 +
.../scaladsl/EventSourcedBehaviorTestKit.scala | 4 +-
.../testkit/scaladsl/PersistenceTestKit.scala | 2 +-
.../persistence/testkit/scaladsl/TestOps.scala | 20 +-
.../PersistenceTestKitJournalCompatSpec.scala | 1 +
.../persistence/typed/MyReplicatedBehavior.java | 3 +-
.../typed/ReplicatedAuctionExampleTest.java | 6 +-
.../persistence/typed/ReplicatedMovieExample.java | 33 +-
.../typed/ReplicatedShoppingCartExample.java | 46 +-
.../typed/ReplicatedEventSourcingTest.java | 15 +-
.../javadsl/EventSourcedBehaviorJavaDslTest.java | 32 +-
.../javadsl/RuntimeDurableStateStoreTest.java | 160 ++++
.../ReplicatedEventSourcingCompileOnlySpec.scala | 4 +-
.../scaladsl/EventSourcedBehaviorReplySpec.scala | 8 +-
...urcedBehaviorRetentionOnlyOneSnapshotSpec.scala | 7 +-
.../typed/scaladsl/EventSourcedBehaviorSpec.scala | 16 +-
.../scaladsl/EventSourcedBehaviorStashSpec.scala | 12 +-
.../scaladsl/EventSourcedBehaviorWatchSpec.scala | 4 +-
.../scaladsl/DurableStateBehaviorReplySpec.scala | 6 +-
.../DurableStateBehaviorStashOverflowSpec.scala | 4 +-
.../scaladsl/RuntimeDurableStateStoreSpec.scala | 100 +++
.../persistence/typed/javadsl/package-info.java | 4 +-
.../serialization/ReplicatedEventSourcing.java | 262 ++++---
.../durablestate-runtime-config.excludes | 4 +-
.../pekko/persistence/typed/crdt/ORSet.scala | 2 +-
.../typed/delivery/EventSourcedProducerQueue.scala | 10 +-
.../typed/internal/EventSourcedBehaviorImpl.scala | 8 +-
.../typed/internal/EventSourcedSettings.scala | 7 +-
.../typed/internal/ExternalInteractions.scala | 8 +-
.../typed/internal/ReplayingEvents.scala | 8 +-
.../typed/internal/ReplayingSnapshot.scala | 4 +-
.../typed/internal/RequestingRecoveryPermit.scala | 2 +-
.../pekko/persistence/typed/internal/Running.scala | 14 +-
.../persistence/typed/javadsl/CommandHandler.scala | 6 +-
.../typed/javadsl/CommandHandlerWithReply.scala | 6 +-
.../pekko/persistence/typed/javadsl/Effect.scala | 2 +-
.../persistence/typed/javadsl/EventHandler.scala | 4 +-
.../typed/javadsl/EventSourcedBehavior.scala | 11 +-
.../typed/javadsl/PersistentFSMMigration.scala | 2 +-
.../typed/javadsl/RetentionCriteria.scala | 2 +-
.../typed/scaladsl/EventSourcedBehavior.scala | 16 +-
.../typed/scaladsl/PersistentFSMMigration.scala | 2 +-
.../typed/scaladsl/RetentionCriteria.scala | 2 +-
.../ReplicatedEventSourcingSerializer.scala | 28 +-
.../typed/state/internal/BehaviorSetup.scala | 6 +-
.../state/internal/DurableStateBehaviorImpl.scala | 18 +-
.../state/internal/DurableStateSettings.scala | 39 +-
.../internal/DurableStateStoreInteractions.scala | 2 +-
.../typed/state/internal/Recovering.scala | 6 +-
.../state/internal/RequestingRecoveryPermit.scala | 2 +-
.../persistence/typed/state/internal/Running.scala | 12 +-
.../typed/state/javadsl/CommandHandler.scala | 6 +-
.../state/javadsl/CommandHandlerWithReply.scala | 6 +-
.../typed/state/javadsl/DurableStateBehavior.scala | 23 +-
.../state/scaladsl/DurableStateBehavior.scala | 22 +-
.../typed/BasicPersistentBehaviorTest.java | 51 +-
.../pekko/persistence/typed/BlogPostEntity.java | 116 +--
.../typed/BlogPostEntityDurableState.java | 82 +-
.../typed/DurableStatePersistentBehaviorTest.java | 4 +
.../pekko/persistence/typed/MovieWatchList.java | 52 +-
.../pekko/persistence/typed/NullBlogState.java | 107 +--
.../pekko/persistence/typed/OptionalBlogState.java | 108 +--
...rsistentFsmToTypedMigrationCompileOnlyTest.java | 48 +-
.../persistence/typed/WebStoreCustomerFSM.java | 75 +-
.../persistence/typed/auction/AuctionCommand.java | 28 +-
.../persistence/typed/auction/AuctionEntity.java | 16 +-
.../persistence/typed/auction/AuctionState.java | 4 +-
.../pekko/persistence/typed/auction/Bid.java | 37 +-
.../javadsl/PersistentActorCompileOnlyTest.java | 69 +-
...DurableStatePersistentBehaviorCompileOnly.scala | 24 +-
.../typed/StashingWhenSnapshottingSpec.scala | 2 +-
.../typed/internal/StashStateSpec.scala | 2 +-
.../scaladsl/PersistentActorCompileOnlyTest.scala | 2 +-
.../fsm/japi/pf/FSMStateFunctionBuilder.java | 6 +-
.../persistence/fsm/japi/pf/FSMStopBuilder.java | 6 +-
.../journal/japi/AsyncRecoveryPlugin.java | 4 +-
.../persistence/journal/japi/package-info.java | 4 +-
.../persistence/serialization/MessageFormats.java | 160 ++--
.../persistence/snapshot/japi/package-info.java | 2 +-
.../persistence/state/javadsl/package-info.java | 2 +-
.../fsm-transition-listeners.excludes | 4 +-
persistence/src/main/resources/reference.conf | 3 +
.../org/apache/pekko/persistence/TraitOrder.scala | 2 +-
.../org/apache/pekko/persistence/TraitOrder.scala | 2 +-
.../pekko/persistence/AtLeastOnceDelivery.scala | 6 +-
.../apache/pekko/persistence/Eventsourced.scala | 4 +-
.../apache/pekko/persistence/FilteredPayload.scala | 26 +
.../apache/pekko/persistence/JournalProtocol.scala | 4 +-
.../org/apache/pekko/persistence/Persistence.scala | 6 +-
.../pekko/persistence/PersistencePlugin.scala | 10 +-
.../pekko/persistence/fsm/PersistentFSM.scala | 4 +-
.../pekko/persistence/fsm/PersistentFSMBase.scala | 51 +-
.../pekko/persistence/journal/EventAdapters.scala | 21 +-
.../journal/PersistencePluginProxy.scala | 2 +-
.../journal/leveldb/LeveldbCompaction.scala | 2 +-
.../serialization/FilteredPayloadSerializer.scala | 45 ++
.../serialization/MessageSerializer.scala | 2 +-
.../serialization/SnapshotSerializer.scala | 5 +-
.../pekko/persistence/serialization/package.scala | 19 +-
.../snapshot/local/LocalSnapshotStore.scala | 6 +-
.../state/DurableStateStoreProvider.scala | 6 +-
.../state/DurableStateStoreRegistry.scala | 47 +-
.../persistence/fsm/AbstractPersistentFSMTest.java | 84 +--
.../persistence/AtLeastOnceDeliverySpec.scala | 5 +-
.../pekko/persistence/PersistentActorSpec.scala | 2 +-
.../SnapshotRecoveryWithEmptyJournalSpec.scala | 2 +-
.../persistence/SnapshotSerializationSpec.scala | 2 +-
.../apache/pekko/persistence/SnapshotSpec.scala | 2 +-
.../pekko/persistence/fsm/PersistentFSMSpec.scala | 30 +-
.../AsyncWriteJournalResponseOrderSpec.scala | 4 +-
.../persistence/journal/SteppingInmemJournal.scala | 2 +-
.../FilteredPayloadSerializerSpec.scala | 35 +
.../SnapshotSerializerMigrationAkkaSpec.scala | 8 +-
.../SnapshotSerializerNoMigrationSpec.scala | 4 +-
.../serialization/SnapshotSerializerSpec.scala | 15 +-
.../exception/DurableStateExceptionsSpec.scala | 2 +-
project/AddLogTimestamps.scala | 2 +-
project/CopyrightHeaderForBoilerplate.scala | 2 +-
project/CopyrightHeaderForBuild.scala | 2 +-
...ForJdk9.scala => CopyrightHeaderForJdk21.scala} | 8 +-
project/CopyrightHeaderForProtobuf.scala | 2 +-
project/Dependencies.scala | 29 +-
project/JavaFormatter.scala | 2 +-
project/{Jdk9.scala => Jdk21.scala} | 38 +-
project/JdkOptions.scala | 5 +-
project/MultiNode.scala | 2 +-
project/PekkoBuild.scala | 25 +-
project/PekkoDevelocityPlugin.scala | 10 +-
project/ProjectFileIgnoreSupport.scala | 4 +-
project/Protobuf.scala | 2 +-
project/SbtMultiJvmPlugin.scala | 4 +-
project/ScalaFixExtraRulesPlugin.scala | 2 +-
...k9Plugin.scala => ScalaFixForJdk21Plugin.scala} | 12 +-
project/ScalafixForMultiNodePlugin.scala | 4 +-
project/ScalafixIgnoreFilePlugin.scala | 4 +-
project/StreamOperatorsIndexGenerator.scala | 3 +-
project/TestExtras.scala | 14 +-
project/ValidatePullRequest.scala | 2 +-
project/VersionGenerator.scala | 2 +-
project/build.properties | 2 +-
project/plugins.sbt | 19 +-
.../aeron/AeronStreamMaxThroughputSpec.scala | 4 +-
.../pekko/remote/artery/protobuf/TestMessages.java | 50 +-
.../apache/pekko/remote/ArteryControlFormats.java | 236 +++---
.../org/apache/pekko/remote/ContainerFormats.java | 272 ++++---
.../apache/pekko/remote/SystemMessageFormats.java | 120 +--
.../java/org/apache/pekko/remote/WireFormats.java | 444 ++++++-----
.../pekko/remote/artery/aeron/AeronErrorLog.java | 12 +-
.../remote/artery/compress/CountMinSketch.java | 4 +-
.../future-lazy-vals.excludes | 6 +-
remote/src/main/resources/reference.conf | 13 +-
.../org/apache/pekko/remote/AckedDelivery.scala | 26 +-
.../scala/org/apache/pekko/remote/Endpoint.scala | 125 ++-
.../org/apache/pekko/remote/RemoteDaemon.scala | 2 +-
.../pekko/remote/RemoteMetricsExtension.scala | 2 +-
.../org/apache/pekko/remote/RemoteSettings.scala | 6 +-
.../scala/org/apache/pekko/remote/Remoting.scala | 89 ++-
.../pekko/remote/RemotingLifecycleEvent.scala | 4 +-
.../pekko/remote/artery/ArterySettings.scala | 6 +
.../pekko/remote/artery/ArteryTransport.scala | 119 ++-
.../apache/pekko/remote/artery/Association.scala | 42 +-
.../pekko/remote/artery/EnvelopeBufferPool.scala | 2 +-
.../org/apache/pekko/remote/artery/Handshake.scala | 2 +-
.../pekko/remote/artery/ImmutableLongMap.scala | 5 +-
.../pekko/remote/artery/LruBoundedCache.scala | 12 +-
.../pekko/remote/artery/RemoteInstrument.scala | 2 +-
.../remote/artery/RemotingFlightRecorder.scala | 2 +-
.../remote/artery/SystemMessageDelivery.scala | 14 +-
.../artery/aeron/ArteryAeronUdpTransport.scala | 14 +-
.../pekko/remote/artery/aeron/TaskRunner.scala | 8 +-
.../remote/artery/compress/CompressionTable.scala | 2 +-
.../artery/compress/DecompressionTable.scala | 4 +-
.../artery/compress/InboundCompressions.scala | 84 ++-
.../remote/artery/compress/TopHeavyHitters.scala | 110 +--
.../artery/tcp/ConfigSSLEngineProvider.scala | 2 +-
.../pekko/remote/artery/tcp/TcpFraming.scala | 2 +-
.../artery/tcp/ssl/PemManagersProvider.scala | 3 +-
.../pekko/remote/artery/tcp/ssl/X509Readers.scala | 2 +-
.../serialization/DaemonMsgCreateSerializer.scala | 2 +-
.../serialization/MessageContainerSerializer.scala | 2 +-
.../serialization/MiscMessageSerializer.scala | 14 +-
.../remote/serialization/ProtobufSerializer.scala | 16 +-
.../serialization/SystemMessageSerializer.scala | 2 +-
.../transport/ThrottlerTransportAdapter.scala | 63 +-
.../remote/transport/netty/NettySSLSupport.scala | 34 +-
.../remote/transport/netty/NettyTransport.scala | 15 +-
.../remote/transport/netty/SSLEngineProvider.scala | 46 +-
.../org/apache/pekko/remote/ProtobufProtocol.java | 28 +-
.../remote/protobuf/v3/ProtobufProtocolV3.java | 28 +-
.../transport/ThrottlerTransportAdapterTest.java | 1 +
.../apache/pekko/remote/AckedDeliverySpec.scala | 11 +
.../pekko/remote/ConfigSSLEngineProviderSpec.scala | 47 ++
.../org/apache/pekko/remote/DaemonicSpec.scala | 7 +-
.../remote/ReliableDeliverySupervisorSpec.scala | 225 ++++++
.../apache/pekko/remote/Ticket1978ConfigSpec.scala | 10 +
.../remote/TransientSerializationErrorSpec.scala | 7 +-
.../ActorSelectionQueueDistributionSpec.scala | 288 +++++++
.../remote/artery/BindCanonicalAddressSpec.scala | 2 +-
.../pekko/remote/artery/MetadataCarryingSpec.scala | 2 +-
.../RemoteInstrumentsSerializationSpec.scala | 8 +-
.../remote/artery/RemoteSendConsistencySpec.scala | 4 +-
.../remote/artery/SystemMessageDeliverySpec.scala | 44 ++
.../compress/CompressionIntegrationSpec.scala | 48 ++
.../remote/artery/compress/HeavyHittersSpec.scala | 34 +
.../tcp/ArteryConfigSSLEngineProviderSpec.scala | 43 ++
.../ssl/RotatingKeysSSLEngineProviderSpec.scala | 3 +-
.../apache/pekko/remote/classic/RemotingSpec.scala | 124 ++-
...reateSerializerAllowJavaSerializationSpec.scala | 4 +-
.../remote/transport/ThrottlerHandleSpec.scala | 130 ++++
.../transport/netty/NettySSLSupportSpec.scala | 25 +-
.../serialization/jackson/JacksonModule.scala | 14 +-
.../jackson/JacksonObjectMapperProvider.scala | 14 +-
.../serialization/jackson/JacksonSerializer.scala | 35 +-
.../serialization/jackson/StreamRefModule.scala | 24 +-
.../jackson/TypedActorRefModule.scala | 12 +-
.../jackson/JacksonSerializerSpec.scala | 14 +-
.../serialization/jackson3/JacksonModule.scala | 16 +-
.../jackson3/JacksonObjectMapperProvider.scala | 18 +-
.../serialization/jackson3/JacksonSerializer.scala | 38 +-
.../serialization/jackson3/StreamRefModule.scala | 24 +-
.../jackson3/TypedActorRefModule.scala | 12 +-
.../jackson3/JacksonSerializerSpec.scala | 4 +-
.../org/apache/pekko/event/slf4j/Slf4jLogger.scala | 18 +-
.../pekko/stream/testkit/javadsl/package-info.java | 4 +-
.../pekko/stream/testkit/StreamTestKit.scala | 14 +-
.../impl/fusing/GraphInterpreterSpecKit.scala | 34 +-
.../apache/pekko/stream/testkit/ChainSetup.scala | 6 +-
.../apache/pekko/stream/testkit/ScriptedTest.scala | 2 +-
.../apache/pekko/stream/testkit/StreamSpec.scala | 2 +-
.../pekko/stream/testkit/TwoStreamsSetup.scala | 2 +-
.../tck/FlatMapConcatDoubleSubscriberTest.scala | 2 +-
.../pekko/stream/tck/InputStreamSourceTest.scala | 21 +-
.../tck/PekkoIdentityProcessorVerification.scala | 2 +-
.../pekko/stream/io/InputStreamSinkTest.java | 3 +-
.../pekko/stream/io/SinkAsJavaSourceTest.java | 6 +-
.../pekko/stream/javadsl/AttributesTest.java | 7 +-
.../org/apache/pekko/stream/javadsl/FlowTest.java | 229 +++---
.../pekko/stream/javadsl/FlowThrottleTest.java | 7 +-
.../pekko/stream/javadsl/FlowWithContextTest.java | 134 ++++
.../apache/pekko/stream/javadsl/GraphDslTest.java | 56 +-
.../stream/javadsl/LazyAndFutureFlowTest.java | 7 +-
.../stream/javadsl/LazyAndFutureSourcesTest.java | 13 +-
.../org/apache/pekko/stream/javadsl/SinkTest.java | 28 +-
.../apache/pekko/stream/javadsl/SourceTest.java | 271 ++++---
.../stream/javadsl/SourceWithContextTest.java | 158 ++++
.../javadsl/SourceWithContextThrottleTest.java | 4 +-
.../org/apache/pekko/stream/javadsl/TcpTest.java | 10 +-
.../org/apache/pekko/stream/stage/StageTest.java | 6 +-
.../apache/pekko/stream/DslConsistencySpec.scala | 41 +-
.../pekko/stream/DslFactoriesConsistencySpec.scala | 46 +-
.../scala/org/apache/pekko/stream/FusingSpec.scala | 202 +++++
.../stream/impl/FanoutPublisherBehaviorSpec.scala | 6 +-
.../apache/pekko/stream/impl/FixedBufferSpec.scala | 16 +-
.../apache/pekko/stream/impl/TimeoutsSpec.scala | 93 +++
.../pekko/stream/impl/TraversalBuilderSpec.scala | 6 +-
.../impl/fusing/ActorGraphInterpreterSpec.scala | 63 +-
.../org/apache/pekko/stream/io/FileSinkSpec.scala | 2 +-
.../apache/pekko/stream/io/FileSourceSpec.scala | 46 +-
.../scala/org/apache/pekko/stream/io/TcpSpec.scala | 25 +
.../stream/io/TlsGraphStageEdgeCasesSpec.scala | 7 +-
.../io/compression/DeflateAutoFlushSpec.scala | 4 +-
.../stream/io/compression/GzipAutoFlushSpec.scala | 4 +-
.../pekko/stream/scaladsl/ActorRefSourceSpec.scala | 89 ++-
.../scaladsl/AggregateWithBoundarySpec.scala | 114 +++
.../pekko/stream/scaladsl/AttributesSpec.scala | 5 +-
.../scaladsl/CoupledTerminationFlowSpec.scala | 4 +-
.../pekko/stream/scaladsl/FlowAlsoToSpec.scala | 288 +++++++
.../pekko/stream/scaladsl/FlowBatchSpec.scala | 275 +++++++
.../pekko/stream/scaladsl/FlowCompileSpec.scala | 32 +-
.../pekko/stream/scaladsl/FlowConcatAllSpec.scala | 2 +-
.../pekko/stream/scaladsl/FlowConcatSpec.scala | 14 +-
.../pekko/stream/scaladsl/FlowDelaySpec.scala | 128 ++++
.../pekko/stream/scaladsl/FlowDoOnFirstSpec.scala | 83 +-
.../pekko/stream/scaladsl/FlowExpandSpec.scala | 197 +++++
.../stream/scaladsl/FlowExtrapolateSpec.scala | 108 +++
.../pekko/stream/scaladsl/FlowGroupBySpec.scala | 14 +-
.../FlowGroupedAdjacentByWeightedSpec.scala | 87 +++
.../stream/scaladsl/FlowGroupedWithinSpec.scala | 127 ++++
.../pekko/stream/scaladsl/FlowIdleInjectSpec.scala | 42 ++
.../scaladsl/FlowMapAsyncPartitionedSpec.scala | 55 ++
.../pekko/stream/scaladsl/FlowMapAsyncSpec.scala | 47 ++
.../scaladsl/FlowMapAsyncUnorderedSpec.scala | 53 ++
.../pekko/stream/scaladsl/FlowOnCompleteSpec.scala | 2 +-
.../stream/scaladsl/FlowPrefixAndTailSpec.scala | 12 +-
.../stream/scaladsl/FlowRecoverWithSpec.scala | 10 +-
.../pekko/stream/scaladsl/FlowScanAsyncSpec.scala | 2 +-
.../pekko/stream/scaladsl/FlowSlidingSpec.scala | 11 +-
.../apache/pekko/stream/scaladsl/FlowSpec.scala | 61 +-
.../pekko/stream/scaladsl/FlowSplitAfterSpec.scala | 6 +-
.../pekko/stream/scaladsl/FlowSplitWhenSpec.scala | 6 +-
.../pekko/stream/scaladsl/FlowTakeSpec.scala | 2 +-
.../pekko/stream/scaladsl/FlowThrottleSpec.scala | 63 +-
.../stream/scaladsl/FlowWithContextSpec.scala | 92 +++
.../stream/scaladsl/FlowZipWithIndexSpec.scala | 10 +-
.../pekko/stream/scaladsl/FlowZipWithSpec.scala | 41 +-
.../stream/scaladsl/GraphBackedFlowSpec.scala | 10 +-
.../pekko/stream/scaladsl/GraphBroadcastSpec.scala | 45 ++
.../pekko/stream/scaladsl/GraphConcatSpec.scala | 2 +-
.../stream/scaladsl/GraphMergeLatestSpec.scala | 2 +-
.../stream/scaladsl/GraphMergePreferredSpec.scala | 2 +-
.../scaladsl/GraphMergePrioritizedSpec.scala | 2 +-
.../stream/scaladsl/GraphMergeSequenceSpec.scala | 2 +-
.../stream/scaladsl/GraphMergeSortedSpec.scala | 2 +-
.../pekko/stream/scaladsl/GraphMergeSpec.scala | 2 +-
.../stream/scaladsl/GraphOpsIntegrationSpec.scala | 6 +-
.../pekko/stream/scaladsl/GraphUnzipWithSpec.scala | 8 +-
.../stream/scaladsl/GraphZipLatestWithSpec.scala | 2 +-
.../pekko/stream/scaladsl/GraphZipNSpec.scala | 2 +-
.../pekko/stream/scaladsl/GraphZipSpec.scala | 2 +-
.../pekko/stream/scaladsl/GraphZipWithNSpec.scala | 46 +-
.../pekko/stream/scaladsl/GraphZipWithSpec.scala | 7 +-
.../org/apache/pekko/stream/scaladsl/HubSpec.scala | 45 ++
.../apache/pekko/stream/scaladsl/RestartSpec.scala | 6 +-
.../apache/pekko/stream/scaladsl/SinkSpec.scala | 3 +-
.../stream/scaladsl/SourceWithContextSpec.scala | 93 +++
.../pekko/stream/scaladsl/StageActorRefSpec.scala | 120 ++-
.../pekko/stream/scaladsl/StreamRefsSpec.scala | 64 +-
.../SubstreamSubscriptionTimeoutSpec.scala | 6 +-
.../pekko/stream/scaladsl/TakeLastSinkSpec.scala | 10 +-
.../scaladsl/UnfoldResourceAsyncSourceSpec.scala | 34 +
.../stream/scaladsl/UnfoldResourceSourceSpec.scala | 71 +-
.../pekko/stream/MapAsyncPartitionedSpec.scala | 2 +-
.../pekko/stream/typed/javadsl/package-info.java | 4 +-
.../pekko/stream/typed/javadsl/ActorFlow.scala | 12 +-
.../pekko/stream/typed/javadsl/ActorSource.scala | 2 +-
.../pekko/stream/typed/scaladsl/ActorFlow.scala | 16 +-
.../pekko/stream/typed/scaladsl/ActorSource.scala | 2 +-
.../typed/ActorSourceWithBackpressureExample.java | 18 +-
.../stream/typed/ActorSourceSinkExample.scala | 8 +-
.../scaladsl/ZipLatestWithApply.scala.template | 2 +-
.../stream/scaladsl/ZipWithApply.scala.template | 21 +-
.../org/apache/pekko/stream/StreamRefMessages.java | 174 +++--
.../pekko/stream/javadsl/FramingTruncation.java | 2 +-
.../pekko/stream/javadsl/JavaFlowSupport.java | 9 +-
.../apache/pekko/stream/javadsl/package-info.java | 4 +-
.../pr-2916-boundary-event-allocation.excludes | 6 +
.../remove-dead-actorprocessor.excludes} | 14 +-
...move-deprecated-materializer-settings.excludes} | 24 +-
.../remove-deprecated-methods.excludes | 4 -
.../remove-fanout-publisher-sink.excludes} | 4 +-
.../remove-substream-cancel-strategy.excludes | 32 +-
stream/src/main/resources/reference.conf | 7 +
.../apache/pekko/stream/ActorMaterializer.scala | 112 +--
.../scala/org/apache/pekko/stream/Attributes.scala | 4 +-
.../scala/org/apache/pekko/stream/FanInShape.scala | 12 +-
.../org/apache/pekko/stream/FanOutShape.scala | 12 +-
.../main/scala/org/apache/pekko/stream/Graph.scala | 2 +-
.../scala/org/apache/pekko/stream/KillSwitch.scala | 8 +-
.../apache/pekko/stream/MapAsyncPartitioned.scala | 41 +-
.../main/scala/org/apache/pekko/stream/Shape.scala | 50 +-
.../apache/pekko/stream/StreamRefSettings.scala | 70 --
.../pekko/stream/SubstreamCancelStrategy.scala | 57 --
.../apache/pekko/stream/TypePreservingFanIn.scala | 4 +-
.../apache/pekko/stream/TypePreservingFanOut.scala | 2 +-
.../apache/pekko/stream/impl/ActorProcessor.scala | 105 +--
.../apache/pekko/stream/impl/ActorPublisher.scala | 12 +-
.../stream/impl/ActorRefBackpressureSource.scala | 2 +-
.../pekko/stream/impl/ActorRefSinkStage.scala | 2 +-
.../apache/pekko/stream/impl/ActorRefSource.scala | 174 +++--
.../pekko/stream/impl/CompletedPublishers.scala | 6 +-
.../stream/impl/ExposedPublisherReceive.scala | 43 --
.../scala/org/apache/pekko/stream/impl/FanIn.scala | 26 +-
.../stream/impl/FanoutPublisherBridgeStage.scala | 16 +-
.../stream/impl/JavaFlowAndRsConverters.scala | 8 +-
.../pekko/stream/impl/JavaStreamConcat.scala | 4 +-
.../pekko/stream/impl/JsonObjectParser.scala | 14 +-
.../impl/PhasedFusingActorMaterializer.scala | 90 ++-
.../org/apache/pekko/stream/impl/QueueSource.scala | 2 +-
.../impl/ResizableMultiReaderRingBuffer.scala | 8 +-
.../pekko/stream/impl/SinkholeSubscriber.scala | 2 +-
.../scala/org/apache/pekko/stream/impl/Sinks.scala | 28 +-
.../apache/pekko/stream/impl/StreamLayout.scala | 32 +-
.../stream/impl/StreamSubscriptionTimeout.scala | 12 +-
.../pekko/stream/impl/SubscriberManagement.scala | 18 +-
.../org/apache/pekko/stream/impl/Throttle.scala | 19 +-
.../org/apache/pekko/stream/impl/Timers.scala | 25 +-
.../pekko/stream/impl/TraversalBuilder.scala | 34 +-
.../org/apache/pekko/stream/impl/Unfold.scala | 10 +-
.../pekko/stream/impl/UnfoldResourceSource.scala | 5 +-
.../stream/impl/UnfoldResourceSourceAsync.scala | 21 +-
.../stream/impl/fusing/ActorGraphInterpreter.scala | 256 +++++--
.../stream/impl/fusing/AggregateWithBoundary.scala | 122 ++-
.../pekko/stream/impl/fusing/DoOnFirst.scala | 42 +-
.../pekko/stream/impl/fusing/FlattenConcat.scala | 4 +-
.../stream/impl/fusing/GraphInterpreter.scala | 108 +--
.../pekko/stream/impl/fusing/GraphStages.scala | 4 +-
.../impl/fusing/GroupedAdjacentByWeighted.scala | 26 +-
.../pekko/stream/impl/fusing/InflightSources.scala | 4 +-
.../pekko/stream/impl/fusing/IteratorSource.scala | 2 +-
.../org/apache/pekko/stream/impl/fusing/Ops.scala | 348 +++++++--
.../pekko/stream/impl/fusing/RangeSource.scala | 10 +-
.../pekko/stream/impl/fusing/StreamOfStreams.scala | 30 +-
.../pekko/stream/impl/io/ByteStringParser.scala | 2 +-
.../stream/impl/io/InputStreamSinkStage.scala | 10 +-
.../pekko/stream/impl/io/InputStreamSource.scala | 2 +-
.../stream/impl/io/OutputStreamGraphStage.scala | 2 +-
.../stream/impl/io/OutputStreamSourceStage.scala | 2 +-
.../org/apache/pekko/stream/impl/io/TLSActor.scala | 20 +-
.../apache/pekko/stream/impl/io/TcpStages.scala | 26 +-
.../pekko/stream/impl/io/TlsEngineHelpers.scala | 3 +-
.../pekko/stream/impl/io/TlsGraphStage.scala | 4 +-
.../io/compression/DeflateDecompressorBase.scala | 2 +-
.../org/apache/pekko/stream/impl/package.scala | 16 +-
.../pekko/stream/impl/streamref/SinkRefImpl.scala | 70 +-
.../stream/impl/streamref/SourceRefImpl.scala | 44 +-
.../impl/streamref/StreamRefDefaultSettings.scala | 20 +-
.../impl/streamref/StreamRefSettingsImpl.scala | 41 -
.../stream/impl/streamref/StreamRefsMaster.scala | 4 +-
.../org/apache/pekko/stream/javadsl/BidiFlow.scala | 2 +-
.../pekko/stream/javadsl/DelayStrategy.scala | 2 +-
.../org/apache/pekko/stream/javadsl/Flow.scala | 269 ++++---
.../pekko/stream/javadsl/FlowWithContext.scala | 159 +++-
.../org/apache/pekko/stream/javadsl/Graph.scala | 57 +-
.../apache/pekko/stream/javadsl/RestartFlow.scala | 4 +-
.../apache/pekko/stream/javadsl/RestartSink.scala | 2 +-
.../pekko/stream/javadsl/RestartSource.scala | 4 +-
.../org/apache/pekko/stream/javadsl/Sink.scala | 19 +-
.../org/apache/pekko/stream/javadsl/Source.scala | 272 ++++---
.../pekko/stream/javadsl/SourceWithContext.scala | 155 +++-
.../pekko/stream/javadsl/StreamConverters.scala | 4 +-
.../org/apache/pekko/stream/javadsl/SubFlow.scala | 185 +++--
.../apache/pekko/stream/javadsl/SubSource.scala | 187 +++--
.../org/apache/pekko/stream/javadsl/Tcp.scala | 2 +-
.../org/apache/pekko/stream/javadsl/package.scala | 2 +-
.../apache/pekko/stream/scaladsl/BidiFlow.scala | 4 +-
.../apache/pekko/stream/scaladsl/Compression.scala | 8 +-
.../pekko/stream/scaladsl/DelayStrategy.scala | 2 +-
.../org/apache/pekko/stream/scaladsl/FileIO.scala | 8 +-
.../org/apache/pekko/stream/scaladsl/Flow.scala | 376 +++++----
.../pekko/stream/scaladsl/FlowWithContext.scala | 12 +-
.../pekko/stream/scaladsl/FlowWithContextOps.scala | 269 +++++--
.../org/apache/pekko/stream/scaladsl/Framing.scala | 2 +-
.../org/apache/pekko/stream/scaladsl/Graph.scala | 200 +++--
.../org/apache/pekko/stream/scaladsl/Hub.scala | 116 +--
.../pekko/stream/scaladsl/JavaFlowSupport.scala | 2 +-
.../org/apache/pekko/stream/scaladsl/Queue.scala | 14 +-
.../apache/pekko/stream/scaladsl/RestartFlow.scala | 12 +-
.../apache/pekko/stream/scaladsl/RestartSink.scala | 4 +-
.../pekko/stream/scaladsl/RestartSource.scala | 6 +-
.../apache/pekko/stream/scaladsl/RetryFlow.scala | 2 +-
.../org/apache/pekko/stream/scaladsl/Sink.scala | 30 +-
.../org/apache/pekko/stream/scaladsl/Source.scala | 27 +-
.../pekko/stream/scaladsl/SourceWithContext.scala | 12 +-
.../pekko/stream/scaladsl/StreamConverters.scala | 4 +-
.../org/apache/pekko/stream/scaladsl/TLS.scala | 6 +-
.../org/apache/pekko/stream/scaladsl/Tcp.scala | 6 +-
.../org/apache/pekko/stream/scaladsl/package.scala | 2 +-
.../stream/serialization/StreamRefSerializer.scala | 36 +-
.../org/apache/pekko/stream/stage/GraphStage.scala | 315 ++++++--
.../apache/pekko/stream/stage/StageLogging.scala | 4 +-
.../pekko/testkit/CallingThreadDispatcher.scala | 2 +-
.../apache/pekko/testkit/TestEventListener.scala | 8 +-
.../apache/pekko/testkit/TestJavaSerializer.scala | 2 +-
.../scala/org/apache/pekko/testkit/TestKit.scala | 21 +-
.../org/apache/pekko/testkit/TestKitUtils.scala | 12 +-
.../apache/pekko/testkit/javadsl/EventFilter.scala | 10 +-
.../org/apache/pekko/testkit/javadsl/TestKit.scala | 4 +-
.../scala/org/apache/pekko/testkit/PekkoSpec.scala | 8 +-
.../org/apache/pekko/testkit/PekkoSpecSpec.scala | 2 +-
.../apache/pekko/testkit/TestActorRefSpec.scala | 2 +-
.../metrics/reporter/PekkoConsoleReporter.scala | 6 +-
1179 files changed, 20943 insertions(+), 10989 deletions(-)
create mode 100644 .github/scripts/resolve-scala-version.sh
create mode 100644
actor-typed-tests/src/test/scala/org/apache/pekko/actor/typed/CancelReceiveTimeoutSpec.scala
rename
actor-typed-tests/src/test/scala/org/apache/pekko/actor/typed/internal/{adpater
=> adapter}/PropsAdapterSpec.scala (97%)
copy
actor/src/main/mima-filters/{1.2.x.backwards.excludes/bytestring-indexOf-overload.excludes
=> 2.0.x.backwards.excludes/fsm-transition-listeners.excludes} (88%)
create mode 100644
bench-jmh/src/main/scala/org/apache/pekko/stream/ActorRefSourceBenchmark.scala
create mode 100644
bench-jmh/src/main/scala/org/apache/pekko/stream/AsyncBoundaryThroughputBenchmark.scala
create mode 100644
bench-jmh/src/main/scala/org/apache/pekko/stream/BroadcastHubBenchRunner.scala
create mode 100644
bench-jmh/src/main/scala/org/apache/pekko/stream/MaterializerWiringBenchmark.scala
create mode 100644
bench-jmh/src/main/scala/org/apache/pekko/stream/StageActorRefBenchmark.scala
copy
cluster-tools/src/main/mima-filters/2.0.x.backwards.excludes/remove-deprecated-methods.excludes
=>
cluster-sharding-typed/src/main/mima-filters/2.0.x.backwards.excludes/remove-old-passivation-strategy.excludes
(80%)
create mode 100644
cluster-sharding-typed/src/test/java/jdocs/org/apache/pekko/cluster/sharding/typed/ConsistentHashingShardAllocationCompileOnlyTest.java
create mode 100644
cluster-sharding-typed/src/test/scala/docs/org/apache/pekko/cluster/sharding/typed/ConsistentHashingShardAllocationCompileOnlySpec.scala
copy
cluster-sharding/src/main/mima-filters/{1.0.x.backwards.excludes/jdk-11-specific-classes.excludes
=> 2.0.x.backwards.excludes/consistent-hashing-allocation.excludes} (69%)
copy
stream/src/main/mima-filters/2.0.x.backwards.excludes/remove-gunzip.excludes =>
cluster-sharding/src/main/mima-filters/2.0.x.backwards.excludes/remove-old-passivation-strategy.excludes
(76%)
create mode 100644
cluster-sharding/src/main/scala/org/apache/pekko/cluster/sharding/ConsistentHashingShardAllocationStrategy.scala
copy
cluster-sharding/src/main/scala/org/apache/pekko/cluster/sharding/internal/{AbstractLeastShardAllocationStrategy.scala
=> ClusterShardAllocationMixin.scala} (63%)
create mode 100644
cluster-sharding/src/test/scala/org/apache/pekko/cluster/sharding/ConsistentHashingShardAllocationStrategySpec.scala
create mode 100644
cluster-sharding/src/test/scala/org/apache/pekko/cluster/sharding/ShardRegionDataTypesSpec.scala
create mode 100644
docs/src/main/java/docs/persistence/state/MyJavaStateStore.java
create mode 100644
docs/src/main/java/docs/persistence/state/MyJavaStateStoreProvider.java
create mode 100644 docs/src/main/paradox/durable-state/state-store-plugin.md
create mode 100644
docs/src/main/scala/docs/persistence/state/MyStateStore.scala
create mode 100644
docs/src/test/scala/docs/persistence/state/PersistenceStatePluginDocSpec.scala
create mode 100644
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/EventsBySliceFirehoseReadJournalProvider.scala
create mode 100644
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/internal/EventsBySliceFirehose.scala
copy
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/javadsl/{CurrentEventsByPersistenceIdTypedQuery.scala
=> CurrentEventsByPersistenceIdStartingFromSnapshotQuery.scala} (58%)
copy
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/javadsl/{CurrentEventsBySliceQuery.scala
=> CurrentEventsBySliceStartingFromSnapshotsQuery.scala} (51%)
copy
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/javadsl/{CurrentEventsByPersistenceIdTypedQuery.scala
=> EventsByPersistenceIdStartingFromSnapshotQuery.scala} (58%)
create mode 100644
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/javadsl/EventsBySliceFirehoseQuery.scala
copy
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/javadsl/{CurrentEventsBySliceQuery.scala
=> EventsBySliceStartingFromSnapshotsQuery.scala} (53%)
copy
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/scaladsl/{CurrentEventsByPersistenceIdTypedQuery.scala
=> CurrentEventsByPersistenceIdStartingFromSnapshotQuery.scala} (58%)
copy
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/scaladsl/{CurrentEventsBySliceQuery.scala
=> CurrentEventsBySliceStartingFromSnapshotsQuery.scala} (52%)
copy
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/scaladsl/{CurrentEventsByPersistenceIdTypedQuery.scala
=> EventsByPersistenceIdStartingFromSnapshotQuery.scala} (58%)
create mode 100644
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/scaladsl/EventsBySliceFirehoseQuery.scala
copy
persistence-query/src/main/scala/org/apache/pekko/persistence/query/typed/scaladsl/{CurrentEventsBySliceQuery.scala
=> EventsBySliceStartingFromSnapshotsQuery.scala} (53%)
create mode 100644
persistence-query/src/test/scala/org/apache/pekko/persistence/query/TestClock.scala
create mode 100644
persistence-query/src/test/scala/org/apache/pekko/persistence/query/typed/EventEnvelopeSpec.scala
create mode 100644
persistence-query/src/test/scala/org/apache/pekko/persistence/query/typed/internal/EventsBySliceFirehoseSpec.scala
create mode 100644
persistence-typed-tests/src/test/java/org/apache/pekko/persistence/typed/state/javadsl/RuntimeDurableStateStoreTest.java
create mode 100644
persistence-typed-tests/src/test/scala/org/apache/pekko/persistence/typed/state/scaladsl/RuntimeDurableStateStoreSpec.scala
copy
stream/src/main/mima-filters/1.0.x.backwards.excludes/pr-1374-boundedsourcequeue-iscompleted-classes.backwards.excludes
=>
persistence-typed/src/main/mima-filters/2.0.x.backwards.excludes/durablestate-runtime-config.excludes
(84%)
copy
actor/src/main/mima-filters/1.0.x.backwards.excludes/bytestring-inputstream.excludes
=>
persistence/src/main/mima-filters/2.0.x.backwards.excludes/fsm-transition-listeners.excludes
(84%)
create mode 100644
persistence/src/main/scala/org/apache/pekko/persistence/FilteredPayload.scala
create mode 100644
persistence/src/main/scala/org/apache/pekko/persistence/serialization/FilteredPayloadSerializer.scala
create mode 100644
persistence/src/test/scala/org/apache/pekko/persistence/serialization/FilteredPayloadSerializerSpec.scala
rename project/{CopyrightHeaderForJdk9.scala => CopyrightHeaderForJdk21.scala}
(82%)
rename project/{Jdk9.scala => Jdk21.scala} (63%)
rename project/{ScalaFixForJdk9Plugin.scala => ScalaFixForJdk21Plugin.scala}
(70%)
copy
actor-testkit-typed/src/main/mima-filters/2.0.x.backwards.excludes/renamed-methods.excludes
=>
remote/src/main/mima-filters/2.0.x.backwards.excludes/future-lazy-vals.excludes
(76%)
create mode 100644
remote/src/test/scala/org/apache/pekko/remote/ConfigSSLEngineProviderSpec.scala
create mode 100644
remote/src/test/scala/org/apache/pekko/remote/ReliableDeliverySupervisorSpec.scala
create mode 100644
remote/src/test/scala/org/apache/pekko/remote/artery/ActorSelectionQueueDistributionSpec.scala
create mode 100644
remote/src/test/scala/org/apache/pekko/remote/artery/tcp/ArteryConfigSSLEngineProviderSpec.scala
create mode 100644
remote/src/test/scala/org/apache/pekko/remote/transport/ThrottlerHandleSpec.scala
copy
stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/TlsEngineSelectionSpec.scala
=>
remote/src/test/scala/org/apache/pekko/remote/transport/netty/NettySSLSupportSpec.scala
(52%)
create mode 100644
stream-tests/src/test/java/org/apache/pekko/stream/javadsl/SourceWithContextTest.java
create mode 100644
stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/FlowAlsoToSpec.scala
copy
stream/src/main/mima-filters/{1.0.x.backwards.excludes/28324-jdk9-specific-classes.backwards.excludes
=> 2.0.x.backwards.excludes/remove-dead-actorprocessor.excludes} (63%)
copy
stream/src/main/mima-filters/{1.0.x.backwards.excludes/pr-991-remove-unused-attributes.backwards.excludes
=> 2.0.x.backwards.excludes/remove-deprecated-materializer-settings.excludes}
(53%)
copy
stream/src/main/mima-filters/{1.0.x.backwards.excludes/pr-1619-remove-supervised-group-stage-logic.backwards.excludes
=> 2.0.x.backwards.excludes/remove-fanout-publisher-sink.excludes} (87%)
copy
cluster/src/main/mima-filters/2.0.x.backwards.excludes/remove-deprecated-methods.excludes
=>
stream/src/main/mima-filters/2.0.x.backwards.excludes/remove-substream-cancel-strategy.excludes
(58%)
delete mode 100644
stream/src/main/scala/org/apache/pekko/stream/StreamRefSettings.scala
delete mode 100644
stream/src/main/scala/org/apache/pekko/stream/SubstreamCancelStrategy.scala
delete mode 100644
stream/src/main/scala/org/apache/pekko/stream/impl/ExposedPublisherReceive.scala
copy
actor/src/main/scala/org/apache/pekko/io/dns/internal/AsyncDnsProvider.scala =>
stream/src/main/scala/org/apache/pekko/stream/impl/streamref/StreamRefDefaultSettings.scala
(54%)
delete mode 100644
stream/src/main/scala/org/apache/pekko/stream/impl/streamref/StreamRefSettingsImpl.scala
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]