in core/src/main/scala/org/apache/pekko/kafka/internal/DefaultProducerStage.scala [89:96]
private def checkForCompletion(): Unit =
if (isClosed(stage.in) && awaitingConfirmation == 0) {
completionState match {
case Some(Success(_)) => onCompletionSuccess()
case Some(Failure(ex)) => onCompletionFailure(ex)
case None => failStage(new IllegalStateException("Stage completed, but there is no info about status"))
}
}