@@ -8,7 +8,8 @@ import RxSwift
88
99/// Server implements the server-side portion of the protocol, allowing a few callbacks for customization.
1010public class Server {
11- let messenger : Messenger
11+ // We keep this weak because we strongly inject this object into the messenger callback
12+ weak var messenger : Messenger ?
1213
1314 let onExecute : ( GraphQLRequest ) -> EventLoopFuture < GraphQLResult >
1415 let onSubscribe : ( GraphQLRequest ) -> EventLoopFuture < SubscriptionResult >
@@ -38,8 +39,8 @@ public class Server {
3839 self . onExecute = onExecute
3940 self . onSubscribe = onSubscribe
4041
41- self . messenger. onRecieve { [ weak self ] message in
42- guard let self = self else { return }
42+ messenger. onRecieve { message in
43+ guard let messenger = self . messenger else { return }
4344
4445 self . onMessage ( message)
4546
@@ -51,7 +52,7 @@ public class Server {
5152
5253 guard let json = message. data ( using: . utf8) else {
5354 let error = GraphQLWSError . invalidEncoding ( )
54- self . messenger. error ( error. message, code: error. code)
55+ messenger. error ( error. message, code: error. code)
5556 return
5657 }
5758
@@ -61,42 +62,42 @@ public class Server {
6162 }
6263 catch {
6364 let error = GraphQLWSError . noType ( )
64- self . messenger. error ( error. message, code: error. code)
65+ messenger. error ( error. message, code: error. code)
6566 return
6667 }
6768
6869 switch request. type {
6970 case . GQL_CONNECTION_INIT:
7071 guard let connectionInitRequest = try ? self . decoder. decode ( ConnectionInitRequest . self, from: json) else {
7172 let error = GraphQLWSError . invalidRequestFormat ( messageType: . GQL_CONNECTION_INIT)
72- self . messenger. error ( error. message, code: error. code)
73+ messenger. error ( error. message, code: error. code)
7374 return
7475 }
75- self . onConnectionInit ( connectionInitRequest)
76+ self . onConnectionInit ( connectionInitRequest, messenger )
7677 case . GQL_START:
7778 guard let startRequest = try ? self . decoder. decode ( StartRequest . self, from: json) else {
7879 let error = GraphQLWSError . invalidRequestFormat ( messageType: . GQL_START)
79- self . messenger. error ( error. message, code: error. code)
80+ messenger. error ( error. message, code: error. code)
8081 return
8182 }
82- self . onStart ( startRequest)
83+ self . onStart ( startRequest, messenger )
8384 case . GQL_STOP:
8485 guard let stopRequest = try ? self . decoder. decode ( StopRequest . self, from: json) else {
8586 let error = GraphQLWSError . invalidRequestFormat ( messageType: . GQL_STOP)
86- self . messenger. error ( error. message, code: error. code)
87+ messenger. error ( error. message, code: error. code)
8788 return
8889 }
89- self . onStop ( stopRequest, self . messenger)
90+ self . onStop ( stopRequest, messenger)
9091 case . GQL_CONNECTION_TERMINATE:
9192 guard let connectionTerminateRequest = try ? self . decoder. decode ( ConnectionTerminateRequest . self, from: json) else {
9293 let error = GraphQLWSError . invalidRequestFormat ( messageType: . GQL_CONNECTION_TERMINATE)
93- self . messenger. error ( error. message, code: error. code)
94+ messenger. error ( error. message, code: error. code)
9495 return
9596 }
96- self . onConnectionTerminate ( connectionTerminateRequest)
97+ self . onConnectionTerminate ( connectionTerminateRequest, messenger )
9798 case . unknown:
9899 let error = GraphQLWSError . invalidType ( )
99- self . messenger. error ( error. message, code: error. code)
100+ messenger. error ( error. message, code: error. code)
100101 }
101102 }
102103
@@ -126,7 +127,7 @@ public class Server {
126127 self . onMessage = callback
127128 }
128129
129- private func onConnectionInit( _ connectionInitRequest: ConnectionInitRequest ) {
130+ private func onConnectionInit( _ connectionInitRequest: ConnectionInitRequest , _ messenger : Messenger ) {
130131 guard !initialized else {
131132 let error = GraphQLWSError . tooManyInitializations ( )
132133 messenger. error ( error. message, code: error. code)
@@ -148,7 +149,7 @@ public class Server {
148149 // TODO: Should we send the `ka` message?
149150 }
150151
151- private func onStart( _ startRequest: StartRequest ) {
152+ private func onStart( _ startRequest: StartRequest , _ messenger : Messenger ) {
152153 guard initialized else {
153154 let error = GraphQLWSError . notInitialized ( )
154155 messenger. error ( error. message, code: error. code)
@@ -169,48 +170,50 @@ public class Server {
169170
170171 if isStreaming {
171172 let subscribeFuture = onSubscribe ( graphQLRequest)
172- subscribeFuture. whenSuccess { [ weak self] result in
173- guard let self = self else { return }
173+ subscribeFuture. whenSuccess { result in
174174 guard let streamOpt = result. stream else {
175175 // API issue - subscribe resolver isn't stream
176176 let error = GraphQLWSError . internalAPIStreamIssue ( )
177- self . messenger. error ( error. message, code: error. code)
177+ messenger. error ( error. message, code: error. code)
178178 return
179179 }
180180 let stream = streamOpt as! ObservableSubscriptionEventStream
181181 let observable = stream. observable
182182 observable. subscribe (
183- onNext: { resultFuture in
183+ onNext: { [ weak self] resultFuture in
184+ guard let self = self , let messenger = self . messenger else { return }
184185 resultFuture. whenSuccess { result in
185- self . messenger. send ( DataResponse ( result, id: id) . toJSON ( self . encoder) )
186+ messenger. send ( DataResponse ( result, id: id) . toJSON ( self . encoder) )
186187 }
187188 resultFuture. whenFailure { error in
188- self . messenger. send ( ErrorResponse ( error, id: id) . toJSON ( self . encoder) )
189+ messenger. send ( ErrorResponse ( error, id: id) . toJSON ( self . encoder) )
189190 }
190191 } ,
191- onError: { error in
192- self . messenger. send ( ErrorResponse ( error, id: id) . toJSON ( self . encoder) )
192+ onError: { [ weak self] error in
193+ guard let self = self , let messenger = self . messenger else { return }
194+ messenger. send ( ErrorResponse ( error, id: id) . toJSON ( self . encoder) )
193195 } ,
194- onCompleted: {
195- self . messenger. send ( CompleteResponse ( id: id) . toJSON ( self . encoder) )
196- _ = self . messenger. close ( )
196+ onCompleted: { [ weak self] in
197+ guard let self = self , let messenger = self . messenger else { return }
198+ messenger. send ( CompleteResponse ( id: id) . toJSON ( self . encoder) )
199+ _ = messenger. close ( )
197200 }
198201 ) . disposed ( by: self . disposeBag)
199202 }
200203 subscribeFuture. whenFailure { error in
201204 let error = GraphQLWSError . graphQLError ( error)
202- _ = self . messenger. error ( error. message, code: error. code)
205+ _ = messenger. error ( error. message, code: error. code)
203206 }
204207 }
205208 else {
206209 let executeFuture = onExecute ( graphQLRequest)
207210 executeFuture. whenSuccess { result in
208- self . messenger. send ( DataResponse ( result, id: id) . toJSON ( self . encoder) )
209- self . messenger. send ( CompleteResponse ( id: id) . toJSON ( self . encoder) )
211+ messenger. send ( DataResponse ( result, id: id) . toJSON ( self . encoder) )
212+ messenger. send ( CompleteResponse ( id: id) . toJSON ( self . encoder) )
210213 }
211214 executeFuture. whenFailure { error in
212- self . messenger. send ( ErrorResponse ( error, id: id) . toJSON ( self . encoder) )
213- self . messenger. send ( CompleteResponse ( id: id) . toJSON ( self . encoder) )
215+ messenger. send ( ErrorResponse ( error, id: id) . toJSON ( self . encoder) )
216+ messenger. send ( CompleteResponse ( id: id) . toJSON ( self . encoder) )
214217 }
215218 }
216219 }
@@ -224,7 +227,7 @@ public class Server {
224227 onExit ( )
225228 }
226229
227- private func onConnectionTerminate( _: ConnectionTerminateRequest ) {
230+ private func onConnectionTerminate( _: ConnectionTerminateRequest , _ messenger : Messenger ) {
228231 onExit ( )
229232 _ = messenger. close ( )
230233 }
0 commit comments