blob: 3a5af6ec94513e288d78637f859ab3347f14b5f7 [file] [edit]
/* Copyright 2020 The OpenXLA Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
==============================================================================*/
#ifndef XLA_PJRT_GPU_SE_GPU_PJRT_CLIENT_H_
#define XLA_PJRT_GPU_SE_GPU_PJRT_CLIENT_H_
#include <cstdint>
#include <map>
#include <memory>
#include <optional>
#include <string>
#include <utility>
#include <vector>
#include "absl/base/thread_annotations.h"
#include "absl/container/flat_hash_map.h"
#include "absl/status/status.h"
#include "absl/status/statusor.h"
#include "absl/strings/string_view.h"
#include "absl/synchronization/mutex.h"
#include "absl/time/time.h"
#include "absl/types/span.h"
#include "mlir/IR/BuiltinOps.h"
#include "xla/backends/gpu/collectives/gpu_clique_key.h"
#include "xla/backends/gpu/collectives/gpu_cliques.h"
#include "xla/client/local_client.h"
#include "xla/executable_run_options.h"
#include "xla/future.h"
#include "xla/hlo/builder/xla_computation.h"
#include "xla/layout.h"
#include "xla/pjrt/buffer_sequencing_event.h"
#include "xla/pjrt/device_event.h"
#include "xla/pjrt/distributed/client.h"
#include "xla/pjrt/distributed/key_value_store_interface.h"
#include "xla/pjrt/gpu/se_gpu_topology_description.h"
#include "xla/pjrt/host_memory_allocator.h"
#include "xla/pjrt/local_device_state.h"
#include "xla/pjrt/maybe_owning_mlir_module.h"
#include "xla/pjrt/pjrt_client.h"
#include "xla/pjrt/pjrt_compiler.h"
#include "xla/pjrt/pjrt_executable.h"
#include "xla/pjrt/pjrt_stream_executor_client.h"
#include "xla/pjrt/plugin/xla_gpu/xla_gpu_client_options.h"
#include "xla/pjrt/raw_buffer.h"
#include "xla/pjrt/se_raw_buffer.h"
#include "xla/pjrt/tracked_device_buffer.h"
#include "xla/runtime/device_id.h"
#include "xla/service/computation_placer.h"
#include "xla/service/gpu/gpu_executable_run_options.h"
#include "xla/service/gpu_topology.h"
#include "xla/service/gpu_topology.pb.h"
#include "xla/shape.h"
#include "xla/shape_tree.h"
#include "xla/stream_executor/device_address_allocator.h"
#include "xla/stream_executor/device_description.h"
#include "xla/tsl/concurrency/async_value.h"
#include "xla/tsl/concurrency/ref_count.h"
#include "xla/tsl/framework/allocator.h"
#include "xla/tsl/protobuf/coordination_service.pb.h"
#include "xla/xla_data.pb.h"
#include "tsl/platform/numa.h"
namespace xla {
using DeviceTopologyPair =
std::pair<std::vector<std::unique_ptr<PjRtStreamExecutorDevice>>,
GpuTopologyProto>;
class StreamExecutorGpuDevice : public PjRtStreamExecutorDevice {
public:
StreamExecutorGpuDevice(int id,
std::unique_ptr<LocalDeviceState> local_device_state,
std::string device_kind, std::string device_vendor,
std::string compute_capability, int core_count,
int shared_memory_per_block_optin,
int local_device_id, int process_index,
int process_index_in_partition = 0,
int partition_index = 0,
int numa_node = tsl::port::kNUMANoAffinity);
absl::string_view device_vendor() const;
absl::StatusOr<tsl::AllocatorStats> GetAllocatorStats() const override;
absl::Span<int const> coords() const;
absl::StatusOr<PjRtMemorySpace*> default_memory_space() const override;
private:
std::string device_vendor_;
};
class StreamExecutorGpuHbmMemorySpace : public PjRtStreamExecutorMemorySpace {
public:
static constexpr absl::string_view kKind = "device";
static const int kKindId;
StreamExecutorGpuHbmMemorySpace(int id, PjRtDevice* device);
};
// A custom PjRtClient that overrides the device assignment method.
class StreamExecutorGpuClient : public xla::PjRtStreamExecutorClient {
public:
using xla::PjRtStreamExecutorClient::PjRtStreamExecutorClient;
StreamExecutorGpuClient(
std::string platform_name, LocalClient* client,
std::vector<std::unique_ptr<PjRtStreamExecutorDevice>> devices,
int process_index, std::unique_ptr<se::DeviceAddressAllocator> allocator,
std::unique_ptr<HostMemoryAllocator> host_memory_allocator,
bool should_stage_host_to_device_transfers,
std::unique_ptr<gpu::GpuExecutableRunOptions> gpu_run_options,
std::shared_ptr<KeyValueStoreInterface> kv_store,
bool abort_collectives_on_failure,
std::shared_ptr<const GpuTopology> gpu_topology,
std::optional<int> num_processes);
std::optional<std::shared_ptr<KeyValueStoreInterface>> key_value_store()
const override {
return kv_store_;
}
gpu::GpuExecutableRunOptions* gpu_run_options(
const ExecuteOptions& options) override;
absl::StatusOr<xla::DeviceAssignment> GetDefaultDeviceAssignment(
int num_replicas, int num_partitions) const override;
absl::string_view platform_version() const override;
std::optional<PjRtPluginAttributes> plugin_attributes() const override;
void UpdateGlobalProcessInfo(
absl::Span<tensorflow::CoordinatedTaskStateInfo> infos) override;
using PjRtStreamExecutorClient::CreateBuffersForAsyncHostToDevice;
absl::StatusOr<std::unique_ptr<PjRtClient::AsyncHostToDeviceTransferManager>>
CreateBuffersForAsyncHostToDevice(
absl::Span<const PjRtClient::ShapeSpec> shape_specs,
std::optional<absl::Span<const std::optional<Layout>>> device_layouts,
PjRtMemorySpace* memory_space) override;
// CrossHostSendBuffers and CrossHostReceiveBuffers are part of the new
// cross-host transfers API.
absl::StatusOr<std::vector<Future<>>> CrossHostSendBuffers(
absl::Span<PjRtBuffer* const> buffers,
absl::Span<const GlobalDeviceId> dst_global_device_ids,
std::vector<CrossHostTransferKey> transfer_keys) override;
absl::StatusOr<std::vector<std::unique_ptr<PjRtBuffer>>>
CrossHostReceiveBuffers(
xla::PjRtDevice* device, absl::Span<const xla::Shape> shapes,
absl::Span<const GlobalDeviceId> src_global_device_ids,
std::vector<CrossHostTransferKey> transfer_keys) override;
// ScheduleRemoteSend and MakeCrossHostReceiveBuffers are methods implemented
// to support the legacy cross-host transfers API.
void ScheduleRemoteSend(
PjRtMemorySpace* memory_space,
tsl::RCReference<CommonPjRtRawBuffer> raw_buffer,
std::vector<tsl::RCReference<tsl::AsyncValue>> definition_events,
tsl::RCReference<PjRtDeviceEventPromise> usage_event_promise,
Future<std::string> serialized_descriptor,
PjRtBuffer::RemoteSendCallback on_done) override;
absl::StatusOr<std::vector<std::unique_ptr<PjRtBuffer>>>
MakeCrossHostReceiveBuffers(absl::Span<const Shape> shapes,
PjRtDevice* device,
PjRtCrossHostRecvNotifier notifier) override;
absl::StatusOr<const xla::PjRtTopologyDescription*> GetTopologyDescription()
const override;
absl::StatusOr<Layout> GetDefaultLayout(
PrimitiveType element_type, absl::Span<const int64_t> dims) override;
absl::StatusOr<std::unique_ptr<PjRtLoadedExecutable>> LoadSerialized(
absl::string_view serialized, std::optional<CompileOptions> options,
const LoadOptions& load_options);
absl::StatusOr<std::unique_ptr<PjRtLoadedExecutable>> CompileAndLoad(
const XlaComputation& computation, CompileOptions options) override;
absl::StatusOr<std::unique_ptr<PjRtLoadedExecutable>> CompileAndLoad(
mlir::ModuleOp module, CompileOptions options) override {
return CompileAndLoad(MaybeOwningMlirModule(std::move(module)), options);
}
absl::StatusOr<std::unique_ptr<PjRtLoadedExecutable>> CompileAndLoad(
MaybeOwningMlirModule module, CompileOptions options) override;
absl::StatusOr<PjRtStreamExecutorExecutionOutput> RunAsync(
LocalExecutable& exec, PjRtDevice* device,
absl::Span<const tsl::RCReference<CommonPjRtRawBuffer>> flat_arguments,
absl::Span<const tsl::RCReference<CommonPjRtRawBuffer>> results,
ExecutableRunOptions run_options_inp, bool parameter_is_tupled_arguments,
absl::Span<const Shape> executable_parameter_shapes) override;
absl::Status UpdateCompileOptionsInternal(
CompileOptions* options, ExecutableExtras* returned_extras,
bool lookup_addressable_devices) override;
HostMemoryAllocator* GetHostMemoryAllocator() const override {
return host_memory_allocator_.get();
}
private:
absl::flat_hash_map<GlobalDeviceId, IncarnationId> GetLatestIncarnations(
const ExecuteOptions& options);
std::optional<int> num_nodes_;
const bool abort_collectives_on_failure_ = false;
std::optional<xla::StreamExecutorGpuTopologyDescription> topology_;
std::shared_ptr<KeyValueStoreInterface> kv_store_;
// Helpers for cross host transfers.
absl::Duration cross_host_transfer_timeout_ = absl::Minutes(3);
void ScheduleSendsOnLocalDevice(
PjRtDevice* device, std::vector<PjRtBuffer*> buffers,
std::vector<GlobalDeviceId> dst_global_device_ids,
std::vector<CrossHostTransferKey> transfer_keys,
std::vector<std::shared_ptr<Promise<>>> promises);
struct PrepareReceiveBufferResult {
std::unique_ptr<PjRtBuffer> buffer;
tsl::RCReference<CommonPjRtRawBuffer> raw_buffer;
LocalDeviceState* local_device;
se::Stream* stream;
BufferSequencingEventRef definition_event;
};
absl::StatusOr<PrepareReceiveBufferResult> PrepareReceiveBuffer(
PjRtDevice* device, Shape shape);
};
std::vector<std::unique_ptr<PjRtStreamExecutorDevice>> BuildLocalDevices(
std::map<int, std::unique_ptr<LocalDeviceState>> local_device_states,
int process_id);
std::string MakeComputeCapabilityString(const se::DeviceDescription* desc);
absl::StatusOr<DeviceTopologyPair> BuildDistributedDevices(
absl::string_view platform_name,
std::map<int, std::unique_ptr<LocalDeviceState>> local_device_states,
int process_id, int num_nodes,
gpu::GpuExecutableRunOptions* gpu_executable_run_options,
std::shared_ptr<KeyValueStoreInterface> kv_store, bool enable_mock_nccl,
std::optional<absl::string_view> mock_gpu_topology = std::nullopt,
std::optional<int> partition_index = std::nullopt,
absl::Duration get_local_topology_timeout = absl::Minutes(2),
absl::Duration get_global_topology_timeout = absl::Minutes(5));
absl::StatusOr<std::unique_ptr<PjRtClient>> GetStreamExecutorGpuClient(
const GpuClientOptions& options);
// Get the fabric info of a local device ordinal in the format of
// "clusterUuid/cliqueId". Empty on SM90 or lower.
absl::StatusOr<std::string> GetDeviceFabricInfo(int device_ordinal);
} // namespace xla
#endif // XLA_PJRT_GPU_SE_GPU_PJRT_CLIENT_H_