@@ -504,7 +504,9 @@ CompletableFuture<Void> emitEvent(Event event) {
504504 assert event != Event .EOF : "must not use synthetic EOF event" ;
505505 assert !(event instanceof Event .StreamHangup ) : "must not use synthetic StreamHangup event" ;
506506
507- return CompletableFuture .runAsync (() -> stream .onNext (event ), eventThread );
507+ log .debug ("emit event {}" , event );
508+ return CompletableFuture .runAsync (() -> stream .onNext (event ), eventThread )
509+ .thenAccept (r -> log .debug ("ran {}" , r ));
508510 }
509511
510512 /** Terminate the server half of the stream abruptly. */
@@ -517,6 +519,7 @@ CompletableFuture<Void> hangup() {
517519
518520 /** Emit events before closing the server half of the stream. */
519521 void beforeEof (Event ... events ) {
522+ log .debug ("add {} events for before-EOF" , events .length );
520523 this .pendingEvents .addAll (Arrays .asList (events ));
521524 }
522525
@@ -527,6 +530,8 @@ void beforeEof(Event... events) {
527530 * stream is being hung up by the client, not by us.
528531 */
529532 CompletableFuture <Void > eof (boolean ok ) {
533+ log .debug ("emit EOF, ok={}" , ok );
534+ log .debug ("have {} events pre-hooked" , this .pendingEvents .size ());
530535 CompletableFuture <?> await = ok
531536 ? CompletableFuture .allOf (pendingEvents .stream ().map (this ::emitEvent ).toArray (CompletableFuture []::new ))
532537 : CompletableFuture .completedFuture (null );
0 commit comments