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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
238 changes: 205 additions & 33 deletions src/substruct/substruct_library.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -9,9 +9,12 @@
#include <omp.h>

#include <algorithm>
#include <condition_variable>
#include <cstdint>
#include <limits>
#include <mutex>
#include <semaphore>
#include <shared_mutex>
#include <stdexcept>
#include <unordered_set>
#include <utility>
Expand Down Expand Up @@ -63,6 +66,65 @@ struct GpuTargets {
std::vector<const RDKit::ROMol*> molecules;
std::vector<unsigned int> ids;
std::shared_ptr<const PersistentDeviceTargets> searchTargets;

//! Approximate device memory held by the uploaded batch.
[[nodiscard]] std::size_t deviceBytes() const {
return host->batchAtomStarts.size() * sizeof(int) + host->atomDataPacked.size() * sizeof(AtomDataPacked) +
host->bondTypeCounts.size() * sizeof(BondTypeCounts) +
host->targetAtomBonds.size() * sizeof(TargetAtomBonds);
}
};

//! A search workspace and the per-target match flags its queries fill.
struct QuerySlot {
std::shared_ptr<SubstructSearchWorkspace> workspace;
std::vector<std::uint8_t> matchFlags;
};

//! One GPU's query slots; a query waits for a free slot.
class SlotPool {
public:
explicit SlotPool(std::vector<std::unique_ptr<QuerySlot>> slots) : slots_(std::move(slots)) {
for (auto& slot : slots_) {
free_.push_back(slot.get());
}
}

[[nodiscard]] QuerySlot* take() {
std::unique_lock lock(mutex_);
available_.wait(lock, [this] { return !free_.empty(); });
QuerySlot* slot = free_.back();
free_.pop_back();
return slot;
}

void give(QuerySlot* slot) {
{
const std::lock_guard lock(mutex_);
free_.push_back(slot);
}
available_.notify_one();
}

private:
std::vector<std::unique_ptr<QuerySlot>> slots_;
std::mutex mutex_;
std::condition_variable available_;
std::vector<QuerySlot*> free_;
};

//! Holds a slot taken from a pool for the duration of one GPU search.
class TakenSlot {
public:
explicit TakenSlot(SlotPool& pool) : pool_(pool), slot_(pool.take()) {}
~TakenSlot() { pool_.give(slot_); }
TakenSlot(const TakenSlot&) = delete;
TakenSlot& operator=(const TakenSlot&) = delete;
[[nodiscard]] QuerySlot& operator*() const { return *slot_; }

private:
SlotPool& pool_;
QuerySlot* slot_;
};

} // namespace
Expand Down Expand Up @@ -90,21 +152,18 @@ class SubstructLibrary::Impl {
gpus_.size(),
[&](std::size_t index) { return gpus_[index] != nullptr; },
[&](std::size_t index) { gpus_[index].reset(); });
releaseOnDevices(
workspaces_.size(),
[&](std::size_t index) { return workspaces_[index] != nullptr; },
[&](std::size_t index) { workspaces_[index].reset(); });
destroySlotPools();
}

unsigned int addMol(const RDKit::ROMol& molecule) {
const std::lock_guard lock(mutex_);
const auto lock = lockForWriting();
requireIdsAvailable(1);
molecules_.push_back(copyWithRingInfo(molecule));
return static_cast<unsigned int>(molecules_.size() - 1);
}

