Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
27 commits
Select commit Hold shift + click to select a range
6d3ac59
fix: restart remote configuration worker after fork
pavlokhrebto Sep 14, 2026
eb1c92b
fix: restart agentless delivery after fork
pavlokhrebto Sep 14, 2026
de43b97
fix: stop feature flag delivery on provider replacement
pavlokhrebto Sep 15, 2026
8b04e9b
fix: synchronize remote configuration client registration
pavlokhrebto Sep 15, 2026
9c6b44d
fix: emit provider lifecycle events when configuration changes
pavlokhrebto Sep 15, 2026
81ac9dc
test: verify runtime capability encoding preserves boot capabilities
pavlokhrebto Sep 15, 2026
1291b08
fix(openfeature): preserve configuration across provider lifecycles
pavlokhrebto Sep 15, 2026
ccb9894
ci: test against agentless system-tests branch
pavlokhrebto Sep 10, 2026
e08869f
ci: pin agentless system-tests workflow revision
pavlokhrebto Sep 14, 2026
dde41ac
ci: restore system-tests updater markers
pavlokhrebto Sep 14, 2026
4a96023
fix(openfeature): register remote capabilities by selected source
pavlokhrebto Sep 15, 2026
60c6d41
ci: enable final Ruby agentless system tests
pavlokhrebto Sep 15, 2026
3b438b7
fix(openfeature): cancel provider initialization on shutdown
pavlokhrebto Sep 21, 2026
eac8af4
fix(openfeature): report configuration applied during shutdown as lost
pavlokhrebto Sep 21, 2026
3aed870
fix(openfeature): preserve configuration loss during initialization
pavlokhrebto Sep 21, 2026
8ad3e85
fix(openfeature): preserve delivery across multiple providers
pavlokhrebto Sep 23, 2026
4cc6ca8
fix(openfeature): invalidate pending readiness on configuration loss
pavlokhrebto Sep 23, 2026
30aa91e
fix(open_feature): isolate post-fork delivery failures
pavlokhrebto Sep 24, 2026
8c31157
fix(openfeature): address remaining lifecycle review feedback
pavlokhrebto Sep 25, 2026
af56fb9
refactor(remote): make capabilities telemetry explicit
pavlokhrebto Sep 25, 2026
444836d
style(openfeature): use explicit boolean coercion
pavlokhrebto Sep 25, 2026
38b8536
refactor(openfeature): simplify shutdown handler capture
pavlokhrebto Sep 25, 2026
6d4b426
refactor(openfeature): clarify READY dispatch bookkeeping
pavlokhrebto Sep 25, 2026
507e7b5
test(openfeature): clean up lifecycle concurrency specs
pavlokhrebto Sep 25, 2026
323f71d
fix(remote-config): avoid restarting polling after fork
pavlokhrebto Sep 28, 2026
a2a8cc3
refactor(openfeature): reuse provider event handler type
pavlokhrebto Sep 30, 2026
a1607ad
fix(openfeature): restore lifecycle behavior after stack rebase
pavlokhrebto Sep 30, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .github/workflows/system-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,7 @@ jobs:
test:
needs:
- build
uses: DataDog/system-tests/.github/workflows/system-tests.yml@04861df56371940cd3d0c6bb22392acf75a4624d # Automated: Updated by .github/workflows/update-system-tests.yml.
uses: DataDog/system-tests/.github/workflows/system-tests.yml@275c94f119df1e2cd74d62f90868df4e20c8adc2 # Automated: Updated by .github/workflows/update-system-tests.yml.
permissions:
contents: read
id-token: write
Expand All @@ -80,7 +80,7 @@ jobs:
binaries_artifact: dd-trace-rb
desired_execution_time: 300 # 5 minutes
scenarios_groups: tracer_release
ref: 04861df56371940cd3d0c6bb22392acf75a4624d # Automated: Updated by .github/workflows/update-system-tests.yml.
ref: 275c94f119df1e2cd74d62f90868df4e20c8adc2 # Automated: Updated by .github/workflows/update-system-tests.yml.
force_execute: ${{ needs.build.outputs.forced_tests }}
parametric_job_count: 8
push_to_test_optimization: true
Expand Down
10 changes: 8 additions & 2 deletions lib/datadog/core/configuration/components.rb
Original file line number Diff line number Diff line change
Expand Up @@ -187,13 +187,11 @@ def initialize(settings)

