import Foundation import NIO import NIOHTTP1 import NIOWebSocket import ForgeCamera import ForgeGrade import ForgeColor /// RCP2 param IDs (wire form: RCP_PARAM_ prefix stripped). /// Verified against RED parameter docs via docs/research/red-rcp2-protocol.md. enum RedParams { static let recordState = "RECORD_STATE" // cur.val: 0 idle, 1 recording, set "2" = toggle static let clipName = "CLIP_NAME" static let iso = "ISO" static let colorTemp = "COLOR_TEMPERATURE" static let tint = "TINT" static let sensorFPS = "SENSOR_FRAME_RATE" static let timecode = "TIMECODE" static let cdlEnable = "CDL_ENABLE" static let cdlSlopeR = "CDL_SLOPE_RED" static let cdlSlopeG = "CDL_SLOPE_GREEN" static let cdlSlopeB = "CDL_SLOPE_BLUE" static let cdlOffsetR = "CDL_OFFSET_RED" static let cdlOffsetG = "CDL_OFFSET_GREEN" static let cdlOffsetB = "CDL_OFFSET_BLUE" static let cdlPowerR = "CDL_POWER_RED" static let cdlPowerG = "CDL_POWER_GREEN" static let cdlPowerB = "CDL_POWER_BLUE" static let cdlSaturation = "CDL_SATURATION" } /// DSMC3 driver — real RCP2: WebSocket ws://:9998/rcp, JSON text frames. /// Handshake: rcp_config w/ client name. CDL pushed as per-channel rcp_set ops. /// NO live LUT upload (RCP2 limitation) — capabilities.lut3dSizes empty; grade /// pipeline keeps camera-native CDL only, remainder must stay identity. public actor RedDriver: CameraDriver { public nonisolated let capabilities = CameraCapabilities( vendor: .red, supportsNativeCDL: true, lut3dSizes: [], metadataFields: [.clipName, .exposureIndex, .whiteBalance, .tint, .fps, .timecode]) private let host: String private let port: Int private let pollInterval: Duration private var group: MultiThreadedEventLoopGroup? private var channel: Channel? private var pollTask: Task? /// Latest values per param, merged into CameraState on any change. private var paramCache: [String: String] = [:] private var lastState = CameraState(connection: .disconnected) private var stateContinuations: [UUID: AsyncStream.Continuation] = [:] public init(host: String, port: Int = 9998, pollInterval: Duration = .milliseconds(150)) { self.host = host self.port = port self.pollInterval = pollInterval } public nonisolated var state: AsyncStream { AsyncStream { continuation in let id = UUID() Task { await self.register(id: id, continuation: continuation) } continuation.onTermination = { _ in Task { await self.unregister(id: id) } } } } private func register(id: UUID, continuation: AsyncStream.Continuation) { stateContinuations[id] = continuation } private func unregister(id: UUID) { stateContinuations.removeValue(forKey: id) } private func emit(_ s: CameraState) { lastState = s for c in stateContinuations.values { c.yield(s) } } // MARK: WebSocket transport public func connect() async throws { await teardown() let group = MultiThreadedEventLoopGroup(numberOfThreads: 1) self.group = group let upgradeResult: EventLoopPromise = group.next().makePromise(of: Channel.self) let driver = self let upgrader = NIOWebSocketClientUpgrader( requestKey: Data((0..<16).map { _ in UInt8.random(in: 0...255) }).base64EncodedString(), upgradePipelineHandler: { channel, _ in channel.pipeline.addHandler(RedClientWSHandler(driver: driver)).map { upgradeResult.succeed(channel) } }) let bootstrap = ClientBootstrap(group: group) .channelInitializer { channel in let upgradeConfig = NIOHTTPClientUpgradeConfiguration( upgraders: [upgrader], completionHandler: { _ in channel.pipeline.removeHandler(name: "sendInitialRequest", promise: nil) }) return channel.pipeline.addHTTPClientHandlers(withClientUpgrade: upgradeConfig).flatMap { channel.pipeline.addHandler( SendUpgradeRequestHandler(host: self.host, port: self.port, path: "/rcp"), name: "sendInitialRequest") } } do { let ch = try await bootstrap.connect(host: host, port: port).get() _ = ch // Wait for upgrade completion (or fail after timeout). let upgraded: Channel = try await withThrowingTaskGroup(of: Channel.self) { grp in grp.addTask { try await upgradeResult.futureResult.get() } grp.addTask { try await Task.sleep(for: .seconds(3)) throw CameraError.connectionFailed("WS upgrade timeout") } let first = try await grp.next()! grp.cancelAll() return first } channel = upgraded } catch { upgradeResult.fail(error) await teardown() throw CameraError.connectionFailed("connect \(host):\(port): \(error)") } // RCP2 handshake. try send([ "type": "rcp_config", "strings_decoded": 0, "json_minified": 1, "include_cacheable_flags": 0, "encoding_type": "html", "client": ["name": "Forge"], ]) emit(CameraState(connection: .connected)) startPolling() } public func disconnect() async { pollTask?.cancel() pollTask = nil await teardown() emit(CameraState(connection: .disconnected)) } private func teardown() async { if let ch = channel { _ = try? await ch.close() } channel = nil if let g = group { try? await g.shutdownGracefully() } group = nil } private func send(_ obj: [String: Any]) throws { guard let ch = channel else { throw CameraError.connectionFailed("not connected") } let data = try JSONSerialization.data(withJSONObject: obj, options: [.sortedKeys]) var buf = ch.allocator.buffer(capacity: data.count) buf.writeBytes(data) let frame = WebSocketFrame(fin: true, opcode: .text, maskKey: randomMask(), data: buf) ch.writeAndFlush(NIOAny(frame), promise: nil) } private func randomMask() -> WebSocketMaskingKey { WebSocketMaskingKey([UInt8.random(in: 0...255), .random(in: 0...255), .random(in: 0...255), .random(in: 0...255)])! } // MARK: Inbound (called from WS handler) func handleInbound(_ data: Data) { guard let obj = try? JSONSerialization.jsonObject(with: data) as? [String: Any], let id = obj["id"] as? String else { return } // Value under cur.val (int/double/string). guard let cur = obj["cur"] as? [String: Any], let val = cur["val"] else { return } let stringValue: String if let s = val as? String { stringValue = s } else if let i = val as? Int { stringValue = String(i) } else if let d = val as? Double { stringValue = String(d) } else { return } paramCache[id] = stringValue rebuildState() } private func rebuildState() { let newState = CameraState( connection: .connected, isRecording: paramCache[RedParams.recordState] == "1", clipName: (paramCache[RedParams.clipName]?.isEmpty ?? true) ? nil : paramCache[RedParams.clipName], metadata: CameraMetadata( exposureIndex: paramCache[RedParams.iso].flatMap(Int.init), whiteBalance: paramCache[RedParams.colorTemp].flatMap(Int.init), tint: paramCache[RedParams.tint].flatMap(Int.init), fps: paramCache[RedParams.sensorFPS].flatMap(Double.init), timecode: paramCache[RedParams.timecode])) if newState != lastState { emit(newState) } } func handleChannelClosed() { if lastState.connection == .connected { emit(CameraState(connection: .disconnected)) } } // MARK: Polling (rcp_get; responses arrive via handleInbound) private func startPolling() { pollTask?.cancel() pollTask = Task { let polled = [ RedParams.recordState, RedParams.clipName, RedParams.iso, RedParams.colorTemp, RedParams.tint, RedParams.sensorFPS, RedParams.timecode, ] while !Task.isCancelled { for p in polled { try? self.send(["type": "rcp_get", "id": p]) } try? await Task.sleep(for: pollInterval) } } } // MARK: Look push public func push(look: FlattenedLook) async throws { // RCP2 cannot upload LUTs. Non-identity LUT = grade won't match on camera. if let lut = look.lut, !lut.isApproximatelyIdentity() { throw CameraError.unsupportedLook( "RED RCP2 cannot receive live 3D LUTs — use CDL-only grade or pre-load LUT on camera") } guard let cdl = look.cdl else { throw CameraError.unsupportedLook("nothing to push (no CDL)") } let sets: [(String, Any)] = [ (RedParams.cdlSlopeR, Double(cdl.slope.x)), (RedParams.cdlSlopeG, Double(cdl.slope.y)), (RedParams.cdlSlopeB, Double(cdl.slope.z)), (RedParams.cdlOffsetR, Double(cdl.offset.x)), (RedParams.cdlOffsetG, Double(cdl.offset.y)), (RedParams.cdlOffsetB, Double(cdl.offset.z)), (RedParams.cdlPowerR, Double(cdl.power.x)), (RedParams.cdlPowerG, Double(cdl.power.y)), (RedParams.cdlPowerB, Double(cdl.power.z)), (RedParams.cdlSaturation, Double(cdl.saturation)), (RedParams.cdlEnable, 1), ] for (param, value) in sets { try send(["type": "rcp_set", "id": param, "value": value]) } } } // MARK: - NIO handlers /// Sends the initial HTTP upgrade request when the channel becomes active. final class SendUpgradeRequestHandler: ChannelInboundHandler, RemovableChannelHandler, @unchecked Sendable { typealias InboundIn = HTTPClientResponsePart typealias OutboundOut = HTTPClientRequestPart private let host: String private let port: Int private let path: String init(host: String, port: Int, path: String) { self.host = host self.port = port self.path = path } func channelActive(context: ChannelHandlerContext) { var headers = HTTPHeaders() headers.add(name: "Host", value: "\(host):\(port)") headers.add(name: "Content-Length", value: "0") let head = HTTPRequestHead(version: .http1_1, method: .GET, uri: path, headers: headers) context.write(wrapOutboundOut(.head(head)), promise: nil) context.writeAndFlush(wrapOutboundOut(.end(nil)), promise: nil) context.fireChannelActive() } } /// Inbound WS frames -> driver actor. final class RedClientWSHandler: ChannelInboundHandler, @unchecked Sendable { typealias InboundIn = WebSocketFrame private let driver: RedDriver init(driver: RedDriver) { self.driver = driver } func channelRead(context: ChannelHandlerContext, data: NIOAny) { let frame = unwrapInboundIn(data) switch frame.opcode { case .text: var buf = frame.unmaskedData guard let bytes = buf.readBytes(length: buf.readableBytes) else { return } let payload = Data(bytes) Task { await driver.handleInbound(payload) } case .connectionClose: context.close(promise: nil) default: break } } func channelInactive(context: ChannelHandlerContext) { Task { await driver.handleChannelClosed() } context.fireChannelInactive() } }