Skip to content

Commit e789328

Browse files
author
Jay Herron
committed
Changes observer to resolve to a future
1 parent 82f3c21 commit e789328

2 files changed

Lines changed: 19 additions & 15 deletions

File tree

Sources/GraphQL/Subscription/Subscribe.swift

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -59,15 +59,15 @@ func subscribe(
5959
return sourceFuture.map{ subscriptionResult -> SubscriptionResult in
6060
do {
6161
let subscriptionObserver = try subscriptionResult.get()
62-
let eventObserver = subscriptionObserver.map { eventPayload -> GraphQLResult in
62+
let eventObserver = subscriptionObserver.map { eventPayload -> Future<GraphQLResult> in
6363

6464
// For each payload yielded from a subscription, map it over the normal
6565
// GraphQL `execute` function, with `payload` as the rootValue.
6666
// This implements the "MapSourceToResponseEvent" algorithm described in
6767
// the GraphQL specification. The `execute` function provides the
6868
// "ExecuteSubscriptionEvent" algorithm, as it is nearly identical to the
6969
// "ExecuteQuery" algorithm, for which `execute` is also used.
70-
let eventResolved = try execute(
70+
return execute(
7171
queryStrategy: queryStrategy,
7272
mutationStrategy: mutationStrategy,
7373
subscriptionStrategy: subscriptionStrategy,
@@ -79,10 +79,8 @@ func subscribe(
7979
eventLoopGroup: eventLoopGroup,
8080
variableValues: variableValues,
8181
operationName: operationName
82-
).wait() // TODO remove this wait
83-
return eventResolved
82+
)
8483
}
85-
// TODO Making a future here feels it indicates a mistake...
8684
return SubscriptionResult.success(eventObserver)
8785
} catch let graphQLError as GraphQLError {
8886
return SubscriptionResult.failure(graphQLError)
@@ -253,19 +251,21 @@ func executeSubscription(
253251
return SourceEventStreamResult.failure(context.errors.first!)
254252
} else if let error = resolved as? GraphQLError {
255253
return SourceEventStreamResult.failure(error)
256-
} else if let observable = resolved as? Observable<Any> {
254+
} else if let observable = resolved as? SourceEventStreamObservable {
257255
return SourceEventStreamResult.success(observable)
258256
} else if resolved == nil {
259257
return SourceEventStreamResult.failure(
260258
GraphQLError(message: "Resolved subscription was nil")
261259
)
262260
} else {
263261
return SourceEventStreamResult.failure(
264-
GraphQLError(message: "Subscription field resolver must return an Observable<Any>, not \(Swift.type(of:resolved))")
262+
GraphQLError(message: "Subscription field resolver must return an SourceEventStreamObservable, not \(Swift.type(of:resolved))")
265263
)
266264
}
267265
}
268266
}
269267

270-
typealias SubscriptionResult = Result<Observable<GraphQLResult>, GraphQLError>
271-
typealias SourceEventStreamResult = Result<Observable<Any>, GraphQLError>
268+
typealias SubscriptionObservable = Observable<Future<GraphQLResult>>
269+
typealias SubscriptionResult = Result<SubscriptionObservable, GraphQLError>
270+
typealias SourceEventStreamObservable = Observable<Any>
271+
typealias SourceEventStreamResult = Result<SourceEventStreamObservable, GraphQLError>

Tests/GraphQLTests/Subscription/SubscriptionTests.swift

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,8 @@ class SubscriptionTests : XCTestCase {
2525
]]
2626
)
2727
let _ = subscription.subscribe { event in
28-
XCTAssertEqual(event.element, expected)
28+
let payload = try! event.element!.wait()
29+
XCTAssertEqual(payload, expected)
2930
}.disposed(by: disposeBag)
3031
pubsub.onNext(Email(
3132
from: "yuzhi@graphql.org",
@@ -105,7 +106,8 @@ class SubscriptionTests : XCTestCase {
105106
]]
106107
)
107108
let _ = subscription.subscribe { event in
108-
XCTAssertEqual(event.element, expected)
109+
let payload = try! event.element!.wait()
110+
XCTAssertEqual(payload, expected)
109111
}.disposed(by: disposeBag)
110112
pubsub.onNext(Email(
111113
from: "yuzhi@graphql.org",
@@ -158,7 +160,9 @@ class SubscriptionTests : XCTestCase {
158160
}
159161
""")
160162

161-
let _ = subscription.subscribe().disposed(by: disposeBag)
163+
let _ = subscription.subscribe{ event in
164+
let _ = try! event.element!.wait()
165+
}.disposed(by: disposeBag)
162166
pubsub.onNext(Email(
163167
from: "yuzhi@graphql.org",
164168
subject: "Alright",
@@ -278,7 +282,7 @@ let defaultEmails = [
278282
private func createDbAndSubscription(
279283
pubsub:Observable<Any>,
280284
query:String
281-
) throws -> Observable<GraphQLResult> {
285+
) throws -> SubscriptionObservable {
282286

283287
var emails = defaultEmails
284288

@@ -304,7 +308,7 @@ private func createSubscription(
304308
pubsub: Observable<Any>,
305309
schema: GraphQLSchema,
306310
query: String
307-
) throws -> Observable<GraphQLResult> {
311+
) throws -> SubscriptionObservable {
308312
let document = try parse(source: query)
309313
let eventLoopGroup = MultiThreadedEventLoopGroup(numberOfThreads: 1)
310314
let subscriptionOrError = try subscribe(
@@ -344,7 +348,7 @@ private func emailSchemaWithResolvers(resolve: GraphQLFieldResolve?, subscribe:
344348
)
345349
}
346350

347-
private func extractSubscription(_ subscriptionResult: SubscriptionResult) throws -> Observable<GraphQLResult> {
351+
private func extractSubscription(_ subscriptionResult: SubscriptionResult) throws -> SubscriptionObservable {
348352
switch subscriptionResult {
349353
case .success(let subscription):
350354
return subscription

0 commit comments

Comments
 (0)