diff --git a/.gitignore b/.gitignore index 4ad6f92ff..1280d573e 100644 --- a/.gitignore +++ b/.gitignore @@ -1,4 +1,5 @@ build/ +build-*/ .cache/ .vscode/ @@ -8,3 +9,6 @@ build/ __pycache__/ /data/ + +# Submission attachments are distributed separately from the framework PR. +/delivery/ diff --git a/CMakeLists.txt b/CMakeLists.txt index 6bd8069d4..6605533ec 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -12,6 +12,17 @@ set(CMAKE_CXX_STANDARD 20) set(CMAKE_CXX_STANDARD_REQUIRED ON) set(CMAKE_CXX_EXTENSIONS OFF) +# Prefer the runtime library shipped with a Conda toolchain. Its newer ABI is +# required by current NCCL packages, while the compiler's private GCC runtime +# may be intentionally older than the environment runtime. +get_filename_component(INFINITRAIN_COMPILER_BIN_DIR "${CMAKE_CXX_COMPILER}" DIRECTORY) +get_filename_component(INFINITRAIN_TOOLCHAIN_PREFIX "${INFINITRAIN_COMPILER_BIN_DIR}" DIRECTORY) +set(INFINITRAIN_TOOLCHAIN_LIB_DIR "${INFINITRAIN_TOOLCHAIN_PREFIX}/lib") +if(EXISTS "${INFINITRAIN_TOOLCHAIN_LIB_DIR}/libstdc++.so") + link_directories(BEFORE "${INFINITRAIN_TOOLCHAIN_LIB_DIR}") + add_link_options("-L${INFINITRAIN_TOOLCHAIN_LIB_DIR}" "-Wl,-rpath,${INFINITRAIN_TOOLCHAIN_LIB_DIR}") +endif() + # Generate compile_commands.json set(CMAKE_EXPORT_COMPILE_COMMANDS ON) @@ -62,7 +73,7 @@ endif() # ------------------------------------------------------------------------------ # Framework core sources (*.cc), excluding cpu kernels (they are built separately) -file(GLOB_RECURSE SRC ${PROJECT_SOURCE_DIR}/infini_train/src/*.cc) +file(GLOB_RECURSE SRC CONFIGURE_DEPENDS ${PROJECT_SOURCE_DIR}/infini_train/src/*.cc) list(FILTER SRC EXCLUDE REGEX ".*kernels/cpu/.*") if(NOT USE_CUDA) list(FILTER SRC EXCLUDE REGEX ".*runtime/cuda/.*") @@ -73,7 +84,7 @@ if(NOT USE_NCCL) endif() # CPU kernels (*.cc) -file(GLOB_RECURSE CPU_KERNELS ${PROJECT_SOURCE_DIR}/infini_train/src/kernels/cpu/*.cc) +file(GLOB_RECURSE CPU_KERNELS CONFIGURE_DEPENDS ${PROJECT_SOURCE_DIR}/infini_train/src/kernels/cpu/*.cc) # ------------------------------------------------------------------------------ # CPU kernels library @@ -96,12 +107,18 @@ if(USE_CUDA) enable_language(CUDA) find_package(CUDAToolkit REQUIRED) include_directories(${CUDAToolkit_INCLUDE_DIRS}) + # Conda's CUDA packages keep headers below targets/x86_64-linux/include, + # while FindCUDAToolkit may only report the prefix include directory. + set(CUDATOOLKIT_TARGET_INCLUDE_DIR "${CUDAToolkit_ROOT}/targets/x86_64-linux/include") + if(EXISTS "${CUDATOOLKIT_TARGET_INCLUDE_DIR}/cuda_runtime.h") + include_directories(${CUDATOOLKIT_TARGET_INCLUDE_DIR}) + endif() # CUDA compilation options set(CMAKE_CUDA_FLAGS "${CMAKE_CUDA_FLAGS} --expt-extended-lambda --expt-relaxed-constexpr") # Only compile CUDA kernels / cuda sources here (your original used src/*.cu) - file(GLOB_RECURSE CUDA_KERNELS ${PROJECT_SOURCE_DIR}/infini_train/src/*.cu) + file(GLOB_RECURSE CUDA_KERNELS CONFIGURE_DEPENDS ${PROJECT_SOURCE_DIR}/infini_train/src/*.cu) add_library(infini_train_cuda_kernels STATIC ${CUDA_KERNELS}) set_target_properties(infini_train_cuda_kernels PROPERTIES CUDA_ARCHITECTURES "75;80;90") diff --git a/README.md b/README.md index abd8070b2..1ced9f5ec 100644 --- a/README.md +++ b/README.md @@ -145,6 +145,10 @@ The generated files can be passed directly to the corresponding executables: --dataset data/mnist ``` +The MNIST example supports CPU, CUDA, and single-process DDP. Use +`--device cuda` for GPU training and `--nthread_per_process 2` for two-GPU DDP. +The `--bs` option specifies the batch size per rank. + ##### GPT-2 124M ```bash diff --git a/example/mnist/dataset.cc b/example/mnist/dataset.cc index ee683f6d4..f1974fd38 100644 --- a/example/mnist/dataset.cc +++ b/example/mnist/dataset.cc @@ -89,7 +89,8 @@ MNISTDataset::MNISTDataset(const std::string &dataset, bool train) std::format("{}/{}-labels-idx1-ubyte", dataset, train ? kTrainPrefix : kTestPrefix))), image_dims_(image_file_.dims.begin() + 1, image_file_.dims.end()), label_dims_(label_file_.dims.begin() + 1, label_file_.dims.end()), - image_size_in_bytes_(kSN3TypeToSize.at(image_file_.type) + // Images are normalized to float32 below, so sample views must use the post-conversion byte stride. + image_size_in_bytes_(sizeof(float) * std::accumulate(image_dims_.begin(), image_dims_.end(), 1, std::multiplies())), label_size_in_bytes_(kSN3TypeToSize.at(label_file_.type) * std::accumulate(label_dims_.begin(), label_dims_.end(), 1, std::multiplies())) { @@ -110,6 +111,7 @@ MNISTDataset::MNISTDataset(const std::string &dataset, bool train) } } image_file_.tensor = std::move(transposed_tensor); + image_dims_.insert(image_dims_.begin(), 1); } std::pair, std::shared_ptr> diff --git a/example/mnist/main.cc b/example/mnist/main.cc index 7744e0947..2de9e44bc 100644 --- a/example/mnist/main.cc +++ b/example/mnist/main.cc @@ -1,27 +1,40 @@ +#include +#include #include #include #include -#include #include -#include +#include #include #include "gflags/gflags.h" #include "glog/logging.h" +#if defined(USE_CUDA) +#include +#endif + +#include "infini_train/include/autograd/grad_mode.h" #include "infini_train/include/dataloader.h" #include "infini_train/include/device.h" #include "infini_train/include/nn/modules/loss.h" +#include "infini_train/include/nn/parallel/ddp/distributed_data_parallel.h" +#include "infini_train/include/nn/parallel/global.h" +#include "infini_train/include/nn/parallel/parallel_functional.h" +#include "infini_train/include/nn/parallel/process_group.h" +#include "infini_train/include/nn/parallel/rank.h" +#include "infini_train/include/nn/parallel/utils.h" #include "infini_train/include/optimizer.h" #include "example/mnist/dataset.h" #include "example/mnist/net.h" -DEFINE_string(dataset, "", "mnist dataset path"); -DEFINE_int32(bs, 64, "batch size"); -DEFINE_int32(num_epoch, 1, "num epochs"); -DEFINE_double(lr, 0.01, "learning rate"); -DEFINE_string(device, "cpu", "device type (cpu/cuda)"); +DEFINE_string(dataset, "", "MNIST dataset path"); +DEFINE_int32(bs, 64, "Per-rank batch size"); +DEFINE_int32(num_epoch, 5, "Number of training epochs"); +DEFINE_double(lr, 0.01, "Learning rate"); +DEFINE_string(device, "cpu", "Device type (cpu/cuda)"); +DEFINE_int32(nthread_per_process, 1, "Training threads per process; values greater than one enable CUDA DDP"); using namespace infini_train; @@ -31,102 +44,194 @@ constexpr int kNumClasses = 10; constexpr char kDeviceCPU[] = "cpu"; constexpr char kDeviceCUDA[] = "cuda"; -}; // namespace -DEFINE_validator(device, - [](const char *, const std::string &value) { return value == kDeviceCPU || value == kDeviceCUDA; }); +std::array ReduceMetrics(std::array metrics, Device device, + const nn::parallel::ProcessGroup *ddp_pg) { + if (ddp_pg == nullptr) { + return metrics; + } -int main(int argc, char *argv[]) { - gflags::ParseCommandLineFlags(&argc, &argv, true); - google::InitGoogleLogging(argv[0]); + auto device_metrics = std::make_shared(metrics.data(), std::vector{3}, DataType::kFLOAT32, device); + nn::parallel::function::AllReduce(device_metrics, nn::parallel::function::ReduceOpType::kSum, ddp_pg); + const Tensor host_metrics = device_metrics->To(Device()); + const auto *values = static_cast(host_metrics.DataPtr()); + return {values[0], values[1], values[2]}; +} - auto train_dataset = std::make_shared(FLAGS_dataset, true); - DataLoader train_dataloader(train_dataset, FLAGS_bs); +void Train(const nn::parallel::Rank &rank, const std::shared_ptr &network, + const std::shared_ptr &train_dataset, const std::shared_ptr &test_dataset) { + using namespace nn::parallel; + + global::thread_global_rank = rank.GlobalRank(); + const int ddp_world_size = global::GetDataParallelSize(); + int ddp_rank = 0; + const ProcessGroup *ddp_pg = nullptr; + + Device device; + if (rank.IsParallel()) { + device = Device(Device::DeviceType::kCUDA, global::GetDeviceIndex(rank.thread_rank())); + if (ddp_world_size > 1) { + auto *pg_factory = ProcessGroupFactory::Instance(device.type()); + ddp_pg = pg_factory->GetOrCreate(GetDataParallelProcessGroupName(rank.GlobalRank()), + GetDataParallelGroupRanks(rank.GlobalRank())); + ddp_rank = ddp_pg->GetGroupRank(rank.GlobalRank()); + } + } else { + device = FLAGS_device == kDeviceCPU ? Device() : Device(Device::DeviceType::kCUDA, 0); + } + const Device cpu_device; - // TODO(dcj): Add sampler & eval dataloader later. - auto test_dataset = std::make_shared(FLAGS_dataset, false); - DataLoader test_dataloader(test_dataset, FLAGS_bs); + network->To(device); + if (ddp_pg != nullptr) { + // Independent replicas are constructed serially in main; synchronize their initial parameters before DDP hooks. + ddp_pg->Broadcast(network->Parameters(), /*root_rank_in_group=*/0); + } - auto network = MNIST(); - Device device = FLAGS_device == kDeviceCPU ? Device() : Device(Device::DeviceType::kCUDA, 0); - Device cpu_device = Device(); - network.To(device); + std::shared_ptr model = network; + if (ddp_pg != nullptr) { + model = std::make_shared(model, rank, DistributedDataParallelConfig{}); + } - auto loss_fn = nn::CrossEntropyLoss(); + nn::CrossEntropyLoss loss_fn; loss_fn.To(device); - auto optimizer = optimizers::SGD(network.Parameters(), FLAGS_lr); + auto optimizer = optimizers::SGD(model->Parameters(), FLAGS_lr); + DistributedDataLoader train_dataloader(train_dataset, FLAGS_bs, ddp_rank, ddp_world_size); + DistributedDataLoader test_dataloader(test_dataset, FLAGS_bs, ddp_rank, ddp_world_size); - for (int epoch = 0; epoch < FLAGS_num_epoch; ++epoch) { - int train_idx = 0; - float total_loss = 0.0; + const size_t full_train_batches = train_dataset->Size() / FLAGS_bs; + const size_t train_steps + = ddp_world_size > 1 ? full_train_batches / ddp_world_size : (train_dataset->Size() + FLAGS_bs - 1) / FLAGS_bs; + CHECK_GT(train_steps, 0) << "The training set is smaller than one global batch."; + for (int epoch = 0; epoch < FLAGS_num_epoch; ++epoch) { + float train_loss_sum = 0.0f; + float train_samples = 0.0f; const auto epoch_start = std::chrono::high_resolution_clock::now(); - for (const auto &[image, label] : train_dataloader) { + auto train_iterator = train_dataloader.begin(); + for (size_t train_idx = 0; train_idx < train_steps; ++train_idx, ++train_iterator) { + const auto [image, label] = *train_iterator; + const auto batch_size = static_cast(image->Dims()[0]); auto new_image = std::make_shared(image->To(device)); auto new_label = std::make_shared(label->To(device)); - auto outputs = network.Forward({new_image}); optimizer.ZeroGrad(); - - auto loss = loss_fn.Forward({outputs[0], new_label}); + const auto outputs = (*model)({new_image}); + const auto loss = loss_fn.Forward({outputs[0], new_label}); loss[0]->Backward(); - // Defer the loss D2H copy until after backward; reading it earlier would synchronize CUDA - // between forward and backward. - auto loss_cpu = loss[0]->To(cpu_device); - float current_loss = static_cast(loss_cpu.DataPtr())[0]; - total_loss += current_loss; - if (train_idx % kNumItersOfOutputDuration == 0) { - LOG(ERROR) << "epoch: " << epoch << ", [" << train_idx * FLAGS_bs << "/" << train_dataset->Size() - << "] " - << " loss: " << current_loss; + // Reading after backward keeps CUDA forward and backward asynchronous with respect to the host. + const Tensor loss_cpu = loss[0]->To(cpu_device); + const float current_loss = static_cast(loss_cpu.DataPtr())[0]; + train_loss_sum += current_loss * batch_size; + train_samples += batch_size; + if (rank.IsMainRank() && train_idx % kNumItersOfOutputDuration == 0) { + LOG(INFO) << "epoch: " << epoch << ", [" << train_idx * FLAGS_bs * ddp_world_size << "/" + << train_dataset->Size() << "] loss: " << current_loss; } - optimizer.Step(); - train_idx += 1; } const auto epoch_end = std::chrono::high_resolution_clock::now(); const double duration_us = std::chrono::duration(epoch_end - epoch_start).count(); + const auto global_train_metrics = ReduceMetrics({train_loss_sum, 0.0f, train_samples}, device, ddp_pg); + if (rank.IsMainRank()) { + LOG(INFO) << std::format("epoch {:2d}/{} | train loss {:.6f} | lr {:.2e} | ({:.2f} ms | {:.0f} samples/s)", + epoch, FLAGS_num_epoch - 1, global_train_metrics[0] / global_train_metrics[2], + FLAGS_lr, duration_us / 1e3f, global_train_metrics[2] / (duration_us / 1e6)); + } - LOG(ERROR) << std::format("epoch {:2d}/{} | train loss {:.6f} | lr {:.2e} | ({:.2f} ms | {:.0f} samples/s)", - epoch, FLAGS_num_epoch - 1, total_loss / train_idx, FLAGS_lr, duration_us / 1e3f, - train_dataset->Size() / (duration_us / 1e6)); + float test_loss_sum = 0.0f; + float correct = 0.0f; + float test_samples = 0.0f; + { + autograd::NoGradGuard no_grad; + for (const auto &[image, label] : test_dataloader) { + auto new_image = std::make_shared(image->To(device)); + auto new_label = std::make_shared(label->To(device)); + const auto outputs = (*model)({new_image}); + const auto loss = loss_fn.Forward({outputs[0], new_label}); + + const Tensor output_cpu = outputs[0]->To(cpu_device); + const Tensor label_cpu = label->To(cpu_device); + const Tensor loss_cpu = loss[0]->To(cpu_device); + const int64_t batch_size = output_cpu.Dims()[0]; + const auto *labels = static_cast(label_cpu.DataPtr()); + const auto *output_values = static_cast(output_cpu.DataPtr()); + for (int64_t batch_idx = 0; batch_idx < batch_size; ++batch_idx) { + const auto *scores = output_values + batch_idx * kNumClasses; + const int prediction = std::max_element(scores, scores + kNumClasses) - scores; + correct += prediction == labels[batch_idx]; + } + test_loss_sum += static_cast(loss_cpu.DataPtr())[0] * batch_size; + test_samples += batch_size; + } + } + + const auto global_test_metrics = ReduceMetrics({test_loss_sum, correct, test_samples}, device, ddp_pg); + if (rank.IsMainRank()) { + LOG(INFO) << "epoch: " << epoch << " | test loss: " << global_test_metrics[0] / global_test_metrics[2] + << " | test accuracy: " << global_test_metrics[1] / global_test_metrics[2] << " (" + << global_test_metrics[1] << "/" << global_test_metrics[2] << ")"; + } } +} +} // namespace - // TODO(dcj): Add no_grad() context manager later. - std::vector test_losses; - int correct = 0; - int total = 0; - for (const auto &[image, label] : test_dataloader) { - auto new_image = std::make_shared(image->To(device)); - auto new_label = std::make_shared(label->To(device)); - - auto label_cpu = label->To(cpu_device); - auto outputs = network.Forward({new_image}); - auto output_cpu = outputs[0]->To(cpu_device); - auto loss = loss_fn.Forward({outputs[0], new_label}); - auto loss_cpu = loss[0]->To(cpu_device); - - const int batch_size = output_cpu.Dims()[0]; - for (int batch_idx = 0; batch_idx < batch_size; ++batch_idx) { - auto label_index = reinterpret_cast(label_cpu.DataPtr())[batch_idx]; - const auto *output_values = static_cast(output_cpu.DataPtr()) + batch_idx * kNumClasses; - const int output_index = std::max_element(output_values, output_values + kNumClasses) - output_values; - if (output_index == label_index) { - ++correct; - } +DEFINE_validator(device, + [](const char *, const std::string &value) { return value == kDeviceCPU || value == kDeviceCUDA; }); +DEFINE_validator(nthread_per_process, [](const char *, int32_t value) { return value > 0; }); + +int main(int argc, char *argv[]) { + gflags::ParseCommandLineFlags(&argc, &argv, true); + google::InitGoogleLogging(argv[0]); + CHECK(!FLAGS_dataset.empty()) << "--dataset must point to the MNIST files."; + CHECK_GT(FLAGS_bs, 0); + if (FLAGS_nthread_per_process > 1) { + CHECK_EQ(FLAGS_device, kDeviceCUDA) << "DDP requires --device=cuda."; +#if defined(USE_CUDA) + const int visible_gpu_count = [] { + int count = 0; + const cudaError_t status = cudaGetDeviceCount(&count); + CHECK_EQ(status, cudaSuccess) << "Unable to query CUDA-visible GPUs: " << cudaGetErrorString(status); + return count; + }(); + CHECK_LE(FLAGS_nthread_per_process, visible_gpu_count) + << "--nthread_per_process exceeds the number of GPUs visible through CUDA_VISIBLE_DEVICES."; +#else + LOG(FATAL) << "DDP requires a CUDA-enabled build."; +#endif + } + + nn::parallel::global::InitAllEnv(FLAGS_nthread_per_process, /*tensor_parallel_size=*/1, + /*sequence_parallel_enabled=*/false, /*pipeline_parallel_size=*/1, + /*virtual_pipeline_parallel=*/1); + LOG(INFO) << nn::parallel::global::ProcessGroupOverview(); + + auto train_dataset = std::make_shared(FLAGS_dataset, true); + auto test_dataset = std::make_shared(FLAGS_dataset, false); + + // Parameter initialization uses a shared RNG, so build each replica before launching training threads. + std::vector> networks; + networks.reserve(FLAGS_nthread_per_process); + for (int idx = 0; idx < FLAGS_nthread_per_process; ++idx) { networks.emplace_back(std::make_shared()); } + + if (FLAGS_nthread_per_process > 1) { + std::vector threads; + threads.reserve(FLAGS_nthread_per_process); + for (int idx = 0; idx < FLAGS_nthread_per_process; ++idx) { + nn::parallel::Rank rank(nn::parallel::global::GetGlobalProcRank(), idx, + nn::parallel::global::GetNprocPerNode(), FLAGS_nthread_per_process); + threads.emplace_back(Train, rank, networks[idx], train_dataset, test_dataset); } - total += batch_size; - test_losses.push_back(static_cast(loss_cpu.DataPtr())[0]); + for (auto &thread : threads) { thread.join(); } + } else { + nn::parallel::Rank rank(nn::parallel::global::GetGlobalProcRank(), 0, nn::parallel::global::GetNprocPerNode(), + FLAGS_nthread_per_process); + Train(rank, networks[0], train_dataset, test_dataset); } - const auto avg_loss = std::accumulate(test_losses.begin(), test_losses.end(), 0.0) / test_losses.size(); - LOG(ERROR) << "Total: " << total << ", Correct: " << correct - << ", Accuracy: " << static_cast(correct) / total << ", AverageLoss: " << avg_loss; gflags::ShutDownCommandLineFlags(); google::ShutdownGoogleLogging(); - return 0; } diff --git a/example/mnist/net.cc b/example/mnist/net.cc index 501fee7ef..8ace86430 100644 --- a/example/mnist/net.cc +++ b/example/mnist/net.cc @@ -7,7 +7,7 @@ #include "glog/logging.h" #include "infini_train/include/nn/modules/activations.h" -#include "infini_train/include/nn/modules/container.h" +#include "infini_train/include/nn/modules/conv2d.h" #include "infini_train/include/nn/modules/linear.h" #include "infini_train/include/nn/modules/module.h" #include "infini_train/include/tensor.h" @@ -15,17 +15,19 @@ namespace nn = infini_train::nn; MNIST::MNIST() { - std::vector> layers; - layers.push_back(std::make_shared(784, 30)); - layers.push_back(std::make_shared()); - modules_["sequential"] = std::make_shared(std::move(layers)); - modules_["linear2"] = std::make_shared(30, 10); + modules_["conv1"] = std::make_shared(1, 8, 3, 1, 1); + modules_["relu1"] = std::make_shared(); + modules_["conv2"] = std::make_shared(8, 16, 3, 2, 1); + modules_["relu2"] = std::make_shared(); + modules_["classifier"] = std::make_shared(16 * 14 * 14, 10); } std::vector> MNIST::Forward(const std::vector> &x) { CHECK_EQ(x.size(), 1); - auto x1 = (*modules_["sequential"])(x); - auto x2 = (*modules_["linear2"])(x1); - return x2; + auto output = (*modules_["conv1"])(x); + output = (*modules_["relu1"])(output); + output = (*modules_["conv2"])(output); + output = (*modules_["relu2"])(output); + return (*modules_["classifier"])({output[0]->Flatten(1)}); } diff --git a/infini_train/include/autograd/activations.h b/infini_train/include/autograd/activations.h index a63977263..809a7b67f 100644 --- a/infini_train/include/autograd/activations.h +++ b/infini_train/include/autograd/activations.h @@ -21,4 +21,16 @@ class Sigmoid : public Function { const std::vector> &output_tensors) override; std::vector> Backward(const std::vector> &grad_outputs) override; }; + +class ReLU : public Function { +public: + static constexpr char kType[] = "ReLUFunction"; + + ReLU() : Function(kType) {} + + std::vector> Forward(const std::vector> &input_tensors) override; + void SetupContext(const std::vector> &input_tensors, + const std::vector> &output_tensors) override; + std::vector> Backward(const std::vector> &grad_outputs) override; +}; } // namespace infini_train::autograd diff --git a/infini_train/include/autograd/conv2d.h b/infini_train/include/autograd/conv2d.h new file mode 100644 index 000000000..fcbc1af4f --- /dev/null +++ b/infini_train/include/autograd/conv2d.h @@ -0,0 +1,35 @@ +#pragma once + +#include +#include +#include + +#include "infini_train/include/autograd/function.h" + +namespace infini_train { +class Tensor; +} + +namespace infini_train::autograd { + +// Differentiable 2D cross-correlation for NCHW tensors with groups fixed to one. +class Conv2d : public Function { +public: + static constexpr char kType[] = "Conv2dFunction"; + + Conv2d(int64_t stride, int64_t padding) : Function(kType), stride_(stride), padding_(padding) {} + + std::vector> Forward(const std::vector> &input_tensors) override; + void SetupContext(const std::vector> &input_tensors, + const std::vector> &output_tensors) override; + std::vector> Backward(const std::vector> &grad_outputs) override; + +private: + int64_t stride_ = 1; + int64_t padding_ = 0; + bool bias_ = false; + std::vector input_dims_; + std::vector weight_dims_; +}; + +} // namespace infini_train::autograd diff --git a/infini_train/include/common/cuda/kernel_helper.cuh b/infini_train/include/common/cuda/kernel_helper.cuh index 6b532afb2..2783f73c7 100644 --- a/infini_train/include/common/cuda/kernel_helper.cuh +++ b/infini_train/include/common/cuda/kernel_helper.cuh @@ -109,8 +109,10 @@ template __device__ __forceinline__ T Cos(const T &x) { } template __device__ __forceinline__ T Tanh(const T &x) { - if constexpr (std::is_same_v || std::is_same_v) { - return htanh(x); + if constexpr (std::is_same_v) { + return __float2bfloat16(tanhf(__bfloat162float(x))); + } else if constexpr (std::is_same_v) { + return __float2half(tanhf(__half2float(x))); } else if constexpr (std::is_same_v) { return tanhf(x); } else { diff --git a/infini_train/include/nn/modules/activations.h b/infini_train/include/nn/modules/activations.h index deb029576..ee440059a 100644 --- a/infini_train/include/nn/modules/activations.h +++ b/infini_train/include/nn/modules/activations.h @@ -17,6 +17,14 @@ class Sigmoid : public CloneableModule { std::vector> Forward(const std::vector> &input_tensors) override; }; +class ReLU : public CloneableModule { +public: + static constexpr char kType[] = "ReLU"; + ReLU() : CloneableModule(kType) {} + + std::vector> Forward(const std::vector> &input_tensors) override; +}; + class NewGELU : public CloneableModule { public: static constexpr char kType[] = "NewGELU"; diff --git a/infini_train/include/nn/modules/conv2d.h b/infini_train/include/nn/modules/conv2d.h new file mode 100644 index 000000000..47cf40adf --- /dev/null +++ b/infini_train/include/nn/modules/conv2d.h @@ -0,0 +1,37 @@ +#pragma once + +#include +#include +#include + +#include "infini_train/include/device.h" +#include "infini_train/include/nn/modules/module.h" + +namespace infini_train { +class Tensor; +} + +namespace infini_train::nn { + +// A minimal Conv2d module for NCHW FP32 inputs. It supports square kernels, +// symmetric padding, a shared stride, and groups fixed to one. +class Conv2d : public CloneableModule { +public: + static constexpr char kType[] = "Conv2d"; + static constexpr char kParamWeightName[] = "weight"; + static constexpr char kParamBiasName[] = "bias"; + + Conv2d(int64_t in_channels, int64_t out_channels, int64_t kernel_size, int64_t stride = 1, int64_t padding = 0, + bool bias = true, Device device = Device()); + + std::vector> Forward(const std::vector> &input_tensors) override; + +private: + void ResetParameters(); + + int64_t stride_ = 1; + int64_t padding_ = 0; + bool bias_ = true; +}; + +} // namespace infini_train::nn diff --git a/infini_train/src/autograd/activations.cc b/infini_train/src/autograd/activations.cc index bb8b8e5ea..6fd9f911c 100644 --- a/infini_train/src/autograd/activations.cc +++ b/infini_train/src/autograd/activations.cc @@ -30,4 +30,22 @@ std::vector> Sigmoid::Backward(const std::vectorGetDevice().type(); return {Dispatcher::Instance().Call>({device, "SigmoidBackward"}, output, grad_output)}; } + +std::vector> ReLU::Forward(const std::vector> &input_tensors) { + CHECK_EQ(input_tensors.size(), 1); + const auto &input = input_tensors[0]; + return {Dispatcher::Instance().Call>({input->GetDevice().type(), "ReLUForward"}, input)}; +} + +void ReLU::SetupContext(const std::vector> &input_tensors, + const std::vector> &) { + ctx_.SaveForBackward({input_tensors[0]}); +} + +std::vector> ReLU::Backward(const std::vector> &grad_outputs) { + CHECK_EQ(grad_outputs.size(), 1); + const auto input = ctx_.GetSavedTensors()[0]; + return {Dispatcher::Instance().Call>({input->GetDevice().type(), "ReLUBackward"}, input, + grad_outputs[0])}; +} } // namespace infini_train::autograd diff --git a/infini_train/src/autograd/conv2d.cc b/infini_train/src/autograd/conv2d.cc new file mode 100644 index 000000000..af00e3e09 --- /dev/null +++ b/infini_train/src/autograd/conv2d.cc @@ -0,0 +1,57 @@ +#include "infini_train/include/autograd/conv2d.h" + +#include "glog/logging.h" + +#include "infini_train/include/dispatcher.h" +#include "infini_train/include/tensor.h" + +namespace infini_train::autograd { + +std::vector> Conv2d::Forward(const std::vector> &input_tensors) { + CHECK(input_tensors.size() == 2 || input_tensors.size() == 3); + const auto &input = input_tensors[0]; + const auto &weight = input_tensors[1]; + const auto bias = input_tensors.size() == 3 ? input_tensors[2] : nullptr; + return {Dispatcher::Instance().Call>({input->GetDevice().type(), "Conv2dForward"}, input, + weight, bias, stride_, padding_)}; +} + +void Conv2d::SetupContext(const std::vector> &input_tensors, + const std::vector> &) { + const auto &needs_input_grad = ctx_.needs_input_grad(); + const bool need_input = !needs_input_grad.empty() && needs_input_grad[0]; + const bool need_weight = needs_input_grad.size() > 1 && needs_input_grad[1]; + ctx_.SaveForBackward({need_weight ? input_tensors[0] : nullptr, need_input ? input_tensors[1] : nullptr}); + bias_ = input_tensors.size() == 3; + input_dims_ = input_tensors[0]->Dims(); + weight_dims_ = input_tensors[1]->Dims(); +} + +std::vector> Conv2d::Backward(const std::vector> &grad_outputs) { + CHECK_EQ(grad_outputs.size(), 1); + const auto &grad_output = grad_outputs[0]; + const auto saved = ctx_.GetSavedTensors(); + const bool need_input = ctx_.needs_input_grad()[0]; + const bool need_weight = ctx_.needs_input_grad()[1]; + const bool need_bias = bias_ && ctx_.needs_input_grad()[2]; + const auto device = grad_output->GetDevice().type(); + + std::shared_ptr grad_input; + std::shared_ptr grad_weight; + std::shared_ptr grad_bias; + if (need_input) { + grad_input = Dispatcher::Instance().Call>({device, "Conv2dBackwardInput"}, saved[1], + grad_output, input_dims_, stride_, padding_); + } + if (need_weight) { + grad_weight = Dispatcher::Instance().Call>( + {device, "Conv2dBackwardWeight"}, saved[0], grad_output, weight_dims_, stride_, padding_); + } + if (need_bias) { + grad_bias = Dispatcher::Instance().Call>({device, "Conv2dBackwardBias"}, grad_output); + } + return bias_ ? std::vector>{grad_input, grad_weight, grad_bias} + : std::vector>{grad_input, grad_weight}; +} + +} // namespace infini_train::autograd diff --git a/infini_train/src/dataloader.cc b/infini_train/src/dataloader.cc index 322df553a..36d25cc01 100644 --- a/infini_train/src/dataloader.cc +++ b/infini_train/src/dataloader.cc @@ -15,14 +15,17 @@ namespace infini_train { namespace { // TODO(dcj): Use official stack implementation later. std::shared_ptr Stack(const std::vector> &tensors) { - const int batch_size = tensors.size(); - const auto &dims = tensors[0]->Dims(); - const int stacked_dim = std::accumulate(dims.begin(), dims.end(), 1, std::multiplies()); - auto stacked_tensor = std::make_shared(std::vector{batch_size, stacked_dim}, tensors[0]->Dtype()); + CHECK(!tensors.empty()); + const int64_t batch_size = tensors.size(); + const auto &sample_dims = tensors[0]->Dims(); + std::vector stacked_dims; + stacked_dims.reserve(sample_dims.size() + 1); + stacked_dims.push_back(batch_size); + stacked_dims.insert(stacked_dims.end(), sample_dims.begin(), sample_dims.end()); + auto stacked_tensor = std::make_shared(stacked_dims, tensors[0]->Dtype()); for (const auto &tensor : tensors) { - CHECK_EQ(static_cast(tensors[0]->Dtype()), static_cast(tensor->Dtype())); - const auto &dims = tensor->Dims(); - CHECK_EQ(stacked_dim, std::accumulate(dims.begin(), dims.end(), 1, std::multiplies())); + CHECK(tensors[0]->Dtype() == tensor->Dtype()); + CHECK(tensor->Dims() == sample_dims) << "All samples in a batch must have the same shape."; } size_t offset = 0; diff --git a/infini_train/src/kernels/cpu/conv2d.cc b/infini_train/src/kernels/cpu/conv2d.cc new file mode 100644 index 000000000..573cbcb6f --- /dev/null +++ b/infini_train/src/kernels/cpu/conv2d.cc @@ -0,0 +1,221 @@ +#include +#include +#include + +#include "glog/logging.h" + +#include "infini_train/include/dispatcher.h" +#include "infini_train/include/tensor.h" + +namespace infini_train::kernels::cpu { +namespace { + +struct Conv2dShape { + int64_t batch; + int64_t in_channels; + int64_t out_channels; + int64_t input_height; + int64_t input_width; + int64_t kernel_height; + int64_t kernel_width; + int64_t output_height; + int64_t output_width; +}; + +Conv2dShape ValidateShapes(const std::shared_ptr &input, const std::shared_ptr &weight, int64_t stride, + int64_t padding) { + CHECK(input->Dtype() == DataType::kFLOAT32); + CHECK(weight->Dtype() == DataType::kFLOAT32); + CHECK(input->GetDevice() == weight->GetDevice()); + CHECK_EQ(input->Dims().size(), 4) << "Conv2d expects NCHW input"; + CHECK_EQ(weight->Dims().size(), 4) << "Conv2d expects OIHW weight"; + CHECK_GT(stride, 0); + CHECK_GE(padding, 0); + const auto &input_dims = input->Dims(); + const auto &weight_dims = weight->Dims(); + CHECK_EQ(input_dims[1], weight_dims[1]); + CHECK_GE(input_dims[2] + 2 * padding, weight_dims[2]); + CHECK_GE(input_dims[3] + 2 * padding, weight_dims[3]); + const int64_t output_height = (input_dims[2] + 2 * padding - weight_dims[2]) / stride + 1; + const int64_t output_width = (input_dims[3] + 2 * padding - weight_dims[3]) / stride + 1; + CHECK_GT(output_height, 0); + CHECK_GT(output_width, 0); + return {input_dims[0], input_dims[1], weight_dims[0], input_dims[2], input_dims[3], + weight_dims[2], weight_dims[3], output_height, output_width}; +} + +size_t InputOffset(const Conv2dShape &shape, int64_t batch, int64_t channel, int64_t height, int64_t width) { + return ((batch * shape.in_channels + channel) * shape.input_height + height) * shape.input_width + width; +} + +size_t WeightOffset(const Conv2dShape &shape, int64_t out_channel, int64_t in_channel, int64_t kernel_height, + int64_t kernel_width) { + return ((out_channel * shape.in_channels + in_channel) * shape.kernel_height + kernel_height) * shape.kernel_width + + kernel_width; +} + +size_t OutputOffset(const Conv2dShape &shape, int64_t batch, int64_t channel, int64_t height, int64_t width) { + return ((batch * shape.out_channels + channel) * shape.output_height + height) * shape.output_width + width; +} + +} // namespace + +std::shared_ptr Conv2dForward(const std::shared_ptr &input, const std::shared_ptr &weight, + const std::shared_ptr &bias, int64_t stride, int64_t padding) { + const Conv2dShape shape = ValidateShapes(input, weight, stride, padding); + if (bias) { + CHECK(bias->Dtype() == DataType::kFLOAT32); + CHECK(bias->GetDevice() == input->GetDevice()); + CHECK(bias->Dims() == (std::vector{shape.out_channels})); + } + auto output = std::make_shared( + std::vector{shape.batch, shape.out_channels, shape.output_height, shape.output_width}, + DataType::kFLOAT32); + const auto *input_data = static_cast(input->DataPtr()); + const auto *weight_data = static_cast(weight->DataPtr()); + const auto *bias_data = bias ? static_cast(bias->DataPtr()) : nullptr; + auto *output_data = static_cast(output->DataPtr()); + for (int64_t batch = 0; batch < shape.batch; ++batch) { + for (int64_t out_channel = 0; out_channel < shape.out_channels; ++out_channel) { + for (int64_t out_height = 0; out_height < shape.output_height; ++out_height) { + for (int64_t out_width = 0; out_width < shape.output_width; ++out_width) { + float value = bias_data ? bias_data[out_channel] : 0.0f; + for (int64_t in_channel = 0; in_channel < shape.in_channels; ++in_channel) { + for (int64_t kernel_height = 0; kernel_height < shape.kernel_height; ++kernel_height) { + const int64_t input_height = out_height * stride + kernel_height - padding; + if (input_height < 0 || input_height >= shape.input_height) { + continue; + } + for (int64_t kernel_width = 0; kernel_width < shape.kernel_width; ++kernel_width) { + const int64_t input_width = out_width * stride + kernel_width - padding; + if (input_width >= 0 && input_width < shape.input_width) { + value + += input_data[InputOffset(shape, batch, in_channel, input_height, input_width)] + * weight_data[WeightOffset(shape, out_channel, in_channel, kernel_height, + kernel_width)]; + } + } + } + } + output_data[OutputOffset(shape, batch, out_channel, out_height, out_width)] = value; + } + } + } + } + return output; +} + +std::shared_ptr Conv2dBackwardInput(const std::shared_ptr &weight, + const std::shared_ptr &grad_output, + const std::vector &input_dims, int64_t stride, int64_t padding) { + CHECK(weight->GetDevice() == grad_output->GetDevice()); + CHECK(grad_output->Dtype() == DataType::kFLOAT32); + auto input = std::make_shared(input_dims, DataType::kFLOAT32); + const Conv2dShape shape = ValidateShapes(input, weight, stride, padding); + CHECK(grad_output->Dims() + == (std::vector{shape.batch, shape.out_channels, shape.output_height, shape.output_width})); + input->Fill(0.0f); + auto *grad_input_data = static_cast(input->DataPtr()); + const auto *weight_data = static_cast(weight->DataPtr()); + const auto *grad_output_data = static_cast(grad_output->DataPtr()); + for (int64_t batch = 0; batch < shape.batch; ++batch) { + for (int64_t out_channel = 0; out_channel < shape.out_channels; ++out_channel) { + for (int64_t out_height = 0; out_height < shape.output_height; ++out_height) { + for (int64_t out_width = 0; out_width < shape.output_width; ++out_width) { + const float grad = grad_output_data[OutputOffset(shape, batch, out_channel, out_height, out_width)]; + for (int64_t in_channel = 0; in_channel < shape.in_channels; ++in_channel) { + for (int64_t kernel_height = 0; kernel_height < shape.kernel_height; ++kernel_height) { + const int64_t input_height = out_height * stride + kernel_height - padding; + if (input_height < 0 || input_height >= shape.input_height) { + continue; + } + for (int64_t kernel_width = 0; kernel_width < shape.kernel_width; ++kernel_width) { + const int64_t input_width = out_width * stride + kernel_width - padding; + if (input_width >= 0 && input_width < shape.input_width) { + grad_input_data[InputOffset(shape, batch, in_channel, input_height, input_width)] + += grad + * weight_data[WeightOffset(shape, out_channel, in_channel, kernel_height, + kernel_width)]; + } + } + } + } + } + } + } + } + return input; +} + +std::shared_ptr Conv2dBackwardWeight(const std::shared_ptr &input, + const std::shared_ptr &grad_output, + const std::vector &weight_dims, int64_t stride, int64_t padding) { + CHECK(input->GetDevice() == grad_output->GetDevice()); + CHECK(grad_output->Dtype() == DataType::kFLOAT32); + auto weight = std::make_shared(weight_dims, DataType::kFLOAT32); + const Conv2dShape shape = ValidateShapes(input, weight, stride, padding); + CHECK(grad_output->Dims() + == (std::vector{shape.batch, shape.out_channels, shape.output_height, shape.output_width})); + auto *grad_weight_data = static_cast(weight->DataPtr()); + const auto *input_data = static_cast(input->DataPtr()); + const auto *grad_output_data = static_cast(grad_output->DataPtr()); + for (int64_t out_channel = 0; out_channel < shape.out_channels; ++out_channel) { + for (int64_t in_channel = 0; in_channel < shape.in_channels; ++in_channel) { + for (int64_t kernel_height = 0; kernel_height < shape.kernel_height; ++kernel_height) { + for (int64_t kernel_width = 0; kernel_width < shape.kernel_width; ++kernel_width) { + float value = 0.0f; + for (int64_t batch = 0; batch < shape.batch; ++batch) { + for (int64_t out_height = 0; out_height < shape.output_height; ++out_height) { + const int64_t input_height = out_height * stride + kernel_height - padding; + if (input_height < 0 || input_height >= shape.input_height) { + continue; + } + for (int64_t out_width = 0; out_width < shape.output_width; ++out_width) { + const int64_t input_width = out_width * stride + kernel_width - padding; + if (input_width >= 0 && input_width < shape.input_width) { + value + += input_data[InputOffset(shape, batch, in_channel, input_height, input_width)] + * grad_output_data[OutputOffset(shape, batch, out_channel, out_height, + out_width)]; + } + } + } + } + grad_weight_data[WeightOffset(shape, out_channel, in_channel, kernel_height, kernel_width)] = value; + } + } + } + } + return weight; +} + +std::shared_ptr Conv2dBackwardBias(const std::shared_ptr &grad_output) { + CHECK(grad_output->Dtype() == DataType::kFLOAT32); + CHECK_EQ(grad_output->Dims().size(), 4); + const auto &dims = grad_output->Dims(); + auto grad_bias = std::make_shared(std::vector{dims[1]}, DataType::kFLOAT32); + const auto *grad_output_data = static_cast(grad_output->DataPtr()); + auto *grad_bias_data = static_cast(grad_bias->DataPtr()); + const size_t spatial_size = dims[2] * dims[3]; + for (int64_t channel = 0; channel < dims[1]; ++channel) { + float value = 0.0f; + for (int64_t batch = 0; batch < dims[0]; ++batch) { + const size_t offset = (batch * dims[1] + channel) * spatial_size; + for (size_t index = 0; index < spatial_size; ++index) { value += grad_output_data[offset + index]; } + } + grad_bias_data[channel] = value; + } + return grad_bias; +} + +} // namespace infini_train::kernels::cpu + +#define REGISTER_CPU_CONV2D_KERNEL(kernel_name) \ + REGISTER_KERNEL(infini_train::Device::DeviceType::kCPU, kernel_name, infini_train::kernels::cpu::kernel_name) + +REGISTER_CPU_CONV2D_KERNEL(Conv2dForward) +REGISTER_CPU_CONV2D_KERNEL(Conv2dBackwardInput) +REGISTER_CPU_CONV2D_KERNEL(Conv2dBackwardWeight) +REGISTER_CPU_CONV2D_KERNEL(Conv2dBackwardBias) + +#undef REGISTER_CPU_CONV2D_KERNEL diff --git a/infini_train/src/kernels/cpu/relu.cc b/infini_train/src/kernels/cpu/relu.cc new file mode 100644 index 000000000..12f343cc1 --- /dev/null +++ b/infini_train/src/kernels/cpu/relu.cc @@ -0,0 +1,42 @@ +#include + +#include "infini_train/include/dispatcher.h" +#include "infini_train/include/tensor.h" + +namespace infini_train::kernels::cpu { + +std::shared_ptr ReLUForward(const std::shared_ptr &input) { + CHECK(input->Dtype() == DataType::kFLOAT32); + auto output = std::make_shared(input->Dims(), DataType::kFLOAT32); + const auto *input_data = static_cast(input->DataPtr()); + auto *output_data = static_cast(output->DataPtr()); + for (size_t index = 0; index < input->NumElements(); ++index) { + output_data[index] = input_data[index] > 0.0f ? input_data[index] : 0.0f; + } + return output; +} + +std::shared_ptr ReLUBackward(const std::shared_ptr &input, const std::shared_ptr &grad_output) { + CHECK(input->Dtype() == DataType::kFLOAT32); + CHECK(grad_output->Dtype() == DataType::kFLOAT32); + CHECK(input->GetDevice() == grad_output->GetDevice()); + CHECK(input->Dims() == grad_output->Dims()); + auto grad_input = std::make_shared(input->Dims(), DataType::kFLOAT32); + const auto *input_data = static_cast(input->DataPtr()); + const auto *grad_output_data = static_cast(grad_output->DataPtr()); + auto *grad_input_data = static_cast(grad_input->DataPtr()); + for (size_t index = 0; index < input->NumElements(); ++index) { + grad_input_data[index] = input_data[index] > 0.0f ? grad_output_data[index] : 0.0f; + } + return grad_input; +} + +} // namespace infini_train::kernels::cpu + +#define REGISTER_CPU_RELU_KERNEL(kernel_name) \ + REGISTER_KERNEL(infini_train::Device::DeviceType::kCPU, kernel_name, infini_train::kernels::cpu::kernel_name) + +REGISTER_CPU_RELU_KERNEL(ReLUForward) +REGISTER_CPU_RELU_KERNEL(ReLUBackward) + +#undef REGISTER_CPU_RELU_KERNEL diff --git a/infini_train/src/kernels/cuda/conv2d.cu b/infini_train/src/kernels/cuda/conv2d.cu new file mode 100644 index 000000000..424759c65 --- /dev/null +++ b/infini_train/src/kernels/cuda/conv2d.cu @@ -0,0 +1,282 @@ +#include +#include + +#include "infini_train/include/common/cuda/common_cuda.h" +#include "infini_train/include/core/runtime/device_guard.h" +#include "infini_train/include/dispatcher.h" +#include "infini_train/include/tensor.h" + +#include "infini_train/src/core/runtime/cuda/cuda_runtime_common.h" + +namespace infini_train::kernels::cuda { +namespace { + +struct Conv2dShape { + int batch; + int in_channels; + int out_channels; + int input_height; + int input_width; + int kernel_height; + int kernel_width; + int output_height; + int output_width; +}; + +Conv2dShape ValidateShapes(const std::shared_ptr &input, const std::shared_ptr &weight, int64_t stride, + int64_t padding) { + CHECK(input->Dtype() == DataType::kFLOAT32); + CHECK(weight->Dtype() == DataType::kFLOAT32); + CHECK(input->GetDevice() == weight->GetDevice()); + CHECK_EQ(input->Dims().size(), 4) << "Conv2d expects NCHW input"; + CHECK_EQ(weight->Dims().size(), 4) << "Conv2d expects OIHW weight"; + CHECK_GT(stride, 0); + CHECK_GE(padding, 0); + const auto &input_dims = input->Dims(); + const auto &weight_dims = weight->Dims(); + CHECK_EQ(input_dims[1], weight_dims[1]); + CHECK_GE(input_dims[2] + 2 * padding, weight_dims[2]); + CHECK_GE(input_dims[3] + 2 * padding, weight_dims[3]); + const int64_t output_height = (input_dims[2] + 2 * padding - weight_dims[2]) / stride + 1; + const int64_t output_width = (input_dims[3] + 2 * padding - weight_dims[3]) / stride + 1; + CHECK_GT(output_height, 0); + CHECK_GT(output_width, 0); + return {static_cast(input_dims[0]), static_cast(input_dims[1]), static_cast(weight_dims[0]), + static_cast(input_dims[2]), static_cast(input_dims[3]), static_cast(weight_dims[2]), + static_cast(weight_dims[3]), static_cast(output_height), static_cast(output_width)}; +} + +cudaStream_t GetCudaStream(const Device &device) { + return dynamic_cast(core::GetDeviceGuardImpl(device.type())->GetStream(device)) + ->cuda_stream(); +} + +__device__ size_t InputOffset(const Conv2dShape &shape, int batch, int channel, int height, int width) { + return ((batch * shape.in_channels + channel) * shape.input_height + height) * shape.input_width + width; +} + +__device__ size_t WeightOffset(const Conv2dShape &shape, int out_channel, int in_channel, int kernel_height, + int kernel_width) { + return ((out_channel * shape.in_channels + in_channel) * shape.kernel_height + kernel_height) * shape.kernel_width + + kernel_width; +} + +__device__ size_t OutputOffset(const Conv2dShape &shape, int batch, int channel, int height, int width) { + return ((batch * shape.out_channels + channel) * shape.output_height + height) * shape.output_width + width; +} + +__global__ void Conv2dForwardKernel(const float *input, const float *weight, const float *bias, float *output, + Conv2dShape shape, int stride, int padding) { + const size_t index = blockIdx.x * blockDim.x + threadIdx.x; + const size_t num_elements + = static_cast(shape.batch) * shape.out_channels * shape.output_height * shape.output_width; + if (index >= num_elements) { + return; + } + int temporary = static_cast(index); + const int out_width = temporary % shape.output_width; + temporary /= shape.output_width; + const int out_height = temporary % shape.output_height; + temporary /= shape.output_height; + const int out_channel = temporary % shape.out_channels; + const int batch = temporary / shape.out_channels; + float value = bias ? bias[out_channel] : 0.0f; + for (int in_channel = 0; in_channel < shape.in_channels; ++in_channel) { + for (int kernel_height = 0; kernel_height < shape.kernel_height; ++kernel_height) { + const int input_height = out_height * stride + kernel_height - padding; + if (input_height < 0 || input_height >= shape.input_height) { + continue; + } + for (int kernel_width = 0; kernel_width < shape.kernel_width; ++kernel_width) { + const int input_width = out_width * stride + kernel_width - padding; + if (input_width >= 0 && input_width < shape.input_width) { + value += input[InputOffset(shape, batch, in_channel, input_height, input_width)] + * weight[WeightOffset(shape, out_channel, in_channel, kernel_height, kernel_width)]; + } + } + } + } + output[index] = value; +} + +__global__ void Conv2dBackwardInputKernel(const float *weight, const float *grad_output, float *grad_input, + Conv2dShape shape, int stride, int padding) { + const size_t index = blockIdx.x * blockDim.x + threadIdx.x; + const size_t num_elements + = static_cast(shape.batch) * shape.in_channels * shape.input_height * shape.input_width; + if (index >= num_elements) { + return; + } + int temporary = static_cast(index); + const int input_width = temporary % shape.input_width; + temporary /= shape.input_width; + const int input_height = temporary % shape.input_height; + temporary /= shape.input_height; + const int in_channel = temporary % shape.in_channels; + const int batch = temporary / shape.in_channels; + float value = 0.0f; + for (int out_channel = 0; out_channel < shape.out_channels; ++out_channel) { + for (int kernel_height = 0; kernel_height < shape.kernel_height; ++kernel_height) { + const int numerator_height = input_height + padding - kernel_height; + if (numerator_height < 0 || numerator_height % stride != 0) { + continue; + } + const int out_height = numerator_height / stride; + if (out_height >= shape.output_height) { + continue; + } + for (int kernel_width = 0; kernel_width < shape.kernel_width; ++kernel_width) { + const int numerator_width = input_width + padding - kernel_width; + if (numerator_width < 0 || numerator_width % stride != 0) { + continue; + } + const int out_width = numerator_width / stride; + if (out_width < shape.output_width) { + value += grad_output[OutputOffset(shape, batch, out_channel, out_height, out_width)] + * weight[WeightOffset(shape, out_channel, in_channel, kernel_height, kernel_width)]; + } + } + } + } + grad_input[index] = value; +} + +__global__ void Conv2dBackwardWeightKernel(const float *input, const float *grad_output, float *grad_weight, + Conv2dShape shape, int stride, int padding) { + const size_t index = blockIdx.x * blockDim.x + threadIdx.x; + const size_t num_elements + = static_cast(shape.out_channels) * shape.in_channels * shape.kernel_height * shape.kernel_width; + if (index >= num_elements) { + return; + } + int temporary = static_cast(index); + const int kernel_width = temporary % shape.kernel_width; + temporary /= shape.kernel_width; + const int kernel_height = temporary % shape.kernel_height; + temporary /= shape.kernel_height; + const int in_channel = temporary % shape.in_channels; + const int out_channel = temporary / shape.in_channels; + float value = 0.0f; + for (int batch = 0; batch < shape.batch; ++batch) { + for (int out_height = 0; out_height < shape.output_height; ++out_height) { + const int input_height = out_height * stride + kernel_height - padding; + if (input_height < 0 || input_height >= shape.input_height) { + continue; + } + for (int out_width = 0; out_width < shape.output_width; ++out_width) { + const int input_width = out_width * stride + kernel_width - padding; + if (input_width >= 0 && input_width < shape.input_width) { + value += input[InputOffset(shape, batch, in_channel, input_height, input_width)] + * grad_output[OutputOffset(shape, batch, out_channel, out_height, out_width)]; + } + } + } + } + grad_weight[index] = value; +} + +__global__ void Conv2dBackwardBiasKernel(const float *grad_output, float *grad_bias, Conv2dShape shape) { + const int out_channel = blockIdx.x * blockDim.x + threadIdx.x; + if (out_channel >= shape.out_channels) { + return; + } + float value = 0.0f; + for (int batch = 0; batch < shape.batch; ++batch) { + for (int out_height = 0; out_height < shape.output_height; ++out_height) { + for (int out_width = 0; out_width < shape.output_width; ++out_width) { + value += grad_output[OutputOffset(shape, batch, out_channel, out_height, out_width)]; + } + } + } + grad_bias[out_channel] = value; +} + +} // namespace + +std::shared_ptr Conv2dForward(const std::shared_ptr &input, const std::shared_ptr &weight, + const std::shared_ptr &bias, int64_t stride, int64_t padding) { + const Conv2dShape shape = ValidateShapes(input, weight, stride, padding); + if (bias) { + CHECK(bias->Dtype() == DataType::kFLOAT32); + CHECK(bias->GetDevice() == input->GetDevice()); + CHECK(bias->Dims() == (std::vector{shape.out_channels})); + } + auto output = std::make_shared( + std::vector{shape.batch, shape.out_channels, shape.output_height, shape.output_width}, + DataType::kFLOAT32, input->GetDevice()); + constexpr int kThreads = 256; + const size_t num_elements = output->NumElements(); + const int blocks = (num_elements + kThreads - 1) / kThreads; + Conv2dForwardKernel<<GetDevice())>>>( + static_cast(input->DataPtr()), static_cast(weight->DataPtr()), + bias ? static_cast(bias->DataPtr()) : nullptr, static_cast(output->DataPtr()), shape, + static_cast(stride), static_cast(padding)); + CUDA_CHECK(cudaGetLastError()); + return output; +} + +std::shared_ptr Conv2dBackwardInput(const std::shared_ptr &weight, + const std::shared_ptr &grad_output, + const std::vector &input_dims, int64_t stride, int64_t padding) { + CHECK(weight->GetDevice() == grad_output->GetDevice()); + CHECK(grad_output->Dtype() == DataType::kFLOAT32); + auto grad_input = std::make_shared(input_dims, DataType::kFLOAT32, grad_output->GetDevice()); + auto input_for_shape = std::make_shared(input_dims, DataType::kFLOAT32, grad_output->GetDevice()); + const Conv2dShape shape = ValidateShapes(input_for_shape, weight, stride, padding); + CHECK(grad_output->Dims() + == (std::vector{shape.batch, shape.out_channels, shape.output_height, shape.output_width})); + constexpr int kThreads = 256; + const int blocks = (grad_input->NumElements() + kThreads - 1) / kThreads; + Conv2dBackwardInputKernel<<GetDevice())>>>( + static_cast(weight->DataPtr()), static_cast(grad_output->DataPtr()), + static_cast(grad_input->DataPtr()), shape, static_cast(stride), static_cast(padding)); + CUDA_CHECK(cudaGetLastError()); + return grad_input; +} + +std::shared_ptr Conv2dBackwardWeight(const std::shared_ptr &input, + const std::shared_ptr &grad_output, + const std::vector &weight_dims, int64_t stride, int64_t padding) { + CHECK(input->GetDevice() == grad_output->GetDevice()); + CHECK(grad_output->Dtype() == DataType::kFLOAT32); + auto grad_weight = std::make_shared(weight_dims, DataType::kFLOAT32, input->GetDevice()); + auto weight_for_shape = std::make_shared(weight_dims, DataType::kFLOAT32, input->GetDevice()); + const Conv2dShape shape = ValidateShapes(input, weight_for_shape, stride, padding); + CHECK(grad_output->Dims() + == (std::vector{shape.batch, shape.out_channels, shape.output_height, shape.output_width})); + constexpr int kThreads = 256; + const int blocks = (grad_weight->NumElements() + kThreads - 1) / kThreads; + Conv2dBackwardWeightKernel<<GetDevice())>>>( + static_cast(input->DataPtr()), static_cast(grad_output->DataPtr()), + static_cast(grad_weight->DataPtr()), shape, static_cast(stride), static_cast(padding)); + CUDA_CHECK(cudaGetLastError()); + return grad_weight; +} + +std::shared_ptr Conv2dBackwardBias(const std::shared_ptr &grad_output) { + CHECK(grad_output->Dtype() == DataType::kFLOAT32); + CHECK_EQ(grad_output->Dims().size(), 4); + const auto &dims = grad_output->Dims(); + Conv2dShape shape{static_cast(dims[0]), 0, static_cast(dims[1]), 0, 0, 0, 0, static_cast(dims[2]), + static_cast(dims[3])}; + auto grad_bias + = std::make_shared(std::vector{dims[1]}, DataType::kFLOAT32, grad_output->GetDevice()); + constexpr int kThreads = 256; + const int blocks = (shape.out_channels + kThreads - 1) / kThreads; + Conv2dBackwardBiasKernel<<GetDevice())>>>( + static_cast(grad_output->DataPtr()), static_cast(grad_bias->DataPtr()), shape); + CUDA_CHECK(cudaGetLastError()); + return grad_bias; +} + +} // namespace infini_train::kernels::cuda + +#define REGISTER_CUDA_CONV2D_KERNEL(kernel_name) \ + REGISTER_KERNEL(infini_train::Device::DeviceType::kCUDA, kernel_name, infini_train::kernels::cuda::kernel_name) + +REGISTER_CUDA_CONV2D_KERNEL(Conv2dForward) +REGISTER_CUDA_CONV2D_KERNEL(Conv2dBackwardInput) +REGISTER_CUDA_CONV2D_KERNEL(Conv2dBackwardWeight) +REGISTER_CUDA_CONV2D_KERNEL(Conv2dBackwardBias) + +#undef REGISTER_CUDA_CONV2D_KERNEL diff --git a/infini_train/src/kernels/cuda/linear.cu b/infini_train/src/kernels/cuda/linear.cu index 1b4c18190..8f7a2925a 100644 --- a/infini_train/src/kernels/cuda/linear.cu +++ b/infini_train/src/kernels/cuda/linear.cu @@ -155,6 +155,24 @@ __global__ void ReduceColumnsKernel(const TIn *__restrict__ input, TOut *__restr } } +template +__global__ void ReduceRowsKernel(const TIn *__restrict__ input, TOut *__restrict__ output, int num_rows, int num_cols) { + using BlockReduce = cub::BlockReduce; + __shared__ typename BlockReduce::TempStorage temp_storage; + + const int col = blockIdx.x; + float sum = 0.0f; + + for (int row = threadIdx.x; row < num_rows; row += blockDim.x) { + sum += common::cuda::Cast(input[row * num_cols + col]); + } + + const float reduced = BlockReduce(temp_storage).Sum(sum); + if (threadIdx.x == 0) { + output[col] = reduced; + } +} + std::shared_ptr LinearBackwardInput(const std::shared_ptr &weight, const std::shared_ptr &grad_output, bool transpose, int64_t in_features, int64_t out_features, @@ -302,23 +320,24 @@ std::shared_ptr LinearBackwardBias(const std::shared_ptr &grad_o infini_train::core::GetDeviceGuardImpl(device.type())->GetStream(device)) ->cuda_stream(); - // d_bias = \sum_i(i=0, bs-1) d_output[i] - // TODO(dcj): use thrust::fill or reduce kernel do this + // d_bias[j] = sum_i d_output[i, j]. The gradient is row-major [bs, out_features], + // so each block reduces one feature across all batch rows. constexpr int BLOCK_SIZE = 256; switch (compute_dtype) { DISPATCH_CASE(WRAP({ - ReduceColumnsKernel<<>>( + ReduceRowsKernel<<>>( static_cast(grad_output->DataPtr()), - static_cast(grad_bias->DataPtr()), out_features, bs); + static_cast(grad_bias->DataPtr()), bs, out_features); }), DataType::kFLOAT32) DISPATCH_CASE(WRAP({ - ReduceColumnsKernel<<>>( + ReduceRowsKernel<<>>( static_cast(grad_output->DataPtr()), - static_cast(grad_bias->DataPtr()), out_features, bs); + static_cast(grad_bias->DataPtr()), bs, out_features); }), DataType::kBFLOAT16) } + CUDA_CHECK(cudaGetLastError()); return grad_bias; } diff --git a/infini_train/src/kernels/cuda/relu.cu b/infini_train/src/kernels/cuda/relu.cu new file mode 100644 index 000000000..e41339082 --- /dev/null +++ b/infini_train/src/kernels/cuda/relu.cu @@ -0,0 +1,68 @@ +#include + +#include "infini_train/include/common/cuda/common_cuda.h" +#include "infini_train/include/core/runtime/device_guard.h" +#include "infini_train/include/dispatcher.h" +#include "infini_train/include/tensor.h" + +#include "infini_train/src/core/runtime/cuda/cuda_runtime_common.h" + +namespace infini_train::kernels::cuda { +namespace { + +template +__global__ void ReLUKernel(const float *input, const float *grad_output, float *output, size_t num_elements) { + const size_t index = blockIdx.x * blockDim.x + threadIdx.x; + if (index >= num_elements) { + return; + } + if constexpr (kBackward) { + output[index] = input[index] > 0.0f ? grad_output[index] : 0.0f; + } else { + output[index] = input[index] > 0.0f ? input[index] : 0.0f; + } +} + +cudaStream_t GetCudaStream(const Device &device) { + return dynamic_cast(core::GetDeviceGuardImpl(device.type())->GetStream(device)) + ->cuda_stream(); +} + +} // namespace + +std::shared_ptr ReLUForward(const std::shared_ptr &input) { + CHECK(input->Dtype() == DataType::kFLOAT32); + auto output = std::make_shared(input->Dims(), DataType::kFLOAT32, input->GetDevice()); + constexpr int kThreads = 256; + const int blocks = (input->NumElements() + kThreads - 1) / kThreads; + ReLUKernel<<GetDevice())>>>( + static_cast(input->DataPtr()), nullptr, static_cast(output->DataPtr()), + input->NumElements()); + CUDA_CHECK(cudaGetLastError()); + return output; +} + +std::shared_ptr ReLUBackward(const std::shared_ptr &input, const std::shared_ptr &grad_output) { + CHECK(input->Dtype() == DataType::kFLOAT32); + CHECK(grad_output->Dtype() == DataType::kFLOAT32); + CHECK(input->GetDevice() == grad_output->GetDevice()); + CHECK(input->Dims() == grad_output->Dims()); + auto grad_input = std::make_shared(input->Dims(), DataType::kFLOAT32, input->GetDevice()); + constexpr int kThreads = 256; + const int blocks = (input->NumElements() + kThreads - 1) / kThreads; + ReLUKernel<<GetDevice())>>>( + static_cast(input->DataPtr()), static_cast(grad_output->DataPtr()), + static_cast(grad_input->DataPtr()), input->NumElements()); + CUDA_CHECK(cudaGetLastError()); + return grad_input; +} + +} // namespace infini_train::kernels::cuda + +#define REGISTER_CUDA_RELU_KERNEL(kernel_name) \ + REGISTER_KERNEL(infini_train::Device::DeviceType::kCUDA, kernel_name, infini_train::kernels::cuda::kernel_name) + +REGISTER_CUDA_RELU_KERNEL(ReLUForward) +REGISTER_CUDA_RELU_KERNEL(ReLUBackward) + +#undef REGISTER_CUDA_RELU_KERNEL diff --git a/infini_train/src/nn/modules/activations.cc b/infini_train/src/nn/modules/activations.cc index d1bbc9da8..a53738478 100644 --- a/infini_train/src/nn/modules/activations.cc +++ b/infini_train/src/nn/modules/activations.cc @@ -12,6 +12,10 @@ std::vector> Sigmoid::Forward(const std::vector()->Apply(input_tensors); } +std::vector> ReLU::Forward(const std::vector> &input_tensors) { + return std::make_shared()->Apply(input_tensors); +} + std::vector> NewGELU::Forward(const std::vector> &x) { auto &input = x[0]; return {0.5 * input diff --git a/infini_train/src/nn/modules/conv2d.cc b/infini_train/src/nn/modules/conv2d.cc new file mode 100644 index 000000000..c1a288f3c --- /dev/null +++ b/infini_train/src/nn/modules/conv2d.cc @@ -0,0 +1,51 @@ +#include "infini_train/include/nn/modules/conv2d.h" + +#include + +#include "glog/logging.h" + +#include "infini_train/include/autograd/conv2d.h" +#include "infini_train/include/nn/init.h" +#include "infini_train/include/tensor.h" + +namespace infini_train::nn { + +Conv2d::Conv2d(int64_t in_channels, int64_t out_channels, int64_t kernel_size, int64_t stride, int64_t padding, + bool bias, Device device) + : CloneableModule(kType), stride_(stride), padding_(padding), bias_(bias) { + CHECK_GT(in_channels, 0); + CHECK_GT(out_channels, 0); + CHECK_GT(kernel_size, 0); + CHECK_GT(stride, 0); + CHECK_GE(padding, 0); + device_ = device; + parameters_[kParamWeightName] + = std::make_shared(std::vector{out_channels, in_channels, kernel_size, kernel_size}, + DataType::kFLOAT32, device_) + ->RequiresGrad(); + if (bias_) { + parameters_[kParamBiasName] + = std::make_shared(std::vector{out_channels}, DataType::kFLOAT32, device_)->RequiresGrad(); + } + ResetParameters(); +} + +std::vector> Conv2d::Forward(const std::vector> &input_tensors) { + CHECK_EQ(input_tensors.size(), 1); + return std::make_shared(stride_, padding_) + ->Apply(bias_ ? std::vector>{input_tensors[0], parameters_[kParamWeightName], + parameters_[kParamBiasName]} + : std::vector>{input_tensors[0], parameters_[kParamWeightName]}); +} + +void Conv2d::ResetParameters() { + init::KaimingUniform(parameters_[kParamWeightName], 0.0f, init::KaimingMode::kFanIn, init::NonLinearityType::kReLU); + if (bias_) { + const auto [fan_in, fan_out] = init::CalculateFanInAndFanOut(parameters_[kParamWeightName]); + static_cast(fan_out); + const float bound = 1.0f / std::sqrt(static_cast(fan_in)); + init::Uniform(parameters_[kParamBiasName], -bound, bound); + } +} + +} // namespace infini_train::nn diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 3bfaa548a..af14d924e 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -15,6 +15,9 @@ add_subdirectory(module) # Tensor tests add_subdirectory(tensor) +# Example data pipeline tests +add_subdirectory(example) + # Optimizer tests add_subdirectory(optimizer) diff --git a/tests/autograd/CMakeLists.txt b/tests/autograd/CMakeLists.txt index 50bd60964..a21d2e0d5 100644 --- a/tests/autograd/CMakeLists.txt +++ b/tests/autograd/CMakeLists.txt @@ -2,7 +2,7 @@ # Autograd tests # ============================================================================ -file(GLOB AUTOGRAD_SOURCES ${CMAKE_CURRENT_SOURCE_DIR}/test_*.cc) +file(GLOB AUTOGRAD_SOURCES CONFIGURE_DEPENDS ${CMAKE_CURRENT_SOURCE_DIR}/test_*.cc) infini_train_add_test_suite(test_autograd SOURCES ${AUTOGRAD_SOURCES} diff --git a/tests/autograd/test_autograd_cnn_training_smoke.cc b/tests/autograd/test_autograd_cnn_training_smoke.cc new file mode 100644 index 000000000..27ae5c9a6 --- /dev/null +++ b/tests/autograd/test_autograd_cnn_training_smoke.cc @@ -0,0 +1,100 @@ +#include +#include +#include +#include +#include +#include + +#include "gtest/gtest.h" + +#include "infini_train/include/nn/modules/activations.h" +#include "infini_train/include/nn/modules/conv2d.h" +#include "infini_train/include/nn/modules/linear.h" +#include "infini_train/include/nn/modules/loss.h" +#include "infini_train/include/optimizer.h" +#include "infini_train/include/tensor.h" + +#include "tests/common/test_utils.h" + +using namespace infini_train; + +namespace { + +std::vector ToHostValues(const std::shared_ptr &tensor) { + const Tensor host_tensor = tensor->To(Device()); + const auto *data = static_cast(host_tensor.DataPtr()); + return {data, data + host_tensor.NumElements()}; +} + +float L1Norm(const std::shared_ptr &tensor) { + const auto values = ToHostValues(tensor); + return std::accumulate(values.begin(), values.end(), 0.0f, + [](float total, float value) { return total + std::abs(value); }); +} + +float L1Difference(const std::vector &before, const std::shared_ptr &after) { + const auto after_values = ToHostValues(after); + CHECK_EQ(before.size(), after_values.size()); + float difference = 0.0f; + for (size_t idx = 0; idx < before.size(); ++idx) { difference += std::abs(before[idx] - after_values[idx]); } + return difference; +} + +} // namespace + +class AutogradCnnTrainingSmokeTest : public infini_train::test::InfiniTrainTest {}; + +TEST_P(AutogradCnnTrainingSmokeTest, BackwardPopulatesGradientsAndSGDUpdatesParameters) { + constexpr int64_t kBatchSize = 4; + constexpr int64_t kImageSize = 28; + + std::vector image_values(kBatchSize * kImageSize * kImageSize); + for (size_t idx = 0; idx < image_values.size(); ++idx) { + image_values[idx] = static_cast((idx * 17) % 256) / 255.0f; + } + auto images + = std::make_shared(image_values.data(), std::vector{kBatchSize, 1, kImageSize, kImageSize}, + DataType::kFLOAT32, GetDevice()); + + const std::vector label_values{0, 1, 2, 3}; + auto host_labels = std::make_shared(std::vector{kBatchSize}, DataType::kUINT8); + std::copy(label_values.begin(), label_values.end(), static_cast(host_labels->DataPtr())); + auto labels = std::make_shared(host_labels->To(GetDevice())); + + auto conv1 = std::make_shared(1, 4, 3, 1, 1, true, GetDevice()); + auto relu1 = std::make_shared(); + auto conv2 = std::make_shared(4, 8, 3, 2, 1, true, GetDevice()); + auto relu2 = std::make_shared(); + auto classifier = std::make_shared(8 * 14 * 14, 10, true, GetDevice()); + + auto hidden = (*conv1)({images}); + hidden = (*relu1)(hidden); + hidden = (*conv2)(hidden); + hidden = (*relu2)(hidden); + const auto logits = (*classifier)({hidden[0]->Flatten(1)}); + + nn::CrossEntropyLoss loss_fn; + const auto loss = loss_fn.Forward({logits[0], labels}); + + std::vector> parameters = conv1->Parameters(); + const auto append_parameters = [¶meters](const std::shared_ptr &module) { + const auto module_parameters = module->Parameters(); + parameters.insert(parameters.end(), module_parameters.begin(), module_parameters.end()); + }; + append_parameters(conv2); + append_parameters(classifier); + ASSERT_EQ(parameters.size(), 6U); + const auto classifier_weight_before = ToHostValues(classifier->parameter("weight")); + + loss[0]->Backward(); + for (const auto ¶meter : parameters) { + ASSERT_NE(parameter->grad(), nullptr); + EXPECT_GT(L1Norm(parameter->grad()), 0.0f); + } + + optimizers::SGD optimizer(parameters, 0.1f); + optimizer.Step(); + EXPECT_GT(L1Difference(classifier_weight_before, classifier->parameter("weight")), 0.0f); +} + +INFINI_TRAIN_REGISTER_TEST(AutogradCnnTrainingSmokeTest); diff --git a/tests/autograd/test_autograd_conv2d.cc b/tests/autograd/test_autograd_conv2d.cc new file mode 100644 index 000000000..c931e378a --- /dev/null +++ b/tests/autograd/test_autograd_conv2d.cc @@ -0,0 +1,95 @@ +#include + +#include "gtest/gtest.h" + +#include "infini_train/include/autograd/conv2d.h" +#include "infini_train/include/autograd/grad_mode.h" +#include "infini_train/include/dispatcher.h" +#include "infini_train/include/nn/parallel/global.h" +#include "infini_train/include/tensor.h" + +#include "tests/common/test_utils.h" + +using namespace infini_train; + +class AutogradConv2dTest : public test::InfiniTrainTest {}; + +TEST_P(AutogradConv2dTest, ForwardWithoutBias) { + const std::vector input_values{1.0f, 2.0f, 3.0f, 4.0f, 5.0f, 6.0f, 7.0f, 8.0f, 9.0f}; + const std::vector weight_values{1.0f, 0.0f, 0.0f, -1.0f}; + auto input = std::make_shared(input_values.data(), std::vector{1, 1, 3, 3}, DataType::kFLOAT32, + GetDevice()); + auto weight = std::make_shared(weight_values.data(), std::vector{1, 1, 2, 2}, DataType::kFLOAT32, + GetDevice()); + + auto output = std::make_shared(1, 0)->Apply({input, weight})[0]; + EXPECT_EQ(output->Dims(), (std::vector{1, 1, 2, 2})); + test::ExpectTensorFloatEqual(output, std::vector{-4.0f, -4.0f, -4.0f, -4.0f}); +} + +TEST_P(AutogradConv2dTest, BackwardWithBias) { + const std::vector input_values{1.0f, 2.0f, 3.0f, 4.0f, 5.0f, 6.0f, 7.0f, 8.0f, 9.0f}; + const std::vector weight_values{1.0f, 0.0f, 0.0f, -1.0f}; + const std::vector bias_values{2.0f}; + auto input = std::make_shared(input_values.data(), std::vector{1, 1, 3, 3}, DataType::kFLOAT32, + GetDevice()); + auto weight = std::make_shared(weight_values.data(), std::vector{1, 1, 2, 2}, DataType::kFLOAT32, + GetDevice()); + auto bias = std::make_shared(bias_values.data(), std::vector{1}, DataType::kFLOAT32, GetDevice()); + input->RequiresGrad(); + weight->RequiresGrad(); + bias->RequiresGrad(); + + auto output = std::make_shared(1, 0)->Apply({input, weight, bias})[0]; + test::ExpectTensorFloatEqual(output, std::vector{-2.0f, -2.0f, -2.0f, -2.0f}); + auto grad_output = std::make_shared(output->Dims(), DataType::kFLOAT32, GetDevice()); + grad_output->Fill(1.0f); + output->Backward(grad_output); + + test::ExpectTensorFloatEqual(input->grad(), + std::vector{1.0f, 1.0f, 0.0f, 1.0f, 0.0f, -1.0f, 0.0f, -1.0f, -1.0f}); + test::ExpectTensorFloatEqual(weight->grad(), std::vector{12.0f, 16.0f, 24.0f, 28.0f}); + test::ExpectTensorFloatEqual(bias->grad(), std::vector{4.0f}); +} + +TEST_P(AutogradConv2dTest, BackwardSupportsPaddingAndStride) { + const std::vector input_values{1.0f, 2.0f, 3.0f, 4.0f}; + const std::vector weight_values{1.0f, 1.0f, 1.0f, 1.0f}; + auto input = std::make_shared(input_values.data(), std::vector{1, 1, 2, 2}, DataType::kFLOAT32, + GetDevice()); + auto weight = std::make_shared(weight_values.data(), std::vector{1, 1, 2, 2}, DataType::kFLOAT32, + GetDevice()); + input->RequiresGrad(); + weight->RequiresGrad(); + + auto output = std::make_shared(2, 1)->Apply({input, weight})[0]; + EXPECT_EQ(output->Dims(), (std::vector{1, 1, 2, 2})); + test::ExpectTensorFloatEqual(output, std::vector{1.0f, 2.0f, 3.0f, 4.0f}); + auto grad_output = std::make_shared(output->Dims(), DataType::kFLOAT32, GetDevice()); + grad_output->Fill(1.0f); + auto direct_grad_input = Dispatcher::Instance().Call>( + {GetDevice().type(), "Conv2dBackwardInput"}, weight, grad_output, input->Dims(), 2, 1); + test::ExpectTensorFloatEqual(direct_grad_input, std::vector{1.0f, 1.0f, 1.0f, 1.0f}); + output->Backward(grad_output); + + test::ExpectTensorFloatEqual(input->grad(), std::vector{1.0f, 1.0f, 1.0f, 1.0f}); + test::ExpectTensorFloatEqual(weight->grad(), std::vector{4.0f, 3.0f, 2.0f, 1.0f}); +} + +TEST_P(AutogradConv2dTest, SupportsNoGradMode) { + const std::vector input_values{1.0f, 2.0f, 3.0f, 4.0f}; + const std::vector weight_values{1.0f}; + auto input = std::make_shared(input_values.data(), std::vector{1, 1, 2, 2}, DataType::kFLOAT32, + GetDevice()); + auto weight = std::make_shared(weight_values.data(), std::vector{1, 1, 1, 1}, DataType::kFLOAT32, + GetDevice()); + input->RequiresGrad(); + weight->RequiresGrad(); + + autograd::NoGradGuard no_grad; + auto output = std::make_shared(1, 0)->Apply({input, weight})[0]; + EXPECT_EQ(output->grad_fn(), nullptr); + test::ExpectTensorFloatEqual(output, input_values); +} + +INFINI_TRAIN_REGISTER_TEST(AutogradConv2dTest); diff --git a/tests/autograd/test_autograd_linear_backward.cc b/tests/autograd/test_autograd_linear_backward.cc index ba0f6fe1b..089288076 100644 --- a/tests/autograd/test_autograd_linear_backward.cc +++ b/tests/autograd/test_autograd_linear_backward.cc @@ -13,18 +13,29 @@ using namespace infini_train; class AutogradLinearBackwardTest : public infini_train::test::InfiniTrainTest {}; TEST_P(AutogradLinearBackwardTest, LinearBackward) { - auto input = std::make_shared(std::vector{2, 3}, DataType::kFLOAT32, GetDevice(), true); - input->Fill(1.0f); - auto weight = std::make_shared(std::vector{4, 3}, DataType::kFLOAT32, GetDevice(), true); - weight->Fill(1.0f); - auto bias = std::make_shared(std::vector{4}, DataType::kFLOAT32, GetDevice(), true); - bias->Fill(0.0f); + const std::vector input_data{1.0f, 2.0f, 3.0f, 4.0f, 5.0f, 6.0f}; + auto input + = std::make_shared(input_data.data(), std::vector{2, 3}, DataType::kFLOAT32, GetDevice()); + input->set_requires_grad(true); + const std::vector weight_data{ + 1.0f, 0.0f, -1.0f, 2.0f, 1.0f, 0.0f, 0.0f, 1.0f, 2.0f, -1.0f, 2.0f, 1.0f, + }; + auto weight + = std::make_shared(weight_data.data(), std::vector{4, 3}, DataType::kFLOAT32, GetDevice()); + weight->set_requires_grad(true); + const std::vector bias_data{0.0f, 0.0f, 0.0f, 0.0f}; + auto bias = std::make_shared(bias_data.data(), std::vector{4}, DataType::kFLOAT32, GetDevice()); + bias->set_requires_grad(true); auto linear_fn = std::make_shared(); auto result = linear_fn->Apply({input, weight, bias}); - auto grad = std::make_shared(std::vector{2, 4}, DataType::kFLOAT32, GetDevice(), true); - grad->Fill(1.0f); + const std::vector grad_data{1.0f, 2.0f, 3.0f, 4.0f, 5.0f, 6.0f, 7.0f, 8.0f}; + auto grad = std::make_shared(grad_data.data(), std::vector{2, 4}, DataType::kFLOAT32, GetDevice()); auto grad_inputs = linear_fn->Backward({grad}); EXPECT_EQ(grad_inputs.size(), 3); + test::ExpectTensorFloatEqual(grad_inputs[0], std::vector{1.0f, 13.0f, 9.0f, 9.0f, 29.0f, 17.0f}); + test::ExpectTensorFloatEqual(grad_inputs[1], std::vector{21.0f, 27.0f, 33.0f, 26.0f, 34.0f, 42.0f, 31.0f, + 41.0f, 51.0f, 36.0f, 48.0f, 60.0f}); + test::ExpectTensorFloatEqual(grad_inputs[2], std::vector{6.0f, 8.0f, 10.0f, 12.0f}); } TEST_P(AutogradLinearBackwardTest, LinearBackwardNoBias) { @@ -38,6 +49,8 @@ TEST_P(AutogradLinearBackwardTest, LinearBackwardNoBias) { grad->Fill(1.0f); auto grad_inputs = linear_fn->Backward({grad}); EXPECT_EQ(grad_inputs.size(), 2); + test::ExpectTensorFloatEqual(grad_inputs[0], 4.0f); + test::ExpectTensorFloatEqual(grad_inputs[1], 2.0f); } INFINI_TRAIN_REGISTER_TEST(AutogradLinearBackwardTest); diff --git a/tests/autograd/test_autograd_relu.cc b/tests/autograd/test_autograd_relu.cc new file mode 100644 index 000000000..5685e147e --- /dev/null +++ b/tests/autograd/test_autograd_relu.cc @@ -0,0 +1,29 @@ +#include + +#include "gtest/gtest.h" + +#include "infini_train/include/autograd/activations.h" +#include "infini_train/include/nn/parallel/global.h" +#include "infini_train/include/tensor.h" + +#include "tests/common/test_utils.h" + +using namespace infini_train; + +class AutogradReLUTest : public test::InfiniTrainTest {}; + +TEST_P(AutogradReLUTest, ForwardAndBackward) { + const std::vector values{-2.0f, -0.0f, 1.5f, 3.0f}; + auto input = std::make_shared(values.data(), std::vector{2, 2}, DataType::kFLOAT32, GetDevice()); + input->RequiresGrad(); + + auto output = std::make_shared()->Apply({input})[0]; + test::ExpectTensorFloatEqual(output, std::vector{0.0f, 0.0f, 1.5f, 3.0f}); + + auto grad_output = std::make_shared(output->Dims(), DataType::kFLOAT32, GetDevice()); + grad_output->Fill(1.0f); + output->Backward(grad_output); + test::ExpectTensorFloatEqual(input->grad(), std::vector{0.0f, 0.0f, 1.0f, 1.0f}); +} + +INFINI_TRAIN_REGISTER_TEST(AutogradReLUTest); diff --git a/tests/example/CMakeLists.txt b/tests/example/CMakeLists.txt new file mode 100644 index 000000000..bc75a23dd --- /dev/null +++ b/tests/example/CMakeLists.txt @@ -0,0 +1,13 @@ +infini_train_add_test(test_mnist_dataset + SOURCES + test_mnist_dataset.cc + ${CMAKE_SOURCE_DIR}/example/mnist/dataset.cc + LABELS cpu +) + +target_include_directories(test_mnist_dataset PRIVATE ${CMAKE_SOURCE_DIR}) + +# External fixture producer/comparator is delivered separately from the framework PR. +add_executable(mnist_parity mnist_parity.cc ${CMAKE_SOURCE_DIR}/example/mnist/net.cc) +target_include_directories(mnist_parity PRIVATE ${CMAKE_SOURCE_DIR}) +link_infini_train_exe(mnist_parity) diff --git a/tests/example/mnist_parity.cc b/tests/example/mnist_parity.cc new file mode 100644 index 000000000..a5850aad3 --- /dev/null +++ b/tests/example/mnist_parity.cc @@ -0,0 +1,113 @@ +// Binary fixture bridge for independent CNN and DDP numerical validation. +#include +#include +#include +#include +#include + +#include "example/mnist/net.h" +#include "gflags/gflags.h" +#include "glog/logging.h" +#include "infini_train/include/nn/modules/loss.h" +#include "infini_train/include/nn/parallel/ddp/distributed_data_parallel.h" +#include "infini_train/include/nn/parallel/global.h" +#include "infini_train/include/nn/parallel/process_group.h" +#include "infini_train/include/nn/parallel/rank.h" +#include "infini_train/include/nn/parallel/utils.h" +#include "infini_train/include/optimizer.h" + +DEFINE_string(fixture_dir, "", "Directory of float32 parameter/input and uint8 label fixtures"); +DEFINE_string(output_dir, "", "Directory for per-step binary snapshots"); +DEFINE_string(device, "cpu", "cpu or cuda"); +DEFINE_int32(world_size, 1, "Number of CUDA DDP replicas"); +DEFINE_int32(batch_size, 8, "Global fixture batch size"); +DEFINE_int32(steps, 3, "Consecutive SGD updates"); +DEFINE_double(lr, 0.1, "SGD learning rate"); + +using namespace infini_train; +namespace fs = std::filesystem; + +void Read(const fs::path &path, Tensor &tensor) { + CHECK_EQ(fs::file_size(path), tensor.SizeInBytes()) << path; + std::ifstream file(path, std::ios::binary); + file.read(static_cast(tensor.DataPtr()), tensor.SizeInBytes()); + CHECK(file.good()) << path; +} +void Write(const fs::path &path, const std::shared_ptr &tensor) { + CHECK(tensor != nullptr) << path; + auto host = tensor->To(Device()); + std::ofstream file(path, std::ios::binary); + file.write(static_cast(host.DataPtr()), host.SizeInBytes()); + CHECK(file.good()) << path; +} +void Run(int index, std::shared_ptr network) { + using namespace nn::parallel; + global::thread_global_rank = index; + const Device device = FLAGS_device == "cpu" ? Device() : Device(Device::DeviceType::kCUDA, index); + network->To(device); + std::shared_ptr model = network; + if (FLAGS_world_size > 1) { + ProcessGroupFactory::Instance(device.type()) + ->GetOrCreate(GetDataParallelProcessGroupName(index), GetDataParallelGroupRanks(index)); + model = std::make_shared(network, Rank(0, index, 1, FLAGS_world_size), + DistributedDataParallelConfig{}); + } + nn::CrossEntropyLoss criterion; + optimizers::SGD optimizer(model->Parameters(), FLAGS_lr); + const fs::path output = fs::path(FLAGS_output_dir) / ("rank" + std::to_string(index)); + fs::create_directories(output); + const int local_batch = FLAGS_batch_size / FLAGS_world_size; + for (int step = 0; step < FLAGS_steps; ++step) { + const std::string prefix = "step" + std::to_string(step) + "."; + Tensor all_images({FLAGS_batch_size, 1, 28, 28}, DataType::kFLOAT32); + Tensor all_labels({FLAGS_batch_size}, DataType::kUINT8); + Read(fs::path(FLAGS_fixture_dir) / (prefix + "input.bin"), all_images); + Read(fs::path(FLAGS_fixture_dir) / (prefix + "labels.bin"), all_labels); + Tensor images(all_images, index * local_batch * 28 * 28 * sizeof(float), {local_batch, 1, 28, 28}); + Tensor labels(all_labels, index * local_batch * sizeof(uint8_t), {local_batch}); + optimizer.ZeroGrad(); + auto logits = (*model)({std::make_shared(images.To(device))})[0]; + auto loss = criterion.Forward({logits, std::make_shared(labels.To(device))})[0]; + loss->Backward(); + Write(output / (prefix + "logits.bin"), logits); + Write(output / (prefix + "loss.bin"), loss); + for (const auto &[name, parameter] : network->NamedParameters()) { + Write(output / (prefix + name + ".grad.bin"), parameter->grad()); + } + optimizer.Step(); + for (const auto &[name, parameter] : network->NamedParameters()) { + Write(output / (prefix + name + ".parameter.bin"), parameter); + } + } +} +int main(int argc, char **argv) { + gflags::ParseCommandLineFlags(&argc, &argv, true); + google::InitGoogleLogging(argv[0]); + CHECK(FLAGS_device == "cpu" || FLAGS_device == "cuda"); + CHECK_GT(FLAGS_world_size, 0); + CHECK_GT(FLAGS_batch_size, 0); + CHECK_GT(FLAGS_steps, 0); + CHECK_EQ(FLAGS_batch_size % FLAGS_world_size, 0); + CHECK(FLAGS_world_size == 1 || FLAGS_device == "cuda"); + CHECK(!FLAGS_fixture_dir.empty() && !FLAGS_output_dir.empty()); + nn::parallel::global::InitAllEnv(FLAGS_world_size, 1, false, 1, 1); + if (FLAGS_world_size > 1) { + nn::parallel::ProcessGroupFactory::Instance(Device::DeviceType::kCUDA); + } + std::vector> networks; + for (int index = 0; index < FLAGS_world_size; ++index) { + auto network = std::make_shared(); + for (const auto &[name, parameter] : network->NamedParameters()) { + Read(fs::path(FLAGS_fixture_dir) / (name + ".bin"), *parameter); + } + networks.push_back(network); + } + if (FLAGS_world_size == 1) { + Run(0, networks[0]); + } else { + std::vector threads; + for (int index = 0; index < FLAGS_world_size; ++index) { threads.emplace_back(Run, index, networks[index]); } + for (auto &thread : threads) { thread.join(); } + } + return 0; +} diff --git a/tests/example/test_mnist_dataset.cc b/tests/example/test_mnist_dataset.cc new file mode 100644 index 000000000..233890e2b --- /dev/null +++ b/tests/example/test_mnist_dataset.cc @@ -0,0 +1,79 @@ +#include +#include +#include +#include +#include +#include +#include + +#include "gtest/gtest.h" + +#include "example/mnist/dataset.h" + +namespace { + +void WriteBigEndianU32(std::ofstream *stream, uint32_t value) { + for (int shift = 24; shift >= 0; shift -= 8) { stream->put(static_cast((value >> shift) & 0xffU)); } +} + +void WriteMnistFixture(const std::filesystem::path &directory) { + std::filesystem::create_directories(directory); + + std::ofstream images(directory / "train-images-idx3-ubyte", std::ios::binary); + ASSERT_TRUE(images.is_open()); + WriteBigEndianU32(&images, 0x00000803U); + WriteBigEndianU32(&images, 2); + WriteBigEndianU32(&images, 28); + WriteBigEndianU32(&images, 28); + for (int sample = 0; sample < 2; ++sample) { + for (int pixel = 0; pixel < 28 * 28; ++pixel) { images.put(static_cast((sample * 128 + pixel) % 256)); } + } + + std::ofstream labels(directory / "train-labels-idx1-ubyte", std::ios::binary); + ASSERT_TRUE(labels.is_open()); + WriteBigEndianU32(&labels, 0x00000801U); + WriteBigEndianU32(&labels, 2); + labels.put(static_cast(3)); + labels.put(static_cast(7)); +} + +class TemporaryDirectory { +public: + TemporaryDirectory() + : path_(std::filesystem::temp_directory_path() + / ("infinitrain-mnist-dataset-" + + std::to_string(std::chrono::steady_clock::now().time_since_epoch().count()))) {} + + ~TemporaryDirectory() { std::filesystem::remove_all(path_); } + + const std::filesystem::path &path() const { return path_; } + +private: + std::filesystem::path path_; +}; + +} // namespace + +TEST(MNISTDatasetTest, UsesFloat32StrideAfterImageNormalization) { + TemporaryDirectory fixture; + WriteMnistFixture(fixture.path()); + + MNISTDataset dataset(fixture.path().string(), true); + ASSERT_EQ(dataset.Size(), 2U); + + const auto [first_image, first_label] = dataset[0]; + const auto [second_image, second_label] = dataset[1]; + EXPECT_EQ(first_image->Dims(), (std::vector{1, 28, 28})); + EXPECT_EQ(second_image->Dims(), (std::vector{1, 28, 28})); + EXPECT_EQ(first_label->Dims(), std::vector{}); + EXPECT_EQ(second_label->Dims(), std::vector{}); + + const auto *first_pixels = static_cast(first_image->DataPtr()); + const auto *second_pixels = static_cast(second_image->DataPtr()); + EXPECT_FLOAT_EQ(first_pixels[0], 0.0f); + EXPECT_FLOAT_EQ(first_pixels[1], 1.0f / 255.0f); + EXPECT_FLOAT_EQ(second_pixels[0], 128.0f / 255.0f); + EXPECT_FLOAT_EQ(second_pixels[1], 129.0f / 255.0f); + EXPECT_EQ(*static_cast(first_label->DataPtr()), 3); + EXPECT_EQ(*static_cast(second_label->DataPtr()), 7); +} diff --git a/tests/tensor/CMakeLists.txt b/tests/tensor/CMakeLists.txt index d697c7eef..5c9f872f5 100644 --- a/tests/tensor/CMakeLists.txt +++ b/tests/tensor/CMakeLists.txt @@ -8,11 +8,11 @@ # tests/tensor/cuda_only/*.cc CUDA-only tests (test_tensor_cuda_only, USE_CUDA=ON only) # Shared parameterized tests — produces test_tensor_cpu / test_tensor_cuda. -file(GLOB TENSOR_SHARED_SOURCES ${CMAKE_CURRENT_SOURCE_DIR}/test_*.cc) +file(GLOB TENSOR_SHARED_SOURCES CONFIGURE_DEPENDS ${CMAKE_CURRENT_SOURCE_DIR}/test_*.cc) infini_train_add_test_suite(test_tensor SOURCES ${TENSOR_SHARED_SOURCES}) # CPU-only tests — independent target. -file(GLOB TENSOR_CPU_ONLY_SOURCES ${CMAKE_CURRENT_SOURCE_DIR}/cpu_only/test_*.cc) +file(GLOB TENSOR_CPU_ONLY_SOURCES CONFIGURE_DEPENDS ${CMAKE_CURRENT_SOURCE_DIR}/cpu_only/test_*.cc) if(TENSOR_CPU_ONLY_SOURCES) infini_train_add_test(test_tensor_cpu_only SOURCES ${TENSOR_CPU_ONLY_SOURCES} @@ -22,7 +22,7 @@ endif() # CUDA-only tests — only built when USE_CUDA is enabled. if(USE_CUDA) - file(GLOB TENSOR_CUDA_ONLY_SOURCES ${CMAKE_CURRENT_SOURCE_DIR}/cuda_only/test_*.cc) + file(GLOB TENSOR_CUDA_ONLY_SOURCES CONFIGURE_DEPENDS ${CMAKE_CURRENT_SOURCE_DIR}/cuda_only/test_*.cc) if(TENSOR_CUDA_ONLY_SOURCES) infini_train_add_test(test_tensor_cuda_only SOURCES ${TENSOR_CUDA_ONLY_SOURCES} diff --git a/tests/tensor/test_dataloader.cc b/tests/tensor/test_dataloader.cc new file mode 100644 index 000000000..cf94af250 --- /dev/null +++ b/tests/tensor/test_dataloader.cc @@ -0,0 +1,46 @@ +#include +#include +#include + +#include "gtest/gtest.h" + +#include "infini_train/include/dataloader.h" +#include "infini_train/include/dataset.h" +#include "infini_train/include/tensor.h" + +#include "tests/common/test_utils.h" + +namespace infini_train { +namespace { + +class ImageLikeDataset final : public Dataset { +public: + std::pair, std::shared_ptr> operator[](size_t index) const override { + const float first = static_cast(index * 4 + 1); + const std::vector image_values{first, first + 1.0f, first + 2.0f, first + 3.0f}; + const std::vector label_values{static_cast(index)}; + return {std::make_shared(image_values.data(), std::vector{1, 2, 2}, DataType::kFLOAT32), + std::make_shared(label_values.data(), std::vector{}, DataType::kFLOAT32)}; + } + + size_t Size() const override { return 3; } +}; + +class DataLoaderShapeTest : public test::InfiniTrainTest {}; + +TEST_P(DataLoaderShapeTest, PreservesImageSampleDimensionsWhenStacking) { + ONLY_CPU(); + auto dataset = std::make_shared(); + DataLoader loader(dataset, 2); + + const auto [images, labels] = *loader.begin(); + EXPECT_EQ(images->Dims(), (std::vector{2, 1, 2, 2})); + EXPECT_EQ(labels->Dims(), (std::vector{2})); + test::ExpectTensorFloatEqual(images, std::vector{1.0f, 2.0f, 3.0f, 4.0f, 5.0f, 6.0f, 7.0f, 8.0f}); + test::ExpectTensorFloatEqual(labels, std::vector{0.0f, 1.0f}); +} + +INFINI_TRAIN_REGISTER_TEST(DataLoaderShapeTest); + +} // namespace +} // namespace infini_train