forked from quarkusio/quarkus
-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
WebSockets Next: fire CDI event when a connection is opened/closed
- resolves quarkusio#40217
- Loading branch information
Showing
5 changed files
with
224 additions
and
0 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
98 changes: 98 additions & 0 deletions
98
...t/src/test/java/io/quarkus/websockets/next/test/openconnections/ConnectionEventsTest.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,98 @@ | ||
package io.quarkus.websockets.next.test.openconnections; | ||
|
||
import static org.junit.jupiter.api.Assertions.assertEquals; | ||
import static org.junit.jupiter.api.Assertions.assertNotNull; | ||
import static org.junit.jupiter.api.Assertions.assertTrue; | ||
|
||
import java.net.URI; | ||
import java.util.concurrent.CountDownLatch; | ||
import java.util.concurrent.TimeUnit; | ||
import java.util.concurrent.atomic.AtomicReference; | ||
|
||
import jakarta.enterprise.event.ObservesAsync; | ||
import jakarta.inject.Inject; | ||
import jakarta.inject.Singleton; | ||
|
||
import org.junit.jupiter.api.Test; | ||
import org.junit.jupiter.api.extension.RegisterExtension; | ||
|
||
import io.quarkus.test.QuarkusUnitTest; | ||
import io.quarkus.test.common.http.TestHTTPResource; | ||
import io.quarkus.websockets.next.Closed; | ||
import io.quarkus.websockets.next.OnClose; | ||
import io.quarkus.websockets.next.OnOpen; | ||
import io.quarkus.websockets.next.Open; | ||
import io.quarkus.websockets.next.WebSocket; | ||
import io.quarkus.websockets.next.WebSocketConnection; | ||
import io.quarkus.websockets.next.test.utils.WSClient; | ||
import io.vertx.core.Vertx; | ||
|
||
public class ConnectionEventsTest { | ||
|
||
@RegisterExtension | ||
public static final QuarkusUnitTest test = new QuarkusUnitTest() | ||
.withApplicationRoot(root -> { | ||
root.addClasses(Endpoint.class, ObservingBean.class, WSClient.class); | ||
}); | ||
|
||
@Inject | ||
Vertx vertx; | ||
|
||
@TestHTTPResource("endpoint") | ||
URI endUri; | ||
|
||
@Test | ||
void testEvents() throws Exception { | ||
String client1ConnectionId; | ||
try (WSClient client1 = WSClient.create(vertx).connect(endUri)) { | ||
client1.waitForMessages(1); | ||
client1ConnectionId = client1.getMessages().get(0).toString(); | ||
assertTrue(ObservingBean.OPEN_LATCH.await(5, TimeUnit.SECONDS)); | ||
assertNotNull(ObservingBean.OPEN_CONN.get()); | ||
assertEquals(client1ConnectionId, ObservingBean.OPEN_CONN.get().id()); | ||
} | ||
assertTrue(Endpoint.CLOSED_LATCH.await(5, TimeUnit.SECONDS)); | ||
assertTrue(ObservingBean.CLOSED_LATCH.await(5, TimeUnit.SECONDS)); | ||
assertNotNull(ObservingBean.CLOSED_CONN.get()); | ||
assertEquals(client1ConnectionId, ObservingBean.CLOSED_CONN.get().id()); | ||
} | ||
|
||
@WebSocket(path = "/endpoint") | ||
public static class Endpoint { | ||
|
||
static final CountDownLatch CLOSED_LATCH = new CountDownLatch(1); | ||
|
||
@OnOpen | ||
String open(WebSocketConnection connection) { | ||
return connection.id(); | ||
} | ||
|
||
@OnClose | ||
void close() { | ||
CLOSED_LATCH.countDown(); | ||
} | ||
|
||
} | ||
|
||
@Singleton | ||
public static class ObservingBean { | ||
|
||
static final CountDownLatch OPEN_LATCH = new CountDownLatch(1); | ||
static final CountDownLatch CLOSED_LATCH = new CountDownLatch(1); | ||
|
||
static final AtomicReference<WebSocketConnection> OPEN_CONN = new AtomicReference<>(); | ||
static final AtomicReference<WebSocketConnection> CLOSED_CONN = new AtomicReference<>(); | ||
|
||
void onOpen(@ObservesAsync @Open WebSocketConnection connection) { | ||
OPEN_CONN.set(connection); | ||
OPEN_LATCH.countDown(); | ||
} | ||
|
||
void onClose(@ObservesAsync @Closed WebSocketConnection connection) { | ||
CLOSED_CONN.set(connection); | ||
CLOSED_LATCH.countDown(); | ||
} | ||
|
||
} | ||
|
||
} |
41 changes: 41 additions & 0 deletions
41
extensions/websockets-next/runtime/src/main/java/io/quarkus/websockets/next/Closed.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,41 @@ | ||
package io.quarkus.websockets.next; | ||
|
||
import static java.lang.annotation.ElementType.FIELD; | ||
import static java.lang.annotation.ElementType.METHOD; | ||
import static java.lang.annotation.ElementType.PARAMETER; | ||
import static java.lang.annotation.ElementType.TYPE; | ||
import static java.lang.annotation.RetentionPolicy.RUNTIME; | ||
|
||
import java.lang.annotation.Documented; | ||
import java.lang.annotation.Retention; | ||
import java.lang.annotation.Target; | ||
|
||
import jakarta.enterprise.event.Event; | ||
import jakarta.enterprise.event.ObservesAsync; | ||
import jakarta.enterprise.util.AnnotationLiteral; | ||
import jakarta.inject.Qualifier; | ||
|
||
/** | ||
* A CDI event of type {@link WebSocketConnection} with this qualifier is fired asynchronously when a connection is closed. | ||
* | ||
* @see ObservesAsync | ||
* @see Event#fireAsync(Object) | ||
*/ | ||
@Qualifier | ||
@Documented | ||
@Retention(RUNTIME) | ||
@Target({ METHOD, FIELD, PARAMETER, TYPE }) | ||
public @interface Closed { | ||
|
||
/** | ||
* Supports inline instantiation of the {@link Closed} qualifier. | ||
*/ | ||
public static final class Literal extends AnnotationLiteral<Closed> implements Closed { | ||
|
||
public static final Literal INSTANCE = new Literal(); | ||
|
||
private static final long serialVersionUID = 1L; | ||
|
||
} | ||
|
||
} |
41 changes: 41 additions & 0 deletions
41
extensions/websockets-next/runtime/src/main/java/io/quarkus/websockets/next/Open.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,41 @@ | ||
package io.quarkus.websockets.next; | ||
|
||
import static java.lang.annotation.ElementType.FIELD; | ||
import static java.lang.annotation.ElementType.METHOD; | ||
import static java.lang.annotation.ElementType.PARAMETER; | ||
import static java.lang.annotation.ElementType.TYPE; | ||
import static java.lang.annotation.RetentionPolicy.RUNTIME; | ||
|
||
import java.lang.annotation.Documented; | ||
import java.lang.annotation.Retention; | ||
import java.lang.annotation.Target; | ||
|
||
import jakarta.enterprise.event.Event; | ||
import jakarta.enterprise.event.ObservesAsync; | ||
import jakarta.enterprise.util.AnnotationLiteral; | ||
import jakarta.inject.Qualifier; | ||
|
||
/** | ||
* A CDI event of type {@link WebSocketConnection} with this qualifier is fired asynchronously when a new connection is opened. | ||
* | ||
* @see ObservesAsync | ||
* @see Event#fireAsync(Object) | ||
*/ | ||
@Qualifier | ||
@Documented | ||
@Retention(RUNTIME) | ||
@Target({ METHOD, FIELD, PARAMETER, TYPE }) | ||
public @interface Open { | ||
|
||
/** | ||
* Supports inline instantiation of the {@link Open} qualifier. | ||
*/ | ||
public static final class Literal extends AnnotationLiteral<Open> implements Open { | ||
|
||
public static final Literal INSTANCE = new Literal(); | ||
|
||
private static final long serialVersionUID = 1L; | ||
|
||
} | ||
|
||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters