He-Pin opened a new pull request, #3409:
URL: https://github.com/apache/pekko/pull/3409

   ### Motivation
   Sometimes you want to wait for a `Sink` to fully complete, including any 
cleanup work or final commit it performs in `postStop`, but the sink does not 
materialize a `Future[Done]`. The existing `watchTermination` operator is 
placed *before* the sink and therefore only signals when the upstream of the 
sink has terminated.
   
   ### Modification
   Add `Sink.watchTermination` to the Scala and Java DSLs. It wraps sinks that 
consist of a single `GraphStage` with a delegating stage whose materialized 
`Future[Done]`/`CompletionStage<Done>`:
   - completes with success only after the wrapped sink's `postStop` has run,
   - fails with the upstream failure when the stream failed,
   - fails with an `AbruptStreamTerminationException` when the stream was 
abruptly terminated.
   
   The original materialized value, including `mapMaterializedValue` 
transforms, is preserved. Composite sinks consisting of multiple stages are 
rejected with an `IllegalArgumentException` (documented in the operator page 
and scaladoc).
   
   Implementation details:
   - `WatchedSink` rewrites the sink traversal, replacing the single terminal 
stage with a `WatchedSinkStage` and replaying the trailing materialized-value 
composition steps, preserving attributes and islands (`.async`).
   - `WatchedSinkLogic` delegates all port handlers and lifecycle hooks to the 
wrapped logic, mirroring the interpreter, port wiring, `stageId` and 
attributes, and records termination causes from the delegated handlers, handler 
exceptions, and the connection failure slot (covering wrapped stages that swap 
their inlet handler after materialization, e.g. lazy sinks).
   - `GraphStageLogic` gains an internal termination hook fired from 
`afterPostStop` so the promise also completes when the interpreter finalizes 
the wrapped logic directly (async-callback self-termination, e.g. `Sink.queue` 
cancellation).
   - Operator docs page, operator index entries and Scala/Java doc snippets are 
included.
   
   ### Result
   Users can await full sink termination, including `postStop` cleanup, via a 
materialized `Future[Done]` / `CompletionStage<Done>`.
   
   ### Tests
   - `sbt "stream-tests/Test/testOnly 
org.apache.pekko.stream.scaladsl.SinkWatchTerminationSpec"` - 17/17 passed
   - `sbt "stream-tests/Test/testOnly org.apache.pekko.stream.scaladsl.*Sink*"` 
- 159 passed
   - `sbt "stream-tests/Test/testOnly 
org.apache.pekko.stream.scaladsl.FlowWatchTerminationSpec 
org.apache.pekko.stream.scaladsl.QueueSinkSpec 
org.apache.pekko.stream.scaladsl.GraphStageTimersSpec 
org.apache.pekko.stream.impl.GraphStageLogicSpec 
org.apache.pekko.stream.impl.SubInletOutletSpec 
org.apache.pekko.stream.impl.LinearTraversalBuilderSpec 
org.apache.pekko.stream.DslConsistencySpec"` - 124 passed
   - `sbt "stream-tests/Test/testOnly 
org.apache.pekko.stream.javadsl.SinkTest"` - passed
   - `sbt stream/mimaReportBinaryIssues` - no issues
   - `sbt "++3.3.8" stream/compile` - passed
   - `sbt docs/paradox` - passed
   - `sbt headerCreateAll scalafmtAll scalafmtSbt javafmtCheckAll` - passed
   - `scalafmt --mode diff-ref=origin/main` - no changes
   - `git diff --check` - clean
   - `sbt sortImports` - failed with scalafix plugin `NoSuchMethodError` 
(environment issue), imports kept consistent manually
   
   ### References
   Fixes #2377


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to