Skip to main content

relay_server/services/projects/source/
redis.rs

1use relay_base_schema::project::ProjectKey;
2use relay_config::{Config, ConfigSnapshot};
3use relay_redis::{AsyncRedisClient, RedisError};
4use relay_statsd::metric;
5use std::fmt::Debug;
6use std::sync::Arc;
7
8use crate::services::projects::project::{IncomingProjectState, ProjectState, Revision};
9use crate::services::projects::source::SourceProjectState;
10use crate::statsd::{RelayCounters, RelayDistributions, RelayTimers};
11use relay_redis::redis::cmd;
12
13#[derive(Clone, Debug)]
14pub struct RedisProjectSource {
15    config: Arc<Config>,
16    redis: AsyncRedisClient,
17}
18
19#[derive(Debug, thiserror::Error)]
20pub enum RedisProjectError {
21    #[error("failed to parse projectconfig from redis")]
22    Parsing(#[from] serde_json::Error),
23
24    #[error("failed to talk to redis")]
25    Redis(#[from] RedisError),
26}
27
28fn parse_redis_response(raw_response: &[u8]) -> Result<IncomingProjectState, RedisProjectError> {
29    let decompression_result = metric!(timer(RelayTimers::ProjectStateDecompression), {
30        zstd::decode_all(raw_response)
31    });
32
33    let decoded_response = match &decompression_result {
34        Ok(decoded) => {
35            metric!(
36                distribution(RelayDistributions::ProjectStateSizeBytesCompressed) =
37                    raw_response.len() as f64
38            );
39            metric!(
40                distribution(RelayDistributions::ProjectStateSizeBytesDecompressed) =
41                    decoded.len() as f64
42            );
43            decoded.as_slice()
44        }
45        // If decoding fails, assume uncompressed payload and try again
46        Err(_) => raw_response,
47    };
48
49    Ok(serde_json::from_slice(decoded_response)?)
50}
51
52impl RedisProjectSource {
53    pub fn new(config: Arc<Config>, redis: AsyncRedisClient) -> Self {
54        RedisProjectSource { config, redis }
55    }
56
57    /// Fetches a project config from Redis.
58    ///
59    /// The returned project state is [`ProjectState::Pending`] if the requested project config is not
60    /// stored in Redis.
61    pub async fn get_config_if_changed(
62        &self,
63        key: ProjectKey,
64        revision: Revision,
65    ) -> Result<SourceProjectState, RedisProjectError> {
66        let mut connection = self.redis.get_connection().await?;
67        let config = self.config.current();
68
69        // Only check for the revision if we were passed a revision.
70        if let Some(revision) = revision.as_str() {
71            let current_revision: Option<String> = cmd("GET")
72                .arg(get_redis_rev_key(&config, key))
73                .query_async(&mut connection)
74                .await
75                .map_err(RedisError::Redis)?;
76
77            relay_log::trace!(
78                "Redis revision {current_revision:?}, requested revision {revision:?}"
79            );
80            if current_revision.as_deref() == Some(revision) {
81                metric!(
82                    counter(RelayCounters::ProjectStateRedis) += 1,
83                    hit = "revision",
84                );
85                return Ok(SourceProjectState::NotModified);
86            }
87        }
88
89        let raw_response_opt: Option<Vec<u8>> = cmd("GET")
90            .arg(get_redis_project_config_key(&config, key))
91            .query_async(&mut connection)
92            .await
93            .map_err(RedisError::Redis)?;
94
95        let Some(response) = raw_response_opt else {
96            metric!(
97                counter(RelayCounters::ProjectStateRedis) += 1,
98                hit = "false"
99            );
100            return Ok(SourceProjectState::New(ProjectState::Pending));
101        };
102
103        let response = ProjectState::from(parse_redis_response(response.as_slice())?);
104
105        // If we were passed a revision, check if we just loaded the same revision from Redis.
106        //
107        // We always want to keep the old revision alive if possible, since the already loaded
108        // version has already initialized caches.
109        //
110        // While this is theoretically possible this should always been handled using the above revision
111        // check using the additional Redis key.
112        if response.revision() == revision {
113            metric!(
114                counter(RelayCounters::ProjectStateRedis) += 1,
115                hit = "project_config_revision"
116            );
117            Ok(SourceProjectState::NotModified)
118        } else {
119            metric!(
120                counter(RelayCounters::ProjectStateRedis) += 1,
121                hit = "project_config"
122            );
123            Ok(SourceProjectState::New(response))
124        }
125    }
126}
127
128fn get_redis_project_config_key(config: &ConfigSnapshot, key: ProjectKey) -> String {
129    let prefix = config.projectconfig_cache_prefix();
130    format!("{prefix}:{key}")
131}
132
133fn get_redis_rev_key(config: &ConfigSnapshot, key: ProjectKey) -> String {
134    let prefix = config.projectconfig_cache_prefix();
135    format!("{prefix}:{key}.rev")
136}
137
138#[cfg(test)]
139mod tests {
140    use super::*;
141
142    #[test]
143    fn test_parse_redis_response() {
144        let raw_response = b"{}";
145        let result = parse_redis_response(raw_response);
146        assert!(result.is_ok());
147    }
148
149    #[test]
150    fn test_parse_redis_response_compressed() {
151        let raw_response = b"(\xb5/\xfd \x02\x11\x00\x00{}"; // As dumped by python zstandard library
152        let result = parse_redis_response(raw_response);
153        assert!(result.is_ok(), "{result:?}");
154    }
155}