diff --git a/CHANGELOG.md b/CHANGELOG.md index eedbf8ffe..6881b89b3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -88,7 +88,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed - Retry inbound webhook auth from local broker -- Escape glob in @agent-relay/* to satisfy prettier +- Escape glob in @agent-relay/\* to satisfy prettier ## [9.1.0] - 2026-06-24 diff --git a/packages/sdk-swift/Package.resolved b/packages/sdk-swift/Package.resolved new file mode 100644 index 000000000..47cdadbdf --- /dev/null +++ b/packages/sdk-swift/Package.resolved @@ -0,0 +1,14 @@ +{ + "pins" : [ + { + "identity" : "relaycast", + "kind" : "remoteSourceControl", + "location" : "https://github.com/AgentWorkforce/relaycast.git", + "state" : { + "revision" : "21830a01a62defcff261d2de59843b6afff7f534", + "version" : "4.2.0" + } + } + ], + "version" : 2 +} diff --git a/packages/sdk-swift/Package.swift b/packages/sdk-swift/Package.swift index 0bc2547e5..914cce95e 100644 --- a/packages/sdk-swift/Package.swift +++ b/packages/sdk-swift/Package.swift @@ -1,6 +1,16 @@ // swift-tools-version: 5.9 import PackageDescription +// The hosted-participant transport (`AgentRelaySDK`) is a thin facade over the +// relaycast Swift engine SDK (product `Relaycast`). +// +// relaycast's Swift SDK lives in a subdirectory of the relaycast monorepo +// (`packages/sdk-swift`), so it cannot be consumed as a plain git-URL SwiftPM +// dependency on its own (git dependencies require `Package.swift` at the +// repository root). A root-level manifest vending the `Relaycast` library was +// added in the relaycast monorepo (AgentWorkforce/relaycast#208) and published +// as v4.2.0; this package depends on it via that repository's git URL. +// let package = Package( name: "AgentRelaySDK", platforms: [ @@ -13,9 +23,18 @@ let package = Package( .library(name: "AgentRelaySDK", targets: ["AgentRelaySDK"]), .library(name: "AgentRelayBrokerSDK", targets: ["AgentRelayBrokerSDK"]) ], + dependencies: [ + .package( + url: "https://github.com/AgentWorkforce/relaycast.git", + from: "4.2.0" + ) + ], targets: [ .target( name: "AgentRelaySDK", + dependencies: [ + .product(name: "Relaycast", package: "relaycast") + ], path: "Sources/AgentRelaySDK" ), .target( diff --git a/packages/sdk-swift/README.md b/packages/sdk-swift/README.md index bb51a1d3c..b1008e487 100644 --- a/packages/sdk-swift/README.md +++ b/packages/sdk-swift/README.md @@ -21,6 +21,30 @@ Add the package in Swift Package Manager: Then depend on either `AgentRelaySDK` or `AgentRelayBrokerSDK`. +### `AgentRelaySDK` wraps the relaycast engine SDK + +`AgentRelaySDK`'s hosted transport is a thin facade over the published relaycast +Swift engine SDK (product `Relaycast`, package `relaycast-swift`, which lives in +the relaycast monorepo under `packages/sdk-swift`). All HTTP and realtime +WebSocket work is delegated to relaycast; `AgentRelaySDK` keeps only the +relay-specific glue (action-dispatch loop, `RelayChannelEvent` shape, and the +`AsyncStream`-based public API). + +relaycast's Swift SDK lives in a subdirectory of the relaycast monorepo +(`packages/sdk-swift`), so it cannot be consumed as a plain git-URL SwiftPM +dependency on its own (git dependencies require `Package.swift` at the +repository root). A root-level manifest that vends the `Relaycast` library is +added to the relaycast monorepo (see +[AgentWorkforce/relaycast#208](https://github.com/AgentWorkforce/relaycast/pull/208)), +and this package depends on it via that repository's git URL: + +```swift +.package(url: "https://github.com/AgentWorkforce/relaycast.git", from: "4.2.0") +``` + +The root manifest landed in relaycast#208 and was published as v4.2.0, so this +package depends on it by version. + ## Quick start ```swift diff --git a/packages/sdk-swift/Sources/AgentRelaySDK/AgentRelayClient.swift b/packages/sdk-swift/Sources/AgentRelaySDK/AgentRelayClient.swift index a6d4a9d3f..b2b6ab0f5 100644 --- a/packages/sdk-swift/Sources/AgentRelaySDK/AgentRelayClient.swift +++ b/packages/sdk-swift/Sources/AgentRelaySDK/AgentRelayClient.swift @@ -1,18 +1,26 @@ import Foundation - +import Relaycast + +/// Hosted-participant client. This is a thin facade over the relaycast Swift +/// engine SDK (`Relaycast`): registration, reconnect, workspace lookup, channel +/// posting, DMs, action handling, and the realtime event stream are all served +/// by `Relaycast.RelayCast` / `Relaycast.AgentClient` / `Relaycast.WsClient`. +/// +/// The public surface (types, method signatures, AsyncStream APIs) is preserved +/// so existing callers keep working; only the transport implementation changed. public final class AgentRelay: @unchecked Sendable { private let core: HostedWorkspaceCore public let workspaceKey: String public let baseURL: URL - public init(workspaceKey: String, baseURL: URL) { + public init(workspaceKey: String, baseURL: URL? = nil) { self.workspaceKey = workspaceKey - self.baseURL = baseURL - let http = HostedHTTP(baseURL: baseURL, apiKey: workspaceKey) - self.core = HostedWorkspaceCore(workspaceKey: workspaceKey, baseURL: baseURL, http: http) + let resolved = Self.resolveBaseURL(from: baseURL) + self.baseURL = resolved + self.core = HostedWorkspaceCore(workspaceKey: workspaceKey, baseURL: resolved) } - public convenience init(apiKey: String, baseURL: URL) { + public convenience init(apiKey: String, baseURL: URL? = nil) { self.init(workspaceKey: apiKey, baseURL: baseURL) } @@ -45,6 +53,13 @@ public final class AgentRelay: @unchecked Sendable { public func workspaceInfo() async throws -> JSONValue { try await core.workspaceInfo() } + + private static func resolveBaseURL(from baseURL: URL?) -> URL { + if let baseURL { + return baseURL + } + return URL(string: "https://cast.agentrelay.com")! + } } /// Compatibility alias for existing Swift consumers. In this module the client @@ -54,143 +69,144 @@ public typealias AgentRelayClient = AgentRelay final class HostedWorkspaceCore: @unchecked Sendable { let workspaceKey: String let baseURL: URL - let http: any HostedHTTPClient - let encoder = JSONEncoder() - - init(workspaceKey: String, baseURL: URL, http: any HostedHTTPClient) { + // `Relaycast.RelayCast(options:)` can throw (e.g. empty apiKey, invalid + // baseURL). The public `AgentRelay` initializers are non-throwing, so we + // capture the construction result eagerly and rethrow a translated + // `RelayError` on first use instead of force-unwrapping (which would crash + // the process on bad configuration). + private let relayResult: Result + + init(workspaceKey: String, baseURL: URL) { self.workspaceKey = workspaceKey self.baseURL = baseURL - self.http = http + // PRESERVE the configured host: pass it explicitly into relaycast. + self.relayResult = Result { + try Relaycast.RelayCast( + options: Relaycast.RelayCastOptions( + apiKey: workspaceKey, + baseURL: baseURL.absoluteString + ) + ) + } + } + + /// Resolve the wrapped relaycast engine, surfacing configuration errors as + /// `RelayError` rather than crashing. + func relayCast() throws -> Relaycast.RelayCast { + do { + return try relayResult.get() + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } func register(name: String, type: RelayAgentType) async throws -> AgentRegistration { - let body = try encode(RegisterAgentRequest(name: name, type: type)) - let response = try await http.post(path: "/v1/agents", body: body) - let registration = try decodeAPIData(response, as: AgentRegistrationResponse.self) - return makeRegistration(registration) + let relay = try relayCast() + do { + let created = try await relay.agents.register( + Relaycast.CreateAgentRequest(name: name, type: type.relaycastType) + ) + return makeRegistration(created) + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } func registerOrRotate(name: String, type: RelayAgentType) async throws -> AgentRegistration { + let relay = try relayCast() do { - return try await register(name: name, type: type) - } catch RelayError.protocolError(code: let code, message: let message, retryable: _) where isNameConflict(code: code, message: message) { - let agent = try await getAgent(name: name) - let token = try await rotateToken(name: agent.name) - return AgentRegistration( - id: agent.id, - name: agent.name, - token: token, - status: agent.status, - createdAt: agent.createdAt ?? agent.lastSeenAt - ) { [baseURL, http] id, agentName, token in - let agentHTTP = HostedHTTP(baseURL: baseURL, apiKey: token) - let transport = RelayEventTransport(baseURL: baseURL, token: token) - let core = HostedParticipantCore( - agentId: id, - agentName: agentName, - token: token, - baseURL: baseURL, - workspaceHTTP: http, - agentHTTP: agentHTTP, - transport: transport - ) - return AgentClient(core: core, id: id, name: agentName, token: token) - } + let created = try await relay.registerOrRotate( + Relaycast.RegisterAgentRequest(name: name, type: type.relaycastType) + ) + return makeRegistration(created) + } catch let error as Relaycast.RelayError { + throw RelayError(error) } } func reconnect(apiToken: String) async throws -> AgentClient { - let agentHTTP = HostedHTTP(baseURL: baseURL, apiKey: apiToken) - let data = try await agentHTTP.get(path: "/v1/agent", query: nil) - let agent = try decodeAPIData(data, as: RelayAgent.self) - return agentClient(id: agent.id, name: agent.name, token: apiToken) + let relay = try relayCast() + do { + let engine = try await relay.reconnect(Relaycast.AgentReconnectOptions(apiToken: apiToken)) + let me = try await engine.me() + let core = HostedParticipantCore(engineSource: .ready(engine: engine, relay: relay), agentId: me.id, agentName: me.name, token: apiToken, baseURL: baseURL) + return AgentClient(core: core, id: me.id, name: me.name, token: apiToken) + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } func agentClient(id: String, name: String, token: String) -> AgentClient { - let agentHTTP = HostedHTTP(baseURL: baseURL, apiKey: token) - let transport = RelayEventTransport(baseURL: baseURL, token: token) - let core = HostedParticipantCore( - agentId: id, - agentName: name, - token: token, - baseURL: baseURL, - workspaceHTTP: http, - agentHTTP: agentHTTP, - transport: transport - ) + let core = makeParticipantCore(id: id, name: name, token: token) return AgentClient(core: core, id: id, name: name, token: token) } func workspaceInfo() async throws -> JSONValue { - let data = try await http.get(path: "/v1/workspace", query: nil) - return try decodeAPIData(data, as: JSONValue.self) + let relay = try relayCast() + do { + let workspace = try await relay.workspace.info() + return Self.workspaceJSON(workspace) + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } - private func getAgent(name: String) async throws -> RelayAgent { - let data = try await http.get(path: "/v1/agents/\(Self.escapePath(name))", query: nil) - return try decodeAPIData(data, as: RelayAgent.self) + func makeParticipantCore(id: String, name: String, token: String) -> HostedParticipantCore { + Self.makeParticipantCore(relayResult: relayResult, baseURL: baseURL, id: id, name: name, token: token) } - private func rotateToken(name: String) async throws -> String { - let data = try await http.post(path: "/v1/agents/\(Self.escapePath(name))/rotate-token", body: encodeEmptyObject()) - return try decodeAPIData(data, as: RotateTokenResponse.self).token + static func makeParticipantCore(relayResult: Result, baseURL: URL, id: String, name: String, token: String) -> HostedParticipantCore { + // Defer the per-agent engine build (`relay.asAgent`) into the actor so + // that configuration/handshake errors propagate through the first async + // call as `RelayError` instead of force-unwrapping and crashing. + return HostedParticipantCore(engineSource: .deferred(relayResult: relayResult, token: token), agentId: id, agentName: name, token: token, baseURL: baseURL) } - private func makeRegistration(_ response: AgentRegistrationResponse) -> AgentRegistration { - AgentRegistration( + func makeRegistration(_ response: Relaycast.CreateAgentResponse) -> AgentRegistration { + // Capture the transport state (`relayResult`, `baseURL`) strongly so a + // persisted `AgentRegistration` stays usable even if the owning + // `AgentRelay`/`HostedWorkspaceCore` is released before `asClient()` is + // called. `Relaycast.RelayCast` is the only shared, reusable state; the + // per-agent engine is created lazily from it in the closure. + let relayResult = self.relayResult + let baseURL = self.baseURL + return AgentRegistration( id: response.id, name: response.name, token: response.token, - status: response.status, + status: RelayAgentStatus(response.status), createdAt: response.createdAt - ) { [baseURL, http] id, agentName, token in - let agentHTTP = HostedHTTP(baseURL: baseURL, apiKey: token) - let transport = RelayEventTransport(baseURL: baseURL, token: token) - let core = HostedParticipantCore( - agentId: id, - agentName: agentName, - token: token, + ) { id, agentName, token in + let core = HostedWorkspaceCore.makeParticipantCore( + relayResult: relayResult, baseURL: baseURL, - workspaceHTTP: http, - agentHTTP: agentHTTP, - transport: transport + id: id, + name: agentName, + token: token ) return AgentClient(core: core, id: id, name: agentName, token: token) } } - private func encode(_ value: T) throws -> Data { - do { - return try encoder.encode(value) - } catch { - throw RelayError.encodingFailed(String(describing: error)) + private static func workspaceJSON(_ workspace: Relaycast.Workspace) -> JSONValue { + var object: [String: JSONValue] = [ + "id": .string(workspace.id), + "name": .string(workspace.name), + "created_at": .string(workspace.createdAt) + ] + if let systemPrompt = workspace.systemPrompt { + object["system_prompt"] = .string(systemPrompt) } - } - - private func encodeEmptyObject() -> Data { - Data("{}".utf8) - } - - private static func escapePath(_ value: String) -> String { - var allowed = CharacterSet.urlPathAllowed - allowed.remove(charactersIn: "/") - return value.addingPercentEncoding(withAllowedCharacters: allowed) ?? value - } - - private func isNameConflict(code: String, message: String) -> Bool { - let normalizedCode = code.trimmingCharacters(in: .whitespacesAndNewlines).lowercased() - if ["agent_already_exists", "name_conflict", "name_taken", "agent_exists", "conflict", "duplicate", "http_409"].contains(normalizedCode) { - return true + if let plan = workspace.plan { + object["plan"] = .string(plan) } - return message.lowercased().contains("already exists") + if let metadata = workspace.metadata { + object["metadata"] = .object(metadata.mapValues { JSONValue($0) }) + } + return .object(object) } } -private struct RegisterAgentRequest: Encodable { - let name: String - let type: RelayAgentType -} - public final class AgentClient: @unchecked Sendable { private let core: HostedParticipantCore public let id: String @@ -294,65 +310,121 @@ private struct RegisteredAction: Sendable { let handler: RelayActionHandler } +/// Higher-level glue kept ON TOP of the relaycast engine SDK: the +/// action-dispatch loop, channel-event normalization into `RelayChannelEvent`, +/// AsyncStream fan-out, and subscription bookkeeping. The realtime socket and +/// HTTP calls are delegated to the wrapped `Relaycast.AgentClient` / +/// `Relaycast.RelayCast`. actor HostedParticipantCore { + /// How the per-agent engine is obtained. `.ready` is used by `reconnect`, + /// which already has a live engine; `.deferred` builds the engine lazily via + /// `relay.asAgent(token)` on first connect so that `asAgent`/configuration + /// errors propagate as `RelayError` instead of crashing at construction. + enum EngineSource { + case ready(engine: Relaycast.AgentClient, relay: Relaycast.RelayCast) + case deferred(relayResult: Result, token: String) + } + let agentId: String let agentName: String let token: String let baseURL: URL - let workspaceHTTP: any HostedHTTPClient - let agentHTTP: any HostedHTTPClient - let transport: any HostedEventTransportClient - let encoder = JSONEncoder() - let decoder = JSONDecoder() - - private var routerTask: Task? + private let engineSource: EngineSource + private var resolvedEngine: Relaycast.AgentClient? + private var resolvedRelay: Relaycast.RelayCast? + private let encoder = JSONEncoder() + private let decoder = JSONDecoder() + + private var connected = false + private var listenersInstalled = false private var subscribedChannels: Set = [] private var channelContinuations: [String: [AsyncStream.Continuation]] = [:] private var inboundMessageContinuations: [AsyncStream.Continuation] = [] private var eventContinuations: [AsyncStream.Continuation] = [] private var connectionStateContinuations: [AsyncStream.Continuation] = [] private var actionHandlers: [String: RegisteredAction] = [:] - - init( - agentId: String, - agentName: String, - token: String, - baseURL: URL, - workspaceHTTP: any HostedHTTPClient, - agentHTTP: any HostedHTTPClient, - transport: any HostedEventTransportClient - ) { + private var unsubscribeHandlers: [() -> Void] = [] + + // Serialized inbound-event pipeline. The engine delivers events in order + // from a single receive loop; we yield them (synchronously, FIFO) into this + // stream and drain them through `routeEvent` on one consumer task so that + // ordering is preserved end-to-end. (Spawning an unstructured `Task` per + // event would let the scheduler reorder closely-spaced events.) + private var eventBuffer: AsyncStream.Continuation? + private var eventPump: Task? + + init(engineSource: EngineSource, agentId: String, agentName: String, token: String, baseURL: URL) { + self.engineSource = engineSource self.agentId = agentId self.agentName = agentName self.token = token self.baseURL = baseURL - self.workspaceHTTP = workspaceHTTP - self.agentHTTP = agentHTTP - self.transport = transport } - func ensureConnected() async throws { - if routerTask == nil || routerTask?.isCancelled == true { - routerTask = Task { [weak self] in await self?.routeFrames() } + /// Resolve (and cache) the relaycast `RelayCast` engine wrapper. Surfaces + /// configuration errors as `RelayError`. + private func relayCast() throws -> Relaycast.RelayCast { + if let resolvedRelay { return resolvedRelay } + switch engineSource { + case .ready(_, let relay): + resolvedRelay = relay + return relay + case .deferred(let relayResult, _): + do { + let relay = try relayResult.get() + resolvedRelay = relay + return relay + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } - await transport.setOnConnect { [weak self] in - await self?.transportDidReconnect() + } + + /// Resolve (and cache) the per-agent engine, building it lazily for the + /// `.deferred` source. Surfaces `asAgent`/configuration errors as `RelayError`. + private func engine() throws -> Relaycast.AgentClient { + if let resolvedEngine { return resolvedEngine } + switch engineSource { + case .ready(let engine, _): + resolvedEngine = engine + return engine + case .deferred(_, let token): + let relay = try relayCast() + do { + let engine = try relay.asAgent(token) + resolvedEngine = engine + return engine + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } - try await transport.connect() - notifyConnectionState(.connected) - try await syncSubscriptions() } - func transportDidReconnect() async { + func ensureConnected() async throws { + let engine = try engine() + installListenersIfNeeded(engine: engine) + if !connected { + engine.connect() + connected = true + } notifyConnectionState(.connected) - try? await syncSubscriptions() + syncSubscriptions(engine: engine) } func disconnect() async { - routerTask?.cancel() - routerTask = nil - await transport.disconnect() - _ = try? await agentHTTP.post(path: "/v1/agents/disconnect", body: Data("{}".utf8)) + // Only tear down an engine that was actually built/connected; building + // one here just to disconnect it would be pointless (and could throw). + if let resolvedEngine { + await resolvedEngine.disconnect() + } + connected = false + for unsubscribe in unsubscribeHandlers { unsubscribe() } + unsubscribeHandlers.removeAll() + listenersInstalled = false + eventBuffer?.finish() + eventBuffer = nil + eventPump?.cancel() + eventPump = nil notifyConnectionState(.disconnected) for continuations in channelContinuations.values { for continuation in continuations { continuation.finish() } @@ -388,14 +460,21 @@ actor HostedParticipantCore { } func post(channel: String, text: String) async throws { - let path = "/v1/channels/\(Self.escapePath(Self.normalizeChannel(channel)))/messages" - let body = try encode(SendChannelMessageRequest(text: text, mode: "wait")) - _ = try await agentHTTP.post(path: path, body: body) + let engine = try engine() + do { + _ = try await engine.send(Self.normalizeChannel(channel), text: text, options: Relaycast.SendMessageOptions(mode: .wait)) + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } func dm(to target: String, text: String) async throws { - let body = try encode(SendDirectMessageRequest(to: Self.stripSigil(target), text: text, mode: "wait")) - _ = try await agentHTTP.post(path: "/v1/dm", body: body) + let engine = try engine() + do { + _ = try await engine.dm(Self.stripSigil(target), text: text, options: Relaycast.DMOptions(mode: .wait)) + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } func registerAction( @@ -408,7 +487,7 @@ actor HostedParticipantCore { guard !actionName.isEmpty else { throw RelayError.protocolError(code: "invalid_action_name", message: "Action name cannot be empty", retryable: false) } - let inputSchema = try decodeJSONValue(inputSchemaJSON) + let inputSchema = try decodeRelaycastObject(inputSchemaJSON) let registrationId = UUID().uuidString actionHandlers[actionName] = RegisteredAction(id: registrationId, handler: handler) @@ -437,37 +516,104 @@ actor HostedParticipantCore { } } - private func registerActionDescriptor(name: String, description: String, inputSchema: JSONValue) async throws { - let request = RegisterActionDescriptorRequest( - name: name, - description: description, - handlerAgent: agentName, - inputSchema: inputSchema - ) - let body = try encode(request) - _ = try await workspaceHTTP.post(path: "/v1/actions", body: body) + private func registerActionDescriptor(name: String, description: String, inputSchema: [String: Relaycast.JSONValue]) async throws { + let relay = try relayCast() + do { + _ = try await relay.actions.register( + Relaycast.RegisterActionRequest( + name: name, + description: description, + handlerAgent: agentName, + inputSchema: inputSchema + ) + ) + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } private func unregisterActionDescriptor(name: String) async throws { - _ = try await workspaceHTTP.delete(path: "/v1/actions/\(Self.escapePath(name))") + let relay = try relayCast() + do { + try await relay.actions.delete(name) + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } - private func syncSubscriptions() async throws { + private func syncSubscriptions(engine: Relaycast.AgentClient) { guard !subscribedChannels.isEmpty else { return } - let body = try encode(SocketSubscribeMessage(channels: Array(subscribedChannels).sorted())) - try await transport.send(body) + engine.subscribe(Array(subscribedChannels).sorted()) } - private func routeFrames() async { - for await data in transport.inbound { - guard let event = try? decoder.decode(RelayEvent.self, from: data) else { - continue + // MARK: - Realtime listeners + + private func installListenersIfNeeded(engine: Relaycast.AgentClient) { + guard !listenersInstalled else { return } + listenersInstalled = true + + // Build the serialized event pipeline before wiring engine callbacks so + // the first delivered event already has a place to queue. + var continuation: AsyncStream.Continuation! + let stream = AsyncStream { continuation = $0 } + eventBuffer = continuation + eventPump = Task { [weak self] in + for await event in stream { + await self?.routeEvent(event) } - routeEvent(event) } + + // `engine.on.*` fire in order from the engine's single receive loop; + // yielding into `buffer` (synchronously) preserves that order, and the + // single `eventPump` consumer drains them sequentially. + let buffer = continuation! + let ingest: @Sendable (Relaycast.WsEvent) -> Void = { event in + buffer.yield(RelayEvent(event)) + } + unsubscribeHandlers.append(engine.on.messageCreated(ingest)) + unsubscribeHandlers.append(engine.on.threadReply(ingest)) + unsubscribeHandlers.append(engine.on.dmReceived(ingest)) + unsubscribeHandlers.append(engine.on.groupDMReceived(ingest)) + unsubscribeHandlers.append(engine.on.actionInvoked(ingest)) + + unsubscribeHandlers.append(engine.on.connected { [weak self] in + guard let self else { return } + Task { await self.transportDidConnect() } + }) + unsubscribeHandlers.append(engine.on.disconnected { [weak self] in + guard let self else { return } + Task { await self.handleEngineDisconnect() } + }) + unsubscribeHandlers.append(engine.on.reconnecting { [weak self] attempt in + guard let self else { return } + Task { await self.notifyConnectionStateAsync(.reconnecting(attempt: attempt)) } + }) + } + + private func transportDidConnect() async { + notifyConnectionState(.connected) + // The engine is necessarily resolved here: this fires from the engine's + // own `connected` callback, which is only installed after `engine()` ran. + if let resolvedEngine { + syncSubscriptions(engine: resolvedEngine) + } + } + + private func notifyConnectionStateAsync(_ state: ConnectionStateChange) async { + notifyConnectionState(state) + } + + /// Handle an engine-initiated disconnect. Reset `connected` so that a later + /// `ensureConnected()` will actually re-issue `engine.connect()`, matching + /// the manual `disconnect()` path (which also clears the flag). Without this + /// the flag stays `true` after an engine drop and reconnection is skipped. + private func handleEngineDisconnect() async { + connected = false notifyConnectionState(.disconnected) } + // MARK: - Event routing (glue kept on top of the engine) + private func routeEvent(_ event: RelayEvent) { for continuation in eventContinuations { continuation.yield(event) @@ -535,7 +681,7 @@ actor HostedParticipantCore { ) async { do { let invocation = try await loadInvocation(actionName: actionName, invocationId: invocationId) - let input = invocation.input ?? .object([:]) + let input = invocation.input ?? [:] let inputString = actionInputString(input) let output = await registration.handler(inputString) try await completeInvocation(actionName: actionName, invocationId: invocationId, output: parseHandlerOutput(output)) @@ -544,63 +690,79 @@ actor HostedParticipantCore { } } - private func loadInvocation(actionName: String, invocationId: String) async throws -> RelayActionInvocation { - let path = "/v1/actions/\(Self.escapePath(actionName))/invocations/\(Self.escapePath(invocationId))" - let data = try await agentHTTP.get(path: path, query: nil) - return try decodeAPIData(data, as: RelayActionInvocation.self) + private func loadInvocation(actionName: String, invocationId: String) async throws -> Relaycast.ActionInvocation { + let engine = try engine() + do { + return try await engine.actions.getInvocation(name: actionName, invocationID: invocationId) + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } - private func completeInvocation(actionName: String, invocationId: String, output: JSONValue) async throws { - try await completeInvocation( - actionName: actionName, - invocationId: invocationId, - body: CompleteInvocationRequest(output: Self.outputRecord(output), error: nil) - ) + private func completeInvocation(actionName: String, invocationId: String, output: [String: Relaycast.JSONValue]) async throws { + let engine = try engine() + do { + _ = try await engine.actions.completeInvocation( + name: actionName, + invocationID: invocationId, + data: Relaycast.CompleteInvocationRequest(output: output) + ) + } catch let error as Relaycast.RelayError { + throw RelayError(error) + } } private func completeInvocation(actionName: String, invocationId: String, error: String) async throws { - try await completeInvocation(actionName: actionName, invocationId: invocationId, body: CompleteInvocationRequest(output: nil, error: error)) + let engine = try engine() + do { + _ = try await engine.actions.completeInvocation( + name: actionName, + invocationID: invocationId, + data: Relaycast.CompleteInvocationRequest(error: error) + ) + } catch let relayError as Relaycast.RelayError { + throw RelayError(relayError) + } } - private func completeInvocation(actionName: String, invocationId: String, body value: CompleteInvocationRequest) async throws { - let path = "/v1/actions/\(Self.escapePath(actionName))/invocations/\(Self.escapePath(invocationId))/complete" - _ = try await agentHTTP.post(path: path, body: try encode(value)) - } + // MARK: - JSON helpers - private func decodeJSONValue(_ json: String) throws -> JSONValue { + private func decodeRelaycastObject(_ json: String) throws -> [String: Relaycast.JSONValue] { guard let data = json.data(using: .utf8) else { throw RelayError.encodingFailed("Input schema is not valid UTF-8") } do { - return try decoder.decode(JSONValue.self, from: data) + let value = try decoder.decode(Relaycast.JSONValue.self, from: data) + guard case .object(let object) = value else { + throw RelayError.decodingFailed("Input schema must be a JSON object") + } + return object + } catch let error as RelayError { + throw error } catch { throw RelayError.decodingFailed("Invalid inputSchemaJSON: \(error)") } } - private func actionInputString(_ input: JSONValue) -> String { - guard let data = try? encoder.encode(input), let string = String(data: data, encoding: .utf8) else { + private func actionInputString(_ input: [String: Relaycast.JSONValue]) -> String { + guard let data = try? encoder.encode(Relaycast.JSONValue.object(input)), + let string = String(data: data, encoding: .utf8) else { return "{}" } return string } - private func parseHandlerOutput(_ output: String) -> JSONValue { + private func parseHandlerOutput(_ output: String) -> [String: Relaycast.JSONValue] { let trimmed = output.trimmingCharacters(in: .whitespacesAndNewlines) if !trimmed.isEmpty, let data = trimmed.data(using: .utf8), - let value = try? decoder.decode(JSONValue.self, from: data) { - return value - } - return .string(output) - } - - private func encode(_ value: T) throws -> Data { - do { - return try encoder.encode(value) - } catch { - throw RelayError.encodingFailed(String(describing: error)) + let value = try? decoder.decode(Relaycast.JSONValue.self, from: data) { + if case .object(let object) = value { + return object + } + return ["value": value] } + return ["value": .string(output)] } private func notifyConnectionState(_ state: ConnectionStateChange) { @@ -620,12 +782,6 @@ actor HostedParticipantCore { stripSigil(value).trimmingCharacters(in: .whitespacesAndNewlines) } - private static func escapePath(_ value: String) -> String { - var allowed = CharacterSet.urlPathAllowed - allowed.remove(charactersIn: "/") - return value.addingPercentEncoding(withAllowedCharacters: allowed) ?? value - } - private static func date(from timestamp: String?) -> Date { guard let timestamp, let date = ISO8601DateFormatter().date(from: timestamp) @@ -652,64 +808,4 @@ actor HostedParticipantCore { } return error.localizedDescription } - - private static func outputRecord(_ value: JSONValue) -> JSONValue { - if case .object = value { - return value - } - return .object(["value": value]) - } -} - -private struct SendChannelMessageRequest: Encodable { - let text: String - let mode: String -} - -private struct SendDirectMessageRequest: Encodable { - let to: String - let text: String - let mode: String -} - -private struct SocketSubscribeMessage: Encodable { - let type = "subscribe" - let channels: [String] -} - -private struct RegisterActionDescriptorRequest: Encodable { - let name: String - let description: String - let handlerAgent: String - let inputSchema: JSONValue - - enum CodingKeys: String, CodingKey { - case name, description - case handlerAgent = "handler_agent" - case inputSchema = "input_schema" - } -} - -private struct CompleteInvocationRequest: Encodable { - let output: JSONValue? - let error: String? - - func encode(to encoder: Encoder) throws { - var container = encoder.container(keyedBy: CodingKeys.self) - switch (output, error) { - case (.some(let output), .none): - try container.encode(output, forKey: .output) - case (.none, .some(let error)): - try container.encode(error, forKey: .error) - default: - throw EncodingError.invalidValue( - self, - EncodingError.Context(codingPath: encoder.codingPath, debugDescription: "completion requires output or error") - ) - } - } - - enum CodingKeys: String, CodingKey { - case output, error - } } diff --git a/packages/sdk-swift/Sources/AgentRelaySDK/RelayHTTP.swift b/packages/sdk-swift/Sources/AgentRelaySDK/RelayHTTP.swift deleted file mode 100644 index 28cf4a39c..000000000 --- a/packages/sdk-swift/Sources/AgentRelaySDK/RelayHTTP.swift +++ /dev/null @@ -1,142 +0,0 @@ -import Foundation - -#if canImport(FoundationNetworking) -import FoundationNetworking -#endif - -protocol HostedHTTPClient: Sendable { - func get(path: String, query: [String: String]?) async throws -> Data - func post(path: String, body: Data?) async throws -> Data - func delete(path: String) async throws -> Data -} - -actor HostedHTTP: HostedHTTPClient { - private let baseURL: URL - private let apiKey: String - private let session: URLSession - - init(baseURL: URL, apiKey: String, session: URLSession = .shared) { - self.baseURL = baseURL - self.apiKey = apiKey - self.session = session - } - - func get(path: String, query: [String: String]? = nil) async throws -> Data { - try await request(method: "GET", path: path, query: query, body: nil) - } - - func post(path: String, body: Data?) async throws -> Data { - try await request(method: "POST", path: path, query: nil, body: body) - } - - func delete(path: String) async throws -> Data { - try await request(method: "DELETE", path: path, query: nil, body: nil) - } - - private func request(method: String, path: String, query: [String: String]?, body: Data?) async throws -> Data { - guard let url = Self.resolveAPIURL(baseURL: baseURL, path: path, query: query) else { - throw RelayError.invalidBaseURL("Could not resolve URL for path \(path)") - } - var request = URLRequest(url: url) - request.httpMethod = method - request.setValue("Bearer \(apiKey)", forHTTPHeaderField: "Authorization") - request.setValue("agent-relay-swift", forHTTPHeaderField: "X-Relaycast-Origin-Client") - request.setValue("swift-sdk-split", forHTTPHeaderField: "X-Relaycast-Origin-Version") - if let body { - request.httpBody = body - request.setValue("application/json", forHTTPHeaderField: "Content-Type") - } - - let (data, response): (Data, URLResponse) - do { - (data, response) = try await session.data(for: request) - } catch { - throw RelayError.connectionFailed(String(describing: error)) - } - - guard let http = response as? HTTPURLResponse else { - throw RelayError.connectionFailed("Non-HTTP response from Relaycast") - } - - if !(200..<300).contains(http.statusCode) { - let (code, message) = decodeErrorBody(data, fallbackStatus: http.statusCode) - throw RelayError.protocolError( - code: code, - message: message, - retryable: http.statusCode == 429 || http.statusCode >= 500 - ) - } - - return data - } - - static func resolveAPIURL(baseURL: URL, path: String, query: [String: String]? = nil) -> URL? { - guard var components = URLComponents(url: baseURL, resolvingAgainstBaseURL: false) else { - return nil - } - if components.scheme == "ws" { components.scheme = "http" } - if components.scheme == "wss" { components.scheme = "https" } - - var basePath = components.path - while basePath.hasSuffix("/") { basePath = String(basePath.dropLast()) } - if basePath.hasSuffix("/v1/ws") { - basePath = String(basePath.dropLast("/v1/ws".count)) - } - - let normalizedPath = path.hasPrefix("/") ? path : "/" + path - components.path = basePath + normalizedPath - components.queryItems = query?.map { URLQueryItem(name: $0.key, value: $0.value) } - components.fragment = nil - return components.url - } - - private func decodeErrorBody(_ data: Data, fallbackStatus: Int) -> (code: String, message: String) { - if let envelope = try? JSONDecoder().decode(APIEnvelope.self, from: data), - let error = envelope.error { - return (error.code, error.message) - } - if let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any] { - let error = json["error"] as? [String: Any] - let code = (error?["code"] as? String) ?? (json["code"] as? String) ?? "http_\(fallbackStatus)" - let message = - (error?["message"] as? String) - ?? (json["message"] as? String) - ?? "HTTP \(fallbackStatus)" - return (code, message) - } - return ("http_\(fallbackStatus)", "HTTP \(fallbackStatus)") - } -} - -struct APIEnvelope: Decodable { - let ok: Bool - let data: T? - let error: APIErrorPayload? -} - -struct APIErrorPayload: Decodable { - let code: String - let message: String -} - -struct EmptyAPIData: Decodable {} - -func decodeAPIData(_ data: Data, as type: T.Type = T.self) throws -> T { - do { - let envelope = try JSONDecoder().decode(APIEnvelope.self, from: data) - if envelope.ok, let value = envelope.data { - return value - } - if envelope.ok, T.self == EmptyAPIData.self { - return EmptyAPIData() as! T - } - if let error = envelope.error { - throw RelayError.protocolError(code: error.code, message: error.message, retryable: false) - } - throw RelayError.decodingFailed("Relay API response did not include data") - } catch let error as RelayError { - throw error - } catch { - throw RelayError.decodingFailed(String(describing: error)) - } -} diff --git a/packages/sdk-swift/Sources/AgentRelaySDK/RelayTransport.swift b/packages/sdk-swift/Sources/AgentRelaySDK/RelayTransport.swift deleted file mode 100644 index acb415477..000000000 --- a/packages/sdk-swift/Sources/AgentRelaySDK/RelayTransport.swift +++ /dev/null @@ -1,214 +0,0 @@ -import Foundation - -#if canImport(FoundationNetworking) -import FoundationNetworking -#endif - -protocol HostedEventTransportClient: Sendable { - var inbound: AsyncStream { get } - - func setOnConnect(_ handler: @escaping @Sendable () async -> Void) async - func connect() async throws - func disconnect() async - func send(_ message: Data) async throws -} - -public actor RelayEventTransport: HostedEventTransportClient { - public enum ConnectionState: Sendable { - case disconnected - case connecting - case connected - case reconnecting - } - - public enum TransportError: Error, Sendable { - case notConnected - case sendFailed(String) - case connectionFailed(String) - } - - public nonisolated let inbound: AsyncStream - - private let baseURL: URL - private let token: String - private let session: URLSession - private var webSocketTask: URLSessionWebSocketTask? - private var receiveTask: Task? - private var pingTask: Task? - private var reconnectTask: Task? - private var inboundContinuation: AsyncStream.Continuation? - private var state: ConnectionState = .disconnected - private var manuallyDisconnected = false - private var reconnectAttempt = 0 - private var onConnect: (@Sendable () async -> Void)? - - public init(baseURL: URL, token: String, session: URLSession = .shared) { - self.baseURL = baseURL - self.token = token - self.session = session - var continuationRef: AsyncStream.Continuation? - self.inbound = AsyncStream { continuation in - continuationRef = continuation - } - self.inboundContinuation = continuationRef - } - - public func setOnConnect(_ handler: @escaping @Sendable () async -> Void) async { - self.onConnect = handler - } - - public func connect() async throws { - switch state { - case .connected, .connecting: - return - case .disconnected, .reconnecting: - break - } - - manuallyDisconnected = false - state = reconnectAttempt == 0 ? .connecting : .reconnecting - let isReconnect = reconnectAttempt > 0 - let request = URLRequest(url: Self.resolveWebSocketURL(baseURL: baseURL, token: token) ?? baseURL) - let task = session.webSocketTask(with: request) - webSocketTask = task - task.resume() - state = .connected - reconnectAttempt = 0 - startReceiveLoop() - startPingLoop() - if isReconnect, let onConnect { - await onConnect() - } - } - - public func disconnect() async { - manuallyDisconnected = true - receiveTask?.cancel() - pingTask?.cancel() - reconnectTask?.cancel() - receiveTask = nil - pingTask = nil - reconnectTask = nil - webSocketTask?.cancel(with: .goingAway, reason: nil) - webSocketTask = nil - state = .disconnected - } - - public func send(_ message: Data) async throws { - guard let task = webSocketTask, state == .connected else { - throw TransportError.notConnected - } - do { - if let string = String(data: message, encoding: .utf8) { - try await task.send(.string(string)) - } else { - try await task.send(.data(message)) - } - } catch { - throw TransportError.sendFailed(String(describing: error)) - } - } - - static func resolveWebSocketURL(baseURL: URL, token: String) -> URL? { - var components = URLComponents(url: baseURL, resolvingAgainstBaseURL: false) - if components?.scheme == "http" { components?.scheme = "ws" } - if components?.scheme == "https" { components?.scheme = "wss" } - - var path = components?.path ?? "" - while path.hasSuffix("/") { path = String(path.dropLast()) } - if path.hasSuffix("/v1/ws") { - path = String(path.dropLast("/v1/ws".count)) - } - components?.path = path + "/v1/ws" - - var queryItems = components?.queryItems?.filter { $0.name != "token" } ?? [] - queryItems.append(URLQueryItem(name: "token", value: token)) - queryItems.append(URLQueryItem(name: "origin_client", value: "agent-relay-swift")) - queryItems.append(URLQueryItem(name: "origin_version", value: "swift-sdk-split")) - components?.queryItems = queryItems - - return components?.url - } - - private func startReceiveLoop() { - receiveTask?.cancel() - receiveTask = Task { [weak self] in - guard let self else { return } - while !Task.isCancelled { - do { - guard let task = await self.webSocketTask else { return } - let message = try await task.receive() - await self.handle(message) - } catch { - await self.handleDisconnect(error: error) - return - } - } - } - } - - private func startPingLoop() { - pingTask?.cancel() - pingTask = Task { [weak self] in - guard let self else { return } - while !Task.isCancelled { - try? await Task.sleep(for: .seconds(30)) - if Task.isCancelled { return } - try? await self.send(Self.encodeSocketMessage(["type": "ping"])) - } - } - } - - private func handle(_ message: URLSessionWebSocketTask.Message) { - switch message { - case .data(let data): - inboundContinuation?.yield(data) - case .string(let string): - if let data = string.data(using: .utf8) { - inboundContinuation?.yield(data) - } - @unknown default: - break - } - } - - private func handleDisconnect(error: Error) async { - receiveTask?.cancel() - pingTask?.cancel() - webSocketTask?.cancel(with: .goingAway, reason: nil) - webSocketTask = nil - - guard !manuallyDisconnected else { - state = .disconnected - return - } - - state = .reconnecting - let delay = reconnectDelay(for: reconnectAttempt) - reconnectAttempt += 1 - reconnectTask?.cancel() - reconnectTask = Task { [weak self] in - guard let self else { return } - try? await Task.sleep(for: .milliseconds(delay)) - do { - try await self.connect() - } catch { - await self.handleDisconnect(error: error) - } - } - } - - private func reconnectDelay(for attempt: Int) -> Int { - switch attempt { - case 0: return 1_000 - case 1: return 2_000 - case 2: return 4_000 - case 3: return 8_000 - default: return 30_000 - } - } - - static func encodeSocketMessage(_ value: [String: Any]) throws -> Data { - try JSONSerialization.data(withJSONObject: value) - } -} diff --git a/packages/sdk-swift/Sources/AgentRelaySDK/RelayTypes.swift b/packages/sdk-swift/Sources/AgentRelaySDK/RelayTypes.swift index 4e6a764d8..ff84ecf7d 100644 --- a/packages/sdk-swift/Sources/AgentRelaySDK/RelayTypes.swift +++ b/packages/sdk-swift/Sources/AgentRelaySDK/RelayTypes.swift @@ -125,34 +125,6 @@ public struct RelayAgent: Decodable, Sendable { } } -struct AgentRegistrationResponse: Decodable, Sendable { - let id: String - let name: String - let token: String - let status: RelayAgentStatus - let createdAt: String? - - enum CodingKeys: String, CodingKey { - case id, name, token, status - case createdAt = "created_at" - case createdAtCamel = "createdAt" - } - - init(from decoder: Decoder) throws { - let container = try decoder.container(keyedBy: CodingKeys.self) - id = try container.decode(String.self, forKey: .id) - name = try container.decode(String.self, forKey: .name) - token = try container.decode(String.self, forKey: .token) - status = (try? container.decode(RelayAgentStatus.self, forKey: .status)) ?? .unknown - createdAt = try container.decodeIfPresent(String.self, forKey: .createdAt) - ?? container.decodeIfPresent(String.self, forKey: .createdAtCamel) - } -} - -struct RotateTokenResponse: Decodable, Sendable { - let token: String -} - public struct RelayChannelEvent: Sendable { public let from: String public let body: String @@ -270,6 +242,30 @@ public struct RelayEvent: Decodable, Sendable { public let status: String? public let rawJSON: JSONValue? + init( + type: String, + id: String? = nil, + channel: String? = nil, + message: RelayMessage? = nil, + invocationId: String? = nil, + actionName: String? = nil, + callerName: String? = nil, + agentName: String? = nil, + status: String? = nil, + rawJSON: JSONValue? = nil + ) { + self.type = type + self.id = id + self.channel = channel + self.message = message + self.invocationId = invocationId + self.actionName = actionName + self.callerName = callerName + self.agentName = agentName + self.status = status + self.rawJSON = rawJSON + } + enum CodingKeys: String, CodingKey { case type, id, channel, message, status, payload case invocationId = "invocation_id" @@ -332,36 +328,6 @@ private struct RelaycastEventAgent: Decodable { let name: String? } -struct RelayActionInvocation: Decodable, Sendable { - let invocationId: String - let actionName: String - let callerName: String? - let input: JSONValue? - let status: String - - enum CodingKeys: String, CodingKey { - case invocationId = "invocation_id" - case invocationIdCamel = "invocationId" - case actionName = "action_name" - case actionNameCamel = "actionName" - case callerName = "caller_name" - case callerNameCamel = "callerName" - case input, status - } - - init(from decoder: Decoder) throws { - let container = try decoder.container(keyedBy: CodingKeys.self) - invocationId = try container.decodeIfPresent(String.self, forKey: .invocationId) - ?? container.decode(String.self, forKey: .invocationIdCamel) - actionName = try container.decodeIfPresent(String.self, forKey: .actionName) - ?? container.decode(String.self, forKey: .actionNameCamel) - callerName = try container.decodeIfPresent(String.self, forKey: .callerName) - ?? container.decodeIfPresent(String.self, forKey: .callerNameCamel) - input = try container.decodeIfPresent(JSONValue.self, forKey: .input) - status = (try? container.decode(String.self, forKey: .status)) ?? "invoked" - } -} - public actor ActionHandle { public nonisolated let name: String diff --git a/packages/sdk-swift/Sources/AgentRelaySDK/RelaycastBridge.swift b/packages/sdk-swift/Sources/AgentRelaySDK/RelaycastBridge.swift new file mode 100644 index 000000000..7d23eb159 --- /dev/null +++ b/packages/sdk-swift/Sources/AgentRelaySDK/RelaycastBridge.swift @@ -0,0 +1,96 @@ +import Foundation +import Relaycast + +// MARK: - Type bridging between AgentRelaySDK's public surface and the +// relaycast engine SDK (`Relaycast`). These conversions let AgentRelaySDK keep +// its existing public types while delegating all transport to relaycast. + +extension RelayAgentType { + var relaycastType: Relaycast.AgentType { + switch self { + case .agent: return .agent + case .human: return .human + case .system: return .system + } + } +} + +extension RelayAgentStatus { + init(_ status: Relaycast.AgentStatus) { + switch status { + case .online: self = .online + case .offline: self = .offline + case .away: self = .away + } + } +} + +extension JSONValue { + /// Convert a relaycast `JSONValue` (which distinguishes int/double) into the + /// AgentRelaySDK `JSONValue` (number-based). + init(_ value: Relaycast.JSONValue) { + switch value { + case .null: + self = .null + case .bool(let bool): + self = .bool(bool) + case .int(let int): + self = .number(Double(int)) + case .double(let double): + self = .number(double) + case .string(let string): + self = .string(string) + case .array(let array): + self = .array(array.map { JSONValue($0) }) + case .object(let object): + self = .object(object.mapValues { JSONValue($0) }) + } + } +} + +extension RelayError { + /// Map a relaycast `RelayError` onto AgentRelaySDK's `RelayError` so callers + /// continue to see the existing error surface. + init(_ error: Relaycast.RelayError) { + switch error { + case .api(let code, let message, let statusCode, let retryable): + let normalizedCode = statusCode == 409 ? "agent_already_exists" : code + self = .protocolError(code: normalizedCode, message: message, retryable: retryable) + case .transport(let message, _, let retryable, _): + self = .protocolError(code: "transport_error", message: message, retryable: retryable) + case .invalidRequest(let message): + self = .encodingFailed(message) + case .invalidResponse(let message, _): + self = .decodingFailed(message) + case .missingData(let message): + self = .decodingFailed(message) + case .notConnected: + self = .notConnected + } + } +} + +extension RelayEvent { + /// Build an AgentRelaySDK `RelayEvent` from a relaycast realtime `WsEvent`. + /// + /// relaycast emits flat events (`type` plus payload fields at the top level), + /// while `RelayEvent` already understands both the flat and nested envelope + /// shapes via its `Decodable` implementation. Re-encode the relaycast event + /// to JSON and decode it through that flexible path so every existing field + /// extraction rule (snake/camel, nested `payload`, `agent.name`, etc.) is + /// reused verbatim. + init(_ event: Relaycast.WsEvent) { + var object: [String: Relaycast.JSONValue] = event.payload + object["type"] = .string(event.type) + + let encoder = JSONEncoder() + if let data = try? encoder.encode(Relaycast.JSONValue.object(object)), + let decoded = try? JSONDecoder().decode(RelayEvent.self, from: data) { + self = decoded + return + } + + // Fallback: preserve at least the type if re-decoding ever fails. + self = RelayEvent(type: event.type) + } +} diff --git a/packages/sdk-swift/Tests/AgentRelaySDKTests/AgentRelaySDKTests.swift b/packages/sdk-swift/Tests/AgentRelaySDKTests/AgentRelaySDKTests.swift index d1f9f0f85..1c1cd08b3 100644 --- a/packages/sdk-swift/Tests/AgentRelaySDKTests/AgentRelaySDKTests.swift +++ b/packages/sdk-swift/Tests/AgentRelaySDKTests/AgentRelaySDKTests.swift @@ -1,436 +1,173 @@ import Foundation import XCTest +import Relaycast @testable import AgentRelaySDK -private actor MockHostedHTTP: HostedHTTPClient { - struct Request: Sendable { - let method: String - let path: String - let body: Data? - } - - private var requests: [Request] = [] - private var getResponses: [String: Data] - private var postResponses: [String: Data] - private var postErrors: [String: Error] - - init(getResponses: [String: Data] = [:], postResponses: [String: Data] = [:], postErrors: [String: Error] = [:]) { - self.getResponses = getResponses - self.postResponses = postResponses - self.postErrors = postErrors - } - - func get(path: String, query: [String: String]?) async throws -> Data { - requests.append(Request(method: "GET", path: path, body: nil)) - return getResponses[path] ?? Self.envelope(#"{}"#) - } - - func post(path: String, body: Data?) async throws -> Data { - requests.append(Request(method: "POST", path: path, body: body)) - if let error = postErrors[path] { - throw error - } - return postResponses[path] ?? Self.envelope(#"{}"#) - } - - func delete(path: String) async throws -> Data { - requests.append(Request(method: "DELETE", path: path, body: nil)) - return Self.envelope(#"{}"#) - } - - func allRequests() -> [Request] { - requests - } - - static func envelope(_ dataJSON: String) -> Data { - Data(#"{"ok":true,"data":\#(dataJSON)}"#.utf8) - } -} - -private actor MockHostedTransport: HostedEventTransportClient { - nonisolated let inbound: AsyncStream - - private let continuation: AsyncStream.Continuation - private var sent: [Data] = [] - private var connectCount = 0 - private var onConnect: (@Sendable () async -> Void)? - - init() { - var continuationRef: AsyncStream.Continuation? - self.inbound = AsyncStream { continuation in - continuationRef = continuation - } - self.continuation = continuationRef! - } - - func setOnConnect(_ handler: @escaping @Sendable () async -> Void) async { - onConnect = handler - } - - func connect() async throws { - connectCount += 1 - } - - func disconnect() async { - continuation.finish() - } - - func send(_ message: Data) async throws { - sent.append(message) - } +/// These tests cover the hosted-participant facade that now wraps the relaycast +/// Swift engine SDK. Network-dependent behaviour (registration, posting, the +/// realtime socket) is exercised by the relaycast package's own test suite; here +/// we verify the facade configuration and the bridging glue that keeps +/// AgentRelaySDK's public surface intact on top of relaycast. +final class HostedParticipantSDKTests: XCTestCase { - func emit(_ json: String) { - continuation.yield(Data(json.utf8)) - } + // MARK: - Facade configuration - func sentMessages() -> [Data] { - sent + func testClientInitDefaultsToHostedGateway() { + let client = AgentRelayClient(apiKey: "rk_test_key") + XCTAssertEqual(client.workspaceKey, "rk_test_key") + XCTAssertEqual(client.baseURL.absoluteString, "https://cast.agentrelay.com") } - func connections() -> Int { - connectCount + func testClientHonoursExplicitBaseURL() { + let client = AgentRelayClient(apiKey: "rk_test_key", baseURL: URL(string: "https://example.test")!) + XCTAssertEqual(client.baseURL.absoluteString, "https://example.test") } - func triggerReconnect() async { - await onConnect?() + func testRelayCastInitUsesConfiguredHost() throws { + // Sanity-check that relaycast accepts the preserved host explicitly. + let relay = try Relaycast.RelayCast( + options: Relaycast.RelayCastOptions(apiKey: "rk_test", baseURL: "https://gateway.relaycast.dev") + ) + XCTAssertEqual(relay.client.baseURL.absoluteString, "https://gateway.relaycast.dev") } -} -private actor StringRecorder { - private var values: [String] = [] + // MARK: - Type bridging - func append(_ value: String) { - values.append(value) + func testAgentTypeBridging() { + XCTAssertEqual(RelayAgentType.agent.relaycastType, .agent) + XCTAssertEqual(RelayAgentType.human.relaycastType, .human) + XCTAssertEqual(RelayAgentType.system.relaycastType, .system) } - func all() -> [String] { - values + func testAgentStatusBridging() { + XCTAssertEqual(RelayAgentStatus(Relaycast.AgentStatus.online), .online) + XCTAssertEqual(RelayAgentStatus(Relaycast.AgentStatus.offline), .offline) + XCTAssertEqual(RelayAgentStatus(Relaycast.AgentStatus.away), .away) } -} -private func jsonObject(_ data: Data) throws -> [String: Any] { - try XCTUnwrap(JSONSerialization.jsonObject(with: data) as? [String: Any]) -} - -private func waitForRequest( - in http: MockHostedHTTP, - timeout: TimeInterval = 1.0, - matching predicate: (MockHostedHTTP.Request) throws -> Bool -) async throws -> MockHostedHTTP.Request { - let deadline = Date().addingTimeInterval(timeout) - while Date() < deadline { - let requests = await http.allRequests() - for request in requests.reversed() { - if try predicate(request) { - return request - } + func testJSONValueBridgingPreservesShape() { + let relaycastValue: Relaycast.JSONValue = .object([ + "text": .string("hi"), + "count": .int(3), + "ratio": .double(1.5), + "flag": .bool(true), + "items": .array([.string("a"), .int(2)]), + "nothing": .null + ]) + let bridged = JSONValue(relaycastValue) + guard case .object(let object) = bridged else { + return XCTFail("Expected object") } - try await Task.sleep(nanoseconds: 10_000_000) - } - XCTFail("Timed out waiting for matching request") - return MockHostedHTTP.Request(method: "GET", path: "", body: nil) -} - -final class HostedParticipantSDKTests: XCTestCase { - func testClientInitUsesProvidedBaseURL() { - let client = AgentRelayClient(apiKey: "rk_test_key", baseURL: URL(string: "https://relay.example.com")!) - XCTAssertEqual(client.workspaceKey, "rk_test_key") - XCTAssertEqual(client.baseURL.absoluteString, "https://relay.example.com") - } - - func testHostedWebSocketURLUsesV1WSAndToken() { - let url = RelayEventTransport.resolveWebSocketURL( - baseURL: URL(string: "https://relay.example.com")!, - token: "at_test" + XCTAssertEqual(object["text"], .string("hi")) + XCTAssertEqual(object["count"], .number(3)) + XCTAssertEqual(object["ratio"], .number(1.5)) + XCTAssertEqual(object["flag"], .bool(true)) + XCTAssertEqual(object["items"], .array([.string("a"), .number(2)])) + XCTAssertEqual(object["nothing"], .null) + } + + func testErrorBridgingMapsConflictToAlreadyExists() { + let relaycastError = Relaycast.RelayError.api( + code: "some_code", + message: "name_taken", + statusCode: 409, + retryable: false ) - XCTAssertEqual(url?.scheme, "wss") - XCTAssertEqual(url?.path, "/v1/ws") - XCTAssertTrue(url?.query?.contains("token=at_test") == true) + let bridged = RelayError(relaycastError) + guard case .protocolError(let code, let message, _) = bridged else { + return XCTFail("Expected protocolError") + } + XCTAssertEqual(code, "agent_already_exists") + XCTAssertEqual(message, "name_taken") } - func testHostedAPIURLAppendsV1Path() { - let url = HostedHTTP.resolveAPIURL(baseURL: URL(string: "https://relay.example.com")!, path: "/v1/dm") - XCTAssertEqual(url?.absoluteString, "https://relay.example.com/v1/dm") + func testErrorBridgingMapsNotConnected() { + if case .notConnected = RelayError(Relaycast.RelayError.notConnected) { + // ok + } else { + XCTFail("Expected notConnected") + } } - func testRegisterOrRotateTreatsAgentAlreadyExistsAsConflict() async throws { - let workspaceHTTP = MockHostedHTTP( - getResponses: [ - "/v1/agents/swift-agent": MockHostedHTTP.envelope( - #"{"id":"ag_existing","name":"swift-agent","type":"agent","status":"online","created_at":"2026-01-01T00:00:00Z"}"# - ) - ], - postResponses: [ - "/v1/agents/swift-agent/rotate-token": MockHostedHTTP.envelope(#"{"token":"at_rotated"}"#) - ], - postErrors: [ - "/v1/agents": RelayError.protocolError( - code: "agent_already_exists", - message: "name_taken", - retryable: false - ) - ] - ) - let core = HostedWorkspaceCore( - workspaceKey: "rk_test", - baseURL: URL(string: "https://relay.example.com")!, - http: workspaceHTTP - ) + // MARK: - Realtime event glue - let registration = try await core.registerOrRotate(name: "swift-agent", type: .agent) + func testRelayEventFromWsEventExtractsMessageFields() { + // relaycast emits flat events: type plus payload fields at the top level. + let wsEvent = Relaycast.WsEvent(type: "message.created", payload: [ + "channel": .string("general"), + "message": .object([ + "id": .string("msg_1"), + "message_id": .string("msg_1"), + "body": .string("hello"), + "from": .object(["name": .string("alice")]), + "channel": .object(["name": .string("general")]) + ]) + ]) - XCTAssertEqual(registration.id, "ag_existing") - XCTAssertEqual(registration.name, "swift-agent") - XCTAssertEqual(registration.token, "at_rotated") - let requests = await workspaceHTTP.allRequests() - XCTAssertEqual( - requests.map { "\($0.method) \($0.path)" }, - [ - "POST /v1/agents", - "GET /v1/agents/swift-agent", - "POST /v1/agents/swift-agent/rotate-token" - ] - ) + let event = RelayEvent(wsEvent) + XCTAssertEqual(event.type, "message.created") + XCTAssertEqual(event.channel, "general") + XCTAssertEqual(event.message?.text, "hello") + XCTAssertEqual(event.message?.from.name, "alice") + XCTAssertEqual(event.message?.channel?.name, "general") } - func testAgentClientPostsChannelAndDirectMessagesToHostedEndpoints() async throws { - let workspaceHTTP = MockHostedHTTP() - let agentHTTP = MockHostedHTTP() - let transport = MockHostedTransport() - let core = HostedParticipantCore( - agentId: "ag_1", - agentName: "swift-agent", - token: "at_test", - baseURL: URL(string: "https://relay.example.com")!, - workspaceHTTP: workspaceHTTP, - agentHTTP: agentHTTP, - transport: transport - ) - let client = AgentClient(core: core, id: "ag_1", name: "swift-agent", token: "at_test") - - try await client.post(to: "#general", message: "hello") - try await client.dm(to: "@reviewer", message: "ping") + func testRelayEventFromWsEventExtractsActionInvocationFields() { + let wsEvent = Relaycast.WsEvent(type: "action.invoked", payload: [ + "invocation_id": .string("inv_1"), + "action_name": .string("echo"), + "caller_name": .string("alice") + ]) - let requests = await agentHTTP.allRequests() - XCTAssertEqual(requests.map(\.path), ["/v1/channels/general/messages", "/v1/dm"]) - let channelBody = try jsonObject(try XCTUnwrap(requests[0].body)) - XCTAssertEqual(channelBody["text"] as? String, "hello") - XCTAssertEqual(channelBody["mode"] as? String, "wait") - let dmBody = try jsonObject(try XCTUnwrap(requests[1].body)) - XCTAssertEqual(dmBody["to"] as? String, "reviewer") - XCTAssertEqual(dmBody["text"] as? String, "ping") + let event = RelayEvent(wsEvent) + XCTAssertEqual(event.type, "action.invoked") + XCTAssertEqual(event.invocationId, "inv_1") + XCTAssertEqual(event.actionName, "echo") + XCTAssertEqual(event.callerName, "alice") } - func testChannelSubscribeSendsHostedSubscribeFrameAndRoutesMessages() async throws { - let transport = MockHostedTransport() - let core = HostedParticipantCore( - agentId: "ag_1", - agentName: "swift-agent", - token: "at_test", - baseURL: URL(string: "https://relay.example.com")!, - workspaceHTTP: MockHostedHTTP(), - agentHTTP: MockHostedHTTP(), - transport: transport - ) - let client = AgentClient(core: core, id: "ag_1", name: "swift-agent", token: "at_test") - let channel = client.channel("general") - try await channel.subscribe() - var sentMessages = await transport.sentMessages() - XCTAssertEqual(sentMessages.count, 1) - await transport.triggerReconnect() - sentMessages = await transport.sentMessages() - XCTAssertEqual(sentMessages.count, 2) - - let messages = channel.events - let task = Task { () -> RelayChannelEvent? in - for await event in messages { - return event - } - return nil - } - - await transport.emit( - """ - { - "type": "message.received", - "payload": { - "channel": "general", - "message": { - "id": "msg_1", - "body": "hello", - "from": {"name": "alice"}, - "channel": {"name": "general"} - } - } - } - """ - ) - - let maybeEvent = await task.value - let event = try XCTUnwrap(maybeEvent) - XCTAssertEqual(event.from, "alice") - XCTAssertEqual(event.body, "hello") - XCTAssertEqual(event.channel, "general") + func testRelayEventFromWsEventPreservesTypeWithEmptyPayload() { + let event = RelayEvent(Relaycast.WsEvent(type: "pong")) + XCTAssertEqual(event.type, "pong") + XCTAssertNil(event.message) } - func testRegisterActionPostsDescriptorAndCompletesHostedInvocation() async throws { - let invocationPath = "/v1/actions/echo/invocations/inv_1" - let completionPath = "\(invocationPath)/complete" - let workspaceHTTP = MockHostedHTTP() - let agentHTTP = MockHostedHTTP( - getResponses: [ - invocationPath: MockHostedHTTP.envelope( - #"{"invocation_id":"inv_1","action_name":"echo","caller_name":"alice","input":{"text":"hello"},"status":"invoked"}"# - ) - ], - postResponses: [ - completionPath: MockHostedHTTP.envelope( - #"{"invocation_id":"inv_1","action_name":"echo","status":"completed","output":{"ok":true},"error":null,"duration_ms":1,"completed_at":"2026-01-01T00:00:00Z"}"# - ) - ] - ) - let transport = MockHostedTransport() - let core = HostedParticipantCore( - agentId: "ag_1", - agentName: "swift-agent", - token: "at_test", - baseURL: URL(string: "https://relay.example.com")!, - workspaceHTTP: workspaceHTTP, - agentHTTP: agentHTTP, - transport: transport - ) - let client = AgentClient(core: core, id: "ag_1", name: "swift-agent", token: "at_test") - let recorder = StringRecorder() + // MARK: - Registration lifecycle - _ = try await client.registerAction( - name: "echo", - description: "Echo input", - inputSchemaJSON: #"{"type":"object","properties":{"text":{"type":"string"}}}"# - ) { input in - await recorder.append(input) - return #"{"ok":true}"# - } + /// A persisted `AgentRegistration` must stay usable even if the owning + /// `AgentRelay`/`HostedWorkspaceCore` is released before `asClient()` is + /// called. The factory captures the relay/baseURL transport state strongly, + /// so this must not trap. + func testRegistrationAsClientSurvivesCoreRelease() throws { + let response = try makeCreateAgentResponse(id: "agent_1", name: "swift-agent", token: "rk_token") - let descriptor = try await waitForRequest(in: workspaceHTTP) { request in - request.method == "POST" && request.path == "/v1/actions" - } - let descriptorBody = try jsonObject(try XCTUnwrap(descriptor.body)) - XCTAssertEqual(descriptorBody["name"] as? String, "echo") - XCTAssertEqual(descriptorBody["description"] as? String, "Echo input") - XCTAssertEqual(descriptorBody["handler_agent"] as? String, "swift-agent") - XCTAssertEqual((descriptorBody["input_schema"] as? [String: Any])?["type"] as? String, "object") - - await transport.emit( - """ - { - "type": "action.invoked", - "payload": { - "type": "action.invoked", - "invocation_id": "inv_1", - "action_name": "echo", - "caller_name": "alice" - } - } - """ + var core: HostedWorkspaceCore? = HostedWorkspaceCore( + workspaceKey: "rk_test_key", + baseURL: URL(string: "https://gateway.relaycast.dev")! ) + let registration = core!.makeRegistration(response) - let completion = try await waitForRequest(in: agentHTTP) { request in - request.method == "POST" && request.path == completionPath - } - let completionBody = try jsonObject(try XCTUnwrap(completion.body)) - XCTAssertEqual((completionBody["output"] as? [String: Any])?["ok"] as? Bool, true) + // Drop the owning core before using the persisted registration. + core = nil - let inputs = await recorder.all() - XCTAssertEqual(inputs.count, 1) - let input = try jsonObject(Data(inputs[0].utf8)) - XCTAssertEqual(input["text"] as? String, "hello") + let client = registration.asClient() + XCTAssertEqual(client.id, "agent_1") + XCTAssertEqual(client.name, "swift-agent") + XCTAssertEqual(client.token, "rk_token") } - func testRegisterActionWrapsScalarHandlerOutput() async throws { - let invocationPath = "/v1/actions/plain/invocations/inv_scalar" - let completionPath = "\(invocationPath)/complete" - let workspaceHTTP = MockHostedHTTP() - let agentHTTP = MockHostedHTTP( - getResponses: [ - invocationPath: MockHostedHTTP.envelope( - #"{"invocation_id":"inv_scalar","action_name":"plain","caller_name":"alice","input":{},"status":"invoked"}"# - ) - ], - postResponses: [ - completionPath: MockHostedHTTP.envelope( - #"{"invocation_id":"inv_scalar","action_name":"plain","status":"completed","output":{"value":"done"},"error":null}"# - ) - ] - ) - let transport = MockHostedTransport() - let core = HostedParticipantCore( - agentId: "ag_1", - agentName: "swift-agent", - token: "at_test", - baseURL: URL(string: "https://relay.example.com")!, - workspaceHTTP: workspaceHTTP, - agentHTTP: agentHTTP, - transport: transport - ) - let client = AgentClient(core: core, id: "ag_1", name: "swift-agent", token: "at_test") - - _ = try await client.registerAction( - name: "plain", - description: "Plain output", - inputSchemaJSON: #"{"type":"object"}"# - ) { _ in - "done" - } - - await transport.emit( - """ - { - "type": "action.invoked", - "payload": { - "action_name": "plain", - "invocation_id": "inv_scalar", - "caller_name": "alice" - } - } - """ - ) - - let completion = try await waitForRequest(in: agentHTTP) { request in - request.method == "POST" && request.path == completionPath - } - let completionBody = try jsonObject(try XCTUnwrap(completion.body)) - XCTAssertEqual((completionBody["output"] as? [String: Any])?["value"] as? String, "done") + private func makeCreateAgentResponse(id: String, name: String, token: String) throws -> Relaycast.CreateAgentResponse { + let json = """ + {"id":"\(id)","name":"\(name)","token":"\(token)","status":"online","createdAt":"2026-06-20T00:00:00Z"} + """ + return try JSONDecoder().decode(Relaycast.CreateAgentResponse.self, from: Data(json.utf8)) } - func testUnregisterActionDeletesHostedDescriptor() async throws { - let workspaceHTTP = MockHostedHTTP() - let core = HostedParticipantCore( - agentId: "ag_1", - agentName: "swift-agent", - token: "at_test", - baseURL: URL(string: "https://relay.example.com")!, - workspaceHTTP: workspaceHTTP, - agentHTTP: MockHostedHTTP(), - transport: MockHostedTransport() - ) - let client = AgentClient(core: core, id: "ag_1", name: "swift-agent", token: "at_test") + // MARK: - Action handle lifecycle - let handle = try await client.registerAction( - name: "cleanup", - description: "Cleanup", - inputSchemaJSON: #"{"type":"object"}"# - ) { _ in - "{}" - } + func testActionHandleExposesName() async { + let handle = ActionHandle(name: "echo") { } + XCTAssertEqual(handle.name, "echo") await handle.unregister() - - let requests = await workspaceHTTP.allRequests() - XCTAssertTrue(requests.contains { $0.method == "DELETE" && $0.path == "/v1/actions/cleanup" }) - } -} - -private extension HostedParticipantCore { - func transportAsMock() -> MockHostedTransport? { - transport as? MockHostedTransport } }