diff --git a/scripts/!test.sh b/scripts/!test.sh deleted file mode 100755 index 6c779757a0..0000000000 --- a/scripts/!test.sh +++ /dev/null @@ -1,102 +0,0 @@ -#!/bin/bash - -clear -source ./pulse_lib.sh - -TOKEN=$(./token.sh claims.json) -ZP="00000000-0000-0000-0000-000000000001/TESTS" -# /AnyKey" - -# put ${ZP} "one text" - -# put "00000000-0000-0000-0000-000000000001/TESTS" "text 1" "If-None-Match: *" "Blooooooooo: blya" - -#exit - -#put "00000000-0000-0000-0000-000000000001/TESTS" "Value_1" "HULY-TTL: 3" -#echo "sleep 1 sec" -#sleep 1 -#get "00000000-0000-0000-0000-000000000001/TESTS" -#echo "sleep 3 sec" -#sleep 2 -#get "00000000-0000-0000-0000-000000000001/TESTS" - -put "00000000-0000-0000-0000-000000000001/TESTS1" "Value_1" "HULY-TTL: 3" -put "00000000-0000-0000-0000-000000000001/TESTS2" "Value_1" "HULY-TTL: 3" -put "00000000-0000-0000-0000-000000000001/HREST2" "Value_1" "HULY-TTL: 3" -get "00000000-0000-0000-0000-000000000001?prefix=TES" -sleep 1 -get "00000000-0000-0000-0000-000000000001?prefix=" - -exit - -echo "--------- delete ----------" -put "00000000-0000-0000-0000-000000000001/TESTS" "Value_2" "HULY-TTL: 3" -get "00000000-0000-0000-0000-000000000001/TESTS" -delete "00000000-0000-0000-0000-000000000001/TESTS" -get "00000000-0000-0000-0000-000000000001/TESTS" -#HULY-EXPIRE-AT: - - -# put "00000000-0000-0000-0000-000000000001/TESTS/some" "text 2" -# put "00000000-0000-0000-0000-000000000001/TESTS/some/gogo/" "text 3" - - -exit - - -echo "================> LIST" - put "00000000-0000-0000-0000-000000000001/Huome2/MyKey1" "value1" - put "00000000-0000-0000-0000-000000000001/Huome2/MyKey2" "value2" - get "00000000-0000-0000-0000-000000000001/Huome2" - delete "00000000-0000-0000-0000-000000000001/Huome2/MyKey1" - delete "00000000-0000-0000-0000-000000000001/Huome2/MyKey2" - -echo "================> WRONG UUID" - get "WrongUUID/TESTS/AnyKey" - -echo "================> INSERT If-None-Match" - - echo "-- Expected Error: 400 Bad Request (If-None-Match may be only *)" - put ${ZP} "enother text" "If-None-Match" "552e21cd4cd9918678e3c1a0df491bc3" - - delete ${ZP} - - echo "-- Expected OK: 201 Created (key was not exist)" - put ${ZP} "enother text" "If-None-Match" "*" - - put ${ZP} "some text" - echo "-- Expected Error: 412 Precondition Failed (key was exist)" - put ${ZP} "enother text" "If-None-Match" "*" - -echo "================> UPDATE PUT If-Match" - - get ${ZP} - - echo "-- Expected OK: 204 No Content (right hash)" - put ${ZP} "some text" "If-Match" "552e21cd4cd9918678e3c1a0df491bc3" - get ${ZP} - - echo "-- Expected OK: 204 No Content (hash still right)" - put ${ZP} "enother version" "If-Match" "552e21cd4cd9918678e3c1a0df491bc3" - get ${ZP} - - echo "-- Expected OK: 204 No Content (any hash)" - put ${ZP} "enother version2" "If-Match" "*" - get ${ZP} - - echo "-- Expected Error: 412 Precondition Failed (wrong hash)" - put ${ZP} "enother version3" "If-Match" "552e21cd4cd9918678e3c1a0df491bc3" - - delete ${ZP} - - echo "-- Expected Error: 412 Precondition Failed (any hash not found)" - put ${ZP} "enother version2" "If-Match" "*" - -echo "================> UPSERT (Expected OK)" - put ${ZP} "my value" - get ${ZP} - put ${ZP} "my new value" - get ${ZP} - -exit diff --git a/scripts/!ws.sh b/scripts/!ws.sh deleted file mode 100755 index 2d111e134c..0000000000 --- a/scripts/!ws.sh +++ /dev/null @@ -1,58 +0,0 @@ -#!/bin/bash - -clear -#source ./pulse_lib.sh - -websocat ws://127.0.0.1:8095/ws/testworkspace - -exit - - -let ws = new WebSocket("ws://localhost:8095/ws/testworkspace"); -ws.onmessage = e => console.log("Message from server:", e.data); -ws.onopen = () => ws.send("Hello from browser!"); - - - - - - - - - - - - - - - - - - - -TOKEN=$(./token.sh claims.json) -ZP="00000000-0000-0000-0000-000000000001/TESTS" -# /AnyKey" - -# put ${ZP} "one text" - -# put "00000000-0000-0000-0000-000000000001/TESTS" "text 1" "If-None-Match: *" "Blooooooooo: blya" - -#exit - -#put "00000000-0000-0000-0000-000000000001/TESTS" "Value_1" "HULY-TTL: 3" -#echo "sleep 1 sec" -#sleep 1 -#get "00000000-0000-0000-0000-000000000001/TESTS" -#echo "sleep 3 sec" -#sleep 2 -#get "00000000-0000-0000-0000-000000000001/TESTS" - -put "00000000-0000-0000-0000-000000000001/TESTS1" "Value_1" "HULY-TTL: 3" -put "00000000-0000-0000-0000-000000000001/TESTS2" "Value_1" "HULY-TTL: 3" -put "00000000-0000-0000-0000-000000000001/HREST2" "Value_1" "HULY-TTL: 3" -get "00000000-0000-0000-0000-000000000001?prefix=TES" -sleep 1 -get "00000000-0000-0000-0000-000000000001?prefix=" - -exit diff --git a/scripts/TEST.html b/scripts/TEST.html new file mode 100644 index 0000000000..8cd50f79b5 --- /dev/null +++ b/scripts/TEST.html @@ -0,0 +1,125 @@ + + + + + WebSocket JSON Tester + + + + +

