1919import io .reactivesocket .Frame .Lease ;
2020import io .reactivesocket .Frame .Request ;
2121import io .reactivesocket .Frame .Response ;
22+ import io .reactivesocket .events .EventListener ;
23+ import io .reactivesocket .events .EventListener .RequestType ;
24+ import io .reactivesocket .events .EventPublishingSocket ;
25+ import io .reactivesocket .events .EventPublishingSocketImpl ;
2226import io .reactivesocket .exceptions .ApplicationException ;
2327import io .reactivesocket .exceptions .SetupException ;
2428import io .reactivesocket .frame .FrameHeaderFlyweight ;
29+ import io .reactivesocket .internal .DisabledEventPublisher ;
30+ import io .reactivesocket .internal .EventPublisher ;
2531import io .reactivesocket .internal .KnownErrorFilter ;
2632import io .reactivesocket .internal .RemoteReceiver ;
2733import io .reactivesocket .internal .RemoteSender ;
2834import io .reactivesocket .lease .LeaseEnforcingSocket ;
2935import io .reactivesocket .reactivestreams .extensions .Px ;
3036import io .reactivesocket .reactivestreams .extensions .internal .subscribers .Subscribers ;
37+ import io .reactivesocket .util .Clock ;
3138import org .agrona .collections .Int2ObjectHashMap ;
3239import org .reactivestreams .Publisher ;
3340import org .reactivestreams .Subscription ;
3441
3542import java .util .Collection ;
3643import java .util .function .Consumer ;
3744
45+ import static io .reactivesocket .events .EventListener .RequestType .*;
46+
3847/**
3948 * Server side ReactiveSocket. Receives {@link Frame}s from a
4049 * {@link ClientReactiveSocket}
@@ -44,20 +53,28 @@ public class ServerReactiveSocket implements ReactiveSocket {
4453 private final DuplexConnection connection ;
4554 private final Publisher <Frame > serverInput ;
4655 private final Consumer <Throwable > errorConsumer ;
56+ private final EventPublisher <? extends EventListener > eventPublisher ;
4757
4858 private final Int2ObjectHashMap <Subscription > subscriptions ;
4959 private final Int2ObjectHashMap <RemoteReceiver > channelProcessors ;
5060
5161 private final ReactiveSocket requestHandler ;
5262 private Subscription receiversSubscription ;
63+ private final EventPublishingSocket eventPublishingSocket ;
64+
5365 public ServerReactiveSocket (DuplexConnection connection , ReactiveSocket requestHandler ,
54- boolean clientHonorsLease , Consumer <Throwable > errorConsumer ) {
66+ boolean clientHonorsLease , Consumer <Throwable > errorConsumer ,
67+ EventPublisher <? extends EventListener > eventPublisher ) {
5568 this .requestHandler = requestHandler ;
5669 this .connection = connection ;
5770 serverInput = connection .receive ();
5871 this .errorConsumer = new KnownErrorFilter (errorConsumer );
72+ this .eventPublisher = eventPublisher ;
5973 subscriptions = new Int2ObjectHashMap <>();
6074 channelProcessors = new Int2ObjectHashMap <>();
75+ eventPublishingSocket = eventPublisher .isEventPublishingEnabled ()?
76+ new EventPublishingSocketImpl (eventPublisher , false ) : EventPublishingSocket .DISABLED ;
77+
6178 Px .from (connection .onClose ()).subscribe (Subscribers .cleanup (() -> {
6279 cleanup ();
6380 }));
@@ -74,6 +91,10 @@ public ServerReactiveSocket(DuplexConnection connection, ReactiveSocket requestH
7491 });
7592 }
7693 }
94+ public ServerReactiveSocket (DuplexConnection connection , ReactiveSocket requestHandler ,
95+ boolean clientHonorsLease , Consumer <Throwable > errorConsumer ) {
96+ this (connection , requestHandler , clientHonorsLease , errorConsumer , new DisabledEventPublisher <>());
97+ }
7798
7899 public ServerReactiveSocket (DuplexConnection connection , ReactiveSocket requestHandler ,
79100 Consumer <Throwable > errorConsumer ) {
@@ -165,19 +186,19 @@ private Publisher<Void> handleFrame(Frame frame) {
165186 case SETUP :
166187 return Px .error (new IllegalStateException ("Setup frame received post setup." ));
167188 case REQUEST_RESPONSE :
168- return handleReceive (streamId , requestResponse (frame ));
189+ return handleRequestResponse (streamId , requestResponse (frame ));
169190 case CANCEL :
170191 return handleCancelFrame (streamId );
171192 case KEEPALIVE :
172193 return handleKeepAliveFrame (frame );
173194 case REQUEST_N :
174195 return handleRequestN (streamId , frame );
175196 case REQUEST_STREAM :
176- return doReceive (streamId , requestStream (frame ));
197+ return doReceive (streamId , requestStream (frame ), RequestStream );
177198 case FIRE_AND_FORGET :
178199 return handleFireAndForget (streamId , fireAndForget (frame ));
179200 case REQUEST_SUBSCRIPTION :
180- return doReceive (streamId , requestSubscription (frame ));
201+ return doReceive (streamId , requestSubscription (frame ), RequestStream );
181202 case REQUEST_CHANNEL :
182203 return handleChannel (streamId , frame );
183204 case RESPONSE :
@@ -251,13 +272,14 @@ private synchronized void cleanup() {
251272 requestHandler .close ().subscribe (Subscribers .empty ());
252273 }
253274
254- private Publisher <Void > handleReceive (int streamId , Publisher <Payload > response ) {
275+ private Publisher <Void > handleRequestResponse (int streamId , Publisher <Payload > response ) {
255276 final Runnable cleanup = () -> {
256277 synchronized (this ) {
257278 subscriptions .remove (streamId );
258279 }
259280
260281 };
282+ long now = publishSingleFrameReceiveEvents (streamId , RequestResponse );
261283
262284 Px <Frame > frames =
263285 Px
@@ -282,26 +304,29 @@ private Publisher<Void> handleReceive(int streamId, Publisher<Payload> response)
282304 return Frame .Error .from (streamId , throwable );
283305 });
284306
285- return Px .from (connection .send (frames ));
307+ return Px .from (eventPublishingSocket . decorateSend ( streamId , connection .send (frames ), now , RequestResponse ));
286308
287309 }
288310
289- private Publisher <Void > doReceive (int streamId , Publisher <Payload > response ) {
311+ private Publisher <Void > doReceive (int streamId , Publisher <Payload > response , RequestType requestType ) {
312+ long now = publishSingleFrameReceiveEvents (streamId , requestType );
290313 Px <Frame > resp = Px .from (response )
291314 .map (payload -> Response .from (streamId , FrameType .RESPONSE , payload ));
292315 RemoteSender sender = new RemoteSender (resp , () -> subscriptions .remove (streamId ), streamId , 2 );
293316 subscriptions .put (streamId , sender );
294- return connection .send (sender );
317+ return eventPublishingSocket . decorateSend ( streamId , connection .send (sender ), now , requestType );
295318 }
296319
297320 private Publisher <Void > handleChannel (int streamId , Frame firstFrame ) {
321+ long now = publishSingleFrameReceiveEvents (streamId , RequestChannel );
298322 int initialRequestN = Request .initialRequestN (firstFrame );
299323 Frame firstAsNext = Request .from (streamId , FrameType .NEXT , firstFrame , initialRequestN );
300324 RemoteReceiver receiver = new RemoteReceiver (connection , streamId , () -> removeChannelProcessor (streamId ),
301325 firstAsNext , receiversSubscription , true );
302326 channelProcessors .put (streamId , receiver );
303327
304- Px <Frame > response = Px .from (requestChannel (receiver ))
328+ Px <Frame > response = Px .from (requestChannel (eventPublishingSocket .decorateReceive (streamId , receiver ,
329+ RequestChannel )))
305330 .map (payload -> Response .from (streamId , FrameType .RESPONSE , payload ));
306331
307332 RemoteSender sender = new RemoteSender (response , () -> removeSubscriptions (streamId ), streamId ,
@@ -310,7 +335,7 @@ private Publisher<Void> handleChannel(int streamId, Frame firstFrame) {
310335 subscriptions .put (streamId , sender );
311336 }
312337
313- return connection .send (sender );
338+ return eventPublishingSocket . decorateSend ( streamId , connection .send (sender ), now , RequestChannel );
314339 }
315340
316341 private Publisher <Void > handleFireAndForget (int streamId , Publisher <Void > result ) {
@@ -368,4 +393,14 @@ private synchronized void addSubscription(int streamId, Subscription subscriptio
368393 private synchronized void removeSubscription (int streamId ) {
369394 subscriptions .remove (streamId );
370395 }
396+
397+ private long publishSingleFrameReceiveEvents (int streamId , RequestType requestType ) {
398+ long now = Clock .now ();
399+ if (eventPublisher .isEventPublishingEnabled ()) {
400+ EventListener eventListener = eventPublisher .getEventListener ();
401+ eventListener .requestReceiveStart (streamId , requestType );
402+ eventListener .requestReceiveComplete (streamId , requestType , Clock .elapsedSince (now ), Clock .unit ());
403+ }
404+ return now ;
405+ }
371406}
0 commit comments