From 9e36e8dbb4591391b5bac0cabdfe8ba5340e62ae Mon Sep 17 00:00:00 2001 From: Yunyu Lin Date: Wed, 24 Sep 2025 12:42:21 -0400 Subject: [PATCH 1/7] Migrate to Vert.x 4 and JDK 17 --- pom.xml | 20 +++- .../yunyulin/speechdrop/Broadcaster.java | 14 +-- .../yunyulin/speechdrop/PurgeTask.java | 10 +- .../speechdrop/SpeechDropApplication.java | 111 +++++++++--------- .../speechdrop/SpeechDropVerticle.java | 36 ++++-- .../speechdrop/handlers/IndexHandler.java | 16 +-- .../speechdrop/handlers/RoomHandler.java | 49 ++++---- .../yunyulin/speechdrop/room/Room.java | 87 +++++++------- 8 files changed, 190 insertions(+), 153 deletions(-) diff --git a/pom.xml b/pom.xml index ee69d17..68275b8 100644 --- a/pom.xml +++ b/pom.xml @@ -12,10 +12,12 @@ frontend edu.vanderbilt.yunyulin.speechdrop.SpeechDropVerticle UTF-8 - 1.8 - 1.8 - 1.8 - 1.8 + 17 + 17 + 17 + 17 + 17 + 4.5.9 @@ -89,7 +91,7 @@ io.vertx vertx-stack-depchain - 3.9.16 + ${vertx.version} pom import @@ -105,6 +107,14 @@ io.vertx vertx-web + + com.fasterxml.jackson.core + jackson-databind + + + com.fasterxml.jackson.core + jackson-annotations + com.google.guava guava diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java index 8cd6cda..f340f9c 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java @@ -7,7 +7,7 @@ import io.vertx.core.json.JsonObject; import io.vertx.ext.bridge.BridgeEventType; import io.vertx.ext.bridge.PermittedOptions; -import io.vertx.ext.web.handler.sockjs.BridgeOptions; +import io.vertx.ext.web.handler.sockjs.SockJSBridgeOptions; import io.vertx.ext.web.handler.sockjs.SockJSHandler; import static edu.vanderbilt.yunyulin.speechdrop.SpeechDropApplication.LOGGER; @@ -21,17 +21,17 @@ public class Broadcaster { public Broadcaster(Vertx vertx, RoomHandler roomHandler) { this.vertx = vertx; this.sockJSHandler = SockJSHandler.create(vertx); - sockJSHandler.bridge(new BridgeOptions().addOutboundPermitted( - new PermittedOptions().setAddressRegex("speechdrop\\.room\\..+") - ), be -> { + sockJSHandler.bridge(new SockJSBridgeOptions().addOutboundPermitted( + new PermittedOptions().setAddressRegex("speechdrop\\.room\\..+") + ), be -> { if (be.type() == BridgeEventType.REGISTER) { String address = be.getRawMessage().getString("address"); String roomId = getRoomId(address); if (roomId != null && roomHandler.roomExists(roomId)) { - roomHandler.getRoom(roomId).getIndex().setHandler(ar -> { + roomHandler.getRoom(roomId).getIndex().onComplete(ar -> { // Copies envelope structure from EventBusBridgeImpl##deliverMessage - be.socket().write(Buffer.buffer(new JsonObject() - .put("type", "rec") + be.socket().write(Buffer.buffer(new JsonObject() + .put("type", "rec") .put("address", address) .put("body", ar.result()) .encode() diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java index 57cce74..580b0d7 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java @@ -25,11 +25,11 @@ public void schedule() { .filter(el -> System.currentTimeMillis() - el.getValue().ctime > purgeIntervalInSeconds * 1000) .map(Map.Entry::getKey) .collect(Collectors.toList()); - CompositeFuture.all(toRemove.stream().map(roomHandler::queueRoomDeletion).collect(Collectors.toList())) - .setHandler(res -> { - roomHandler.writeRooms(); - LOGGER.info("Purged " + toRemove.size() + " rooms"); - }); + CompositeFuture.all(toRemove.stream().map(roomHandler::queueRoomDeletion).collect(Collectors.toList())) + .onComplete(res -> { + roomHandler.writeRooms(); + LOGGER.info("Purged " + toRemove.size() + " rooms"); + }); }; vertx.setPeriodic(3 * 60 * 60 * 1000, runPurge); runPurge.handle(0L); diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java index f46f0f4..ad427f7 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java @@ -5,10 +5,10 @@ import io.vertx.core.Handler; import io.vertx.core.Vertx; import io.vertx.core.buffer.Buffer; -import io.vertx.core.json.JsonArray; -import io.vertx.core.json.JsonObject; -import io.vertx.core.logging.Logger; -import io.vertx.core.logging.LoggerFactory; +import io.vertx.core.json.JsonArray; +import io.vertx.core.json.JsonObject; +import io.vertx.core.impl.logging.Logger; +import io.vertx.core.impl.logging.LoggerFactory; import io.vertx.ext.web.Router; import io.vertx.ext.web.RoutingContext; import io.vertx.ext.web.handler.*; @@ -88,12 +88,11 @@ public void mount(Router router) throws IOException { Files.createDirectories(BASE_PATH.toPath()); new PurgeTask(roomHandler, vertx, config.getInteger("purgeIntervalInSeconds")).schedule(); - router.route().handler(BodyHandler.create().setBodyLimit(maxUploadSize).setDeleteUploadedFilesOnEnd(true)); - router.route().handler(CookieHandler.create()); - router.route().handler(SessionHandler.create( - LocalSessionStore.create(vertx, "speechdrop-sessions", 10000L) - ).setSessionTimeout(6 * 60 * 60 * 1000)); - router.route().handler(CSRFHandler.create(config.getString("csrfSecret"))); + router.route().handler(BodyHandler.create().setBodyLimit(maxUploadSize).setDeleteUploadedFilesOnEnd(true)); + router.route().handler(SessionHandler.create( + LocalSessionStore.create(vertx, "speechdrop-sessions", 10000L) + ).setSessionTimeout(6 * 60 * 60 * 1000)); + router.route().handler(CSRFHandler.create(vertx, config.getString("csrfSecret"))); router.route("/").method(GET).handler(ctx -> ctx.response().putHeader(CONTENT_TYPE, TEXT_HTML).end(mainPage) @@ -125,11 +124,11 @@ public void mount(Router router) throws IOException { sendEmptyIndex(ctx, 404); } else { Room r = roomHandler.getRoom(roomId); - r.getIndex().setHandler(ar -> - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()) - ); - } - }); + r.getIndex().onComplete(ar -> + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()) + ); + } + }); router.route("/:roomid/archive").method(GET).handler(ctx -> { String roomId = ctx.request().getParam("roomid"); @@ -137,11 +136,11 @@ public void mount(Router router) throws IOException { ctx.response().setStatusCode(404).end(); } else { Room r = roomHandler.getRoom(roomId); - r.getFiles().setHandler(ar -> { - Collection files = ar.result(); - String outFile = r.getData().name.trim() + ".zip"; - - Handler writeBufferToResponse = buf -> ctx.response() + r.getFiles().onComplete(ar -> { + Collection files = ar.result(); + String outFile = r.getData().name.trim() + ".zip"; + + Handler writeBufferToResponse = buf -> ctx.response() .setChunked(true) // See https://stackoverflow.com/a/38324508 .putHeader("Content-Disposition", "attachment; filename=\"" + outFile @@ -151,26 +150,26 @@ public void mount(Router router) throws IOException { .putHeader(CONTENT_TYPE, "application/octet-stream") .end(buf); - if (files.size() == 0) { - writeBufferToResponse.handle(Util.getEmptyZipBuffer()); - } else { - vertx.executeBlocking(fut -> { - try { - fut.complete(Util.zip(files)); - } catch (IOException e) { - fut.fail(e); - } - }, false, res -> { - if (res.succeeded()) { - writeBufferToResponse.handle(res.result()); - } else { - ctx.response().setStatusCode(500).end(); - } - }); - } - }); - } - }); + if (files.size() == 0) { + writeBufferToResponse.handle(Util.getEmptyZipBuffer()); + } else { + vertx.executeBlocking(promise -> { + try { + promise.complete(Util.zip(files)); + } catch (IOException e) { + promise.fail(e); + } + }, false).onComplete(res -> { + if (res.succeeded()) { + writeBufferToResponse.handle(res.result()); + } else { + ctx.response().setStatusCode(500).end(); + } + }); + } + }); + } + }); router.route("/:roomid/upload").method(POST).produces(APPLICATION_JSON_PRODUCES).handler(ctx -> { String roomId = ctx.request().getParam("roomid"); @@ -179,11 +178,11 @@ public void mount(Router router) throws IOException { ctx.response().setStatusCode(404).end(EMPTY_INDEX); } else { Room r = roomHandler.getRoom(roomId); - r.handleUpload(ctx).setHandler(ar -> { - if (ar.succeeded()) { - String index = ar.result(); - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON_PRODUCES).end(index); - broadcaster.publishUpdate(r.getId(), index); + r.handleUpload(ctx).onComplete(ar -> { + if (ar.succeeded()) { + String index = ar.result(); + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON_PRODUCES).end(index); + broadcaster.publishUpdate(r.getId(), index); } else { ctx.response().setStatusCode(400).end( new JsonObject().put("err", ar.cause().getMessage()).toString() @@ -204,13 +203,13 @@ public void mount(Router router) throws IOException { if (fileIndex == null) { sendEmptyIndex(ctx, 400); } else { - r.deleteFile(Integer.parseInt(fileIndex)).setHandler(ar -> { - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()); - broadcaster.publishUpdate(r.getId(), ar.result()); - }); - } - } - }); + r.deleteFile(Integer.parseInt(fileIndex)).onComplete(ar -> { + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()); + broadcaster.publishUpdate(r.getId(), ar.result()); + }); + } + } + }); router.route("/:roomid/").method(GET).handler(ctx -> { String roomId = ctx.request().getParam("roomid"); @@ -230,10 +229,10 @@ public void mount(Router router) throws IOException { redirect(ctx, "/"); } else { Room r = roomHandler.getRoom(roomId); - r.getIndex().setHandler(ar -> { - String escapedRoomName = HTML_ESCAPER.escape(r.getData().name); - JsonObject configPayload = new JsonObject() - .put("mediaUrl", mediaUrl) + r.getIndex().onComplete(ar -> { + String escapedRoomName = HTML_ESCAPER.escape(r.getData().name); + JsonObject configPayload = new JsonObject() + .put("mediaUrl", mediaUrl) .put("roomName", escapedRoomName) .put("roomId", r.getId()) .put("allowedMimes", allowedMimeTypesCsv) diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropVerticle.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropVerticle.java index 7def8b1..f0b4cb9 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropVerticle.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropVerticle.java @@ -1,28 +1,46 @@ package edu.vanderbilt.yunyulin.speechdrop; import io.vertx.core.AbstractVerticle; +import io.vertx.core.Promise; import io.vertx.core.http.HttpServer; import io.vertx.core.http.HttpServerOptions; import io.vertx.ext.web.Router; -import io.vertx.ext.web.impl.Utils; public class SpeechDropVerticle extends AbstractVerticle { public static String VERSION; @Override - public void start() throws Exception { - VERSION = Utils.readFileToString(vertx, "VERSION"); + public void start(Promise startPromise) { HttpServerOptions serverOptions = new HttpServerOptions(); serverOptions.setCompressionSupported(true); HttpServer httpServer = vertx.createHttpServer(serverOptions); Router router = Router.router(vertx); - new SpeechDropApplication(vertx, config(), - Utils.readFileToString(vertx, "main.html"), - Utils.readFileToString(vertx, "room.html"), - Utils.readFileToString(vertx, "about.html") - ).mount(router); + try { + VERSION = readFile("VERSION"); + new SpeechDropApplication(vertx, config(), + readFile("main.html"), + readFile("room.html"), + readFile("about.html") + ).mount(router); + } catch (Exception e) { + startPromise.fail(e); + return; + } - httpServer.requestHandler(router::accept).listen(config().getInteger("port"), config().getString("host")); + httpServer + .requestHandler(router) + .listen(config().getInteger("port"), config().getString("host")) + .onComplete(ar -> { + if (ar.succeeded()) { + startPromise.complete(); + } else { + startPromise.fail(ar.cause()); + } + }); + } + + private String readFile(String path) { + return vertx.fileSystem().readFileBlocking(path).toString(); } } \ No newline at end of file diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java index 73b3f1d..a0d07a1 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java @@ -78,9 +78,9 @@ public void deleteFile(int index, Handler indexHandler) { if (entries.get(index) != null) { entries.set(index, null); completedIndex = writeIndex(); - vertx.fileSystem().deleteRecursive( - new File(uploadDirectory, Integer.toString(index)).getPath(), true, null - ); + vertx.fileSystem().deleteRecursive( + new File(uploadDirectory, Integer.toString(index)).getPath(), true + ).onFailure(err -> LOGGER.error("Failed to delete uploaded file directory", err)); } else { completedIndex = getIndexString(); } @@ -113,11 +113,11 @@ public String getIndexString() { private String writeIndex() { String indexString = getIndexString(); - vertx.fileSystem().writeFile(indexFile.getPath(), - Buffer.buffer(indexString), - null); - return indexString; - } + vertx.fileSystem().writeFile(indexFile.getPath(), + Buffer.buffer(indexString)) + .onFailure(err -> LOGGER.error("Failed to write index file", err)); + return indexString; + } private void checkLoad() { if (!loaded) throw new IllegalStateException("Index not loaded"); diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java index 960f6a7..b8beeb0 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java @@ -9,8 +9,9 @@ import edu.vanderbilt.yunyulin.speechdrop.room.Room; import edu.vanderbilt.yunyulin.speechdrop.room.RoomData; import io.vertx.core.Future; -import io.vertx.core.Vertx; -import io.vertx.core.buffer.Buffer; +import io.vertx.core.Promise; +import io.vertx.core.Vertx; +import io.vertx.core.buffer.Buffer; import lombok.Getter; import java.io.File; @@ -90,21 +91,29 @@ public Room getRoom(String id) { public void writeRooms() { try { - vertx.fileSystem().writeFile(roomsFile.getPath(), Buffer.buffer(mapper.writeValueAsString(dataStore)), null); - } catch (JsonProcessingException e) { // This should never happen - e.printStackTrace(); - } - } - - public Future queueRoomDeletion(String id) { - Future fut = Future.future(); - if (dataStore.remove(id) != null) { - roomCache.invalidate(id); - File toDelete = new File(SpeechDropApplication.BASE_PATH, id); - vertx.fileSystem().deleteRecursive(toDelete.getPath(), true, fut.completer()); - } else { - fut.complete(); - } - return fut; - } -} + vertx.fileSystem() + .writeFile(roomsFile.getPath(), Buffer.buffer(mapper.writeValueAsString(dataStore))) + .onFailure(err -> LOGGER.error("Failed to persist rooms metadata", err)); + } catch (JsonProcessingException e) { // This should never happen + e.printStackTrace(); + } + } + + public Future queueRoomDeletion(String id) { + Promise promise = Promise.promise(); + if (dataStore.remove(id) != null) { + roomCache.invalidate(id); + File toDelete = new File(SpeechDropApplication.BASE_PATH, id); + vertx.fileSystem().deleteRecursive(toDelete.getPath(), true).onComplete(ar -> { + if (ar.succeeded()) { + promise.complete(); + } else { + promise.fail(ar.cause()); + } + }); + } else { + promise.complete(); + } + return promise.future(); + } +} diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java index c61c322..e7dcebe 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java @@ -2,9 +2,10 @@ import edu.vanderbilt.yunyulin.speechdrop.SpeechDropApplication; import edu.vanderbilt.yunyulin.speechdrop.handlers.IndexHandler; -import io.vertx.core.Future; -import io.vertx.core.Handler; -import io.vertx.core.Vertx; +import io.vertx.core.Future; +import io.vertx.core.Handler; +import io.vertx.core.Promise; +import io.vertx.core.Vertx; import io.vertx.ext.web.FileUpload; import io.vertx.ext.web.RoutingContext; import lombok.Getter; @@ -42,43 +43,43 @@ private void scheduleOperation(Handler operation) { } public Future handleUpload(RoutingContext ctx) { - Future uploadFuture = Future.future(); - - Iterator itr = ctx.fileUploads().iterator(); - if (!itr.hasNext()) { - uploadFuture.fail(new Exception("no_file")); - } - - FileUpload uploadedFile = itr.next(); - Date now = new Date(); - - String mimeType = uploadedFile.contentType(); - if (!SpeechDropApplication.allowedMimeTypes.contains(mimeType)) { - uploadFuture.fail(new Exception("bad_type")); - } - if (uploadedFile.size() > SpeechDropApplication.maxUploadSize) { - uploadFuture.fail(new Exception("too_large")); - } - - scheduleOperation(index -> index.addFile(uploadedFile, now, uploadFuture::complete)); - return uploadFuture; - } - - public Future getIndex() { - Future indexFuture = Future.future(); - scheduleOperation(index -> indexFuture.complete(index.getIndexString())); - return indexFuture; - } - - public Future> getFiles() { - Future> fileFuture = Future.future(); - scheduleOperation(index -> fileFuture.complete(index.getFiles())); - return fileFuture; - } - - public Future deleteFile(int fileIndex) { - Future deleteFuture = Future.future(); - scheduleOperation(index -> index.deleteFile(fileIndex, deleteFuture::complete)); - return deleteFuture; - } -} + Promise uploadPromise = Promise.promise(); + + Iterator itr = ctx.fileUploads().iterator(); + if (!itr.hasNext()) { + uploadPromise.fail(new Exception("no_file")); + } + + FileUpload uploadedFile = itr.next(); + Date now = new Date(); + + String mimeType = uploadedFile.contentType(); + if (!SpeechDropApplication.allowedMimeTypes.contains(mimeType)) { + uploadPromise.fail(new Exception("bad_type")); + } + if (uploadedFile.size() > SpeechDropApplication.maxUploadSize) { + uploadPromise.fail(new Exception("too_large")); + } + + scheduleOperation(index -> index.addFile(uploadedFile, now, uploadPromise::complete)); + return uploadPromise.future(); + } + + public Future getIndex() { + Promise indexPromise = Promise.promise(); + scheduleOperation(index -> indexPromise.complete(index.getIndexString())); + return indexPromise.future(); + } + + public Future> getFiles() { + Promise> filesPromise = Promise.promise(); + scheduleOperation(index -> filesPromise.complete(index.getFiles())); + return filesPromise.future(); + } + + public Future deleteFile(int fileIndex) { + Promise deletePromise = Promise.promise(); + scheduleOperation(index -> index.deleteFile(fileIndex, deletePromise::complete)); + return deletePromise.future(); + } +} From 5371d0c74c749ebedf73b2e4660291f4de2b6452 Mon Sep 17 00:00:00 2001 From: Yunyu Lin Date: Wed, 24 Sep 2025 14:01:24 -0400 Subject: [PATCH 2/7] Harden Vert.x 4 migration handlers --- .../yunyulin/speechdrop/Broadcaster.java | 23 ++-- .../yunyulin/speechdrop/PurgeTask.java | 16 ++- .../speechdrop/SpeechDropApplication.java | 118 +++++++++++------- .../yunyulin/speechdrop/room/Room.java | 3 + 4 files changed, 98 insertions(+), 62 deletions(-) diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java index f340f9c..f712936 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java @@ -29,14 +29,21 @@ public Broadcaster(Vertx vertx, RoomHandler roomHandler) { String roomId = getRoomId(address); if (roomId != null && roomHandler.roomExists(roomId)) { roomHandler.getRoom(roomId).getIndex().onComplete(ar -> { - // Copies envelope structure from EventBusBridgeImpl##deliverMessage - be.socket().write(Buffer.buffer(new JsonObject() - .put("type", "rec") - .put("address", address) - .put("body", ar.result()) - .encode() - )); - }); + if (ar.succeeded()) { + // Copies envelope structure from EventBusBridgeImpl##deliverMessage + be.socket().write(Buffer.buffer(new JsonObject() + .put("type", "rec") + .put("address", address) + .put("body", ar.result()) + .encode() + )); + be.complete(true); + } else { + LOGGER.error("Failed to deliver initial index for room " + roomId, ar.cause()); + be.complete(false); + } + }); + return; } else { be.complete(false); return; diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java index 580b0d7..5e2f55e 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java @@ -28,10 +28,14 @@ public void schedule() { CompositeFuture.all(toRemove.stream().map(roomHandler::queueRoomDeletion).collect(Collectors.toList())) .onComplete(res -> { roomHandler.writeRooms(); - LOGGER.info("Purged " + toRemove.size() + " rooms"); + if (res.succeeded()) { + LOGGER.info("Purged " + toRemove.size() + " rooms"); + } else { + LOGGER.error("Failed to purge one or more rooms", res.cause()); + } }); - }; - vertx.setPeriodic(3 * 60 * 60 * 1000, runPurge); - runPurge.handle(0L); - } -} + }; + vertx.setPeriodic(3 * 60 * 60 * 1000, runPurge); + runPurge.handle(0L); + } +} diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java index ad427f7..5bbd187 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java @@ -123,25 +123,35 @@ public void mount(Router router) throws IOException { if (!roomHandler.roomExists(roomId)) { sendEmptyIndex(ctx, 404); } else { - Room r = roomHandler.getRoom(roomId); - r.getIndex().onComplete(ar -> - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()) - ); + Room r = roomHandler.getRoom(roomId); + r.getIndex().onComplete(ar -> { + if (ar.succeeded()) { + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()); + } else { + LOGGER.error("Failed to fetch index for room " + roomId, ar.cause()); + sendEmptyIndex(ctx, 500); + } + }); } }); - - router.route("/:roomid/archive").method(GET).handler(ctx -> { - String roomId = ctx.request().getParam("roomid"); - if (!roomHandler.roomExists(roomId)) { - ctx.response().setStatusCode(404).end(); - } else { - Room r = roomHandler.getRoom(roomId); + + router.route("/:roomid/archive").method(GET).handler(ctx -> { + String roomId = ctx.request().getParam("roomid"); + if (!roomHandler.roomExists(roomId)) { + ctx.response().setStatusCode(404).end(); + } else { + Room r = roomHandler.getRoom(roomId); r.getFiles().onComplete(ar -> { + if (ar.failed()) { + LOGGER.error("Failed to gather files for room " + roomId, ar.cause()); + ctx.response().setStatusCode(500).end(); + return; + } Collection files = ar.result(); String outFile = r.getData().name.trim() + ".zip"; Handler writeBufferToResponse = buf -> ctx.response() - .setChunked(true) + .setChunked(true) // See https://stackoverflow.com/a/38324508 .putHeader("Content-Disposition", "attachment; filename=\"" + outFile .replace("\\", "\\\\") @@ -164,6 +174,7 @@ public void mount(Router router) throws IOException { writeBufferToResponse.handle(res.result()); } else { ctx.response().setStatusCode(500).end(); + LOGGER.error("Failed to archive files for room " + roomId, res.cause()); } }); } @@ -177,35 +188,41 @@ public void mount(Router router) throws IOException { LOGGER.warn("(Upload) Nonexist " + roomId); ctx.response().setStatusCode(404).end(EMPTY_INDEX); } else { - Room r = roomHandler.getRoom(roomId); + Room r = roomHandler.getRoom(roomId); r.handleUpload(ctx).onComplete(ar -> { if (ar.succeeded()) { String index = ar.result(); ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON_PRODUCES).end(index); broadcaster.publishUpdate(r.getId(), index); - } else { - ctx.response().setStatusCode(400).end( - new JsonObject().put("err", ar.cause().getMessage()).toString() - ); - } - }); - } - }); - - router.route("/:roomid/delete").method(POST).produces(APPLICATION_JSON_PRODUCES).handler(ctx -> { - String roomId = ctx.request().getParam("roomid"); - if (!roomHandler.roomExists(roomId)) { - LOGGER.warn("(Upload) Nonexist " + roomId); - sendEmptyIndex(ctx, 404); - } else { - Room r = roomHandler.getRoom(roomId); - String fileIndex = ctx.request().getFormAttribute("fileIndex"); - if (fileIndex == null) { - sendEmptyIndex(ctx, 400); - } else { + } else { + LOGGER.warn("[" + roomId + "] Upload failed", ar.cause()); + ctx.response().setStatusCode(400).end( + new JsonObject().put("err", ar.cause().getMessage()).toString() + ); + } + }); + } + }); + + router.route("/:roomid/delete").method(POST).produces(APPLICATION_JSON_PRODUCES).handler(ctx -> { + String roomId = ctx.request().getParam("roomid"); + if (!roomHandler.roomExists(roomId)) { + LOGGER.warn("(Upload) Nonexist " + roomId); + sendEmptyIndex(ctx, 404); + } else { + Room r = roomHandler.getRoom(roomId); + String fileIndex = ctx.request().getFormAttribute("fileIndex"); + if (fileIndex == null) { + sendEmptyIndex(ctx, 400); + } else { r.deleteFile(Integer.parseInt(fileIndex)).onComplete(ar -> { - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()); - broadcaster.publishUpdate(r.getId(), ar.result()); + if (ar.succeeded()) { + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()); + broadcaster.publishUpdate(r.getId(), ar.result()); + } else { + LOGGER.error("Failed to delete file in room " + roomId, ar.cause()); + sendEmptyIndex(ctx, 500); + } }); } } @@ -228,22 +245,27 @@ public void mount(Router router) throws IOException { if (!roomHandler.roomExists(roomId)) { redirect(ctx, "/"); } else { - Room r = roomHandler.getRoom(roomId); + Room r = roomHandler.getRoom(roomId); r.getIndex().onComplete(ar -> { + if (ar.failed()) { + LOGGER.error("Failed to render room " + roomId, ar.cause()); + ctx.response().setStatusCode(500).end(); + return; + } String escapedRoomName = HTML_ESCAPER.escape(r.getData().name); JsonObject configPayload = new JsonObject() .put("mediaUrl", mediaUrl) - .put("roomName", escapedRoomName) - .put("roomId", r.getId()) - .put("allowedMimes", allowedMimeTypesCsv) - .put("initialFiles", new JsonArray(ar.result())); - - ctx.response().putHeader(CONTENT_TYPE, TEXT_HTML).end( - roomTemplate - .replace("<%= ROOM_CONFIG %>", configPayload.encode()) - ); - }); - } - }); + .put("roomName", escapedRoomName) + .put("roomId", r.getId()) + .put("allowedMimes", allowedMimeTypesCsv) + .put("initialFiles", new JsonArray(ar.result())); + + ctx.response().putHeader(CONTENT_TYPE, TEXT_HTML).end( + roomTemplate + .replace("<%= ROOM_CONFIG %>", configPayload.encode()) + ); + }); + } + }); } } diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java index e7dcebe..dc0bcd7 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java @@ -48,6 +48,7 @@ public Future handleUpload(RoutingContext ctx) { Iterator itr = ctx.fileUploads().iterator(); if (!itr.hasNext()) { uploadPromise.fail(new Exception("no_file")); + return uploadPromise.future(); } FileUpload uploadedFile = itr.next(); @@ -56,9 +57,11 @@ public Future handleUpload(RoutingContext ctx) { String mimeType = uploadedFile.contentType(); if (!SpeechDropApplication.allowedMimeTypes.contains(mimeType)) { uploadPromise.fail(new Exception("bad_type")); + return uploadPromise.future(); } if (uploadedFile.size() > SpeechDropApplication.maxUploadSize) { uploadPromise.fail(new Exception("too_large")); + return uploadPromise.future(); } scheduleOperation(index -> index.addFile(uploadedFile, now, uploadPromise::complete)); From 3fa825b5300b5748498b82bbe93bed7fb1d78764 Mon Sep 17 00:00:00 2001 From: Yunyu Lin Date: Wed, 24 Sep 2025 16:58:59 -0400 Subject: [PATCH 3/7] Remove Vert.x 4 failure wrappers --- .../yunyulin/speechdrop/Broadcaster.java | 20 +++---- .../yunyulin/speechdrop/PurgeTask.java | 6 +-- .../speechdrop/SpeechDropApplication.java | 48 +++++------------ .../speechdrop/handlers/IndexHandler.java | 53 +++++++++---------- .../speechdrop/handlers/RoomHandler.java | 13 ++--- .../yunyulin/speechdrop/room/Room.java | 3 -- 6 files changed, 50 insertions(+), 93 deletions(-) diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java index f712936..0957b75 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java @@ -29,19 +29,13 @@ public Broadcaster(Vertx vertx, RoomHandler roomHandler) { String roomId = getRoomId(address); if (roomId != null && roomHandler.roomExists(roomId)) { roomHandler.getRoom(roomId).getIndex().onComplete(ar -> { - if (ar.succeeded()) { - // Copies envelope structure from EventBusBridgeImpl##deliverMessage - be.socket().write(Buffer.buffer(new JsonObject() - .put("type", "rec") - .put("address", address) - .put("body", ar.result()) - .encode() - )); - be.complete(true); - } else { - LOGGER.error("Failed to deliver initial index for room " + roomId, ar.cause()); - be.complete(false); - } + // Copies envelope structure from EventBusBridgeImpl##deliverMessage + be.socket().write(Buffer.buffer(new JsonObject() + .put("type", "rec") + .put("address", address) + .put("body", ar.result()) + .encode() + )); }); return; } else { diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java index 5e2f55e..c40f7f3 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java @@ -28,11 +28,7 @@ public void schedule() { CompositeFuture.all(toRemove.stream().map(roomHandler::queueRoomDeletion).collect(Collectors.toList())) .onComplete(res -> { roomHandler.writeRooms(); - if (res.succeeded()) { - LOGGER.info("Purged " + toRemove.size() + " rooms"); - } else { - LOGGER.error("Failed to purge one or more rooms", res.cause()); - } + LOGGER.info("Purged " + toRemove.size() + " rooms"); }); }; vertx.setPeriodic(3 * 60 * 60 * 1000, runPurge); diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java index 5bbd187..aa43a47 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java @@ -118,20 +118,15 @@ public void mount(Router router) throws IOException { } }); - router.route("/:roomid/index").method(GET).produces(APPLICATION_JSON_PRODUCES).handler(ctx -> { - String roomId = ctx.request().getParam("roomid"); - if (!roomHandler.roomExists(roomId)) { - sendEmptyIndex(ctx, 404); - } else { + router.route("/:roomid/index").method(GET).produces(APPLICATION_JSON_PRODUCES).handler(ctx -> { + String roomId = ctx.request().getParam("roomid"); + if (!roomHandler.roomExists(roomId)) { + sendEmptyIndex(ctx, 404); + } else { Room r = roomHandler.getRoom(roomId); - r.getIndex().onComplete(ar -> { - if (ar.succeeded()) { - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()); - } else { - LOGGER.error("Failed to fetch index for room " + roomId, ar.cause()); - sendEmptyIndex(ctx, 500); - } - }); + r.getIndex().onComplete(ar -> + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()) + ); } }); @@ -142,11 +137,6 @@ public void mount(Router router) throws IOException { } else { Room r = roomHandler.getRoom(roomId); r.getFiles().onComplete(ar -> { - if (ar.failed()) { - LOGGER.error("Failed to gather files for room " + roomId, ar.cause()); - ctx.response().setStatusCode(500).end(); - return; - } Collection files = ar.result(); String outFile = r.getData().name.trim() + ".zip"; @@ -174,7 +164,6 @@ public void mount(Router router) throws IOException { writeBufferToResponse.handle(res.result()); } else { ctx.response().setStatusCode(500).end(); - LOGGER.error("Failed to archive files for room " + roomId, res.cause()); } }); } @@ -195,7 +184,6 @@ public void mount(Router router) throws IOException { ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON_PRODUCES).end(index); broadcaster.publishUpdate(r.getId(), index); } else { - LOGGER.warn("[" + roomId + "] Upload failed", ar.cause()); ctx.response().setStatusCode(400).end( new JsonObject().put("err", ar.cause().getMessage()).toString() ); @@ -216,13 +204,8 @@ public void mount(Router router) throws IOException { sendEmptyIndex(ctx, 400); } else { r.deleteFile(Integer.parseInt(fileIndex)).onComplete(ar -> { - if (ar.succeeded()) { - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()); - broadcaster.publishUpdate(r.getId(), ar.result()); - } else { - LOGGER.error("Failed to delete file in room " + roomId, ar.cause()); - sendEmptyIndex(ctx, 500); - } + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()); + broadcaster.publishUpdate(r.getId(), ar.result()); }); } } @@ -242,16 +225,11 @@ public void mount(Router router) throws IOException { } router.route("/:roomid").method(GET).handler(ctx -> { String roomId = ctx.request().getParam("roomid"); - if (!roomHandler.roomExists(roomId)) { - redirect(ctx, "/"); - } else { + if (!roomHandler.roomExists(roomId)) { + redirect(ctx, "/"); + } else { Room r = roomHandler.getRoom(roomId); r.getIndex().onComplete(ar -> { - if (ar.failed()) { - LOGGER.error("Failed to render room " + roomId, ar.cause()); - ctx.response().setStatusCode(500).end(); - return; - } String escapedRoomName = HTML_ESCAPER.escape(r.getData().name); JsonObject configPayload = new JsonObject() .put("mediaUrl", mediaUrl) diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java index a0d07a1..e207e84 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java @@ -56,36 +56,36 @@ public void addFile(FileUpload uploadedFile, Date creationTime, Handler LOGGER.info("[" + uploadDirectory.getName() + "] Processing upload " + uploadedFile.fileName() + " (" + uploadedFile.size() + ")"); - vertx.fileSystem().mkdir(uploadDirectory.getPath(), uploadDirRes -> { - File destDir = new File(uploadDirectory, Integer.toString(entries.size())); - vertx.fileSystem().mkdir(destDir.getPath(), mkdirRes -> { - File dest = new File(destDir, uploadedFile.fileName()); - vertx.fileSystem().copy(uploadedFile.uploadedFileName(), dest.getPath(), - res -> { - entries.add(new FileEntry(dest.getName(), creationTime.getTime())); - indexHandler.handle(writeIndex()); - } - ); - }); - }); - } + vertx.fileSystem().mkdir(uploadDirectory.getPath(), uploadDirRes -> { + File destDir = new File(uploadDirectory, Integer.toString(entries.size())); + vertx.fileSystem().mkdir(destDir.getPath(), mkdirRes -> { + File dest = new File(destDir, uploadedFile.fileName()); + vertx.fileSystem().copy(uploadedFile.uploadedFileName(), dest.getPath(), + res -> { + entries.add(new FileEntry(dest.getName(), creationTime.getTime())); + indexHandler.handle(writeIndex()); + } + ); + }); + }); + } public void deleteFile(int index, Handler indexHandler) { checkLoad(); LOGGER.info("[" + uploadDirectory.getName() + "] Processing delete for index " + index); - String completedIndex; - if (entries.get(index) != null) { - entries.set(index, null); - completedIndex = writeIndex(); + String completedIndex; + if (entries.get(index) != null) { + entries.set(index, null); + completedIndex = writeIndex(); vertx.fileSystem().deleteRecursive( new File(uploadDirectory, Integer.toString(index)).getPath(), true - ).onFailure(err -> LOGGER.error("Failed to delete uploaded file directory", err)); - } else { - completedIndex = getIndexString(); - } - indexHandler.handle(completedIndex); - } + ); + } else { + completedIndex = getIndexString(); + } + indexHandler.handle(completedIndex); + } public Collection getFiles() { checkLoad(); @@ -111,11 +111,10 @@ public String getIndexString() { } } - private String writeIndex() { - String indexString = getIndexString(); + private String writeIndex() { + String indexString = getIndexString(); vertx.fileSystem().writeFile(indexFile.getPath(), - Buffer.buffer(indexString)) - .onFailure(err -> LOGGER.error("Failed to write index file", err)); + Buffer.buffer(indexString)); return indexString; } diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java index b8beeb0..82d02c4 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java @@ -90,10 +90,9 @@ public Room getRoom(String id) { } public void writeRooms() { - try { + try { vertx.fileSystem() - .writeFile(roomsFile.getPath(), Buffer.buffer(mapper.writeValueAsString(dataStore))) - .onFailure(err -> LOGGER.error("Failed to persist rooms metadata", err)); + .writeFile(roomsFile.getPath(), Buffer.buffer(mapper.writeValueAsString(dataStore))); } catch (JsonProcessingException e) { // This should never happen e.printStackTrace(); } @@ -104,13 +103,7 @@ public Future queueRoomDeletion(String id) { if (dataStore.remove(id) != null) { roomCache.invalidate(id); File toDelete = new File(SpeechDropApplication.BASE_PATH, id); - vertx.fileSystem().deleteRecursive(toDelete.getPath(), true).onComplete(ar -> { - if (ar.succeeded()) { - promise.complete(); - } else { - promise.fail(ar.cause()); - } - }); + vertx.fileSystem().deleteRecursive(toDelete.getPath(), true).onComplete(promise); } else { promise.complete(); } diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java index dc0bcd7..e7dcebe 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java @@ -48,7 +48,6 @@ public Future handleUpload(RoutingContext ctx) { Iterator itr = ctx.fileUploads().iterator(); if (!itr.hasNext()) { uploadPromise.fail(new Exception("no_file")); - return uploadPromise.future(); } FileUpload uploadedFile = itr.next(); @@ -57,11 +56,9 @@ public Future handleUpload(RoutingContext ctx) { String mimeType = uploadedFile.contentType(); if (!SpeechDropApplication.allowedMimeTypes.contains(mimeType)) { uploadPromise.fail(new Exception("bad_type")); - return uploadPromise.future(); } if (uploadedFile.size() > SpeechDropApplication.maxUploadSize) { uploadPromise.fail(new Exception("too_large")); - return uploadPromise.future(); } scheduleOperation(index -> index.addFile(uploadedFile, now, uploadPromise::complete)); From 212dd30caf4db3ad181352cea79685639ce27511 Mon Sep 17 00:00:00 2001 From: Yunyu Lin Date: Wed, 24 Sep 2025 20:03:38 -0400 Subject: [PATCH 4/7] Use JDK 17 in CI build --- .github/workflows/jar-build.yml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/.github/workflows/jar-build.yml b/.github/workflows/jar-build.yml index 6f37414..cca9c62 100644 --- a/.github/workflows/jar-build.yml +++ b/.github/workflows/jar-build.yml @@ -15,11 +15,11 @@ jobs: - name: Check out repository uses: actions/checkout@v4 - - name: Set up Temurin JDK 8 + - name: Set up Temurin JDK 17 uses: actions/setup-java@v4 with: distribution: temurin - java-version: '8' + java-version: '17' cache: maven - name: Build JAR From df1f9716c1332e18a0360648a3684e7e25ba20fe Mon Sep 17 00:00:00 2001 From: Yunyu Lin Date: Wed, 24 Sep 2025 20:12:04 -0400 Subject: [PATCH 5/7] Return promises from room operations --- .../yunyulin/speechdrop/Broadcaster.java | 2 +- .../yunyulin/speechdrop/PurgeTask.java | 16 ++++++++++------ .../speechdrop/SpeechDropApplication.java | 10 +++++----- .../speechdrop/handlers/RoomHandler.java | 5 ++--- .../yunyulin/speechdrop/room/Room.java | 17 ++++++++--------- 5 files changed, 26 insertions(+), 24 deletions(-) diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java index 0957b75..a8993f8 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java @@ -28,7 +28,7 @@ public Broadcaster(Vertx vertx, RoomHandler roomHandler) { String address = be.getRawMessage().getString("address"); String roomId = getRoomId(address); if (roomId != null && roomHandler.roomExists(roomId)) { - roomHandler.getRoom(roomId).getIndex().onComplete(ar -> { + roomHandler.getRoom(roomId).getIndex().future().onComplete(ar -> { // Copies envelope structure from EventBusBridgeImpl##deliverMessage be.socket().write(Buffer.buffer(new JsonObject() .put("type", "rec") diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java index c40f7f3..fb158a3 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java @@ -1,10 +1,10 @@ package edu.vanderbilt.yunyulin.speechdrop; -import edu.vanderbilt.yunyulin.speechdrop.handlers.RoomHandler; -import io.vertx.core.CompositeFuture; -import io.vertx.core.Future; -import io.vertx.core.Handler; -import io.vertx.core.Vertx; +import edu.vanderbilt.yunyulin.speechdrop.handlers.RoomHandler; +import io.vertx.core.CompositeFuture; +import io.vertx.core.Promise; +import io.vertx.core.Handler; +import io.vertx.core.Vertx; import lombok.AllArgsConstructor; import java.util.List; @@ -25,7 +25,11 @@ public void schedule() { .filter(el -> System.currentTimeMillis() - el.getValue().ctime > purgeIntervalInSeconds * 1000) .map(Map.Entry::getKey) .collect(Collectors.toList()); - CompositeFuture.all(toRemove.stream().map(roomHandler::queueRoomDeletion).collect(Collectors.toList())) + CompositeFuture.all( + toRemove.stream() + .map(roomHandler::queueRoomDeletion) + .map(Promise::future) + .collect(Collectors.toList())) .onComplete(res -> { roomHandler.writeRooms(); LOGGER.info("Purged " + toRemove.size() + " rooms"); diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java index aa43a47..0c110c8 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java @@ -124,7 +124,7 @@ public void mount(Router router) throws IOException { sendEmptyIndex(ctx, 404); } else { Room r = roomHandler.getRoom(roomId); - r.getIndex().onComplete(ar -> + r.getIndex().future().onComplete(ar -> ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()) ); } @@ -136,7 +136,7 @@ public void mount(Router router) throws IOException { ctx.response().setStatusCode(404).end(); } else { Room r = roomHandler.getRoom(roomId); - r.getFiles().onComplete(ar -> { + r.getFiles().future().onComplete(ar -> { Collection files = ar.result(); String outFile = r.getData().name.trim() + ".zip"; @@ -178,7 +178,7 @@ public void mount(Router router) throws IOException { ctx.response().setStatusCode(404).end(EMPTY_INDEX); } else { Room r = roomHandler.getRoom(roomId); - r.handleUpload(ctx).onComplete(ar -> { + r.handleUpload(ctx).future().onComplete(ar -> { if (ar.succeeded()) { String index = ar.result(); ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON_PRODUCES).end(index); @@ -203,7 +203,7 @@ public void mount(Router router) throws IOException { if (fileIndex == null) { sendEmptyIndex(ctx, 400); } else { - r.deleteFile(Integer.parseInt(fileIndex)).onComplete(ar -> { + r.deleteFile(Integer.parseInt(fileIndex)).future().onComplete(ar -> { ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()); broadcaster.publishUpdate(r.getId(), ar.result()); }); @@ -229,7 +229,7 @@ public void mount(Router router) throws IOException { redirect(ctx, "/"); } else { Room r = roomHandler.getRoom(roomId); - r.getIndex().onComplete(ar -> { + r.getIndex().future().onComplete(ar -> { String escapedRoomName = HTML_ESCAPER.escape(r.getData().name); JsonObject configPayload = new JsonObject() .put("mediaUrl", mediaUrl) diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java index 82d02c4..721598d 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java @@ -8,7 +8,6 @@ import edu.vanderbilt.yunyulin.speechdrop.SpeechDropApplication; import edu.vanderbilt.yunyulin.speechdrop.room.Room; import edu.vanderbilt.yunyulin.speechdrop.room.RoomData; -import io.vertx.core.Future; import io.vertx.core.Promise; import io.vertx.core.Vertx; import io.vertx.core.buffer.Buffer; @@ -98,7 +97,7 @@ public void writeRooms() { } } - public Future queueRoomDeletion(String id) { + public Promise queueRoomDeletion(String id) { Promise promise = Promise.promise(); if (dataStore.remove(id) != null) { roomCache.invalidate(id); @@ -107,6 +106,6 @@ public Future queueRoomDeletion(String id) { } else { promise.complete(); } - return promise.future(); + return promise; } } diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java index e7dcebe..abb4715 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java @@ -2,7 +2,6 @@ import edu.vanderbilt.yunyulin.speechdrop.SpeechDropApplication; import edu.vanderbilt.yunyulin.speechdrop.handlers.IndexHandler; -import io.vertx.core.Future; import io.vertx.core.Handler; import io.vertx.core.Promise; import io.vertx.core.Vertx; @@ -42,7 +41,7 @@ private void scheduleOperation(Handler operation) { } } - public Future handleUpload(RoutingContext ctx) { + public Promise handleUpload(RoutingContext ctx) { Promise uploadPromise = Promise.promise(); Iterator itr = ctx.fileUploads().iterator(); @@ -62,24 +61,24 @@ public Future handleUpload(RoutingContext ctx) { } scheduleOperation(index -> index.addFile(uploadedFile, now, uploadPromise::complete)); - return uploadPromise.future(); + return uploadPromise; } - public Future getIndex() { + public Promise getIndex() { Promise indexPromise = Promise.promise(); scheduleOperation(index -> indexPromise.complete(index.getIndexString())); - return indexPromise.future(); + return indexPromise; } - public Future> getFiles() { + public Promise> getFiles() { Promise> filesPromise = Promise.promise(); scheduleOperation(index -> filesPromise.complete(index.getFiles())); - return filesPromise.future(); + return filesPromise; } - public Future deleteFile(int fileIndex) { + public Promise deleteFile(int fileIndex) { Promise deletePromise = Promise.promise(); scheduleOperation(index -> index.deleteFile(fileIndex, deletePromise::complete)); - return deletePromise.future(); + return deletePromise; } } From fe128e3fe6850b2f3f5742179d8559d2aa5c5817 Mon Sep 17 00:00:00 2001 From: Yunyu Lin Date: Wed, 24 Sep 2025 20:18:09 -0400 Subject: [PATCH 6/7] Remove future chaining from room operations --- .../yunyulin/speechdrop/Broadcaster.java | 4 +- .../yunyulin/speechdrop/PurgeTask.java | 4 +- .../speechdrop/SpeechDropApplication.java | 19 +++++----- .../speechdrop/handlers/RoomHandler.java | 11 ++---- .../yunyulin/speechdrop/room/Room.java | 38 +++++++++---------- 5 files changed, 33 insertions(+), 43 deletions(-) diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java index a8993f8..d866619 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java @@ -28,12 +28,12 @@ public Broadcaster(Vertx vertx, RoomHandler roomHandler) { String address = be.getRawMessage().getString("address"); String roomId = getRoomId(address); if (roomId != null && roomHandler.roomExists(roomId)) { - roomHandler.getRoom(roomId).getIndex().future().onComplete(ar -> { + roomHandler.getRoom(roomId).getIndex(index -> { // Copies envelope structure from EventBusBridgeImpl##deliverMessage be.socket().write(Buffer.buffer(new JsonObject() .put("type", "rec") .put("address", address) - .put("body", ar.result()) + .put("body", index) .encode() )); }); diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java index fb158a3..f21bfb2 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java @@ -2,7 +2,6 @@ import edu.vanderbilt.yunyulin.speechdrop.handlers.RoomHandler; import io.vertx.core.CompositeFuture; -import io.vertx.core.Promise; import io.vertx.core.Handler; import io.vertx.core.Vertx; import lombok.AllArgsConstructor; @@ -24,11 +23,10 @@ public void schedule() { List toRemove = roomHandler.getDataStore().entrySet().stream() .filter(el -> System.currentTimeMillis() - el.getValue().ctime > purgeIntervalInSeconds * 1000) .map(Map.Entry::getKey) - .collect(Collectors.toList()); + .collect(Collectors.toList()); CompositeFuture.all( toRemove.stream() .map(roomHandler::queueRoomDeletion) - .map(Promise::future) .collect(Collectors.toList())) .onComplete(res -> { roomHandler.writeRooms(); diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java index 0c110c8..f70b148 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java @@ -124,8 +124,8 @@ public void mount(Router router) throws IOException { sendEmptyIndex(ctx, 404); } else { Room r = roomHandler.getRoom(roomId); - r.getIndex().future().onComplete(ar -> - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()) + r.getIndex(index -> + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(index) ); } }); @@ -136,8 +136,7 @@ public void mount(Router router) throws IOException { ctx.response().setStatusCode(404).end(); } else { Room r = roomHandler.getRoom(roomId); - r.getFiles().future().onComplete(ar -> { - Collection files = ar.result(); + r.getFiles(files -> { String outFile = r.getData().name.trim() + ".zip"; Handler writeBufferToResponse = buf -> ctx.response() @@ -178,7 +177,7 @@ public void mount(Router router) throws IOException { ctx.response().setStatusCode(404).end(EMPTY_INDEX); } else { Room r = roomHandler.getRoom(roomId); - r.handleUpload(ctx).future().onComplete(ar -> { + r.handleUpload(ctx, ar -> { if (ar.succeeded()) { String index = ar.result(); ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON_PRODUCES).end(index); @@ -203,9 +202,9 @@ public void mount(Router router) throws IOException { if (fileIndex == null) { sendEmptyIndex(ctx, 400); } else { - r.deleteFile(Integer.parseInt(fileIndex)).future().onComplete(ar -> { - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()); - broadcaster.publishUpdate(r.getId(), ar.result()); + r.deleteFile(Integer.parseInt(fileIndex), index -> { + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(index); + broadcaster.publishUpdate(r.getId(), index); }); } } @@ -229,14 +228,14 @@ public void mount(Router router) throws IOException { redirect(ctx, "/"); } else { Room r = roomHandler.getRoom(roomId); - r.getIndex().future().onComplete(ar -> { + r.getIndex(index -> { String escapedRoomName = HTML_ESCAPER.escape(r.getData().name); JsonObject configPayload = new JsonObject() .put("mediaUrl", mediaUrl) .put("roomName", escapedRoomName) .put("roomId", r.getId()) .put("allowedMimes", allowedMimeTypesCsv) - .put("initialFiles", new JsonArray(ar.result())); + .put("initialFiles", new JsonArray(index)); ctx.response().putHeader(CONTENT_TYPE, TEXT_HTML).end( roomTemplate diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java index 721598d..85e9982 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java @@ -8,7 +8,7 @@ import edu.vanderbilt.yunyulin.speechdrop.SpeechDropApplication; import edu.vanderbilt.yunyulin.speechdrop.room.Room; import edu.vanderbilt.yunyulin.speechdrop.room.RoomData; -import io.vertx.core.Promise; +import io.vertx.core.Future; import io.vertx.core.Vertx; import io.vertx.core.buffer.Buffer; import lombok.Getter; @@ -97,15 +97,12 @@ public void writeRooms() { } } - public Promise queueRoomDeletion(String id) { - Promise promise = Promise.promise(); + public Future queueRoomDeletion(String id) { if (dataStore.remove(id) != null) { roomCache.invalidate(id); File toDelete = new File(SpeechDropApplication.BASE_PATH, id); - vertx.fileSystem().deleteRecursive(toDelete.getPath(), true).onComplete(promise); - } else { - promise.complete(); + return vertx.fileSystem().deleteRecursive(toDelete.getPath(), true); } - return promise; + return Future.succeededFuture(); } } diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java index abb4715..f182389 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java @@ -2,8 +2,9 @@ import edu.vanderbilt.yunyulin.speechdrop.SpeechDropApplication; import edu.vanderbilt.yunyulin.speechdrop.handlers.IndexHandler; +import io.vertx.core.AsyncResult; +import io.vertx.core.Future; import io.vertx.core.Handler; -import io.vertx.core.Promise; import io.vertx.core.Vertx; import io.vertx.ext.web.FileUpload; import io.vertx.ext.web.RoutingContext; @@ -41,12 +42,11 @@ private void scheduleOperation(Handler operation) { } } - public Promise handleUpload(RoutingContext ctx) { - Promise uploadPromise = Promise.promise(); - + public void handleUpload(RoutingContext ctx, Handler> handler) { Iterator itr = ctx.fileUploads().iterator(); if (!itr.hasNext()) { - uploadPromise.fail(new Exception("no_file")); + handler.handle(Future.failedFuture("no_file")); + return; } FileUpload uploadedFile = itr.next(); @@ -54,31 +54,27 @@ public Promise handleUpload(RoutingContext ctx) { String mimeType = uploadedFile.contentType(); if (!SpeechDropApplication.allowedMimeTypes.contains(mimeType)) { - uploadPromise.fail(new Exception("bad_type")); + handler.handle(Future.failedFuture("bad_type")); + return; } if (uploadedFile.size() > SpeechDropApplication.maxUploadSize) { - uploadPromise.fail(new Exception("too_large")); + handler.handle(Future.failedFuture("too_large")); + return; } - scheduleOperation(index -> index.addFile(uploadedFile, now, uploadPromise::complete)); - return uploadPromise; + scheduleOperation(index -> index.addFile(uploadedFile, now, + newIndex -> handler.handle(Future.succeededFuture(newIndex)))); } - public Promise getIndex() { - Promise indexPromise = Promise.promise(); - scheduleOperation(index -> indexPromise.complete(index.getIndexString())); - return indexPromise; + public void getIndex(Handler handler) { + scheduleOperation(index -> handler.handle(index.getIndexString())); } - public Promise> getFiles() { - Promise> filesPromise = Promise.promise(); - scheduleOperation(index -> filesPromise.complete(index.getFiles())); - return filesPromise; + public void getFiles(Handler> handler) { + scheduleOperation(index -> handler.handle(index.getFiles())); } - public Promise deleteFile(int fileIndex) { - Promise deletePromise = Promise.promise(); - scheduleOperation(index -> index.deleteFile(fileIndex, deletePromise::complete)); - return deletePromise; + public void deleteFile(int fileIndex, Handler handler) { + scheduleOperation(index -> index.deleteFile(fileIndex, handler)); } } From 3adfaa1d66fd6530a33610d9865d4276e88c5985 Mon Sep 17 00:00:00 2001 From: Yunyu Lin Date: Wed, 24 Sep 2025 21:23:25 -0400 Subject: [PATCH 7/7] Migrate to Vert.x 5 and Java 21 --- .github/workflows/jar-build.yml | 4 +- pom.xml | 16 +- .../yunyulin/speechdrop/Broadcaster.java | 35 ++-- .../yunyulin/speechdrop/PurgeTask.java | 4 +- .../speechdrop/SpeechDropApplication.java | 150 ++++++++---------- .../speechdrop/handlers/IndexHandler.java | 128 +++++++-------- .../speechdrop/handlers/RoomHandler.java | 2 +- .../yunyulin/speechdrop/room/Room.java | 87 ++++------ 8 files changed, 202 insertions(+), 224 deletions(-) diff --git a/.github/workflows/jar-build.yml b/.github/workflows/jar-build.yml index cca9c62..61479de 100644 --- a/.github/workflows/jar-build.yml +++ b/.github/workflows/jar-build.yml @@ -15,11 +15,11 @@ jobs: - name: Check out repository uses: actions/checkout@v4 - - name: Set up Temurin JDK 17 + - name: Set up Temurin JDK 21 uses: actions/setup-java@v4 with: distribution: temurin - java-version: '17' + java-version: '21' cache: maven - name: Build JAR diff --git a/pom.xml b/pom.xml index 68275b8..176f298 100644 --- a/pom.xml +++ b/pom.xml @@ -12,12 +12,12 @@ frontend edu.vanderbilt.yunyulin.speechdrop.SpeechDropVerticle UTF-8 - 17 - 17 - 17 - 17 - 17 - 4.5.9 + 21 + 21 + 21 + 21 + 21 + 5.0.4 @@ -126,5 +126,9 @@ 1.18.42 provided + + org.slf4j + slf4j-api + diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java index d866619..044474a 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java @@ -7,6 +7,7 @@ import io.vertx.core.json.JsonObject; import io.vertx.ext.bridge.BridgeEventType; import io.vertx.ext.bridge.PermittedOptions; +import io.vertx.ext.web.Router; import io.vertx.ext.web.handler.sockjs.SockJSBridgeOptions; import io.vertx.ext.web.handler.sockjs.SockJSHandler; @@ -14,29 +15,33 @@ public class Broadcaster { private final Vertx vertx; - private final SockJSHandler sockJSHandler; + private final SockJSHandler sockJSHandler; + private final Router sockJSRouter; private static final String ADDR_PREFIX = "speechdrop.room."; public Broadcaster(Vertx vertx, RoomHandler roomHandler) { this.vertx = vertx; - this.sockJSHandler = SockJSHandler.create(vertx); - sockJSHandler.bridge(new SockJSBridgeOptions().addOutboundPermitted( + this.sockJSHandler = SockJSHandler.create(vertx); + this.sockJSRouter = sockJSHandler.bridge(new SockJSBridgeOptions().addOutboundPermitted( new PermittedOptions().setAddressRegex("speechdrop\\.room\\..+") ), be -> { if (be.type() == BridgeEventType.REGISTER) { String address = be.getRawMessage().getString("address"); String roomId = getRoomId(address); if (roomId != null && roomHandler.roomExists(roomId)) { - roomHandler.getRoom(roomId).getIndex(index -> { - // Copies envelope structure from EventBusBridgeImpl##deliverMessage - be.socket().write(Buffer.buffer(new JsonObject() - .put("type", "rec") - .put("address", address) - .put("body", index) - .encode() - )); - }); + roomHandler.getRoom(roomId).getIndex() + .onSuccess(index -> { + // Copies envelope structure from EventBusBridgeImpl##deliverMessage + be.socket().write(Buffer.buffer(new JsonObject() + .put("type", "rec") + .put("address", address) + .put("body", index) + .encode() + )); + be.complete(true); + }) + .onFailure(err -> be.complete(false)); return; } else { be.complete(false); @@ -62,9 +67,9 @@ private String getRoomId(String address) { } } - public SockJSHandler getSockJSHandler() { - return sockJSHandler; - } + public Router getSockJSRouter() { + return sockJSRouter; + } public void publishUpdate(String room, String data) { vertx.eventBus().publish(ADDR_PREFIX + room, data); diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java index f21bfb2..ee1f1c3 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java @@ -1,7 +1,7 @@ package edu.vanderbilt.yunyulin.speechdrop; import edu.vanderbilt.yunyulin.speechdrop.handlers.RoomHandler; -import io.vertx.core.CompositeFuture; +import io.vertx.core.Future; import io.vertx.core.Handler; import io.vertx.core.Vertx; import lombok.AllArgsConstructor; @@ -24,7 +24,7 @@ public void schedule() { .filter(el -> System.currentTimeMillis() - el.getValue().ctime > purgeIntervalInSeconds * 1000) .map(Map.Entry::getKey) .collect(Collectors.toList()); - CompositeFuture.all( + Future.all( toRemove.stream() .map(roomHandler::queueRoomDeletion) .collect(Collectors.toList())) diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java index f70b148..b25ac05 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/SpeechDropApplication.java @@ -2,23 +2,22 @@ import edu.vanderbilt.yunyulin.speechdrop.handlers.RoomHandler; import edu.vanderbilt.yunyulin.speechdrop.room.Room; -import io.vertx.core.Handler; -import io.vertx.core.Vertx; +import io.vertx.core.Future; +import io.vertx.core.Vertx; import io.vertx.core.buffer.Buffer; import io.vertx.core.json.JsonArray; import io.vertx.core.json.JsonObject; -import io.vertx.core.impl.logging.Logger; -import io.vertx.core.impl.logging.LoggerFactory; -import io.vertx.ext.web.Router; -import io.vertx.ext.web.RoutingContext; -import io.vertx.ext.web.handler.*; -import io.vertx.ext.web.sstore.LocalSessionStore; +import io.vertx.ext.web.Router; +import io.vertx.ext.web.RoutingContext; +import io.vertx.ext.web.handler.*; +import io.vertx.ext.web.sstore.LocalSessionStore; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.io.File; import java.io.IOException; import java.nio.file.Files; import java.util.Arrays; -import java.util.Collection; import java.util.List; import static edu.vanderbilt.yunyulin.speechdrop.Util.HTML_ESCAPER; @@ -98,7 +97,7 @@ public void mount(Router router) throws IOException { ctx.response().putHeader(CONTENT_TYPE, TEXT_HTML).end(mainPage) ); - router.route("/sock/*").handler(broadcaster.getSockJSHandler()); + router.route("/sock/*").subRouter(broadcaster.getSockJSRouter()); router.route("/static/*").handler(StaticHandler.create("static")); @@ -124,9 +123,9 @@ public void mount(Router router) throws IOException { sendEmptyIndex(ctx, 404); } else { Room r = roomHandler.getRoom(roomId); - r.getIndex(index -> - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(index) - ); + r.getIndex() + .onSuccess(index -> ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(index)) + .onFailure(err -> ctx.response().setStatusCode(500).end()); } }); @@ -136,58 +135,43 @@ public void mount(Router router) throws IOException { ctx.response().setStatusCode(404).end(); } else { Room r = roomHandler.getRoom(roomId); - r.getFiles(files -> { - String outFile = r.getData().name.trim() + ".zip"; - - Handler writeBufferToResponse = buf -> ctx.response() - .setChunked(true) - // See https://stackoverflow.com/a/38324508 - .putHeader("Content-Disposition", "attachment; filename=\"" + outFile - .replace("\\", "\\\\") - .replace("\"", "\\\"") + "\"" - ) - .putHeader(CONTENT_TYPE, "application/octet-stream") - .end(buf); - - if (files.size() == 0) { - writeBufferToResponse.handle(Util.getEmptyZipBuffer()); - } else { - vertx.executeBlocking(promise -> { - try { - promise.complete(Util.zip(files)); - } catch (IOException e) { - promise.fail(e); - } - }, false).onComplete(res -> { - if (res.succeeded()) { - writeBufferToResponse.handle(res.result()); - } else { - ctx.response().setStatusCode(500).end(); + r.getFiles() + .compose(files -> { + if (files.isEmpty()) { + return Future.succeededFuture(Util.getEmptyZipBuffer()); } - }); - } - }); + return vertx.executeBlocking(() -> Util.zip(files), false); + }) + .onSuccess(buf -> { + String outFile = r.getData().name.trim() + ".zip"; + ctx.response() + .setChunked(true) + // See https://stackoverflow.com/a/38324508 + .putHeader("Content-Disposition", "attachment; filename=\"" + outFile + .replace("\\", "\\\\") + .replace("\"", "\\\"") + "\"") + .putHeader(CONTENT_TYPE, "application/octet-stream") + .end(buf); + }) + .onFailure(err -> ctx.response().setStatusCode(500).end()); } }); - router.route("/:roomid/upload").method(POST).produces(APPLICATION_JSON_PRODUCES).handler(ctx -> { - String roomId = ctx.request().getParam("roomid"); - if (!roomHandler.roomExists(roomId)) { - LOGGER.warn("(Upload) Nonexist " + roomId); - ctx.response().setStatusCode(404).end(EMPTY_INDEX); - } else { + router.route("/:roomid/upload").method(POST).produces(APPLICATION_JSON_PRODUCES).handler(ctx -> { + String roomId = ctx.request().getParam("roomid"); + if (!roomHandler.roomExists(roomId)) { + LOGGER.warn("(Upload) Nonexist " + roomId); + ctx.response().setStatusCode(404).end(EMPTY_INDEX); + } else { Room r = roomHandler.getRoom(roomId); - r.handleUpload(ctx, ar -> { - if (ar.succeeded()) { - String index = ar.result(); - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON_PRODUCES).end(index); - broadcaster.publishUpdate(r.getId(), index); - } else { - ctx.response().setStatusCode(400).end( - new JsonObject().put("err", ar.cause().getMessage()).toString() - ); - } - }); + r.handleUpload(ctx) + .onSuccess(index -> { + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON_PRODUCES).end(index); + broadcaster.publishUpdate(r.getId(), index); + }) + .onFailure(err -> ctx.response().setStatusCode(400).end( + new JsonObject().put("err", err.getMessage()).toString() + )); } }); @@ -202,10 +186,12 @@ public void mount(Router router) throws IOException { if (fileIndex == null) { sendEmptyIndex(ctx, 400); } else { - r.deleteFile(Integer.parseInt(fileIndex), index -> { - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(index); - broadcaster.publishUpdate(r.getId(), index); - }); + r.deleteFile(Integer.parseInt(fileIndex)) + .onSuccess(index -> { + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(index); + broadcaster.publishUpdate(r.getId(), index); + }) + .onFailure(err -> ctx.response().setStatusCode(500).end()); } } }); @@ -222,27 +208,29 @@ public void mount(Router router) throws IOException { } else { mediaUrl = config.getString("mediaUrl"); } - router.route("/:roomid").method(GET).handler(ctx -> { - String roomId = ctx.request().getParam("roomid"); + router.route("/:roomid").method(GET).handler(ctx -> { + String roomId = ctx.request().getParam("roomid"); if (!roomHandler.roomExists(roomId)) { redirect(ctx, "/"); } else { Room r = roomHandler.getRoom(roomId); - r.getIndex(index -> { - String escapedRoomName = HTML_ESCAPER.escape(r.getData().name); - JsonObject configPayload = new JsonObject() - .put("mediaUrl", mediaUrl) - .put("roomName", escapedRoomName) - .put("roomId", r.getId()) - .put("allowedMimes", allowedMimeTypesCsv) - .put("initialFiles", new JsonArray(index)); + r.getIndex() + .onSuccess(index -> { + String escapedRoomName = HTML_ESCAPER.escape(r.getData().name); + JsonObject configPayload = new JsonObject() + .put("mediaUrl", mediaUrl) + .put("roomName", escapedRoomName) + .put("roomId", r.getId()) + .put("allowedMimes", allowedMimeTypesCsv) + .put("initialFiles", new JsonArray(index)); - ctx.response().putHeader(CONTENT_TYPE, TEXT_HTML).end( - roomTemplate - .replace("<%= ROOM_CONFIG %>", configPayload.encode()) - ); - }); + ctx.response().putHeader(CONTENT_TYPE, TEXT_HTML).end( + roomTemplate + .replace("<%= ROOM_CONFIG %>", configPayload.encode()) + ); + }) + .onFailure(err -> ctx.response().setStatusCode(500).end()); } }); - } -} + } +} diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java index e207e84..fa44615 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/IndexHandler.java @@ -4,14 +4,14 @@ import com.fasterxml.jackson.annotation.JsonProperty; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; -import io.vertx.core.Handler; -import io.vertx.core.Vertx; -import io.vertx.core.buffer.Buffer; -import io.vertx.ext.web.FileUpload; -import lombok.Data; - -import java.io.File; -import java.io.IOException; +import io.vertx.core.Future; +import io.vertx.core.Vertx; +import io.vertx.core.buffer.Buffer; +import io.vertx.ext.web.FileUpload; +import lombok.Data; + +import java.io.File; +import java.io.IOException; import java.util.*; import static edu.vanderbilt.yunyulin.speechdrop.SpeechDropApplication.LOGGER; @@ -33,63 +33,63 @@ public IndexHandler(Vertx vertx, File uploadDirectory) { private List entries; private boolean loaded = false; - public void load(Handler onComplete) { - vertx.fileSystem().readFile(indexFile.getPath(), res -> { - if (res.succeeded()) { - try { - entries = new ArrayList<>(Arrays.asList( - mapper.readValue(res.result().toString(), FileEntry[].class) - )); - } catch (IOException e) { // This should never happen - e.printStackTrace(); - } - } else { - entries = new ArrayList<>(); - } - loaded = true; - onComplete.handle(this); - }); - } - - public void addFile(FileUpload uploadedFile, Date creationTime, Handler indexHandler) { - checkLoad(); - LOGGER.info("[" + uploadDirectory.getName() + "] Processing upload " - + uploadedFile.fileName() - + " (" + uploadedFile.size() + ")"); - vertx.fileSystem().mkdir(uploadDirectory.getPath(), uploadDirRes -> { - File destDir = new File(uploadDirectory, Integer.toString(entries.size())); - vertx.fileSystem().mkdir(destDir.getPath(), mkdirRes -> { - File dest = new File(destDir, uploadedFile.fileName()); - vertx.fileSystem().copy(uploadedFile.uploadedFileName(), dest.getPath(), - res -> { - entries.add(new FileEntry(dest.getName(), creationTime.getTime())); - indexHandler.handle(writeIndex()); - } - ); - }); - }); + public Future load() { + return vertx.fileSystem().readFile(indexFile.getPath()) + .compose(buffer -> { + try { + entries = new ArrayList<>(Arrays.asList( + mapper.readValue(buffer.toString(), FileEntry[].class) + )); + loaded = true; + return Future.succeededFuture(this); + } catch (IOException e) { + return Future.failedFuture(e); + } + }) + .recover(err -> { + entries = new ArrayList<>(); + loaded = true; + return Future.succeededFuture(this); + }); } - - public void deleteFile(int index, Handler indexHandler) { - checkLoad(); - LOGGER.info("[" + uploadDirectory.getName() + "] Processing delete for index " - + index); - String completedIndex; + + public Future addFile(FileUpload uploadedFile, Date creationTime) { + checkLoad(); + LOGGER.info("[" + uploadDirectory.getName() + "] Processing upload " + + uploadedFile.fileName() + + " (" + uploadedFile.size() + ")"); + File destDir = new File(uploadDirectory, Integer.toString(entries.size())); + File dest = new File(destDir, uploadedFile.fileName()); + return vertx.fileSystem() + .mkdirs(uploadDirectory.getPath()) + .compose(v -> vertx.fileSystem().mkdirs(destDir.getPath())) + .compose(v -> vertx.fileSystem().copy(uploadedFile.uploadedFileName(), dest.getPath())) + .compose(v -> { + entries.add(new FileEntry(dest.getName(), creationTime.getTime())); + return writeIndex(); + }); + } + + public Future deleteFile(int index) { + checkLoad(); + LOGGER.info("[" + uploadDirectory.getName() + "] Processing delete for index " + + index); if (entries.get(index) != null) { entries.set(index, null); - completedIndex = writeIndex(); - vertx.fileSystem().deleteRecursive( - new File(uploadDirectory, Integer.toString(index)).getPath(), true - ); - } else { - completedIndex = getIndexString(); + String dirToDelete = new File(uploadDirectory, Integer.toString(index)).getPath(); + return writeIndex() + .compose(indexString -> vertx.fileSystem() + .deleteRecursive(dirToDelete) + .recover(err -> Future.succeededFuture()) + .map(v -> indexString) + ); } - indexHandler.handle(completedIndex); + return Future.succeededFuture(getIndexString()); } - - public Collection getFiles() { - checkLoad(); - List files = new ArrayList<>(entries.size()); + + public Collection getFiles() { + checkLoad(); + List files = new ArrayList<>(entries.size()); int index = 0; for (FileEntry entry : entries) { if (entry != null) { @@ -111,11 +111,11 @@ public String getIndexString() { } } - private String writeIndex() { + private Future writeIndex() { String indexString = getIndexString(); - vertx.fileSystem().writeFile(indexFile.getPath(), - Buffer.buffer(indexString)); - return indexString; + return vertx.fileSystem() + .writeFile(indexFile.getPath(), Buffer.buffer(indexString)) + .map(indexString); } private void checkLoad() { diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java index 85e9982..2e64fc5 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java @@ -101,7 +101,7 @@ public Future queueRoomDeletion(String id) { if (dataStore.remove(id) != null) { roomCache.invalidate(id); File toDelete = new File(SpeechDropApplication.BASE_PATH, id); - return vertx.fileSystem().deleteRecursive(toDelete.getPath(), true); + return vertx.fileSystem().deleteRecursive(toDelete.getPath()); } return Future.succeededFuture(); } diff --git a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java index f182389..17756ef 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/room/Room.java @@ -2,51 +2,35 @@ import edu.vanderbilt.yunyulin.speechdrop.SpeechDropApplication; import edu.vanderbilt.yunyulin.speechdrop.handlers.IndexHandler; -import io.vertx.core.AsyncResult; import io.vertx.core.Future; -import io.vertx.core.Handler; import io.vertx.core.Vertx; -import io.vertx.ext.web.FileUpload; -import io.vertx.ext.web.RoutingContext; -import lombok.Getter; - -import java.io.File; -import java.util.*; - -public class Room { - @Getter - private final String id; - @Getter - private final RoomData data; - - private final Queue> queuedOperations = new ArrayDeque<>(2); - private IndexHandler indexHandler; - - public Room(Vertx vertx, String id, RoomData data) { - this.id = id; - this.data = data; - - new IndexHandler(vertx, new File(SpeechDropApplication.BASE_PATH, id)).load(loadedIndex -> { - this.indexHandler = loadedIndex; - while (!queuedOperations.isEmpty()) { - queuedOperations.poll().handle(loadedIndex); - } - }); - } - - private void scheduleOperation(Handler operation) { - if (indexHandler != null) { - operation.handle(indexHandler); - } else { - queuedOperations.offer(operation); - } - } - - public void handleUpload(RoutingContext ctx, Handler> handler) { +import io.vertx.ext.web.FileUpload; +import io.vertx.ext.web.RoutingContext; +import lombok.Getter; + +import java.io.File; +import java.util.*; + +public class Room { + @Getter + private final String id; + @Getter + private final RoomData data; + + private final Future indexHandlerFuture; + + public Room(Vertx vertx, String id, RoomData data) { + this.id = id; + this.data = data; + + IndexHandler handler = new IndexHandler(vertx, new File(SpeechDropApplication.BASE_PATH, id)); + this.indexHandlerFuture = handler.load(); + } + + public Future handleUpload(RoutingContext ctx) { Iterator itr = ctx.fileUploads().iterator(); if (!itr.hasNext()) { - handler.handle(Future.failedFuture("no_file")); - return; + return Future.failedFuture("no_file"); } FileUpload uploadedFile = itr.next(); @@ -54,27 +38,24 @@ public void handleUpload(RoutingContext ctx, Handler> handle String mimeType = uploadedFile.contentType(); if (!SpeechDropApplication.allowedMimeTypes.contains(mimeType)) { - handler.handle(Future.failedFuture("bad_type")); - return; + return Future.failedFuture("bad_type"); } if (uploadedFile.size() > SpeechDropApplication.maxUploadSize) { - handler.handle(Future.failedFuture("too_large")); - return; + return Future.failedFuture("too_large"); } - scheduleOperation(index -> index.addFile(uploadedFile, now, - newIndex -> handler.handle(Future.succeededFuture(newIndex)))); + return indexHandlerFuture.compose(index -> index.addFile(uploadedFile, now)); } - public void getIndex(Handler handler) { - scheduleOperation(index -> handler.handle(index.getIndexString())); + public Future getIndex() { + return indexHandlerFuture.map(IndexHandler::getIndexString); } - public void getFiles(Handler> handler) { - scheduleOperation(index -> handler.handle(index.getFiles())); + public Future> getFiles() { + return indexHandlerFuture.map(index -> new ArrayList<>(index.getFiles())); } - public void deleteFile(int fileIndex, Handler handler) { - scheduleOperation(index -> index.deleteFile(fileIndex, handler)); + public Future deleteFile(int fileIndex) { + return indexHandlerFuture.compose(index -> index.deleteFile(fileIndex)); } }