-
Notifications
You must be signed in to change notification settings - Fork 4
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
5 changed files
with
162 additions
and
1 deletion.
There are no files selected for viewing
82 changes: 82 additions & 0 deletions
82
plugin/src/software/aws/toolkits/eclipse/amazonq/broker/EventBroker.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,82 @@ | ||
// Copyright 2024 Amazon.com, Inc. or its affiliates. All Rights Reserved. | ||
// SPDX-License-Identifier: Apache-2.0 | ||
|
||
package software.aws.toolkits.eclipse.amazonq.broker; | ||
|
||
import java.util.Map; | ||
import java.util.concurrent.ConcurrentHashMap; | ||
import java.util.concurrent.Flow.Subscription; | ||
import java.util.concurrent.SubmissionPublisher; | ||
import java.util.concurrent.atomic.AtomicReference; | ||
|
||
import software.aws.toolkits.eclipse.amazonq.subscriber.Subscriber; | ||
|
||
public final class EventBroker { | ||
|
||
private static final EventBroker INSTANCE; | ||
private final Map<Class<?>, SubmissionPublisher<?>> publishers; | ||
|
||
static { | ||
INSTANCE = new EventBroker(); | ||
} | ||
|
||
private EventBroker() { | ||
publishers = new ConcurrentHashMap<>(); | ||
} | ||
|
||
public static EventBroker getInstance() { | ||
return INSTANCE; | ||
} | ||
|
||
public <T> void post(final T event) { | ||
if (event == null) { | ||
return; | ||
} | ||
|
||
@SuppressWarnings("unchecked") | ||
SubmissionPublisher<T> publisher = (SubmissionPublisher<T>) getPublisher(event.getClass()); | ||
publisher.submit(event); | ||
} | ||
|
||
public <T> Subscription subscribe(final Subscriber<T> subscriber) { | ||
SubmissionPublisher<T> publisher = getPublisher(subscriber.getSubscriptionEventClass()); | ||
AtomicReference<Subscription> subscriptionReference = new AtomicReference<>(); | ||
|
||
java.util.concurrent.Flow.Subscriber<T> subscriberWrapper = new java.util.concurrent.Flow.Subscriber<>() { | ||
private java.util.concurrent.Flow.Subscription subscription; | ||
|
||
@Override | ||
public void onSubscribe(final java.util.concurrent.Flow.Subscription subscription) { | ||
this.subscription = subscription; | ||
subscriptionReference.set(subscription); | ||
this.subscription.request(1); | ||
} | ||
|
||
@Override | ||
public void onNext(final T event) { | ||
subscriber.handleEvent(event); | ||
this.subscription.request(1); | ||
} | ||
|
||
@Override | ||
public void onError(final Throwable throwable) { | ||
subscriber.handleError(throwable); | ||
} | ||
|
||
@Override | ||
public void onComplete() { | ||
// TODO: add if required | ||
} | ||
}; | ||
|
||
publisher.subscribe(subscriberWrapper); | ||
return subscriptionReference.get(); | ||
} | ||
|
||
@SuppressWarnings("unchecked") | ||
private <T> SubmissionPublisher<T> getPublisher(final Class<T> eventType) { | ||
return (SubmissionPublisher<T>) publishers.computeIfAbsent(eventType, | ||
key -> new SubmissionPublisher<>()); | ||
} | ||
|
||
} |
17 changes: 17 additions & 0 deletions
17
plugin/src/software/aws/toolkits/eclipse/amazonq/events/TestEvent.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,17 @@ | ||
// Copyright 2024 Amazon.com, Inc. or its affiliates. All Rights Reserved. | ||
// SPDX-License-Identifier: Apache-2.0 | ||
|
||
package software.aws.toolkits.eclipse.amazonq.events; | ||
|
||
public final class TestEvent { | ||
private final String message; | ||
|
||
public TestEvent(final String message) { | ||
this.message = message; | ||
} | ||
|
||
public String getMessage() { | ||
return message; | ||
} | ||
|
||
} |
28 changes: 28 additions & 0 deletions
28
plugin/src/software/aws/toolkits/eclipse/amazonq/publishers/TestPublisher.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,28 @@ | ||
// Copyright 2024 Amazon.com, Inc. or its affiliates. All Rights Reserved. | ||
// SPDX-License-Identifier: Apache-2.0 | ||
|
||
package software.aws.toolkits.eclipse.amazonq.publishers; | ||
|
||
import software.aws.toolkits.eclipse.amazonq.broker.EventBroker; | ||
import software.aws.toolkits.eclipse.amazonq.events.TestEvent; | ||
|
||
public final class TestPublisher { | ||
|
||
public TestPublisher() { | ||
Thread publisherThread = new Thread(() -> { | ||
try { | ||
Thread.sleep(5000); | ||
EventBroker eventBroker = EventBroker.getInstance(); | ||
|
||
for (int i = 0; i < 10; i++) { | ||
eventBroker.post(new TestEvent("Test Event " + i)); | ||
} | ||
} catch (InterruptedException e) { | ||
Thread.currentThread().interrupt(); | ||
} | ||
}, "TestPublisher-Thread"); | ||
|
||
publisherThread.start(); | ||
} | ||
|
||
} |
13 changes: 13 additions & 0 deletions
13
plugin/src/software/aws/toolkits/eclipse/amazonq/subscriber/Subscriber.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,13 @@ | ||
// Copyright 2024 Amazon.com, Inc. or its affiliates. All Rights Reserved. | ||
// SPDX-License-Identifier: Apache-2.0 | ||
|
||
package software.aws.toolkits.eclipse.amazonq.subscriber; | ||
|
||
public interface Subscriber<T> { | ||
|
||
Class<T> getSubscriptionEventClass(); | ||
|
||
void handleEvent(T event); | ||
void handleError(Throwable error); | ||
|
||
} |
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