@telemetry = self.class.build_telemetry(settings, agent_settings, @logger)

# Bind Remote Configuration dispatch to this tree, which starts before it becomes the global Components instance.
@remote = Remote::Component.build(
settings,
agent_settings,
logger: @logger,
telemetry: telemetry,
open_feature_component_provider: -> { @open_feature },
)
@tracer = Datadog::Tracing::Component.build_tracer(settings, agent_settings, logger: @logger)
@crashtracker = self.class.build_crashtracker(settings, agent_settings, logger: @logger)
Expand Down Expand Up @@ -262,6 +260,14 @@ def after_fork
ProcessDiscovery.after_fork
symbol_database&.after_fork!
data_streams&.restart_flush_thread
begin
@open_feature_activation.after_fork
rescue => e
# Feature Flags is optional and must never interrupt other post-fork handlers.
description = "Feature Flags delivery failed to restart after fork"
logger.error("#{description}: #{e.class}: #{e.message}")
telemetry.report(e, description: description)
end
end

# Hot-swaps with a new sampler.
Expand Down
18 changes: 3 additions & 15 deletions lib/datadog/core/remote/client/capabilities.rb
Original file line number Diff line number Diff line change
Expand Up @@ -6,23 +6,21 @@
require_relative "../../../di/remote"
require_relative "../../../symbol_database"
require_relative "../../../symbol_database/remote"
require_relative "../../../open_feature/remote"

module Datadog
module Core
module Remote
class Client
# Capabilities
class Capabilities
def initialize(settings, telemetry, open_feature_component_provider: nil)
open_feature_component_provider ||= -> {}
def initialize(settings, telemetry:)
@capabilities = []
@products = []
@receivers = []
@telemetry = telemetry
@mutex = Mutex.new

register(settings, open_feature_component_provider)
register(settings)

@base64_capabilities = capabilities_to_base64(@capabilities)
end
Expand Down Expand Up @@ -65,7 +63,7 @@ def remove_products(*products)

private

def register(settings, open_feature_component_provider)
def register(settings)
if settings.respond_to?(:appsec) && settings.appsec.enabled
register_capabilities(Datadog::AppSec::Remote.capabilities)
register_products(Datadog::AppSec::Remote.products)
Expand Down Expand Up @@ -113,16 +111,6 @@ def register(settings, open_feature_component_provider)
end
end
end
if settings.respond_to?(:open_feature) && settings.open_feature.enabled
register_capabilities(Datadog::OpenFeature::Remote.capabilities)
register_products(Datadog::OpenFeature::Remote.products)
register_receivers(
Datadog::OpenFeature::Remote.receivers(
@telemetry,
component_provider: open_feature_component_provider,
),
)
end
end

def register_capabilities(capabilities)
Expand Down
45 changes: 25 additions & 20 deletions lib/datadog/core/remote/component.rb
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ module Remote
#
# @api private
class Component
attr_reader :logger, :client, :healthy, :worker
attr_reader :logger, :healthy, :worker

def initialize(settings, capabilities, agent_settings, logger:)
@logger = logger
Expand All @@ -28,6 +28,7 @@ def initialize(settings, capabilities, agent_settings, logger:)

@barrier = Barrier.new(settings.remote.boot_timeout_seconds)

@client_mutex = Mutex.new
@client = Client.new(@transport, @capabilities, settings: settings, logger: logger)
@healthy = false
logger.debug { "new remote configuration client: #{@client.id} products: #{@capabilities.products.sort.join(", ")}" }
Expand All @@ -40,7 +41,7 @@ def initialize(settings, capabilities, agent_settings, logger:)
end

