Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 22 additions & 18 deletions Sources/GraphQLTransportWS/Client.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,43 +50,52 @@ public actor Client<InitPayload: Equatable & Codable> {
do {
response = try decoder.decode(Response.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

switch response.type {
case .connectionAck:
guard
let connectionAckResponse = try? decoder.decode(
let connectionAckResponse: ConnectionAckResponse
do {
connectionAckResponse = try decoder.decode(
ConnectionAckResponse.self,
from: message
)
else {
try await error(.invalidResponseFormat(messageType: .connectionAck))
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .connectionAck, error: error))
return
}
try await onConnectionAck(connectionAckResponse, self)
case .next:
guard let nextResponse = try? decoder.decode(NextResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .next))
let nextResponse: NextResponse
do {
nextResponse = try decoder.decode(NextResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .next, error: error))
return
}
try await onNext(nextResponse, self)
case .error:
guard let errorResponse = try? decoder.decode(ErrorResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .error))
let errorResponse: ErrorResponse
do {
errorResponse = try decoder.decode(ErrorResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .error, error: error))
return
}
try await onError(errorResponse, self)
case .complete:
guard let completeResponse = try? decoder.decode(CompleteResponse.self, from: message)
else {
try await error(.invalidResponseFormat(messageType: .complete))
let completeResponse: CompleteResponse
do {
completeResponse = try decoder.decode(CompleteResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .complete, error: error))
return
}
try await onComplete(completeResponse, self)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand DownExpand Up@@ -123,9 +132,4 @@ public actor Client<InitPayload: Equatable & Codable> {
)
)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
8 changes: 4 additions & 4 deletions Sources/GraphQLTransportWS/GraphqlTransportWSError.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,16 +58,16 @@ struct GraphQLTransportWSError: Error {
)
}

