2026-07-10 16:00:36 -04:00
|
|
|
import Foundation
|
|
|
|
|
import NIO
|
2026-07-11 00:33:26 -04:00
|
|
|
import NIOHTTP1
|
|
|
|
|
import NIOWebSocket
|
2026-07-10 16:00:36 -04:00
|
|
|
|
2026-07-11 00:33:26 -04:00
|
|
|
/// RCP2-faithful RED camera simulator.
|
|
|
|
|
/// Wire format verified against real Komodo clients (docs/research/red-rcp2-protocol.md):
|
|
|
|
|
/// - WebSocket server at /rcp (real camera: port 9998)
|
|
|
|
|
/// - Text frames, one JSON object each
|
|
|
|
|
/// - Handshake: client sends {"type":"rcp_config",...}, sim replies ack
|
|
|
|
|
/// - {"type":"rcp_get","id":P} -> {"id":P,"cur":{"val":...}}
|
|
|
|
|
/// - {"type":"rcp_set","id":P,"value":V} -> stores, replies current value
|
|
|
|
|
/// - pushNotification() sends unsolicited {"id":P,"cur":{"val":...}} to all clients
|
2026-07-10 16:00:36 -04:00
|
|
|
public actor RedSimulator {
|
2026-07-11 00:33:26 -04:00
|
|
|
|
|
|
|
|
/// Param values are ints, doubles, or strings on the wire.
|
|
|
|
|
public enum ParamValue: Sendable, Equatable {
|
|
|
|
|
case int(Int)
|
|
|
|
|
case double(Double)
|
|
|
|
|
case string(String)
|
|
|
|
|
|
|
|
|
|
var jsonValue: Any {
|
|
|
|
|
switch self {
|
|
|
|
|
case .int(let i): return i
|
|
|
|
|
case .double(let d): return d
|
|
|
|
|
case .string(let s): return s
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
static func from(_ any: Any) -> ParamValue? {
|
|
|
|
|
if let i = any as? Int { return .int(i) }
|
|
|
|
|
if let d = any as? Double { return .double(d) }
|
|
|
|
|
if let s = any as? String { return .string(s) }
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-10 16:00:36 -04:00
|
|
|
private var group: MultiThreadedEventLoopGroup?
|
|
|
|
|
private var channel: Channel?
|
|
|
|
|
public private(set) var boundPort: Int?
|
|
|
|
|
|
2026-07-11 00:33:26 -04:00
|
|
|
public private(set) var params: [String: ParamValue] = [
|
|
|
|
|
"RECORD_STATE": .int(0),
|
|
|
|
|
"CLIP_NAME": .string(""),
|
|
|
|
|
"ISO": .int(800),
|
|
|
|
|
"COLOR_TEMPERATURE": .int(5600),
|
|
|
|
|
"TINT": .int(0),
|
|
|
|
|
"SENSOR_FRAME_RATE": .string("24"),
|
|
|
|
|
"TIMECODE": .string("00:00:00:00"),
|
|
|
|
|
"CDL_ENABLE": .int(0),
|
|
|
|
|
"CAMERA_TYPE": .string("V-RAPTOR-SIM"),
|
2026-07-10 16:00:36 -04:00
|
|
|
]
|
2026-07-11 00:33:26 -04:00
|
|
|
/// Raw rcp_config messages received (handshake verification; decode in consumers).
|
|
|
|
|
public private(set) var receivedConfigsData: [Data] = []
|
|
|
|
|
|
|
|
|
|
private var clients: [ObjectIdentifier: Channel] = [:]
|
2026-07-10 16:00:36 -04:00
|
|
|
|
|
|
|
|
public init() {}
|
|
|
|
|
|
2026-07-11 00:33:26 -04:00
|
|
|
public func setParam(_ key: String, value: ParamValue) {
|
2026-07-10 16:00:36 -04:00
|
|
|
params[key] = value
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-11 00:33:26 -04:00
|
|
|
/// Send unsolicited notification to all connected clients (camera-initiated update).
|
|
|
|
|
public func pushNotification(param: String, value: ParamValue) {
|
|
|
|
|
params[param] = value
|
|
|
|
|
let payload: [String: Any] = ["id": param, "type": "rcp_cur_int", "cur": ["val": value.jsonValue]]
|
|
|
|
|
guard let data = try? JSONSerialization.data(withJSONObject: payload, options: [.sortedKeys]) else { return }
|
|
|
|
|
for ch in clients.values {
|
|
|
|
|
Self.sendText(data, over: ch)
|
2026-07-10 16:00:36 -04:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-11 00:33:26 -04:00
|
|
|
// MARK: Frame handling
|
|
|
|
|
|
|
|
|
|
func clientConnected(_ ch: Channel) {
|
|
|
|
|
clients[ObjectIdentifier(ch)] = ch
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func clientDisconnected(_ ch: Channel) {
|
|
|
|
|
clients.removeValue(forKey: ObjectIdentifier(ch))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func handleText(_ data: Data, from ch: Channel) {
|
|
|
|
|
guard let obj = try? JSONSerialization.jsonObject(with: data) as? [String: Any],
|
|
|
|
|
let type = obj["type"] as? String else { return }
|
|
|
|
|
|
|
|
|
|
switch type {
|
|
|
|
|
case "rcp_config":
|
|
|
|
|
receivedConfigsData.append(data)
|
|
|
|
|
reply(["type": "rcp_config_ack", "ok": true], to: ch)
|
|
|
|
|
case "rcp_get":
|
|
|
|
|
guard let id = obj["id"] as? String else { return }
|
|
|
|
|
if let value = params[id] {
|
|
|
|
|
reply(["id": id, "type": "rcp_cur_int", "cur": ["val": value.jsonValue]], to: ch)
|
|
|
|
|
} else {
|
|
|
|
|
reply(["id": id, "type": "rcp_error", "error": "unknown param"], to: ch)
|
2026-07-10 16:00:36 -04:00
|
|
|
}
|
2026-07-11 00:33:26 -04:00
|
|
|
case "rcp_set":
|
|
|
|
|
guard let id = obj["id"] as? String,
|
|
|
|
|
let raw = obj["value"], let value = ParamValue.from(raw) else { return }
|
|
|
|
|
params[id] = value
|
|
|
|
|
reply(["id": id, "type": "rcp_cur_int", "cur": ["val": value.jsonValue]], to: ch)
|
2026-07-10 16:00:36 -04:00
|
|
|
default:
|
2026-07-11 00:33:26 -04:00
|
|
|
reply(["type": "rcp_error", "error": "unsupported type \(type)"], to: ch)
|
2026-07-10 16:00:36 -04:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-11 00:33:26 -04:00
|
|
|
private func reply(_ obj: [String: Any], to ch: Channel) {
|
|
|
|
|
guard let data = try? JSONSerialization.data(withJSONObject: obj, options: [.sortedKeys]) else { return }
|
|
|
|
|
Self.sendText(data, over: ch)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private static func sendText(_ data: Data, over ch: Channel) {
|
|
|
|
|
var buf = ch.allocator.buffer(capacity: data.count)
|
|
|
|
|
buf.writeBytes(data)
|
|
|
|
|
let frame = WebSocketFrame(fin: true, opcode: .text, data: buf)
|
|
|
|
|
ch.writeAndFlush(NIOAny(frame), promise: nil)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// MARK: Lifecycle
|
|
|
|
|
|
2026-07-10 16:00:36 -04:00
|
|
|
public func start(port: Int) async throws {
|
|
|
|
|
let group = MultiThreadedEventLoopGroup(numberOfThreads: 1)
|
|
|
|
|
self.group = group
|
|
|
|
|
let sim = self
|
2026-07-11 00:33:26 -04:00
|
|
|
|
|
|
|
|
let upgrader = NIOWebSocketServerUpgrader(
|
|
|
|
|
shouldUpgrade: { channel, head in
|
|
|
|
|
// Real camera serves /rcp.
|
|
|
|
|
guard head.uri == "/rcp" else {
|
|
|
|
|
return channel.eventLoop.makeSucceededFuture(nil)
|
|
|
|
|
}
|
|
|
|
|
return channel.eventLoop.makeSucceededFuture(HTTPHeaders())
|
|
|
|
|
},
|
|
|
|
|
upgradePipelineHandler: { channel, _ in
|
|
|
|
|
Task { await sim.clientConnected(channel) }
|
|
|
|
|
return channel.pipeline.addHandler(RedWSHandler(sim: sim))
|
|
|
|
|
})
|
|
|
|
|
|
2026-07-10 16:00:36 -04:00
|
|
|
let bootstrap = ServerBootstrap(group: group)
|
|
|
|
|
.serverChannelOption(ChannelOptions.socketOption(.so_reuseaddr), value: 1)
|
|
|
|
|
.childChannelInitializer { channel in
|
2026-07-11 00:33:26 -04:00
|
|
|
let config = NIOHTTPServerUpgradeConfiguration(
|
|
|
|
|
upgraders: [upgrader],
|
|
|
|
|
completionHandler: { _ in })
|
|
|
|
|
return channel.pipeline.configureHTTPServerPipeline(withServerUpgrade: config)
|
2026-07-10 16:00:36 -04:00
|
|
|
}
|
|
|
|
|
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 {
|
2026-07-11 00:33:26 -04:00
|
|
|
for ch in clients.values {
|
|
|
|
|
_ = try? await ch.close()
|
|
|
|
|
}
|
|
|
|
|
clients.removeAll()
|
2026-07-10 16:00:36 -04:00
|
|
|
try await channel?.close()
|
|
|
|
|
try await group?.shutdownGracefully()
|
|
|
|
|
channel = nil
|
|
|
|
|
group = nil
|
|
|
|
|
boundPort = nil
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-11 00:33:26 -04:00
|
|
|
/// Inbound WebSocket frames -> sim actor.
|
|
|
|
|
final class RedWSHandler: ChannelInboundHandler, @unchecked Sendable {
|
|
|
|
|
typealias InboundIn = WebSocketFrame
|
|
|
|
|
typealias OutboundOut = WebSocketFrame
|
2026-07-10 16:00:36 -04:00
|
|
|
|
|
|
|
|
private let sim: RedSimulator
|
|
|
|
|
|
|
|
|
|
init(sim: RedSimulator) {
|
|
|
|
|
self.sim = sim
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func channelRead(context: ChannelHandlerContext, data: NIOAny) {
|
2026-07-11 00:33:26 -04:00
|
|
|
let frame = unwrapInboundIn(data)
|
|
|
|
|
switch frame.opcode {
|
|
|
|
|
case .text:
|
|
|
|
|
var buf = frame.unmaskedData
|
|
|
|
|
guard let bytes = buf.readBytes(length: buf.readableBytes) else { return }
|
2026-07-10 16:00:36 -04:00
|
|
|
let payload = Data(bytes)
|
2026-07-11 00:33:26 -04:00
|
|
|
let ch = context.channel
|
|
|
|
|
Task { await sim.handleText(payload, from: ch) }
|
|
|
|
|
case .connectionClose:
|
|
|
|
|
let ch = context.channel
|
|
|
|
|
Task { await sim.clientDisconnected(ch) }
|
|
|
|
|
context.close(promise: nil)
|
|
|
|
|
case .ping:
|
|
|
|
|
var buf = frame.unmaskedData
|
|
|
|
|
let pong = WebSocketFrame(fin: true, opcode: .pong, data: buf)
|
|
|
|
|
context.writeAndFlush(wrapOutboundOut(pong), promise: nil)
|
|
|
|
|
_ = buf
|
|
|
|
|
default:
|
|
|
|
|
break
|
2026-07-10 16:00:36 -04:00
|
|
|
}
|
|
|
|
|
}
|
2026-07-11 00:33:26 -04:00
|
|
|
|
|
|
|
|
func channelInactive(context: ChannelHandlerContext) {
|
|
|
|
|
let ch = context.channel
|
|
|
|
|
Task { await sim.clientDisconnected(ch) }
|
|
|
|
|
context.fireChannelInactive()
|
|
|
|
|
}
|
2026-07-10 16:00:36 -04:00
|
|
|
}
|