Skip to content

Commit 43cc28c

Browse files
committed
improve script api
1 parent d0f53fa commit 43cc28c

3 files changed

Lines changed: 95 additions & 37 deletions

File tree

‎api/src/redis/client.rs‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ use fred::prelude::*;
22
use futures::StreamExt;
33
use itertools::Itertools;
44

5-
use crate::redis::{AddEvent, StreamService, constants, scripts, types::RedisStr};
5+
use crate::redis::{AddEvent, StreamService, constants, scripts::RedisScripts, types::RedisStr};
66

77
/// Redis client from static pool. Used for quick operations like retrieving stream status and
88
/// initializing a stream, not long-running / blocking commands.
@@ -53,7 +53,7 @@ impl RedisClient {
5353
let stream_key = self.stream.stream_key(key);
5454
let meta_key = self.stream.meta_key(key);
5555

56-
scripts::SCRIPTS
56+
RedisScripts
5757
.start_stream(&self.client, &stream_key, &meta_key, ttl)
5858
.await
5959
}
@@ -68,7 +68,7 @@ impl RedisClient {
6868
let stream_key = self.stream.stream_key(key);
6969
let meta_key = self.stream.meta_key(key);
7070

71-
scripts::SCRIPTS
71+
RedisScripts
7272
.write_events(&self.client, &stream_key, &meta_key, self.max_len, events)
7373
.await
7474
}
@@ -94,7 +94,7 @@ impl RedisClient {
9494
let stream_key = self.stream.stream_key(key);
9595
let meta_key = self.stream.meta_key(key);
9696

97-
scripts::SCRIPTS
97+
RedisScripts
9898
.finish_stream(&self.client, &stream_key, &meta_key, status, event)
9999
.await
100100
}

‎api/src/redis/scripts.rs‎

Lines changed: 87 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -4,26 +4,17 @@ use fred::{clients::Client, prelude::FredResult, types::scripts::Script};
44

55
use crate::redis::{AddEvent, StreamStatus, constants, types::RedisStr};
66

7-
/// Lua scripts for atomic Redis stream mutations.
8-
pub(super) static SCRIPTS: LazyLock<RedisScripts> = LazyLock::new(RedisScripts::new);
9-
10-
pub(super) struct RedisScripts {
11-
start_stream: Script,
12-
write_events: Script,
13-
finish_stream: Script,
14-
}
7+
/// Lua scripts for atomic Redis stream mutations. The scripts return
8+
/// `nil` (i.e. `None`) when the stream state does not allow the mutation.
9+
pub(super) struct RedisScripts;
1510

1611
impl RedisScripts {
17-
fn new() -> Self {
18-
Self {
19-
start_stream: Script::from_lua(START_STREAM_SCRIPT),
20-
write_events: Script::from_lua(WRITE_EVENTS_SCRIPT),
21-
finish_stream: Script::from_lua(FINISH_STREAM_SCRIPT),
22-
}
23-
}
24-
25-
/// Start/activate the stream and returns the ID of the start event. Returns
26-
/// `None` if stream is already active.
12+
/// Start and activate a stream.
13+
///
14+
/// Returns the Redis stream ID for the start event. Returns `None` if the
15+
/// stream is already active. If an inactive stream exists at the same key,
16+
/// the script deletes the old stream and metadata before creating the new
17+
/// stream.
2718
pub(super) async fn start_stream(
2819
&self,
2920
client: &Client,
@@ -40,12 +31,15 @@ impl RedisScripts {
4031
constants::START,
4132
];
4233

43-
self.start_stream
34+
START_STREAM_SCRIPT
4435
.evalsha_with_reload(&client, (stream_key, meta_key), args)
4536
.await
4637
}
4738

48-
/// Write events to the stream. Returns `None` if stream is not active.
39+
/// Write a batch of events to an active stream.
40+
//
41+
/// Returns the Redis stream IDs for all written events. Returns `None` if
42+
/// the stream is not active, without writing any events.
4943
pub(super) async fn write_events(
5044
&self,
5145
client: &Client,
@@ -67,12 +61,15 @@ impl RedisScripts {
6761
None => [&ev.event, "0", ""],
6862
}));
6963

70-
self.write_events
64+
WRITE_EVENTS_SCRIPT
7165
.evalsha_with_reload(client, (stream_key, meta_key), args)
7266
.await
7367
}
7468

75-
/// Write end event and mark stream as inactive
69+
/// Write a terminal event and mark the stream inactive.
70+
///
71+
/// Returns the Redis stream ID for the terminal event. Returns `None` if
72+
/// the stream is not active, without appending a terminal event.
7673
pub(super) async fn finish_stream(
7774
&self,
7875
client: &Client,
@@ -89,14 +86,30 @@ impl RedisScripts {
8986
event,
9087
];
9188

92-
self.finish_stream
89+
FINISH_STREAM_SCRIPT
9390
.evalsha_with_reload(&client, (stream_key, meta_key), args)
9491
.await
9592
}
9693
}
9794

