136 lines
5.1 KiB
Swift
136 lines
5.1 KiB
Swift
|
|
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)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|