Skip to content

Commit 107ff72

Browse files
committed
add endpoint to retrieve stream event history
1 parent cd66056 commit 107ff72

6 files changed

Lines changed: 176 additions & 9 deletions

File tree

‎.github/workflows/rust-client.yml‎

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -34,15 +34,17 @@ jobs:
3434
echo "version=$VERSION" >> $GITHUB_OUTPUT
3535
3636
- name: Generate Rust client
37+
id: generate
3738
run: |
38-
cargo progenitor \
39+
output=$(cargo progenitor \
3940
-i spec/openapi.json \
4041
-o clients/rust/ \
4142
-n tinistream-client \
4243
--interface builder \
4344
--tags separate \
4445
--license-name MIT \
45-
--version ${{ steps.spec-version.outputs.version }} || true
46+
--version ${{ steps.spec-version.outputs.version }}) || true
47+
echo "output=$output" >> $GITHUB_OUTPUT
4648
4749
- name: Check for changes
4850
id: changes
@@ -72,6 +74,11 @@ jobs:
7274
- Updated Rust client to match OpenAPI spec version `${{ steps.spec-version.outputs.version }}`
7375
- Generated using `cargo-progenitor`
7476
77+
**Output of `cargo-progenitor`**
78+
```
79+
${{ steps.generate.outputs.output }}
80+
```
81+
7582
**Generated on:** ${{ steps.date.outputs.date }}
7683
7784
Please review the changes and merge if they look correct.

‎Cargo.lock‎

Lines changed: 1 addition & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎api/Cargo.toml‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
[package]
22
name = "tinistream"
3-
version = "0.1.2"
3+
version = "0.1.3"
44
edition = "2021"
55
publish = false
66

@@ -22,7 +22,7 @@ itertools = "0.14.0"
2222
rocket = { version = "0.5.1", features = ["json", "secrets"] }
2323
rocket_okapi = { version = "0.9.0", features = ["rapidoc", "rocket_ws"] }
2424
rocket_ws = { version = "0.1.1" }
25-
schemars = "0.8.22"
25+
schemars = { version = "0.8.22" }
2626
serde = { version = "1.0", features = ["derive"] }
2727
serde_json = "1.0"
2828
thiserror = "2.0"

‎api/src/api/stream.rs‎

Lines changed: 45 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,14 +2,14 @@ use rocket::{futures::StreamExt, get, post, serde::json::Json, Route, State};
22
use rocket_okapi::{okapi::openapi3::OpenApi, openapi, openapi_get_routes_spec};
33
use schemars::JsonSchema;
44
use serde::{Deserialize, Serialize};
5-
use time::ext::NumericalDuration;
5+
use time::{ext::NumericalDuration, format_description::well_known, UtcDateTime};
66

77
use crate::{
88
auth::{create_client_token, ApiKeyAuth, Crypto},
99
config::AppConfig,
1010
data::JsonStream,
1111
errors::ApiError,
12-
redis::{stream_sse_url, RedisClient, StreamStatus, DATA_KEY, EVENT_KEY},
12+
redis::*,
1313
};
1414

1515
pub fn get_routes() -> (Vec<Route>, OpenApi) {
@@ -18,6 +18,7 @@ pub fn get_routes() -> (Vec<Route>, OpenApi) {
1818
create_stream,
1919
create_token,
2020
get_stream_info,
21+
get_stream_events,
2122
add_events,
2223
add_events_json_stream,
2324
cancel_stream,
@@ -64,6 +65,35 @@ async fn get_stream_info(
6465
}))
6566
}
6667

68+
/// # Get stream events
69+
/// Get all events so far in a stream
70+
#[openapi(tag = "Stream")]
71+
#[get("/events?<key>")]
72+
async fn get_stream_events(
73+
_api_key: ApiKeyAuth,
74+
key: &str,
75+
reader: RedisReader,
76+
) -> Result<Json<Vec<StreamEvent>>, ApiError> {
77+
let (prev_events, _, _) = reader.get_prev_events(key, None).await?;
78+
let events = prev_events
79+
.into_iter()
80+
.filter_map(|(id, mut data)| {
81+
let unix_millis: i64 = id.split('-').next().unwrap_or_default().parse().ok()?;
82+
let date_time = UtcDateTime::from_unix_timestamp(unix_millis / 1000).ok()?;
83+
let iso_time = date_time.format(&well_known::Rfc3339).ok()?;
84+
let event = StreamEvent {
85+
id: (*id).to_owned(),
86+
time: iso_time,
87+
event: data.remove(EVENT_KEY).as_deref().map(|e| e.to_owned())?,
88+
data: data.remove(DATA_KEY).as_deref().map(|d| d.to_owned()),
89+
};
90+
Some(event)
91+
})
92+
.collect();
93+
94+
Ok(Json(events))
95+
}
96+
6797
/// # Create stream
6898
/// Create a new stream, and get a client URL and token to connect to the stream
6999
#[openapi(tag = "Stream")]
@@ -217,6 +247,19 @@ pub struct StreamInfo {
217247
ttl: i64,
218248
}
219249

