Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ members = [
"nodes/sink-node",
"nodes/status-node",
"nodes/sink-dynamic-node",
"examples/rust-cuda-dataflow/receiver",
]

[package]
Expand Down
2 changes: 2 additions & 0 deletions examples/cmake-arrowcuda-dataflow/.gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
build/
out/
33 changes: 33 additions & 0 deletions examples/cmake-arrowcuda-dataflow/CMakeLists.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
cmake_minimum_required(VERSION 3.15)
project(cmake-arrowcuda-dataflow LANGUAGES CXX)

set(CMAKE_CXX_STANDARD 20)
set(CMAKE_CXX_FLAGS "-fPIC")

find_package(Arrow REQUIRED)
find_package(ArrowCUDA REQUIRED)

include(DoraTargets.cmake)

link_directories(${dora_link_dirs})

add_executable(node node/main.cc ${node_bridge})
add_dependencies(node Dora_cxx)
target_include_directories(node PRIVATE ${dora_cxx_include_dir})
target_link_libraries(node dora_node_api_cxx Arrow::arrow_shared ArrowCUDA::arrow_cuda_shared)

# Platform-specific system libraries required by the dora static library
if(APPLE)
target_link_libraries(node
"-framework CoreServices" "-framework CoreFoundation" "-framework Security"
pthread m c resolv)
elseif(UNIX)
target_link_libraries(node pthread dl m rt)
elseif(WIN32)
target_link_libraries(node
advapi32 userenv kernel32 ws2_32 bcrypt ncrypt schannel ntdll iphlpapi
cfgmgr32 credui crypt32 cryptnet fwpuclnt gdi32 msimg32 mswsock ole32
oleaut32 opengl32 secur32 shell32 synchronization user32 winspool)
endif()

install(TARGETS node DESTINATION ${CMAKE_CURRENT_SOURCE_DIR}/build)
107 changes: 107 additions & 0 deletions examples/cmake-arrowcuda-dataflow/DoraTargets.cmake
Original file line number Diff line number Diff line change
@@ -0,0 +1,107 @@
set(DORA_ROOT_DIR "" CACHE FILEPATH "Path to the root of dora")

set(dora_c_include_dir "${CMAKE_CURRENT_BINARY_DIR}/include/c")

set(dora_cxx_include_dir "${CMAKE_CURRENT_BINARY_DIR}/include/cxx")
set(node_bridge "${CMAKE_CURRENT_BINARY_DIR}/node_bridge.cc")
set(operator_bridge "${CMAKE_CURRENT_BINARY_DIR}/operator_bridge.cc")

if(DORA_ROOT_DIR)
include(ExternalProject)
ExternalProject_Add(Dora
SOURCE_DIR ${DORA_ROOT_DIR}
BUILD_IN_SOURCE True
CONFIGURE_COMMAND ""
BUILD_COMMAND
cargo build
--package dora-node-api-c
&&
cargo build
--package dora-operator-api-c
&&
cargo build
--package dora-node-api-cxx
&&
cargo build
--package dora-operator-api-cxx
INSTALL_COMMAND ""
)

add_custom_command(OUTPUT ${node_bridge} ${dora_cxx_include_dir} ${operator_bridge} ${dora_c_include_dir}
WORKING_DIRECTORY ${DORA_ROOT_DIR}
DEPENDS Dora
COMMAND
mkdir -p ${dora_cxx_include_dir}
&&
mkdir -p ${CMAKE_CURRENT_BINARY_DIR}/include/c
&&
cp target/cxxbridge/dora-node-api-cxx/src/lib.rs.cc ${node_bridge}
&&
cp target/cxxbridge/dora-node-api-cxx/src/lib.rs.h ${dora_cxx_include_dir}/dora-node-api.h
&&
cp target/cxxbridge/dora-operator-api-cxx/src/lib.rs.cc ${operator_bridge}
&&
cp target/cxxbridge/dora-operator-api-cxx/src/lib.rs.h ${dora_cxx_include_dir}/dora-operator-api.h
&&
cp -r apis/c/node ${CMAKE_CURRENT_BINARY_DIR}/include/c
&&
cp -r apis/c/operator ${CMAKE_CURRENT_BINARY_DIR}/include/c

)

add_custom_target(Dora_c DEPENDS ${dora_c_include_dir})
add_custom_target(Dora_cxx DEPENDS ${node_bridge} ${operator_bridge} ${dora_cxx_include_dir})
set(dora_link_dirs ${DORA_ROOT_DIR}/target/debug)
else()
include(ExternalProject)
ExternalProject_Add(Dora
PREFIX ${CMAKE_CURRENT_BINARY_DIR}/dora
GIT_REPOSITORY https://github.com/dora-rs/dora.git
GIT_TAG main
BUILD_IN_SOURCE True
CONFIGURE_COMMAND ""
BUILD_COMMAND
cargo build
--package dora-node-api-c
--target-dir ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target
&&
cargo build
--package dora-operator-api-c
--target-dir ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target
&&
cargo build
--package dora-node-api-cxx
--target-dir ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target
&&
cargo build
--package dora-operator-api-cxx
--target-dir ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target
INSTALL_COMMAND ""
)