98-
/// Lua script to atomically create a stream unless it is already active.
99-
const START_STREAM_SCRIPT: &str = r#"
95+
/// Atomically create a stream unless it is already active.
96+
///
97+
/// Key contract:
98+
/// - `KEYS[1]`: Redis stream key
99+
/// - `KEYS[2]`: stream metadata hash key
100+
///
101+
/// Argument contract:
102+
/// - `ARGV[1]`: metadata status field name
103+
/// - `ARGV[2]`: active status value
104+
/// - `ARGV[3]`: stream/meta TTL in seconds
105+
/// - `ARGV[4]`: stream entry event field name
106+
/// - `ARGV[5]`: start event value
107+
///
108+
/// Return contract:
109+
/// - stream ID for the start event when created
110+
/// - `nil` when the stream is already active
111+
static START_STREAM_SCRIPT: LazyLock<Script> = LazyLock::new(|| {
112+
let lua = r#"
100113
if redis.call('HGET', KEYS[2], ARGV[1]) == ARGV[2] then
101114
return nil
102115
end
@@ -109,9 +122,32 @@ redis.call('EXPIRE', KEYS[2], ARGV[3])
109122
110123
return id
111124
"#;
112-
113-
/// Lua script to atomically check for an active stream and write events.
114-
const WRITE_EVENTS_SCRIPT: &str = r#"
125+
Script::from_lua(lua)
126+
});
127+
128+
/// Atomically write a batch of events if the stream is active.
129+
///
130+
/// Key contract:
131+
/// - `KEYS[1]`: Redis stream key
132+
/// - `KEYS[2]`: stream metadata hash key
133+
///
134+
/// Fixed argument contract:
135+
/// - `ARGV[1]`: metadata status field name
136+
/// - `ARGV[2]`: active status value
137+
/// - `ARGV[3]`: approximate stream max length
138+
/// - `ARGV[4]`: stream entry event field name
139+
/// - `ARGV[5]`: stream entry data field name
140+
///
141+
/// Repeated event argument contract, starting at `ARGV[6]`:
142+
/// - event name
143+
/// - data flag: `"1"` means include the data field, `"0"` means omit it
144+
/// - data value, or an empty placeholder when the flag is `"0"`
145+
///
146+
/// Return contract:
147+
/// - array of stream IDs for the written events
148+
/// - `nil` when the stream is not active
149+
static WRITE_EVENTS_SCRIPT: LazyLock<Script> = LazyLock::new(|| {
150+
let lua = r#"
115151
if redis.call('HGET', KEYS[2], ARGV[1]) ~= ARGV[2] then
116152
return nil
117153
end
@@ -135,9 +171,27 @@ end
135171
136172
return ids
137173
"#;
138-
139-
/// Lua script to atomically append a terminal event and mark a stream inactive.
140-
const FINISH_STREAM_SCRIPT: &str = r#"
174+
Script::from_lua(lua)
175+
});
176+
177+
/// Atomically append a terminal event and mark a stream inactive.
178+
///
179+
/// Key contract:
180+
/// - `KEYS[1]`: Redis stream key
181+
/// - `KEYS[2]`: stream metadata hash key
182+
///
183+
/// Argument contract:
184+
/// - `ARGV[1]`: metadata status field name
185+
/// - `ARGV[2]`: active status value
186+
/// - `ARGV[3]`: final status value
187+
/// - `ARGV[4]`: stream entry event field name
188+
/// - `ARGV[5]`: terminal event value
189+
///
190+
/// Return contract:
191+
/// - stream ID for the terminal event
192+
/// - `nil` when the stream is not active
193+
static FINISH_STREAM_SCRIPT: LazyLock<Script> = LazyLock::new(|| {
194+
let lua = r#"
141195
if redis.call('HGET', KEYS[2], ARGV[1]) ~= ARGV[2] then
142196
return nil
143197
end
@@ -147,3 +201,5 @@ redis.call('HSET', KEYS[2], ARGV[1], ARGV[3])
147201
148202
return id
149203
"#;
204+
Script::from_lua(lua)
205+
});

‎api/src/redis/writer.rs‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
use fred::prelude::FredResult;
22

3-
use crate::redis::{AddEvent, ExclusiveClient, StreamService, scripts, types::RedisStr};
3+
use crate::redis::{
4+
AddEvent, ExclusiveClient, StreamService, scripts::RedisScripts, types::RedisStr,
5+
};
46

57
/// A stream writer with an exclusive lock on a Redis connection, for
68
/// long-running write operations (e.g. for ingesting events into Redis)
@@ -29,7 +31,7 @@ impl RedisWriter {
2931
let stream_key = self.stream.stream_key(key);
3032
let meta_key = self.stream.meta_key(key);
3133

32-
scripts::SCRIPTS
34+
RedisScripts
3335
.write_events(&self.client, &stream_key, &meta_key, self.max_len, events)
3436
.await
3537
}

0 commit comments

Comments
 (0)