BB.Estimator.Server (bb v0.31.2)

Copy Markdown View Source

Wrapper GenServer for estimator callback modules.

Responsibilities:

  • Resolves parameter references in opts at startup and on parameter changes (mirroring BB.Sensor.Server / BB.Controller.Server).
  • Subscribes to the estimator's declared input paths.
  • Rejects any envelope stamped with a node other than the local one before it is compared against another envelope's timestamp, since monotonic clocks are node-local. A rejected envelope is dropped with reason :cross_node and is not retained, so a multi-input alias whose newest envelope came from another node stays missing until a local one arrives.
  • Dispatches incoming input messages to the callback module via BB.Estimator.handle_input/2. Single-input estimators receive the bare %BB.Message{}; multi-input estimators receive a %{input_name => %BB.Message{}} map gathered when the driver input arrives.
  • Enforces max_input_age: an envelope older than its input's budget is discarded before handle_input/2 is called, and the estimator transitions to :degraded with reason :stale_input. Retained non-driver envelopes are re-checked when the driver arrives, since they can age past their budget while they wait. This is independent of latency_budget, which measures callback execution time.
  • For multi-input estimators, enforces sync_tolerance: if any non-driver input is stale relative to the driver by more than the configured tolerance, the dispatch is dropped (with [:bb, :estimator, :dropped] telemetry) instead of fired with a stale snapshot.
  • Publishes each {output_name, message} returned from a callback's {:reply, outputs, state} reply to that output's configured path.
  • Runs the :healthy / :degraded / :lost health state machine and dispatches the configured on_degraded / on_lost / on_recovered commands on each transition. For multi-input estimators only driver arrivals reset the lost_after timer, so a live auxiliary stream cannot mask a driver that has stopped publishing.
  • Emits :input, :output, :latency, :dropped, and :transition telemetry. Every timing measurement is in nanoseconds, the same unit as %BB.Message{}.monotonic_time.

Init args

The framework supplies the following keys in the start-link init arg. Internal keys (double-underscored) are stripped by the server before calling BB.Estimator.init/1; public keys (:bb, :estimator_context) are passed through unchanged so user code can read them.

  • :__callback_module__ - the user's estimator module.
  • :__estimator_inputs__ - `%{mode: :single:multi, inputs: [...],
    sync_tolerance_ns: integer()nil}` describing the input wiring. Each
    input spec carries :name, :path, :driver? and a resolved :max_input_age_ns.
  • :__estimator_outputs__ - %{output_name => [atom()]} mapping output names to their full pubsub paths.
  • :bb - %{robot: module, path: [atom]}, the per-process context.
  • :estimator_context - the BB.Estimator.Context.t() exposed to the callback module's init/1.

Plus any user options declared via the estimator's options_schema/0.

Summary

Functions

Returns a specification to start this module under a supervisor.

Types

health_state()

@type health_state() :: :healthy | :degraded | :lost

input_spec()

@type input_spec() :: %{
  :name => atom(),
  :path => [atom()],
  :driver? => boolean(),
  optional(:max_input_age_ns) => integer() | nil
}

t()

@type t() :: %BB.Estimator.Server{
  bb: %{robot: module(), path: [atom()]},
  callback_module: module(),
  consecutive_ok: non_neg_integer(),
  context: BB.Estimator.Context.t(),
  driver_input: atom() | nil,
  health_state: health_state(),
  input_name_by_path: %{required([atom()]) => atom()},
  inputs: [input_spec()],
  last_messages: %{required(atom()) => BB.Message.t()},
  latency_budget_ns: integer() | nil,
  lost_after_ns: integer() | nil,
  lost_timer_ref: reference() | nil,
  max_input_age_ns: %{required(atom()) => integer() | nil},
  mode: :single | :multi,
  on_degraded: atom() | nil,
  on_lost: atom() | nil,
  on_recovered: atom() | nil,
  outputs: %{required(atom()) => [atom()]},
  param_subscriptions: %{required([atom()]) => atom()},
  raw_opts: keyword(),
  recover_after: pos_integer(),
  resolved_opts: keyword(),
  sync_tolerance_ns: integer() | nil,
  user_state: term()
}

Functions

child_spec(init_arg)

Returns a specification to start this module under a supervisor.

See Supervisor.