add_custom_command(OUTPUT ${node_bridge} ${dora_cxx_include_dir} ${operator_bridge} ${dora_c_include_dir}
WORKING_DIRECTORY ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target
DEPENDS Dora
COMMAND
mkdir -p ${CMAKE_CURRENT_BINARY_DIR}/include/c
&&
mkdir -p ${dora_cxx_include_dir}
&&
cp cxxbridge/dora-node-api-cxx/src/lib.rs.cc ${node_bridge}
&&
cp cxxbridge/dora-node-api-cxx/src/lib.rs.h ${dora_cxx_include_dir}/dora-node-api.h
&&
cp cxxbridge/dora-operator-api-cxx/src/lib.rs.cc ${operator_bridge}
&&
cp cxxbridge/dora-operator-api-cxx/src/lib.rs.h ${dora_cxx_include_dir}/dora-operator-api.h
&&
cp ../apis/c/node ${CMAKE_CURRENT_BINARY_DIR}/include/c -r
&&
cp ../apis/c/operator ${CMAKE_CURRENT_BINARY_DIR}/include/c -r
)

set(dora_link_dirs ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target/debug)

add_custom_target(Dora_c DEPENDS ${dora_c_include_dir})
add_custom_target(Dora_cxx DEPENDS ${node_bridge} ${operator_bridge} ${dora_cxx_include_dir})
endif()
14 changes: 14 additions & 0 deletions examples/cmake-arrowcuda-dataflow/dataflow.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
nodes:
- id: sender
path: sender.py
inputs:
next: receiver/next
outputs:
- cuda_data

- id: receiver
path: build/node
inputs:
cuda_data: sender/cuda_data
outputs:
- next
178 changes: 178 additions & 0 deletions examples/cmake-arrowcuda-dataflow/node/main.cc
Original file line number Diff line number Diff line change
@@ -0,0 +1,178 @@
#include "dora-node-api.h"
#include <arrow/api.h>
#include <arrow/c/bridge.h>
#include <arrow/gpu/cuda_api.h>
#include <cstring>
#include <iostream>
#include <string>

// ---------------------------------------------------------------------------
// Arrow FFI helpers
// ---------------------------------------------------------------------------

struct ArrowInput {
std::shared_ptr<arrow::Array> array;
std::string id;
rust::cxxbridge1::Box<Metadata> metadata;
};

ArrowInput receive_arrow(rust::cxxbridge1::Box<DoraEvent> event) {
struct ArrowArray c_array;
struct ArrowSchema c_schema;

auto info = event_as_arrow_input_with_info(
std::move(event),
reinterpret_cast<uint8_t *>(&c_array),
reinterpret_cast<uint8_t *>(&c_schema));

if (!info.error.empty()) {
std::cerr << "receive_arrow error: " << std::string(info.error)
<< std::endl;
return {nullptr, "", std::move(info.metadata)};
}

auto result = arrow::ImportArray(&c_array, &c_schema);
if (!result.ok()) {
std::cerr << "ImportArray failed: " << result.status().ToString()
<< std::endl;
return {nullptr, "", std::move(info.metadata)};
}
return {result.MoveValueUnsafe(), std::string(info.id),
std::move(info.metadata)};
}

bool send_arrow(rust::cxxbridge1::Box<OutputSender> &sender,
const std::string &output_id,
const std::shared_ptr<arrow::Array> &array) {
struct ArrowArray c_array;
struct ArrowSchema c_schema;

auto status = arrow::ExportArray(*array, &c_array, &c_schema);
if (!status.ok()) {
std::cerr << "ExportArray failed: " << status.ToString() << std::endl;
return false;
}

auto result = send_arrow_output(
sender, output_id,
reinterpret_cast<uint8_t *>(&c_array),
reinterpret_cast<uint8_t *>(&c_schema));

if (!result.error.empty()) {
std::cerr << "send_arrow error: " << std::string(result.error)
<< std::endl;
return false;
}
return true;
}

// ---------------------------------------------------------------------------

