relay_server/services/projects/source/
redis.rs1use 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 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 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 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 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{}"; let result = parse_redis_response(raw_response);
153 assert!(result.is_ok(), "{result:?}");
154 }
155}