Packages
statsig_elixir
0.20.3-beta.2607240300
0.20.3-beta.2607250300
0.20.3-beta.2607240300
0.20.2
0.20.2-rc.2607201629
0.20.2-beta.2607230301
0.20.2-beta.2607220301
0.20.2-beta.2607180300
0.20.2-beta.2607170300
0.20.1
0.20.1-rc.2607140134
0.20.1-rc.2607132133
0.20.1-beta.2607150300
0.20.0
0.19.9-rc.2607090505
0.19.9-rc.2607010924
0.19.9-beta.2607110300
0.19.9-beta.2607100302
0.19.9-beta.2607090302
0.19.9-beta.2607080301
0.19.9-beta.2607020303
0.19.9-beta.2607010304
0.19.9-beta.2606260304
0.19.8
0.19.7-rc.2606232103
0.19.7-beta.2606240303
0.19.7-beta.2606190306
0.19.7-beta.2606180304
0.19.6
0.19.6-rc.2606161624
0.19.6-rc.2606151903
0.19.6-beta.2606170305
0.19.6-beta.2606100304
0.19.5
0.19.5-rc.2606050349
0.19.5-beta.2606060303
0.19.4
0.19.4-rc.2605212136
0.19.3
0.19.3-beta.2604211739
0.19.2
0.19.2-rc.2604140138
0.19.2-beta.2604150312
0.19.1
0.19.1-rc.2604121707
0.19.1-beta.2604130314
0.19.1-beta.2604110309
0.19.1-beta.2604100313
0.19.0
0.18.2-rc.2604080018
0.18.2-beta.2604080312
0.18.1
0.18.1-rc.2604070138
0.18.1-beta.2604072116
0.18.0
0.17.3-rc.2603272156
0.17.3-beta.2603290312
0.17.3-beta.2603280308
0.17.3-beta.2603250310
0.17.2
0.17.2-rc.2603240138
0.17.2-beta.2603241852
0.17.1
0.17.1-rc.2603201903
0.17.1-rc.2603192323
0.17.1-beta.2603210301
0.17.1-beta.2603200305
0.17.1-beta.2603181954
0.17.0
0.16.6-rc.2603180529
0.16.6-rc.2603170138
0.16.6-beta.2603181807
0.16.5
0.16.5-rc.2603131810
0.16.5-beta.2603130304
0.16.4
0.16.4-rc.2603092140
0.16.4-beta.2603110302
0.16.4-beta.2603060304
0.16.4-beta.2603051742
0.16.4-beta.2603040814
0.16.4-beta.2603040303
0.16.3
0.16.3-rc.2602271954
0.16.3-rc.2602270324
0.16.3-rc.2602240138
0.16.3-beta.2603040747
0.16.3-beta.2603021916
0.16.3-beta.2602280254
0.16.3-beta.2602262056
0.16.3-beta.2602250307
0.16.3-beta.2602240529
0.16.2
0.16.2-rc.2602202152
0.16.2-rc.2602200042
0.16.2-beta.2602210300
0.16.2-beta.2602200305
0.16.1
0.16.1-rc.2602180406
0.16.1-beta.2602190307
0.16.0
0.15.2-rc.2602172254
0.15.2-rc.2602140111
0.15.2-rc.2602130016
0.15.2-rc.2602100139
0.15.2-rc.2602030138
0.15.2-beta.2602180309
0.15.2-beta.2602150309
0.15.2-beta.2602140304
0.15.2-beta.2602130310
0.15.2-beta.2602122245
0.15.2-beta.2602122035
0.15.2-beta.2602120040
0.15.2-beta.2602110310
0.15.2-beta.2602050305
0.15.2-beta.2602040304
0.15.2-beta.2602031515
0.15.1
0.15.1-rc.2601270136
0.15.1-beta.2601310133
0.15.1-beta.2601301825
0.15.1-beta.2601300156
0.15.1-beta.2601292313
0.15.1-beta.2601280249
0.15.0
0.14.2-rc.2601230113
0.14.2-rc.2601200136
0.14.2-rc.2601130135
0.14.2-rc.2601062017
0.14.2-rc.2512230135
0.14.2-beta.2601262335
0.14.2-beta.2601240243
0.14.2-beta.2601230247
0.14.2-beta.2601220251
0.14.2-beta.2601210246
0.14.2-beta.2601170240
0.14.2-beta.2601160247
0.14.2-beta.2601140250
0.14.2-beta.2601130116
0.14.2-beta.2601100241
0.14.2-beta.2601070246
0.14.2-beta.2512310244
0.14.2-beta.2512240241
0.14.1
0.14.1-rc.2512190311
0.14.1-rc.2512162207
0.14.1-rc.2512160223
0.14.1-beta.2512200239
0.14.1-beta.2512190242
0.14.1-beta.2512170241
0.14.0
0.13.1-beta.2512140246
0.13.1-beta.2512130201
0.13.1-beta.2512120242
0.13.1-beta.2512110242
0.13.1-beta.2512110007
0.13.1-beta.2512102250
0.13.1-beta.2512102139
0.13.1-beta.2512100241
0.13.1-beta.2512051931
0.13.1-beta.2512050240
0.12.1
0.12.1-rc.2511121816
0.12.1-rc.2511110014
0.12.1-beta.2511110238
0.12.0
0.11.2-rc.2511040014
0.11.2-beta.2511050238
0.11.2-beta.2511032240
0.11.2-beta.2511032212
0.11.2-beta.2511030239
0.11.2-beta.2510310237
0.11.1
0.11.1-rc.2510292219
0.11.1-rc.2510280013
0.11.1-beta.2510300237
0.11.1-beta.2510290239
0.11.0
0.10.3-rc.2510210409
0.10.3-rc.2510210014
0.10.3-rc.6
0.10.3-beta.2510250234
0.10.3-beta.2510240234
0.10.3-beta.2510230235
0.10.3-beta.2510220237
0.10.2
0.10.2-rc.2510141845
0.10.2-beta.2510180231
0.10.2-beta.2510172041
0.10.2-beta.2510170235
0.10.2-beta.2510160235
0.10.2-beta.2510152259
0.10.2-beta.2510152125
0.10.2-beta.2510151849
0.10.2-beta.2510150236
0.10.2-beta.2510142155
0.10.2-beta.2510142023
0.10.1
0.10.0
0.10.0-rc.2510090432
0.10.0-rc.2510070136
0.10.0-beta.2510120235
0.10.0-beta.2510102126
0.10.0-beta.2510100234
0.10.0-beta.2510090234
0.10.0-beta.2510072144
0.9.6-beta.2510020232
0.9.5-rc.2510030411
0.9.5-rc.2510020138
0.9.5-beta.2510010238
0.9.4-rc.2509300113
0.9.4-beta.2509252256
0.9.3
0.9.3-beta.2509250234
0.9.3-beta.2509240232
0.9.3-beta.2509231844
0.9.2-rc.2509221315
0.9.2-beta.2509190233
0.9.2-beta.2509180233
0.9.2-beta.2509180231
0.9.2-beta.2509170231
0.9.1
0.9.1-rc.2509190115
0.9.1-rc.2509181925
0.9.1-rc.2509171909
0.9.0-rc.1
0.8.9-beta.2509161804
0.8.9-beta.2509160231
0.8.8-beta.2509130225
0.8.8-beta.2509112152
0.8.8-beta.2509112049
0.8.8-beta.2509110233
0.8.7
0.8.7-rc.2509102057
0.8.7-rc.2509102017
0.8.7-rc.2509100003
0.8.7-rc.2509092342
0.8.7-beta.2509092208
0.8.6
0.8.6-rc.2509082013
0.8.6-beta.2509040231
0.8.6-beta.2509020235
0.8.5
0.8.5-beta.2508300230
0.8.4
0.8.3
0.8.3-beta.2508280233
0.8.2
0.8.2-rc.2508282224
0.8.1
0.8.1-rc.2508271913
0.8.1-rc.2508262341
0.8.1-rc.2508251314
0.8.0
0.7.4-rc.2508222145
0.7.4-rc.2508220107
0.7.4-rc.2508220104
0.7.4-rc.2508220101
0.7.4-rc.2508220057
0.7.4-rc.2508220049
0.7.4-rc.2508220046
0.7.4-rc.2508220041
0.7.4-rc.2508220033
0.7.4-rc.2508220032
0.7.4-rc.2508220027
0.7.4-rc.2508220022
0.7.4-rc.2508220017
0.7.4-rc.2508220010
0.7.4-rc.2508220006
0.7.4-rc.2508220000
0.7.4-rc.2508212357
0.7.4-rc.2508212349
0.7.4-rc.2508202203
0.7.4-rc.2508202034
0.7.4-rc.2508191829
0.7.4-beta.2508220235
0.7.4-beta.2508210235
0.7.4-beta.2508200235
0.7.4-beta.2508190236
0.7.4-beta.2508160236
0.7.4-beta.2508150240
0.0.6-beta.7
0.0.6-beta.6
0.0.2
retired
0.0.1
retired
A performant elixir SDK for Statsig feature gates and experiments using Rustler
Current section
Files
Jump to
Current section
Files
native/statsig_elixir/src/data_store_nfi.rs
use async_trait::async_trait;
use parking_lot::Mutex;
use rustler::{
env::OwnedEnv,
types::atom::Atom,
types::binary::{Binary, OwnedBinary},
types::local_pid::LocalPid,
Encoder, Env, Error, ResourceArc, Term,
};
use statsig_rust::{
data_store_interface::{
DataStoreBytesResponse, DataStoreResponse, DataStoreTrait, RequestPath,
},
log_d, StatsigErr,
};
use std::{cell::RefCell, mem};
use tokio::sync::oneshot;
const TAG: &str = "[DataStore NFI] ";
// Track the active BEAM Env on this scheduler thread so async work can reuse it.
thread_local! {
static MANAGED_ENVS: RefCell<Vec<Env<'static>>> = const { RefCell::new(Vec::new()) };
}
pub struct ManagedEnvGuard {
active: bool,
}
impl ManagedEnvGuard {
pub fn new(env: Env<'_>) -> Self {
// SAFETY: The Env outlives the guard and is only accessed on this thread.
let env_static = unsafe { mem::transmute::<Env<'_>, Env<'static>>(env) };
MANAGED_ENVS.with(|stack| stack.borrow_mut().push(env_static));
ManagedEnvGuard { active: true }
}
}
impl Drop for ManagedEnvGuard {
fn drop(&mut self) {
if self.active {
MANAGED_ENVS.with(|stack| {
stack.borrow_mut().pop();
});
self.active = false;
}
}
}
fn current_managed_env() -> Option<Env<'static>> {
MANAGED_ENVS.with(|stack| stack.borrow().last().copied())
}
mod atoms {
rustler::atoms! {
data_store_request = "statsig_data_store_request",
initialize,
shutdown,
get,
get_bytes,
set,
set_bytes,
support_polling_updates_for,
no_payload
}
}
#[derive(rustler::NifStruct)]
#[module = "Statsig.DataStore.Reference"]
pub struct StatsigDataStoreReference {
pub pid: LocalPid,
}
#[derive(rustler::NifStruct)]
#[module = "Statsig.DataStore.Response"]
pub struct StatsigDataStoreResponse {
pub result: Option<String>,
pub time: Option<u64>,
}
#[derive(rustler::NifStruct)]
#[module = "Statsig.DataStore.BytesResponse"]
pub struct StatsigDataStoreBytesResponse<'a> {
pub result: Option<Binary<'a>>,
pub time: Option<u64>,
}
pub struct ElixirDataStore {
pid: LocalPid,
}
impl ElixirDataStore {
pub fn new(pid: LocalPid) -> Self {
Self { pid }
}
async fn request_unit(&self, request: RequestKind) -> Result<(), StatsigErr> {
match self.send_request(request).await? {
ResponsePayload::Unit => Ok(()),
_ => Err(StatsigErr::DataStoreFailure(
"Unexpected reply from Elixir data store".to_string(),
)),
}
}
async fn request_bool(&self, request: RequestKind) -> Result<bool, StatsigErr> {
match self.send_request(request).await? {
ResponsePayload::Bool(value) => Ok(value),
_ => Err(StatsigErr::DataStoreFailure(
"Unexpected reply from Elixir data store".to_string(),
)),
}
}
async fn request_data(&self, request: RequestKind) -> Result<DataStoreResponse, StatsigErr> {
match self.send_request(request).await? {
ResponsePayload::Data(value) => Ok(value),
_ => Err(StatsigErr::DataStoreFailure(
"Unexpected reply from Elixir data store".to_string(),
)),
}
}
async fn request_bytes_data(
&self,
request: RequestKind,
) -> Result<DataStoreBytesResponse, StatsigErr> {
match self.send_request(request).await? {
ResponsePayload::BytesData(value) => Ok(value),
_ => Err(StatsigErr::DataStoreFailure(
"Unexpected reply from Elixir data store".to_string(),
)),
}
}
async fn send_request(&self, request: RequestKind) -> Result<ResponsePayload, StatsigErr> {
log_d!(TAG, "Sending data store request to Elixir");
let (sender, receiver) = oneshot::channel();
let resource = ResourceArc::new(DataStoreRequestResource::new(
request.response_kind(),
request.is_bytes_request(),
sender,
));
if let Some(env) = current_managed_env() {
env.send(
&self.pid,
(
atoms::data_store_request(),
resource.clone(),
request.atom(),
request.encode_payload(env),
),
)
.map_err(|_| {
StatsigErr::DataStoreFailure("Failed to message Elixir data store".to_string())
})?;
} else {
let mut env = OwnedEnv::new();
env.send_and_clear(&self.pid, |env| {
(
atoms::data_store_request(),
resource.clone(),
request.atom(),
request.encode_payload(env),
)
})
.map_err(|_| {
StatsigErr::DataStoreFailure("Failed to message Elixir data store".to_string())
})?;
}
receiver.await.map_err(|_| {
StatsigErr::DataStoreFailure("Elixir data store did not reply to request".to_string())
})?
}
}
#[async_trait]
impl DataStoreTrait for ElixirDataStore {
async fn initialize(&self) -> Result<(), StatsigErr> {
self.request_unit(RequestKind::Initialize).await
}
async fn shutdown(&self) -> Result<(), StatsigErr> {
self.request_unit(RequestKind::Shutdown).await
}
async fn get(&self, key: &str) -> Result<DataStoreResponse, StatsigErr> {
self.request_data(RequestKind::Get(key.to_string())).await
}
async fn get_bytes(&self, key: &str) -> Result<DataStoreBytesResponse, StatsigErr> {
self.request_bytes_data(RequestKind::GetBytes(key.to_string()))
.await
}
async fn set(&self, key: &str, value: &str, time: Option<u64>) -> Result<(), StatsigErr> {
log_d!(TAG, "Setting key in data store");
self.request_unit(RequestKind::Set {
key: key.to_string(),
value: value.to_string(),
time,
})
.await
}
async fn set_bytes(
&self,
key: &str,
value: &[u8],
time: Option<u64>,
) -> Result<(), StatsigErr> {
log_d!(TAG, "Setting bytes key in data store");
self.request_unit(RequestKind::SetBytes {
key: key.to_string(),
value: value.to_vec(),
time,
})
.await
}
async fn support_polling_updates_for(&self, path: RequestPath) -> bool {
self.request_bool(RequestKind::SupportPolling(path.to_string()))
.await
.unwrap_or(false)
}
}
enum RequestKind {
Initialize,
Shutdown,
Get(String),
GetBytes(String),
Set {
key: String,
value: String,
time: Option<u64>,
},
SetBytes {
key: String,
value: Vec<u8>,
time: Option<u64>,
},
SupportPolling(String),
}
impl RequestKind {
fn atom(&self) -> Atom {
match self {
RequestKind::Initialize => atoms::initialize(),
RequestKind::Shutdown => atoms::shutdown(),
RequestKind::Get(_) => atoms::get(),
RequestKind::GetBytes(_) => atoms::get_bytes(),
RequestKind::Set { .. } => atoms::set(),
RequestKind::SetBytes { .. } => atoms::set_bytes(),
RequestKind::SupportPolling(_) => atoms::support_polling_updates_for(),
}
}
fn encode_payload<'a>(&self, env: Env<'a>) -> Term<'a> {
match self {
RequestKind::Initialize | RequestKind::Shutdown => atoms::no_payload().encode(env),
RequestKind::Get(key) => key.encode(env),
RequestKind::GetBytes(key) => key.encode(env),
RequestKind::Set { key, value, time } => {
(key.clone(), value.clone(), *time).encode(env)
}
RequestKind::SetBytes { key, value, time } => {
let mut binary =
OwnedBinary::new(value.len()).expect("failed to allocate Elixir binary");
binary.copy_from_slice(value);
let binary = Binary::from_owned(binary, env);
(key.clone(), binary, *time).encode(env)
}
RequestKind::SupportPolling(path) => path.encode(env),
}
}
fn response_kind(&self) -> ResponseKind {
match self {
RequestKind::Initialize
| RequestKind::Shutdown
| RequestKind::Set { .. }
| RequestKind::SetBytes { .. } => ResponseKind::Unit,
RequestKind::Get(_) => ResponseKind::Data,
RequestKind::GetBytes(_) => ResponseKind::BytesData,
RequestKind::SupportPolling(_) => ResponseKind::Bool,
}
}
fn is_bytes_request(&self) -> bool {
matches!(
self,
RequestKind::GetBytes(_) | RequestKind::SetBytes { .. }
)
}
}
#[derive(Clone, Copy)]
pub enum ResponseKind {
Unit,
Data,
BytesData,
Bool,
}
pub enum ResponsePayload {
Unit,
Data(DataStoreResponse),
BytesData(DataStoreBytesResponse),
Bool(bool),
}
pub struct DataStoreRequestResource {
response_kind: ResponseKind,
is_bytes_request: bool,
sender: Mutex<Option<oneshot::Sender<Result<ResponsePayload, StatsigErr>>>>,
}
impl DataStoreRequestResource {
pub fn new(
response_kind: ResponseKind,
is_bytes_request: bool,
sender: oneshot::Sender<Result<ResponsePayload, StatsigErr>>,
) -> Self {
Self {
response_kind,
is_bytes_request,
sender: Mutex::new(Some(sender)),
}
}
fn fulfill_ok(&self, payload: Term) -> Result<(), Error> {
let decoded: Result<ResponsePayload, Error> = match self.response_kind {
ResponseKind::Unit => Ok(ResponsePayload::Unit),
ResponseKind::Bool => {
payload
.decode::<bool>()
.map(ResponsePayload::Bool)
.map_err(|err| {
self.send_err(format!("Failed to decode boolean payload: {err:?}"));
err
})
}
ResponseKind::Data => payload
.decode::<StatsigDataStoreResponse>()
.map(|resp| ResponsePayload::Data(resp.into()))
.map_err(|err| {
self.send_err(format!("Failed to decode data payload: {err:?}"));
err
}),
ResponseKind::BytesData => payload
.decode::<StatsigDataStoreBytesResponse>()
.map(|resp| ResponsePayload::BytesData(resp.into()))
.map_err(|err| {
self.send_err(format!("Failed to decode bytes data payload: {err:?}"));
err
}),
};
match decoded {
Ok(value) => {
self.send_ok(value);
Ok(())
}
Err(err) => Err(err),
}
}
fn fulfill_err(&self, reason: String) -> Result<(), Error> {
self.send_err(reason);
Ok(())
}
fn send_ok(&self, payload: ResponsePayload) {
if let Some(sender) = self.sender.lock().take() {
let _ = sender.send(Ok(payload));
}
}
fn send_err(&self, reason: String) {
if let Some(sender) = self.sender.lock().take() {
let error = if self.is_bytes_request && reason == "BytesNotImplemented" {
StatsigErr::BytesNotImplemented
} else {
StatsigErr::DataStoreFailure(reason)
};
let _ = sender.send(Err(error));
}
}
}
impl From<StatsigDataStoreResponse> for DataStoreResponse {
fn from(value: StatsigDataStoreResponse) -> Self {
DataStoreResponse {
result: value.result,
time: value.time,
}
}
}
impl From<StatsigDataStoreBytesResponse<'_>> for DataStoreBytesResponse {
fn from(value: StatsigDataStoreBytesResponse) -> Self {
DataStoreBytesResponse {
result: value.result.map(|result| result.as_slice().to_vec()),
time: value.time,
}
}
}
#[rustler::nif]
pub fn data_store_reply(
request: ResourceArc<DataStoreRequestResource>,
payload: Term,
) -> Result<(), Error> {
request.fulfill_ok(payload)
}
#[rustler::nif]
pub fn data_store_reply_error(
request: ResourceArc<DataStoreRequestResource>,
reason: String,
) -> Result<(), Error> {
request.fulfill_err(reason)
}