WebSocket JSON Tester

+ + + +
+ + + +

+ + + + + + +

Waiting for server response...
+ + + + + diff --git a/scripts/TEST_HTTP_API.sh b/scripts/TEST_HTTP_API.sh index 4e5614c065..3e1a69e1f5 100755 --- a/scripts/TEST_HTTP_API.sh +++ b/scripts/TEST_HTTP_API.sh @@ -6,6 +6,53 @@ source ./pulse_lib.sh TOKEN=$(./token.sh claims.json) ZP="00000000-0000-0000-0000-000000000001/TESTS" +echo "--------- if-match ----------" + + delete ${ZP} + put "00000000-0000-0000-0000-000000000001/TESTS" "Value_1" "HULY-TTL: 1" "If-Match: *" + get ${ZP} + put "00000000-0000-0000-0000-000000000001/TESTS" "Value_2" "HULY-TTL: 1" + get ${ZP} + put "00000000-0000-0000-0000-000000000001/TESTS" "Value_3" "HULY-TTL: 1" "If-Match: dd358c74cb9cb897424838fbcb69c933" + put "00000000-0000-0000-0000-000000000001/TESTS" "Value_4" "HULY-TTL: 1" "If-Match: *" + put "00000000-0000-0000-0000-000000000001/TESTS" "Value_5" "HULY-TTL: 1" "If-Match: c7bcabf6b98a220f2f4888a18d01568d" + put "00000000-0000-0000-0000-000000000001/TESTS" "Value_6" "HULY-TTL: 1" "If-None-Match: *" + +echo "-- Expected OK: 201 Created (key was not exist)" + + put ${ZP} "enother text" "If-None-Match" "*" + + put ${ZP} "some text" + echo "-- Expected Error: 412 Precondition Failed (key was exist)" + put ${ZP} "enother text" "If-None-Match" "*" + +echo "================> UPDATE PUT If-Match" + + get ${ZP} + + echo "-- Expected OK: 204 No Content (right hash)" + put ${ZP} "some text" "If-Match" "552e21cd4cd9918678e3c1a0df491bc3" + get ${ZP} + + echo "-- Expected OK: 204 No Content (hash still right)" + put ${ZP} "enother version" "If-Match" "552e21cd4cd9918678e3c1a0df491bc3" + + + + + + + + + + + + + + + + + put "00000000-0000-0000-0000-000000000001/TESTS" "Value_1" "HULY-TTL: 3" echo "sleep 1 sec" sleep 1 diff --git a/src/handlers_http.rs b/src/handlers_http.rs index eddaac3e20..7786627bc8 100644 --- a/src/handlers_http.rs +++ b/src/handlers_http.rs @@ -8,7 +8,6 @@ use tracing::{error, trace}; use uuid::Uuid; use crate::ws_owner; -// type BucketPath = web::Path<(String)>; type ObjectPath = web::Path<(String, String)>; use crate::redis::{ @@ -21,15 +20,34 @@ use crate::redis::{ }; use actix_web::{ - Error, HttpMessage, HttpRequest, HttpResponse, error, + HttpRequest, HttpResponse, error, Error, web::{self, Data, Json, Query}, }; -/// list -// #[derive(Deserialize)] -// pub struct ListInfo { prefix: Option } +pub fn map_handler_error(err: impl std::fmt::Display) -> Error { + let msg = err.to_string(); + + if let Some(detail) = msg.split(" - ExtensionError: ").nth(1) { + if let Some((code, text)) = detail.split_once(": ") { + let text = format!("{} {}", code, text); + return match code { + "400" => actix_web::error::ErrorBadRequest(text), + "404" => actix_web::error::ErrorNotFound(text), + "412" => actix_web::error::ErrorPreconditionFailed(text), + "500" => actix_web::error::ErrorInternalServerError(text), + _ => actix_web::error::ErrorInternalServerError("unexpected error"), + }; + } + } + actix_web::error::ErrorInternalServerError("internal error") +} + + +/// list + +// #[derive(Deserialize)] pub async fn list( req: HttpRequest, path: web::Path, @@ -52,66 +70,8 @@ pub async fn list( Ok(HttpResponse::Ok().json(entries)) - }() - .await - .map_err(|err| { - tracing::error!(error = %err, "Internal error in GET handler"); - actix_web::error::ErrorInternalServerError("internal error") - }) - + }().await.map_err(map_handler_error) } -/* - path: BucketPath, - query: Query, - redis: web::Data>>, -) -> Result, actix_web::error::Error> { - - ws_owner::workspace_owner(&req)?; // Check workspace - - let (workspace) = path.into_inner(); - trace!(workspace, prefix = ?query.prefix, "list request"); - - // ... - - async move || -> anyhow::Result> { - let connection = pool.get().await?; - - let response = if let Some(prefix) = &query.prefix { - let pattern = format!("{}%", prefix); - let statement = r#" - select key from kvs where workspace=$1 and namespace=$2 and key like $3 - "#; - - connection - .query(statement, &[&wsuuid, &nsstr, &pattern]) - .await? - } else { - let statement = r#" - select key from kvs where workspace=$1 and namespace=$2 - "#; - - connection.query(statement, &[&wsuuid, &nsstr]).await? - }; - - let count = response.len(); - - let keys = response.into_iter().map(|row| row.get(0)).collect(); - - Ok(Json(ListResponse { - keys, - count, - namespace: nsstr.to_owned(), - workspace: wsstr.to_owned(), - })) - }() - .await - .map_err(|error| { - error!(op = "list", workspace, namespace, ?error, "internal error"); - error::ErrorInternalServerError("") - }) -} -*/ - /// get / (test) @@ -124,8 +84,6 @@ pub async fn get( ws_owner::workspace_owner(&req)?; // Check workspace let (workspace, key) = path.into_inner(); - // println!("\nworkspace = {}", workspace); - // println!("key = {}\n", key); trace!(workspace, key, "get request"); @@ -135,16 +93,13 @@ pub async fn get( Ok( redis_read(&mut *conn, &workspace, &key).await? - .map(|entry| HttpResponse::Ok().json(entry)) + .map(|entry| HttpResponse::Ok() + .insert_header(("ETag", &*entry.etag)) + .json(entry)) .unwrap_or_else(|| HttpResponse::NotFound().body("empty")) ) - }() - .await - .map_err(|err| { - tracing::error!(error = %err, "Internal error in GET handler"); - actix_web::error::ErrorInternalServerError("internal error") - }) + }().await.map_err(map_handler_error) } @@ -161,8 +116,6 @@ pub async fn put( let (workspace, key) = path.into_inner(); - trace!(workspace, key, "put request"); - async move || -> anyhow::Result { let mut conn = redis.lock().await; @@ -183,11 +136,9 @@ pub async fn put( let mut mode = Some(SaveMode::Upsert); if let Some(h) = req.headers().get("If-Match") { // `If-Match: *` - update only if the key exists let s = h.to_str().map_err(|_| anyhow!("Invalid If-Match header"))?; - if s == "*" { mode = Some(SaveMode::Update); } else { - // TODO: `If-Match: ` — update only if current value's MD5 matches - return Err(anyhow!("TODO: Only '*' suported now")); - } - } else if let Some(h) = req.headers().get("If-None-Match") { // `If-None-Match: *` — insert only if the key does not exist + if s == "*" { mode = Some(SaveMode::Update); } // `If-Match: *` — update only if exist + else { mode = Some(SaveMode::Equal(s.to_string())); } // `If-Match: ` — update only if current + } else if let Some(h) = req.headers().get("If-None-Match") { // `If-None-Match: *` — insert only if does not exist let s = h.to_str().map_err(|_| anyhow!("Invalid If-None-Match header"))?; if s == "*" { mode = Some(SaveMode::Insert); } else { return Err(anyhow!("If-None-Match must be '*'")); } } @@ -195,12 +146,7 @@ pub async fn put( redis_save(&mut *conn, &workspace, &key, &body[..], ttl, mode).await?; return Ok(HttpResponse::Ok().body("DONE")); - }() - .await - .map_err(|err| { - tracing::error!(error = %err, "Internal error in GET handler"); - actix_web::error::ErrorInternalServerError("internal error") - }) + }().await.map_err(map_handler_error) } @@ -218,9 +164,6 @@ pub async fn delete( let (workspace, key) = path.into_inner(); trace!(workspace, key, "delete request"); -// let wsuuid = Uuid::parse_str(workspace.as_str()) -// .map_err(|e| error::ErrorBadRequest(format!("Invalid UUID in workspace: {}", e)))?; - async move || -> anyhow::Result { let mut conn = redis.lock().await; @@ -232,10 +175,6 @@ pub async fn delete( }; Ok(response) - }() - .await - .map_err(|err| { - tracing::error!(error = %err, "Internal error in DELETE handler"); - actix_web::error::ErrorInternalServerError("internal error") - }) + }().await.map_err(map_handler_error) } + diff --git a/src/main.rs b/src/main.rs index a30f297e9f..99d2b56c4a 100644 --- a/src/main.rs +++ b/src/main.rs @@ -47,8 +47,6 @@ use config::CONFIG; use hulyrs::services::jwt::actix::ServiceRequestExt; use secrecy::SecretString; -// pub type Pool = bb8::Pool>; - fn initialize_tracing(level: tracing::Level) { use tracing_subscriber::{filter::targets::Targets, prelude::*}; diff --git a/src/redis.rs b/src/redis.rs index c17f152a91..29590a5b84 100644 --- a/src/redis.rs +++ b/src/redis.rs @@ -7,10 +7,12 @@ pub enum Ttl { At(u64), // EXAT (timestamp in seconds) } +#[derive(Debug)] pub enum SaveMode { Upsert, // default: set or overwrite Insert, // only if not exists (NX) Update, // only if exists (XX) + Equal(String), // only if md5 matches provided } use redis::{ @@ -27,14 +29,16 @@ pub struct RedisArray { pub key: String, pub data: String, pub expires_at: u64, // sec to expire TTL + pub etag: String, // md5 hash (data) } -fn error(msg: &'static str) -> RedisResult { - Err(redis::RedisError::from((redis::ErrorKind::ExtensionError, msg))) +fn error(code: u16, msg: impl Into) -> redis::RedisResult { + let msg = msg.into(); + let full = format!("{}: {}", code, msg); + Err(redis::RedisError::from(( redis::ErrorKind::ExtensionError, "", full ))) } /// redis_list(&connection,workspace,prefix) - pub async fn redis_list( conn: &mut MultiplexedConnection, workspace: &str, @@ -70,8 +74,9 @@ pub async fn redis_list( results.push(RedisArray { workspace: workspace.to_string(), key, - data: value, + data: value.clone(), expires_at: ttl as u64, + etag: hex::encode(md5::compute(&value).0), }); } } @@ -84,9 +89,7 @@ pub async fn redis_list( } - /// redis_read(&connection,workspace,key) - #[allow(dead_code)] pub async fn redis_read( conn: &mut MultiplexedConnection, @@ -97,22 +100,23 @@ pub async fn redis_read( let data: Option = redis::cmd("HGET").arg(workspace).arg(key).query_async(conn).await?; let Some(data) = data else { return Ok(None); }; - // let ttl: i64 = redis::cmd("HTTL").arg(workspace).arg("FIELDS").arg(1).arg(key).query_async(conn).await?; let ttl_vec: Vec = redis::cmd("HTTL").arg(workspace).arg("FIELDS").arg(1).arg(key).query_async(conn).await?; let ttl = ttl_vec.get(0).copied().unwrap_or(-3); // -3 unknown error - if ttl == -1 { return error("TTL not setL"); } - if ttl == -2 { return error("Key not found"); } - if ttl < 0 { return error("Unknown TTL error"); } + 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"); } Ok(Some(RedisArray { workspace: workspace.to_string(), key: key.to_string(), - data, + data: data.clone(), expires_at: ttl as u64, + etag: hex::encode(md5::compute(&data).0), })) } + /// TTL sec /// redis_save(&mut conn, "workspace", "key", "val", Some(Ttl::Sec(300)), Some(SaveMode::Insert)).await?; /// @@ -139,34 +143,50 @@ pub async fn redis_save( Some(Ttl::At(timestamp)) => { let now = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs(); if timestamp <= now { - return error("TTL timestamp exceeds MAX_TTL limit"); + return error(400, "TTL timestamp exceeds MAX_TTL limit"); } (timestamp - now) as usize } None => CONFIG.max_ttl, }; - if sec == 0 { return error("TTL must be > 0"); } - if sec > CONFIG.max_ttl { return error("TTL exceeds MAX_TTL"); } + if sec == 0 { return error(400, "TTL must be > 0"); } + if sec > CONFIG.max_ttl { return error(412, "TTL exceeds MAX_TTL"); } let mut cmd = redis::cmd("HSET"); cmd.arg(workspace).arg(key).arg(value); // Mode variants match mode.unwrap_or(SaveMode::Upsert) { - SaveMode::Upsert => {} // none - SaveMode::Insert => { cmd.arg("NX"); } - SaveMode::Update => { cmd.arg("XX"); } + + SaveMode::Upsert => {} // none + + SaveMode::Insert => { + let exists: bool = redis::cmd("HEXISTS").arg(workspace).arg(key).query_async(conn).await?; + if exists { return error(412, "Insert: key already exists"); } + } + + SaveMode::Update => { + let exists: bool = redis::cmd("HEXISTS").arg(workspace).arg(key).query_async(conn).await?; + if !exists { return error(404, "Update: key does not exist"); } + } + + SaveMode::Equal(md5) => { + let current_value: Option = redis::cmd("HGET").arg(workspace).arg(key).query_async(conn).await?; + if let Some(existing) = current_value { + let actual_md5 = hex::encode(md5::compute(&existing).0); + if actual_md5 != md5 { return error(412, format!("md5 mismatch, current: {} expected: {}", actual_md5, md5)); } + } else { return error(404, "Equal: key does not exist"); } + } + } // 1) HSET execute - if cmd.query_async::>(&mut *conn).await?.is_none() { - return error("SET failed: NX/XX condition not met"); - } + cmd.query_async::(&mut *conn).await?; // 2) HEXPIRE execute let res: Vec = redis::cmd("HEXPIRE").arg(workspace).arg(sec).arg("FIELDS").arg(1).arg(key).query_async(&mut *conn).await?; if res.get(0).copied().unwrap_or(0) == 0 { - return error("HEXPIRE failed: field not found or TTL not set"); + return error(404, "HEXPIRE field not found or TTL not set"); } Ok(()) diff --git a/src/redis.rs.ok b/src/redis.rs.ok deleted file mode 100644 index bd57538c0d..0000000000 --- a/src/redis.rs.ok +++ /dev/null @@ -1,231 +0,0 @@ -// hget hset TODO -// статистику TODO - -// ------------------------------- - -use crate::config::{CONFIG, RedisMode}; - -use std::time::{SystemTime, UNIX_EPOCH}; - -pub enum Ttl { - Sec(usize), // EX - At(u64), // EXAT (timestamp in seconds) -} - -pub enum SaveMode { - Upsert, // default: set or overwrite - Insert, // only if not exists (NX) - Update, // only if exists (XX) -} - -use redis::{ - AsyncCommands, RedisResult, - ToRedisArgs, - Client, ConnectionInfo, ProtocolVersion, RedisConnectionInfo, aio::MultiplexedConnection }; -use url::Url; - -use serde::{Deserialize, Serialize}; - -#[derive(Debug, Serialize)] -pub struct RedisArray { - pub workspace: String, - pub key: String, - pub data: String, - pub expires_at: Option, // секунды до истечения TTL -} - -fn error(msg: &'static str) -> redis::RedisResult<()> { - Err(redis::RedisError::from(( redis::ErrorKind::ExtensionError, msg ))) -} - -/// redis_read(&connection,key) - -#[allow(dead_code)] -pub async fn redis_read( - conn: &mut MultiplexedConnection, - workspace: &str, - key: &str, -) -> redis::RedisResult> { - - let data: Option = redis::cmd("HGET").arg(workspace).arg(key).query_async(conn).await?; - let Some(data) = data else { return Ok(None); }; - - // let ttl: i64 = redis::cmd("TTL").arg(redis_key).query_async(conn).await?; - let ttl: i64 = redis::cmd("TTL").arg(workspace).arg(key).query_async(conn).await?; - let expires_at = if ttl >= 0 { Some(ttl as u64) } else { None }; // -1 (нет TTL), -2 (нет ключа) - - Ok(Some(RedisArray { - workspace: workspace.to_string(), - key: key.to_string(), - data, - expires_at, - })) -} - - -/* -EX — срок жизни в секундах (e.g. EX 60 = 1 минута). -EXAT — дата истечения в секундах с эпохи Unix. -KEEPTTL — сохраняет текущий TTL ключа при перезаписи. - -Нет, несложно — Redis уже поддерживает это с помощью флагов NX и XX: - NX — записать только если ключ не существует - XX — перезаписать только если ключ уже существует - -Ты просто добавляешь .arg("NX") или .arg("XX") в команду SET. -Варианты: - SET key val EX 60 NX — с TTL, только если не существует - SET key val XX — только если уже существует, без TTL - SET key val — просто перезаписать, без TTL - -*/ - -/// 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?; - -#[allow(dead_code)] -pub async fn redis_save( - conn: &mut MultiplexedConnection, - workspace: &str, - key: &str, - value: T, - ttl: Option, - mode: Option, -) -> RedisResult<()> { - - // TTL variants - match ttl { - Some(Ttl::Sec(secs)) => { - if secs == 0 { - return error("TTL must be > 0"); - } - if secs > CONFIG.max_ttl { - return error("TTL exceeds MAX_TTL"); - } - cmd.arg("EX").arg(secs); - } - Some(Ttl::At(timestamp)) => { - let now = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs(); - if timestamp <= now { - return error("TTL timestamp is in the past"); - } - if timestamp - now > CONFIG.max_ttl as u64 { - return error("TTL timestamp exceeds MAX_TTL limit"); - } - cmd.arg("EXAT").arg(timestamp); - } - None => { cmd.arg("EX").arg(CONFIG.max_ttl); } - } - - - let mut cmd = redis::cmd("HSET"); - cmd.arg(workspace).arg(key).arg(value).query_async(conn).await?; - -redis::cmd("HEXPIRE") - .arg(workspace) - .arg(ttl_seconds) - .arg("FIELDS") - .arg(1).arg(key) - .query_async(conn).await?; - - - - // Mode variants - match mode.unwrap_or(SaveMode::Upsert) { - SaveMode::Upsert => { } // nothing - SaveMode::Insert => { cmd.arg("NX"); } - SaveMode::Update => { cmd.arg("XX"); } - } - - let res: Option = cmd.query_async(&mut *conn).await?; - - if res.is_none() { // nil - if NX/XX error - return error("SET failed: NX/XX condition not met"); - } else { - Ok(()) - } -} - - -#[allow(dead_code)] -pub async fn redis_delete( - conn: &mut MultiplexedConnection, - workspace: &str, - key: &str, -) -> redis::RedisResult { - - let deleted: i32 = redis::cmd("HDEL") - .arg(workspace) - .arg(key) - .query_async(conn) - .await?; - - Ok(deleted > 0) -} - - - -/// redis_connect() -pub async fn redis_connect() -> anyhow::Result { - let default_port = match CONFIG.redis_mode { - RedisMode::Sentinel => 6379, - RedisMode::Direct => 6380, - }; - - let urls = CONFIG - .redis_urls - .iter() - .map(|url| { - redis::ConnectionAddr::Tcp( - url.host().unwrap().to_string(), - url.port().unwrap_or(default_port), - ) - }) - .collect::>(); - - let conn = if CONFIG.redis_mode == RedisMode::Sentinel { - use redis::sentinel::{SentinelClientBuilder, SentinelServerType}; - - let mut sentinel = SentinelClientBuilder::new( - urls, - CONFIG.redis_service.to_owned(), - SentinelServerType::Master, - ) - .unwrap() - .set_client_to_redis_protocol(ProtocolVersion::RESP3) - .set_client_to_redis_db(0) - .set_client_to_redis_password(CONFIG.redis_password.clone()) - .set_client_to_sentinel_password(CONFIG.redis_password.clone()) - .build()?; - - sentinel.get_async_connection().await? - } else { - let single = urls - .first() - .ok_or_else(|| anyhow::anyhow!("No redis URL provided"))?; - - let redis_connection_info = RedisConnectionInfo { - db: 0, - username: None, - password: Some(CONFIG.redis_password.clone()), - protocol: ProtocolVersion::RESP3, - }; - - let connection_info = ConnectionInfo { - addr: single.clone(), - redis: redis_connection_info, - }; - - let client = Client::open(connection_info)?; - client.get_multiplexed_async_connection().await? - }; - - Ok(conn) -} diff --git a/src/ws_owner.rs b/src/ws_owner.rs index 941ded604c..d97035802a 100644 --- a/src/ws_owner.rs +++ b/src/ws_owner.rs @@ -2,7 +2,6 @@ use hulyrs::services::jwt::Claims; use uuid::Uuid; use actix_web::{ Error, HttpMessage, HttpRequest, error }; - /// Checking workspace in Authorization pub fn workspace_owner(req: &HttpRequest) -> Result<(), Error> { let extensions = req.extensions(); @@ -36,4 +35,3 @@ pub fn workspace_owner(req: &HttpRequest) -> Result<(), Error> { Ok(()) } -