std::vector<unsigned int> addMols(const std::vector<const RDKit::ROMol*>& molecules) {
const std::lock_guard lock(mutex_);
const auto lock = lockForWriting();
requireIdsAvailable(molecules.size());
const std::size_t firstId = molecules_.size();
molecules_.resize(firstId + molecules.size());
Expand Down Expand Up @@ -135,7 +194,7 @@ class SubstructLibrary::Impl {
}

void finalize(cudaStream_t stream) {
const std::lock_guard lock(mutex_);
const auto lock = lockForWriting();
useCurrentDeviceIfNoneConfigured();
if (deviceIds_.size() > 1 && stream != nullptr) {
throw std::invalid_argument("A single external CUDA stream cannot be used with a multi-GPU substructure library");
Expand Down Expand Up @@ -208,16 +267,29 @@ class SubstructLibrary::Impl {
exceptions.store(std::current_exception());
}
}
bool slotsRebuilt = false;
try {
exceptions.rethrow();
createWorkspaces();
rdkitIds_.reserve(rdkitIds_.size() + newRdkitIds.size());
// Query slots are sized against memory that already holds the new batches. The batches they replace are
// freed below, so their memory counts as available.
std::vector<std::size_t> freedBytes(deviceIds_.size(), 0);
for (std::size_t gpu = 0; gpu < deviceIds_.size(); ++gpu) {
if (replacements[gpu] != nullptr && gpus_[gpu] != nullptr) {
freedBytes[gpu] = gpus_[gpu]->deviceBytes();
}
}
slotsRebuilt = true;
createSlotPools(freedBytes);
} catch (...) {
static_cast<void>(cudaGetLastError());
releaseOnDevices(
replacements.size(),
[&](std::size_t index) { return replacements[index] != nullptr; },
[&](std::size_t index) { replacements[index].reset(); });
if (slotsRebuilt && finalized_) {
restoreSlotPools();
}
throw;
}

Expand All @@ -231,19 +303,24 @@ class SubstructLibrary::Impl {
}

[[nodiscard]] std::size_t size() const {
const std::lock_guard lock(mutex_);
const auto lock = lockForQuery();
return numFinalized_;
}

[[nodiscard]] std::size_t pendingSize() const {
const std::lock_guard lock(mutex_);
const auto lock = lockForQuery();
return molecules_.size() - numFinalized_;
}

[[nodiscard]] std::size_t maxConcurrentQueries() const {
const auto lock = lockForQuery();
return maxConcurrentQueries_;
}

[[nodiscard]] std::vector<unsigned int> getMatches(const RDKit::ROMol& query,
int maxResults,
cudaStream_t stream) const {
const std::lock_guard lock(mutex_);
const auto lock = lockForQuery();
requireFinalized();
if (maxResults == 0) {
return {};
Expand All @@ -252,13 +329,13 @@ class SubstructLibrary::Impl {
}

[[nodiscard]] std::size_t countMatches(const RDKit::ROMol& query, cudaStream_t stream) const {
const std::lock_guard lock(mutex_);
const auto lock = lockForQuery();
requireFinalized();
return matches(query, stream, -1).size();
}

[[nodiscard]] bool hasMatch(const RDKit::ROMol& query, cudaStream_t stream) const {
const std::lock_guard lock(mutex_);
const auto lock = lockForQuery();
requireFinalized();
return !matches(query, stream, 1).empty();
}
Expand All @@ -274,6 +351,21 @@ class SubstructLibrary::Impl {
if (!finalized_) {
throw std::logic_error("Substructure library must be finalized before querying");
}
if (slotPools_.empty()) {
throw std::runtime_error("Substructure library has no query slots after a failed finalize(); retry finalize()");
}
}

// Writers hold writerGate_ while waiting for exclusive access and queries pass through it first, so a waiting
// writer is not starved by queries that keep arriving.
[[nodiscard]] std::unique_lock<std::shared_mutex> lockForWriting() {
const std::lock_guard gate(writerGate_);
return std::unique_lock(mutex_);
}

[[nodiscard]] std::shared_lock<std::shared_mutex> lockForQuery() const {
{ const std::lock_guard gate(writerGate_); }
return std::shared_lock(mutex_);
}

[[nodiscard]] int threads() const {
Expand Down Expand Up @@ -334,17 +426,77 @@ class SubstructLibrary::Impl {
return next;
}

// Search state for each GPU, created on the first finalize() and reused by every query.
void createWorkspaces() {
if (workspaces_.size() == deviceIds_.size()) {
return;
// Give each GPU as many query slots as fit below 85% of its memory, counting freedBytes[gpu] as free. Queries
// with recursive SMARTS also need scratch memory, so fewer of them may run at once.
void createSlotPools(const std::vector<std::size_t>& freedBytes) {
destroySlotPools();
constexpr std::size_t memoryPercent = 85;
std::size_t slots = std::numeric_limits<std::size_t>::max();
std::vector<std::size_t> availableBytes(deviceIds_.size());
std::vector<std::size_t> slotBytes(deviceIds_.size());
std::vector<std::size_t> recursiveBytes(deviceIds_.size());
for (std::size_t gpu = 0; gpu < deviceIds_.size(); ++gpu) {
const WithDevice device(deviceIds_[gpu]);
std::size_t freeBytes = 0;
std::size_t totalBytes = 0;
cudaCheckError(cudaMemGetInfo(&freeBytes, &totalBytes));
const std::size_t usedBytes = totalBytes - std::min(totalBytes, freeBytes + freedBytes[gpu]);
const std::size_t budgetBytes = totalBytes * memoryPercent / 100;
availableBytes[gpu] = budgetBytes > usedBytes ? budgetBytes - usedBytes : 0;
slotBytes[gpu] = estimateSubstructSearchWorkspaceBytes(gpuConfig(gpu));
recursiveBytes[gpu] = estimateRecursiveScratchBytes(gpuConfig(gpu));
if (availableBytes[gpu] < slotBytes[gpu] + recursiveBytes[gpu]) {
throw std::runtime_error("Substructure library cannot fit one query below 85% of GPU memory");
}
slots = std::min(slots, (availableBytes[gpu] - recursiveBytes[gpu]) / slotBytes[gpu]);
}
std::vector<std::shared_ptr<SubstructSearchWorkspace>> created(deviceIds_.size());
std::size_t recursiveSlots = slots;
for (std::size_t gpu = 0; gpu < deviceIds_.size(); ++gpu) {
created[gpu] = makeSubstructSearchWorkspace(deviceIds_[gpu]);
recursiveSlots = std::min(recursiveSlots, (availableBytes[gpu] - slots * slotBytes[gpu]) / recursiveBytes[gpu]);
}

std::vector<std::unique_ptr<SlotPool>> pools(deviceIds_.size());
try {
for (std::size_t gpu = 0; gpu < deviceIds_.size(); ++gpu) {
const WithDevice device(deviceIds_[gpu]);
std::vector<std::unique_ptr<QuerySlot>> gpuSlots(slots);
for (auto& slot : gpuSlots) {
slot = std::make_unique<QuerySlot>();
slot->workspace = makeSubstructSearchWorkspace(deviceIds_[gpu]);
}
pools[gpu] = std::make_unique<SlotPool>(std::move(gpuSlots));
}
} catch (...) {
releaseOnDevices(
pools.size(),
[&](std::size_t gpu) { return pools[gpu] != nullptr; },
[&](std::size_t gpu) { pools[gpu].reset(); });
throw;
}
workspaces_ = std::move(created);
matchFlags_.resize(deviceIds_.size());
slotPools_ = std::move(pools);
recursiveQuerySlots_ = std::make_unique<std::counting_semaphore<>>(
static_cast<std::ptrdiff_t>(std::max<std::size_t>(1, recursiveSlots)));
maxConcurrentQueries_ = slots;
}

// After a failed finalize(), give the previously finalized molecules their query slots back. If that fails too,
// queries are refused until finalize() succeeds.
void restoreSlotPools() noexcept {
try {
createSlotPools(std::vector<std::size_t>(deviceIds_.size(), 0));
} catch (...) {
destroySlotPools();
}
}

void destroySlotPools() noexcept {
releaseOnDevices(
slotPools_.size(),
[&](std::size_t gpu) { return slotPools_[gpu] != nullptr; },
[&](std::size_t gpu) { slotPools_[gpu].reset(); });
slotPools_.clear();
recursiveQuerySlots_.reset();
maxConcurrentQueries_ = 0;
}

[[nodiscard]] SubstructSearchConfig gpuConfig(std::size_t gpu) const {
Expand All @@ -367,6 +519,19 @@ class SubstructLibrary::Impl {
if (deviceIds_.size() > 1 && stream != nullptr) {
throw std::invalid_argument("A single external CUDA stream cannot be used with a multi-GPU substructure library");
}
// Recursive SMARTS hold extra scratch memory on every GPU while they run.
struct RecursiveSlot {
std::counting_semaphore<>* slots = nullptr;
~RecursiveSlot() {
if (slots != nullptr) {
slots->release();
}
}
} recursiveSlot;
if (hasRecursiveSmarts(&query)) {
recursiveQuerySlots_->acquire();
recursiveSlot.slots = recursiveQuerySlots_.get();
}
std::vector<std::vector<unsigned int>> perGpu(deviceIds_.size());
detail::OpenMPExceptionRegistry exceptions;
#pragma omp parallel for num_threads(static_cast<int>(deviceIds_.size())) schedule(static)
Expand Down Expand Up @@ -417,9 +582,10 @@ class SubstructLibrary::Impl {
cudaStream_t stream) const {
ScopedNvtxRange searchRange("SubstructLibrary GPU search");
const GpuTargets& targets = *gpus_[gpu];
std::vector<std::uint8_t>& flags = matchFlags_[gpu];
const SubstructSearchConfig config = gpuConfig(gpu);
hasSubstructMatch(*targets.searchTargets, query, flags, config.algorithm, stream, config, workspaces_[gpu].get());
const TakenSlot slot(*slotPools_[gpu]);
std::vector<std::uint8_t>& flags = (*slot).matchFlags;
const SubstructSearchConfig config = gpuConfig(gpu);
hasSubstructMatch(*targets.searchTargets, query, flags, config.algorithm, stream, config, (*slot).workspace.get());
std::vector<unsigned int> result;
for (std::size_t target = 0; target < flags.size(); ++target) {
if (flags[target] != 0) {
Expand Down Expand Up @@ -451,22 +617,24 @@ class SubstructLibrary::Impl {
}
}

SubstructSearchConfig config_;
std::vector<int> deviceIds_;
// Held by every call, so queries, additions, and finalize() run one at a time.
mutable std::mutex mutex_;
SubstructSearchConfig config_;
std::vector<int> deviceIds_;
// Queries share mutex_; additions and finalize() hold it exclusively.
mutable std::shared_mutex mutex_;
mutable std::mutex writerGate_;

//! Every added molecule, indexed by ID; the first numFinalized_ are searchable.
std::vector<std::unique_ptr<RDKit::ROMol>> molecules_;
std::size_t numFinalized_ = 0;
bool finalized_ = false;

std::vector<std::unique_ptr<GpuTargets>> gpus_;
std::vector<std::unique_ptr<GpuTargets>> gpus_;
//! Finalized molecules the GPU format cannot represent, ascending; matched with RDKit.
std::vector<unsigned int> rdkitIds_;
std::vector<std::shared_ptr<SubstructSearchWorkspace>> workspaces_;
//! Per-GPU match flags, reused across queries to avoid reallocating them.
mutable std::vector<std::vector<std::uint8_t>> matchFlags_;
std::vector<unsigned int> rdkitIds_;
mutable std::vector<std::unique_ptr<SlotPool>> slotPools_;
//! Bounds how many queries with recursive SMARTS hold scratch memory at once.
std::unique_ptr<std::counting_semaphore<>> recursiveQuerySlots_;
std::size_t maxConcurrentQueries_ = 0;
};

SubstructLibrary::SubstructLibrary(SubstructSearchConfig config) : impl_(std::make_unique<Impl>(std::move(config))) {}
Expand All @@ -493,6 +661,10 @@ std::size_t SubstructLibrary::pendingSize() const {
return impl_->pendingSize();
}

std::size_t SubstructLibrary::maxConcurrentQueries() const {
return impl_->maxConcurrentQueries();
}

std::vector<unsigned int> SubstructLibrary::getMatches(const RDKit::ROMol& query,
int maxResults,
cudaStream_t stream) const {
Expand Down
6 changes: 5 additions & 1 deletion src/substruct/substruct_library.h
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,8 @@ namespace nvMolKit {
* configured, molecules are split across those GPUs and results are merged in
* ID order. Molecules the GPU format cannot represent are matched with RDKit.
*
* Calls are serialized: queries, additions, and finalize() run one at a time.
* Queries may run concurrently, as many as fit in GPU memory; additions and
* finalize() wait for running queries and hold off new ones.
*/
class SubstructLibrary {
public:
Expand Down Expand Up @@ -59,6 +60,9 @@ class SubstructLibrary {
/** Number of molecules added since the last finalize(). */
[[nodiscard]] std::size_t pendingSize() const;

/** How many queries can run at once, set by finalize() from the GPU memory left after the molecules. */
[[nodiscard]] std::size_t maxConcurrentQueries() const;

/** IDs of matching molecules, ascending. maxResults limits the count; -1 returns all, 0 returns none. */
[[nodiscard]] std::vector<unsigned int> getMatches(const RDKit::ROMol& query,
int maxResults = -1,
Expand Down
Loading
Loading