forked from gunner95/vertx-rest
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathAbstractRestVerticle.java
More file actions
154 lines (136 loc) · 6 KB
/
Copy pathAbstractRestVerticle.java
File metadata and controls
154 lines (136 loc) · 6 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
package com.dream11.rest;
import com.dream11.rest.exception.mapper.GenericExceptionMapper;
import com.dream11.rest.exception.mapper.ValidationExceptionMapper;
import com.dream11.rest.exception.mapper.WebApplicationExceptionMapper;
import com.dream11.rest.filter.LoggerFilter;
import com.dream11.rest.filter.RequestResponseFilter;
import com.dream11.rest.filter.TimeoutFilter;
import com.dream11.rest.provider.JsonProvider;
import com.dream11.rest.provider.ParamConverterProvider;
import com.dream11.rest.provider.impl.JacksonProvider;
import com.dream11.rest.util.AnnotationUtil;
import io.reactivex.rxjava3.core.Completable;
import io.reactivex.rxjava3.core.Single;
import io.vertx.core.http.HttpServerOptions;
import io.vertx.core.json.jackson.DatabindCodec;
import io.vertx.rxjava3.core.AbstractVerticle;
import io.vertx.rxjava3.core.Context;
import io.vertx.rxjava3.core.RxHelper;
import io.vertx.rxjava3.core.http.HttpServer;
import io.vertx.rxjava3.core.http.HttpServerRequest;
import io.vertx.rxjava3.ext.web.Router;
import io.vertx.rxjava3.ext.web.handler.BodyHandler;
import io.vertx.rxjava3.ext.web.handler.ResponseContentTypeHandler;
import io.vertx.rxjava3.ext.web.handler.StaticHandler;
import jakarta.ws.rs.Path;
import jakarta.ws.rs.ext.Provider;
import java.util.ArrayList;
import java.util.List;
import lombok.extern.slf4j.Slf4j;
import lombok.val;
import org.jboss.resteasy.plugins.server.vertx.VertxRequestHandler;
import org.jboss.resteasy.plugins.server.vertx.VertxResteasyDeployment;
import org.jboss.resteasy.spi.ResteasyProviderFactory;
/**
* Starts backpressure enabled HTTP server and registers providers in resteasy deployment.
*/
@Slf4j
public abstract class AbstractRestVerticle extends AbstractVerticle {
final String packageName;
final HttpServerOptions httpServerOptions;
HttpServer httpServer;
protected AbstractRestVerticle(String packageName) {
this(packageName, new HttpServerOptions());
}
protected AbstractRestVerticle(String packageName, HttpServerOptions httpServerOptions) {
this.packageName = packageName;
this.httpServerOptions = httpServerOptions;
}
protected abstract ClassInjector getInjector();
protected RequestResponseFilter getReqResFilter() {
return new LoggerFilter();
}
protected JsonProvider getJsonProvider() {
return new JacksonProvider(DatabindCodec.mapper());
}
@Override
public Completable rxStart() {
return this.startHttpServer().doOnSuccess(server -> this.httpServer = server).ignoreElement();
}
private Single<HttpServer> startHttpServer() {
VertxResteasyDeployment deployment = this.buildResteasyDeployment();
Router router = this.getRouter();
VertxRequestHandler vertxRequestHandler = new VertxRequestHandler(vertx.getDelegate(), deployment);
val server = vertx.createHttpServer(this.httpServerOptions);
val handleRequests = server.requestStream()
.toFlowable()
.map(HttpServerRequest::pause)
.onBackpressureDrop(req -> {
log.error("Dropping request with status 503");
req.getDelegate().response().setStatusCode(503).end();
})
.observeOn(RxHelper.scheduler(new Context(this.context))) // For backpressure
.doOnNext(req -> {
if (req.path().matches("/swagger(.*)")) {
router.handle(req);
} else {
vertxRequestHandler.handle(req.getDelegate());
}
})
.map(HttpServerRequest::resume)
.doOnError(error -> log.error("Uncaught ERROR while handling request", error))
.ignoreElements();
return server
.rxListen()
.doOnSuccess(res -> log.info("Started http server at port: {} for package: {}", this.httpServerOptions.getPort(), packageName))
.doOnError(error -> log.error(
"Failed to start http server at port : {} with error: {}", this.httpServerOptions.getPort(), error.getMessage()))
.doOnSubscribe(disposable -> handleRequests.subscribe());
}
protected Router getRouter() {
Router router = Router.router(vertx);
router.route().handler(BodyHandler.create());
router.route().handler(ResponseContentTypeHandler.create());
router.route().handler(StaticHandler.create());
return router;
}
@Override
public Completable rxStop() {
if (this.httpServer != null) {
return this.httpServer.rxClose()
.doOnComplete(() -> log.info("http server stopped successfully"))
.doOnError(err -> log.info("Failed to stop http server", err));
}
return Completable.complete();
}
protected VertxResteasyDeployment buildResteasyDeployment() {
VertxResteasyDeployment deployment = new VertxResteasyDeployment();
deployment.start();
List<Class<?>> routes = AnnotationUtil.getClassesWithAnnotation(packageName, Path.class);
log.info("JAX-RS routes : " + routes.size());
this.registerProviders(deployment.getProviderFactory());
// not using deployment.getRegistry().addPerInstanceResource because it creates new instance of resource for each request
routes.forEach(route -> deployment.getRegistry().addSingletonResource(this.getInjector().getInstance(route)));
return deployment;
}
private void registerProviders(ResteasyProviderFactory resteasyProviderFactory) {
this.getProviderObjects().forEach(resteasyProviderFactory::register);
this.getProviders().forEach(resteasyProviderFactory::register);
}
protected List<Class<?>> getProviders() {
List<Class<?>> providers = new ArrayList<>();
providers.add(TimeoutFilter.class);
providers.add(ValidationExceptionMapper.class);
providers.add(GenericExceptionMapper.class);
providers.add(WebApplicationExceptionMapper.class);
providers.add(ParamConverterProvider.class);
providers.add(this.getReqResFilter().getClass());
providers.addAll(AnnotationUtil.getClassesWithAnnotation(packageName, Provider.class));
return providers;
}
protected List<Object> getProviderObjects() {
List<Object> providerObjects = new ArrayList<>();
providerObjects.add(this.getJsonProvider());
return providerObjects;
}
}