Initial support for scripting (#4016)
* Add ZeroMQ external submodule * ZeroMQ libzmq building on macOS * Added RPC namespace, settings and logging * Added request queue handling and new classes * Add C++ interface to ZeroMQ * Added start of ZeroMQ RPC Server implementation. * Request construction and callback request handling * Read and write memory implementation * Add ID to request format and send reply * Add RPC setting to macOS UI * Fixed initialization order bug and added exception handling * Working read-write through Python * Update CMakeLists for libzmq to resolve target name conflict on Windows * Platform-specific CMake definitions for Windows/non-Windows * Add comments * Revert "Add RPC setting to macOS UI" * Always run RPC server instead of configurable * Add Python scripting example. Updated .gitignore * Rename member variables to remove trailing underscore * Finally got libzmq external project building on macOS * Add missing dependency during libzmq build * Adding more missing dependencies [skip ci] * Only build what is required from libzmq * Extra length checks on client input * Call InvalidateCacheRange after memory write * Revert MinGW change. Fix clang-format. Improve error handling in request/reply. Allow any length of data read/write in Python. * Re-organized RPC static global state into a proper class. [skip ci] * Make sure libzmq always builds in Release mode * Renamed Request to Packet since Request and Reply are the same thing * Moved request fulfillment out of Packet and into RPCServer * Change request thread from sleep to condition variable * Remove non-blocking polling from ZMQ server code. Receive now blocks and terminates properly without sleeping. This change significantly improves script speed. * Move scripting files to dist/ instead of src/ * C++ code review changes for jroweboy [skip ci] * Python code review changes for jroweboy [skip ci] * Add docstrings and tests to citra.py [skip ci] * Add host OS check for libzmq build * Revert "Add host OS check for libzmq build" * Fixed a hang when emulation is stopped and restarted due to improper destruction order of ZMQ objects [skip ci] * Add scripting directory to archive packaging [skip ci] * Specify C/CXX compiler variables on MinGW build * Only specify compiler on Linux mingw * Use gcc and g++ on Windows mingw * Specify generator for mingw * Don't specify toolchain on windows mingw * Changed citra.py to support Python 3 instead of Python 2 * Fix bug where RPC wouldn't restart after Stop/Start emulation * Added copyright to headers and reorganized includes and forward declarations
This commit is contained in:
parent
abeee6859e
commit
04dd91be82
20 changed files with 670 additions and 1 deletions
|
@ -410,6 +410,14 @@ add_library(core STATIC
|
|||
movie.h
|
||||
perf_stats.cpp
|
||||
perf_stats.h
|
||||
rpc/packet.cpp
|
||||
rpc/packet.h
|
||||
rpc/rpc_server.cpp
|
||||
rpc/rpc_server.h
|
||||
rpc/server.cpp
|
||||
rpc/server.h
|
||||
rpc/zmq_server.cpp
|
||||
rpc/zmq_server.h
|
||||
settings.cpp
|
||||
settings.h
|
||||
telemetry_session.cpp
|
||||
|
@ -436,3 +444,5 @@ if (ARCHITECTURE_x86_64)
|
|||
)
|
||||
target_link_libraries(core PRIVATE dynarmic)
|
||||
endif()
|
||||
|
||||
target_link_libraries(core PUBLIC libzmq-headers cppzmq-headers libzmq)
|
||||
|
|
|
@ -25,6 +25,7 @@
|
|||
#include "core/loader/loader.h"
|
||||
#include "core/memory_setup.h"
|
||||
#include "core/movie.h"
|
||||
#include "core/rpc/rpc_server.h"
|
||||
#include "core/settings.h"
|
||||
#include "network/network.h"
|
||||
#include "video_core/video_core.h"
|
||||
|
@ -173,6 +174,7 @@ System::ResultStatus System::Init(EmuWindow& emu_window, u32 system_mode) {
|
|||
dsp_core->EnableStretching(Settings::values.enable_audio_stretching);
|
||||
|
||||
telemetry_session = std::make_unique<Core::TelemetrySession>();
|
||||
rpc_server = std::make_unique<RPC::RPCServer>();
|
||||
service_manager = std::make_shared<Service::SM::ServiceManager>();
|
||||
shared_page_handler = std::make_shared<SharedPage::Handler>();
|
||||
|
||||
|
@ -224,6 +226,7 @@ void System::Shutdown() {
|
|||
Kernel::Shutdown();
|
||||
HW::Shutdown();
|
||||
telemetry_session.reset();
|
||||
rpc_server.reset();
|
||||
service_manager.reset();
|
||||
dsp_core.reset();
|
||||
cpu_core.reset();
|
||||
|
|
|
@ -21,6 +21,10 @@ namespace AudioCore {
|
|||
class DspInterface;
|
||||
}
|
||||
|
||||
namespace RPC {
|
||||
class RPCServer;
|
||||
}
|
||||
|
||||
namespace Service {
|
||||
namespace SM {
|
||||
class ServiceManager;
|
||||
|
@ -202,6 +206,9 @@ private:
|
|||
/// Frontend applets
|
||||
std::shared_ptr<Frontend::SoftwareKeyboard> registered_swkbd;
|
||||
|
||||
/// RPC Server for scripting support
|
||||
std::unique_ptr<RPC::RPCServer> rpc_server;
|
||||
|
||||
/// Shared Page
|
||||
std::shared_ptr<SharedPage::Handler> shared_page_handler;
|
||||
|
||||
|
|
15
src/core/rpc/packet.cpp
Normal file
15
src/core/rpc/packet.cpp
Normal file
|
@ -0,0 +1,15 @@
|
|||
#include <algorithm>
|
||||
#include <cstring>
|
||||
|
||||
#include "core/rpc/packet.h"
|
||||
|
||||
namespace RPC {
|
||||
|
||||
Packet::Packet(const PacketHeader& header, u8* data,
|
||||
std::function<void(Packet&)> send_reply_callback)
|
||||
: header(header), send_reply_callback(std::move(send_reply_callback)) {
|
||||
|
||||
std::memcpy(packet_data.data(), data, std::min(header.packet_size, MAX_PACKET_DATA_SIZE));
|
||||
}
|
||||
|
||||
}; // namespace RPC
|
78
src/core/rpc/packet.h
Normal file
78
src/core/rpc/packet.h
Normal file
|
@ -0,0 +1,78 @@
|
|||
// Copyright 2018 Citra Emulator Project
|
||||
// Licensed under GPLv2 or any later version
|
||||
// Refer to the license.txt file included.
|
||||
|
||||
#pragma once
|
||||
|
||||
#include <array>
|
||||
#include <functional>
|
||||
#include "common/common_types.h"
|
||||
|
||||
namespace RPC {
|
||||
|
||||
enum class PacketType {
|
||||
Undefined = 0,
|
||||
ReadMemory,
|
||||
WriteMemory,
|
||||
};
|
||||
|
||||
struct PacketHeader {
|
||||
u32 version;
|
||||
u32 id;
|
||||
PacketType packet_type;
|
||||
u32 packet_size;
|
||||
};
|
||||
|
||||
constexpr u32 CURRENT_VERSION = 1;
|
||||
constexpr u32 MIN_PACKET_SIZE = sizeof(PacketHeader);
|
||||
constexpr u32 MAX_PACKET_DATA_SIZE = 32;
|
||||
constexpr u32 MAX_PACKET_SIZE = MIN_PACKET_SIZE + MAX_PACKET_DATA_SIZE;
|
||||
constexpr u32 MAX_READ_SIZE = MAX_PACKET_DATA_SIZE;
|
||||
|
||||
class Packet {
|
||||
public:
|
||||
Packet(const PacketHeader& header, u8* data, std::function<void(Packet&)> send_reply_callback);
|
||||
|
||||
u32 GetVersion() const {
|
||||
return header.version;
|
||||
}
|
||||
|
||||
u32 GetId() const {
|
||||
return header.id;
|
||||
}
|
||||
|
||||
PacketType GetPacketType() const {
|
||||
return header.packet_type;
|
||||
}
|
||||
|
||||
u32 GetPacketDataSize() const {
|
||||
return header.packet_size;
|
||||
}
|
||||
|
||||
const PacketHeader& GetHeader() const {
|
||||
return header;
|
||||
}
|
||||
|
||||
std::array<u8, MAX_PACKET_DATA_SIZE>& GetPacketData() {
|
||||
return packet_data;
|
||||
}
|
||||
|
||||
void SetPacketDataSize(u32 size) {
|
||||
header.packet_size = size;
|
||||
}
|
||||
|
||||
void SendReply() {
|
||||
send_reply_callback(*this);
|
||||
}
|
||||
|
||||
private:
|
||||
void HandleReadMemory(u32 address, u32 data_size);
|
||||
void HandleWriteMemory(u32 address, const u8* data, u32 data_size);
|
||||
|
||||
struct PacketHeader header;
|
||||
std::array<u8, MAX_PACKET_DATA_SIZE> packet_data;
|
||||
|
||||
std::function<void(Packet&)> send_reply_callback;
|
||||
};
|
||||
|
||||
} // namespace RPC
|
139
src/core/rpc/rpc_server.cpp
Normal file
139
src/core/rpc/rpc_server.cpp
Normal file
|
@ -0,0 +1,139 @@
|
|||
#include "common/logging/log.h"
|
||||
#include "core/arm/arm_interface.h"
|
||||
#include "core/core.h"
|
||||
#include "core/memory.h"
|
||||
#include "core/rpc/packet.h"
|
||||
#include "core/rpc/rpc_server.h"
|
||||
|
||||
namespace RPC {
|
||||
|
||||
RPCServer::RPCServer() : server(*this) {
|
||||
LOG_INFO(RPC_Server, "Starting RPC server ...");
|
||||
|
||||
Start();
|
||||
|
||||
LOG_INFO(RPC_Server, "RPC started.");
|
||||
}
|
||||
|
||||
RPCServer::~RPCServer() {
|
||||
LOG_INFO(RPC_Server, "Stopping RPC ...");
|
||||
|
||||
Stop();
|
||||
|
||||
LOG_INFO(RPC_Server, "RPC stopped.");
|
||||
}
|
||||
|
||||
void RPCServer::HandleReadMemory(Packet& packet, u32 address, u32 data_size) {
|
||||
if (data_size > MAX_READ_SIZE) {
|
||||
return;
|
||||
}
|
||||
|
||||
// Note: Memory read occurs asynchronously from the state of the emulator
|
||||
Memory::ReadBlock(address, packet.GetPacketData().data(), data_size);
|
||||
packet.SetPacketDataSize(data_size);
|
||||
packet.SendReply();
|
||||
}
|
||||
|
||||
void RPCServer::HandleWriteMemory(Packet& packet, u32 address, const u8* data, u32 data_size) {
|
||||
// Only allow writing to certain memory regions
|
||||
if ((address >= Memory::PROCESS_IMAGE_VADDR && address <= Memory::PROCESS_IMAGE_VADDR_END) ||
|
||||
(address >= Memory::HEAP_VADDR && address <= Memory::HEAP_VADDR_END) ||
|
||||
(address >= Memory::N3DS_EXTRA_RAM_VADDR && address <= Memory::N3DS_EXTRA_RAM_VADDR_END)) {
|
||||
// Note: Memory write occurs asynchronously from the state of the emulator
|
||||
Memory::WriteBlock(address, data, data_size);
|
||||
// If the memory happens to be executable code, make sure the changes become visible
|
||||
Core::CPU().InvalidateCacheRange(address, data_size);
|
||||
}
|
||||
packet.SetPacketDataSize(0);
|
||||
packet.SendReply();
|
||||
}
|
||||
|
||||
bool RPCServer::ValidatePacket(const PacketHeader& packet_header) {
|
||||
if (packet_header.version <= CURRENT_VERSION) {
|
||||
switch (packet_header.packet_type) {
|
||||
case PacketType::ReadMemory:
|
||||
case PacketType::WriteMemory:
|
||||
if (packet_header.packet_size >= (sizeof(u32) * 2)) {
|
||||
return true;
|
||||
}
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
void RPCServer::HandleSingleRequest(std::unique_ptr<Packet> request_packet) {
|
||||
bool success = false;
|
||||
|
||||
if (ValidatePacket(request_packet->GetHeader())) {
|
||||
// Currently, all request types use the address/data_size wire format
|
||||
u32 address = 0;
|
||||
u32 data_size = 0;
|
||||
std::memcpy(&address, request_packet->GetPacketData().data(), sizeof(address));
|
||||
std::memcpy(&data_size, request_packet->GetPacketData().data() + sizeof(address),
|
||||
sizeof(data_size));
|
||||
|
||||
switch (request_packet->GetPacketType()) {
|
||||
case PacketType::ReadMemory:
|
||||
if (data_size > 0 && data_size <= MAX_READ_SIZE) {
|
||||
HandleReadMemory(*request_packet, address, data_size);
|
||||
success = true;
|
||||
}
|
||||
break;
|
||||
case PacketType::WriteMemory:
|
||||
if (data_size > 0 && data_size <= MAX_PACKET_DATA_SIZE - (sizeof(u32) * 2)) {
|
||||
const u8* data = request_packet->GetPacketData().data() + (sizeof(u32) * 2);
|
||||
HandleWriteMemory(*request_packet, address, data, data_size);
|
||||
success = true;
|
||||
}
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if (!success) {
|
||||
// Send an empty reply, so as not to hang the client
|
||||
request_packet->SetPacketDataSize(0);
|
||||
request_packet->SendReply();
|
||||
}
|
||||
}
|
||||
|
||||
void RPCServer::HandleRequestsLoop() {
|
||||
std::unique_ptr<RPC::Packet> request_packet;
|
||||
|
||||
LOG_INFO(RPC_Server, "Request handler started.");
|
||||
|
||||
while (true) {
|
||||
std::unique_lock<std::mutex> lock(request_queue_mutex);
|
||||
request_queue_cv.wait(lock, [&] { return !running || request_queue.Pop(request_packet); });
|
||||
if (!running) {
|
||||
break;
|
||||
}
|
||||
HandleSingleRequest(std::move(request_packet));
|
||||
}
|
||||
}
|
||||
|
||||
void RPCServer::QueueRequest(std::unique_ptr<RPC::Packet> request) {
|
||||
std::unique_lock<std::mutex> lock(request_queue_mutex);
|
||||
request_queue.Push(std::move(request));
|
||||
request_queue_cv.notify_one();
|
||||
}
|
||||
|
||||
void RPCServer::Start() {
|
||||
running = true;
|
||||
const auto threadFunction = [this]() { HandleRequestsLoop(); };
|
||||
request_handler_thread = std::thread(threadFunction);
|
||||
server.Start();
|
||||
}
|
||||
|
||||
void RPCServer::Stop() {
|
||||
running = false;
|
||||
request_queue_cv.notify_one();
|
||||
request_handler_thread.join();
|
||||
server.Stop();
|
||||
}
|
||||
|
||||
}; // namespace RPC
|
40
src/core/rpc/rpc_server.h
Normal file
40
src/core/rpc/rpc_server.h
Normal file
|
@ -0,0 +1,40 @@
|
|||
// Copyright 2018 Citra Emulator Project
|
||||
// Licensed under GPLv2 or any later version
|
||||
// Refer to the license.txt file included.
|
||||
|
||||
#pragma once
|
||||
|
||||
#include <condition_variable>
|
||||
#include <memory>
|
||||
#include <mutex>
|
||||
#include <thread>
|
||||
#include "common/threadsafe_queue.h"
|
||||
#include "core/rpc/server.h"
|
||||
|
||||
namespace RPC {
|
||||
|
||||
class RPCServer {
|
||||
public:
|
||||
RPCServer();
|
||||
~RPCServer();
|
||||
|
||||
void QueueRequest(std::unique_ptr<RPC::Packet> request);
|
||||
|
||||
private:
|
||||
void Start();
|
||||
void Stop();
|
||||
void HandleReadMemory(Packet& packet, u32 address, u32 data_size);
|
||||
void HandleWriteMemory(Packet& packet, u32 address, const u8* data, u32 data_size);
|
||||
bool ValidatePacket(const PacketHeader& packet_header);
|
||||
void HandleSingleRequest(std::unique_ptr<Packet> request);
|
||||
void HandleRequestsLoop();
|
||||
|
||||
Server server;
|
||||
Common::SPSCQueue<std::unique_ptr<Packet>> request_queue;
|
||||
bool running = false;
|
||||
std::thread request_handler_thread;
|
||||
std::mutex request_queue_mutex;
|
||||
std::condition_variable request_queue_cv;
|
||||
};
|
||||
|
||||
} // namespace RPC
|
35
src/core/rpc/server.cpp
Normal file
35
src/core/rpc/server.cpp
Normal file
|
@ -0,0 +1,35 @@
|
|||
#include <functional>
|
||||
|
||||
#include "common/threadsafe_queue.h"
|
||||
#include "core/core.h"
|
||||
#include "core/rpc/rpc_server.h"
|
||||
#include "core/rpc/server.h"
|
||||
|
||||
namespace RPC {
|
||||
|
||||
Server::Server(RPCServer& rpc_server) : rpc_server(rpc_server) {}
|
||||
|
||||
void Server::Start() {
|
||||
const auto callback = [this](std::unique_ptr<RPC::Packet> new_request) {
|
||||
NewRequestCallback(std::move(new_request));
|
||||
};
|
||||
|
||||
try {
|
||||
zmq_server = std::make_unique<ZMQServer>(callback);
|
||||
} catch (...) {
|
||||
LOG_ERROR(RPC_Server, "Error starting ZeroMQ server");
|
||||
}
|
||||
}
|
||||
|
||||
void Server::Stop() {
|
||||
zmq_server.reset();
|
||||
}
|
||||
|
||||
void Server::NewRequestCallback(std::unique_ptr<RPC::Packet> new_request) {
|
||||
LOG_INFO(RPC_Server, "Received request version={} id={} type={} size={}",
|
||||
new_request->GetVersion(), new_request->GetId(),
|
||||
static_cast<u32>(new_request->GetPacketType()), new_request->GetPacketDataSize());
|
||||
rpc_server.QueueRequest(std::move(new_request));
|
||||
}
|
||||
|
||||
}; // namespace RPC
|
27
src/core/rpc/server.h
Normal file
27
src/core/rpc/server.h
Normal file
|
@ -0,0 +1,27 @@
|
|||
// Copyright 2018 Citra Emulator Project
|
||||
// Licensed under GPLv2 or any later version
|
||||
// Refer to the license.txt file included.
|
||||
|
||||
#pragma once
|
||||
|
||||
#include "core/rpc/packet.h"
|
||||
#include "core/rpc/zmq_server.h"
|
||||
|
||||
namespace RPC {
|
||||
|
||||
class RPCServer;
|
||||
class ZMQServer;
|
||||
|
||||
class Server {
|
||||
public:
|
||||
Server(RPCServer& rpc_server);
|
||||
void Start();
|
||||
void Stop();
|
||||
void NewRequestCallback(std::unique_ptr<RPC::Packet> new_request);
|
||||
|
||||
private:
|
||||
RPCServer& rpc_server;
|
||||
std::unique_ptr<ZMQServer> zmq_server;
|
||||
};
|
||||
|
||||
} // namespace RPC
|
78
src/core/rpc/zmq_server.cpp
Normal file
78
src/core/rpc/zmq_server.cpp
Normal file
|
@ -0,0 +1,78 @@
|
|||
#include "common/common_types.h"
|
||||
#include "core/core.h"
|
||||
#include "core/rpc/packet.h"
|
||||
#include "core/rpc/zmq_server.h"
|
||||
|
||||
namespace RPC {
|
||||
|
||||
ZMQServer::ZMQServer(std::function<void(std::unique_ptr<Packet>)> new_request_callback)
|
||||
: zmq_context(std::move(std::make_unique<zmq::context_t>(1))),
|
||||
zmq_socket(std::move(std::make_unique<zmq::socket_t>(*zmq_context, ZMQ_REP))),
|
||||
new_request_callback(std::move(new_request_callback)) {
|
||||
// Use a random high port
|
||||
// TODO: Make configurable or increment port number on failure
|
||||
zmq_socket->bind("tcp://127.0.0.1:45987");
|
||||
LOG_INFO(RPC_Server, "ZeroMQ listening on port 45987");
|
||||
|
||||
worker_thread = std::thread(&ZMQServer::WorkerLoop, this);
|
||||
}
|
||||
|
||||
ZMQServer::~ZMQServer() {
|
||||
// Triggering the zmq_context destructor will cancel
|
||||
// any blocking calls to zmq_socket->recv()
|
||||
running = false;
|
||||
zmq_context.reset();
|
||||
worker_thread.join();
|
||||
|
||||
LOG_INFO(RPC_Server, "ZeroMQ stopped");
|
||||
}
|
||||
|
||||
void ZMQServer::WorkerLoop() {
|
||||
zmq::message_t request;
|
||||
while (running) {
|
||||
try {
|
||||
if (zmq_socket->recv(&request, 0)) {
|
||||
if (request.size() >= MIN_PACKET_SIZE && request.size() <= MAX_PACKET_SIZE) {
|
||||
u8* request_buffer = static_cast<u8*>(request.data());
|
||||
PacketHeader header;
|
||||
std::memcpy(&header, request_buffer, sizeof(header));
|
||||
if ((request.size() - MIN_PACKET_SIZE) == header.packet_size) {
|
||||
u8* data = request_buffer + MIN_PACKET_SIZE;
|
||||
std::function<void(Packet&)> send_reply_callback =
|
||||
std::bind(&ZMQServer::SendReply, this, std::placeholders::_1);
|
||||
std::unique_ptr<Packet> new_packet =
|
||||
std::make_unique<Packet>(header, data, send_reply_callback);
|
||||
|
||||
// Send the request to the upper layer for handling
|
||||
new_request_callback(std::move(new_packet));
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (...) {
|
||||
LOG_WARNING(RPC_Server, "Failed to receive data on ZeroMQ socket");
|
||||
}
|
||||
}
|
||||
|
||||
// Destroying the socket must be done by this thread.
|
||||
zmq_socket.reset();
|
||||
}
|
||||
|
||||
void ZMQServer::SendReply(Packet& reply_packet) {
|
||||
if (running) {
|
||||
auto reply_buffer =
|
||||
std::make_unique<u8[]>(MIN_PACKET_SIZE + reply_packet.GetPacketDataSize());
|
||||
auto reply_header = reply_packet.GetHeader();
|
||||
|
||||
std::memcpy(reply_buffer.get(), &reply_header, sizeof(reply_header));
|
||||
std::memcpy(reply_buffer.get() + (4 * sizeof(u32)), reply_packet.GetPacketData().data(),
|
||||
reply_packet.GetPacketDataSize());
|
||||
|
||||
zmq_socket->send(reply_buffer.get(), MIN_PACKET_SIZE + reply_packet.GetPacketDataSize());
|
||||
|
||||
LOG_INFO(RPC_Server, "Sent reply version({}) id=({}) type=({}) size=({})",
|
||||
reply_packet.GetVersion(), reply_packet.GetId(),
|
||||
static_cast<u32>(reply_packet.GetPacketType()), reply_packet.GetPacketDataSize());
|
||||
}
|
||||
}
|
||||
|
||||
}; // namespace RPC
|
34
src/core/rpc/zmq_server.h
Normal file
34
src/core/rpc/zmq_server.h
Normal file
|
@ -0,0 +1,34 @@
|
|||
// Copyright 2018 Citra Emulator Project
|
||||
// Licensed under GPLv2 or any later version
|
||||
// Refer to the license.txt file included.
|
||||
|
||||
#pragma once
|
||||
|
||||
#include <functional>
|
||||
#include <thread>
|
||||
#define ZMQ_STATIC
|
||||
#include <zmq.hpp>
|
||||
|
||||
namespace RPC {
|
||||
|
||||
class Packet;
|
||||
|
||||
class ZMQServer {
|
||||
public:
|
||||
explicit ZMQServer(std::function<void(std::unique_ptr<Packet>)> new_request_callback);
|
||||
~ZMQServer();
|
||||
|
||||
private:
|
||||
void WorkerLoop();
|
||||
void SendReply(Packet& request);
|
||||
|
||||
std::thread worker_thread;
|
||||
std::atomic_bool running = true;
|
||||
|
||||
std::unique_ptr<zmq::context_t> zmq_context;
|
||||
std::unique_ptr<zmq::socket_t> zmq_socket;
|
||||
|
||||
std::function<void(std::unique_ptr<Packet>)> new_request_callback;
|
||||
};
|
||||
|
||||
} // namespace RPC
|
Loading…
Add table
Add a link
Reference in a new issue