diff --git a/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/ObservableSseEmitter.java b/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/ObservableSseEmitter.java index 14a5ae58..a267c54c 100644 --- a/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/ObservableSseEmitter.java +++ b/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/ObservableSseEmitter.java @@ -18,6 +18,7 @@ package org.springframework.cloud.netflix.rx; import org.springframework.http.MediaType; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; + import rx.Observable; /** @@ -28,18 +29,17 @@ import rx.Observable; */ class ObservableSseEmitter extends SseEmitter { - private final ResponseBodyEmitterSubscriber subscriber; + public ObservableSseEmitter(Observable observable) { + this(null, observable); + } - public ObservableSseEmitter(Observable observable) { - this(null, observable); - } + public ObservableSseEmitter(MediaType mediaType, Observable observable) { + this(null, mediaType, observable); + } - public ObservableSseEmitter(MediaType mediaType, Observable observable) { - this(null, mediaType, observable); - } - - public ObservableSseEmitter(Long timeout, MediaType mediaType, Observable observable) { - super(timeout); - this.subscriber = new ResponseBodyEmitterSubscriber<>(mediaType, observable, this); - } + public ObservableSseEmitter(Long timeout, MediaType mediaType, + Observable observable) { + super(timeout); + new ResponseBodyEmitterSubscriber<>(mediaType, observable, this); + } } diff --git a/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/SingleDeferredResult.java b/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/SingleDeferredResult.java index af313392..232a053e 100644 --- a/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/SingleDeferredResult.java +++ b/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/SingleDeferredResult.java @@ -29,22 +29,19 @@ import rx.Single; */ class SingleDeferredResult extends DeferredResult { - private static final Object EMPTY_RESULT = new Object(); + private static final Object EMPTY_RESULT = new Object(); - private final DeferredResultSubscriber subscriber; + public SingleDeferredResult(Single single) { + this(null, EMPTY_RESULT, single); + } - public SingleDeferredResult(Single single) { - this(null, EMPTY_RESULT, single); - } + public SingleDeferredResult(long timeout, Single single) { + this(timeout, EMPTY_RESULT, single); + } - public SingleDeferredResult(long timeout, Single single) { - this(timeout, EMPTY_RESULT, single); - } - - public SingleDeferredResult(Long timeout, Object timeoutResult, Single single) { - super(timeout, timeoutResult); - Assert.notNull(single, "single can not be null"); - - subscriber = new DeferredResultSubscriber<>(single.toObservable(), this); - } + public SingleDeferredResult(Long timeout, Object timeoutResult, Single single) { + super(timeout, timeoutResult); + Assert.notNull(single, "single can not be null"); + new DeferredResultSubscriber<>(single.toObservable(), this); + } } diff --git a/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/SingleReturnValueHandler.java b/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/SingleReturnValueHandler.java index 190a06d7..1db91b76 100644 --- a/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/SingleReturnValueHandler.java +++ b/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/rx/SingleReturnValueHandler.java @@ -31,42 +31,46 @@ import rx.Single; import rx.functions.Func1; /** - * A specialized {@link AsyncHandlerMethodReturnValueHandler} that handles {@link Single} return types. + * A specialized {@link AsyncHandlerMethodReturnValueHandler} that handles {@link Single} + * return types. * * @author Spencer Gibb * @author Jakub Narloch */ public class SingleReturnValueHandler implements AsyncHandlerMethodReturnValueHandler { - @Override - public boolean isAsyncReturnValue(Object returnValue, MethodParameter returnType) { - return returnValue != null && supportsReturnType(returnType); - } + @Override + public boolean isAsyncReturnValue(Object returnValue, MethodParameter returnType) { + return returnValue != null && supportsReturnType(returnType); + } - @Override - public boolean supportsReturnType(MethodParameter returnType) { - return Single.class.isAssignableFrom(returnType.getParameterType()) || isResponseEntity(returnType); - } + @Override + public boolean supportsReturnType(MethodParameter returnType) { + return Single.class.isAssignableFrom(returnType.getParameterType()) + || isResponseEntity(returnType); + } private boolean isResponseEntity(MethodParameter returnType) { - if(ResponseEntity.class.isAssignableFrom(returnType.getParameterType())) { - Class bodyType = ResolvableType.forMethodParameter(returnType).getGeneric(0).resolve(); + if (ResponseEntity.class.isAssignableFrom(returnType.getParameterType())) { + Class bodyType = ResolvableType.forMethodParameter(returnType) + .getGeneric(0).resolve(); return bodyType != null && Single.class.isAssignableFrom(bodyType); } return false; } - @SuppressWarnings("unchecked") - @Override - public void handleReturnValue(Object returnValue, MethodParameter returnType, ModelAndViewContainer mavContainer, NativeWebRequest webRequest) throws Exception { + @Override + public void handleReturnValue(Object returnValue, MethodParameter returnType, + ModelAndViewContainer mavContainer, NativeWebRequest webRequest) + throws Exception { - if (returnValue == null) { - mavContainer.setRequestHandled(true); - return; - } + if (returnValue == null) { + mavContainer.setRequestHandled(true); + return; + } ResponseEntity> responseEntity = getResponseEntity(returnValue); - if(responseEntity != null) { + if (responseEntity != null) { returnValue = responseEntity.getBody(); if (returnValue == null) { mavContainer.setRequestHandled(true); @@ -74,10 +78,10 @@ public class SingleReturnValueHandler implements AsyncHandlerMethodReturnValueHa } } - final Single single = Single.class.cast(returnValue); - WebAsyncUtils.getAsyncManager(webRequest) - .startDeferredResultProcessing(convertToDeferredResult(responseEntity, single), mavContainer); - } + final Single single = Single.class.cast(returnValue); + WebAsyncUtils.getAsyncManager(webRequest).startDeferredResultProcessing( + convertToDeferredResult(responseEntity, single), mavContainer); + } @SuppressWarnings("unchecked") private ResponseEntity> getResponseEntity(Object returnValue) { @@ -88,28 +92,32 @@ public class SingleReturnValueHandler implements AsyncHandlerMethodReturnValueHa return null; } - protected DeferredResult convertToDeferredResult(final ResponseEntity> responseEntity, Single single) { + protected DeferredResult convertToDeferredResult( + final ResponseEntity> responseEntity, Single single) { - //TODO: fix when java8 :-) - Single singleResponse = single.map(new Func1() { - @Override - public ResponseEntity call(Object object) { - return new ResponseEntity<>(object, getHttpHeaders(responseEntity), getHttpStatus(responseEntity)); - } - }); + // TODO: use lambda when java8 :-) + Single> singleResponse = single + .map(new Func1>() { + @Override + public ResponseEntity call(Object object) { + return new ResponseEntity(object, + getHttpHeaders(responseEntity), + getHttpStatus(responseEntity)); + } + }); return new SingleDeferredResult<>(singleResponse); } private HttpStatus getHttpStatus(ResponseEntity responseEntity) { - if(responseEntity == null) { + if (responseEntity == null) { return HttpStatus.OK; } return responseEntity.getStatusCode(); } private HttpHeaders getHttpHeaders(ResponseEntity responseEntity) { - if(responseEntity == null) { + if (responseEntity == null) { return new HttpHeaders(); } return responseEntity.getHeaders();