diff --git a/taskbus/README.md b/taskbus/README.md index 3a167ca..2ab6e9e 100644 --- a/taskbus/README.md +++ b/taskbus/README.md @@ -1,2 +1,13 @@ # taskbus + A task distribution, which assign tasks to specific executors + +## How to use + +Just view tests in [`test`](src/test/java) folder for example. + +## Caution + +This project is just demo a simple task distribution. +In production there will be more complicated cases that you need to handle, like executor type, +monitor, etc. diff --git a/taskbus/build.gradle.kts b/taskbus/build.gradle.kts index e6eecba..0d6c1f0 100644 --- a/taskbus/build.gradle.kts +++ b/taskbus/build.gradle.kts @@ -1,19 +1,40 @@ plugins { - id("java") + java + idea } group = "com.miti99" version = "1.0-SNAPSHOT" +java { + toolchain { + languageVersion.set(JavaLanguageVersion.of(21)) + } +} + +idea { + module { + isDownloadJavadoc = true + isDownloadSources = true + } +} + repositories { mavenCentral() } dependencies { - testImplementation(platform("org.junit:junit-bom:5.10.0")) + annotationProcessor("org.projectlombok:lombok:1.18.36") + compileOnly("org.projectlombok:lombok:1.18.36") + + implementation("org.slf4j:slf4j-simple:2.0.16") + + testAnnotationProcessor("org.projectlombok:lombok:1.18.36") + testCompileOnly("org.projectlombok:lombok:1.18.36") + testImplementation(platform("org.junit:junit-bom:5.11.4")) testImplementation("org.junit.jupiter:junit-jupiter") } tasks.test { useJUnitPlatform() -} \ No newline at end of file +} diff --git a/taskbus/src/main/java/com/miti99/Main.java b/taskbus/src/main/java/com/miti99/Main.java deleted file mode 100644 index 540f964..0000000 --- a/taskbus/src/main/java/com/miti99/Main.java +++ /dev/null @@ -1,7 +0,0 @@ -package com.miti99; - -public class Main { - public static void main(String[] args) { - System.out.println("Hello, World!"); - } -} \ No newline at end of file diff --git a/taskbus/src/main/java/com/miti99/taskbus/bus/Bus.java b/taskbus/src/main/java/com/miti99/taskbus/bus/Bus.java new file mode 100644 index 0000000..9eb9b8d --- /dev/null +++ b/taskbus/src/main/java/com/miti99/taskbus/bus/Bus.java @@ -0,0 +1,7 @@ +package com.miti99.taskbus.bus; + +import com.miti99.taskbus.task.Task; + +public interface Bus { + void submit(Task task); +} diff --git a/taskbus/src/main/java/com/miti99/taskbus/bus/impl/ManualBus.java b/taskbus/src/main/java/com/miti99/taskbus/bus/impl/ManualBus.java new file mode 100644 index 0000000..13f3682 --- /dev/null +++ b/taskbus/src/main/java/com/miti99/taskbus/bus/impl/ManualBus.java @@ -0,0 +1,29 @@ +package com.miti99.taskbus.bus.impl; + +import com.miti99.taskbus.bus.Bus; +import com.miti99.taskbus.task.Task; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; + +public class ManualBus implements Bus { + private final ExecutorService[] executors; + private final int poolSize; + + public ManualBus(int poolSize) { + this.poolSize = poolSize; + executors = new ExecutorService[poolSize]; + for (int i = 0; i < poolSize; i++) { + executors[i] = Executors.newSingleThreadExecutor(); + } + } + + public ManualBus() { + this(Runtime.getRuntime().availableProcessors()); + } + + @Override + public void submit(Task task) { + int executorId = task.hash() % poolSize; + executors[executorId].submit(task::execute); + } +} diff --git a/taskbus/src/main/java/com/miti99/taskbus/bus/impl/NormalBus.java b/taskbus/src/main/java/com/miti99/taskbus/bus/impl/NormalBus.java new file mode 100644 index 0000000..9549faa --- /dev/null +++ b/taskbus/src/main/java/com/miti99/taskbus/bus/impl/NormalBus.java @@ -0,0 +1,16 @@ +package com.miti99.taskbus.bus.impl; + +import com.miti99.taskbus.bus.Bus; +import com.miti99.taskbus.task.Task; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; + +public class NormalBus implements Bus { + private final ExecutorService executors = + Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors()); + + @Override + public void submit(Task task) { + executors.submit(task::execute); + } +} diff --git a/taskbus/src/main/java/com/miti99/taskbus/task/Task.java b/taskbus/src/main/java/com/miti99/taskbus/task/Task.java new file mode 100644 index 0000000..d808f95 --- /dev/null +++ b/taskbus/src/main/java/com/miti99/taskbus/task/Task.java @@ -0,0 +1,7 @@ +package com.miti99.taskbus.task; + +public interface Task { + int hash(); + + void execute(); +} diff --git a/taskbus/src/main/java/com/miti99/taskbus/task/event/Event.java b/taskbus/src/main/java/com/miti99/taskbus/task/event/Event.java new file mode 100644 index 0000000..56e3d03 --- /dev/null +++ b/taskbus/src/main/java/com/miti99/taskbus/task/event/Event.java @@ -0,0 +1,12 @@ +package com.miti99.taskbus.task.event; + +import com.miti99.taskbus.task.Task; + +public interface Event extends Task { + int userId(); + + @Override + default int hash() { + return userId(); + } +} diff --git a/taskbus/src/main/java/com/miti99/taskbus/task/event/UserLoginEvent.java b/taskbus/src/main/java/com/miti99/taskbus/task/event/UserLoginEvent.java new file mode 100644 index 0000000..919f85e --- /dev/null +++ b/taskbus/src/main/java/com/miti99/taskbus/task/event/UserLoginEvent.java @@ -0,0 +1,12 @@ +package com.miti99.taskbus.task.event; + +import lombok.extern.slf4j.Slf4j; + +@Slf4j +public record UserLoginEvent(int userId) implements Event { + + @Override + public void execute() { + log.info("User {} logged in", userId); + } +} diff --git a/taskbus/src/main/java/com/miti99/taskbus/task/request/Request.java b/taskbus/src/main/java/com/miti99/taskbus/task/request/Request.java new file mode 100644 index 0000000..e75efc2 --- /dev/null +++ b/taskbus/src/main/java/com/miti99/taskbus/task/request/Request.java @@ -0,0 +1,12 @@ +package com.miti99.taskbus.task.request; + +import com.miti99.taskbus.task.Task; + +public interface Request extends Task { + int userId(); + + @Override + default int hash() { + return userId(); + } +} diff --git a/taskbus/src/main/java/com/miti99/taskbus/task/request/SendMessageRequest.java b/taskbus/src/main/java/com/miti99/taskbus/task/request/SendMessageRequest.java new file mode 100644 index 0000000..ca234ea --- /dev/null +++ b/taskbus/src/main/java/com/miti99/taskbus/task/request/SendMessageRequest.java @@ -0,0 +1,12 @@ +package com.miti99.taskbus.task.request; + +import lombok.extern.slf4j.Slf4j; + +@Slf4j +public record SendMessageRequest(int userId, String message) implements Request { + + @Override + public void execute() { + log.info("User {} sent message: {}", userId, message); + } +} diff --git a/taskbus/src/test/java/com/miti99/taskbus/bus/impl/BusTest.java b/taskbus/src/test/java/com/miti99/taskbus/bus/impl/BusTest.java new file mode 100644 index 0000000..6a6846c --- /dev/null +++ b/taskbus/src/test/java/com/miti99/taskbus/bus/impl/BusTest.java @@ -0,0 +1,24 @@ +package com.miti99.taskbus.bus.impl; + +import com.miti99.taskbus.bus.Bus; +import com.miti99.taskbus.task.event.UserLoginEvent; +import com.miti99.taskbus.task.request.SendMessageRequest; +import lombok.RequiredArgsConstructor; + +@RequiredArgsConstructor +public abstract class BusTest { + private final Bus bus; + + void submitTasks() { + bus.submit(new UserLoginEvent(1)); + bus.submit(new SendMessageRequest(1, "msg1")); + bus.submit(new SendMessageRequest(1, "msg2")); + bus.submit(new UserLoginEvent(2)); + bus.submit(new SendMessageRequest(2, "msg1")); + bus.submit(new SendMessageRequest(2, "msg2")); + bus.submit(new SendMessageRequest(1, "msg3")); + bus.submit(new SendMessageRequest(1, "msg4")); + bus.submit(new SendMessageRequest(2, "msg3")); + bus.submit(new SendMessageRequest(2, "msg4")); + } +} diff --git a/taskbus/src/test/java/com/miti99/taskbus/bus/impl/ManualBusTest.java b/taskbus/src/test/java/com/miti99/taskbus/bus/impl/ManualBusTest.java new file mode 100644 index 0000000..2f8fd2a --- /dev/null +++ b/taskbus/src/test/java/com/miti99/taskbus/bus/impl/ManualBusTest.java @@ -0,0 +1,15 @@ +package com.miti99.taskbus.bus.impl; + +import org.junit.jupiter.api.Test; + +class ManualBusTest extends BusTest { + + public ManualBusTest() { + super(new ManualBus()); + } + + @Test + void runTest() { + submitTasks(); + } +} diff --git a/taskbus/src/test/java/com/miti99/taskbus/bus/impl/NormalBusTest.java b/taskbus/src/test/java/com/miti99/taskbus/bus/impl/NormalBusTest.java new file mode 100644 index 0000000..0bb0d42 --- /dev/null +++ b/taskbus/src/test/java/com/miti99/taskbus/bus/impl/NormalBusTest.java @@ -0,0 +1,17 @@ +package com.miti99.taskbus.bus.impl; + +import static org.junit.jupiter.api.Assertions.*; + +import org.junit.jupiter.api.Test; + +class NormalBusTest extends BusTest { + + public NormalBusTest() { + super(new NormalBus()); + } + + @Test + void runTest() { + submitTasks(); + } +}