From 2b6515a4bcdb78c966e61f9f13b365d62f8bab3f Mon Sep 17 00:00:00 2001 From: Leonid Kaganov Date: Thu, 18 Sep 2025 21:41:17 +0300 Subject: [PATCH] fix: resolve ambiguity in Redis write conditions Signed-off-by: Leonid Kaganov --- .gitignore | 1 + Cargo.lock | 261 ++++++++++++++++++++++++++++++++++++++- Cargo.toml | 7 +- scripts/TEST_HTTP_API.sh | 2 +- src/config/default.toml | 2 +- src/main.rs | 9 ++ src/redis.rs | 209 ++++++++++++++++--------------- tests/rest_api.rs | 43 +++++++ 8 files changed, 431 insertions(+), 103 deletions(-) create mode 100644 tests/rest_api.rs diff --git a/.gitignore b/.gitignore index 20e2e71e9f..67f4d268c9 100644 --- a/.gitignore +++ b/.gitignore @@ -7,6 +7,7 @@ commit.sh /src/GO.sh /src/GOT.sh GO.sh +TEST.sh DROP_DB.sh TODO.txt DOCKER.sh diff --git a/Cargo.lock b/Cargo.lock index 7007dd680d..bf9de1a2ac 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -80,7 +80,7 @@ dependencies = [ "flate2", "foldhash", "futures-core", - "h2", + "h2 0.3.27", "http 0.2.12", "httparse", "httpdate", @@ -618,6 +618,16 @@ dependencies = [ "version_check", ] +[[package]] +name = "core-foundation" +version = "0.9.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "91e195e091a93c46f7102ec7818a2aa394e1e1771c3ab4825963fa03e45afb8f" +dependencies = [ + "core-foundation-sys", + "libc", +] + [[package]] name = "core-foundation-sys" version = "0.8.7" @@ -852,12 +862,28 @@ dependencies = [ "typeid", ] +[[package]] +name = "errno" +version = "0.3.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "976dd42dc7e85965fe702eb8164f21f450704bdde31faefd6471dba214cb594e" +dependencies = [ + "libc", + "windows-sys 0.59.0", +] + [[package]] name = "fallible-iterator" version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4443176a9f2c162692bd3d352d745ef9413eec5782a80d8fd6f8a1ac692a07f7" +[[package]] +name = "fastrand" +version = "2.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "37909eebbb50d72f9059c3b6d82c0463f2ff062c9e95845c43a6c9c0355411be" + [[package]] name = "flate2" version = "1.1.2" @@ -880,6 +906,21 @@ version = "0.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d9c4f5dac5e15c24eb999c26181a6ca40b39fe946cbe4c263c7209467bc83af2" +[[package]] +name = "foreign-types" +version = "0.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f6f339eb8adc052cd2ca78910fda869aefa38d22d5cb648e6485e4d3fc06f3b1" +dependencies = [ + "foreign-types-shared", +] + +[[package]] +name = "foreign-types-shared" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "00b0228411908ca8685dba7fc2cdd70ec9990a6e753e89b6ac91a84c40fbaf4b" + [[package]] name = "form_urlencoded" version = "1.2.2" @@ -1069,6 +1110,25 @@ dependencies = [ "tracing", ] +[[package]] +name = "h2" +version = "0.4.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f3c0b69cfcb4e1b9f1bf2f53f95f766e4661169728ec61cd3fe5a0166f2d1386" +dependencies = [ + "atomic-waker", + "bytes", + "fnv", + "futures-core", + "futures-sink", + "http 1.3.1", + "indexmap 2.10.0", + "slab", + "tokio", + "tokio-util", + "tracing", +] + [[package]] name = "hashbrown" version = "0.12.3" @@ -1197,6 +1257,7 @@ dependencies = [ "md5", "redis", "refinery", + "reqwest", "secrecy", "serde", "serde_json", @@ -1258,6 +1319,7 @@ dependencies = [ "bytes", "futures-channel", "futures-core", + "h2 0.4.12", "http 1.3.1", "http-body", "httparse", @@ -1286,6 +1348,22 @@ dependencies = [ "webpki-roots 1.0.2", ] +[[package]] +name = "hyper-tls" +version = "0.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "70206fc6890eaca9fde8a0bf71caa2ddfc9fe045ac9e5c70df101a7dbde866e0" +dependencies = [ + "bytes", + "http-body-util", + "hyper", + "hyper-util", + "native-tls", + "tokio", + "tokio-native-tls", + "tower-service", +] + [[package]] name = "hyper-util" version = "0.1.16" @@ -1305,9 +1383,11 @@ dependencies = [ "percent-encoding", "pin-project-lite", "socket2 0.6.0", + "system-configuration", "tokio", "tower-service", "tracing", + "windows-registry", ] [[package]] @@ -1595,6 +1675,12 @@ dependencies = [ "redox_syscall 0.5.17", ] +[[package]] +name = "linux-raw-sys" +version = "0.9.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cd945864f07fe9f5371a27ad7b52a172b4b499999f1d97574c9fa68373937e12" + [[package]] name = "litemap" version = "0.8.0" @@ -1689,6 +1775,23 @@ dependencies = [ "windows-sys 0.59.0", ] +[[package]] +name = "native-tls" +version = "0.2.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "87de3442987e9dbec73158d5c715e7ad9072fda936bb03d19d7fa10e00520f0e" +dependencies = [ + "libc", + "log", + "openssl", + "openssl-probe", + "openssl-sys", + "schannel", + "security-framework", + "security-framework-sys", + "tempfile", +] + [[package]] name = "nonzero_ext" version = "0.3.0" @@ -1754,6 +1857,50 @@ version = "1.21.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "42f5e15c9953c5e4ccceeb2e7382a716482c34515315f7b03532b8b4e8393d2d" +[[package]] +name = "openssl" +version = "0.10.72" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fedfea7d58a1f73118430a55da6a286e7b044961736ce96a16a17068ea25e5da" +dependencies = [ + "bitflags 2.9.2", + "cfg-if", + "foreign-types", + "libc", + "once_cell", + "openssl-macros", + "openssl-sys", +] + +[[package]] +name = "openssl-macros" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a948666b637a0f465e8564c73e89d4dde00d72d4d473cc972f390fc3dcee7d9c" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "openssl-probe" +version = "0.1.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d05e27ee213611ffe7d6348b942e8f942b37114c00cc03cec254295a4a17852e" + +[[package]] +name = "openssl-sys" +version = "0.9.109" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "90096e2e47630d78b7d1c20952dc621f957103f8bc2c8359ec81290d75238571" +dependencies = [ + "cc", + "libc", + "pkg-config", + "vcpkg", +] + [[package]] name = "ordered-multimap" version = "0.7.3" @@ -2290,15 +2437,20 @@ checksum = "d429f34c8092b2d42c7c93cec323bb4adeb7c67698f70839adec842ec10c7ceb" dependencies = [ "base64 0.22.1", "bytes", + "encoding_rs", "futures-core", + "h2 0.4.12", "http 1.3.1", "http-body", "http-body-util", "hyper", "hyper-rustls", + "hyper-tls", "hyper-util", "js-sys", "log", + "mime", + "native-tls", "percent-encoding", "pin-project-lite", "quinn", @@ -2309,6 +2461,7 @@ dependencies = [ "serde_urlencoded", "sync_wrapper", "tokio", + "tokio-native-tls", "tokio-rustls 0.26.2", "tower", "tower-http", @@ -2462,6 +2615,19 @@ version = "2.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "357703d41365b4b27c590e3ed91eabb1b663f07c4c084095e60cbed4362dff0d" +[[package]] +name = "rustix" +version = "1.0.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d97817398dd4bb2e6da002002db259209759911da105da92bec29ccb12cf58bf" +dependencies = [ + "bitflags 2.9.2", + "errno", + "libc", + "linux-raw-sys", + "windows-sys 0.59.0", +] + [[package]] name = "rustls" version = "0.20.9" @@ -2530,6 +2696,15 @@ dependencies = [ "winapi-util", ] +[[package]] +name = "schannel" +version = "0.1.27" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1f29ebaa345f945cec9fbbc532eb307f0fdad8161f281b6369539c8d84876b3d" +dependencies = [ + "windows-sys 0.59.0", +] + [[package]] name = "schemars" version = "0.9.0" @@ -2580,6 +2755,29 @@ dependencies = [ "zeroize", ] +[[package]] +name = "security-framework" +version = "2.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "897b2245f0b511c87893af39b033e5ca9cce68824c4d7e7630b5a1d339658d02" +dependencies = [ + "bitflags 2.9.2", + "core-foundation", + "core-foundation-sys", + "libc", + "security-framework-sys", +] + +[[package]] +name = "security-framework-sys" +version = "2.14.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "49db231d56a190491cb4aeda9527f1ad45345af50b0851622a7adb8c03b01c32" +dependencies = [ + "core-foundation-sys", + "libc", +] + [[package]] name = "serde" version = "1.0.219" @@ -2892,6 +3090,40 @@ dependencies = [ "syn", ] +[[package]] +name = "system-configuration" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3c879d448e9d986b661742763247d3693ed13609438cf3d006f51f5368a5ba6b" +dependencies = [ + "bitflags 2.9.2", + "core-foundation", + "system-configuration-sys", +] + +[[package]] +name = "system-configuration-sys" +version = "0.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8e1d1b10ced5ca923a1fcb8d03e96b8d3268065d724548c0211415ff6ac6bac4" +dependencies = [ + "core-foundation-sys", + "libc", +] + +[[package]] +name = "tempfile" +version = "3.19.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7437ac7763b9b123ccf33c338a5cc1bac6f69b45a136c19bdd8a65e3916435bf" +dependencies = [ + "fastrand", + "getrandom 0.3.3", + "once_cell", + "rustix", + "windows-sys 0.59.0", +] + [[package]] name = "thiserror" version = "1.0.69" @@ -3037,6 +3269,16 @@ dependencies = [ "syn", ] +[[package]] +name = "tokio-native-tls" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bbae76ab933c85776efabc971569dd6119c580d8f5d448769dec1764bf796ef2" +dependencies = [ + "native-tls", + "tokio", +] + [[package]] name = "tokio-postgres" version = "0.7.13" @@ -3449,6 +3691,12 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" +[[package]] +name = "vcpkg" +version = "0.2.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426" + [[package]] name = "version_check" version = "0.9.5" @@ -3726,6 +3974,17 @@ version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5e6ad25900d524eaabdbbb96d20b4311e1e7ae1699af4fb28c17ae66c80d798a" +[[package]] +name = "windows-registry" +version = "0.5.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5b8a9ed28765efc97bbc954883f4e6796c33a06546ebafacbabee9696967499e" +dependencies = [ + "windows-link", + "windows-result", + "windows-strings", +] + [[package]] name = "windows-result" version = "0.3.4" diff --git a/Cargo.toml b/Cargo.toml index e6ebdc268e..058bf5516f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,15 +1,15 @@ [package] name = "hulypulse" -version = "0.1.21" +version = "0.1.23" edition = "2024" [dependencies] -tokio = { version = "1", features = ["full"] } +tokio = { version = "1", features = ["full", "macros", "rt-multi-thread"] } tracing = "0.1.41" tracing-subscriber = "0.3.19" anyhow = "1.0.97" config = "0.15.4" -serde = "1.0.219" +serde = { version = "1.0.219", features = ["derive"] } actix = "0.13.5" actix-web = "4.10.2" actix-cors = "0.7.1" @@ -32,6 +32,7 @@ hulyrs = { git = "https://github.com/hcengineering/hulyrs.git", features = [ secrecy = "0.10.3" tokio-stream = "0.1" strum = { version = "0.27.2", features = ["derive"] } +reqwest = { version = "0.12.23", features = ["json"] } [[bin]] name = "hulypulse" diff --git a/scripts/TEST_HTTP_API.sh b/scripts/TEST_HTTP_API.sh index a07394a694..1b8b9816dd 100755 --- a/scripts/TEST_HTTP_API.sh +++ b/scripts/TEST_HTTP_API.sh @@ -21,7 +21,7 @@ get "00000000-0000-0000-0000-000000000001/TESTS/" -exit +#exit diff --git a/src/config/default.toml b/src/config/default.toml index 9050797929..b2039f57f6 100644 --- a/src/config/default.toml +++ b/src/config/default.toml @@ -1,4 +1,4 @@ -bind_port = 8095 +bind_port = 8099 bind_host = "0.0.0.0" token_secret = "secret" diff --git a/src/main.rs b/src/main.rs index e132eeac6c..123175bda0 100644 --- a/src/main.rs +++ b/src/main.rs @@ -67,7 +67,9 @@ async fn extract_claims( token: Option, } + println!("CONFIG.no_authorization = {}", CONFIG.no_authorization); if !CONFIG.no_authorization { + println!("Extracting claims..."); let query = request.extract::>().await?.into_inner(); let claims = if let Some(token) = query.token { @@ -175,7 +177,14 @@ async fn main() -> anyhow::Result<()> { ) // WebSocket .route( "/status", + // web::get().to(|req: actix_web::HttpRequest| async move { + // for (name, value) in req.headers() { + // println!("HEADER {:?}: {:?}", name, value); + // } + // HttpResponse::Ok().finish() + // }), web::get().to({ + // println!("S T A T U S"); move |hub_state: web::Data>>, db_backend: web::Data| { let hub_state = hub_state.clone(); async move { diff --git a/src/redis.rs b/src/redis.rs index 8bfd409c24..ccd2ef1732 100644 --- a/src/redis.rs +++ b/src/redis.rs @@ -48,6 +48,8 @@ use redis::{ }; use serde::Serialize; +static MAX_LOOP_COUNT: usize = 1000; // to avoid infinite loops + #[derive(Debug, Serialize)] pub struct RedisArray { pub key: String, @@ -85,28 +87,16 @@ pub fn deprecated_symbol_error(s: &str) -> redis::RedisResult<()> { } } -// if CONFIG.memory_mode == Some(true) { - -// // memory_status -// let map = hub_state.read().await; -// let memory_keys = format!("{} keys in memory", map.len()); -// let memory_bytes = format!("{} bytes used", map.values().map(|v| v.data.len()).sum::()); -// format!("{} keys, {} bytes", memory_keys, memory_bytes) -// } else { -// let mut conn = db_backend.redis_connection.lock().await; -// }; - /// redis_info(&connection) pub async fn redis_info(conn: &mut MultiplexedConnection) -> redis::RedisResult { let info: String = redis::cmd("INFO").query_async(conn).await?; - // Разбираем её построчно let mut redis_keys: Option = None; let mut redis_bytes: Option = None; for line in info.lines() { if line.starts_with("db0:") { - // db0:keys=152,expires=10,avg_ttl=456789 + // parsing: db0:keys=152,expires=10,avg_ttl=456789 if let Some(keys_part) = line.split(',').find(|s| s.starts_with("keys=")) { if let Some(val) = keys_part.strip_prefix("keys=") { redis_keys = val.parse::().ok(); @@ -200,14 +190,11 @@ pub async fn redis_read( }; let ttl: i64 = redis::cmd("TTL").arg(key).query_async(conn).await?; - if ttl == -1 { - return error(500, "TTL not set"); - } - if ttl == -2 { - return error(500, "Key not found"); - } - if ttl < 0 { - return error(500, "Unknown TTL error"); + match ttl { + -1 => return error(500, "TTL not set"), + -2 => return error(500, "Key not found"), + x if x < 0 => return error(500, "Unknown TTL error"), + _ => {} // ttl >= 0, ок } Ok(Some(RedisArray { @@ -218,16 +205,7 @@ pub async fn redis_read( })) } -/// TTL sec -/// redis_save(&mut conn, "key", "val", Some(Ttl::Sec(300)), Some(SaveMode::Insert)).await?; -/// -/// TTL at -/// let at_unixtime: u64 = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs() + 600; -/// redis_save(&mut conn, "key", "val", Some(Ttl::At(at_unixtime)), Some(SaveMode::Update)).await?; -/// -/// w/o TTL (CONFIG.max_ttl) -/// redis_save(&mut conn, "key", "val", None, None).await?; - +/// redis_save(&connection,key,value,[ttl?],[mode?]) pub async fn redis_save( conn: &mut MultiplexedConnection, key: &str, @@ -272,30 +250,49 @@ pub async fn redis_save( return error(412, "TTL exceeds MAX_TTL"); } - let mut cmd = redis::cmd("SET"); - cmd.arg(key).arg(value).arg("EX").arg(sec); - - // Mode variants let mode = mode.unwrap_or(SaveMode::Upsert); match mode { - SaveMode::Upsert => {} // none - - SaveMode::Insert => { - cmd.arg("NX"); - } // if NOT Exist - - SaveMode::Update => { - cmd.arg("XX"); - } // if Exist + SaveMode::Upsert | SaveMode::Insert | SaveMode::Update => { + let mut cmd = redis::cmd("SET"); + cmd.arg(key).arg(value).arg("EX").arg(sec); + match mode { + SaveMode::Insert => { + cmd.arg("NX"); + } // if NOT Exist + SaveMode::Update => { + cmd.arg("XX"); + } // if Exist + _ => {} + }; + let result: Option = cmd.query_async(conn).await?; + if result.is_none() { + return match mode { + SaveMode::Insert => error(412, "Insert: key already exists"), + SaveMode::Update => error(404, "Update: key does not exist"), + _ => error(500, "Unexpected Redis SET failure"), + }; + } + Ok(()) + } SaveMode::Equal(ref expected_md5) => { - // if md5 === actual_md5 - let current_value: Option = - redis::cmd("GET").arg(key).query_async(conn).await?; - if let Some(existing) = current_value { + let mut loop_count = 0; + loop { + let _: () = redis::cmd("WATCH").arg(key).query_async(conn).await?; + + let current: Option = redis::cmd("GET").arg(key).query_async(conn).await?; + let existing = match current { + None => { + let _: () = redis::cmd("UNWATCH").query_async(conn).await?; + return error(404, "Equal: key does not exist"); + } + Some(v) => v, + }; + // check md5 let actual_md5 = hex::encode(md5::compute(&existing).0); if &actual_md5 != expected_md5 { + let _: () = redis::cmd("UNWATCH").query_async(conn).await?; return error( 412, format!( @@ -304,31 +301,38 @@ pub async fn redis_save( ), ); } - } else { - return error(404, "Equal: key does not exist"); + + // MULTI/EXEC + let mut pipe = redis::pipe(); + pipe.atomic() + .cmd("SET") + .arg(key) + .arg(value.to_redis_args()) + .arg("EX") + .arg(sec); + + let result: Option = pipe.query_async(conn).await?; + if result.is_some() { + break; + } + // None -> key was changed -> repeat loop + loop_count += 1; + if loop_count > MAX_LOOP_COUNT { + let _: () = redis::cmd("UNWATCH").query_async(conn).await?; + return error(500, "Something wrong: too many retries on Equal mode"); + } } + + Ok(()) } } - - // execute - let result: Option = cmd.query_async(conn).await?; - - if result.is_none() { - match mode { - SaveMode::Insert => return error(412, "Insert: key already exists"), - SaveMode::Update => return error(404, "Update: key does not exist"), - _ => return error(500, "Unexpected Redis SET failure"), - } - } - - Ok(()) } /// redis_delete(&connection,key) pub async fn redis_delete( conn: &mut MultiplexedConnection, key: &str, - mode: Option, // <— добавили + mode: Option, ) -> RedisResult { deprecated_symbol_error(key)?; @@ -339,37 +343,58 @@ pub async fn redis_delete( let mode = mode.unwrap_or(SaveMode::Upsert); match mode { + SaveMode::Update | SaveMode::Upsert => { + let deleted: i32 = redis::cmd("DEL").arg(key).query_async(conn).await?; + return Ok(deleted > 0); + } + SaveMode::Equal(ref expected_md5) => { - let current: Option = redis::cmd("GET").arg(key).query_async(conn).await?; - match current { - None => return error(404, "Equal: key does not exist"), - Some(val) => { - let actual_md5 = hex::encode(md5::compute(&val).0); - if &actual_md5 != expected_md5 { - return error( - 412, - format!( - "md5 mismatch, current: {}, expected: {}", - actual_md5, expected_md5 - ), - ); + let mut loop_count = 0; + loop { + let _: () = redis::cmd("WATCH").arg(key).query_async(conn).await?; + + let current: Option = redis::cmd("GET").arg(key).query_async(conn).await?; + let existing = match current { + None => { + let _: () = redis::cmd("UNWATCH").query_async(conn).await?; + return error(404, "Equal: key does not exist"); } + Some(val) => val, + }; + + // check md5 + let actual_md5 = hex::encode(md5::compute(&existing).0); + if &actual_md5 != expected_md5 { + let _: () = redis::cmd("UNWATCH").query_async(conn).await?; + return error( + 412, + format!( + "md5 mismatch, current: {}, expected: {}", + actual_md5, expected_md5 + ), + ); + } + + let mut pipe = redis::pipe(); + pipe.atomic().cmd("DEL").arg(key); + + let deleted: Option = pipe.query_async(conn).await?; + if let Some(n) = deleted { + return Ok(n > 0); + } + // None -> key was changed -> repeat loop + loop_count += 1; + if loop_count > MAX_LOOP_COUNT { + let _: () = redis::cmd("UNWATCH").query_async(conn).await?; + return error(500, "Something wrong: too many retries on Equal mode"); } } } + SaveMode::Insert => { return error(412, "Insert mode is not supported for delete"); } - SaveMode::Update | SaveMode::Upsert => {} } - - let deleted: i32 = redis::cmd("DEL").arg(key).query_async(conn).await?; - - if deleted == 0 && matches!(mode, SaveMode::Equal(_)) { - return error(404, "Delete: key does not exist"); - } - - Ok(deleted > 0) } impl TryFrom for RedisEvent { @@ -389,7 +414,7 @@ impl TryFrom for RedisEvent { } }; - // "__keyevent@0__:set" → event="set", db=0; payload = key + // parsing: "__keyevent@0__:set" → event="set", db=0; payload = key let event = channel.rsplit(':').next().unwrap_or(""); let message = match event { "set" => RedisEventAction::Set, @@ -399,13 +424,6 @@ impl TryFrom for RedisEvent { other => RedisEventAction::Other(other.to_string()), }; - // let db = channel - // .find('@') - // .and_then(|at| channel.get(at + 1..)) - // .and_then(|rest| rest.find("__:").map(|end| &rest[..end])) - // .and_then(|s| s.parse::().ok()) - // .unwrap_or(0); - Ok(RedisEvent { // db, key: payload.clone(), @@ -416,7 +434,6 @@ impl TryFrom for RedisEvent { pub async fn receiver( redis_client: Client, - // hub: HubServiceHandle hub_state: Arc>, ) -> anyhow::Result<()> { let mut redis = redis_client.get_multiplexed_async_connection().await?; @@ -443,8 +460,6 @@ pub async fn receiver( while let Some(message) = messages.next().await { match RedisEvent::try_from(message) { Ok(ev) => { - // debug!("redis event: {ev:#?}"); - push_event(&hub_state, &mut redis, ev).await; } Err(e) => { diff --git a/tests/rest_api.rs b/tests/rest_api.rs new file mode 100644 index 0000000000..80c5d9e82d --- /dev/null +++ b/tests/rest_api.rs @@ -0,0 +1,43 @@ +use reqwest::StatusCode; +use std::env; + +fn server_url() -> String { + env::var("TEST_SERVER_URL").unwrap_or_else(|_| "http://127.0.0.1/api".to_string()) +} + +#[tokio::test] +async fn put_and_get() { + let base = server_url(); + println!("Use url: {}", &base); + + let client = reqwest::Client::new(); + + // PUT значение + let put_resp = client + .put(&format!( + "{}/00000000-0000-0000-0000-000000000001/TESTS/key1", + base + )) + .header("HULY-TTL", "2") + .body("Value_1") + .send() + .await + .unwrap(); + + assert_eq!(put_resp.status(), StatusCode::OK); + + // GET значение + let get_resp = client + .get(&format!( + "{}/00000000-0000-0000-0000-000000000001/TESTS/key1", + base + )) + .send() + .await + .unwrap(); + + assert_eq!(get_resp.status(), StatusCode::OK); + + let text = get_resp.text().await.unwrap(); + assert!(text.contains("Value_1")); +}