diff --git a/.github/workflows/jar-build.yml b/.github/workflows/jar-build.yml index 6f37414..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 8 + - name: Set up Temurin JDK 21 uses: actions/setup-java@v4 with: distribution: temurin - java-version: '8' + java-version: '21' cache: maven - name: Build JAR diff --git a/pom.xml b/pom.xml index ee69d17..176f298 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 + 21 + 21 + 21 + 21 + 21 + 5.0.4 @@ -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 @@ -116,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 8cd6cda..044474a 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/Broadcaster.java @@ -7,36 +7,42 @@ 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.Router; +import io.vertx.ext.web.handler.sockjs.SockJSBridgeOptions; import io.vertx.ext.web.handler.sockjs.SockJSHandler; import static edu.vanderbilt.yunyulin.speechdrop.SpeechDropApplication.LOGGER; 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 BridgeOptions().addOutboundPermitted( - new PermittedOptions().setAddressRegex("speechdrop\\.room\\..+") - ), be -> { + 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().setHandler(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() - )); - }); + 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); return; @@ -61,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 57cce74..ee1f1c3 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/PurgeTask.java @@ -1,10 +1,9 @@ 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.Future; +import io.vertx.core.Handler; +import io.vertx.core.Vertx; import lombok.AllArgsConstructor; import java.util.List; @@ -24,14 +23,17 @@ 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()); - CompositeFuture.all(toRemove.stream().map(roomHandler::queueRoomDeletion).collect(Collectors.toList())) - .setHandler(res -> { - roomHandler.writeRooms(); - LOGGER.info("Purged " + toRemove.size() + " rooms"); - }); - }; - vertx.setPeriodic(3 * 60 * 60 * 1000, runPurge); - runPurge.handle(0L); - } -} + .collect(Collectors.toList()); + Future.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..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.logging.Logger; -import io.vertx.core.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.core.json.JsonArray; +import io.vertx.core.json.JsonObject; +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; @@ -88,18 +87,17 @@ 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) ); - router.route("/sock/*").handler(broadcaster.getSockJSHandler()); + router.route("/sock/*").subRouter(broadcaster.getSockJSRouter()); router.route("/static/*").handler(StaticHandler.create("static")); @@ -119,98 +117,84 @@ 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 { - Room r = roomHandler.getRoom(roomId); - r.getIndex().setHandler(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"); - if (!roomHandler.roomExists(roomId)) { - 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() - .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(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(); - } - }); - } - }); - } - }); - - 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).setHandler(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 { - r.deleteFile(Integer.parseInt(fileIndex)).setHandler(ar -> { - ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(ar.result()); - broadcaster.publishUpdate(r.getId(), ar.result()); - }); - } - } - }); + 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() + .onSuccess(index -> ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(index)) + .onFailure(err -> ctx.response().setStatusCode(500).end()); + } + }); + + 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() + .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 { + Room r = roomHandler.getRoom(roomId); + 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() + )); + } + }); + + 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)) + .onSuccess(index -> { + ctx.response().putHeader(CONTENT_TYPE, APPLICATION_JSON).end(index); + broadcaster.publishUpdate(r.getId(), index); + }) + .onFailure(err -> ctx.response().setStatusCode(500).end()); + } + } + }); router.route("/:roomid/").method(GET).handler(ctx -> { String roomId = ctx.request().getParam("roomid"); @@ -224,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"); - if (!roomHandler.roomExists(roomId)) { - 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) - .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()) - ); - }); - } - }); - } -} + 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() + .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()) + ); + }) + .onFailure(err -> ctx.response().setStatusCode(500).end()); + } + }); + } +} 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..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 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(); - vertx.fileSystem().deleteRecursive( - new File(uploadDirectory, Integer.toString(index)).getPath(), true, null - ); - } else { - completedIndex = getIndexString(); - } - indexHandler.handle(completedIndex); - } - - public Collection getFiles() { - checkLoad(); - List files = new ArrayList<>(entries.size()); + 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 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); + String dirToDelete = new File(uploadDirectory, Integer.toString(index)).getPath(); + return writeIndex() + .compose(indexString -> vertx.fileSystem() + .deleteRecursive(dirToDelete) + .recover(err -> Future.succeededFuture()) + .map(v -> indexString) + ); + } + return Future.succeededFuture(getIndexString()); + } + + public Collection getFiles() { + checkLoad(); + List files = new ArrayList<>(entries.size()); int index = 0; for (FileEntry entry : entries) { if (entry != null) { @@ -111,13 +111,12 @@ public String getIndexString() { } } - private String writeIndex() { - String indexString = getIndexString(); - vertx.fileSystem().writeFile(indexFile.getPath(), - Buffer.buffer(indexString), - null); - return indexString; - } + private Future writeIndex() { + String indexString = getIndexString(); + return vertx.fileSystem() + .writeFile(indexFile.getPath(), Buffer.buffer(indexString)) + .map(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..2e64fc5 100644 --- a/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java +++ b/src/main/java/edu/vanderbilt/yunyulin/speechdrop/handlers/RoomHandler.java @@ -8,9 +8,9 @@ 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.Vertx; -import io.vertx.core.buffer.Buffer; +import io.vertx.core.Future; +import io.vertx.core.Vertx; +import io.vertx.core.buffer.Buffer; import lombok.Getter; import java.io.File; @@ -89,22 +89,20 @@ 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; - } -} + try { + vertx.fileSystem() + .writeFile(roomsFile.getPath(), Buffer.buffer(mapper.writeValueAsString(dataStore))); + } catch (JsonProcessingException e) { // This should never happen + e.printStackTrace(); + } + } + + 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()); + } + 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 c61c322..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,83 +2,60 @@ 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.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 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; - } -} +import io.vertx.core.Future; +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 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()) { + return Future.failedFuture("no_file"); + } + + FileUpload uploadedFile = itr.next(); + Date now = new Date(); + + String mimeType = uploadedFile.contentType(); + if (!SpeechDropApplication.allowedMimeTypes.contains(mimeType)) { + return Future.failedFuture("bad_type"); + } + if (uploadedFile.size() > SpeechDropApplication.maxUploadSize) { + return Future.failedFuture("too_large"); + } + + return indexHandlerFuture.compose(index -> index.addFile(uploadedFile, now)); + } + + public Future getIndex() { + return indexHandlerFuture.map(IndexHandler::getIndexString); + } + + public Future> getFiles() { + return indexHandlerFuture.map(index -> new ArrayList<>(index.getFiles())); + } + + public Future deleteFile(int fileIndex) { + return indexHandlerFuture.compose(index -> index.deleteFile(fileIndex)); + } +}