11package studio .o7 .octopus .plugin ;
22
3+ import gentle .Error ;
4+ import gentle .Result ;
35import io .grpc .stub .StreamObserver ;
4- import it .unimi .dsi .fastutil .Pair ;
5- import it .unimi .dsi .fastutil .objects .Object2ObjectArrayMap ;
6- import it .unimi .dsi .fastutil .objects .Object2ObjectMap ;
7- import lombok .NonNull ;
86import lombok .extern .slf4j .Slf4j ;
9- import org .jetbrains .annotations .NotNull ;
10- import org .jetbrains .annotations .Nullable ;
117import studio .o7 .octopus .plugin .api .Octopus ;
12- import studio .o7 .octopus .plugin .api .listener . Listener ;
13- import studio .o7 .octopus .plugin .observer . EmptyObserver ;
14- import studio .o7 .octopus .plugin .utils . ProtoUtils ;
15- import studio .o7 .octopus .sdk . OctopusSDK ;
16- import studio .o7 .octopus .sdk . gen . api . v1 .* ;
17- import studio .o7 .octopus .sdk . gen . api . v1 . Object ;
18-
19- import java . time . Instant ;
20- import java . util . Collection ;
8+ import studio .o7 .octopus .plugin .api .EventHandler ;
9+ import studio .o7 .octopus .plugin .api . OctopusError ;
10+ import studio .o7 .octopus .plugin .api . QueryParameter ;
11+ import studio .o7 .octopus .plugin . authentication . OctopusCredentials ;
12+ import studio .o7 .octopus .plugin . channel . OctopusChannelFactory ;
13+ import studio .o7 .octopus .plugin . observer . OctopusObserver ;
14+ import studio . o7 . octopus . sdk . v1 .*;
15+ import studio . o7 . octopus . sdk . v1 . Object ;
16+
2117import java .util .UUID ;
22- import java .util .concurrent .atomic .AtomicReference ;
2318
24- @ Slf4j (topic = "OctopusPlugin" )
19+ @ Slf4j (topic = "OctopusPlugin" )
2520public final class OctopusImpl implements Octopus {
26- private static final EmptyObserver EMPTY_OBSERVER = new EmptyObserver ();
2721
28- private final OctopusGrpc .OctopusStub stub = OctopusSDK .stub ();
29- private final OctopusGrpc .OctopusBlockingStub blockingStub = OctopusSDK .blockingStub ();
30- private final Object2ObjectMap <UUID , Pair <Listener , StreamObserver <ListenMessage >>> listeners = new Object2ObjectArrayMap <>();
22+ private final OctopusObserver streamObserver ;
23+ private StreamObserver <ListenMessage > responseObserver ;
24+
25+ private final OctopusGrpc .OctopusBlockingStub blocking ;
26+ private final OctopusGrpc .OctopusStub async ;
27+
28+ public OctopusImpl (String token , String host , int port ) {
29+ var channel = OctopusChannelFactory .getOrCreate (host , port );
30+
31+ var blockingStub = OctopusGrpc .newBlockingStub (channel );
32+ this .blocking = blockingStub .withCallCredentials (new OctopusCredentials (token ));
33+
34+ var asyncStub = OctopusGrpc .newStub (channel );
35+ async = asyncStub .withCallCredentials (new OctopusCredentials (token ));
36+
37+ this .streamObserver = new OctopusObserver ();
38+ }
3139
3240 @ Override
33- public @ NotNull Collection <Entry > get (@ NonNull String keyPattern , boolean includeExpired , @ Nullable Instant createdRangeStart , @ Nullable Instant createdRangeEnd ) {
34- var builder = GetRequest .newBuilder ();
41+ public Result <Object , Error > get (String key ) {
42+ var request = GetRequest .newBuilder ().setKey (key ).build ();
43+ try {
44+ var response = blocking .get (request );
45+ return Result .ok (response .getObject ());
46+ } catch (RuntimeException e ) {
47+ log .error ("Failed to send a get request: {}" , e .getMessage ());
48+ return Result .err (OctopusError .GET_REQUEST_FAILED );
49+ }
50+ }
3551
36- builder .setKeyPattern (keyPattern );
37- builder .setIncludeExpired (includeExpired );
52+ @ Override
53+ public Result <QueryResponse , Error > query (QueryParameter queryParameter ) {
54+ var request = QueryRequest .newBuilder ();
55+ request .setKeyPattern (queryParameter .getKeyPattern ());
56+ request .setDataFilter (queryParameter .getDataFilter ());
57+
58+ request .setIncludeExpired (queryParameter .isIncludeExpired ());
59+
60+ var paginator = Paginator .newBuilder ().setPage (queryParameter .getPage ()).setPageSize (queryParameter .getPageSize ()).build ();
61+
62+ request .setPaginator (paginator );
3863
39- if (createdRangeStart != null )
40- builder .setCreatedAtRangeStart (ProtoUtils .toProto (createdRangeStart ));
64+ if (queryParameter .getCreatedAtStart () != null ) {
65+ request .setCreatedAtRangeStart (queryParameter .getCreatedAtStart ());
66+ }
4167
42- if (createdRangeEnd != null )
43- builder .setCreatedAtRangeEnd (ProtoUtils .toProto (createdRangeEnd ));
68+ if (queryParameter .getCreatedAtEnd () != null ) {
69+ request .setCreatedAtRangeEnd (queryParameter .getCreatedAtEnd ());
70+ }
4471
45- return this .blockingStub .get (builder .build ()).getEntriesList ();
72+ try {
73+ var response = blocking .query (request .build ());
74+ return Result .ok (response );
75+ } catch (RuntimeException e ) {
76+ log .error ("Failed to send a query request: {}" , e .getMessage ());
77+ return Result .err (OctopusError .QUERY_REQUEST_FAILED );
78+ }
4679 }
4780
4881 @ Override
49- public void registerListener (@ NonNull Listener listener ) {
50- var requestRef = new AtomicReference <StreamObserver <ListenMessage >>();
51-
52- var observer = stub .listen (new StreamObserver <>() {
53- @ Override
54- public void onNext (EventCall value ) {
55- var start = System .currentTimeMillis ();
56- var request = requestRef .get ();
57- if (request == null ) return ;
58-
59- listener .onCall (value .getObject ());
60-
61- var msg = ListenMessage .newBuilder ()
62- .setCallback (value )
63- .build ();
64-
65- request = requestRef .get ();
66- if (request == null ) return ;
67-
68- request .onNext (msg );
69- log .debug ("Finished EventCall `{}` in {}ms" , value .getCallId (), System .currentTimeMillis () - start );
70- }
71-
72- @ Override
73- public void onError (Throwable t ) {
74- requestRef .set (null );
75- log .error ("Cannot call event on listener {} with key-pattern {}" , listener .getListenerUniqueId (), listener .getKeyPattern (), t );
76- unregisterListener (listener );
77- }
78-
79- @ Override
80- public void onCompleted () {
81- requestRef .set (null );
82- unregisterListener (listener );
83- log .debug ("Completed listener `{}`" , listener .getListenerUniqueId ());
84- }
85- });
86-
87- requestRef .set (observer );
88-
89- observer .onNext (ListenMessage .newBuilder ()
90- .setRegister (ListenRegister .newBuilder ()
91- .setKeyPattern (listener .getKeyPattern ())
92- .setPriority (listener .getPriority ())
93- .build ()).build ());
94-
95- this .listeners .put (listener .getListenerUniqueId (), new Pair <>() {
96- @ Override
97- public Listener left () {
98- return listener ;
99- }
100-
101- @ Override
102- public StreamObserver <ListenMessage > right () {
103- return observer ;
104- }
105- });
82+ public Result <Entry , Error > call (studio .o7 .octopus .sdk .v1 .Object obj ) {
83+ try {
84+ var response = blocking .call (obj );
85+ return Result .ok (response );
86+ } catch (RuntimeException e ) {
87+ log .error ("Failed to send a call request: {}" , e .getMessage ());
88+ return Result .err (OctopusError .CALL_REQUEST_FAILED );
89+ }
10690 }
10791
10892 @ Override
109- public void unregisterListener (@ NonNull Listener listener ) {
110- unregisterListener (listener .getListenerUniqueId ());
93+ public void write (Object obj ) {
94+ try {
95+ blocking .write (obj );
96+ } catch (RuntimeException e ) {
97+ log .error ("Failed to send a write request: {}" , e .getMessage ());
98+ }
11199 }
112100
113101 @ Override
114- public void unregisterListener (@ NonNull UUID listenerUniqueId ) {
115- var pair = this .listeners .get (listenerUniqueId );
116- if (pair == null ) return ;
117- this .listeners .remove (listenerUniqueId );
118- var right = pair .right ();
119- if (right == null ) return ;
120- right .onCompleted ();
102+ public void registerHandler (EventHandler eventHandler ) {
103+ log .debug ("Adding a new handler to the pattern {}" , eventHandler .getKeyPattern ());
104+ streamObserver .addHandler (eventHandler );
105+
106+ if (responseObserver == null ) {
107+ log .debug ("initializing stream to octopus" );
108+ this .responseObserver = async .listen (streamObserver );
109+ }
110+
111+ var req = ListenMessage .newBuilder ().addAllKeyPattern (streamObserver .getKeys ());
112+ responseObserver .onNext (req .build ());
121113 }
122114
123115 @ Override
124- public @ NotNull Entry call ( @ NonNull Object obj ) {
125- return blockingStub . call ( obj );
116+ public void unregisterHandler ( @ org . jspecify . annotations . NonNull EventHandler eventHandler ) {
117+ this . unregisterHandler ( eventHandler . getListenerUniqueId () );
126118 }
127119
128120 @ Override
129- public void callAndForget ( @ NonNull Object obj ) {
130- stub . write ( obj , EMPTY_OBSERVER );
121+ public void unregisterHandler ( UUID listenerUniqueId ) {
122+ this . streamObserver . removeHandler ( listenerUniqueId );
131123 }
124+
132125}
0 commit comments