Class: Datadog::OpenFeature::FlagEvaluation::Writer
- Inherits:
-
Object
- Object
- Datadog::OpenFeature::FlagEvaluation::Writer
- Includes:
- Core::Workers::Async::Thread
- Defined in:
- lib/datadog/open_feature/flag_evaluation/writer.rb
Overview
Background writer that drains the two-tier aggregation maps and POSTs batches to /evp_proxy/v2/api/v2/flagevaluation every FLUSH_INTERVAL_SECONDS.
The writer owns the aggregation cycle:
1. Hook calls enqueue (non-blocking) — never aggregates inline.
2. Background thread wakes, calls aggregator.record for each enqueued event, flushes.
3. flush_once drains aggregation maps, builds payload, sends via transport.
Thread model: MRI Ruby GIL — Mutex + ConditionVariable + SizedQueue + Thread. The flush loop waits on a ConditionVariable (interruptible) rather than a bare sleep, so #stop can wake the worker immediately and still drain + final-flush.
Constant Summary collapse
- FLUSH_INTERVAL_SECONDS =
10- DRAIN_INTERVAL_SECONDS =
0.1- SHUTDOWN_TIMEOUT_SECONDS =
5- QUEUE_SIZE =
4_096- MAX_DRAIN_EVENTS_PER_CYCLE =
1_024- PAYLOAD_SIZE_LIMIT_BYTES =
Core::EVP::PAYLOAD_SIZE_LIMIT_BYTES
- TELEMETRY_NAMESPACE =
"tracers"- ROWS_DROPPED_METRIC =
"flagevaluation.rows.dropped"- ROWS_DEGRADED_METRIC =
"flagevaluation.rows.degraded"- PAYLOAD_SPLITS_METRIC =
"flagevaluation.payload.splits"- CONTEXT_TRUNCATED_METRIC =
"flagevaluation.context.truncated"- REASON_QUEUE_OVERFLOW =
"queue_overflow"- REASON_DEGRADED_CAP =
"degraded_cap"- REASON_CARDINALITY_CAP =
"cardinality_cap"- REASON_PAYLOAD_LIMIT =
"payload_limit"- REASON_PRE_QUEUE_OVERFLOW =
"pre_queue_overflow"- REASON_SERIALIZATION_ERROR =
"serialization_error"- REASON_SNAPSHOT_ERROR =
"snapshot_error"- TARGETING_KEY_FIELD =
Must equal OpenFeature::SDK::EvaluationContext::TARGETING_KEY. Duplicated as a literal rather than referenced because the SDK is an optional dependency and this file loads without it. If the two drift, the targeting key stops being excluded from the context snapshot and lands in context.evaluation as raw PII, so a spec asserts the equality.
"targeting_key"- TARGETING_KEY_HASH_PREFIX =
"sha256_"
Constants included from Core::Workers::Async::Thread
Core::Workers::Async::Thread::FORK_POLICY_RESTART, Core::Workers::Async::Thread::FORK_POLICY_STOP, Core::Workers::Async::Thread::MUTEX_INIT, Core::Workers::Async::Thread::SHUTDOWN_TIMEOUT
Instance Attribute Summary collapse
-
#dropped_queue_overflow ⇒ Object
readonly
Observable count of events dropped because the async hand-off queue was full.
-
#service_context ⇒ Object
readonly
Service context fields for the batch wrapper.
Attributes included from Core::Workers::Async::Thread
Instance Method Summary collapse
-
#enqueue(flag_key:, eval_time_ms:, variant: nil, allocation_key: nil, error_message: nil, runtime_default: nil, targeting_key: nil, attrs: nil, observe_full_evaluation_data: false, **_event) ⇒ Object
Non-blocking enqueue from the finally hook.
-
#initialize(transport:, logger:, telemetry: nil) ⇒ Writer
constructor
A new instance of Writer.
-
#stop ⇒ Object
Stop the background thread and flush remaining events.
Methods included from Core::Workers::Async::Thread
#completed?, #error?, #failed?, #forked?, included, #join, #run_async?, #running?, #started?, #terminate
Constructor Details
#initialize(transport:, logger:, telemetry: nil) ⇒ Writer
Returns a new instance of Writer.
64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 |
# File 'lib/datadog/open_feature/flag_evaluation/writer.rb', line 64 def initialize(transport:, logger:, telemetry: nil) @transport = transport @logger = logger @telemetry = telemetry @aggregator = Aggregator.new @queue = SizedQueue.new(QUEUE_SIZE) @stop_mutex = Mutex.new @stop_cond = ConditionVariable.new @stopped = false @dropped_queue_overflow = 0 @dropped_pre_queue_overflow = 0 @context_truncated_counts = Hash.new(0) @context_snapshot_error_logged = false self.fork_policy = Core::Workers::Async::Thread::FORK_POLICY_RESTART @service_context = build_service_context start_background_thread end |
Instance Attribute Details
#dropped_queue_overflow ⇒ Object (readonly)
Observable count of events dropped because the async hand-off queue was full. Reset to 0 each flush after being emitted, mirroring the aggregator's overflow counter.
62 63 64 |
# File 'lib/datadog/open_feature/flag_evaluation/writer.rb', line 62 def dropped_queue_overflow @dropped_queue_overflow end |
#service_context ⇒ Object (readonly)
Service context fields for the batch wrapper.
58 59 60 |
# File 'lib/datadog/open_feature/flag_evaluation/writer.rb', line 58 def service_context @service_context end |
Instance Method Details
#enqueue(flag_key:, eval_time_ms:, variant: nil, allocation_key: nil, error_message: nil, runtime_default: nil, targeting_key: nil, attrs: nil, observe_full_evaluation_data: false, **_event) ⇒ Object
Non-blocking enqueue from the finally hook. Drops + counts on overflow.
85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 |
# File 'lib/datadog/open_feature/flag_evaluation/writer.rb', line 85 def enqueue( flag_key:, eval_time_ms:, variant: nil, allocation_key: nil, error_message: nil, runtime_default: nil, targeting_key: nil, attrs: nil, observe_full_evaluation_data: false, **_event ) start_background_thread if forked? # Avoid snapshot work when the queue is already full. if @queue.size >= QUEUE_SIZE @stop_mutex.synchronize { @dropped_pre_queue_overflow += 1 } return end observe_full_evaluation_data = observe_full_evaluation_data == true attrs = observe_full_evaluation_data ? snapshot_context(attrs) : nil bounded_event = { flag_key: snapshot_string(flag_key), variant: snapshot_string(variant), allocation_key: snapshot_string(allocation_key), error_message: normalized_error_code(), runtime_default: runtime_default, targeting_key: snapshot_targeting_key(targeting_key), eval_time_ms: snapshot_integer(eval_time_ms), attrs: attrs, observe_full_evaluation_data: observe_full_evaluation_data, } @queue.push(bounded_event, true) @stop_mutex.synchronize { @stop_cond.signal } start_background_thread unless running? rescue ThreadError # Report queue backpressure on the next flush. @stop_mutex.synchronize { @dropped_queue_overflow += 1 } end |
#stop ⇒ Object
Stop the background thread and flush remaining events. Wakes the worker out of its interruptible wait so the drain + final flush happen immediately (no up-to-10s delay).
121 122 123 124 125 126 127 128 129 130 131 |
# File 'lib/datadog/open_feature/flag_evaluation/writer.rb', line 121 def stop @stop_mutex.synchronize do @stopped = true @stop_cond.broadcast end return true if join(SHUTDOWN_TIMEOUT_SECONDS) @logger.debug { "OpenFeature EVP: writer did not stop gracefully; terminating worker thread" } terminate end |