int main() {
auto dora_node = init_dora_node();
auto id = std::string(node_id(dora_node.send_output));
std::cout << "[" << id << "] started" << std::endl;

// Get CUDA context for device 0
auto manager = arrow::cuda::CudaDeviceManager::Instance();
if (!manager.ok()) {
std::cerr << "Failed to get CudaDeviceManager: "
<< manager.status().ToString() << std::endl;
return 1;
}
auto ctx_result = (*manager)->GetContext(0);
if (!ctx_result.ok()) {
std::cerr << "Failed to get CUDA context: "
<< ctx_result.status().ToString() << std::endl;
return 1;
}
auto ctx = *ctx_result;
std::cout << "[" << id << "] CUDA context ready on device 0" << std::endl;

// Build an empty array for "next" output
arrow::Int8Builder empty_builder;
auto empty_result = empty_builder.Finish();
auto empty_array = empty_result.ValueOrDie();

for (int i = 0; i < 100; i++) {
auto event = dora_node.events->next();
auto ty = event_type(event);

if (ty == DoraEventType::Stop ||
ty == DoraEventType::AllInputsClosed) {
break;
}
if (ty != DoraEventType::Input) {
continue;
}

auto input = receive_arrow(std::move(event));
if (!input.array) {
continue;
}

// The Arrow array is Int8 containing the 64-byte CUDA IPC handle
auto int8_array =
std::static_pointer_cast<arrow::Int8Array>(input.array);
const int8_t *raw = int8_array->raw_values();
int64_t handle_len = int8_array->length();

// Read metadata
int64_t size = input.metadata->get_int("size");
auto shape = input.metadata->get_list_int("shape");
auto dtype_str = std::string(input.metadata->get_str("dtype"));

std::cout << "[" << id << "] received CUDA IPC handle (" << handle_len
<< " bytes), buffer size=" << size
<< ", dtype=" << dtype_str << ", shape=[";
for (size_t j = 0; j < shape.size(); j++) {
if (j > 0) std::cout << ", ";
std::cout << shape[j];
}
std::cout << "]" << std::endl;

// Reconstruct IPC handle from raw bytes
auto ipc_handle_result = arrow::cuda::CudaIpcMemHandle::FromBuffer(
reinterpret_cast<const void *>(raw));
if (!ipc_handle_result.ok()) {
std::cerr << "FromBuffer failed: "
<< ipc_handle_result.status().ToString() << std::endl;
continue;
}
auto ipc_handle = *ipc_handle_result;

// Open the IPC buffer — zero-copy access to sender's GPU memory
auto buffer_result = ctx->OpenIpcBuffer(*ipc_handle);
if (!buffer_result.ok()) {
std::cerr << "OpenIpcBuffer failed: "
<< buffer_result.status().ToString() << std::endl;
continue;
}
auto cuda_buffer = *buffer_result;

// Copy first few values to host and print
int64_t copy_bytes = std::min(size, (int64_t)(10 * sizeof(int64_t)));
std::vector<uint8_t> host_buf(copy_bytes);
auto copy_status =
cuda_buffer->CopyToHost(0, copy_bytes, host_buf.data());
if (!copy_status.ok()) {
std::cerr << "CopyToHost failed: " << copy_status.ToString()
<< std::endl;
continue;
}

auto *values = reinterpret_cast<int64_t *>(host_buf.data());
int num_values = copy_bytes / sizeof(int64_t);
std::cout << "[" << id << "] first " << num_values << " values:";
for (int j = 0; j < num_values; j++) {
std::cout << " " << values[j];
}
std::cout << std::endl;

// Send "next" to trigger the sender again
send_arrow(dora_node.send_output, "next", empty_array);
}

std::cout << "[" << id << "] done" << std::endl;
return 0;
}
19 changes: 19 additions & 0 deletions examples/cmake-arrowcuda-dataflow/sender.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
import pyarrow as pa
import torch
from dora import Node
from dora.cuda import torch_to_ipc_buffer

torch.tensor([], device="cuda")

node = Node()

for i in range(10):
tensor = torch.arange(i * 100, (i + 1) * 100, dtype=torch.int64, device="cuda")
ipc_buffer, metadata = torch_to_ipc_buffer(tensor)
node.send_output("cuda_data", ipc_buffer, metadata)

event = node.next()
if event["type"] != "INPUT":
break

print("sender done")
15 changes: 15 additions & 0 deletions examples/rust-cuda-dataflow/dataflow.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
nodes:
- id: sender
path: sender.py
inputs:
next: receiver/next
outputs:
- cuda_data

- id: receiver
build: cargo build -p rust-cuda-dataflow-receiver
path: ../../../target/debug/rust-cuda-dataflow-receiver
inputs:
cuda_data: sender/cuda_data
outputs:
- next
10 changes: 10 additions & 0 deletions examples/rust-cuda-dataflow/receiver/Cargo.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
[package]
name = "rust-cuda-dataflow-receiver"
version = "0.1.0"
edition = "2024"
publish = false

[dependencies]
dora-node-api = { git = "https://github.com/dora-rs/dora.git", rev = "77c277910b0ce87b902faa1ab369a33cbcd555f4" }
arrow-rs-cuda = { git = "https://github.com/haixuantao/arrow-rs-cuda" }
eyre = "0.6.8"
Loading
Loading