Class: Datadog::OpenFeature::FlagEvaluation::Writer

Inherits:
Object
  • Object
show all
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

Attributes included from Core::Workers::Async::Thread

#error, #fork_policy, #result

Instance Method Summary collapse

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