/usr/local/lib64/python3.6/site-packages/pyarrow/include/arrow/util
NameSizeModeActions
algorithm.h12290644editdlrm
aligned_storage.h43020644editdlrm
align_util.h26360644editdlrm
async_generator.h640880644editdlrm
async_util.h96330644editdlrm
atomic_shared_ptr.h36400644editdlrm
base64.h10980644editdlrm
basic_decimal.h203340644editdlrm
benchmark_util.h45840644editdlrm
bitmap.h174630644editdlrm
bitmap_builders.h15630644editdlrm
bitmap_generate.h35630644editdlrm
bitmap_ops.h90840644editdlrm
bitmap_reader.h83470644editdlrm
bitmap_visit.h34600644editdlrm
bitmap_writer.h93600644editdlrm
bitset_stack.h27890644editdlrm
bit_block_counter.h191410644editdlrm
bit_run_reader.h165990644editdlrm
bit_stream_utils.h169860644editdlrm
bit_util.h115680644editdlrm
bpacking.h11750644editdlrm
bpacking64_default.h1959340644editdlrm
bpacking_avx2.h10090644editdlrm
bpacking_avx512.h10110644editdlrm
bpacking_default.h1032320644editdlrm
bpacking_neon.h10090644editdlrm
bpacking_simd128_generated.h985290644editdlrm
bpacking_simd256_generated.h774750644editdlrm
bpacking_simd512_generated.h670810644editdlrm
byte_stream_split.h287840644editdlrm
cancel.h29110644editdlrm
checked_cast.h20760644editdlrm
compare.h19810644editdlrm
compression.h73670644editdlrm
concurrent_map.h17750644editdlrm
config.h16590644editdlrm
converter.h146570644editdlrm
counting_semaphore.h22510644editdlrm
cpu_info.h47240644editdlrm
decimal.h116870644editdlrm
delimiting.h73350644editdlrm
dispatch.h32350644editdlrm
double_conversion.h11950644editdlrm
endian.h80710644editdlrm
formatting.h206120644editdlrm
functional.h56120644editdlrm
future.h361680644editdlrm
future_iterator.h25170644editdlrm
hashing.h306330644editdlrm
hash_util.h19140644editdlrm
int_util.h42840644editdlrm
io_util.h103260644editdlrm
iterator.h181230644editdlrm
key_value_metadata.h35790644editdlrm
launder.h10510644editdlrm
logging.h93380644editdlrm
macros.h73490644editdlrm
make_unique.h14750644editdlrm
map.h24760644editdlrm
math_constants.h11060644editdlrm
memory.h15660644editdlrm
mutex.h18330644editdlrm
optional.h11740644editdlrm
parallel.h36160644editdlrm
pcg_random.h11460644editdlrm
print.h17250644editdlrm
queue.h10170644editdlrm
range.h48340644editdlrm
rle_encoding.h310290644editdlrm
simd.h13330644editdlrm
small_vector.h146600644editdlrm
sort.h24660644editdlrm
spaced.h35670644editdlrm
stopwatch.h14010644editdlrm
string.h25700644editdlrm
string_builder.h24460644editdlrm
string_view.h12690644editdlrm
task_group.h43620644editdlrm
tdigest.h30520644editdlrm
test_common.h28370644editdlrm
thread_pool.h153380644editdlrm
time.h29880644editdlrm
trie.h71570644editdlrm
type_fwd.h14090644editdlrm
type_traits.h28940644editdlrm
ubsan.h27770644editdlrm
unreachable.h9260644editdlrm
uri.h32970644editdlrm
utf8.h187800644editdlrm
value_parsing.h270560644editdlrm
variant.h137280644editdlrm
vector.h56650644editdlrm
visibility.h14630644editdlrm
windows_compatibility.h12600644editdlrm
windows_fixup.h13790644editdlrm
Edit: /usr/local/lib64/python3.6/site-packages/pyarrow/include/arrow/util/async_util.h (9633B)
// Licensed to the Apache Software Foundation (ASF) under one // or more contributor license agreements. See the NOTICE file // distributed with this work for additional information // regarding copyright ownership. The ASF licenses this file // to you 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. #pragma once #include #include "arrow/result.h" #include "arrow/status.h" #include "arrow/util/future.h" #include "arrow/util/mutex.h" namespace arrow { namespace util { /// Custom deleter for AsyncDestroyable objects template struct DestroyingDeleter { void operator()(T* p) { if (p) { p->Destroy(); } } }; /// An object which should be asynchronously closed before it is destroyed /// /// Classes can extend this to ensure that the close method is called and completed /// before the instance is deleted. This provides smart_ptr / delete semantics for /// objects with an asynchronous destructor. /// /// Classes which extend this must be constructed using MakeSharedAsync or MakeUniqueAsync class ARROW_EXPORT AsyncDestroyable { public: AsyncDestroyable(); virtual ~AsyncDestroyable(); /// A future which will complete when the AsyncDestroyable has finished and is ready /// to be deleted. /// /// This can be used to ensure all work done by this object has been completed before /// proceeding. Future<> on_closed() { return on_closed_; } protected: /// Subclasses should override this and perform any cleanup. Once the future returned /// by this method finishes then this object is eligible for destruction and any /// reference to `this` will be invalid virtual Future<> DoDestroy() = 0; private: void Destroy(); Future<> on_closed_; #ifndef NDEBUG bool constructed_correctly_ = false; #endif template friend struct DestroyingDeleter; template friend std::shared_ptr MakeSharedAsync(Args&&... args); template friend std::unique_ptr> MakeUniqueAsync(Args&&... args); }; template std::shared_ptr MakeSharedAsync(Args&&... args) { static_assert(std::is_base_of::value, "Nursery::MakeSharedCloseable only works with AsyncDestroyable types"); std::shared_ptr ptr(new T(std::forward(args)...), DestroyingDeleter()); #ifndef NDEBUG ptr->constructed_correctly_ = true; #endif return ptr; } template std::unique_ptr> MakeUniqueAsync(Args&&... args) { static_assert(std::is_base_of::value, "Nursery::MakeUniqueCloseable only works with AsyncDestroyable types"); std::unique_ptr> ptr(new T(std::forward(args)...), DestroyingDeleter()); #ifndef NDEBUG ptr->constructed_correctly_ = true; #endif return ptr; } /// A utility which keeps track of a collection of asynchronous tasks /// /// This can be used to provide structured concurrency for asynchronous development. /// A task group created at a high level can be distributed amongst low level components /// which register work to be completed. The high level job can then wait for all work /// to be completed before cleaning up. class ARROW_EXPORT AsyncTaskGroup { public: /// Add a task to be tracked by this task group /// /// If a previous task has failed then adding a task will fail /// /// If WaitForTasksToFinish has been called and the returned future has been marked /// completed then adding a task will fail. Status AddTask(std::function>()> task); /// Add a task that has already been started Status AddTask(const Future<>& task); /// Signal that top level tasks are done being added /// /// It is allowed for tasks to be added after this call provided the future has not yet /// completed. This should be safe as long as the tasks being added are added as part /// of a task that is tracked. As soon as the count of running tasks reaches 0 this /// future will be marked complete. /// /// Any attempt to add a task after the returned future has completed will fail. /// /// The returned future that will finish when all running tasks have finsihed. Future<> End(); /// A future that will be finished after End is called and all tasks have completed /// /// This is the same future that is returned by End() but calling this method does /// not indicate that top level tasks are done being added. End() must still be called /// at some point or the future returned will never finish. /// /// This is a utility method for workflows where the finish future needs to be /// referenced before all top level tasks have been queued. Future<> OnFinished() const; private: Status AddTaskUnlocked(const Future<>& task, util::Mutex::Guard guard); bool finished_adding_ = false; int running_tasks_ = 0; Status err_; Future<> all_tasks_done_ = Future<>::Make(); util::Mutex mutex_; }; /// A task group which serializes asynchronous tasks in a push-based workflow /// /// Tasks will be executed in the order they are added /// /// This will buffer results in an unlimited fashion so it should be combined /// with some kind of backpressure class ARROW_EXPORT SerializedAsyncTaskGroup { public: SerializedAsyncTaskGroup(); /// Push an item into the serializer and (eventually) into the consumer /// /// The item will not be delivered to the consumer until all previous items have been /// consumed. /// /// If the consumer returns an error then this serializer will go into an error state /// and all subsequent pushes will fail with that error. Pushes that have been queued /// but not delivered will be silently dropped. /// /// \return True if the item was pushed immediately to the consumer, false if it was /// queued Status AddTask(std::function>()> task); /// Signal that all top level tasks have been added /// /// The returned future that will finish when all tasks have been consumed. Future<> End(); /// A future that finishes when all queued items have been delivered. /// /// This will return the same future returned by End but will not signal /// that all tasks have been finished. End must be called at some point in order for /// this future to finish. Future<> OnFinished() const { return on_finished_; } private: void ConsumeAsMuchAsPossibleUnlocked(util::Mutex::Guard&& guard); bool TryDrainUnlocked(); Future<> on_finished_; std::queue>()>> tasks_; util::Mutex mutex_; bool ended_ = false; Status err_; Future<> processing_; }; class ARROW_EXPORT AsyncToggle { public: /// Get a future that will complete when the toggle next becomes open /// /// If the toggle is open this returns immediately /// If the toggle is closed this future will be unfinished until the next call to Open Future<> WhenOpen(); /// \brief Close the toggle /// /// After this call any call to WhenOpen will be delayed until the next open void Close(); /// \brief Open the toggle /// /// Note: This call may complete a future, triggering any callbacks, and generally /// should not be done while holding any locks. /// /// Note: If Open is called from multiple threads it could lead to a situation where /// callbacks from the second open finish before callbacks on the first open. /// /// All current waiters will be released to enter, even if another close call /// quickly follows void Open(); /// \brief Return true if the toggle is currently open bool IsOpen(); private: Future<> when_open_ = Future<>::MakeFinished(); bool closed_ = false; util::Mutex mutex_; }; /// \brief Options to control backpressure behavior struct ARROW_EXPORT BackpressureOptions { /// \brief Create default options that perform no backpressure BackpressureOptions() : toggle(NULLPTR), resume_if_below(0), pause_if_above(0) {} /// \brief Create options that will perform backpressure /// /// \param toggle A toggle to be shared between the producer and consumer /// \param resume_if_below The producer should resume producing if the backpressure /// queue has fewer than resume_if_below items. /// \param pause_if_above The producer should pause producing if the backpressure /// queue has more than pause_if_above items BackpressureOptions(std::shared_ptr toggle, uint32_t resume_if_below, uint32_t pause_if_above) : toggle(std::move(toggle)), resume_if_below(resume_if_below), pause_if_above(pause_if_above) {} static BackpressureOptions Make(uint32_t resume_if_below = 32, uint32_t pause_if_above = 64); static BackpressureOptions NoBackpressure(); std::shared_ptr toggle; uint32_t resume_if_below; uint32_t pause_if_above; }; } // namespace util } // namespace arrow