Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -37,8 +37,8 @@ public final class FfmTransport implements NativeTransport {
private static final ValueLayout.OfInt FRAME_INT =
ValueLayout.JAVA_INT_UNALIGNED.withOrder(ByteOrder.LITTLE_ENDIAN);

// Wire status codes from the Rust C ABI; anything but these two carries an error payload.
private static final int STATUS_OK = 0;
private static final int STATUS_ERR = 1;
private static final int STATUS_ABSENT = 2;

private static final Linker LINKER = Linker.nativeLinker();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,10 +27,11 @@ void generatedCompanionEnqueuesAndHandles(@TempDir Path dir) throws Exception {
String totalId = queue.enqueue(GreeterTasks.TOTAL, List.of(1, 2, 3, 4));

CountDownLatch done = new CountDownLatch(2);
try (Worker worker = queue.worker()
Worker worker = queue.worker()
.apply(builder -> GreeterTasks.bind(builder, new Greeter()))
.on(EventName.SUCCESS, event -> done.countDown())
.start()) {
.start();
try (worker) {
assertTrue(done.await(20, TimeUnit.SECONDS), "both tasks should complete");
assertEquals("hello ada", queue.getResult(greetId, String.class).orElseThrow());
assertEquals(10, queue.getResult(totalId, Integer.class).orElseThrow());
Expand All @@ -48,10 +49,11 @@ void injectsResourceParameter(@TempDir Path dir) throws Exception {
queue.resource("salutation", ctx -> "hi");
String id = queue.enqueue(ResourceGreeterTasks.GREET, "ada");
CountDownLatch done = new CountDownLatch(1);
try (Worker worker = queue.worker()
Worker worker = queue.worker()
.apply(builder -> ResourceGreeterTasks.bind(builder, new ResourceGreeter()))
.on(EventName.SUCCESS, event -> done.countDown())
.start()) {
.start();
try (worker) {
assertTrue(done.await(20, TimeUnit.SECONDS), "task should complete");
assertEquals("hi ada", queue.getResult(id, String.class).orElseThrow());
}
Expand All @@ -65,10 +67,11 @@ void registerViaGeneratedHandlerRegistry(@TempDir Path dir) throws Exception {
Taskito.builder().sqlite(dir.resolve("hr.db").toString()).open()) {
String id = queue.enqueue(GreeterTasks.GREET, "grace");
CountDownLatch done = new CountDownLatch(1);
try (Worker worker = queue.worker()
Worker worker = queue.worker()
.register(GreeterTasks.handlers(new Greeter())) // generated HandlerRegistry
.on(EventName.SUCCESS, event -> done.countDown())
.start()) {
.start();
try (worker) {
assertTrue(done.await(20, TimeUnit.SECONDS), "task should complete");
assertEquals("hello grace", queue.getResult(id, String.class).orElseThrow());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,8 @@ void legacyStoredSubscriptionStillMatchesOutcomes(@TempDir Path dir) throws Exce
WebhookUpdate.builder().events(List.of("success")).build());

queue.enqueue(task, 1);
try (Worker worker = queue.worker().handle(task, p -> p).start()) {
Worker worker = queue.worker().handle(task, p -> p).start();
try (worker) {
// The loopback URL fails the SSRF guard, but only a MATCHED hook
// records a delivery — which is exactly what we assert on.
List<Delivery> deliveries = List.of();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ void shedsStaleJobsUnderOverload(@TempDir Path dir) throws Exception {
queue.enqueue(slow, "x");
}

try (Worker worker = queue.worker()
Worker worker = queue.worker()
.concurrency(1)
.batchSize(1)
.handle(slow, payload -> {
Expand All @@ -56,7 +56,8 @@ void shedsStaleJobsUnderOverload(@TempDir Path dir) throws Exception {
}
return null;
})
.start()) {
.start();
try (worker) {
// Poll until every job is accounted for and at least one was shed.
long deadline = System.nanoTime() + Duration.ofSeconds(45).toNanos();
long codelDead = 0;
Expand Down Expand Up @@ -110,7 +111,8 @@ void codelConfiguredAfterBuilderIsStillApplied(@TempDir Path dir) throws Excepti
});
queue.codel("default", 1, 30);

try (Worker worker = builder.start()) {
Worker worker = builder.start();
try (worker) {
long deadline = System.nanoTime() + Duration.ofSeconds(45).toNanos();
long codelDead = 0;
QueueStats stats = queue.stats();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,15 +33,16 @@ void listAndPurgeByTask(@TempDir Path dir) throws Exception {
queue.enqueue(beta, "3");

CountDownLatch dead = new CountDownLatch(3);
try (Worker worker = queue.worker()
Worker worker = queue.worker()
.handle(alpha, (String p) -> {
throw new IllegalStateException("boom");
})
.handle(beta, (String p) -> {
throw new IllegalStateException("boom");
})
.on(EventName.DEAD, event -> dead.countDown())
.start()) {
.start();
try (worker) {
assertTrue(dead.await(20, TimeUnit.SECONDS), "all three should dead-letter");
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ void dependentWaitsForDependency(@TempDir Path dir) throws Exception {
String idA = queue.enqueue(a, "a");
queue.enqueue(b, "b", EnqueueOptions.builder().dependsOn(idA).build());

try (Worker worker = queue.worker()
Worker worker = queue.worker()
.handle(a, p -> {
aDone.set(true);
return p;
Expand All @@ -45,7 +45,8 @@ void dependentWaitsForDependency(@TempDir Path dir) throws Exception {
bRan.countDown();
return p;
})
.start()) {
.start();
try (worker) {
assertTrue(bRan.await(20, TimeUnit.SECONDS), "dependent should run once its dependency completes");
assertTrue(bSawADone.get(), "the dependency must complete before the dependent runs");
}
Expand All @@ -65,15 +66,16 @@ void cascadeCancelsDependentWhenDependencyFails(@TempDir Path dir) throws Except
String idB = queue.enqueue(
b, "b", EnqueueOptions.builder().dependsOn(idA).build());

try (Worker worker = queue.worker()
Worker worker = queue.worker()
.handle(a, p -> {
throw new RuntimeException("boom");
})
.handle(b, p -> {
bRan.set(true);
return p;
})
.start()) {
.start();
try (worker) {
Job jobB = queue.awaitJob(idB, Duration.ofSeconds(20)).orElseThrow();
assertEquals(JobStatus.CANCELLED, jobB.status, "a failed dependency cascade-cancels the dependent");
assertFalse(bRan.get(), "the dependent must never run");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,11 +38,12 @@ private static List<Integer> runAndRecord(Taskito queue, int total) throws Excep
for (int i = 0; i < total; i++) {
queue.enqueue(rec, i);
}
try (Worker worker = queue.worker()
Worker worker = queue.worker()
.concurrency(1)
.batchSize(1)
.handle(rec, order::add)
.start()) {
.start();
try (worker) {
long deadline = System.nanoTime() + Duration.ofSeconds(20).toNanos();
while (System.nanoTime() < deadline && order.size() < total) {
Thread.sleep(50);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,8 @@ void awaitJobReturnsTerminalState(@TempDir Path dir) throws Exception {
try (Taskito queue =
Taskito.builder().sqlite(dir.resolve("erg.db").toString()).open()) {
String id = queue.enqueue(echo, 42);
try (Worker worker = queue.worker().handle(echo, p -> p).start()) {
Worker worker = queue.worker().handle(echo, p -> p).start();
try (worker) {
Job job = queue.awaitJob(id, Duration.ofSeconds(20)).orElseThrow();
assertEquals(JobStatus.COMPLETE, job.status);
assertEquals(42, queue.getResult(id, Integer.class).orElseThrow());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,13 +34,14 @@ void convertsPayload(@TempDir Path dir) throws Exception {
queue.intercept((task, payload) -> Interception.convert((Integer) payload * 10));
AtomicInteger seen = new AtomicInteger();
CountDownLatch ran = new CountDownLatch(1);
try (Worker worker = queue.worker()
Worker worker = queue.worker()
.handle(A, p -> {
seen.set(p);
ran.countDown();
return p;
})
.start()) {
.start();
try (worker) {
queue.enqueue(A, 5);
assertTrue(ran.await(20, TimeUnit.SECONDS));
assertEquals(50, seen.get());
Expand All @@ -57,13 +58,14 @@ void redirectsTask(@TempDir Path dir) throws Exception {
task.equals("ic.a") ? Interception.redirect("ic.b", payload) : Interception.pass());
AtomicInteger aRan = new AtomicInteger();
CountDownLatch bRan = new CountDownLatch(1);
try (Worker worker = queue.worker()
Worker worker = queue.worker()
.handle(A, p -> aRan.incrementAndGet())
.handle(B, p -> {
bRan.countDown();
return p;
})
.start()) {
.start();
try (worker) {
queue.enqueue(A, 1);
assertTrue(bRan.await(20, TimeUnit.SECONDS));
assertEquals(0, aRan.get());
Expand Down Expand Up @@ -97,13 +99,14 @@ void batchEnqueueConvertsPayloads(@TempDir Path dir) throws Exception {
queue.intercept((task, payload) -> Interception.convert((Integer) payload * 10));
AtomicInteger sum = new AtomicInteger();
CountDownLatch ran = new CountDownLatch(3);
try (Worker worker = queue.worker()
Worker worker = queue.worker()
.handle(A, p -> {
sum.addAndGet(p);
ran.countDown();
return p;
})
.start()) {
.start();
try (worker) {
queue.enqueueMany(A, List.of(1, 2, 3));
assertTrue(ran.await(20, TimeUnit.SECONDS));
assertEquals(60, sum.get()); // (1+2+3)*10 — interception ran on the batch
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,8 @@ void managedConsumersInOneWorker(@TempDir Path dir) throws Exception {
queue.publish("t-skip", value);
}

try (Worker worker = queue.worker().start()) {
Worker worker = queue.worker().start();
try (worker) {
// Plain delivery: every handled message is acked, so the cursor catches up.
pollUntil(() -> invoke.size() == 3);
assertEquals(List.of("a", "b", "c"), invoke);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -148,8 +148,8 @@ void logAndFanoutSubscribersCoexist(@TempDir Path dir) {
queue.subscribe("events", onEvent); // fan-out: one job per publish
queue.subscribeLog("events", "log"); // log: one stored message per publish

try (Worker worker =
queue.worker().handle(onEvent, payload -> payload).start()) {
Worker worker = queue.worker().handle(onEvent, payload -> payload).start();
try (worker) {
List<Job> deliveries = queue.publish("events", "seen");

// The fan-out subscriber ran its job to completion...
Expand Down
26 changes: 16 additions & 10 deletions sdks/java/src/test/java/org/byteveda/taskito/core/PubSubTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -224,8 +224,9 @@ void ephemeralSubscriptionRegistersWithARunningWorker(@TempDir Path dir) throws
// No owning worker yet, so the subscription is not registered.
assertTrue(queue.publish("orders", "o-1").isEmpty());

try (Worker worker =
queue.worker().handle(SEND_EMAIL, payload -> payload).start()) {
Worker worker =
queue.worker().handle(SEND_EMAIL, payload -> payload).start();
try (worker) {
List<Job> deliveries = queue.publish("orders", "o-2");
assertEquals(1, deliveries.size());
Job done = queue.awaitJob(deliveries.get(0).id, Duration.ofSeconds(10))
Expand All @@ -244,7 +245,8 @@ void reapSparesFreshEphemeralRowsWithinTheRegistrationGrace(@TempDir Path dir) t
Map.of("backend", "sqlite", "dsn", dir.resolve("t.db").toString()));
JniQueueBackend backend = JniQueueBackend.open(options);
try (Taskito queue = Taskito.builder().open(backend)) {
try (Worker worker = queue.worker().start()) {
Worker worker = queue.worker().start();
try (worker) {
String liveWorkerId = queue.listWorkers().get(0).workerId;
backend.registerSubscription("orders", "live", SEND_EMAIL.name(), "default", false, liveWorkerId);
backend.registerSubscription("orders", "ghost", SEND_EMAIL.name(), "default", false, "gone-worker");
Expand All @@ -270,7 +272,8 @@ void ephemeralRegistrationWithoutAnOwnerIsRejected(@TempDir Path dir) throws Exc
.writeValueAsString(
Map.of("backend", "sqlite", "dsn", dir.resolve("t.db").toString()));
JniQueueBackend backend = JniQueueBackend.open(options);
try (Taskito queue = Taskito.builder().open(backend)) {
Taskito queue = Taskito.builder().open(backend);
try (queue) {
TaskitoException rejected = assertThrows(
TaskitoException.class,
() -> backend.registerSubscription("orders", "sink", SEND_EMAIL.name(), "default", false, null));
Expand All @@ -290,8 +293,9 @@ void pauseSurvivesRedeclareAndWorkerStart(@TempDir Path dir) {
assertTrue(queue.publish("orders", "o-1").isEmpty());

// Neither must a worker start re-registering the declaration.
try (Worker worker =
queue.worker().handle(SEND_EMAIL, payload -> payload).start()) {
Worker worker =
queue.worker().handle(SEND_EMAIL, payload -> payload).start();
try (worker) {
assertTrue(queue.publish("orders", "o-2").isEmpty());
}
}
Expand All @@ -303,8 +307,9 @@ void unsubscribeDropsTheLocalDeclarationSoWorkerStartDoesNotResurrectIt(@TempDir
try (Taskito queue = open(dir)) {
queue.subscribe("orders", SEND_EMAIL);
assertTrue(queue.unsubscribe("orders", SEND_EMAIL.name()));
try (Worker worker =
queue.worker().handle(SEND_EMAIL, payload -> payload).start()) {
Worker worker =
queue.worker().handle(SEND_EMAIL, payload -> payload).start();
try (worker) {
assertTrue(queue.listSubscriptions("orders").isEmpty());
assertTrue(queue.publish("orders", "o-1").isEmpty());
}
Expand Down Expand Up @@ -383,12 +388,13 @@ void topicStatsCountsAFailedDeliveryAsDead(@TempDir Path dir) throws Exception {
"orders", flaky, SubscriptionOptions.builder().name("flaky").build());

CountDownLatch dead = new CountDownLatch(1);
try (Worker worker = queue.worker()
Worker worker = queue.worker()
.handle(flaky, (String payload) -> {
throw new IllegalStateException("boom");
})
.on(EventName.DEAD, event -> dead.countDown())
.start()) {
.start();
try (worker) {
queue.publish("orders", "o-1");
assertTrue(dead.await(20, TimeUnit.SECONDS), "the delivery should dead-letter");
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,13 +27,14 @@ void requeuesStuckRunningJob(@TempDir Path dir) throws Exception {
CountDownLatch running = new CountDownLatch(1);
CountDownLatch release = new CountDownLatch(1);
String id = queue.enqueue(TASK, 1);
try (Worker worker = queue.worker()
Worker worker = queue.worker()
.handle(TASK, payload -> {
running.countDown();
release.await();
return payload;
})
.start()) {
.start();
try (worker) {
assertTrue(running.await(20, TimeUnit.SECONDS), "handler did not start");
// Pause first, so the poller can't re-claim the job the moment the
// requeue makes it pending again and flip the status under the assert.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ void mixedBatchReportsEveryJobOnce(@TempDir Path dir) throws Exception {
queue.enqueue(boom, String.valueOf(i));
}

try (Worker worker = queue.worker()
Worker worker = queue.worker()
.concurrency(8)
.batchSize(each * 2)
.handle(ok, (String p) -> p)
Expand All @@ -60,7 +60,8 @@ void mixedBatchReportsEveryJobOnce(@TempDir Path dir) throws Exception {
dead.add(event.jobId);
finished.countDown();
})
.start()) {
.start();
try (worker) {
assertTrue(finished.await(20, TimeUnit.SECONDS), "every job should report an outcome");
}
}
Expand Down
Loading
Loading