From acce32c493cd20a6dc8ed58bc2cf467f10b819c9 Mon Sep 17 00:00:00 2001 From: Forge Dev Date: Fri, 10 Jul 2026 20:00:36 +0000 Subject: [PATCH] =?UTF-8?q?driver-red:=20RCP2-style=20TCP=20driver=20(leng?= =?UTF-8?q?th-prefixed=20JSON)=20=E2=80=94=20handshake,=20param=20polling,?= =?UTF-8?q?=20IPP2=20CDL=20slot=20+=20LUT=20upload;=20RED=20sim=20?= =?UTF-8?q?=E2=80=94=207=20tests?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Sources/ForgeCameraRED/ForgeCameraRED.swift | 1 - Sources/ForgeCameraRED/RedDriver.swift | 245 ++++++++++++++++++ Sources/ForgeSim/RedSimulator.swift | 135 ++++++++++ .../ForgeCameraREDTests/RedDriverTests.swift | 151 +++++++++++ 4 files changed, 531 insertions(+), 1 deletion(-) delete mode 100644 Sources/ForgeCameraRED/ForgeCameraRED.swift create mode 100644 Sources/ForgeCameraRED/RedDriver.swift create mode 100644 Sources/ForgeSim/RedSimulator.swift create mode 100644 Tests/ForgeCameraREDTests/RedDriverTests.swift diff --git a/Sources/ForgeCameraRED/ForgeCameraRED.swift b/Sources/ForgeCameraRED/ForgeCameraRED.swift deleted file mode 100644 index 7ca92e5..0000000 --- a/Sources/ForgeCameraRED/ForgeCameraRED.swift +++ /dev/null @@ -1 +0,0 @@ -// ForgeCameraRED diff --git a/Sources/ForgeCameraRED/RedDriver.swift b/Sources/ForgeCameraRED/RedDriver.swift new file mode 100644 index 0000000..e0adf76 --- /dev/null +++ b/Sources/ForgeCameraRED/RedDriver.swift @@ -0,0 +1,245 @@ +import Foundation +import ForgeCamera +import ForgeGrade +import ForgeColor + +#if canImport(Glibc) +import Glibc +#endif + +/// RED parameter names — single source of truth, correct against RED SDK +/// (RCP2) when hardware/SDK agreement lands. Sim mirrors these. +enum RedParams { + static let recordState = "RECORD_STATE" + static let clipName = "CLIP_NAME" + static let iso = "ISO" + static let colorTemp = "COLOR_TEMP" + static let tint = "TINT" + static let sensorFPS = "SENSOR_FPS" + static let timecode = "TIMECODE" + static let ipp2CDL = "IPP2_CDL" +} + +/// DSMC3 driver: RCP2-style length-prefixed JSON over TCP. +/// Request/response serialized through the actor — one in-flight request. +public actor RedDriver: CameraDriver { + public nonisolated let capabilities = CameraCapabilities( + vendor: .red, + supportsNativeCDL: true, + lut3dSizes: [17, 33], + metadataFields: [.clipName, .exposureIndex, .whiteBalance, .tint, .fps, .timecode]) + + private let host: String + private let port: Int + private let pollInterval: Duration + + private var socketFD: Int32 = -1 + private var pollTask: Task? + private var lastState = CameraState(connection: .disconnected) + private var stateContinuations: [UUID: AsyncStream.Continuation] = [:] + + public init(host: String, port: Int, 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: Socket transport (blocking, actor-serialized, short timeouts) + + private func openSocket() throws { + let fd = socket(AF_INET, Int32(SOCK_STREAM.rawValue), 0) + guard fd >= 0 else { throw CameraError.connectionFailed("socket() failed") } + + var tv = timeval(tv_sec: 3, tv_usec: 0) + setsockopt(fd, SOL_SOCKET, SO_RCVTIMEO, &tv, socklen_t(MemoryLayout.size)) + setsockopt(fd, SOL_SOCKET, SO_SNDTIMEO, &tv, socklen_t(MemoryLayout.size)) + + var addr = sockaddr_in() + addr.sin_family = sa_family_t(AF_INET) + addr.sin_port = in_port_t(UInt16(port).bigEndian) + guard inet_pton(AF_INET, host, &addr.sin_addr) == 1 else { + close(fd) + throw CameraError.connectionFailed("bad host \(host)") + } + let rc = withUnsafePointer(to: &addr) { ptr in + ptr.withMemoryRebound(to: sockaddr.self, capacity: 1) { sa in + Glibc.connect(fd, sa, socklen_t(MemoryLayout.size)) + } + } + guard rc == 0 else { + close(fd) + throw CameraError.connectionFailed("connect to \(host):\(port) failed") + } + socketFD = fd + } + + private func closeSocket() { + if socketFD >= 0 { + close(socketFD) + socketFD = -1 + } + } + + private func sendAll(_ data: Data) throws { + var sent = 0 + try data.withUnsafeBytes { (buf: UnsafeRawBufferPointer) in + while sent < buf.count { + let n = write(socketFD, buf.baseAddress!.advanced(by: sent), buf.count - sent) + guard n > 0 else { throw CameraError.connectionFailed("write failed") } + sent += n + } + } + } + + private func recvExact(_ count: Int) throws -> Data { + var out = Data(capacity: count) + var remaining = count + var chunk = [UInt8](repeating: 0, count: 64 * 1024) + while remaining > 0 { + let n = read(socketFD, &chunk, min(remaining, chunk.count)) + guard n > 0 else { throw CameraError.connectionFailed("read failed/timeout") } + out.append(contentsOf: chunk[0.. [String: Any] { + guard socketFD >= 0 else { throw CameraError.connectionFailed("not connected") } + let payload = try JSONSerialization.data(withJSONObject: obj, options: [.sortedKeys]) + var frame = Data() + var lenBE = UInt32(payload.count).bigEndian + withUnsafeBytes(of: &lenBE) { frame.append(contentsOf: $0) } + frame.append(payload) + try sendAll(frame) + + let header = try recvExact(4) + // Byte-wise assembly — Data slices may be misaligned for a u32 load. + let length = header.reduce(UInt32(0)) { ($0 << 8) | UInt32($1) } + let body = try recvExact(Int(length)) + guard let reply = try JSONSerialization.jsonObject(with: body) as? [String: Any] else { + throw CameraError.connectionFailed("bad reply JSON") + } + return reply + } + + private func getParam(_ name: String) throws -> String { + let reply = try request(["type": "get", "param": name]) + guard reply["ok"] as? Bool == true, let value = reply["value"] as? String else { + throw CameraError.connectionFailed("get \(name) failed") + } + return value + } + + private func setParam(_ name: String, _ value: String) throws { + let reply = try request(["type": "set", "param": name, "value": value]) + guard reply["ok"] as? Bool == true else { + throw CameraError.pushFailed("set \(name) failed") + } + } + + // MARK: CameraDriver + + public func connect() async throws { + closeSocket() + try openSocket() + let reply = try request(["type": "handshake"]) + guard reply["ok"] as? Bool == true else { + closeSocket() + throw CameraError.connectionFailed("handshake rejected") + } + emit(CameraState(connection: .connected)) + startPolling() + } + + public func disconnect() async { + pollTask?.cancel() + pollTask = nil + closeSocket() + emit(CameraState(connection: .disconnected)) + } + + public func push(look: FlattenedLook) async throws { + if let lut = look.lut, !capabilities.lut3dSizes.contains(lut.size) { + throw CameraError.unsupportedLook("LUT size \(lut.size) not in \(capabilities.lut3dSizes)") + } + if let cdl = look.cdl { + let obj: [String: Any] = [ + "slope": [cdl.slope.x, cdl.slope.y, cdl.slope.z], + "offset": [cdl.offset.x, cdl.offset.y, cdl.offset.z], + "power": [cdl.power.x, cdl.power.y, cdl.power.z], + "saturation": cdl.saturation, + ] + let json = String(data: try JSONSerialization.data(withJSONObject: obj, options: [.sortedKeys]), encoding: .utf8)! + try setParam(RedParams.ipp2CDL, json) + } + if let lut = look.lut { + let reply = try request(["type": "lut", "size": lut.size, "table": lut.table]) + guard reply["ok"] as? Bool == true else { + throw CameraError.pushFailed("LUT upload rejected") + } + } + } + + // MARK: Polling + + private func startPolling() { + pollTask?.cancel() + pollTask = Task { + while !Task.isCancelled { + pollOnce() + try? await Task.sleep(for: pollInterval) + } + } + } + + private func pollOnce() { + do { + let rec = try getParam(RedParams.recordState) == "1" + let clip = try getParam(RedParams.clipName) + let newState = CameraState( + connection: .connected, + isRecording: rec, + clipName: clip.isEmpty ? nil : clip, + metadata: CameraMetadata( + exposureIndex: Int(try getParam(RedParams.iso)), + whiteBalance: Int(try getParam(RedParams.colorTemp)), + tint: Int(try getParam(RedParams.tint)), + fps: Double(try getParam(RedParams.sensorFPS)), + timecode: try getParam(RedParams.timecode))) + if newState != lastState { + emit(newState) + } + } catch { + if lastState.connection == .connected { + emit(CameraState(connection: .disconnected)) + } + } + } +} diff --git a/Sources/ForgeSim/RedSimulator.swift b/Sources/ForgeSim/RedSimulator.swift new file mode 100644 index 0000000..609df16 --- /dev/null +++ b/Sources/ForgeSim/RedSimulator.swift @@ -0,0 +1,135 @@ +import Foundation +import NIO + +/// RCP2-style TCP simulator: 4-byte big-endian length prefix + JSON message. +/// Message types: handshake, get {param}, set {param,value}, lut {size,table}. +/// Replies mirror type with {ok:true} or value payloads. +/// NOTE: real RCP2 framing/params verified against RED SDK in hardware phase; +/// this encodes the documented shape (length-prefixed JSON over TCP). +public actor RedSimulator { + private var group: MultiThreadedEventLoopGroup? + private var channel: Channel? + public private(set) var boundPort: Int? + + public private(set) var handshakeCount = 0 + public private(set) var params: [String: String] = [ + "RECORD_STATE": "0", + "CLIP_NAME": "", + "ISO": "800", + "COLOR_TEMP": "5600", + "TINT": "0", + "SENSOR_FPS": "24", + "TIMECODE": "00:00:00:00", + "MODEL": "V-RAPTOR-SIM", + ] + public struct UploadedLut: Sendable { + public let size: Int + public let entryCount: Int + } + public private(set) var uploadedLuts: [UploadedLut] = [] + + public init() {} + + public func setParam(_ key: String, value: String) { + params[key] = value + } + + /// Data-in/Data-out across the actor boundary (JSON dicts aren't Sendable). + func handleFrame(_ payload: Data) -> Data { + guard let obj = try? JSONSerialization.jsonObject(with: payload) as? [String: Any] else { + return (try? JSONSerialization.data(withJSONObject: ["ok": false, "error": "bad JSON"])) ?? Data() + } + let reply = handle(message: obj) + return (try? JSONSerialization.data(withJSONObject: reply, options: [.sortedKeys])) ?? Data() + } + + func handle(message: [String: Any]) -> [String: Any] { + switch message["type"] as? String { + case "handshake": + handshakeCount += 1 + return ["type": "handshake", "ok": true, "model": params["MODEL"] ?? ""] + case "get": + guard let param = message["param"] as? String, let value = params[param] else { + return ["type": "get", "ok": false, "error": "unknown param"] + } + return ["type": "get", "ok": true, "param": param, "value": value] + case "set": + guard let param = message["param"] as? String, let value = message["value"] as? String else { + return ["type": "set", "ok": false, "error": "bad set"] + } + params[param] = value + return ["type": "set", "ok": true, "param": param] + case "lut": + guard let size = message["size"] as? Int, + let table = message["table"] as? [Any], + table.count == size * size * size * 3 else { + return ["type": "lut", "ok": false, "error": "bad lut"] + } + uploadedLuts.append(UploadedLut(size: size, entryCount: table.count)) + return ["type": "lut", "ok": true] + default: + return ["ok": false, "error": "unknown type"] + } + } + + public func start(port: Int) async throws { + let group = MultiThreadedEventLoopGroup(numberOfThreads: 1) + self.group = group + let sim = self + let bootstrap = ServerBootstrap(group: group) + .serverChannelOption(ChannelOptions.socketOption(.so_reuseaddr), value: 1) + .childChannelInitializer { channel in + channel.pipeline.addHandler(RCPFrameHandler(sim: sim)) + } + let channel = try await bootstrap.bind(host: "127.0.0.1", port: port).get() + self.channel = channel + boundPort = channel.localAddress?.port + } + + public func stop() async throws { + try await channel?.close() + try await group?.shutdownGracefully() + channel = nil + group = nil + boundPort = nil + } +} + +/// Length-prefixed JSON framing: u32 BE length + payload. +final class RCPFrameHandler: ChannelInboundHandler, @unchecked Sendable { + typealias InboundIn = ByteBuffer + typealias OutboundOut = ByteBuffer + + private let sim: RedSimulator + private var buffer = ByteBuffer() + + init(sim: RedSimulator) { + self.sim = sim + } + + func channelRead(context: ChannelHandlerContext, data: NIOAny) { + var incoming = unwrapInboundIn(data) + buffer.writeBuffer(&incoming) + + while true { + guard buffer.readableBytes >= 4, + let length = buffer.getInteger(at: buffer.readerIndex, as: UInt32.self), + buffer.readableBytes >= 4 + Int(length) else { break } + buffer.moveReaderIndex(forwardBy: 4) + guard let bytes = buffer.readBytes(length: Int(length)) else { break } + let payload = Data(bytes) + + let channel = context.channel + let loop = context.eventLoop + Task { + let replyData = await self.sim.handleFrame(payload) + loop.execute { + var out = channel.allocator.buffer(capacity: replyData.count + 4) + out.writeInteger(UInt32(replyData.count)) + out.writeBytes(replyData) + channel.writeAndFlush(out, promise: nil) + } + } + } + } +} diff --git a/Tests/ForgeCameraREDTests/RedDriverTests.swift b/Tests/ForgeCameraREDTests/RedDriverTests.swift new file mode 100644 index 0000000..b9687a4 --- /dev/null +++ b/Tests/ForgeCameraREDTests/RedDriverTests.swift @@ -0,0 +1,151 @@ +import XCTest +import Foundation +import ForgeCamera +import ForgeGrade +import ForgeColor +import ForgeSim +@testable import ForgeCameraRED + +final class RedDriverTests: XCTestCase { + + var sim: RedSimulator! + + override func setUp() async throws { + sim = RedSimulator() + try await sim.start(port: 0) + } + + override func tearDown() async throws { + try? await sim.stop() + sim = nil + } + + func makeDriver(pollInterval: Duration = .milliseconds(20)) async -> RedDriver { + let port = await sim.boundPort! + return RedDriver(host: "127.0.0.1", port: port, pollInterval: pollInterval) + } + + // Capabilities: RED = native CDL (IPP2), 17/33 LUTs. + func testCapabilities() async { + let driver = await makeDriver() + XCTAssertEqual(driver.capabilities.vendor, .red) + XCTAssertTrue(driver.capabilities.supportsNativeCDL) + XCTAssertTrue(driver.capabilities.lut3dSizes.contains(33)) + } + + // Connect performs RCP-style handshake, emits connected. + func testConnectHandshake() async throws { + let driver = await makeDriver() + var iterator = driver.state.makeAsyncIterator() + try await driver.connect() + var connected = false + for _ in 0..<5 { + if let s = await iterator.next(), s.connection == .connected { + connected = true + break + } + } + XCTAssertTrue(connected) + let handshakes = await sim.handshakeCount + XCTAssertEqual(handshakes, 1) + await driver.disconnect() + } + + func testConnectFailsOnDeadPort() async { + let driver = RedDriver(host: "127.0.0.1", port: 1, pollInterval: .milliseconds(20)) + do { + try await driver.connect() + XCTFail("expected throw") + } catch {} + } + + // Rec state + clip name via param polling. + func testRecStateEvents() async throws { + let driver = await makeDriver() + var iterator = driver.state.makeAsyncIterator() + try await driver.connect() + + await sim.setParam("RECORD_STATE", value: "1") + await sim.setParam("CLIP_NAME", value: "B001_C001_0710AB") + var sawRec = false + for _ in 0..<30 { + if let s = await iterator.next(), s.isRecording { + XCTAssertEqual(s.clipName, "B001_C001_0710AB") + sawRec = true + break + } + } + XCTAssertTrue(sawRec) + + await sim.setParam("RECORD_STATE", value: "0") + var sawStop = false + for _ in 0..<30 { + if let s = await iterator.next(), !s.isRecording { + sawStop = true + break + } + } + XCTAssertTrue(sawStop) + await driver.disconnect() + } + + // Metadata params surface. + func testMetadata() async throws { + await sim.setParam("ISO", value: "1600") + await sim.setParam("COLOR_TEMP", value: "4500") + await sim.setParam("TINT", value: "3") + await sim.setParam("SENSOR_FPS", value: "23.98") + await sim.setParam("TIMECODE", value: "15:30:00:00") + + let driver = await makeDriver() + var iterator = driver.state.makeAsyncIterator() + try await driver.connect() + + var got = false + for _ in 0..<30 { + if let s = await iterator.next(), let md = s.metadata, md.exposureIndex == 1600 { + XCTAssertEqual(md.whiteBalance, 4500) + XCTAssertEqual(md.tint, 3) + XCTAssertEqual(md.timecode, "15:30:00:00") + got = true + break + } + } + XCTAssertTrue(got) + await driver.disconnect() + } + + // Push: CDL to IPP2 slots + LUT upload. + func testPushLookWithNativeCDL() async throws { + let driver = await makeDriver() + try await driver.connect() + + let cdl = CDL(slope: SIMD3(1.1, 1.0, 0.9), offset: SIMD3(0.02, 0, 0), power: .one, saturation: 1.3) + let look = FlattenedLook(cdl: cdl, lut: Lut3D.identity(size: 33), latticeSize: 33) + try await driver.push(look: look) + + let cdlValue = await sim.params["IPP2_CDL"] + XCTAssertNotNil(cdlValue) + let obj = try JSONSerialization.jsonObject(with: Data(cdlValue!.utf8)) as? [String: Any] + let slope = obj?["slope"] as? [Double] + XCTAssertEqual(slope?[0] ?? 0, 1.1, accuracy: 1e-5) + + let luts = await sim.uploadedLuts + XCTAssertEqual(luts.count, 1) + XCTAssertEqual(luts[0].size, 33) + await driver.disconnect() + } + + // Unsupported LUT size rejected. + func testUnsupportedLutRejected() async throws { + let driver = await makeDriver() + try await driver.connect() + do { + try await driver.push(look: FlattenedLook(cdl: nil, lut: Lut3D.identity(size: 65), latticeSize: 65)) + XCTFail("expected unsupportedLook") + } catch let e as CameraError { + guard case .unsupportedLook = e else { return XCTFail("wrong error") } + } + await driver.disconnect() + } +}