From 1c55d821913abed3dbf4b5d2f10c8d6953db443c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Komor?= Date: Thu, 17 Sep 2026 10:30:58 +0200 Subject: [PATCH 1/2] Add sqs plugin --- .github/workflows/release.yml | 21 +- .gitignore | 1 + README.md | 1 + sqs/CMakeLists.txt | 21 ++ sqs/GraftcodePluginsInterfaces/IServer.h | 22 ++ sqs/GraftcodePluginsInterfaces/ITransport.h | 18 ++ sqs/GraftcodePluginsInterfaces/_common.h | 5 + sqs/Readme.md | 155 +++++++++++ sqs/SqsPlugin/CMakeLists.txt | 21 ++ sqs/SqsPlugin/SqsClient.cpp | 152 +++++++++++ sqs/SqsPlugin/SqsClient.h | 35 +++ sqs/SqsPlugin/SqsConfig.cpp | 220 ++++++++++++++++ sqs/SqsPlugin/SqsConfig.h | 30 +++ sqs/SqsPlugin/SqsServer.cpp | 275 ++++++++++++++++++++ sqs/SqsPlugin/SqsServer.h | 38 +++ sqs/SqsPlugin/SqsUtil.cpp | 183 +++++++++++++ sqs/SqsPlugin/SqsUtil.h | 36 +++ sqs/SqsPlugin/TransportSqs.cpp | 96 +++++++ sqs/SqsPlugin/TransportSqs.h | 37 +++ sqs/SqsPluginTest/CMakeLists.txt | 35 +++ sqs/SqsPluginTest/SqsTransportTests.cpp | 210 +++++++++++++++ sqs/docker-compose.yml | 14 + sqs/vcpkg.json | 16 ++ 23 files changed, 1641 insertions(+), 1 deletion(-) create mode 100644 sqs/CMakeLists.txt create mode 100644 sqs/GraftcodePluginsInterfaces/IServer.h create mode 100644 sqs/GraftcodePluginsInterfaces/ITransport.h create mode 100644 sqs/GraftcodePluginsInterfaces/_common.h create mode 100644 sqs/Readme.md create mode 100644 sqs/SqsPlugin/CMakeLists.txt create mode 100644 sqs/SqsPlugin/SqsClient.cpp create mode 100644 sqs/SqsPlugin/SqsClient.h create mode 100644 sqs/SqsPlugin/SqsConfig.cpp create mode 100644 sqs/SqsPlugin/SqsConfig.h create mode 100644 sqs/SqsPlugin/SqsServer.cpp create mode 100644 sqs/SqsPlugin/SqsServer.h create mode 100644 sqs/SqsPlugin/SqsUtil.cpp create mode 100644 sqs/SqsPlugin/SqsUtil.h create mode 100644 sqs/SqsPlugin/TransportSqs.cpp create mode 100644 sqs/SqsPlugin/TransportSqs.h create mode 100644 sqs/SqsPluginTest/CMakeLists.txt create mode 100644 sqs/SqsPluginTest/SqsTransportTests.cpp create mode 100644 sqs/docker-compose.yml create mode 100644 sqs/vcpkg.json diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 7aefe1e..7a212a3 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -9,6 +9,7 @@ on: - "rabbitmq/**" - "servicebus/**" - "kafka/**" + - "sqs/**" - ".github/workflows/release.yml" pull_request: branches: [ "main" ] @@ -16,6 +17,7 @@ on: - "rabbitmq/**" - "servicebus/**" - "kafka/**" + - "sqs/**" - ".github/workflows/release.yml" workflow_dispatch: inputs: @@ -36,7 +38,7 @@ jobs: strategy: fail-fast: false matrix: - plugin: [rabbitmq, servicebus, kafka] + plugin: [rabbitmq, servicebus, kafka, sqs] platform: [windows-latest, windows-11-arm, ubuntu-22.04, ubuntu-22.04-arm, macos-26-intel, macos-14] steps: @@ -91,6 +93,18 @@ jobs: echo "INSTALL_PACKAGES=cmake" >> "$GITHUB_ENV" fi ;; + sqs) + echo "DIRECTORY=sqs" >> "$GITHUB_ENV" + echo "ARTIFACT_PATH=sqs/build/SqsPlugin" >> "$GITHUB_ENV" + echo "CONFIGURE_ARGS=-DCMAKE_TOOLCHAIN_FILE=${VCPKG_INSTALLATION_ROOT}/scripts/buildsystems/vcpkg.cmake" >> "$GITHUB_ENV" + if [ "${{ runner.os }}" = "Linux" ]; then + echo "INSTALL_PACKAGES=build-essential cmake git ninja-build pkg-config" >> "$GITHUB_ENV" + elif [ "${{ runner.os }}" = "macOS" ]; then + echo "INSTALL_PACKAGES=cmake ninja pkg-config" >> "$GITHUB_ENV" + else + echo "INSTALL_PACKAGES=cmake" >> "$GITHUB_ENV" + fi + ;; kafka) echo "DIRECTORY=kafka" >> "$GITHUB_ENV" echo "ARTIFACT_PATH=kafka/build/KafkaPlugin" >> "$GITHUB_ENV" @@ -140,6 +154,11 @@ jobs: run: | cmake --build build --config Release + - name: Test (${{ matrix.plugin }}) + working-directory: ${{ env.DIRECTORY }} + run: | + ctest --test-dir build -C Release --output-on-failure + - name: Package (${{ matrix.plugin }}) run: | mkdir -p release-payloads diff --git a/.gitignore b/.gitignore index a8ffd37..5cffd1c 100644 --- a/.gitignore +++ b/.gitignore @@ -7,3 +7,4 @@ build/ bin/ obj/ Binaries/ +vcpkg_installed/ diff --git a/README.md b/README.md index b903cea..7c2057f 100644 --- a/README.md +++ b/README.md @@ -9,6 +9,7 @@ Official open-source Graftcode Gateway plugins for carrying Graft calls over ext | [rabbitmq](rabbitmq/) | RabbitMQ (AMQP 0-9-1), request/reply | | [servicebus](servicebus/) | Azure Service Bus (AMQP 1.0), request/reply and one-way | | [kafka](kafka/) | Apache Kafka, request/reply (correlation-id) | +| [sqs](sqs/) | Amazon SQS, request/reply and one-way | | [observability/opentelemetry](observability/opentelemetry/) | OpenTelemetry / Azure Application Insights connector | Each plugin has its own README with build and configuration steps. For how the Gateway loads a plugin, see the "Plugin server config" section of the [Graftcode Gateway](https://github.com/grft-dev/graftcode-gateway) README. diff --git a/sqs/CMakeLists.txt b/sqs/CMakeLists.txt new file mode 100644 index 0000000..a6fbf37 --- /dev/null +++ b/sqs/CMakeLists.txt @@ -0,0 +1,21 @@ +set(CMAKE_MIN 3.22) +cmake_minimum_required(VERSION ${CMAKE_MIN}) +set(CMAKE_POLICY_VERSION_MINIMUM ${CMAKE_MIN}) +cmake_policy(VERSION ${CMAKE_MIN}) + +project("SqsPlugin" VERSION 1.0.0) + +set(CMAKE_CXX_STANDARD 20) +set(CMAKE_CXX_STANDARD_REQUIRED ON) +set(CMAKE_POSITION_INDEPENDENT_CODE ON) +set(CMAKE_POLICY_DEFAULT_CMP0135 NEW) + +add_definitions(-DUNICODE) +enable_testing() + +if(MSVC) + add_compile_options(/utf-8) +endif() + +add_subdirectory("SqsPlugin") +add_subdirectory("SqsPluginTest") diff --git a/sqs/GraftcodePluginsInterfaces/IServer.h b/sqs/GraftcodePluginsInterfaces/IServer.h new file mode 100644 index 0000000..5d3e5ee --- /dev/null +++ b/sqs/GraftcodePluginsInterfaces/IServer.h @@ -0,0 +1,22 @@ +#pragma once + +#include + +namespace GraftcodeGateway +{ + class IServer + { + public: + using byte = unsigned char; + using WriteResponseFn = + void(*)(void* context, const byte* data, std::size_t size); + using ProcessMessageFn = + bool(*)(const byte* requestData, std::size_t requestSize, + WriteResponseFn writeResponse, void* writeContext); + + virtual ~IServer() = default; + virtual void configure(const char* jsonConfig, ProcessMessageFn processMessage) = 0; + virtual void start() = 0; + virtual void stop() = 0; + }; +} diff --git a/sqs/GraftcodePluginsInterfaces/ITransport.h b/sqs/GraftcodePluginsInterfaces/ITransport.h new file mode 100644 index 0000000..335768c --- /dev/null +++ b/sqs/GraftcodePluginsInterfaces/ITransport.h @@ -0,0 +1,18 @@ +#pragma once + +#include "_common.h" + +namespace Hypertube::Native::Interfaces +{ + class ITransport + { + public: + virtual ~ITransport() = default; + virtual int Initialize( + byte callingRuntimeNumber, + byte calledRuntimeNumber, + byte calledRuntimeVersion) = 0; + virtual int SendCommand(byte* messageByteArray, int32_t messageByteArrayLen) = 0; + virtual int ReadResponse(byte* responseByteArray, int32_t responseByteArrayLen) = 0; + }; +} diff --git a/sqs/GraftcodePluginsInterfaces/_common.h b/sqs/GraftcodePluginsInterfaces/_common.h new file mode 100644 index 0000000..483310f --- /dev/null +++ b/sqs/GraftcodePluginsInterfaces/_common.h @@ -0,0 +1,5 @@ +#pragma once + +#include + +typedef unsigned char byte; diff --git a/sqs/Readme.md b/sqs/Readme.md new file mode 100644 index 0000000..8a41ed5 --- /dev/null +++ b/sqs/Readme.md @@ -0,0 +1,155 @@ +# Amazon SQS Plugin + +This native C++ plugin carries Graftcode Gateway calls over Amazon Simple Queue +Service (SQS). It implements the standard Graftcode transport and server +interfaces and exports `CreateTransportChannel` / `DestroyTransportChannel` and +`CreateServer` / `DestroyServer`. + +Binary Graftcode payloads are Base64-encoded into the SQS message body. RPC +metadata is carried in message attributes: + +- `GraftcodeCorrelationId` identifies the request and response. +- `GraftcodeReplyTo` contains the caller's reply queue URL. + +## Delivery model + +SQS provides at-least-once delivery. The server deletes a request only after +`processMessage` succeeds and, for RPC, the response has been sent. Configure a +dead-letter queue and make called operations idempotent. + +SQS cannot filter receives by message attribute. Consequently, an RPC reply +queue must be dedicated to one client/runtime plugin instance. Do not share one +reply queue between independent processes. Calls made through one transport +instance are serialized. + +Standard and FIFO queues are supported. For a FIFO queue URL ending in +`.fifo`, the plugin automatically sets `MessageGroupId` and +`MessageDeduplicationId`. + +## Build + +The plugin uses the AWS SDK for C++ (`sqs`) and `nlohmann-json`, installed in +vcpkg manifest mode. + +```powershell +git clone https://github.com/microsoft/vcpkg.git +.\vcpkg\bootstrap-vcpkg.bat + +cmake -S . -B build -DCMAKE_BUILD_TYPE=Release ` + -DCMAKE_TOOLCHAIN_FILE=.\vcpkg\scripts\buildsystems\vcpkg.cmake +cmake --build build --config Release +ctest --test-dir build -C Release --output-on-failure +``` + +The output is `build/SqsPlugin/SqsPlugin.dll` on Windows or +`build/SqsPlugin/libSqsPlugin.so` / `.dylib` on Linux/macOS. Use +`libSqsPlugin` as the configured plugin name when the library has the `lib` +prefix. + +## Create queues + +Create one request queue and one reply queue per client runtime: + +```bash +aws sqs create-queue --queue-name graft-requests +aws sqs create-queue --queue-name graft-replies-client-1 +``` + +Use the returned queue URLs in configuration. The gateway identity needs +`sqs:ReceiveMessage`, `sqs:DeleteMessage`, and `sqs:SendMessage`; a client needs +`sqs:SendMessage`, `sqs:ReceiveMessage`, and `sqs:DeleteMessage` for its queues. + +## Gateway configuration + +```json +{ + "name": "SqsPlugin", + "region": "eu-central-1", + "requestQueueUrl": "https://sqs.eu-central-1.amazonaws.com/123456789012/graft-requests", + "waitTimeSeconds": 10, + "visibilityTimeoutSeconds": 60 +} +``` + +Run the Gateway: + +```powershell +./gg .\PhysicsCalculator.dll --config .\pluginConfig.json +``` + +## Runtime configuration + +```csharp +string configSource = +""" +{ + "configurations": { + "graft.nuget.PhysicsCalculator": { + "runtime": "netcore", + "stateless": true, + "plugin": { + "name": "SqsPlugin", + "region": "eu-central-1", + "requestQueueUrl": "https://sqs.eu-central-1.amazonaws.com/123456789012/graft-requests", + "replyQueueUrl": "https://sqs.eu-central-1.amazonaws.com/123456789012/graft-replies-client-1", + "rpcTimeoutMs": 30000, + "waitTimeSeconds": 10 + } + } + } +} +"""; + +graft.nuget.PhysicsCalculator.GraftConfig.SetConfig(configSource); +``` + +The AWS SDK default credential provider chain is used. Prefer IAM roles, +environment variables, shared AWS config, or workload identity. Static +credentials are accepted as `accessKeyId`, `secretAccessKey`, and optional +`sessionToken`, primarily for local emulators. + +## One-way calls + +Set `"oneWay": true` in both runtime and gateway configuration. The client sends +without `GraftcodeReplyTo`, and the gateway processes and deletes the request +without sending a response. + +## LocalStack + +Start the included LocalStack setup: + +```bash +docker compose up -d +aws --endpoint-url http://localhost:4566 sqs create-queue --queue-name graft-requests +aws --endpoint-url http://localhost:4566 sqs create-queue --queue-name graft-replies +``` + +Use this additional configuration: + +```json +{ + "region": "us-east-1", + "endpointOverride": "http://localhost:4566", + "accessKeyId": "test", + "secretAccessKey": "test", + "verifySsl": false +} +``` + +## Configuration reference + +- `region` — AWS region; defaults to `us-east-1`. +- `requestQueueUrl` — required request queue URL (`queueUrl` and `queue` aliases + are accepted). +- `replyQueueUrl` — required by an RPC client; optional on the gateway because + each request supplies its reply destination. +- `rpcTimeoutMs` — RPC response timeout; defaults to `30000`. +- `waitTimeSeconds` — SQS long-poll duration, 1–20; defaults to `1`. +- `visibilityTimeoutSeconds` — per-receive visibility timeout, 1–43200; + defaults to `30`. Set it longer than the maximum call execution time. +- `oneWay` — process without a response; defaults to `false`. +- `messageGroupId` — FIFO message group; defaults to `graftcode`. +- `endpointOverride` — custom SQS endpoint, such as LocalStack. +- `verifySsl` — TLS certificate verification; defaults to `true`. +- `accessKeyId`, `secretAccessKey`, `sessionToken` — optional explicit + credentials; otherwise the AWS SDK provider chain is used. diff --git a/sqs/SqsPlugin/CMakeLists.txt b/sqs/SqsPlugin/CMakeLists.txt new file mode 100644 index 0000000..1cc1247 --- /dev/null +++ b/sqs/SqsPlugin/CMakeLists.txt @@ -0,0 +1,21 @@ +set(target_name SqsPlugin) + +add_library(${target_name} SHARED + "SqsConfig.cpp" + "SqsUtil.cpp" + "SqsClient.cpp" + "TransportSqs.cpp" + "SqsServer.cpp" +) + +find_package(aws-cpp-sdk-sqs CONFIG REQUIRED) +find_package(nlohmann_json CONFIG REQUIRED) + +target_include_directories(${target_name} PUBLIC + "${CMAKE_SOURCE_DIR}/GraftcodePluginsInterfaces" +) + +target_link_libraries(${target_name} PUBLIC + aws-cpp-sdk-sqs + nlohmann_json::nlohmann_json +) diff --git a/sqs/SqsPlugin/SqsClient.cpp b/sqs/SqsPlugin/SqsClient.cpp new file mode 100644 index 0000000..384c7da --- /dev/null +++ b/sqs/SqsPlugin/SqsClient.cpp @@ -0,0 +1,152 @@ +#include "SqsClient.h" + +#include "SqsUtil.h" + +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include + +namespace Graftcode::Plugins::Sqs +{ + namespace + { + Aws::SQS::Model::MessageAttributeValue StringAttribute( + const std::string& value) + { + Aws::SQS::Model::MessageAttributeValue attribute; + attribute.SetDataType("String"); + attribute.SetStringValue(value.c_str()); + return attribute; + } + + std::runtime_error AwsError( + const char* operation, + const Aws::Client::AWSError& error) + { + return std::runtime_error( + std::string("SQS ") + operation + " failed: " + + error.GetExceptionName().c_str() + " - " + + error.GetMessage().c_str()); + } + } + + SqsClient::SqsClient(SqsConfig config) + : config_(std::move(config)) + { + ValidateClientConfig(config_); + client_ = CreateAwsSqsClient(config_); + } + + SqsClient::~SqsClient() = default; + + std::vector SqsClient::Call( + const unsigned char* data, + std::size_t size) + { + if (size > 0 && data == nullptr) { + throw std::invalid_argument( + "SQS client: null payload with non-zero length"); + } + + // SQS cannot broker-filter a shared reply queue by correlation id. + // Serializing calls keeps one configured reply queue deterministic for + // this transport instance. Separate runtime instances should use + // separate reply queues. + std::lock_guard lock(callMutex_); + + const std::string correlationId = NewCorrelationId(); + Aws::SQS::Model::SendMessageRequest send; + send.SetQueueUrl(config_.requestQueueUrl.c_str()); + send.SetMessageBody(EncodePayload(data, size)); + send.AddMessageAttributes( + CorrelationIdAttribute, + StringAttribute(correlationId)); + if (!config_.oneWay) { + send.AddMessageAttributes( + ReplyToAttribute, + StringAttribute(config_.replyQueueUrl)); + } + ConfigureFifoMessage( + send, + config_.requestQueueUrl, + config_.messageGroupId, + correlationId); + + const auto sendOutcome = client_->SendMessage(send); + if (!sendOutcome.IsSuccess()) { + throw AwsError("SendMessage", sendOutcome.GetError()); + } + + if (config_.oneWay) { + return {}; + } + + const auto deadline = std::chrono::steady_clock::now() + + std::chrono::milliseconds(config_.rpcTimeoutMs); + while (std::chrono::steady_clock::now() < deadline) { + const auto remainingMs = + std::chrono::duration_cast( + deadline - std::chrono::steady_clock::now()).count(); + const auto remainingSeconds = + static_cast(std::max( + 0, + remainingMs / 1000)); + + Aws::SQS::Model::ReceiveMessageRequest receive; + receive.SetQueueUrl(config_.replyQueueUrl.c_str()); + receive.SetMaxNumberOfMessages(10); + receive.SetWaitTimeSeconds(static_cast( + (std::min)(config_.waitTimeSeconds, remainingSeconds))); + receive.SetVisibilityTimeout(static_cast( + config_.visibilityTimeoutSeconds)); + receive.AddMessageAttributeNames("All"); + + const auto receiveOutcome = client_->ReceiveMessage(receive); + if (!receiveOutcome.IsSuccess()) { + throw AwsError("ReceiveMessage", receiveOutcome.GetError()); + } + + for (const auto& message : receiveOutcome.GetResult().GetMessages()) { + const auto deleteReply = [&]() { + Aws::SQS::Model::DeleteMessageRequest remove; + remove.SetQueueUrl(config_.replyQueueUrl.c_str()); + remove.SetReceiptHandle(message.GetReceiptHandle()); + const auto deleteOutcome = client_->DeleteMessage(remove); + if (!deleteOutcome.IsSuccess()) { + throw AwsError( + "DeleteMessage", + deleteOutcome.GetError()); + } + }; + + if (GetStringAttribute( + message, + CorrelationIdAttribute) != correlationId) { + // Reply queues are required to be client-instance-dedicated. + // Remove stale responses so they cannot block a FIFO group. + deleteReply(); + continue; + } + + std::vector response = + DecodePayload(message.GetBody()); + deleteReply(); + return response; + } + } + + throw std::runtime_error( + "SQS RPC timed out after " + + std::to_string(config_.rpcTimeoutMs) + + " ms waiting for correlation id " + correlationId); + } +} diff --git a/sqs/SqsPlugin/SqsClient.h b/sqs/SqsPlugin/SqsClient.h new file mode 100644 index 0000000..26bdc47 --- /dev/null +++ b/sqs/SqsPlugin/SqsClient.h @@ -0,0 +1,35 @@ +#pragma once + +#include "SqsConfig.h" + +#include +#include +#include +#include + +namespace Aws::SQS +{ + class SQSClient; +} + +namespace Graftcode::Plugins::Sqs +{ + class SqsClient + { + public: + explicit SqsClient(SqsConfig config); + ~SqsClient(); + + SqsClient(const SqsClient&) = delete; + SqsClient& operator=(const SqsClient&) = delete; + + std::vector Call( + const unsigned char* data, + std::size_t size); + + private: + SqsConfig config_; + std::unique_ptr client_; + std::mutex callMutex_; + }; +} diff --git a/sqs/SqsPlugin/SqsConfig.cpp b/sqs/SqsPlugin/SqsConfig.cpp new file mode 100644 index 0000000..d6faf0d --- /dev/null +++ b/sqs/SqsPlugin/SqsConfig.cpp @@ -0,0 +1,220 @@ +#include "SqsConfig.h" + +#include + +#include +#include +#include +#include +#include +#include + +namespace Graftcode::Plugins::Sqs +{ + namespace + { + using Json = nlohmann::json; + + Json ParseJsonSource(const std::string& source) + { + if (source.empty()) { + throw std::runtime_error("SQS plugin: empty configuration"); + } + + const Json inlineJson = Json::parse(source, nullptr, false); + if (!inlineJson.is_discarded()) { + return inlineJson; + } + + const std::filesystem::path path(source); + std::ifstream file(path); + if (!file.good()) { + throw std::runtime_error( + "SQS plugin: configuration is neither valid JSON nor a readable file"); + } + + std::stringstream contents; + contents << file.rdbuf(); + return Json::parse(contents.str()); + } + + const Json* FindSqsNode(const Json& node) + { + if (node.is_object()) { + const auto sqs = node.find("sqs"); + if (sqs != node.end() && sqs->is_object()) { + return &(*sqs); + } + + const bool hasSqsKeys = + node.contains("requestQueueUrl") || + node.contains("replyQueueUrl") || + node.contains("queueUrl") || + node.contains("queue") || + node.contains("replyQueue") || + node.contains("endpointOverride"); + if (hasSqsKeys) { + return &node; + } + + for (auto it = node.begin(); it != node.end(); ++it) { + if (const Json* found = FindSqsNode(*it)) { + return found; + } + } + } + else if (node.is_array()) { + for (const auto& item : node) { + if (const Json* found = FindSqsNode(item)) { + return found; + } + } + } + return nullptr; + } + + void ReadString( + const Json& node, + const char* key, + std::string& target) + { + const auto value = node.find(key); + if (value != node.end() && value->is_string()) { + target = value->get(); + } + } + + void ReadPositiveUint( + const Json& node, + const char* key, + std::uint32_t& target) + { + const auto value = node.find(key); + if (value == node.end()) { + return; + } + if (!value->is_number_integer()) { + throw std::runtime_error( + std::string("SQS plugin: '") + key + "' must be an integer"); + } + + const auto parsed = value->get(); + if (parsed <= 0 || + parsed > static_cast( + (std::numeric_limits::max)())) { + throw std::runtime_error( + std::string("SQS plugin: invalid value for '") + key + "'"); + } + target = static_cast(parsed); + } + + void ReadBool(const Json& node, const char* key, bool& target) + { + const auto value = node.find(key); + if (value == node.end()) { + return; + } + if (!value->is_boolean()) { + throw std::runtime_error( + std::string("SQS plugin: '") + key + "' must be a boolean"); + } + target = value->get(); + } + } + + SqsConfig ParseSqsConfigSource(const std::string& configSource) + { + const Json root = ParseJsonSource(configSource); + const Json* node = FindSqsNode(root); + if (node == nullptr || !node->is_object()) { + throw std::runtime_error("SQS plugin: configuration must be a JSON object"); + } + + SqsConfig config; + ReadString(*node, "name", config.name); + ReadString(*node, "region", config.region); + ReadString(*node, "endpointOverride", config.endpointOverride); + ReadString(*node, "accessKeyId", config.accessKeyId); + ReadString(*node, "secretAccessKey", config.secretAccessKey); + ReadString(*node, "sessionToken", config.sessionToken); + ReadString(*node, "requestQueueUrl", config.requestQueueUrl); + ReadString(*node, "replyQueueUrl", config.replyQueueUrl); + ReadString(*node, "messageGroupId", config.messageGroupId); + + if (config.requestQueueUrl.empty()) { + ReadString(*node, "queueUrl", config.requestQueueUrl); + } + if (config.requestQueueUrl.empty()) { + ReadString(*node, "queue", config.requestQueueUrl); + } + if (config.replyQueueUrl.empty()) { + ReadString(*node, "replyQueue", config.replyQueueUrl); + } + + ReadPositiveUint(*node, "rpcTimeoutMs", config.rpcTimeoutMs); + ReadPositiveUint(*node, "waitTimeSeconds", config.waitTimeSeconds); + ReadPositiveUint( + *node, + "visibilityTimeoutSeconds", + config.visibilityTimeoutSeconds); + ReadBool(*node, "oneWay", config.oneWay); + ReadBool(*node, "verifySsl", config.verifySsl); + + if (config.waitTimeSeconds > 20) { + throw std::runtime_error( + "SQS plugin: 'waitTimeSeconds' cannot exceed 20"); + } + if (config.visibilityTimeoutSeconds > 43200) { + throw std::runtime_error( + "SQS plugin: 'visibilityTimeoutSeconds' cannot exceed 43200"); + } + if (config.region.empty()) { + throw std::runtime_error("SQS plugin: 'region' cannot be empty"); + } + + const bool hasAccessKey = !config.accessKeyId.empty(); + const bool hasSecretKey = !config.secretAccessKey.empty(); + if (hasAccessKey != hasSecretKey) { + throw std::runtime_error( + "SQS plugin: 'accessKeyId' and 'secretAccessKey' must be provided together"); + } + + return config; + } + + void ValidateClientConfig(const SqsConfig& config) + { + if (config.requestQueueUrl.empty()) { + throw std::runtime_error( + "SQS client: missing required field 'requestQueueUrl'"); + } + if (!config.oneWay && config.replyQueueUrl.empty()) { + throw std::runtime_error( + "SQS client: missing required field 'replyQueueUrl'"); + } + } + + void ValidateServerConfig(const SqsConfig& config) + { + if (config.requestQueueUrl.empty()) { + throw std::runtime_error( + "SQS server: missing required field 'requestQueueUrl'"); + } + } + + bool IsFifoQueueUrl(const std::string& queueUrl) + { + std::string normalized = queueUrl; + if (const auto query = normalized.find('?'); + query != std::string::npos) { + normalized.erase(query); + } + while (!normalized.empty() && normalized.back() == '/') { + normalized.pop_back(); + } + + constexpr const char* suffix = ".fifo"; + return normalized.size() >= 5 && + normalized.compare(normalized.size() - 5, 5, suffix) == 0; + } +} diff --git a/sqs/SqsPlugin/SqsConfig.h b/sqs/SqsPlugin/SqsConfig.h new file mode 100644 index 0000000..47865a1 --- /dev/null +++ b/sqs/SqsPlugin/SqsConfig.h @@ -0,0 +1,30 @@ +#pragma once + +#include +#include + +namespace Graftcode::Plugins::Sqs +{ + struct SqsConfig + { + std::string name; + std::string region{ "us-east-1" }; + std::string endpointOverride; + std::string accessKeyId; + std::string secretAccessKey; + std::string sessionToken; + std::string requestQueueUrl; + std::string replyQueueUrl; + std::string messageGroupId{ "graftcode" }; + std::uint32_t rpcTimeoutMs{ 30000 }; + std::uint32_t waitTimeSeconds{ 1 }; + std::uint32_t visibilityTimeoutSeconds{ 30 }; + bool oneWay{ false }; + bool verifySsl{ true }; + }; + + SqsConfig ParseSqsConfigSource(const std::string& configSource); + void ValidateClientConfig(const SqsConfig& config); + void ValidateServerConfig(const SqsConfig& config); + bool IsFifoQueueUrl(const std::string& queueUrl); +} diff --git a/sqs/SqsPlugin/SqsServer.cpp b/sqs/SqsPlugin/SqsServer.cpp new file mode 100644 index 0000000..448ba9e --- /dev/null +++ b/sqs/SqsPlugin/SqsServer.cpp @@ -0,0 +1,275 @@ +#include "SqsServer.h" + +#include "SqsUtil.h" + +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include + +namespace Graftcode::Plugins::Sqs +{ + namespace + { + void LogInfo(const std::string& message) + { + std::cout << "[SqsServer][INFO] " << message << std::endl; + } + + void LogWarn(const std::string& message) + { + std::cout << "[SqsServer][WARN] " << message << std::endl; + } + + Aws::SQS::Model::MessageAttributeValue StringAttribute( + const std::string& value) + { + Aws::SQS::Model::MessageAttributeValue attribute; + attribute.SetDataType("String"); + attribute.SetStringValue(value.c_str()); + return attribute; + } + } + + SqsServer::~SqsServer() + { + stop(); + } + + void SqsServer::configure( + const char* jsonConfig, + ProcessMessageFn processMessage) + { + if (processMessage == nullptr) { + throw std::runtime_error( + "SQS server: processMessage callback is required"); + } + + SqsConfig config = ParseSqsConfigSource( + jsonConfig != nullptr ? jsonConfig : ""); + ValidateServerConfig(config); + + std::lock_guard lock(stateMutex_); + if (running_.load(std::memory_order_acquire)) { + throw std::runtime_error( + "SQS server: cannot configure while running"); + } + config_ = std::move(config); + processMessage_ = processMessage; + } + + void SqsServer::start() + { + bool expected = false; + if (!running_.compare_exchange_strong(expected, true)) { + return; + } + + stopRequested_.store(false, std::memory_order_release); + SqsConfig config; + ProcessMessageFn processMessage = nullptr; + { + std::lock_guard lock(stateMutex_); + config = config_; + processMessage = processMessage_; + } + + try { + ValidateServerConfig(config); + if (processMessage == nullptr) { + throw std::runtime_error( + "SQS server: processMessage callback is not configured"); + } + + auto client = CreateAwsSqsClient(config); + { + std::lock_guard lock(stateMutex_); + activeClient_ = client.get(); + if (stopRequested_.load(std::memory_order_acquire)) { + activeClient_->DisableRequestProcessing(); + } + } + LogInfo( + "receiving from '" + config.requestQueueUrl + + "' in region '" + config.region + "'"); + + while (!stopRequested_.load(std::memory_order_acquire)) { + Aws::SQS::Model::ReceiveMessageRequest receive; + receive.SetQueueUrl(config.requestQueueUrl.c_str()); + receive.SetMaxNumberOfMessages(1); + receive.SetWaitTimeSeconds( + static_cast(config.waitTimeSeconds)); + receive.SetVisibilityTimeout( + static_cast(config.visibilityTimeoutSeconds)); + receive.AddMessageAttributeNames("All"); + + const auto receiveOutcome = client->ReceiveMessage(receive); + if (!receiveOutcome.IsSuccess()) { + if (!stopRequested_.load(std::memory_order_acquire)) { + LogWarn( + "ReceiveMessage failed: " + + std::string( + receiveOutcome.GetError().GetMessage().c_str())); + std::this_thread::sleep_for(std::chrono::seconds(2)); + } + continue; + } + + for (const auto& message : + receiveOutcome.GetResult().GetMessages()) { + if (stopRequested_.load(std::memory_order_acquire)) { + break; + } + + try { + const std::string correlationId = GetStringAttribute( + message, + CorrelationIdAttribute); + const std::string replyQueueUrl = + config.oneWay + ? std::string() + : GetStringAttribute( + message, + ReplyToAttribute); + if (!config.oneWay && + (correlationId.empty() || + replyQueueUrl.empty())) { + LogWarn( + "RPC request is missing GraftcodeCorrelationId " + "or GraftcodeReplyTo; request left for " + "queue retry/DLQ policy"); + continue; + } + + const std::vector request = + DecodePayload(message.GetBody()); + std::vector response; + const auto writeResponse = []( + void* context, + const byte* data, + std::size_t size) { + auto* output = + static_cast*>(context); + if (output == nullptr) { + return; + } + if (data == nullptr || size == 0) { + output->clear(); + return; + } + output->assign(data, data + size); + }; + + if (!processMessage( + request.data(), + request.size(), + writeResponse, + &response)) { + LogWarn( + "processMessage returned false; request left " + "for queue retry/DLQ policy"); + continue; + } + + if (!replyQueueUrl.empty()) { + Aws::SQS::Model::SendMessageRequest send; + send.SetQueueUrl(replyQueueUrl.c_str()); + send.SetMessageBody(EncodePayload( + response.data(), + response.size())); + send.AddMessageAttributes( + CorrelationIdAttribute, + StringAttribute(correlationId)); + ConfigureFifoMessage( + send, + replyQueueUrl, + config.messageGroupId, + correlationId.empty() + ? NewCorrelationId() + : correlationId); + + const auto sendOutcome = client->SendMessage(send); + if (!sendOutcome.IsSuccess()) { + LogWarn( + "SendMessage reply failed: " + + std::string( + sendOutcome.GetError() + .GetMessage().c_str())); + continue; + } + } + + Aws::SQS::Model::DeleteMessageRequest remove; + remove.SetQueueUrl(config.requestQueueUrl.c_str()); + remove.SetReceiptHandle(message.GetReceiptHandle()); + const auto deleteOutcome = + client->DeleteMessage(remove); + if (!deleteOutcome.IsSuccess()) { + LogWarn( + "DeleteMessage failed: " + + std::string( + deleteOutcome.GetError() + .GetMessage().c_str())); + } + } + catch (const std::exception& error) { + LogWarn( + std::string("request processing failed: ") + + error.what()); + } + } + } + LogInfo("stopped"); + } + catch (const std::exception& error) { + LogWarn(std::string("server stopped after error: ") + error.what()); + } + catch (...) { + LogWarn("server stopped after unknown error"); + } + + { + std::lock_guard lock(stateMutex_); + activeClient_ = nullptr; + running_.store(false, std::memory_order_release); + } + stoppedCondition_.notify_all(); + } + + void SqsServer::stop() + { + stopRequested_.store(true, std::memory_order_release); + std::unique_lock lock(stateMutex_); + if (activeClient_ != nullptr) { + activeClient_->DisableRequestProcessing(); + } + stoppedCondition_.wait(lock, [this]() { + return !running_.load(std::memory_order_acquire); + }); + } +} + +#if defined(_WIN32) +#define SQS_PLUGIN_EXPORT extern "C" __declspec(dllexport) +#else +#define SQS_PLUGIN_EXPORT extern "C" +#endif + +SQS_PLUGIN_EXPORT GraftcodeGateway::IServer* CreateServer() +{ + return new Graftcode::Plugins::Sqs::SqsServer(); +} + +SQS_PLUGIN_EXPORT void DestroyServer(GraftcodeGateway::IServer* server) +{ + delete server; +} diff --git a/sqs/SqsPlugin/SqsServer.h b/sqs/SqsPlugin/SqsServer.h new file mode 100644 index 0000000..750cd7b --- /dev/null +++ b/sqs/SqsPlugin/SqsServer.h @@ -0,0 +1,38 @@ +#pragma once + +#include "IServer.h" +#include "SqsConfig.h" + +#include +#include +#include + +namespace Aws::SQS +{ + class SQSClient; +} + +namespace Graftcode::Plugins::Sqs +{ + class SqsServer final : public GraftcodeGateway::IServer + { + public: + SqsServer() = default; + ~SqsServer() override; + + void configure( + const char* jsonConfig, + ProcessMessageFn processMessage) override; + void start() override; + void stop() override; + + private: + SqsConfig config_; + ProcessMessageFn processMessage_{ nullptr }; + std::mutex stateMutex_; + std::condition_variable stoppedCondition_; + Aws::SQS::SQSClient* activeClient_{ nullptr }; + std::atomic_bool stopRequested_{ false }; + std::atomic_bool running_{ false }; + }; +} diff --git a/sqs/SqsPlugin/SqsUtil.cpp b/sqs/SqsPlugin/SqsUtil.cpp new file mode 100644 index 0000000..297efe2 --- /dev/null +++ b/sqs/SqsPlugin/SqsUtil.cpp @@ -0,0 +1,183 @@ +#include "SqsUtil.h" + +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +namespace Graftcode::Plugins::Sqs +{ + namespace + { + class AwsSdkLifetime + { + public: + AwsSdkLifetime() + { + lifecycleThread_ = std::thread([this]() { + Aws::SDKOptions options; + Aws::InitAPI(options); + { + std::lock_guard lock(mutex_); + initialized_ = true; + } + condition_.notify_all(); + + std::unique_lock lock(mutex_); + condition_.wait(lock, [this]() { return shutdown_; }); + lock.unlock(); + Aws::ShutdownAPI(options); + }); + + std::unique_lock lock(mutex_); + condition_.wait(lock, [this]() { return initialized_; }); + } + + ~AwsSdkLifetime() + { + { + std::lock_guard lock(mutex_); + shutdown_ = true; + } + condition_.notify_all(); + lifecycleThread_.join(); + } + + private: + std::mutex mutex_; + std::condition_variable condition_; + bool initialized_{ false }; + bool shutdown_{ false }; + std::thread lifecycleThread_; + }; + + void EnsureAwsSdkInitialized() + { + static AwsSdkLifetime lifetime; + (void)lifetime; + } + } + + std::unique_ptr CreateAwsSqsClient( + const SqsConfig& config) + { + EnsureAwsSdkInitialized(); + + Aws::Client::ClientConfiguration clientConfig; + clientConfig.region = config.region.c_str(); + clientConfig.verifySSL = config.verifySsl; + clientConfig.connectTimeoutMs = 5000; + clientConfig.requestTimeoutMs = static_cast( + (std::max)( + std::uint32_t{ 30000 }, + (config.waitTimeSeconds + 5) * 1000)); + if (!config.endpointOverride.empty()) { + clientConfig.endpointOverride = config.endpointOverride.c_str(); + } + + if (!config.accessKeyId.empty()) { + const Aws::Auth::AWSCredentials credentials( + config.accessKeyId.c_str(), + config.secretAccessKey.c_str(), + config.sessionToken.c_str()); + return std::make_unique( + credentials, + clientConfig); + } + + return std::make_unique(clientConfig); + } + + Aws::String EncodePayload(const unsigned char* data, std::size_t size) + { + if (size > 0 && data == nullptr) { + throw std::invalid_argument( + "SQS plugin: null payload with non-zero length"); + } + + const Aws::Utils::ByteBuffer bytes( + data != nullptr ? data : reinterpret_cast(""), + size); + Aws::String encoded = "graftcode-base64:"; + encoded += Aws::Utils::HashingUtils::Base64Encode(bytes); + if (encoded.size() > 1024 * 1024) { + throw std::runtime_error( + "SQS plugin: encoded payload exceeds the 1 MiB SQS limit"); + } + return encoded; + } + + std::vector DecodePayload(const Aws::String& body) + { + constexpr const char* prefix = "graftcode-base64:"; + constexpr std::size_t prefixSize = 17; + if (body.size() < prefixSize || + body.compare(0, prefixSize, prefix) != 0) { + throw std::runtime_error( + "SQS plugin: message body has an unsupported encoding"); + } + + const Aws::Utils::ByteBuffer decoded = + Aws::Utils::HashingUtils::Base64Decode(body.substr(prefixSize)); + if (decoded.GetLength() == 0) { + return {}; + } + return std::vector( + decoded.GetUnderlyingData(), + decoded.GetUnderlyingData() + decoded.GetLength()); + } + + std::string NewCorrelationId() + { + thread_local std::mt19937_64 randomEngine( + std::random_device{}() ^ + static_cast( + std::chrono::steady_clock::now().time_since_epoch().count())); + + const auto first = randomEngine(); + const auto second = randomEngine(); + std::ostringstream value; + value << std::hex << std::setfill('0') + << std::setw(16) << first + << std::setw(16) << second; + return value.str(); + } + + std::string GetStringAttribute( + const Aws::SQS::Model::Message& message, + const char* name) + { + const auto& attributes = message.GetMessageAttributes(); + const auto value = attributes.find(name); + if (value == attributes.end()) { + return {}; + } + return value->second.GetStringValue().c_str(); + } + + void ConfigureFifoMessage( + Aws::SQS::Model::SendMessageRequest& request, + const std::string& queueUrl, + const std::string& groupId, + const std::string& deduplicationId) + { + if (!IsFifoQueueUrl(queueUrl)) { + return; + } + request.SetMessageGroupId( + groupId.empty() ? "graftcode" : groupId.c_str()); + request.SetMessageDeduplicationId(deduplicationId.c_str()); + } +} diff --git a/sqs/SqsPlugin/SqsUtil.h b/sqs/SqsPlugin/SqsUtil.h new file mode 100644 index 0000000..f13e9b8 --- /dev/null +++ b/sqs/SqsPlugin/SqsUtil.h @@ -0,0 +1,36 @@ +#pragma once + +#include "SqsConfig.h" + +#include +#include + +#include +#include +#include + +namespace Aws::SQS +{ + class SQSClient; +} + +namespace Graftcode::Plugins::Sqs +{ + inline constexpr const char* CorrelationIdAttribute = + "GraftcodeCorrelationId"; + inline constexpr const char* ReplyToAttribute = "GraftcodeReplyTo"; + + std::unique_ptr CreateAwsSqsClient( + const SqsConfig& config); + Aws::String EncodePayload(const unsigned char* data, std::size_t size); + std::vector DecodePayload(const Aws::String& body); + std::string NewCorrelationId(); + std::string GetStringAttribute( + const Aws::SQS::Model::Message& message, + const char* name); + void ConfigureFifoMessage( + Aws::SQS::Model::SendMessageRequest& request, + const std::string& queueUrl, + const std::string& groupId, + const std::string& deduplicationId); +} diff --git a/sqs/SqsPlugin/TransportSqs.cpp b/sqs/SqsPlugin/TransportSqs.cpp new file mode 100644 index 0000000..4250a95 --- /dev/null +++ b/sqs/SqsPlugin/TransportSqs.cpp @@ -0,0 +1,96 @@ +#include "TransportSqs.h" + +#include "SqsConfig.h" + +#include +#include +#include + +using Graftcode::Plugins::Sqs::SqsClient; +using Graftcode::Plugins::Sqs::TransportSqs; + +TransportSqs::TransportSqs(const char* configSource) + : client_(std::make_unique( + Graftcode::Plugins::Sqs::ParseSqsConfigSource( + configSource != nullptr ? configSource : ""))) +{ +} + +int TransportSqs::Initialize(byte, byte, byte) +{ + return 0; +} + +int TransportSqs::SendCommand( + byte* messageByteArray, + int32_t messageByteArrayLen) +{ + if (messageByteArrayLen < 0 || + (messageByteArrayLen > 0 && messageByteArray == nullptr)) { + throw std::invalid_argument("SQS transport: invalid request payload"); + } + + std::vector response = client_->Call( + messageByteArray, + static_cast(messageByteArrayLen)); + if (!response.empty() && response.front() == static_cast(255)) { + throw std::runtime_error(std::string(response.begin() + 1, response.end())); + } + + const int responseSize = static_cast(response.size()); + { + std::lock_guard lock(responsesMutex_); + responses_[std::this_thread::get_id()] = std::move(response); + } + return responseSize; +} + +int TransportSqs::ReadResponse( + byte* responseByteArray, + int32_t responseByteArrayLen) +{ + if (responseByteArrayLen < 0 || + (responseByteArrayLen > 0 && responseByteArray == nullptr)) { + throw std::invalid_argument("SQS transport: invalid response buffer"); + } + + std::lock_guard lock(responsesMutex_); + const auto response = responses_.find(std::this_thread::get_id()); + if (response == responses_.end()) { + throw std::runtime_error( + "SQS transport: response not found for calling thread"); + } + if (responseByteArrayLen != + static_cast(response->second.size())) { + throw std::runtime_error( + "SQS transport: response buffer length mismatch"); + } + + std::copy( + response->second.begin(), + response->second.end(), + responseByteArray); + responses_.erase(response); + return 0; +} + +#if defined(_WIN32) +#define SQS_PLUGIN_EXPORT extern "C" __declspec(dllexport) +#else +#define SQS_PLUGIN_EXPORT extern "C" +#endif + +SQS_PLUGIN_EXPORT Hypertube::Native::Interfaces::ITransport* +CreateTransportChannel( + const char*, + const unsigned short, + const char* configSource) +{ + return new TransportSqs(configSource); +} + +SQS_PLUGIN_EXPORT void DestroyTransportChannel( + Hypertube::Native::Interfaces::ITransport* transport) +{ + delete transport; +} diff --git a/sqs/SqsPlugin/TransportSqs.h b/sqs/SqsPlugin/TransportSqs.h new file mode 100644 index 0000000..5dd66b0 --- /dev/null +++ b/sqs/SqsPlugin/TransportSqs.h @@ -0,0 +1,37 @@ +#pragma once + +#include "ITransport.h" +#include "SqsClient.h" + +#include +#include +#include +#include +#include + +namespace Graftcode::Plugins::Sqs +{ + class TransportSqs final + : public Hypertube::Native::Interfaces::ITransport + { + public: + explicit TransportSqs(const char* configSource); + ~TransportSqs() override = default; + + int Initialize( + byte callingRuntimeNumber, + byte calledRuntimeNumber, + byte calledRuntimeVersion) override; + int SendCommand( + byte* messageByteArray, + int32_t messageByteArrayLen) override; + int ReadResponse( + byte* responseByteArray, + int32_t responseByteArrayLen) override; + + private: + std::unique_ptr client_; + std::map> responses_; + std::mutex responsesMutex_; + }; +} diff --git a/sqs/SqsPluginTest/CMakeLists.txt b/sqs/SqsPluginTest/CMakeLists.txt new file mode 100644 index 0000000..043229f --- /dev/null +++ b/sqs/SqsPluginTest/CMakeLists.txt @@ -0,0 +1,35 @@ +include(FetchContent) + +add_executable(SqsPluginTests + "SqsTransportTests.cpp" + "${CMAKE_SOURCE_DIR}/SqsPlugin/SqsConfig.cpp" + "${CMAKE_SOURCE_DIR}/SqsPlugin/SqsUtil.cpp" + "${CMAKE_SOURCE_DIR}/SqsPlugin/SqsClient.cpp" + "${CMAKE_SOURCE_DIR}/SqsPlugin/TransportSqs.cpp" + "${CMAKE_SOURCE_DIR}/SqsPlugin/SqsServer.cpp" +) + +FetchContent_Declare( + googletest + GIT_REPOSITORY https://github.com/google/googletest.git + GIT_TAG v1.17.0 +) +set(gtest_force_shared_crt ON CACHE BOOL "" FORCE) +FetchContent_MakeAvailable(googletest) + +find_package(aws-cpp-sdk-sqs CONFIG REQUIRED) +find_package(nlohmann_json CONFIG REQUIRED) + +target_include_directories(SqsPluginTests PRIVATE + "${CMAKE_SOURCE_DIR}/SqsPlugin" + "${CMAKE_SOURCE_DIR}/GraftcodePluginsInterfaces" +) + +target_link_libraries(SqsPluginTests PRIVATE + gtest_main + aws-cpp-sdk-sqs + nlohmann_json::nlohmann_json +) + +include(GoogleTest) +gtest_discover_tests(SqsPluginTests) diff --git a/sqs/SqsPluginTest/SqsTransportTests.cpp b/sqs/SqsPluginTest/SqsTransportTests.cpp new file mode 100644 index 0000000..a702f40 --- /dev/null +++ b/sqs/SqsPluginTest/SqsTransportTests.cpp @@ -0,0 +1,210 @@ +#include + +#include "IServer.h" +#include "ITransport.h" +#include "SqsClient.h" +#include "SqsConfig.h" +#include "SqsServer.h" +#include "SqsUtil.h" + +#include +#include +#include +#include +#include +#include +#include + +extern "C" Hypertube::Native::Interfaces::ITransport* +CreateTransportChannel( + const char* ipAddress, + unsigned short port, + const char* configSource); +extern "C" void DestroyTransportChannel( + Hypertube::Native::Interfaces::ITransport* transport); + +using Graftcode::Plugins::Sqs::DecodePayload; +using Graftcode::Plugins::Sqs::EncodePayload; +using Graftcode::Plugins::Sqs::IsFifoQueueUrl; +using Graftcode::Plugins::Sqs::ParseSqsConfigSource; +using Graftcode::Plugins::Sqs::SqsClient; +using Graftcode::Plugins::Sqs::SqsConfig; +using Graftcode::Plugins::Sqs::SqsServer; +using Graftcode::Plugins::Sqs::ValidateClientConfig; + +namespace +{ + const char* kLocalConfig = R"({ + "name": "SqsPlugin", + "region": "us-east-1", + "endpointOverride": "http://localhost:4566", + "accessKeyId": "test", + "secretAccessKey": "test", + "requestQueueUrl": "http://localhost:4566/000000000000/graft-requests", + "replyQueueUrl": "http://localhost:4566/000000000000/graft-replies", + "rpcTimeoutMs": 5000, + "waitTimeSeconds": 1, + "visibilityTimeoutSeconds": 30, + "verifySsl": false + })"; + + const std::vector kPayload{ + 0x00, 0x01, 0x7f, 0x80, 0xfe, 0xff + }; + + bool EchoCallback( + const unsigned char* data, + std::size_t size, + GraftcodeGateway::IServer::WriteResponseFn writeResponse, + void* context) + { + writeResponse(context, data, size); + return true; + } + + std::string EnvOr(const char* name, const char* fallback) + { + const char* value = std::getenv(name); + return value != nullptr ? value : fallback; + } + + std::string EscapeJson(const std::string& value) + { + std::string escaped; + for (const char character : value) { + if (character == '\\' || character == '"') { + escaped.push_back('\\'); + } + escaped.push_back(character); + } + return escaped; + } +} + +TEST(SqsConfig, ParsesQueueAndAwsSettings) +{ + const SqsConfig config = ParseSqsConfigSource(kLocalConfig); + + EXPECT_EQ(config.region, "us-east-1"); + EXPECT_EQ( + config.requestQueueUrl, + "http://localhost:4566/000000000000/graft-requests"); + EXPECT_EQ( + config.replyQueueUrl, + "http://localhost:4566/000000000000/graft-replies"); + EXPECT_EQ(config.rpcTimeoutMs, 5000u); + EXPECT_FALSE(config.verifySsl); + EXPECT_NO_THROW(ValidateClientConfig(config)); +} + +TEST(SqsConfig, FindsNestedSqsNode) +{ + const SqsConfig config = ParseSqsConfigSource(R"({ + "configurations": { + "package": { + "plugin": { + "sqs": { + "region": "eu-west-1", + "requestQueueUrl": "https://example/requests", + "replyQueueUrl": "https://example/replies" + } + } + } + } + })"); + + EXPECT_EQ(config.region, "eu-west-1"); + EXPECT_EQ(config.requestQueueUrl, "https://example/requests"); +} + +TEST(SqsConfig, RequiresReplyQueueForRpc) +{ + SqsConfig config; + config.requestQueueUrl = "https://example/requests"; + EXPECT_THROW(ValidateClientConfig(config), std::runtime_error); + + config.oneWay = true; + EXPECT_NO_THROW(ValidateClientConfig(config)); +} + +TEST(SqsPayload, PreservesArbitraryBinaryData) +{ + const Aws::String encoded = EncodePayload(kPayload.data(), kPayload.size()); + EXPECT_EQ(DecodePayload(encoded), kPayload); +} + +TEST(SqsPayload, SupportsEmptyPayload) +{ + const Aws::String encoded = EncodePayload(nullptr, 0); + EXPECT_FALSE(encoded.empty()); + EXPECT_TRUE(DecodePayload(encoded).empty()); +} + +TEST(SqsConfig, DetectsFifoQueueUrls) +{ + EXPECT_TRUE(IsFifoQueueUrl("https://sqs.us-east-1.amazonaws.com/1/a.fifo")); + EXPECT_TRUE(IsFifoQueueUrl("https://localhost/a.fifo/")); + EXPECT_TRUE(IsFifoQueueUrl("https://localhost/a.fifo?x=1")); + EXPECT_FALSE(IsFifoQueueUrl("https://sqs.us-east-1.amazonaws.com/1/a")); +} + +TEST(SqsTransport, FactoryCreatesAndDestroysTransport) +{ + auto* transport = CreateTransportChannel("", 0, kLocalConfig); + ASSERT_NE(transport, nullptr); + DestroyTransportChannel(transport); +} + +TEST(SqsLive, RpcRoundTrip) +{ + const char* requestQueue = std::getenv("SQS_REQUEST_QUEUE_URL"); + const char* replyQueue = std::getenv("SQS_REPLY_QUEUE_URL"); + if (requestQueue == nullptr || replyQueue == nullptr) { + GTEST_SKIP() + << "SQS_REQUEST_QUEUE_URL and SQS_REPLY_QUEUE_URL are not set"; + } + + const std::string region = EnvOr("AWS_REGION", "us-east-1"); + const std::string endpoint = EnvOr("SQS_ENDPOINT_OVERRIDE", ""); + const std::string accessKey = EnvOr("AWS_ACCESS_KEY_ID", ""); + const std::string secretKey = EnvOr("AWS_SECRET_ACCESS_KEY", ""); + const std::string sessionToken = EnvOr("AWS_SESSION_TOKEN", ""); + + const std::string json = + std::string("{\"region\":\"") + EscapeJson(region) + + "\",\"requestQueueUrl\":\"" + EscapeJson(requestQueue) + + "\",\"replyQueueUrl\":\"" + EscapeJson(replyQueue) + + "\",\"endpointOverride\":\"" + EscapeJson(endpoint) + + "\",\"accessKeyId\":\"" + EscapeJson(accessKey) + + "\",\"secretAccessKey\":\"" + EscapeJson(secretKey) + + "\",\"sessionToken\":\"" + EscapeJson(sessionToken) + + "\",\"rpcTimeoutMs\":20000,\"waitTimeSeconds\":1}"; + + SqsServer server; + server.configure(json.c_str(), &EchoCallback); + std::thread serverThread([&server]() { + try { + server.start(); + } + catch (...) { + } + }); + std::this_thread::sleep_for(std::chrono::seconds(1)); + + std::vector response; + std::exception_ptr callError; + try { + SqsClient client(ParseSqsConfigSource(json)); + response = client.Call(kPayload.data(), kPayload.size()); + } + catch (...) { + callError = std::current_exception(); + } + + server.stop(); + serverThread.join(); + if (callError) { + std::rethrow_exception(callError); + } + EXPECT_EQ(response, kPayload); +} diff --git a/sqs/docker-compose.yml b/sqs/docker-compose.yml new file mode 100644 index 0000000..2d841ce --- /dev/null +++ b/sqs/docker-compose.yml @@ -0,0 +1,14 @@ +services: + localstack: + image: localstack/localstack:4.8 + ports: + - "4566:4566" + environment: + SERVICES: sqs + AWS_DEFAULT_REGION: us-east-1 + DEBUG: 0 + volumes: + - localstack-data:/var/lib/localstack + +volumes: + localstack-data: diff --git a/sqs/vcpkg.json b/sqs/vcpkg.json new file mode 100644 index 0000000..4805652 --- /dev/null +++ b/sqs/vcpkg.json @@ -0,0 +1,16 @@ +{ + "name": "graftcode-sqs-plugin", + "version": "1.0.0", + "description": "Graftcode Amazon SQS transport/server plugin built on the AWS SDK for C++.", + "builtin-baseline": "1460b31b08c42cc2e9ac2c79f45ec8707e2675e2", + "dependencies": [ + { + "name": "aws-sdk-cpp", + "default-features": false, + "features": [ + "sqs" + ] + }, + "nlohmann-json" + ] +} From a4002fd7e3ff6354e3a92c422f9b6a5cf660c74a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Komor?= Date: Thu, 17 Sep 2026 11:06:15 +0200 Subject: [PATCH 2/2] in progress --- sqs/SqsPlugin/CMakeLists.txt | 2 ++ sqs/SqsPluginTest/CMakeLists.txt | 2 ++ 2 files changed, 4 insertions(+) diff --git a/sqs/SqsPlugin/CMakeLists.txt b/sqs/SqsPlugin/CMakeLists.txt index 1cc1247..4e2fccb 100644 --- a/sqs/SqsPlugin/CMakeLists.txt +++ b/sqs/SqsPlugin/CMakeLists.txt @@ -8,6 +8,7 @@ add_library(${target_name} SHARED "SqsServer.cpp" ) +find_package(aws-cpp-sdk-core CONFIG REQUIRED) find_package(aws-cpp-sdk-sqs CONFIG REQUIRED) find_package(nlohmann_json CONFIG REQUIRED) @@ -17,5 +18,6 @@ target_include_directories(${target_name} PUBLIC target_link_libraries(${target_name} PUBLIC aws-cpp-sdk-sqs + aws-cpp-sdk-core nlohmann_json::nlohmann_json ) diff --git a/sqs/SqsPluginTest/CMakeLists.txt b/sqs/SqsPluginTest/CMakeLists.txt index 043229f..96b454f 100644 --- a/sqs/SqsPluginTest/CMakeLists.txt +++ b/sqs/SqsPluginTest/CMakeLists.txt @@ -17,6 +17,7 @@ FetchContent_Declare( set(gtest_force_shared_crt ON CACHE BOOL "" FORCE) FetchContent_MakeAvailable(googletest) +find_package(aws-cpp-sdk-core CONFIG REQUIRED) find_package(aws-cpp-sdk-sqs CONFIG REQUIRED) find_package(nlohmann_json CONFIG REQUIRED) @@ -28,6 +29,7 @@ target_include_directories(SqsPluginTests PRIVATE target_link_libraries(SqsPluginTests PRIVATE gtest_main aws-cpp-sdk-sqs + aws-cpp-sdk-core nlohmann_json::nlohmann_json )