static func invalidRequestFormat(messageType: RequestMessageType) -> Self {
static func invalidRequestFormat(messageType: RequestMessageType, error: Error) -> Self {
return self.init(
"Request message doesn't match '\(messageType.type.rawValue)' JSON format",
"Request message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}

static func invalidResponseFormat(messageType: ResponseMessageType) -> Self {
static func invalidResponseFormat(messageType: ResponseMessageType, error: Error) -> Self {
return self.init(
"Response message doesn't match '\(messageType.type.rawValue)' JSON format",
"Response message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}
Expand Down
7 changes: 7 additions & 0 deletions Sources/GraphQLTransportWS/Messenger.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,3 +15,10 @@ public protocol Messenger: Sendable {
/// - code: An error code
func error(_ message: String, code: Int) async throws
}

extension Messenger {
/// Send an error through the messenger and close the connection
func error(_ error: GraphQLTransportWSError) async throws {
try await self.error(error.message, code: error.code.rawValue)
}
}
105 changes: 52 additions & 53 deletions Sources/GraphQLTransportWS/Server.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ where

private var initialized = false
private var initResult: InitPayloadResult?
private var subscriptionTasks = [String: Task<Void, any Error>]()
private var executionTasks = [String: Task<Void, any Error>]()

/// Create a new server
///
Expand DownExpand Up@@ -53,7 +53,7 @@ where
}

deinit {
subscriptionTasks.values.forEach { $0.cancel() }
executionTasks.values.forEach { $0.cancel() }
}

/// Listen and react to the provided async sequence of client messages. This function will block until the stream is completed.
Expand All@@ -70,39 +70,44 @@ where
do {
request = try decoder.decode(Request.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

// handle incoming message
switch request.type {
case .connectionInit:
guard
let connectionInitRequest = try? decoder.decode(
let connectionInitRequest: ConnectionInitRequest<InitPayload>
do {
connectionInitRequest = try decoder.decode(
ConnectionInitRequest<InitPayload>.self,
from: message
)
else {
try await error(.invalidRequestFormat(messageType: .connectionInit))
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .connectionInit, error: error))
return
}
try await onConnectionInit(connectionInitRequest, messenger)
case .subscribe:
guard let subscribeRequest = try? decoder.decode(SubscribeRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .subscribe))
let subscribeRequest: SubscribeRequest
do {
subscribeRequest = try decoder.decode(SubscribeRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .subscribe, error: error))
return
}
try await onSubscribe(subscribeRequest)
case .complete:
guard let completeRequest = try? decoder.decode(CompleteRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .complete))
let completeRequest: CompleteRequest
do {
completeRequest = try decoder.decode(CompleteRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .complete, error: error))
return
}
try await onOperationComplete(completeRequest)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand All@@ -111,14 +116,14 @@ where
_: Messenger
) async throws {
guard !initialized else {
try await error(.tooManyInitializations())
try await messenger.error(.tooManyInitializations())
return
}

do {
initResult = try await onInit(connectionInitRequest.payload)
} catch {
try await self.error(.forbidden())
try await messenger.error(.forbidden())
return
}
initialized = true
Expand All@@ -128,62 +133,66 @@ where

private func onSubscribe(_ subscribeRequest: SubscribeRequest) async throws {
guard initialized, let initResult else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = subscribeRequest.id
if subscriptionTasks[id] != nil {
try await error(.subscriberAlreadyExists(id: id))
}

let graphQLRequest = subscribeRequest.payload

var isStreaming = false
let isStreaming: Bool
do {
isStreaming = try graphQLRequest.isSubscription()
} catch {
try await sendError(error, id: id)
return
}

if isStreaming {
subscriptionTasks[id] = Task {
guard executionTasks[id] == nil else {
try await messenger.error(.subscriberAlreadyExists(id: id))
return
}
executionTasks[id] = Task {
defer {
executionTasks.removeValue(forKey: id)
}

if isStreaming {
let stream: SubscriptionSequenceType
do {
let stream = try await onSubscribe(graphQLRequest, initResult)
for try await event in stream {
try Task.checkCancellation()
try await self.sendNext(event, id: id)
}
stream = try await onSubscribe(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
subscriptionTasks.removeValue(forKey: id)
throw error
return
}
for try await event in stream {
try await self.sendNext(event, id: id)
}
executionTasks.removeValue(forKey: id)
} else {
let result: GraphQLResult
do {
result = try await onExecute(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
return
}
try await self.sendComplete(id: id)
subscriptionTasks.removeValue(forKey: id)
}
} else {
do {
let result = try await onExecute(graphQLRequest, initResult)
try await sendNext(result, id: id)
try await sendComplete(id: id)
} catch {
try await sendError(error, id: id)
}
try await sendComplete(id: id)
}
}

private func onOperationComplete(_ completeRequest: CompleteRequest) async throws {
guard initialized else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = completeRequest.id
if let task = subscriptionTasks[id] {
if let task = executionTasks[id] {
task.cancel()
subscriptionTasks.removeValue(forKey: id)
executionTasks.removeValue(forKey: id)
}
try await onOperationComplete(id)
}
Expand DownExpand Up@@ -238,14 +247,4 @@ where
private func sendError(_ error: Error, id: String) async throws {
try await sendError([error], id: id)
}

/// Send an `error` response through the messenger
private func sendError(_ errorMessage: String, id: String) async throws {
try await sendError(GraphQLError(message: errorMessage), id: id)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
56 changes: 56 additions & 0 deletions Tests/GraphQLTransportWSTests/GraphQLTransportWSTests.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -268,6 +268,62 @@ struct GraphqlTransportWSTests {
)
}

/// Tests malformed requests include decoder details in the transport error
@Test func malformedRequestIncludesDecodingDetails() async throws {
let api = TestAPI()
let context = TestContext()
let server = Server<TokenInitPayload, Void, AsyncThrowingStream<GraphQLResult, Error>>(
messenger: serverMessenger,
onInit: { _ in },
onExecute: { graphQLRequest, _ in
try await api.execute(
request: graphQLRequest.query,
context: context
)
},
onSubscribe: { graphQLRequest, _ in
try await api.subscribe(
request: graphQLRequest.query,
context: context
).get()
}
)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"complete"}"#.utf8))
continuation.finish()

try await server.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in serverMessenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Request message doesn't match 'complete' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

/// Tests malformed responses include decoder details in the transport error
@Test func malformedResponseIncludesDecodingDetails() async throws {
let messenger = TestMessenger()
let client = Client<TokenInitPayload>(messenger: messenger)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"next"}"#.utf8))
continuation.finish()

try await client.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in messenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Response message doesn't match 'next' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

enum TestError: Error {
case couldBeAnything
}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 22 additions & 18 deletions Sources/GraphQLTransportWS/Client.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,43 +50,52 @@ public actor Client<InitPayload: Equatable & Codable> {
do {
response = try decoder.decode(Response.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

switch response.type {
case .connectionAck:
guard
let connectionAckResponse = try? decoder.decode(
let connectionAckResponse: ConnectionAckResponse
do {
connectionAckResponse = try decoder.decode(
ConnectionAckResponse.self,
from: message
)
else {
try await error(.invalidResponseFormat(messageType: .connectionAck))
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .connectionAck, error: error))
return
}
try await onConnectionAck(connectionAckResponse, self)
case .next:
guard let nextResponse = try? decoder.decode(NextResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .next))
let nextResponse: NextResponse
do {
nextResponse = try decoder.decode(NextResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .next, error: error))
return
}
try await onNext(nextResponse, self)
case .error:
guard let errorResponse = try? decoder.decode(ErrorResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .error))
let errorResponse: ErrorResponse
do {
errorResponse = try decoder.decode(ErrorResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .error, error: error))
return
}
try await onError(errorResponse, self)
case .complete:
guard let completeResponse = try? decoder.decode(CompleteResponse.self, from: message)
else {
try await error(.invalidResponseFormat(messageType: .complete))
let completeResponse: CompleteResponse
do {
completeResponse = try decoder.decode(CompleteResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .complete, error: error))
return
}
try await onComplete(completeResponse, self)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand DownExpand Up@@ -123,9 +132,4 @@ public actor Client<InitPayload: Equatable & Codable> {
)
)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
8 changes: 4 additions & 4 deletions Sources/GraphQLTransportWS/GraphqlTransportWSError.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,16 +58,16 @@ struct GraphQLTransportWSError: Error {
)
}

static func invalidRequestFormat(messageType: RequestMessageType) -> Self {
static func invalidRequestFormat(messageType: RequestMessageType, error: Error) -> Self {
return self.init(
"Request message doesn't match '\(messageType.type.rawValue)' JSON format",
"Request message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}

static func invalidResponseFormat(messageType: ResponseMessageType) -> Self {
static func invalidResponseFormat(messageType: ResponseMessageType, error: Error) -> Self {
return self.init(
"Response message doesn't match '\(messageType.type.rawValue)' JSON format",
"Response message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}
Expand Down
7 changes: 7 additions & 0 deletions Sources/GraphQLTransportWS/Messenger.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,3 +15,10 @@ public protocol Messenger: Sendable {
/// - code: An error code
func error(_ message: String, code: Int) async throws
}

extension Messenger {
/// Send an error through the messenger and close the connection
func error(_ error: GraphQLTransportWSError) async throws {
try await self.error(error.message, code: error.code.rawValue)
}
}
105 changes: 52 additions & 53 deletions Sources/GraphQLTransportWS/Server.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ where

private var initialized = false
private var initResult: InitPayloadResult?
private var subscriptionTasks = [String: Task<Void, any Error>]()
private var executionTasks = [String: Task<Void, any Error>]()

/// Create a new server
///
Expand DownExpand Up@@ -53,7 +53,7 @@ where
}

deinit {
subscriptionTasks.values.forEach { $0.cancel() }
executionTasks.values.forEach { $0.cancel() }
}

/// Listen and react to the provided async sequence of client messages. This function will block until the stream is completed.
Expand All@@ -70,39 +70,44 @@ where
do {
request = try decoder.decode(Request.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

// handle incoming message
switch request.type {
case .connectionInit:
guard
let connectionInitRequest = try? decoder.decode(
let connectionInitRequest: ConnectionInitRequest<InitPayload>
do {
connectionInitRequest = try decoder.decode(
ConnectionInitRequest<InitPayload>.self,
from: message
)
else {
try await error(.invalidRequestFormat(messageType: .connectionInit))
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .connectionInit, error: error))
return
}
try await onConnectionInit(connectionInitRequest, messenger)
case .subscribe:
guard let subscribeRequest = try? decoder.decode(SubscribeRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .subscribe))
let subscribeRequest: SubscribeRequest
do {
subscribeRequest = try decoder.decode(SubscribeRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .subscribe, error: error))
return
}
try await onSubscribe(subscribeRequest)
case .complete:
guard let completeRequest = try? decoder.decode(CompleteRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .complete))
let completeRequest: CompleteRequest
do {
completeRequest = try decoder.decode(CompleteRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .complete, error: error))
return
}
try await onOperationComplete(completeRequest)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand All@@ -111,14 +116,14 @@ where
_: Messenger
) async throws {
guard !initialized else {
try await error(.tooManyInitializations())
try await messenger.error(.tooManyInitializations())
return
}

do {
initResult = try await onInit(connectionInitRequest.payload)
} catch {
try await self.error(.forbidden())
try await messenger.error(.forbidden())
return
}
initialized = true
Expand All@@ -128,62 +133,66 @@ where

private func onSubscribe(_ subscribeRequest: SubscribeRequest) async throws {
guard initialized, let initResult else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = subscribeRequest.id
if subscriptionTasks[id] != nil {
try await error(.subscriberAlreadyExists(id: id))
}

let graphQLRequest = subscribeRequest.payload

var isStreaming = false
let isStreaming: Bool
do {
isStreaming = try graphQLRequest.isSubscription()
} catch {
try await sendError(error, id: id)
return
}

if isStreaming {
subscriptionTasks[id] = Task {
guard executionTasks[id] == nil else {
try await messenger.error(.subscriberAlreadyExists(id: id))
return
}
executionTasks[id] = Task {
defer {
executionTasks.removeValue(forKey: id)
}

if isStreaming {
let stream: SubscriptionSequenceType
do {
let stream = try await onSubscribe(graphQLRequest, initResult)
for try await event in stream {
try Task.checkCancellation()
try await self.sendNext(event, id: id)
}
stream = try await onSubscribe(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
subscriptionTasks.removeValue(forKey: id)
throw error
return
}
for try await event in stream {
try await self.sendNext(event, id: id)
}
executionTasks.removeValue(forKey: id)
} else {
let result: GraphQLResult
do {
result = try await onExecute(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
return
}
try await self.sendComplete(id: id)
subscriptionTasks.removeValue(forKey: id)
}
} else {
do {
let result = try await onExecute(graphQLRequest, initResult)
try await sendNext(result, id: id)
try await sendComplete(id: id)
} catch {
try await sendError(error, id: id)
}
try await sendComplete(id: id)
}
}

private func onOperationComplete(_ completeRequest: CompleteRequest) async throws {
guard initialized else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = completeRequest.id
if let task = subscriptionTasks[id] {
if let task = executionTasks[id] {
task.cancel()
subscriptionTasks.removeValue(forKey: id)
executionTasks.removeValue(forKey: id)
}
try await onOperationComplete(id)
}
Expand DownExpand Up@@ -238,14 +247,4 @@ where
private func sendError(_ error: Error, id: String) async throws {
try await sendError([error], id: id)
}

/// Send an `error` response through the messenger
private func sendError(_ errorMessage: String, id: String) async throws {
try await sendError(GraphQLError(message: errorMessage), id: id)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
56 changes: 56 additions & 0 deletions Tests/GraphQLTransportWSTests/GraphQLTransportWSTests.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -268,6 +268,62 @@ struct GraphqlTransportWSTests {
)
}

/// Tests malformed requests include decoder details in the transport error
@Test func malformedRequestIncludesDecodingDetails() async throws {
let api = TestAPI()
let context = TestContext()
let server = Server<TokenInitPayload, Void, AsyncThrowingStream<GraphQLResult, Error>>(
messenger: serverMessenger,
onInit: { _ in },
onExecute: { graphQLRequest, _ in
try await api.execute(
request: graphQLRequest.query,
context: context
)
},
onSubscribe: { graphQLRequest, _ in
try await api.subscribe(
request: graphQLRequest.query,
context: context
).get()
}
)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"complete"}"#.utf8))
continuation.finish()

try await server.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in serverMessenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Request message doesn't match 'complete' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

/// Tests malformed responses include decoder details in the transport error
@Test func malformedResponseIncludesDecodingDetails() async throws {
let messenger = TestMessenger()
let client = Client<TokenInitPayload>(messenger: messenger)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"next"}"#.utf8))
continuation.finish()

try await client.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in messenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Response message doesn't match 'next' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

enum TestError: Error {
case couldBeAnything
}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 22 additions & 18 deletions Sources/GraphQLTransportWS/Client.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,43 +50,52 @@ public actor Client<InitPayload: Equatable & Codable> {
do {
response = try decoder.decode(Response.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

switch response.type {
case .connectionAck:
guard
let connectionAckResponse = try? decoder.decode(
let connectionAckResponse: ConnectionAckResponse
do {
connectionAckResponse = try decoder.decode(
ConnectionAckResponse.self,
from: message
)
else {
try await error(.invalidResponseFormat(messageType: .connectionAck))
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .connectionAck, error: error))
return
}
try await onConnectionAck(connectionAckResponse, self)
case .next:
guard let nextResponse = try? decoder.decode(NextResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .next))
let nextResponse: NextResponse
do {
nextResponse = try decoder.decode(NextResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .next, error: error))
return
}
try await onNext(nextResponse, self)
case .error:
guard let errorResponse = try? decoder.decode(ErrorResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .error))
let errorResponse: ErrorResponse
do {
errorResponse = try decoder.decode(ErrorResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .error, error: error))
return
}
try await onError(errorResponse, self)
case .complete:
guard let completeResponse = try? decoder.decode(CompleteResponse.self, from: message)
else {
try await error(.invalidResponseFormat(messageType: .complete))
let completeResponse: CompleteResponse
do {
completeResponse = try decoder.decode(CompleteResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .complete, error: error))
return
}
try await onComplete(completeResponse, self)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand DownExpand Up@@ -123,9 +132,4 @@ public actor Client<InitPayload: Equatable & Codable> {
)
)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
8 changes: 4 additions & 4 deletions Sources/GraphQLTransportWS/GraphqlTransportWSError.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,16 +58,16 @@ struct GraphQLTransportWSError: Error {
)
}

static func invalidRequestFormat(messageType: RequestMessageType) -> Self {
static func invalidRequestFormat(messageType: RequestMessageType, error: Error) -> Self {
return self.init(
"Request message doesn't match '\(messageType.type.rawValue)' JSON format",
"Request message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}

static func invalidResponseFormat(messageType: ResponseMessageType) -> Self {
static func invalidResponseFormat(messageType: ResponseMessageType, error: Error) -> Self {
return self.init(
"Response message doesn't match '\(messageType.type.rawValue)' JSON format",
"Response message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}
Expand Down
7 changes: 7 additions & 0 deletions Sources/GraphQLTransportWS/Messenger.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,3 +15,10 @@ public protocol Messenger: Sendable {
/// - code: An error code
func error(_ message: String, code: Int) async throws
}

extension Messenger {
/// Send an error through the messenger and close the connection
func error(_ error: GraphQLTransportWSError) async throws {
try await self.error(error.message, code: error.code.rawValue)
}
}
105 changes: 52 additions & 53 deletions Sources/GraphQLTransportWS/Server.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ where

private var initialized = false
private var initResult: InitPayloadResult?
private var subscriptionTasks = [String: Task<Void, any Error>]()
private var executionTasks = [String: Task<Void, any Error>]()

/// Create a new server
///
Expand DownExpand Up@@ -53,7 +53,7 @@ where
}

deinit {
subscriptionTasks.values.forEach { $0.cancel() }
executionTasks.values.forEach { $0.cancel() }
}

/// Listen and react to the provided async sequence of client messages. This function will block until the stream is completed.
Expand All@@ -70,39 +70,44 @@ where
do {
request = try decoder.decode(Request.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

// handle incoming message
switch request.type {
case .connectionInit:
guard
let connectionInitRequest = try? decoder.decode(
let connectionInitRequest: ConnectionInitRequest<InitPayload>
do {
connectionInitRequest = try decoder.decode(
ConnectionInitRequest<InitPayload>.self,
from: message
)
else {
try await error(.invalidRequestFormat(messageType: .connectionInit))
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .connectionInit, error: error))
return
}
try await onConnectionInit(connectionInitRequest, messenger)
case .subscribe:
guard let subscribeRequest = try? decoder.decode(SubscribeRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .subscribe))
let subscribeRequest: SubscribeRequest
do {
subscribeRequest = try decoder.decode(SubscribeRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .subscribe, error: error))
return
}
try await onSubscribe(subscribeRequest)
case .complete:
guard let completeRequest = try? decoder.decode(CompleteRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .complete))
let completeRequest: CompleteRequest
do {
completeRequest = try decoder.decode(CompleteRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .complete, error: error))
return
}
try await onOperationComplete(completeRequest)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand All@@ -111,14 +116,14 @@ where
_: Messenger
) async throws {
guard !initialized else {
try await error(.tooManyInitializations())
try await messenger.error(.tooManyInitializations())
return
}

do {
initResult = try await onInit(connectionInitRequest.payload)
} catch {
try await self.error(.forbidden())
try await messenger.error(.forbidden())
return
}
initialized = true
Expand All@@ -128,62 +133,66 @@ where

private func onSubscribe(_ subscribeRequest: SubscribeRequest) async throws {
guard initialized, let initResult else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = subscribeRequest.id
if subscriptionTasks[id] != nil {
try await error(.subscriberAlreadyExists(id: id))
}

let graphQLRequest = subscribeRequest.payload

var isStreaming = false
let isStreaming: Bool
do {
isStreaming = try graphQLRequest.isSubscription()
} catch {
try await sendError(error, id: id)
return
}

if isStreaming {
subscriptionTasks[id] = Task {
guard executionTasks[id] == nil else {
try await messenger.error(.subscriberAlreadyExists(id: id))
return
}
executionTasks[id] = Task {
defer {
executionTasks.removeValue(forKey: id)
}

if isStreaming {
let stream: SubscriptionSequenceType
do {
let stream = try await onSubscribe(graphQLRequest, initResult)
for try await event in stream {
try Task.checkCancellation()
try await self.sendNext(event, id: id)
}
stream = try await onSubscribe(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
subscriptionTasks.removeValue(forKey: id)
throw error
return
}
for try await event in stream {
try await self.sendNext(event, id: id)
}
executionTasks.removeValue(forKey: id)
} else {
let result: GraphQLResult
do {
result = try await onExecute(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
return
}
try await self.sendComplete(id: id)
subscriptionTasks.removeValue(forKey: id)
}
} else {
do {
let result = try await onExecute(graphQLRequest, initResult)
try await sendNext(result, id: id)
try await sendComplete(id: id)
} catch {
try await sendError(error, id: id)
}
try await sendComplete(id: id)
}
}

private func onOperationComplete(_ completeRequest: CompleteRequest) async throws {
guard initialized else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = completeRequest.id
if let task = subscriptionTasks[id] {
if let task = executionTasks[id] {
task.cancel()
subscriptionTasks.removeValue(forKey: id)
executionTasks.removeValue(forKey: id)
}
try await onOperationComplete(id)
}
Expand DownExpand Up@@ -238,14 +247,4 @@ where
private func sendError(_ error: Error, id: String) async throws {
try await sendError([error], id: id)
}

/// Send an `error` response through the messenger
private func sendError(_ errorMessage: String, id: String) async throws {
try await sendError(GraphQLError(message: errorMessage), id: id)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
56 changes: 56 additions & 0 deletions Tests/GraphQLTransportWSTests/GraphQLTransportWSTests.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -268,6 +268,62 @@ struct GraphqlTransportWSTests {
)
}

/// Tests malformed requests include decoder details in the transport error
@Test func malformedRequestIncludesDecodingDetails() async throws {
let api = TestAPI()
let context = TestContext()
let server = Server<TokenInitPayload, Void, AsyncThrowingStream<GraphQLResult, Error>>(
messenger: serverMessenger,
onInit: { _ in },
onExecute: { graphQLRequest, _ in
try await api.execute(
request: graphQLRequest.query,
context: context
)
},
onSubscribe: { graphQLRequest, _ in
try await api.subscribe(
request: graphQLRequest.query,
context: context
).get()
}
)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"complete"}"#.utf8))
continuation.finish()

try await server.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in serverMessenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Request message doesn't match 'complete' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

/// Tests malformed responses include decoder details in the transport error
@Test func malformedResponseIncludesDecodingDetails() async throws {
let messenger = TestMessenger()
let client = Client<TokenInitPayload>(messenger: messenger)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"next"}"#.utf8))
continuation.finish()

try await client.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in messenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Response message doesn't match 'next' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

enum TestError: Error {
case couldBeAnything
}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 22 additions & 18 deletions Sources/GraphQLTransportWS/Client.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,43 +50,52 @@ public actor Client<InitPayload: Equatable & Codable> {
do {
response = try decoder.decode(Response.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

switch response.type {
case .connectionAck:
guard
let connectionAckResponse = try? decoder.decode(
let connectionAckResponse: ConnectionAckResponse
do {
connectionAckResponse = try decoder.decode(
ConnectionAckResponse.self,
from: message
)
else {
try await error(.invalidResponseFormat(messageType: .connectionAck))
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .connectionAck, error: error))
return
}
try await onConnectionAck(connectionAckResponse, self)
case .next:
guard let nextResponse = try? decoder.decode(NextResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .next))
let nextResponse: NextResponse
do {
nextResponse = try decoder.decode(NextResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .next, error: error))
return
}
try await onNext(nextResponse, self)
case .error:
guard let errorResponse = try? decoder.decode(ErrorResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .error))
let errorResponse: ErrorResponse
do {
errorResponse = try decoder.decode(ErrorResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .error, error: error))
return
}
try await onError(errorResponse, self)
case .complete:
guard let completeResponse = try? decoder.decode(CompleteResponse.self, from: message)
else {
try await error(.invalidResponseFormat(messageType: .complete))
let completeResponse: CompleteResponse
do {
completeResponse = try decoder.decode(CompleteResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .complete, error: error))
return
}
try await onComplete(completeResponse, self)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand DownExpand Up@@ -123,9 +132,4 @@ public actor Client<InitPayload: Equatable & Codable> {
)
)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
8 changes: 4 additions & 4 deletions Sources/GraphQLTransportWS/GraphqlTransportWSError.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,16 +58,16 @@ struct GraphQLTransportWSError: Error {
)
}

static func invalidRequestFormat(messageType: RequestMessageType) -> Self {
static func invalidRequestFormat(messageType: RequestMessageType, error: Error) -> Self {
return self.init(
"Request message doesn't match '\(messageType.type.rawValue)' JSON format",
"Request message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}

static func invalidResponseFormat(messageType: ResponseMessageType) -> Self {
static func invalidResponseFormat(messageType: ResponseMessageType, error: Error) -> Self {
return self.init(
"Response message doesn't match '\(messageType.type.rawValue)' JSON format",
"Response message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}
Expand Down
7 changes: 7 additions & 0 deletions Sources/GraphQLTransportWS/Messenger.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,3 +15,10 @@ public protocol Messenger: Sendable {
/// - code: An error code
func error(_ message: String, code: Int) async throws
}

extension Messenger {
/// Send an error through the messenger and close the connection
func error(_ error: GraphQLTransportWSError) async throws {
try await self.error(error.message, code: error.code.rawValue)
}
}
105 changes: 52 additions & 53 deletions Sources/GraphQLTransportWS/Server.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ where

private var initialized = false
private var initResult: InitPayloadResult?
private var subscriptionTasks = [String: Task<Void, any Error>]()
private var executionTasks = [String: Task<Void, any Error>]()

/// Create a new server
///
Expand DownExpand Up@@ -53,7 +53,7 @@ where
}

deinit {
subscriptionTasks.values.forEach { $0.cancel() }
executionTasks.values.forEach { $0.cancel() }
}

/// Listen and react to the provided async sequence of client messages. This function will block until the stream is completed.
Expand All@@ -70,39 +70,44 @@ where
do {
request = try decoder.decode(Request.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

// handle incoming message
switch request.type {
case .connectionInit:
guard
let connectionInitRequest = try? decoder.decode(
let connectionInitRequest: ConnectionInitRequest<InitPayload>
do {
connectionInitRequest = try decoder.decode(
ConnectionInitRequest<InitPayload>.self,
from: message
)
else {
try await error(.invalidRequestFormat(messageType: .connectionInit))
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .connectionInit, error: error))
return
}
try await onConnectionInit(connectionInitRequest, messenger)
case .subscribe:
guard let subscribeRequest = try? decoder.decode(SubscribeRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .subscribe))
let subscribeRequest: SubscribeRequest
do {
subscribeRequest = try decoder.decode(SubscribeRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .subscribe, error: error))
return
}
try await onSubscribe(subscribeRequest)
case .complete:
guard let completeRequest = try? decoder.decode(CompleteRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .complete))
let completeRequest: CompleteRequest
do {
completeRequest = try decoder.decode(CompleteRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .complete, error: error))
return
}
try await onOperationComplete(completeRequest)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand All@@ -111,14 +116,14 @@ where
_: Messenger
) async throws {
guard !initialized else {
try await error(.tooManyInitializations())
try await messenger.error(.tooManyInitializations())
return
}

do {
initResult = try await onInit(connectionInitRequest.payload)
} catch {
try await self.error(.forbidden())
try await messenger.error(.forbidden())
return
}
initialized = true
Expand All@@ -128,62 +133,66 @@ where

private func onSubscribe(_ subscribeRequest: SubscribeRequest) async throws {
guard initialized, let initResult else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = subscribeRequest.id
if subscriptionTasks[id] != nil {
try await error(.subscriberAlreadyExists(id: id))
}

let graphQLRequest = subscribeRequest.payload

var isStreaming = false
let isStreaming: Bool
do {
isStreaming = try graphQLRequest.isSubscription()
} catch {
try await sendError(error, id: id)
return
}

if isStreaming {
subscriptionTasks[id] = Task {
guard executionTasks[id] == nil else {
try await messenger.error(.subscriberAlreadyExists(id: id))
return
}
executionTasks[id] = Task {
defer {
executionTasks.removeValue(forKey: id)
}

if isStreaming {
let stream: SubscriptionSequenceType
do {
let stream = try await onSubscribe(graphQLRequest, initResult)
for try await event in stream {
try Task.checkCancellation()
try await self.sendNext(event, id: id)
}
stream = try await onSubscribe(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
subscriptionTasks.removeValue(forKey: id)
throw error
return
}
for try await event in stream {
try await self.sendNext(event, id: id)
}
executionTasks.removeValue(forKey: id)
} else {
let result: GraphQLResult
do {
result = try await onExecute(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
return
}
try await self.sendComplete(id: id)
subscriptionTasks.removeValue(forKey: id)
}
} else {
do {
let result = try await onExecute(graphQLRequest, initResult)
try await sendNext(result, id: id)
try await sendComplete(id: id)
} catch {
try await sendError(error, id: id)
}
try await sendComplete(id: id)
}
}

private func onOperationComplete(_ completeRequest: CompleteRequest) async throws {
guard initialized else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = completeRequest.id
if let task = subscriptionTasks[id] {
if let task = executionTasks[id] {
task.cancel()
subscriptionTasks.removeValue(forKey: id)
executionTasks.removeValue(forKey: id)
}
try await onOperationComplete(id)
}
Expand DownExpand Up@@ -238,14 +247,4 @@ where
private func sendError(_ error: Error, id: String) async throws {
try await sendError([error], id: id)
}

/// Send an `error` response through the messenger
private func sendError(_ errorMessage: String, id: String) async throws {
try await sendError(GraphQLError(message: errorMessage), id: id)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
56 changes: 56 additions & 0 deletions Tests/GraphQLTransportWSTests/GraphQLTransportWSTests.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -268,6 +268,62 @@ struct GraphqlTransportWSTests {
)
}

/// Tests malformed requests include decoder details in the transport error
@Test func malformedRequestIncludesDecodingDetails() async throws {
let api = TestAPI()
let context = TestContext()
let server = Server<TokenInitPayload, Void, AsyncThrowingStream<GraphQLResult, Error>>(
messenger: serverMessenger,
onInit: { _ in },
onExecute: { graphQLRequest, _ in
try await api.execute(
request: graphQLRequest.query,
context: context
)
},
onSubscribe: { graphQLRequest, _ in
try await api.subscribe(
request: graphQLRequest.query,
context: context
).get()
}
)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"complete"}"#.utf8))
continuation.finish()

try await server.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in serverMessenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Request message doesn't match 'complete' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

/// Tests malformed responses include decoder details in the transport error
@Test func malformedResponseIncludesDecodingDetails() async throws {
let messenger = TestMessenger()
let client = Client<TokenInitPayload>(messenger: messenger)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"next"}"#.utf8))
continuation.finish()

try await client.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in messenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Response message doesn't match 'next' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

enum TestError: Error {
case couldBeAnything
}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 22 additions & 18 deletions Sources/GraphQLTransportWS/Client.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,43 +50,52 @@ public actor Client<InitPayload: Equatable & Codable> {
do {
response = try decoder.decode(Response.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

switch response.type {
case .connectionAck:
guard
let connectionAckResponse = try? decoder.decode(
let connectionAckResponse: ConnectionAckResponse
do {
connectionAckResponse = try decoder.decode(
ConnectionAckResponse.self,
from: message
)
else {
try await error(.invalidResponseFormat(messageType: .connectionAck))
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .connectionAck, error: error))
return
}
try await onConnectionAck(connectionAckResponse, self)
case .next:
guard let nextResponse = try? decoder.decode(NextResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .next))
let nextResponse: NextResponse
do {
nextResponse = try decoder.decode(NextResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .next, error: error))
return
}
try await onNext(nextResponse, self)
case .error:
guard let errorResponse = try? decoder.decode(ErrorResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .error))
let errorResponse: ErrorResponse
do {
errorResponse = try decoder.decode(ErrorResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .error, error: error))
return
}
try await onError(errorResponse, self)
case .complete:
guard let completeResponse = try? decoder.decode(CompleteResponse.self, from: message)
else {
try await error(.invalidResponseFormat(messageType: .complete))
let completeResponse: CompleteResponse
do {
completeResponse = try decoder.decode(CompleteResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .complete, error: error))
return
}
try await onComplete(completeResponse, self)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand DownExpand Up@@ -123,9 +132,4 @@ public actor Client<InitPayload: Equatable & Codable> {
)
)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
8 changes: 4 additions & 4 deletions Sources/GraphQLTransportWS/GraphqlTransportWSError.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,16 +58,16 @@ struct GraphQLTransportWSError: Error {
)
}

static func invalidRequestFormat(messageType: RequestMessageType) -> Self {
static func invalidRequestFormat(messageType: RequestMessageType, error: Error) -> Self {
return self.init(
"Request message doesn't match '\(messageType.type.rawValue)' JSON format",
"Request message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}

static func invalidResponseFormat(messageType: ResponseMessageType) -> Self {
static func invalidResponseFormat(messageType: ResponseMessageType, error: Error) -> Self {
return self.init(
"Response message doesn't match '\(messageType.type.rawValue)' JSON format",
"Response message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}
Expand Down
7 changes: 7 additions & 0 deletions Sources/GraphQLTransportWS/Messenger.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,3 +15,10 @@ public protocol Messenger: Sendable {
/// - code: An error code
func error(_ message: String, code: Int) async throws
}

extension Messenger {
/// Send an error through the messenger and close the connection
func error(_ error: GraphQLTransportWSError) async throws {
try await self.error(error.message, code: error.code.rawValue)
}
}
105 changes: 52 additions & 53 deletions Sources/GraphQLTransportWS/Server.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ where

private var initialized = false
private var initResult: InitPayloadResult?
private var subscriptionTasks = [String: Task<Void, any Error>]()
private var executionTasks = [String: Task<Void, any Error>]()

/// Create a new server
///
Expand DownExpand Up@@ -53,7 +53,7 @@ where
}

deinit {
subscriptionTasks.values.forEach { $0.cancel() }
executionTasks.values.forEach { $0.cancel() }
}

/// Listen and react to the provided async sequence of client messages. This function will block until the stream is completed.
Expand All@@ -70,39 +70,44 @@ where
do {
request = try decoder.decode(Request.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

// handle incoming message
switch request.type {
case .connectionInit:
guard
let connectionInitRequest = try? decoder.decode(
let connectionInitRequest: ConnectionInitRequest<InitPayload>
do {
connectionInitRequest = try decoder.decode(
ConnectionInitRequest<InitPayload>.self,
from: message
)
else {
try await error(.invalidRequestFormat(messageType: .connectionInit))
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .connectionInit, error: error))
return
}
try await onConnectionInit(connectionInitRequest, messenger)
case .subscribe:
guard let subscribeRequest = try? decoder.decode(SubscribeRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .subscribe))
let subscribeRequest: SubscribeRequest
do {
subscribeRequest = try decoder.decode(SubscribeRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .subscribe, error: error))
return
}
try await onSubscribe(subscribeRequest)
case .complete:
guard let completeRequest = try? decoder.decode(CompleteRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .complete))
let completeRequest: CompleteRequest
do {
completeRequest = try decoder.decode(CompleteRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .complete, error: error))
return
}
try await onOperationComplete(completeRequest)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand All@@ -111,14 +116,14 @@ where
_: Messenger
) async throws {
guard !initialized else {
try await error(.tooManyInitializations())
try await messenger.error(.tooManyInitializations())
return
}

do {
initResult = try await onInit(connectionInitRequest.payload)
} catch {
try await self.error(.forbidden())
try await messenger.error(.forbidden())
return
}
initialized = true
Expand All@@ -128,62 +133,66 @@ where

private func onSubscribe(_ subscribeRequest: SubscribeRequest) async throws {
guard initialized, let initResult else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = subscribeRequest.id
if subscriptionTasks[id] != nil {
try await error(.subscriberAlreadyExists(id: id))
}

let graphQLRequest = subscribeRequest.payload

var isStreaming = false
let isStreaming: Bool
do {
isStreaming = try graphQLRequest.isSubscription()
} catch {
try await sendError(error, id: id)
return
}

if isStreaming {
subscriptionTasks[id] = Task {
guard executionTasks[id] == nil else {
try await messenger.error(.subscriberAlreadyExists(id: id))
return
}
executionTasks[id] = Task {
defer {
executionTasks.removeValue(forKey: id)
}

if isStreaming {
let stream: SubscriptionSequenceType
do {
let stream = try await onSubscribe(graphQLRequest, initResult)
for try await event in stream {
try Task.checkCancellation()
try await self.sendNext(event, id: id)
}
stream = try await onSubscribe(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
subscriptionTasks.removeValue(forKey: id)
throw error
return
}
for try await event in stream {
try await self.sendNext(event, id: id)
}
executionTasks.removeValue(forKey: id)
} else {
let result: GraphQLResult
do {
result = try await onExecute(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
return
}
try await self.sendComplete(id: id)
subscriptionTasks.removeValue(forKey: id)
}
} else {
do {
let result = try await onExecute(graphQLRequest, initResult)
try await sendNext(result, id: id)
try await sendComplete(id: id)
} catch {
try await sendError(error, id: id)
}
try await sendComplete(id: id)
}
}

private func onOperationComplete(_ completeRequest: CompleteRequest) async throws {
guard initialized else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = completeRequest.id
if let task = subscriptionTasks[id] {
if let task = executionTasks[id] {
task.cancel()
subscriptionTasks.removeValue(forKey: id)
executionTasks.removeValue(forKey: id)
}
try await onOperationComplete(id)
}
Expand DownExpand Up@@ -238,14 +247,4 @@ where
private func sendError(_ error: Error, id: String) async throws {
try await sendError([error], id: id)
}

/// Send an `error` response through the messenger
private func sendError(_ errorMessage: String, id: String) async throws {
try await sendError(GraphQLError(message: errorMessage), id: id)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
56 changes: 56 additions & 0 deletions Tests/GraphQLTransportWSTests/GraphQLTransportWSTests.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -268,6 +268,62 @@ struct GraphqlTransportWSTests {
)
}

/// Tests malformed requests include decoder details in the transport error
@Test func malformedRequestIncludesDecodingDetails() async throws {
let api = TestAPI()
let context = TestContext()
let server = Server<TokenInitPayload, Void, AsyncThrowingStream<GraphQLResult, Error>>(
messenger: serverMessenger,
onInit: { _ in },
onExecute: { graphQLRequest, _ in
try await api.execute(
request: graphQLRequest.query,
context: context
)
},
onSubscribe: { graphQLRequest, _ in
try await api.subscribe(
request: graphQLRequest.query,
context: context
).get()
}
)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"complete"}"#.utf8))
continuation.finish()

try await server.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in serverMessenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Request message doesn't match 'complete' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

/// Tests malformed responses include decoder details in the transport error
@Test func malformedResponseIncludesDecodingDetails() async throws {
let messenger = TestMessenger()
let client = Client<TokenInitPayload>(messenger: messenger)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"next"}"#.utf8))
continuation.finish()

try await client.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in messenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Response message doesn't match 'next' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

enum TestError: Error {
case couldBeAnything
}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 22 additions & 18 deletions Sources/GraphQLTransportWS/Client.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,43 +50,52 @@ public actor Client<InitPayload: Equatable & Codable> {
do {
response = try decoder.decode(Response.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

switch response.type {
case .connectionAck:
guard
let connectionAckResponse = try? decoder.decode(
let connectionAckResponse: ConnectionAckResponse
do {
connectionAckResponse = try decoder.decode(
ConnectionAckResponse.self,
from: message
)
else {
try await error(.invalidResponseFormat(messageType: .connectionAck))
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .connectionAck, error: error))
return
}
try await onConnectionAck(connectionAckResponse, self)
case .next:
guard let nextResponse = try? decoder.decode(NextResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .next))
let nextResponse: NextResponse
do {
nextResponse = try decoder.decode(NextResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .next, error: error))
return
}
try await onNext(nextResponse, self)
case .error:
guard let errorResponse = try? decoder.decode(ErrorResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .error))
let errorResponse: ErrorResponse
do {
errorResponse = try decoder.decode(ErrorResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .error, error: error))
return
}
try await onError(errorResponse, self)
case .complete:
guard let completeResponse = try? decoder.decode(CompleteResponse.self, from: message)
else {
try await error(.invalidResponseFormat(messageType: .complete))
let completeResponse: CompleteResponse
do {
completeResponse = try decoder.decode(CompleteResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .complete, error: error))
return
}
try await onComplete(completeResponse, self)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand DownExpand Up@@ -123,9 +132,4 @@ public actor Client<InitPayload: Equatable & Codable> {
)
)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
8 changes: 4 additions & 4 deletions Sources/GraphQLTransportWS/GraphqlTransportWSError.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,16 +58,16 @@ struct GraphQLTransportWSError: Error {
)
}

static func invalidRequestFormat(messageType: RequestMessageType) -> Self {
static func invalidRequestFormat(messageType: RequestMessageType, error: Error) -> Self {
return self.init(
"Request message doesn't match '\(messageType.type.rawValue)' JSON format",
"Request message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}

static func invalidResponseFormat(messageType: ResponseMessageType) -> Self {
static func invalidResponseFormat(messageType: ResponseMessageType, error: Error) -> Self {
return self.init(
"Response message doesn't match '\(messageType.type.rawValue)' JSON format",
"Response message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}
Expand Down
7 changes: 7 additions & 0 deletions Sources/GraphQLTransportWS/Messenger.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,3 +15,10 @@ public protocol Messenger: Sendable {
/// - code: An error code
func error(_ message: String, code: Int) async throws
}

extension Messenger {
/// Send an error through the messenger and close the connection
func error(_ error: GraphQLTransportWSError) async throws {
try await self.error(error.message, code: error.code.rawValue)
}
}
105 changes: 52 additions & 53 deletions Sources/GraphQLTransportWS/Server.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ where

private var initialized = false
private var initResult: InitPayloadResult?
private var subscriptionTasks = [String: Task<Void, any Error>]()
private var executionTasks = [String: Task<Void, any Error>]()

/// Create a new server
///
Expand DownExpand Up@@ -53,7 +53,7 @@ where
}

deinit {
subscriptionTasks.values.forEach { $0.cancel() }
executionTasks.values.forEach { $0.cancel() }
}

/// Listen and react to the provided async sequence of client messages. This function will block until the stream is completed.
Expand All@@ -70,39 +70,44 @@ where
do {
request = try decoder.decode(Request.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

// handle incoming message
switch request.type {
case .connectionInit:
guard
let connectionInitRequest = try? decoder.decode(
let connectionInitRequest: ConnectionInitRequest<InitPayload>
do {
connectionInitRequest = try decoder.decode(
ConnectionInitRequest<InitPayload>.self,
from: message
)
else {
try await error(.invalidRequestFormat(messageType: .connectionInit))
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .connectionInit, error: error))
return
}
try await onConnectionInit(connectionInitRequest, messenger)
case .subscribe:
guard let subscribeRequest = try? decoder.decode(SubscribeRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .subscribe))
let subscribeRequest: SubscribeRequest
do {
subscribeRequest = try decoder.decode(SubscribeRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .subscribe, error: error))
return
}
try await onSubscribe(subscribeRequest)
case .complete:
guard let completeRequest = try? decoder.decode(CompleteRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .complete))
let completeRequest: CompleteRequest
do {
completeRequest = try decoder.decode(CompleteRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .complete, error: error))
return
}
try await onOperationComplete(completeRequest)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand All@@ -111,14 +116,14 @@ where
_: Messenger
) async throws {
guard !initialized else {
try await error(.tooManyInitializations())
try await messenger.error(.tooManyInitializations())
return
}

do {
initResult = try await onInit(connectionInitRequest.payload)
} catch {
try await self.error(.forbidden())
try await messenger.error(.forbidden())
return
}
initialized = true
Expand All@@ -128,62 +133,66 @@ where

private func onSubscribe(_ subscribeRequest: SubscribeRequest) async throws {
guard initialized, let initResult else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = subscribeRequest.id
if subscriptionTasks[id] != nil {
try await error(.subscriberAlreadyExists(id: id))
}

let graphQLRequest = subscribeRequest.payload

var isStreaming = false
let isStreaming: Bool
do {
isStreaming = try graphQLRequest.isSubscription()
} catch {
try await sendError(error, id: id)
return
}

if isStreaming {
subscriptionTasks[id] = Task {
guard executionTasks[id] == nil else {
try await messenger.error(.subscriberAlreadyExists(id: id))
return
}
executionTasks[id] = Task {
defer {
executionTasks.removeValue(forKey: id)
}

if isStreaming {
let stream: SubscriptionSequenceType
do {
let stream = try await onSubscribe(graphQLRequest, initResult)
for try await event in stream {
try Task.checkCancellation()
try await self.sendNext(event, id: id)
}
stream = try await onSubscribe(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
subscriptionTasks.removeValue(forKey: id)
throw error
return
}
for try await event in stream {
try await self.sendNext(event, id: id)
}
executionTasks.removeValue(forKey: id)
} else {
let result: GraphQLResult
do {
result = try await onExecute(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
return
}
try await self.sendComplete(id: id)
subscriptionTasks.removeValue(forKey: id)
}
} else {
do {
let result = try await onExecute(graphQLRequest, initResult)
try await sendNext(result, id: id)
try await sendComplete(id: id)
} catch {
try await sendError(error, id: id)
}
try await sendComplete(id: id)
}
}

private func onOperationComplete(_ completeRequest: CompleteRequest) async throws {
guard initialized else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = completeRequest.id
if let task = subscriptionTasks[id] {
if let task = executionTasks[id] {
task.cancel()
subscriptionTasks.removeValue(forKey: id)
executionTasks.removeValue(forKey: id)
}
try await onOperationComplete(id)
}
Expand DownExpand Up@@ -238,14 +247,4 @@ where
private func sendError(_ error: Error, id: String) async throws {
try await sendError([error], id: id)
}

/// Send an `error` response through the messenger
private func sendError(_ errorMessage: String, id: String) async throws {
try await sendError(GraphQLError(message: errorMessage), id: id)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
56 changes: 56 additions & 0 deletions Tests/GraphQLTransportWSTests/GraphQLTransportWSTests.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -268,6 +268,62 @@ struct GraphqlTransportWSTests {
)
}

/// Tests malformed requests include decoder details in the transport error
@Test func malformedRequestIncludesDecodingDetails() async throws {
let api = TestAPI()
let context = TestContext()
let server = Server<TokenInitPayload, Void, AsyncThrowingStream<GraphQLResult, Error>>(
messenger: serverMessenger,
onInit: { _ in },
onExecute: { graphQLRequest, _ in
try await api.execute(
request: graphQLRequest.query,
context: context
)
},
onSubscribe: { graphQLRequest, _ in
try await api.subscribe(
request: graphQLRequest.query,
context: context
).get()
}
)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"complete"}"#.utf8))
continuation.finish()

try await server.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in serverMessenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Request message doesn't match 'complete' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

/// Tests malformed responses include decoder details in the transport error
@Test func malformedResponseIncludesDecodingDetails() async throws {
let messenger = TestMessenger()
let client = Client<TokenInitPayload>(messenger: messenger)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"next"}"#.utf8))
continuation.finish()

try await client.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in messenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Response message doesn't match 'next' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

enum TestError: Error {
case couldBeAnything
}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 22 additions & 18 deletions Sources/GraphQLTransportWS/Client.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,43 +50,52 @@ public actor Client<InitPayload: Equatable & Codable> {
do {
response = try decoder.decode(Response.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

switch response.type {
case .connectionAck:
guard
let connectionAckResponse = try? decoder.decode(
let connectionAckResponse: ConnectionAckResponse
do {
connectionAckResponse = try decoder.decode(
ConnectionAckResponse.self,
from: message
)
else {
try await error(.invalidResponseFormat(messageType: .connectionAck))
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .connectionAck, error: error))
return
}
try await onConnectionAck(connectionAckResponse, self)
case .next:
guard let nextResponse = try? decoder.decode(NextResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .next))
let nextResponse: NextResponse
do {
nextResponse = try decoder.decode(NextResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .next, error: error))
return
}
try await onNext(nextResponse, self)
case .error:
guard let errorResponse = try? decoder.decode(ErrorResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .error))
let errorResponse: ErrorResponse
do {
errorResponse = try decoder.decode(ErrorResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .error, error: error))
return
}
try await onError(errorResponse, self)
case .complete:
guard let completeResponse = try? decoder.decode(CompleteResponse.self, from: message)
else {
try await error(.invalidResponseFormat(messageType: .complete))
let completeResponse: CompleteResponse
do {
completeResponse = try decoder.decode(CompleteResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .complete, error: error))
return
}
try await onComplete(completeResponse, self)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand DownExpand Up@@ -123,9 +132,4 @@ public actor Client<InitPayload: Equatable & Codable> {
)
)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
8 changes: 4 additions & 4 deletions Sources/GraphQLTransportWS/GraphqlTransportWSError.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,16 +58,16 @@ struct GraphQLTransportWSError: Error {
)
}

static func invalidRequestFormat(messageType: RequestMessageType) -> Self {
static func invalidRequestFormat(messageType: RequestMessageType, error: Error) -> Self {
return self.init(
"Request message doesn't match '\(messageType.type.rawValue)' JSON format",
"Request message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}

static func invalidResponseFormat(messageType: ResponseMessageType) -> Self {
static func invalidResponseFormat(messageType: ResponseMessageType, error: Error) -> Self {
return self.init(
"Response message doesn't match '\(messageType.type.rawValue)' JSON format",
"Response message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}
Expand Down
7 changes: 7 additions & 0 deletions Sources/GraphQLTransportWS/Messenger.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,3 +15,10 @@ public protocol Messenger: Sendable {
/// - code: An error code
func error(_ message: String, code: Int) async throws
}

extension Messenger {
/// Send an error through the messenger and close the connection
func error(_ error: GraphQLTransportWSError) async throws {
try await self.error(error.message, code: error.code.rawValue)
}
}
105 changes: 52 additions & 53 deletions Sources/GraphQLTransportWS/Server.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ where

private var initialized = false
private var initResult: InitPayloadResult?
private var subscriptionTasks = [String: Task<Void, any Error>]()
private var executionTasks = [String: Task<Void, any Error>]()

/// Create a new server
///
Expand DownExpand Up@@ -53,7 +53,7 @@ where
}

deinit {
subscriptionTasks.values.forEach { $0.cancel() }
executionTasks.values.forEach { $0.cancel() }
}

/// Listen and react to the provided async sequence of client messages. This function will block until the stream is completed.
Expand All@@ -70,39 +70,44 @@ where
do {
request = try decoder.decode(Request.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

// handle incoming message
switch request.type {
case .connectionInit:
guard
let connectionInitRequest = try? decoder.decode(
let connectionInitRequest: ConnectionInitRequest<InitPayload>
do {
connectionInitRequest = try decoder.decode(
ConnectionInitRequest<InitPayload>.self,
from: message
)
else {
try await error(.invalidRequestFormat(messageType: .connectionInit))
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .connectionInit, error: error))
return
}
try await onConnectionInit(connectionInitRequest, messenger)
case .subscribe:
guard let subscribeRequest = try? decoder.decode(SubscribeRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .subscribe))
let subscribeRequest: SubscribeRequest
do {
subscribeRequest = try decoder.decode(SubscribeRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .subscribe, error: error))
return
}
try await onSubscribe(subscribeRequest)
case .complete:
guard let completeRequest = try? decoder.decode(CompleteRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .complete))
let completeRequest: CompleteRequest
do {
completeRequest = try decoder.decode(CompleteRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .complete, error: error))
return
}
try await onOperationComplete(completeRequest)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand All@@ -111,14 +116,14 @@ where
_: Messenger
) async throws {
guard !initialized else {
try await error(.tooManyInitializations())
try await messenger.error(.tooManyInitializations())
return
}

do {
initResult = try await onInit(connectionInitRequest.payload)
} catch {
try await self.error(.forbidden())
try await messenger.error(.forbidden())
return
}
initialized = true
Expand All@@ -128,62 +133,66 @@ where

private func onSubscribe(_ subscribeRequest: SubscribeRequest) async throws {
guard initialized, let initResult else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = subscribeRequest.id
if subscriptionTasks[id] != nil {
try await error(.subscriberAlreadyExists(id: id))
}

let graphQLRequest = subscribeRequest.payload

var isStreaming = false
let isStreaming: Bool
do {
isStreaming = try graphQLRequest.isSubscription()
} catch {
try await sendError(error, id: id)
return
}

if isStreaming {
subscriptionTasks[id] = Task {
guard executionTasks[id] == nil else {
try await messenger.error(.subscriberAlreadyExists(id: id))
return
}
executionTasks[id] = Task {
defer {
executionTasks.removeValue(forKey: id)
}

if isStreaming {
let stream: SubscriptionSequenceType
do {
let stream = try await onSubscribe(graphQLRequest, initResult)
for try await event in stream {
try Task.checkCancellation()
try await self.sendNext(event, id: id)
}
stream = try await onSubscribe(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
subscriptionTasks.removeValue(forKey: id)
throw error
return
}
for try await event in stream {
try await self.sendNext(event, id: id)
}
executionTasks.removeValue(forKey: id)
} else {
let result: GraphQLResult
do {
result = try await onExecute(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
return
}
try await self.sendComplete(id: id)
subscriptionTasks.removeValue(forKey: id)
}
} else {
do {
let result = try await onExecute(graphQLRequest, initResult)
try await sendNext(result, id: id)
try await sendComplete(id: id)
} catch {
try await sendError(error, id: id)
}
try await sendComplete(id: id)
}
}

private func onOperationComplete(_ completeRequest: CompleteRequest) async throws {
guard initialized else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = completeRequest.id
if let task = subscriptionTasks[id] {
if let task = executionTasks[id] {
task.cancel()
subscriptionTasks.removeValue(forKey: id)
executionTasks.removeValue(forKey: id)
}
try await onOperationComplete(id)
}
Expand DownExpand Up@@ -238,14 +247,4 @@ where
private func sendError(_ error: Error, id: String) async throws {
try await sendError([error], id: id)
}

/// Send an `error` response through the messenger
private func sendError(_ errorMessage: String, id: String) async throws {
try await sendError(GraphQLError(message: errorMessage), id: id)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
56 changes: 56 additions & 0 deletions Tests/GraphQLTransportWSTests/GraphQLTransportWSTests.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -268,6 +268,62 @@ struct GraphqlTransportWSTests {
)
}

/// Tests malformed requests include decoder details in the transport error
@Test func malformedRequestIncludesDecodingDetails() async throws {
let api = TestAPI()
let context = TestContext()
let server = Server<TokenInitPayload, Void, AsyncThrowingStream<GraphQLResult, Error>>(
messenger: serverMessenger,
onInit: { _ in },
onExecute: { graphQLRequest, _ in
try await api.execute(
request: graphQLRequest.query,
context: context
)
},
onSubscribe: { graphQLRequest, _ in
try await api.subscribe(
request: graphQLRequest.query,
context: context
).get()
}
)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"complete"}"#.utf8))
continuation.finish()

try await server.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in serverMessenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Request message doesn't match 'complete' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

/// Tests malformed responses include decoder details in the transport error
@Test func malformedResponseIncludesDecodingDetails() async throws {
let messenger = TestMessenger()
let client = Client<TokenInitPayload>(messenger: messenger)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"next"}"#.utf8))
continuation.finish()

try await client.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in messenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Response message doesn't match 'next' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

enum TestError: Error {
case couldBeAnything
}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 22 additions & 18 deletions Sources/GraphQLTransportWS/Client.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,43 +50,52 @@ public actor Client<InitPayload: Equatable & Codable> {
do {
response = try decoder.decode(Response.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

switch response.type {
case .connectionAck:
guard
let connectionAckResponse = try? decoder.decode(
let connectionAckResponse: ConnectionAckResponse
do {
connectionAckResponse = try decoder.decode(
ConnectionAckResponse.self,
from: message
)
else {
try await error(.invalidResponseFormat(messageType: .connectionAck))
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .connectionAck, error: error))
return
}
try await onConnectionAck(connectionAckResponse, self)
case .next:
guard let nextResponse = try? decoder.decode(NextResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .next))
let nextResponse: NextResponse
do {
nextResponse = try decoder.decode(NextResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .next, error: error))
return
}
try await onNext(nextResponse, self)
case .error:
guard let errorResponse = try? decoder.decode(ErrorResponse.self, from: message) else {
try await error(.invalidResponseFormat(messageType: .error))
let errorResponse: ErrorResponse
do {
errorResponse = try decoder.decode(ErrorResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .error, error: error))
return
}
try await onError(errorResponse, self)
case .complete:
guard let completeResponse = try? decoder.decode(CompleteResponse.self, from: message)
else {
try await error(.invalidResponseFormat(messageType: .complete))
let completeResponse: CompleteResponse
do {
completeResponse = try decoder.decode(CompleteResponse.self, from: message)
} catch {
try await messenger.error(.invalidResponseFormat(messageType: .complete, error: error))
return
}
try await onComplete(completeResponse, self)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand DownExpand Up@@ -123,9 +132,4 @@ public actor Client<InitPayload: Equatable & Codable> {
)
)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
8 changes: 4 additions & 4 deletions Sources/GraphQLTransportWS/GraphqlTransportWSError.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,16 +58,16 @@ struct GraphQLTransportWSError: Error {
)
}

static func invalidRequestFormat(messageType: RequestMessageType) -> Self {
static func invalidRequestFormat(messageType: RequestMessageType, error: Error) -> Self {
return self.init(
"Request message doesn't match '\(messageType.type.rawValue)' JSON format",
"Request message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}

static func invalidResponseFormat(messageType: ResponseMessageType) -> Self {
static func invalidResponseFormat(messageType: ResponseMessageType, error: Error) -> Self {
return self.init(
"Response message doesn't match '\(messageType.type.rawValue)' JSON format",
"Response message doesn't match '\(messageType.type.rawValue)' JSON format: \(error)",
code: .miscellaneous
)
}
Expand Down
7 changes: 7 additions & 0 deletions Sources/GraphQLTransportWS/Messenger.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,3 +15,10 @@ public protocol Messenger: Sendable {
/// - code: An error code
func error(_ message: String, code: Int) async throws
}

extension Messenger {
/// Send an error through the messenger and close the connection
func error(_ error: GraphQLTransportWSError) async throws {
try await self.error(error.message, code: error.code.rawValue)
}
}
105 changes: 52 additions & 53 deletions Sources/GraphQLTransportWS/Server.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ where

private var initialized = false
private var initResult: InitPayloadResult?
private var subscriptionTasks = [String: Task<Void, any Error>]()
private var executionTasks = [String: Task<Void, any Error>]()

/// Create a new server
///
Expand DownExpand Up@@ -53,7 +53,7 @@ where
}

deinit {
subscriptionTasks.values.forEach { $0.cancel() }
executionTasks.values.forEach { $0.cancel() }
}

/// Listen and react to the provided async sequence of client messages. This function will block until the stream is completed.
Expand All@@ -70,39 +70,44 @@ where
do {
request = try decoder.decode(Request.self, from: message)
} catch {
try await self.error(.noType())
try await messenger.error(.noType())
return
}

// handle incoming message
switch request.type {
case .connectionInit:
guard
let connectionInitRequest = try? decoder.decode(
let connectionInitRequest: ConnectionInitRequest<InitPayload>
do {
connectionInitRequest = try decoder.decode(
ConnectionInitRequest<InitPayload>.self,
from: message
)
else {
try await error(.invalidRequestFormat(messageType: .connectionInit))
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .connectionInit, error: error))
return
}
try await onConnectionInit(connectionInitRequest, messenger)
case .subscribe:
guard let subscribeRequest = try? decoder.decode(SubscribeRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .subscribe))
let subscribeRequest: SubscribeRequest
do {
subscribeRequest = try decoder.decode(SubscribeRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .subscribe, error: error))
return
}
try await onSubscribe(subscribeRequest)
case .complete:
guard let completeRequest = try? decoder.decode(CompleteRequest.self, from: message)
else {
try await error(.invalidRequestFormat(messageType: .complete))
let completeRequest: CompleteRequest
do {
completeRequest = try decoder.decode(CompleteRequest.self, from: message)
} catch {
try await messenger.error(.invalidRequestFormat(messageType: .complete, error: error))
return
}
try await onOperationComplete(completeRequest)
default:
try await error(.invalidType())
try await messenger.error(.invalidType())
}
}

Expand All@@ -111,14 +116,14 @@ where
_: Messenger
) async throws {
guard !initialized else {
try await error(.tooManyInitializations())
try await messenger.error(.tooManyInitializations())
return
}

do {
initResult = try await onInit(connectionInitRequest.payload)
} catch {
try await self.error(.forbidden())
try await messenger.error(.forbidden())
return
}
initialized = true
Expand All@@ -128,62 +133,66 @@ where

private func onSubscribe(_ subscribeRequest: SubscribeRequest) async throws {
guard initialized, let initResult else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = subscribeRequest.id
if subscriptionTasks[id] != nil {
try await error(.subscriberAlreadyExists(id: id))
}

let graphQLRequest = subscribeRequest.payload

var isStreaming = false
let isStreaming: Bool
do {
isStreaming = try graphQLRequest.isSubscription()
} catch {
try await sendError(error, id: id)
return
}

if isStreaming {
subscriptionTasks[id] = Task {
guard executionTasks[id] == nil else {
try await messenger.error(.subscriberAlreadyExists(id: id))
return
}
executionTasks[id] = Task {
defer {
executionTasks.removeValue(forKey: id)
}

if isStreaming {
let stream: SubscriptionSequenceType
do {
let stream = try await onSubscribe(graphQLRequest, initResult)
for try await event in stream {
try Task.checkCancellation()
try await self.sendNext(event, id: id)
}
stream = try await onSubscribe(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
subscriptionTasks.removeValue(forKey: id)
throw error
return
}
for try await event in stream {
try await self.sendNext(event, id: id)
}
executionTasks.removeValue(forKey: id)
} else {
let result: GraphQLResult
do {
result = try await onExecute(graphQLRequest, initResult)
} catch {
try await sendError(error, id: id)
return
}
try await self.sendComplete(id: id)
subscriptionTasks.removeValue(forKey: id)
}
} else {
do {
let result = try await onExecute(graphQLRequest, initResult)
try await sendNext(result, id: id)
try await sendComplete(id: id)
} catch {
try await sendError(error, id: id)
}
try await sendComplete(id: id)
}
}

private func onOperationComplete(_ completeRequest: CompleteRequest) async throws {
guard initialized else {
try await error(.notInitialized())
try await messenger.error(.notInitialized())
return
}

let id = completeRequest.id
if let task = subscriptionTasks[id] {
if let task = executionTasks[id] {
task.cancel()
subscriptionTasks.removeValue(forKey: id)
executionTasks.removeValue(forKey: id)
}
try await onOperationComplete(id)
}
Expand DownExpand Up@@ -238,14 +247,4 @@ where
private func sendError(_ error: Error, id: String) async throws {
try await sendError([error], id: id)
}

/// Send an `error` response through the messenger
private func sendError(_ errorMessage: String, id: String) async throws {
try await sendError(GraphQLError(message: errorMessage), id: id)
}

/// Send an error through the messenger and close the connection
private func error(_ error: GraphQLTransportWSError) async throws {
try await messenger.error(error.message, code: error.code.rawValue)
}
}
56 changes: 56 additions & 0 deletions Tests/GraphQLTransportWSTests/GraphQLTransportWSTests.swift
Original file line numberDiff line numberDiff line change
Expand Up@@ -268,6 +268,62 @@ struct GraphqlTransportWSTests {
)
}

/// Tests malformed requests include decoder details in the transport error
@Test func malformedRequestIncludesDecodingDetails() async throws {
let api = TestAPI()
let context = TestContext()
let server = Server<TokenInitPayload, Void, AsyncThrowingStream<GraphQLResult, Error>>(
messenger: serverMessenger,
onInit: { _ in },
onExecute: { graphQLRequest, _ in
try await api.execute(
request: graphQLRequest.query,
context: context
)
},
onSubscribe: { graphQLRequest, _ in
try await api.subscribe(
request: graphQLRequest.query,
context: context
).get()
}
)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"complete"}"#.utf8))
continuation.finish()

try await server.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in serverMessenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Request message doesn't match 'complete' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

/// Tests malformed responses include decoder details in the transport error
@Test func malformedResponseIncludesDecodingDetails() async throws {
let messenger = TestMessenger()
let client = Client<TokenInitPayload>(messenger: messenger)
let (incoming, continuation) = AsyncThrowingStream<Data, any Error>.makeStream()

continuation.yield(Data(#"{"type":"next"}"#.utf8))
continuation.finish()

try await client.listen(to: incoming)

let error = await #expect(throws: TestMessengerError.self) {
for try await _ in messenger.stream {}
}
#expect(error?.code == 4400)
#expect(error?.message.contains("Response message doesn't match 'next' JSON format") == true)
#expect(error?.message.contains("keyNotFound") == true)
#expect(error?.message.contains(#""id""#) == true)
}

enum TestError: Error {
case couldBeAnything
}
Expand Down
Loading