diff --git a/.github/workflows/maven-deploy-release.yml b/.github/workflows/maven-deploy-release.yml
index dbf9bc0..b241f52 100644
--- a/.github/workflows/maven-deploy-release.yml
+++ b/.github/workflows/maven-deploy-release.yml
@@ -23,7 +23,6 @@ jobs:
- name: Deploy to JavaWebStack Repository
run: mvn deploy -B -DbuildVersion=${{ github.event.release.tag_name }} -s build/settings.xml -Dmaven.test.skip=true
env:
- DEPLOYMENT_USERNAME: ${{ secrets.DEPLOYMENT_USERNAME }}
- DEPLOYMENT_PASSWORD: ${{ secrets.DEPLOYMENT_PASSWORD }}
+ CENTRAL_USERNAME: ${{ secrets.CENTRAL_USERNAME }}
+ CENTRAL_PASSWORD: ${{ secrets.CENTRAL_PASSWORD }}
GPG_PASSPHRASE: ${{ secrets.GPG_PASSPHRASE }}
- OSSRH_PASSWORD: ${{ secrets.OSSRH_PASSWORD }}
diff --git a/.github/workflows/maven-deploy.yml b/.github/workflows/maven-deploy.yml
index b035e13..fd39229 100644
--- a/.github/workflows/maven-deploy.yml
+++ b/.github/workflows/maven-deploy.yml
@@ -24,7 +24,6 @@ jobs:
- name: Deploy to JavaWebStack Repository
run: mvn deploy -B -s build/settings.xml -Dmaven.test.skip=true
env:
- DEPLOYMENT_USERNAME: ${{ secrets.DEPLOYMENT_USERNAME }}
- DEPLOYMENT_PASSWORD: ${{ secrets.DEPLOYMENT_PASSWORD }}
+ CENTRAL_USERNAME: ${{ secrets.CENTRAL_USERNAME }}
+ CENTRAL_PASSWORD: ${{ secrets.CENTRAL_PASSWORD }}
GPG_PASSPHRASE: ${{ secrets.GPG_PASSPHRASE }}
- OSSRH_PASSWORD: ${{ secrets.OSSRH_PASSWORD }}
diff --git a/README.md b/README.md
index d15dc91..a13635e 100644
--- a/README.md
+++ b/README.md
@@ -23,4 +23,4 @@ JWS Job Scheduling
jobs
0.0.1-SNAPSHOT
-```
\ No newline at end of file
+```
diff --git a/build/settings.xml b/build/settings.xml
index 07e7726..88571fa 100644
--- a/build/settings.xml
+++ b/build/settings.xml
@@ -5,19 +5,9 @@
>
- javawebstack-snapshots
- ${env.DEPLOYMENT_USERNAME}
- ${env.DEPLOYMENT_PASSWORD}
-
-
- javawebstack-releases
- ${env.DEPLOYMENT_USERNAME}
- ${env.DEPLOYMENT_PASSWORD}
-
-
- ossrh
- JavaWebStack
- ${env.OSSRH_PASSWORD}
+ central
+ ${env.CENTRAL_USERNAME}
+ ${env.CENTRAL_PASSWORD}
gpg
diff --git a/pom.xml b/pom.xml
index c3fe7cc..273848f 100644
--- a/pom.xml
+++ b/pom.xml
@@ -40,33 +40,45 @@
https://github.com/JavaWebStack/jobs/tree/master
+
+
+ central-snapshots
+ https://central.sonatype.com/repository/maven-snapshots/
+
+
+
org.javawebstack
abstract-data
- 1.0.4
+ 1.0.7-SNAPSHOT
org.javawebstack
- http-server
- 1.0.1
+ http-router
+ 1.0.3-SNAPSHOT
+
+
+ org.javawebstack
+ http-router-undertow
+ 1.0.0-SNAPSHOT
org.javawebstack
orm
- 1.0.1
+ 1.0.3-SNAPSHOT
true
redis.clients
jedis
- 4.3.0-m1
+ 5.0.0
true
org.projectlombok
lombok
- 1.18.24
+ 1.18.28
provided
@@ -82,19 +94,19 @@
org.junit.jupiter
junit-jupiter-engine
- 5.9.0
+ 5.10.0
test
org.junit.jupiter
junit-jupiter-api
- 5.9.0
+ 5.10.0
test
mysql
mysql-connector-java
- 8.0.30
+ 8.0.32
test
@@ -102,17 +114,13 @@
- maven-deploy-plugin
- 3.0.0-M1
-
-
- default-deploy
- deploy
-
- deploy
-
-
-
+ org.sonatype.central
+ central-publishing-maven-plugin
+ 0.9.0
+ true
+
+ central
+
org.apache.maven.plugins
@@ -157,7 +165,7 @@
sign
- A313520526A8DFE1C2A30399C35A3D43C557B112
+ EC9CCFF8901F0AA22191DCEDD619376246C066D0
gpg
--no-tty
@@ -173,13 +181,9 @@
- ossrh
- https://s01.oss.sonatype.org/content/repositories/snapshots
+ central
+ https://central.sonatype.com/repository/maven-snapshots/
-
- ossrh
- https://s01.oss.sonatype.org/service/local/staging/deploy/maven2/
-
\ No newline at end of file
diff --git a/src/main/java/org/javawebstack/jobs/api/JobApi.java b/src/main/java/org/javawebstack/jobs/api/JobApi.java
index 8db32d0..d5e7080 100644
--- a/src/main/java/org/javawebstack/jobs/api/JobApi.java
+++ b/src/main/java/org/javawebstack/jobs/api/JobApi.java
@@ -1,9 +1,10 @@
package org.javawebstack.jobs.api;
import lombok.Getter;
-import org.javawebstack.httpserver.HTTPServer;
-import org.javawebstack.httpserver.helper.HttpMethod;
-import org.javawebstack.httpserver.transformer.response.JsonResponseTransformer;
+import org.javawebstack.http.router.HTTPMethod;
+import org.javawebstack.http.router.HTTPRouter;
+import org.javawebstack.http.router.transformer.response.JsonResponseTransformer;
+import org.javawebstack.http.router.undertow.UndertowHTTPSocketServer;
import org.javawebstack.jobs.Jobs;
import org.javawebstack.jobs.api.auth.AuthProvider;
import org.javawebstack.jobs.api.controller.*;
@@ -41,31 +42,31 @@ public JobApi dashboard(boolean enableDashboard) {
return this;
}
- public HTTPServer start(int port) {
- HTTPServer server = new HTTPServer()
+ public HTTPRouter start(int port) {
+ HTTPRouter router = new HTTPRouter(new UndertowHTTPSocketServer())
.responseTransformer(new JsonResponseTransformer().ignoreStrings())
.port(port);
- install(server, null);
- server.start();
- server.beforeInterceptor(ex -> {
+ install(router, null);
+ router.start();
+ router.beforeInterceptor(ex -> {
ex.header("Access-Control-Allow-Origin", "*");
ex.header("Access-Control-Allow-Methods", "*");
ex.header("Access-Control-Allow-Headers", "*");
- if(ex.getMethod() == HttpMethod.OPTIONS) {
+ if(ex.getMethod() == HTTPMethod.OPTIONS) {
ex.close();
return true;
}
return false;
});
- return server;
+ return router;
}
- public JobApi install(HTTPServer server, String prefix) {
+ public JobApi install(HTTPRouter router, String prefix) {
if(prefix == null)
prefix = "";
if(prefix.length() > 0 && enableDashboard)
throw new IllegalArgumentException("Prefix can not be set when the dashboard is enabled!");
- server
+ router
.exceptionHandler(new ErrorController())
.controller(prefix, new JobController(jobs))
.controller(prefix, new RecurringJobController(jobs))
@@ -74,13 +75,13 @@ public JobApi install(HTTPServer server, String prefix) {
.afterAny(prefix + "{*:path}", new ResponseMiddleware())
.middleware("jobs_auth", new AuthMiddleware(this));
if(enableDashboard) {
- server.get("/", ex -> {
+ router.get("/", ex -> {
ex.redirect("/overview");
return "";
});
- server.staticResourceDirectory(prefix, JobApi.class.getClassLoader(), "dashboard");
+ router.staticResourceDirectory(prefix, JobApi.class.getClassLoader(), "dashboard");
String html = loadHtml();
- server.get(prefix + "{*:path}", ex -> html);
+ router.get(prefix + "{*:path}", ex -> html);
}
// TODO install dashboard
return this;
diff --git a/src/main/java/org/javawebstack/jobs/api/controller/ErrorController.java b/src/main/java/org/javawebstack/jobs/api/controller/ErrorController.java
index 71101f0..f2123b7 100644
--- a/src/main/java/org/javawebstack/jobs/api/controller/ErrorController.java
+++ b/src/main/java/org/javawebstack/jobs/api/controller/ErrorController.java
@@ -3,8 +3,8 @@
import org.javawebstack.abstractdata.AbstractArray;
import org.javawebstack.abstractdata.AbstractObject;
import org.javawebstack.abstractdata.AbstractPrimitive;
-import org.javawebstack.httpserver.Exchange;
-import org.javawebstack.httpserver.handler.ExceptionHandler;
+import org.javawebstack.http.router.Exchange;
+import org.javawebstack.http.router.handler.ExceptionHandler;
import org.javawebstack.jobs.api.response.Response;
import org.javawebstack.validator.ValidationError;
import org.javawebstack.validator.ValidationException;
diff --git a/src/main/java/org/javawebstack/jobs/api/controller/JobController.java b/src/main/java/org/javawebstack/jobs/api/controller/JobController.java
index 66c9987..9280ea0 100644
--- a/src/main/java/org/javawebstack/jobs/api/controller/JobController.java
+++ b/src/main/java/org/javawebstack/jobs/api/controller/JobController.java
@@ -2,14 +2,14 @@
import org.javawebstack.abstractdata.AbstractElement;
import org.javawebstack.abstractdata.AbstractObject;
-import org.javawebstack.httpserver.Exchange;
-import org.javawebstack.httpserver.router.annotation.PathPrefix;
-import org.javawebstack.httpserver.router.annotation.With;
-import org.javawebstack.httpserver.router.annotation.params.Body;
-import org.javawebstack.httpserver.router.annotation.params.Path;
-import org.javawebstack.httpserver.router.annotation.verbs.Delete;
-import org.javawebstack.httpserver.router.annotation.verbs.Get;
-import org.javawebstack.httpserver.router.annotation.verbs.Post;
+import org.javawebstack.http.router.Exchange;
+import org.javawebstack.http.router.router.annotation.PathPrefix;
+import org.javawebstack.http.router.router.annotation.With;
+import org.javawebstack.http.router.router.annotation.params.Body;
+import org.javawebstack.http.router.router.annotation.params.Path;
+import org.javawebstack.http.router.router.annotation.verbs.Delete;
+import org.javawebstack.http.router.router.annotation.verbs.Get;
+import org.javawebstack.http.router.router.annotation.verbs.Post;
import org.javawebstack.jobs.JobStatus;
import org.javawebstack.jobs.Jobs;
import org.javawebstack.jobs.api.request.CreateJobRequest;
@@ -73,7 +73,7 @@ public Response get(@Path("id") UUID id, Exchange exchange) {
JobInfo info = storage.getJob(id);
if(info == null)
return Response.error(404, "Job not found");
- AbstractObject res = exchange.getServer().getAbstractMapper().toAbstract(info).object();
+ AbstractObject res = exchange.getRouter().getMapper().map(info).object();
if(exchange.getQueryParameters().has("payload") && (exchange.query("payload").length() == 0 || exchange.query("payload").equals("true")))
res.set("payload", AbstractElement.fromJson(storage.getJobPayload(info.getId())));
return Response.success().setData(res);
diff --git a/src/main/java/org/javawebstack/jobs/api/controller/RecurringJobController.java b/src/main/java/org/javawebstack/jobs/api/controller/RecurringJobController.java
index 2b9cd73..909b5a2 100644
--- a/src/main/java/org/javawebstack/jobs/api/controller/RecurringJobController.java
+++ b/src/main/java/org/javawebstack/jobs/api/controller/RecurringJobController.java
@@ -1,21 +1,18 @@
package org.javawebstack.jobs.api.controller;
-import org.javawebstack.abstractdata.AbstractElement;
import org.javawebstack.abstractdata.AbstractObject;
-import org.javawebstack.httpserver.Exchange;
-import org.javawebstack.httpserver.router.annotation.PathPrefix;
-import org.javawebstack.httpserver.router.annotation.With;
-import org.javawebstack.httpserver.router.annotation.params.Body;
-import org.javawebstack.httpserver.router.annotation.params.Path;
-import org.javawebstack.httpserver.router.annotation.verbs.Delete;
-import org.javawebstack.httpserver.router.annotation.verbs.Get;
-import org.javawebstack.httpserver.router.annotation.verbs.Post;
-import org.javawebstack.jobs.JobStatus;
+
+import org.javawebstack.http.router.Exchange;
+import org.javawebstack.http.router.router.annotation.PathPrefix;
+import org.javawebstack.http.router.router.annotation.With;
+import org.javawebstack.http.router.router.annotation.params.Body;
+import org.javawebstack.http.router.router.annotation.params.Path;
+import org.javawebstack.http.router.router.annotation.verbs.Delete;
+import org.javawebstack.http.router.router.annotation.verbs.Get;
+import org.javawebstack.http.router.router.annotation.verbs.Post;
import org.javawebstack.jobs.Jobs;
-import org.javawebstack.jobs.api.request.CreateJobRequest;
import org.javawebstack.jobs.api.request.CreateRecurringJobRequest;
import org.javawebstack.jobs.api.response.Response;
-import org.javawebstack.jobs.storage.model.JobInfo;
import org.javawebstack.jobs.storage.model.RecurringJobInfo;
import org.javawebstack.jobs.storage.model.RecurringJobQuery;
@@ -53,7 +50,7 @@ public Response get(@Path("id") UUID id, Exchange exchange) {
RecurringJobInfo info = storage.getRecurringJob(id);
if(info == null)
return Response.error(404, "Recurring job not found");
- AbstractObject res = exchange.getServer().getAbstractMapper().toAbstract(info).object();
+ AbstractObject res = exchange.getRouter().getMapper().map(info).object();
return Response.success().setData(res);
}
diff --git a/src/main/java/org/javawebstack/jobs/api/controller/StatusController.java b/src/main/java/org/javawebstack/jobs/api/controller/StatusController.java
index 0a6587f..7f7fe35 100644
--- a/src/main/java/org/javawebstack/jobs/api/controller/StatusController.java
+++ b/src/main/java/org/javawebstack/jobs/api/controller/StatusController.java
@@ -1,8 +1,8 @@
package org.javawebstack.jobs.api.controller;
-import org.javawebstack.httpserver.router.annotation.PathPrefix;
-import org.javawebstack.httpserver.router.annotation.With;
-import org.javawebstack.httpserver.router.annotation.verbs.Get;
+import org.javawebstack.http.router.router.annotation.PathPrefix;
+import org.javawebstack.http.router.router.annotation.With;
+import org.javawebstack.http.router.router.annotation.verbs.Get;
import org.javawebstack.jobs.Job;
import org.javawebstack.jobs.Jobs;
import org.javawebstack.jobs.api.response.Response;
diff --git a/src/main/java/org/javawebstack/jobs/api/controller/WorkerController.java b/src/main/java/org/javawebstack/jobs/api/controller/WorkerController.java
index b71a498..f81e2ad 100644
--- a/src/main/java/org/javawebstack/jobs/api/controller/WorkerController.java
+++ b/src/main/java/org/javawebstack/jobs/api/controller/WorkerController.java
@@ -1,8 +1,8 @@
package org.javawebstack.jobs.api.controller;
-import org.javawebstack.httpserver.router.annotation.PathPrefix;
-import org.javawebstack.httpserver.router.annotation.With;
-import org.javawebstack.httpserver.router.annotation.verbs.Get;
+import org.javawebstack.http.router.router.annotation.PathPrefix;
+import org.javawebstack.http.router.router.annotation.With;
+import org.javawebstack.http.router.router.annotation.verbs.Get;
import org.javawebstack.jobs.Jobs;
import org.javawebstack.jobs.api.response.Response;
diff --git a/src/main/java/org/javawebstack/jobs/api/middleware/AuthMiddleware.java b/src/main/java/org/javawebstack/jobs/api/middleware/AuthMiddleware.java
index 03face5..16918ef 100644
--- a/src/main/java/org/javawebstack/jobs/api/middleware/AuthMiddleware.java
+++ b/src/main/java/org/javawebstack/jobs/api/middleware/AuthMiddleware.java
@@ -1,8 +1,8 @@
package org.javawebstack.jobs.api.middleware;
import lombok.AllArgsConstructor;
-import org.javawebstack.httpserver.Exchange;
-import org.javawebstack.httpserver.handler.RequestHandler;
+import org.javawebstack.http.router.Exchange;
+import org.javawebstack.http.router.handler.RequestHandler;
import org.javawebstack.jobs.api.JobApi;
import org.javawebstack.jobs.api.response.Response;
diff --git a/src/main/java/org/javawebstack/jobs/api/middleware/ResponseMiddleware.java b/src/main/java/org/javawebstack/jobs/api/middleware/ResponseMiddleware.java
index 05ee005..26f4afd 100644
--- a/src/main/java/org/javawebstack/jobs/api/middleware/ResponseMiddleware.java
+++ b/src/main/java/org/javawebstack/jobs/api/middleware/ResponseMiddleware.java
@@ -1,7 +1,7 @@
package org.javawebstack.jobs.api.middleware;
-import org.javawebstack.httpserver.Exchange;
-import org.javawebstack.httpserver.handler.AfterRequestHandler;
+import org.javawebstack.http.router.Exchange;
+import org.javawebstack.http.router.handler.AfterRequestHandler;
import org.javawebstack.jobs.api.response.Response;
public class ResponseMiddleware implements AfterRequestHandler {
diff --git a/src/main/java/org/javawebstack/jobs/scheduler/sql/SQLJobScheduler.java b/src/main/java/org/javawebstack/jobs/scheduler/sql/SQLJobScheduler.java
index c342118..8dc2552 100644
--- a/src/main/java/org/javawebstack/jobs/scheduler/sql/SQLJobScheduler.java
+++ b/src/main/java/org/javawebstack/jobs/scheduler/sql/SQLJobScheduler.java
@@ -4,7 +4,8 @@
import org.javawebstack.jobs.scheduler.model.JobScheduleEntry;
import org.javawebstack.jobs.util.MapBuilder;
import org.javawebstack.jobs.util.SQLUtil;
-import org.javawebstack.orm.wrapper.SQL;
+import org.javawebstack.orm.connection.pool.PooledSQL;
+import org.javawebstack.orm.connection.pool.SQLPool;
import java.sql.SQLException;
import java.time.Instant;
@@ -13,15 +14,15 @@
public class SQLJobScheduler implements JobScheduler {
- final SQL sql;
+ final SQLPool pool;
final String tablePrefix;
- public SQLJobScheduler(SQL sql, String tablePrefix) {
+ public SQLJobScheduler(SQLPool pool, String tablePrefix) {
if(tablePrefix == null)
tablePrefix = "";
- this.sql = sql;
+ this.pool = pool;
this.tablePrefix = tablePrefix;
- try {
+ try(PooledSQL sql = pool.get()) {
sql.write("CREATE TABLE IF NOT EXISTS `" + table("queued_jobs") + "` (`id` VARCHAR(36), `ord` BIGINT, `queue` VARCHAR(50) NOT NULL, `job_id` VARCHAR(36) NOT NULL, `created_at` TIMESTAMP NOT NULL, PRIMARY KEY(`id`));");
sql.write("CREATE TABLE IF NOT EXISTS `" + table("scheduled_jobs") + "` (`id` VARCHAR(36), `ord` BIGINT, `queue` VARCHAR(50) NOT NULL, `job_id` VARCHAR(36) NOT NULL, `scheduled_at` TIMESTAMP NOT NULL, `created_at` TIMESTAMP NOT NULL, PRIMARY KEY(`id`));");
} catch (SQLException e) {
@@ -34,53 +35,65 @@ private String table(String name) {
}
public void enqueue(String queue, UUID id) {
- SQLUtil.insert(sql, table("queued_jobs"), new MapBuilder()
- .set("id", UUID.randomUUID())
- .set("ord", System.currentTimeMillis())
- .set("queue", queue)
- .set("job_id", id)
- .set("created_at", Date.from(Instant.now()))
- .build()
- );
+ try(PooledSQL sql = pool.get()) {
+ SQLUtil.insert(sql, table("queued_jobs"), new MapBuilder()
+ .set("id", UUID.randomUUID())
+ .set("ord", System.currentTimeMillis())
+ .set("queue", queue)
+ .set("job_id", id)
+ .set("created_at", Date.from(Instant.now()))
+ .build()
+ );
+ }
}
public void dequeue(UUID id) {
- SQLUtil.delete(sql, table("queued_jobs"), "`id`=?", id);
+ try(PooledSQL sql = pool.get()) {
+ SQLUtil.delete(sql, table("queued_jobs"), "`id`=?", id);
+ }
}
public void schedule(String queue, Date at, UUID id) {
- SQLUtil.insert(sql, table("scheduled_jobs"), new MapBuilder()
- .set("id", UUID.randomUUID())
- .set("ord", System.currentTimeMillis())
- .set("queue", queue)
- .set("job_id", id)
- .set("scheduled_at", at)
- .set("created_at", Date.from(Instant.now()))
- .build()
- );
+ try(PooledSQL sql = pool.get()) {
+ SQLUtil.insert(sql, table("scheduled_jobs"), new MapBuilder()
+ .set("id", UUID.randomUUID())
+ .set("ord", System.currentTimeMillis())
+ .set("queue", queue)
+ .set("job_id", id)
+ .set("scheduled_at", at)
+ .set("created_at", Date.from(Instant.now()))
+ .build()
+ );
+ }
}
public synchronized UUID poll(String queue) {
- Map e = SQLUtil.select(sql, table("queued_jobs"), "`id`,`job_id`", "WHERE `queue`=? ORDER BY `ord` LIMIT 1", queue).stream().findFirst().orElse(null);
- if(e == null)
- return null;
- SQLUtil.delete(sql, table("queued_jobs"), "`id`=?", e.get("id"));
- return UUID.fromString((String) e.get("job_id"));
+ try(PooledSQL sql = pool.get()) {
+ Map e = SQLUtil.select(sql, table("queued_jobs"), "`id`,`job_id`", "WHERE `queue`=? ORDER BY `ord` LIMIT 1", queue).stream().findFirst().orElse(null);
+ if(e == null)
+ return null;
+ SQLUtil.delete(sql, table("queued_jobs"), "`id`=?", e.get("id"));
+ return UUID.fromString((String) e.get("job_id"));
+ }
}
public List processSchedule(String queue) {
List enqueued = new ArrayList<>();
- SQLUtil.select(sql, table("scheduled_jobs"), "`id`,`job_id`", "WHERE `queue`=? AND `scheduled_at`<=? ORDER BY `ord`", queue, Date.from(Instant.now())).forEach(e -> {
- UUID jobId = UUID.fromString((String) e.get("job_id"));
- enqueue(queue, jobId);
- enqueued.add(jobId);
- SQLUtil.delete(sql, table("scheduled_jobs"), "`id`=?", e.get("id"));
- });
+ try(PooledSQL sql = pool.get()) {
+ SQLUtil.select(sql, table("scheduled_jobs"), "`id`,`job_id`", "WHERE `queue`=? AND `scheduled_at`<=? ORDER BY `ord`", queue, Date.from(Instant.now())).forEach(e -> {
+ UUID jobId = UUID.fromString((String) e.get("job_id"));
+ enqueue(queue, jobId);
+ enqueued.add(jobId);
+ SQLUtil.delete(sql, table("scheduled_jobs"), "`id`=?", e.get("id"));
+ });
+ }
return enqueued;
}
public List getScheduleEntries(String queue) {
- return SQLUtil.select(sql, table("scheduled_jobs"), "`job_id`,`scheduled_at`", "WHERE `queue`=? ORDER BY `ord`", queue).stream().map(this::buildScheduleEntry).collect(Collectors.toList());
+ try(PooledSQL sql = pool.get()) {
+ return SQLUtil.select(sql, table("scheduled_jobs"), "`job_id`,`scheduled_at`", "WHERE `queue`=? ORDER BY `ord`", queue).stream().map(this::buildScheduleEntry).collect(Collectors.toList());
+ }
}
public List getScheduleEntries(List jobIds) {
@@ -88,11 +101,15 @@ public List getScheduleEntries(List jobIds) {
return Collections.emptyList();
String whereIn = jobIds.stream().map(UUID::toString).map(s -> "\"" + s + "\"").collect(Collectors.joining(","));
- return SQLUtil.select(sql, table("scheduled_jobs"), "`job_id`,`scheduled_at`", "WHERE `job_id` IN (" + whereIn + ") ORDER BY `ord`").stream().map(this::buildScheduleEntry).collect(Collectors.toList());
+ try(PooledSQL sql = pool.get()) {
+ return SQLUtil.select(sql, table("scheduled_jobs"), "`job_id`,`scheduled_at`", "WHERE `job_id` IN (" + whereIn + ") ORDER BY `ord`").stream().map(this::buildScheduleEntry).collect(Collectors.toList());
+ }
}
public List getQueueEntries(String queue) {
- return SQLUtil.select(sql, table("queued_jobs"), "`job_id`", "WHERE `queue`=? ORDER BY `ord`", queue).stream().map(e -> UUID.fromString((String) e.get("job_id"))).collect(Collectors.toList());
+ try(PooledSQL sql = pool.get()) {
+ return SQLUtil.select(sql, table("queued_jobs"), "`job_id`", "WHERE `queue`=? ORDER BY `ord`", queue).stream().map(e -> UUID.fromString((String) e.get("job_id"))).collect(Collectors.toList());
+ }
}
private JobScheduleEntry buildScheduleEntry(Map values) {
diff --git a/src/main/java/org/javawebstack/jobs/standalone/StandaloneOptions.java b/src/main/java/org/javawebstack/jobs/standalone/StandaloneOptions.java
index e3b36ab..2e747df 100644
--- a/src/main/java/org/javawebstack/jobs/standalone/StandaloneOptions.java
+++ b/src/main/java/org/javawebstack/jobs/standalone/StandaloneOptions.java
@@ -9,7 +9,9 @@
import org.javawebstack.jobs.storage.JobStorage;
import org.javawebstack.jobs.storage.inmemory.InMemoryJobStorage;
import org.javawebstack.jobs.storage.sql.SQLJobStorage;
-import org.javawebstack.orm.wrapper.MySQL;
+import org.javawebstack.orm.connection.MySQL;
+import org.javawebstack.orm.connection.pool.MinMaxScaler;
+import org.javawebstack.orm.connection.pool.SQLPool;
import java.util.HashMap;
import java.util.Locale;
@@ -82,14 +84,14 @@ public boolean isEnabled(String key, boolean orElse) {
}
}
- private MySQL getMySQL() {
- return new MySQL(
+ private SQLPool getMySQL() {
+ return new SQLPool(new MinMaxScaler(1,1), () -> new MySQL(
get("db.host", "127.0.0.1"),
getInt("db.port", 3306),
get("db.name", "jobs"),
get("db.username", "jobs"),
get("db.password", "")
- );
+ ));
}
public JobSerializer getSerializer() {
diff --git a/src/main/java/org/javawebstack/jobs/storage/model/PaginationQuery.java b/src/main/java/org/javawebstack/jobs/storage/model/PaginationQuery.java
index a66dff2..10523b1 100644
--- a/src/main/java/org/javawebstack/jobs/storage/model/PaginationQuery.java
+++ b/src/main/java/org/javawebstack/jobs/storage/model/PaginationQuery.java
@@ -1,6 +1,7 @@
package org.javawebstack.jobs.storage.model;
-import org.javawebstack.httpserver.Exchange;
+
+import org.javawebstack.http.router.Exchange;
public abstract class PaginationQuery> {
diff --git a/src/main/java/org/javawebstack/jobs/storage/sql/SQLJobStorage.java b/src/main/java/org/javawebstack/jobs/storage/sql/SQLJobStorage.java
index 631c0a7..39889f8 100644
--- a/src/main/java/org/javawebstack/jobs/storage/sql/SQLJobStorage.java
+++ b/src/main/java/org/javawebstack/jobs/storage/sql/SQLJobStorage.java
@@ -7,7 +7,8 @@
import org.javawebstack.jobs.storage.JobStorage;
import org.javawebstack.jobs.util.MapBuilder;
import org.javawebstack.jobs.util.SQLUtil;
-import org.javawebstack.orm.wrapper.SQL;
+import org.javawebstack.orm.connection.pool.PooledSQL;
+import org.javawebstack.orm.connection.pool.SQLPool;
import java.sql.SQLException;
import java.time.Instant;
@@ -18,15 +19,15 @@
public class SQLJobStorage implements JobStorage {
- final SQL sql;
+ final SQLPool pool;
final String tablePrefix;
- public SQLJobStorage(SQL sql, String tablePrefix) {
+ public SQLJobStorage(SQLPool pool, String tablePrefix) {
if(tablePrefix == null)
tablePrefix = "";
- this.sql = sql;
+ this.pool = pool;
this.tablePrefix = tablePrefix;
- try {
+ try(PooledSQL sql = pool.get()) {
sql.write("CREATE TABLE IF NOT EXISTS `" + table("jobs") + "` (`id` VARCHAR(36), `ord` BIGINT, `status` ENUM('CREATED', 'SCHEDULED', 'ENQUEUED', 'PROCESSING', 'SUCCESS', 'FAILED', 'DELETED'), `type` VARCHAR(100) NOT NULL, `payload` LONGTEXT NOT NULL, `created_at` TIMESTAMP NOT NULL, PRIMARY KEY(`id`));");
sql.write("CREATE TABLE IF NOT EXISTS `" + table("job_events") + "` (`id` VARCHAR(36), `ord` BIGINT, `job_id` VARCHAR(36) NOT NULL, `type` ENUM('SCHEDULED','ENQUEUED','PROCESSING','FAILED','SUCCESS') NOT NULL, `created_at` TIMESTAMP NOT NULL, PRIMARY KEY(`id`));");
sql.write("CREATE TABLE IF NOT EXISTS `" + table("job_log_entries") + "` (`id` VARCHAR(36), `ord` BIGINT, `event_id` VARCHAR(36) NOT NULL,`level` ENUM('INFO','WARNING','ERROR') NOT NULL, `message` LONGTEXT NOT NULL, `created_at` TIMESTAMP NOT NULL, PRIMARY KEY(`id`));");
@@ -44,15 +45,17 @@ private String table(String name) {
public void createJob(JobInfo info, String payload) {
JobStorage.super.createJob(info, payload);
- SQLUtil.insert(sql, table("jobs"), new MapBuilder()
- .set("id", info.getId())
- .set("ord", System.currentTimeMillis())
- .set("status", info.getStatus())
- .set("type", info.getType())
- .set("payload", payload)
- .set("created_at", info.getCreatedAt())
- .build()
- );
+ try(PooledSQL sql = pool.get()) {
+ SQLUtil.insert(sql, table("jobs"), new MapBuilder()
+ .set("id", info.getId())
+ .set("ord", System.currentTimeMillis())
+ .set("status", info.getStatus())
+ .set("type", info.getType())
+ .set("payload", payload)
+ .set("created_at", info.getCreatedAt())
+ .build()
+ );
+ }
}
public boolean createRecurringJob(RecurringJobInfo info) {
@@ -62,58 +65,72 @@ public boolean createRecurringJob(RecurringJobInfo info) {
if (recurringJobs.stream().anyMatch(r -> r.getPayload().equals(info.getPayload()) && r.getCron().equals(info.getCron())))
return false;
- SQLUtil.insert(sql, table("recurring_jobs"), new MapBuilder()
- .set("id", info.getId())
- .set("last_job_id", info.getLastJobId())
- .set("queue", info.getQueue())
- .set("payload", info.getPayload())
- .set("type", info.getType())
- .set("ord", System.currentTimeMillis())
- .set("cron_expression", info.getCron().serialize())
- .set("created_at", new Date())
- .set("last_execution_at", info.getLastExecutionAt())
- .build()
- );
+ try(PooledSQL sql = pool.get()) {
+ SQLUtil.insert(sql, table("recurring_jobs"), new MapBuilder()
+ .set("id", info.getId())
+ .set("last_job_id", info.getLastJobId())
+ .set("queue", info.getQueue())
+ .set("payload", info.getPayload())
+ .set("type", info.getType())
+ .set("ord", System.currentTimeMillis())
+ .set("cron_expression", info.getCron().serialize())
+ .set("created_at", new Date())
+ .set("last_execution_at", info.getLastExecutionAt())
+ .build()
+ );
+ }
return true;
}
public JobInfo getJob(UUID id) {
- List