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 Types

using BatchType = Batch

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.

Public Static Functions

static inline constexpr std::string_view get_type_signature()
static inline constexpr std::string_view get_name()
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 Types

using Output = O

Public Functions

inline Utility()
template<typename Dummy = void, typename = std::enable_if_t<(sizeof...(Tags) > 0) && std::is_void_v<Dummy>>>
inline explicit Utility(Tags... tags)
virtual coro::CoroTask<O> process(const I &input) = 0
inline coro::CoroTask<O> process(I &&input)

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

Public Functions

inline explicit UtilityAdapter(std::shared_ptr<Utility<I, O, Tags...>> utility)
inline std::shared_ptr<Task> as_task()

Convert the utility to a Task.

Detects whether the utility needs CoroScope and builds the matching task signature, executing via UtilityExecutor.

inline operator std::shared_ptr<Task>()

Implicit conversion to Task for convenience.

Public Static Functions

static inline constexpr bool needs_context()

Check if the utility needs CoroScope at compile time.

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 Types

using Input = I
using TagsTuple = std::tuple<Tags...>

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

inline explicit UtilityExecutor(std::shared_ptr<Utility<I, O, Tags...>> utility)

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:
  • ctxTask context for dynamic task emission

  • input – Input to process

Throws:

Any – exception propagated from the utility’s process()

Returns:

Output result

inline std::shared_ptr<Utility<I, O, Tags...>> get_utility() const
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()
template<typename DerivedUtility>
auto dftracer::utils::utilities::use(std::shared_ptr<DerivedUtility> utility)

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()