Utilities¶
Namespace: dftracer::utils::utilities
-
template<typename I, typename Batch, typename ...Tags>
class StreamingUtility : public dftracer::utils::utilities::UtilityBase<I, Tags...>¶ Streaming utility: process() yields batches via AsyncGenerator<Batch>.
Use when a single input produces an unbounded or large sequence of output batches that should be consumed lazily rather than materialized all at once.
- Template Parameters:
I – Input type
Batch – Element type yielded per iteration
Tags – Variadic tag types for opt-in features
Public Functions
-
inline StreamingUtility()¶
-
template<typename Dummy = void, typename = std::enable_if_t<(sizeof...(Tags) > 0) && std::is_void_v<Dummy>>>
inline explicit StreamingUtility(Tags... tags)¶
-
virtual coro::AsyncGenerator<Batch> process(const I &input) = 0¶
-
inline void bind_context(CoroScope &ctx)¶
Bind context for streaming utilities with NeedsContext tag. Unlike Utility::process which is wrapped by CoroScope::spawn, streaming utilities need explicit context binding since their AsyncGenerator cannot be spawned directly.
-
inline void unbind_context()¶
Unbind context after streaming completes.
-
template<typename I, typename O, typename ...Tags>
class Utility : public dftracer::utils::utilities::UtilityBase<I, Tags...>¶ Materialized utility: process() returns CoroTask<O>.
- Template Parameters:
I – Input type
O – Output type
Tags – Variadic tag types
Public Functions
-
inline Utility()¶
Public Static Functions
-
static inline constexpr std::string_view get_type_signature()¶
-
static inline constexpr std::string_view get_name()¶
Friends
- friend class behaviors::UtilityExecutor< I, O, Tags... >
-
template<typename I, typename O, typename ...Tags>
class UtilityAdapter¶ Adapter that wraps a Utility as a Task.
Provides the fluent
use(utility).as_task()API: it returns a std::shared_ptr<Task> usable with the standard Task API (depends_on, etc.). Execution goes through UtilityExecutor, which injects the CoroScope context for context-needing utilities and applies env-gated monitoring.Usage:
auto task = use(utility).as_task(); task->depends_on(parent_task);
- Template Parameters:
I – Input type
O – Output type
Tags – Variadic tag types
-
template<typename I, typename ...Tags>
class UtilityBase¶ Shared machinery for all utility variants.
Holds the context pointer; tags are compile-time markers (queried via has_tag<>()), not stored per instance. Type signature is generated at compile time and stored as a static constexpr string_view.
- Template Parameters:
I – Input type
Tags – Variadic tag types for opt-in features
Subclassed by dftracer::utils::utilities::StreamingUtility< ViewReaderInput, ViewReaderBatch >, dftracer::utils::utilities::StreamingUtility< AggregatorInput, AggregationBatch, tags::NeedsContext >, dftracer::utils::utilities::Utility< ComparisonUtilityInput, Result< ComparisonUtilityOutput > >, dftracer::utils::utilities::Utility< PerfettoTraceWriterInput, PerfettoTraceWriterOutput, utilities::tags::NeedsContext >, dftracer::utils::utilities::Utility< ChunkAggregatorInput, ChunkAggregationOutput >, dftracer::utils::utilities::Utility< std::string, Hash >, dftracer::utils::utilities::Utility< ChunkDetailScanInput, Result< ChunkDetailScanOutput > >, dftracer::utils::utilities::Utility< FileChunkMapperInput, FileChunkMapperOutput >, dftracer::utils::utilities::Utility< PatternDirectoryScannerUtilityInput, std::vector< FileEntry >, utilities::tags::NeedsContext >, dftracer::utils::utilities::Utility< ChunkPrunerInput, ChunkPrunerOutput >, dftracer::utils::utilities::Utility< ChunkIndexerInput, ChunkIndexerOutput >, dftracer::utils::utilities::Utility< ChunkManifestMapperUtilityInput, ChunkManifestMapperUtilityOutput >, dftracer::utils::utilities::Utility< FilterableLine, std::optional< Line > >, dftracer::utils::utilities::Utility< EventCollectorFromMetadataCollectorUtilityInput, EventCollectorUtilityOutput >, dftracer::utils::utilities::Utility< LineBatchProcessUtilityInput, LineBatchProcessUtilityOutput< LineOutput > >, dftracer::utils::utilities::Utility< DirectoryProcessInput, BatchFileProcessOutput< FileOutput >, utilities::tags::NeedsContext >, dftracer::utils::utilities::Utility< ResolverInput, ResolverResult, utilities::tags::NeedsContext >, dftracer::utils::utilities::Utility< FileMergeValidatorUtilityInput, FileMergeValidatorUtilityOutput >, dftracer::utils::utilities::Utility< Text, Lines >, dftracer::utils::utilities::Utility< StatisticsAggregatorInput, TraceStatistics >, dftracer::utils::utilities::Utility< EventIdExtractionInput, EventIdExtractionOutput >, dftracer::utils::utilities::Utility< FileDecompressionUtilityInput, Result< FileDecompressionUtilityOutput > >, dftracer::utils::utilities::Utility< StatisticsQueryInput, StatisticsQueryOutput >, dftracer::utils::utilities::Utility< SimpleLineBatchProcessUtilityInput, SimpleLineBatchProcessUtilityOutput< LineOutput > >, dftracer::utils::utilities::Utility< FileCompressionUtilityInput, Result< FileCompressionUtilityOutput > >, dftracer::utils::utilities::Utility< IndexBuildConfig, IndexBuildResult, tags::NeedsContext >, dftracer::utils::utilities::Utility< ReconstructorInput, ReconstructorResult, utilities::tags::NeedsContext >, dftracer::utils::utilities::Utility< ViewBuilderInput, Result< ViewBuilderOutput > >, dftracer::utils::utilities::Utility< ReorganizationPlannerInput, ExtractionPlan, utilities::tags::NeedsContext >, dftracer::utils::utilities::Utility< std::vector< ItemInput >, std::vector< ItemOutput >, utilities::tags::NeedsContext >, dftracer::utils::utilities::Utility< Lines, Lines >, dftracer::utils::utilities::Utility< filesystem::FileEntry, text::Text >, dftracer::utils::utilities::Utility< AssociationResolverInput, AssociationResolverOutput >, dftracer::utils::utilities::Utility< JsonParserInput, JsonParserOutput >, dftracer::utils::utilities::Utility< MetadataCollectorUtilityInput, MetadataCollectorUtilityOutput >, dftracer::utils::utilities::Utility< ChunkVerificationUtilityInput< ChunkType, MetadataType >, ChunkVerificationUtilityOutput, utilities::tags::NeedsContext >, dftracer::utils::utilities::Utility< IndexedReadInput, std::shared_ptr< reader::internal::Reader > >, dftracer::utils::utilities::Utility< DirectoryScannerUtilityInput, std::vector< FileEntry >, utilities::tags::NeedsContext >, dftracer::utils::utilities::Utility< ReconstructionPlannerInput, ReconstructionPlan >, dftracer::utils::utilities::Utility< StreamReadInput, ChunkRange >, dftracer::utils::utilities::Utility< ChunkExtractorUtilityInput, ChunkExtractorUtilityOutput >, dftracer::utils::utilities::Utility< FileMergerUtilityInput, Result< FileMergerUtilityOutput > >, dftracer::utils::utilities::Utility< StringJsonParserInput, JsonParserOutput >
Public Functions
-
UtilityBase() = default¶
-
virtual ~UtilityBase() = default¶
-
UtilityBase(const UtilityBase&) = default¶
-
UtilityBase &operator=(const UtilityBase&) = default¶
-
UtilityBase(UtilityBase&&) = default¶
-
UtilityBase &operator=(UtilityBase&&) = default¶
Public Static Functions
-
template<typename Tag>
static inline constexpr bool has_tag()¶
-
static inline constexpr std::string_view get_type_signature()¶
-
static inline constexpr std::string_view get_name()¶
Friends
- friend class ::dftracer::utils::CoroScope
-
template<typename I, typename O, typename ...Tags>
class UtilityExecutor¶ Runs a utility’s process(), injecting the CoroScope context when the utility needs it.
- Template Parameters:
I – Input type
O – Output type
Tags – Variadic tag types
Public Functions
Construct an executor for the given utility.
- Parameters:
utility – The utility to execute
-
inline coro::CoroTask<O> execute(const I &input)¶
Execute the utility’s process() without a CoroScope context.
- Parameters:
input – Input to process
- Throws:
Any – exception propagated from the utility’s process()
- Returns:
Output result
-
inline coro::CoroTask<O> execute(CoroScope &ctx, const I &input)¶
Execute the utility with a CoroScope context.
Sets the context reference before calling process() and clears it afterwards (including on error), so context-needing utilities can emit dynamic tasks for the duration of execution.
- Parameters:
ctx – Task context for dynamic task emission
input – Input to process
- Throws:
Any – exception propagated from the utility’s process()
- Returns:
Output result
-
struct NeedsContext¶
Marker tag indicating that a utility needs CoroScope for dynamic task emission.
Usage:
class MyUtility : public Utility<Input, Output, tags::NeedsContext> { Output process(const Input& input, CoroScope& ctx) override { // Can use ctx.emit() here for dynamic task emission return result; } };
Free Functions¶
-
template<typename Tag, typename UtilityType>
constexpr bool dftracer::utils::utilities::has_tag()¶ Check if a Utility type has a specific tag.
Usage:
if constexpr (has_tag<tags::NeedsContext, MyUtility>()) { // Utility is parallelizable }
-
template<typename I>
consteval auto dftracer::utils::utilities::make_input_signature()¶
-
template<typename I, typename Batch>
consteval auto dftracer::utils::utilities::make_streaming_signature()¶
-
template<typename I, typename O>
consteval auto dftracer::utils::utilities::make_utility_signature()¶
-
inline bool dftracer::utils::utilities::monitor_deep_enabled()¶
-
long long dftracer::utils::utilities::monitor_enqueue(const void *handle, CoroKind kind)¶
-
void dftracer::utils::utilities::monitor_resume_begin(long long id)¶
-
void dftracer::utils::utilities::monitor_resume_end(const void *handle, bool done)¶
-
void dftracer::utils::utilities::monitor_set_worker(int worker_id)¶
-
void dftracer::utils::utilities::monitor_sync_complete(const void *handle)¶
-
bool dftracer::utils::utilities::monitoring_enabled()¶
Factory function to create a UtilityAdapter:
use(utility).as_task().Reads naturally as “use this utility as a task”. Deduces the base Utility type from the utility’s Input/Output/TagsTuple, so it also works with derived utility classes.
auto utility = std::make_shared<MyUtility>(); // Basic usage: convert to a task and schedule it auto task = use(utility).as_task(); scheduler.schedule(task, input); // Implicit conversion to std::shared_ptr<Task> std::shared_ptr<Task> task = use(utility); // Use with the Task API auto parent = make_task([]() { return 42; }); auto child = use(utility).as_task(); child->depends_on(parent);
- Parameters:
utility – Shared pointer to the utility (may be a derived class)
- Returns:
UtilityAdapter ready for conversion to a Task via as_task()