begin
@client.sync
client.sync
@healthy ||= true
rescue Client::SyncError => e
# Transient errors due to network or agent. Logged the error but not via telemetry
Expand All @@ -60,9 +61,11 @@ def initialize(settings, capabilities, agent_settings, logger:)
end

# client state is unknown, state might be corrupted
@client = Client.new(@transport, @capabilities, settings: settings, logger: logger)
new_client = @client_mutex.synchronize do
@client = Client.new(@transport, @capabilities, settings: settings, logger: logger)
end
@healthy = false
logger.debug { "new remote configuration client: #{@client.id} products: #{@capabilities.products.sort.join(", ")}" }
logger.debug { "new remote configuration client: #{new_client.id} products: #{@capabilities.products.sort.join(", ")}" }

# TODO: bail out if too many errors?
end
Expand All @@ -81,6 +84,10 @@ def started?
@worker.started?
end

def client
@client_mutex.synchronize { @client }
end

# If the worker is not initialized, initialize it.
#
# Then, waits for one client sync to be executed if `kind` is `:once`.
Expand All @@ -96,9 +103,11 @@ def shutdown!
# Recreates the remote configuration client after a fork.
# This ensures each forked process has a unique client ID and fresh state.
def after_fork
@client = Client.new(@transport, @capabilities, settings: @settings, logger: @logger)
new_client = @client_mutex.synchronize do
@client = Client.new(@transport, @capabilities, settings: @settings, logger: @logger)
end
@healthy = false
logger.debug { "remote configuration client recreated after fork: #{@client.id} products: #{@capabilities.products.sort.join(", ")}" }
logger.debug { "remote configuration client recreated after fork: #{new_client.id} products: #{@capabilities.products.sort.join(", ")}" }
end

def add_products(*products)
Expand All @@ -110,14 +119,14 @@ def remove_products(*products)
end

def register(capabilities:, products:, receivers:)
@capabilities.register_runtime(
capabilities: capabilities,
products: products,
receivers: receivers,
)
# A client created concurrently after registration receives the new
# receivers from Capabilities; an older client is updated here.
@client.dispatcher.add_receivers(*receivers)
@client_mutex.synchronize do
@capabilities.register_runtime(
capabilities: capabilities,
products: products,
receivers: receivers,
)
@client.dispatcher.add_receivers(*receivers)
end
end

# Barrier provides a mechanism to fence execution until a condition happens
Expand Down Expand Up @@ -198,14 +207,10 @@ class << self
#
# Those checks are instead performed inside the worker loop.
# This allows users to upgrade their agent while keeping their application running.
def build(settings, agent_settings, logger:, telemetry:, open_feature_component_provider: nil)
def build(settings, agent_settings, logger:, telemetry:)
return unless settings.remote.enabled

capabilities = Client::Capabilities.new(
settings,
telemetry,
open_feature_component_provider: open_feature_component_provider,
)
capabilities = Client::Capabilities.new(settings, telemetry: telemetry)
new(settings, capabilities, agent_settings, logger: logger)
end
end
Expand Down
14 changes: 10 additions & 4 deletions lib/datadog/open_feature.rb
Original file line number Diff line number Diff line change
Expand Up @@ -20,23 +20,29 @@ def self.activate_provider(provider)
# Initialize the component tree before taking its reconfiguration lock.
Datadog.send(:components)
Datadog.send(:safely_synchronize) do
return [nil, nil] if provider.send(:shutdown?)

components = Datadog.send(:components, allow_initialization: false)
activation = components&.send(:open_feature_activation)
@adopted_provider = provider
providers = @adopted_providers ||= {}.compare_by_identity
providers[provider] = true
[activation&.activate(provider), activation&.failure]
end
end

