Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,7 @@ def __init__( # noqa: PLR0913, PLR0915
channel_credentials: grpc.ChannelCredentials | None = None,
sync_metadata_disabled: bool | None = None,
fatal_status_codes: list[str] | None = None,
sync_metadata: typing.Sequence[tuple[str, str]] | None = None,
):
self.host = env_or_default(ENV_VAR_HOST, DEFAULT_HOST) if host is None else host

Expand Down Expand Up @@ -278,3 +279,11 @@ def __init__( # noqa: PLR0913, PLR0915
# Disabling will prevent static context from flagd being used in evaluations.
# GetMetadata and this option will be removed.
self.sync_metadata_disabled = sync_metadata_disabled

# Additional gRPC metadata headers sent on every in-process SyncFlags call.
# Useful for injecting infrastructure-specific headers, e.g. disabling a
# proxy/mesh request timeout on the long-lived sync stream. These are merged
# with any headers the provider sets itself (such as ``flagd-selector``).
self.sync_metadata: tuple[tuple[str, str], ...] = (
tuple(sync_metadata) if sync_metadata is not None else ()
)
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@ def __init__( # noqa: PLR0913
channel_credentials: grpc.ChannelCredentials | None = None,
sync_metadata_disabled: bool | None = None,
fatal_status_codes: list[str] | None = None,
sync_metadata: typing.Sequence[tuple[str, str]] | None = None,
):
"""
Create an instance of the FlagdProvider
Expand All @@ -83,6 +84,9 @@ def __init__( # noqa: PLR0913
:param stream_deadline_ms: the maximum time to wait before a request times out
:param keep_alive_time: the number of milliseconds to keep alive
:param resolver_type: the type of resolver to use
:param sync_metadata: additional gRPC metadata headers (key-value tuples) sent on
every in-process SyncFlags call, e.g. to disable a proxy/mesh
request timeout on the long-lived sync stream
"""
if deadline_ms is None and timeout is not None:
deadline_ms = timeout * 1000
Expand Down Expand Up @@ -112,6 +116,7 @@ def __init__( # noqa: PLR0913
default_authority=default_authority,
channel_credentials=channel_credentials,
sync_metadata_disabled=sync_metadata_disabled,
sync_metadata=sync_metadata,
fatal_status_codes=fatal_status_codes,
)
self.enriched_context: dict = {}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -212,17 +212,21 @@ def _create_request_args(self) -> dict:

return request_args

def _create_metadata(self) -> tuple[tuple[str, str]] | None:
def _create_metadata(self) -> tuple[tuple[str, str], ...] | None:
"""Create gRPC metadata headers for the request.

Returns gRPC metadata as a tuples of tuples containing header key-value pairs.
The selector is passed via the 'flagd-selector' header per flagd v0.11.0+ specification,
while also being included in the request body for backward compatibility with older flagd versions.
Any user-configured ``sync_metadata`` headers are appended, allowing callers to inject
infrastructure-specific headers (e.g. proxy/mesh timeout overrides) on the sync stream.
"""
if self.selector is None:
return None
metadata: tuple[tuple[str, str], ...] = ()
if self.selector is not None:
metadata += (("flagd-selector", self.selector),)
metadata += self.config.sync_metadata

return (("flagd-selector", self.selector),)
return metadata if metadata else None

def _fetch_metadata(self) -> sync_pb2.GetMetadataResponse | None:
if self.config.sync_metadata_disabled:
Expand Down
18 changes: 18 additions & 0 deletions providers/openfeature-provider-flagd/tests/test_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,24 @@ def test_return_default_values_rpc():
assert config.retry_backoff_ms == DEFAULT_RETRY_BACKOFF
assert config.stream_deadline_ms == DEFAULT_STREAM_DEADLINE
assert config.tls is DEFAULT_TLS
assert config.sync_metadata == ()


def test_sync_metadata_passthrough():
metadata = [("x-envoy-upstream-rq-timeout-ms", "0")]
config = Config(resolver=ResolverType.IN_PROCESS, sync_metadata=metadata)
assert config.sync_metadata == (("x-envoy-upstream-rq-timeout-ms", "0"),)


def test_positional_fatal_status_codes_backwards_compatible():
# fatal_status_codes must stay the last positional parameter so callers that
# passed it positionally before sync_metadata was added keep working.
# It is the 22nd positional parameter (21 parameters precede it).
leading_args = [None] * 21
config = Config(*leading_args, ["UNAVAILABLE", "DATA_LOSS"])
assert config.fatal_status_codes == ["UNAVAILABLE", "DATA_LOSS"]
# The positional value must not leak into sync_metadata.
assert config.sync_metadata == ()


def test_return_default_values_in_process():
Expand Down
39 changes: 39 additions & 0 deletions providers/openfeature-provider-flagd/tests/test_grpc_watcher.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ def setUp(self):
config.host = "localhost"
config.port = 5000
config.sync_metadata_disabled = False
config.sync_metadata = ()

flag_store = Mock(spec=FlagStore)
flag_store.update.return_value = None
Expand Down Expand Up @@ -160,3 +161,41 @@ def test_selector_passed_via_both_metadata_and_body(self):
self.assertIn("metadata", kwargs)
metadata = kwargs["metadata"]
self.assertEqual(metadata, (("flagd-selector", "test-selector"),))

def test_custom_sync_metadata_appended(self):
"""User-configured sync_metadata headers are sent on the SyncFlags call."""
self.grpc_watcher.selector = "test-selector"
self.grpc_watcher.config.sync_metadata = (
("x-envoy-upstream-rq-timeout-ms", "0"),
)
mock_stream = iter(
[SyncFlagsResponse(flag_configuration='{"flag_key": "flag_value"}')]
)
self.mock_stub.SyncFlags = Mock(return_value=mock_stream)

self.run_listen_and_shutdown_after()

metadata = self.mock_stub.SyncFlags.call_args.kwargs["metadata"]
self.assertEqual(
metadata,
(
("flagd-selector", "test-selector"),
("x-envoy-upstream-rq-timeout-ms", "0"),
),
)

def test_custom_sync_metadata_without_selector(self):
"""sync_metadata is sent even when no selector is configured."""
self.grpc_watcher.selector = None
self.grpc_watcher.config.sync_metadata = (
("x-envoy-upstream-rq-timeout-ms", "0"),
)
mock_stream = iter(
[SyncFlagsResponse(flag_configuration='{"flag_key": "flag_value"}')]
)
self.mock_stub.SyncFlags = Mock(return_value=mock_stream)

self.run_listen_and_shutdown_after()

metadata = self.mock_stub.SyncFlags.call_args.kwargs["metadata"]
self.assertEqual(metadata, (("x-envoy-upstream-rq-timeout-ms", "0"),))
Loading