ProductDecoder

Trait ProductDecoder 

Source
pub trait ProductDecoder:
    Default
    + Send
    + 'static {
    type Snapshot: Send + Sync + 'static;
    type Error: ApplyError + Send + Sync + 'static;

    const PRODUCT: &'static str;

    // Required methods
    fn decode(
        &mut self,
        id: &ConfigId,
        payload: &[u8],
    ) -> Result<(), Self::Error>;
    fn build(self) -> Result<Self::Snapshot, Self::Error>;
}
Expand description

Decodes one product’s assigned configurations into a snapshot.

A decoder names the product it decodes, so subscribing with a decoder cannot attach it to the wrong product. Define a decoder in the crate that owns its snapshot type. A decoder for a product whose type has no other owner belongs in this crate instead; otherwise this crate knows no product’s payload.

The client drives an implementation once per published snapshot: it constructs a decoder with Default, calls decode once per assigned configuration in ascending ConfigId order, and then calls build. Ascending order is part of the contract, so a reduction that keeps the last valid configuration is well defined.

A product is decoded again only when its assigned configurations or their contents change, so a rejected assignment is not retried until it changes. When the Agent reports its configuration expired, every product that had configurations is decoded as an empty assignment; a product that already had none is not decoded again. The client then asks the Agent for the full assignment until a response is not expired, so configurations that return are decoded as new.

A rejected snapshot leaves subscribers with the last accepted one, so an implementation must build a snapshot from an empty assignment; see build.

Decoding runs on the client’s worker task and delays polling for every product while it runs, so implementations MUST NOT block. A panic in decode or build is caught: the client discards the decoder, rejects every configuration in the assignment, and publishes nothing, so subscribers keep the last accepted snapshot and are not notified.

§Design

A product’s assignment is several configurations, each with its own payload, and the protocol carries an apply status for each one separately. Decoding therefore accumulates configuration by configuration rather than consuming the assignment as a whole: a configuration that decode rejects is attributed to itself and skipped, and the decoder still builds from the rest. Because the client drives the loop, that attribution cannot be forgotten or misdirected. Whether a rejected configuration invalidates the whole snapshot remains the decoder’s choice: build may reject the snapshot, for example when a configuration it requires was rejected.

Accumulating into the decoder, rather than returning one decoded value per configuration, is what lets a product whose configurations have different shapes hold each of them in a field of its own type instead of funneling them through a shared enum.

§Examples

A product that accepts at most one configuration records every configuration it is given and rejects more than one in build, since only build sees the whole assignment:

use datadog_agent_remote_config::{ConfigId, ProductDecoder};

/// Holds the product's single configuration, if it has one.
#[derive(Default)]
struct SingleDecoder {
    payloads: Vec<(ConfigId, Vec<u8>)>,
}

impl ProductDecoder for SingleDecoder {
    const PRODUCT: &'static str = "EXAMPLE_SINGLE";

    type Snapshot = Option<Vec<u8>>;
    type Error = String;

    fn decode(&mut self, id: &ConfigId, payload: &[u8]) -> Result<(), Self::Error> {
        self.payloads.push((id.clone(), payload.to_vec()));
        Ok(())
    }

    fn build(mut self) -> Result<Self::Snapshot, Self::Error> {
        if self.payloads.len() > 1 {
            return Err(format!("Expected at most one configuration; got {}.", self.payloads.len()));
        }
        Ok(self.payloads.pop().map(|(_, payload)| payload))
    }
}

let mut decoder = SingleDecoder::default();
decoder.decode(&ConfigId::new("a"), b"one").unwrap();
decoder.decode(&ConfigId::new("b"), b"two").unwrap();
assert!(decoder.build().is_err());

Required Associated Constants§

Source

const PRODUCT: &'static str

The product this decoder decodes, by its protocol name, such as APM_SEMANTIC_CORE_DD.

The client requests this product from the Agent and keys the product’s subscription by this name. The client can only subscribe once per product: attempting a second subscription to the same product while the first is live returns Error::AlreadySubscribed, whichever decoder it uses.

Required Associated Types§

Source

type Snapshot: Send + Sync + 'static

The type that holds the product’s decoded configuration.

build produces it from the product’s configurations, and subscribers receive it.

Source

type Error: ApplyError + Send + Sync + 'static

The error this product’s decoding and validation produces.

Use String when there is nothing richer to report; a product that wants to describe a failure more precisely for its own diagnostics defines its own type and implements ApplyError for it.

Required Methods§

Source

fn decode(&mut self, id: &ConfigId, payload: &[u8]) -> Result<(), Self::Error>

Accumulates one of the product’s assigned configurations.

The client calls this at most once per distinct id within a single snapshot, then calls build once. A decoder may therefore hold one slot per configuration ID without a second configuration silently displacing the first.

§Errors

Returns Self::Error to reject this configuration alone. The client attributes the rejection to it and skips it, and still decodes the product’s remaining configurations.

Source

fn build(self) -> Result<Self::Snapshot, Self::Error>

Validates the accumulated configurations and produces the snapshot to publish.

Check requirements across configurations here, including whether a required configuration is present when the decoder receives a non-empty assignment. The client also calls this without calling decode when no configurations are assigned (after the backend removes the last one or the Agent reports its configuration expired), or when every assigned configuration was rejected for sharing an ID.

An empty assignment to the decoder means “no configuration,” even if files were assigned but all collided. This method MUST build a snapshot representing that state. Rejecting an empty assignment leaves subscribers reading configuration the backend removed or that expired with the Agent’s cache until a non-empty assignment builds successfully, because a rejection keeps the last accepted snapshot.

The client acknowledges successfully decoded configurations only when this method succeeds.

§Errors

Returns Self::Error to reject the snapshot. The client rejects every successfully decoded configuration with this error’s apply reason; configurations rejected by decode keep their own errors. No new snapshot is published, and subscribers retain the last accepted snapshot.

This method cannot reject selected configurations while publishing the rest. Selective rejection belongs in decode.

Dyn Compatibility§

This trait is not dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety", so this trait is not object safe.

Implementors§