def self.deactivate_provider(provider)
Datadog.send(:safely_synchronize) do
@adopted_provider = nil if @adopted_provider&.equal?(provider)
@adopted_providers&.delete(provider)
components = Datadog.send(:components, allow_initialization: false)
components&.send(:open_feature_activation)&.deactivate(provider)
end
end

# Components startup already holds the non-reentrant reconfiguration lock.
def self.reattach(activation)
provider = @adopted_provider
activation.activate(provider) if provider
providers = @adopted_providers&.keys || []
providers.each { |provider| activation.activate(provider) }
nil
end
end
end
58 changes: 54 additions & 4 deletions lib/datadog/open_feature/activation.rb
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,9 @@ def initialize(settings, agent_settings, remote, logger:, telemetry:)
@telemetry = telemetry
@component = nil
@configuration_source = nil
@delivery_source = nil
@failure = nil
@provider = nil
@providers = {}.compare_by_identity
@activated = false
@delivery_started = false
@shutdown = false
Expand All @@ -33,7 +34,7 @@ def activate(provider)
return if @shutdown
return unless open_feature_available?

@provider = provider
@providers[provider] = true
return @component if @activated && @delivery_started
return if @activated

Expand All @@ -58,6 +59,49 @@ def activated?
@mutex.synchronize { @activated }
end

def providers
@mutex.synchronize { @providers.keys }
end

# Timer-driven delivery has no operation in the child that can restart its inherited worker.
def after_fork
configuration_source = @mutex.synchronize do
if @shutdown || !@delivery_started
nil
else
@configuration_source
end
end

configuration_source&.start
nil
end

def deactivate(provider)
configuration_source, component = @mutex.synchronize do
return if @shutdown
return unless @providers.delete(provider)
return unless @providers.empty?
if @delivery_source == Configuration::Source::REMOTE_CONFIG && @delivery_started
# Remote Configuration is process-scoped and may already hold configuration needed by the next provider.
[nil, nil]
else
previous_source = @configuration_source
previous_component = @component
@component = nil
@configuration_source = nil
@delivery_source = nil
@activated = false
@delivery_started = false
@failure = nil
[previous_source, previous_component]
end
end

configuration_source&.stop
component&.shutdown!
end

def shutdown!
configuration_source, component = @mutex.synchronize do
return if @shutdown
Expand All @@ -67,6 +111,10 @@ def shutdown!
end

configuration_source&.stop
configuration_received = @mutex.synchronize do
!@providers.empty? && !!component&.configuration_received?
end
configuration_changed(Component::CONFIGURATION_LOST) if configuration_received
component&.shutdown!
end

Expand Down Expand Up @@ -97,11 +145,13 @@ def activate_delivery(resolution)
end

@component = component
@delivery_source = resolution.source
@delivery_started = start_delivery(resolution.source)
return component if @delivery_started

component.shutdown!
@component = nil
@delivery_source = nil
nil
end

Expand Down Expand Up @@ -173,8 +223,8 @@ def start_remote_configuration
end

def configuration_changed(event)
provider = @mutex.synchronize { @provider }
provider&.send(:configuration_changed, event)
providers = @mutex.synchronize { @providers.keys }
providers.each { |provider| provider.send(:configuration_changed, event) }
end
end
end
Expand Down
3 changes: 3 additions & 0 deletions lib/datadog/open_feature/component.rb
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ class Component
CONFIGURATION_TIMEOUT = :timeout
CONFIGURATION_SHUTDOWN = :shutdown
CONFIGURATION_CHANGED = :changed
CONFIGURATION_LOST = :lost

attr_reader :engine, :flag_eval_metrics_hook, :flag_eval_evp_hook, :span_enrichment_hook

Expand Down Expand Up @@ -92,6 +93,8 @@ def reconfigure!(configuration)

if @configuration_received
previously_received ? CONFIGURATION_CHANGED : CONFIGURATION_READY
elsif previously_received
CONFIGURATION_LOST
end
end
end
Expand Down
Loading
Loading