250+
#[derive(JsonSchema, Serialize)]
251+
pub struct StreamEvent {
252+
/// ID of the event
253+
id: String,
254+
/// Time of the event (ISO 8601 format)
255+
time: String,
256+
/// Name/type of the event
257+
event: String,
258+
/// Event data
259+
#[serde(skip_serializing_if = "Option::is_none")]
260+
data: Option<String>,
261+
}
262+
220263
#[derive(JsonSchema, Deserialize)]
221264
struct StreamRequest {
222265
key: String,

‎api/src/redis/reader.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -81,7 +81,7 @@ impl RedisReader {
8181

8282
/// Returns a tuple containing the previous events in the stream, the last event ID,
8383
/// and a boolean indicating if the stream has already ended.
84-
async fn get_prev_events(
84+
pub async fn get_prev_events(
8585
&self,
8686
key: &str,
8787
start_event_id: Option<&str>,

‎spec/openapi.json‎

Lines changed: 118 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22
"openapi": "3.0.0",
33
"info": {
44
"title": "tinistream",
5-
"version": "0.1.2"
5+
"version": "0.1.3"
66
},
77
"paths": {
88
"/api/health": {
@@ -494,6 +494,96 @@
494494
]
495495
}
496496
},
497+
"/api/stream/events": {
498+
"get": {
499+
"tags": [
500+
"Stream"
501+
],
502+
"summary": "Get stream events",
503+
"description": "Get all events so far in a stream",
504+
"operationId": "get_stream_events",
505+
"parameters": [
506+
{
507+
"name": "key",
508+
"in": "query",
509+
"required": true,
510+
"schema": {
511+
"type": "string"
512+
}
513+
}
514+
],
515+
"responses": {
516+
"200": {
517+
"description": "",
518+
"content": {
519+
"application/json": {
520+
"schema": {
521+
"type": "array",
522+
"items": {
523+
"$ref": "#/components/schemas/StreamEvent"
524+
}
525+
}
526+
}
527+
}
528+
},
529+
"400": {
530+
"description": "Bad request",
531+
"content": {
532+
"application/json": {
533+
"schema": {
534+
"$ref": "#/components/schemas/ErrorMessage"
535+
}
536+
}
537+
}
538+
},
539+
"401": {
540+
"description": "Unauthorized",
541+
"content": {
542+
"application/json": {
543+
"schema": {
544+
"$ref": "#/components/schemas/ErrorMessage"
545+
}
546+
}
547+
}
548+
},
549+
"404": {
550+
"description": "Not found",
551+
"content": {
552+
"application/json": {
553+
"schema": {
554+
"$ref": "#/components/schemas/ErrorMessage"
555+
}
556+
}
557+
}
558+
},
559+
"422": {
560+
"description": "Incorrectly formatted",
561+
"content": {
562+
"application/json": {
563+
"schema": {
564+
"$ref": "#/components/schemas/ErrorMessage"
565+
}
566+
}
567+
}
568+
},
569+
"500": {
570+
"description": "Internal server error",
571+
"content": {
572+
"application/json": {
573+
"schema": {
574+
"$ref": "#/components/schemas/ErrorMessage"
575+
}
576+
}
577+
}
578+
}
579+
},
580+
"security": [
581+
{
582+
"ApiKey": []
583+
}
584+
]
585+
}
586+
},
497587
"/api/stream/add": {
498588
"post": {
499589
"tags": [
@@ -985,6 +1075,33 @@
9851075
}
9861076
}
9871077
},
1078+
"StreamEvent": {
1079+
"type": "object",
1080+
"required": [
1081+
"event",
1082+
"id",
1083+
"time"
1084+
],
1085+
"properties": {
1086+
"id": {
1087+
"description": "ID of the event",
1088+
"type": "string"
1089+
},
1090+
"time": {
1091+
"description": "Time of the event (ISO 8601 format)",
1092+
"type": "string"
1093+
},
1094+
"event": {
1095+
"description": "Name/type of the event",
1096+
"type": "string"
1097+
},
1098+
"data": {
1099+
"description": "Event data",
1100+
"type": "string",
1101+
"nullable": true
1102+
}
1103+
}
1104+
},
9881105
"AddEventsResponse": {
9891106
"type": "object",
9901107
"required": [

0 commit comments

Comments
 (0)