use serde::{Deserialize, Serialize};
use crate::graphql::client_manifest::DISTRIBUTED_CLIENT_PROTOCOL_VERSION;
use super::OpaqueProtocolToken;
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum DistributedCommandState {
InProgress,
Succeeded,
SucceededPendingProjection,
Atomic,
Rejected,
ProjectionFailed,
Expired,
Unknown,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum DistributedCommandConsistency {
Succeeded,
Eventual,
Atomic,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum DistributedProjectionDisposition {
Revalidate,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DistributedProjectionExpectation {
pub(crate) projection: String,
pub(crate) model: String,
pub(crate) scope_token: OpaqueProtocolToken,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DistributedCommandMetadata {
pub(crate) command_id: String,
pub(crate) causation_id: String,
pub(crate) state: DistributedCommandState,
pub(crate) consistency: DistributedCommandConsistency,
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) projection_disposition: Option<DistributedProjectionDisposition>,
pub(crate) expects: Vec<DistributedProjectionExpectation>,
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) projection: Option<super::CommandProjectionMetadataV1>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub(crate) observations: Vec<DistributedProjectionObservation>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub(crate) records: Vec<DistributedRecordRevision>,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DistributedRecordRevision {
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) path: Option<Vec<String>>,
pub(crate) model: String,
pub(crate) scope_token: OpaqueProtocolToken,
pub(crate) incarnation: String,
pub(crate) revision: String,
pub(crate) tombstone: bool,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DistributedProjectionObservation {
pub(crate) causation_id: String,
pub(crate) projection: String,
pub(crate) model: String,
pub(crate) scope_token: OpaqueProtocolToken,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DistributedLiveCursor {
pub(crate) projection: String,
pub(crate) position: String,
pub(crate) token: OpaqueProtocolToken,
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) enum RequestedLiveResume {
#[default]
Absent,
Invalid,
Cursors(Vec<DistributedLiveCursor>),
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DistributedIndexRevision {
pub(crate) projection: String,
pub(crate) scope_token: OpaqueProtocolToken,
pub(crate) position: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) resume: Option<DistributedLiveCursor>,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DistributedQuerySnapshot {
pub(crate) scope_token: OpaqueProtocolToken,
pub(crate) records_complete: bool,
pub(crate) indexes_comparable: bool,
pub(crate) records: Vec<DistributedRecordRevision>,
pub(crate) indexes: Vec<DistributedIndexRevision>,
pub(crate) observations: Vec<DistributedProjectionObservation>,
}
impl DistributedQuerySnapshot {
pub(super) fn discard_incomparable_index_evidence(&mut self) {
if self.indexes_comparable {
return;
}
self.indexes.clear();
self.observations.clear();
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DistributedLiveMetadata {
pub(crate) supported: bool,
pub(crate) reset: bool,
pub(crate) cursors: Vec<DistributedLiveCursor>,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DistributedTrustedPreset {
pub(crate) name: String,
pub(crate) codec: String,
pub(crate) value: serde_json::Value,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DistributedEnvelopeV1 {
pub(crate) protocol_version: u32,
pub(crate) schema_hash: String,
pub(crate) authorization_generation: String,
pub(crate) cache_scope: OpaqueProtocolToken,
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) generation: Option<DistributedGenerationEnvelope>,
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) operation: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) command: Option<DistributedCommandMetadata>,
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) snapshot: Option<DistributedQuerySnapshot>,
#[serde(skip_serializing_if = "Option::is_none")]
pub(crate) live: Option<DistributedLiveMetadata>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub(crate) trusted_presets: Vec<DistributedTrustedPreset>,
}
impl DistributedEnvelopeV1 {
pub(crate) fn new(
schema_hash: impl Into<String>,
authorization_generation: impl Into<String>,
cache_scope: OpaqueProtocolToken,
operation: Option<String>,
) -> Self {
Self {
protocol_version: DISTRIBUTED_CLIENT_PROTOCOL_VERSION,
schema_hash: schema_hash.into(),
authorization_generation: authorization_generation.into(),
cache_scope,
generation: DistributedGenerationEnvelope::from_environment(),
operation,
command: None,
snapshot: None,
live: None,
trusted_presets: Vec::new(),
}
}
pub(crate) fn with_trusted_presets(
mut self,
trusted_presets: Vec<DistributedTrustedPreset>,
) -> Self {
self.trusted_presets = trusted_presets;
self
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DistributedGenerationEnvelope {
version: u32,
generation_id: String,
release_id: String,
#[serde(skip_serializing_if = "Option::is_none")]
topology_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
compatibility_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
member_id: Option<String>,
}
impl DistributedGenerationEnvelope {
fn from_environment() -> Option<Self> {
let generation_id = std::env::var("DISTRIBUTED_GENERATION_ID").ok()?;
let release_id = std::env::var("DISTRIBUTED_RELEASE_ID").ok()?;
if !bounded_identity(&generation_id) || !bounded_identity(&release_id) {
return None;
}
Some(Self {
version: 1,
generation_id,
release_id,
topology_id: optional_environment_identity("DISTRIBUTED_TOPOLOGY_ID"),
compatibility_id: optional_environment_identity("DISTRIBUTED_COMPATIBILITY_ID"),
member_id: optional_environment_identity("DISTRIBUTED_MEMBER_ID"),
})
}
}
fn optional_environment_identity(name: &str) -> Option<String> {
std::env::var(name)
.ok()
.filter(|value| bounded_identity(value))
}
fn bounded_identity(value: &str) -> bool {
!value.is_empty()
&& value.len() <= 512
&& value == value.trim()
&& !value.chars().any(char::is_control)
}
#[cfg(test)]
mod generation_tests {
use super::*;
#[test]
fn generation_envelope_is_camel_case_and_contains_no_process_environment() {
let envelope = DistributedGenerationEnvelope {
version: 1,
generation_id: "sha256:generation".into(),
release_id: "sha256:release".into(),
topology_id: Some("sha256:topology".into()),
compatibility_id: Some("sha256:compatibility".into()),
member_id: Some("api".into()),
};
assert_eq!(
serde_json::to_value(envelope).unwrap(),
serde_json::json!({
"version": 1,
"generationId": "sha256:generation",
"releaseId": "sha256:release",
"topologyId": "sha256:topology",
"compatibilityId": "sha256:compatibility",
"memberId": "api"
})
);
}
#[test]
fn generation_environment_identities_are_bounded() {
assert!(bounded_identity("sha256:stable"));
assert!(!bounded_identity(""));
assert!(!bounded_identity(" leading"));
assert!(!bounded_identity("line\nbreak"));
assert!(!bounded_identity(&"x".repeat(513)));
}
}