1use std::{
17 collections::{self, HashMap},
18 fmt::{self, Debug, Display},
19 future::IntoFuture,
20 io::{self},
21 pin::Pin,
22 str::FromStr,
23 task::Poll,
24 time::Duration,
25};
26
27use crate::datetime::{rfc3339, DateTime};
28use crate::{
29 error::Error, header::HeaderName, is_valid_subject, HeaderMap, HeaderValue, StatusCode,
30};
31use base64::engine::general_purpose::STANDARD;
32use base64::engine::Engine;
33use bytes::Bytes;
34use futures_util::{future::BoxFuture, FutureExt, TryFutureExt};
35use serde::{Deserialize, Deserializer, Serialize};
36use serde_json::json;
37
38use super::{
39 consumer::{self, Consumer, FromConsumer, IntoConsumerConfig},
40 context::{
41 ConsumerInfoError, ConsumerInfoErrorKind, RequestError, RequestErrorKind, StreamsError,
42 StreamsErrorKind,
43 },
44 errors::ErrorCode,
45 is_valid_name,
46 message::{StreamMessage, StreamMessageError},
47 response::Response,
48 Context,
49};
50
51pub type InfoError = RequestError;
52
53#[derive(Clone, Debug, PartialEq)]
54pub enum DirectGetErrorKind {
55 NotFound,
56 InvalidSubject,
57 TimedOut,
58 Request,
59 ErrorResponse(StatusCode, String),
60 Other,
61}
62
63impl Display for DirectGetErrorKind {
64 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
65 match self {
66 Self::InvalidSubject => write!(f, "invalid subject"),
67 Self::NotFound => write!(f, "message not found"),
68 Self::ErrorResponse(status, description) => {
69 write!(f, "unable to get message: {status} {description}")
70 }
71 Self::Other => write!(f, "error getting message"),
72 Self::TimedOut => write!(f, "timed out"),
73 Self::Request => write!(f, "request failed"),
74 }
75 }
76}
77
78pub type DirectGetError = Error<DirectGetErrorKind>;
79
80impl From<crate::RequestError> for DirectGetError {
81 fn from(err: crate::RequestError) -> Self {
82 match err.kind() {
83 crate::RequestErrorKind::TimedOut => DirectGetError::new(DirectGetErrorKind::TimedOut),
84 crate::RequestErrorKind::NoResponders => {
85 DirectGetError::new(DirectGetErrorKind::ErrorResponse(
86 StatusCode::NO_RESPONDERS,
87 "no responders".to_string(),
88 ))
89 }
90 crate::RequestErrorKind::InvalidSubject
91 | crate::RequestErrorKind::MaxPayloadExceeded
92 | crate::RequestErrorKind::Other => {
93 DirectGetError::with_source(DirectGetErrorKind::Other, err)
94 }
95 }
96 }
97}
98
99impl From<serde_json::Error> for DirectGetError {
100 fn from(err: serde_json::Error) -> Self {
101 DirectGetError::with_source(DirectGetErrorKind::Other, err)
102 }
103}
104
105impl From<StreamMessageError> for DirectGetError {
106 fn from(err: StreamMessageError) -> Self {
107 DirectGetError::with_source(DirectGetErrorKind::Other, err)
108 }
109}
110
111#[derive(Clone, Debug, PartialEq)]
112pub enum DeleteMessageErrorKind {
113 Request,
114 TimedOut,
115 JetStream(super::errors::Error),
116}
117
118impl Display for DeleteMessageErrorKind {
119 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
120 match self {
121 Self::Request => write!(f, "request failed"),
122 Self::TimedOut => write!(f, "timed out"),
123 Self::JetStream(err) => write!(f, "JetStream error: {err}"),
124 }
125 }
126}
127
128pub type DeleteMessageError = Error<DeleteMessageErrorKind>;
129
130#[derive(Debug, Clone)]
134pub struct Stream<T = Info> {
135 pub(crate) info: T,
136 pub(crate) context: Context,
137 pub(crate) name: String,
138}
139
140impl Stream<Info> {
141 pub async fn info(&mut self) -> Result<&Info, InfoError> {
159 let subject = format!("STREAM.INFO.{}", self.info.config.name);
160
161 match self.context.request(subject, &json!({})).await? {
162 Response::Ok::<Info>(info) => {
163 self.info = info;
164 Ok(&self.info)
165 }
166 Response::Err { error } => Err(error.into()),
167 }
168 }
169
170 pub fn cached_info(&self) -> &Info {
189 &self.info
190 }
191}
192
193impl<I> Stream<I> {
194 pub async fn get_info(&self) -> Result<Info, InfoError> {
197 let subject = format!("STREAM.INFO.{}", self.name);
198
199 match self.context.request(subject, &json!({})).await? {
200 Response::Ok::<Info>(info) => Ok(info),
201 Response::Err { error } => Err(error.into()),
202 }
203 }
204
205 pub async fn info_with_subjects<F: AsRef<str>>(
228 &self,
229 subjects_filter: F,
230 ) -> Result<InfoWithSubjects, InfoError> {
231 let subjects_filter = subjects_filter.as_ref().to_string();
232 let info = stream_info_with_details(
234 self.context.clone(),
235 self.name.clone(),
236 0,
237 false,
238 subjects_filter.clone(),
239 )
240 .await?;
241
242 Ok(InfoWithSubjects::new(
243 self.context.clone(),
244 info,
245 subjects_filter,
246 ))
247 }
248
249 pub fn info_builder(&self) -> StreamInfoBuilder {
275 StreamInfoBuilder::new(self.context.clone(), self.name.clone())
276 }
277
278 pub fn direct_get_builder(&self) -> DirectGetBuilder<WithHeaders> {
298 DirectGetBuilder::new(self.context.clone(), self.name.clone())
299 }
300
301 pub async fn direct_get_next_for_subject<T: Into<String>>(
336 &self,
337 subject: T,
338 sequence: Option<u64>,
339 ) -> Result<StreamMessage, DirectGetError> {
340 let subject_str = subject.into();
341 if !is_valid_subject(&subject_str) {
342 return Err(DirectGetError::new(DirectGetErrorKind::InvalidSubject));
343 }
344
345 let mut builder = self.direct_get_builder().next_by_subject(subject_str);
346 if let Some(seq) = sequence {
347 builder = builder.sequence(seq);
348 }
349
350 builder.send().await
351 }
352
353 pub async fn direct_get_first_for_subject<T: Into<String>>(
385 &self,
386 subject: T,
387 ) -> Result<StreamMessage, DirectGetError> {
388 let subject_str = subject.into();
389 if !is_valid_subject(&subject_str) {
390 return Err(DirectGetError::new(DirectGetErrorKind::InvalidSubject));
391 }
392
393 self.direct_get_builder()
394 .next_by_subject(subject_str)
395 .send()
396 .await
397 }
398
399 pub async fn direct_get(&self, sequence: u64) -> Result<StreamMessage, DirectGetError> {
431 self.direct_get_builder().sequence(sequence).send().await
432 }
433
434 pub async fn direct_get_last_for_subject<T: Into<String>>(
466 &self,
467 subject: T,
468 ) -> Result<StreamMessage, DirectGetError> {
469 self.direct_get_builder()
470 .last_by_subject(subject)
471 .send()
472 .await
473 }
474 pub fn raw_message_builder(&self) -> RawMessageBuilder<WithHeaders> {
494 RawMessageBuilder::new(self.context.clone(), self.name.clone())
495 }
496
497 pub async fn get_raw_message(&self, sequence: u64) -> Result<StreamMessage, RawMessageError> {
527 self.raw_message_builder().sequence(sequence).send().await
528 }
529
530 pub async fn get_first_raw_message_by_subject<T: AsRef<str>>(
554 &self,
555 subject: T,
556 sequence: u64,
557 ) -> Result<StreamMessage, RawMessageError> {
558 self.raw_message_builder()
559 .sequence(sequence)
560 .next_by_subject(subject.as_ref().to_string())
561 .send()
562 .await
563 }
564
565 pub async fn get_next_raw_message_by_subject<T: AsRef<str>>(
589 &self,
590 subject: T,
591 ) -> Result<StreamMessage, RawMessageError> {
592 self.raw_message_builder()
593 .next_by_subject(subject.as_ref().to_string())
594 .send()
595 .await
596 }
597
598 pub async fn get_last_raw_message_by_subject(
622 &self,
623 stream_subject: &str,
624 ) -> Result<StreamMessage, LastRawMessageError> {
625 self.raw_message_builder()
626 .last_by_subject(stream_subject.to_string())
627 .send()
628 .await
629 }
630
631 pub async fn delete_message(&self, sequence: u64) -> Result<bool, DeleteMessageError> {
655 let subject = format!("STREAM.MSG.DELETE.{}", self.name);
656 let payload = json!({
657 "seq": sequence,
658 });
659
660 let response: Response<DeleteStatus> = self
661 .context
662 .request(subject, &payload)
663 .map_err(|err| match err.kind() {
664 RequestErrorKind::TimedOut => {
665 DeleteMessageError::new(DeleteMessageErrorKind::TimedOut)
666 }
667 _ => DeleteMessageError::with_source(DeleteMessageErrorKind::Request, err),
668 })
669 .await?;
670
671 match response {
672 Response::Err { error } => Err(DeleteMessageError::new(
673 DeleteMessageErrorKind::JetStream(error),
674 )),
675 Response::Ok(value) => Ok(value.success),
676 }
677 }
678
679 pub fn purge(&self) -> Purge<No, No> {
695 Purge::build(self)
696 }
697
698 #[deprecated(
715 since = "0.25.0",
716 note = "Overloads have been replaced with an into_future based builder. Use Stream::purge().filter(subject) instead."
717 )]
718 pub async fn purge_subject<T>(&self, subject: T) -> Result<PurgeResponse, PurgeError>
719 where
720 T: Into<String>,
721 {
722 self.purge().filter(subject).await
723 }
724
725 pub async fn create_consumer<C: IntoConsumerConfig + FromConsumer>(
749 &self,
750 config: C,
751 ) -> Result<Consumer<C>, ConsumerError> {
752 self.context
753 .create_consumer_on_stream(config, self.name.clone())
754 .await
755 }
756
757 #[cfg(feature = "server_2_10")]
781 pub async fn update_consumer<C: IntoConsumerConfig + FromConsumer>(
782 &self,
783 config: C,
784 ) -> Result<Consumer<C>, ConsumerUpdateError> {
785 self.context
786 .update_consumer_on_stream(config, self.name.clone())
787 .await
788 }
789
790 #[cfg(feature = "server_2_10")]
815 pub async fn create_consumer_strict<C: IntoConsumerConfig + FromConsumer>(
816 &self,
817 config: C,
818 ) -> Result<Consumer<C>, ConsumerCreateStrictError> {
819 self.context
820 .create_consumer_strict_on_stream(config, self.name.clone())
821 .await
822 }
823
824 pub async fn consumer_info<T: AsRef<str>>(
841 &self,
842 name: T,
843 ) -> Result<consumer::Info, ConsumerInfoError> {
844 let name = name.as_ref();
845
846 if !is_valid_name(name) {
847 return Err(ConsumerInfoError::new(ConsumerInfoErrorKind::InvalidName));
848 }
849
850 let subject = format!("CONSUMER.INFO.{}.{}", self.name, name);
851
852 match self.context.request(subject, &json!({})).await? {
853 Response::Ok(info) => Ok(info),
854 Response::Err { error } => Err(error.into()),
855 }
856 }
857
858 pub async fn get_consumer<T: FromConsumer + IntoConsumerConfig>(
877 &self,
878 name: &str,
879 ) -> Result<Consumer<T>, crate::Error> {
880 let info = self.consumer_info(name).await?;
881
882 Ok(Consumer::new(
883 T::try_from_consumer_config(info.config.clone())?,
884 info,
885 self.context.clone(),
886 ))
887 }
888
889 pub async fn get_or_create_consumer<T: FromConsumer + IntoConsumerConfig>(
917 &self,
918 name: &str,
919 config: T,
920 ) -> Result<Consumer<T>, ConsumerError> {
921 let subject = format!("CONSUMER.INFO.{}.{}", self.name, name);
922
923 match self.context.request(subject, &json!({})).await? {
924 Response::Err { error } if error.code() == 404 => self.create_consumer(config).await,
925 Response::Err { error } => Err(error.into()),
926 Response::Ok::<consumer::Info>(info) => Ok(Consumer::new(
927 T::try_from_consumer_config(info.config.clone()).map_err(|err| {
928 ConsumerError::with_source(ConsumerErrorKind::InvalidConsumerType, err)
929 })?,
930 info,
931 self.context.clone(),
932 )),
933 }
934 }
935
936 pub async fn delete_consumer(&self, name: &str) -> Result<DeleteStatus, ConsumerError> {
957 let subject = format!("CONSUMER.DELETE.{}.{}", self.name, name);
958
959 match self.context.request(subject, &json!({})).await? {
960 Response::Ok(delete_status) => Ok(delete_status),
961 Response::Err { error } => Err(error.into()),
962 }
963 }
964
965 #[cfg(feature = "server_2_11")]
990 pub async fn pause_consumer(
991 &self,
992 name: &str,
993 pause_until: DateTime,
994 ) -> Result<PauseResponse, ConsumerError> {
995 self.request_pause_consumer(name, Some(pause_until)).await
996 }
997
998 #[cfg(feature = "server_2_11")]
1019 pub async fn resume_consumer(&self, name: &str) -> Result<PauseResponse, ConsumerError> {
1020 self.request_pause_consumer(name, None).await
1021 }
1022
1023 #[cfg(feature = "server_2_11")]
1024 async fn request_pause_consumer(
1025 &self,
1026 name: &str,
1027 pause_until: Option<DateTime>,
1028 ) -> Result<PauseResponse, ConsumerError> {
1029 let subject = format!("CONSUMER.PAUSE.{}.{}", self.name, name);
1030 let payload = &PauseResumeConsumerRequest { pause_until };
1031
1032 match self.context.request(subject, payload).await? {
1033 Response::Ok::<PauseResponse>(resp) => Ok(resp),
1034 Response::Err { error } => Err(error.into()),
1035 }
1036 }
1037
1038 #[cfg(feature = "server_2_14")]
1070 #[cfg_attr(docsrs, doc(cfg(feature = "server_2_14")))]
1071 pub async fn reset_consumer(
1072 &self,
1073 name: &str,
1074 seq: Option<u64>,
1075 ) -> Result<ConsumerResetResponse, ConsumerResetError> {
1076 let subject = format!("CONSUMER.RESET.{}.{}", self.name, name);
1077 let payload = ConsumerResetRequest {
1078 seq: seq.unwrap_or(0),
1079 };
1080
1081 match self.context.request(subject, &payload).await? {
1082 Response::Ok::<ConsumerResetResponse>(resp) => Ok(resp),
1083 Response::Err { error } => Err(error.into()),
1084 }
1085 }
1086
1087 pub fn consumer_names(&self) -> ConsumerNames {
1106 ConsumerNames {
1107 context: self.context.clone(),
1108 stream: self.name.clone(),
1109 offset: 0,
1110 page_request: None,
1111 consumers: Vec::new(),
1112 done: false,
1113 }
1114 }
1115
1116 pub fn consumers(&self) -> Consumers {
1135 Consumers {
1136 context: self.context.clone(),
1137 stream: self.name.clone(),
1138 offset: 0,
1139 page_request: None,
1140 consumers: Vec::new(),
1141 done: false,
1142 }
1143 }
1144}
1145
1146pub struct StreamInfoBuilder {
1147 pub(crate) context: Context,
1148 pub(crate) name: String,
1149 pub(crate) deleted: bool,
1150 pub(crate) subject: String,
1151}
1152
1153impl StreamInfoBuilder {
1154 fn new(context: Context, name: String) -> Self {
1155 Self {
1156 context,
1157 name,
1158 deleted: false,
1159 subject: "".to_string(),
1160 }
1161 }
1162
1163 pub fn with_deleted(mut self, deleted: bool) -> Self {
1164 self.deleted = deleted;
1165 self
1166 }
1167
1168 pub fn subjects<S: Into<String>>(mut self, subject: S) -> Self {
1169 self.subject = subject.into();
1170 self
1171 }
1172
1173 pub async fn fetch(self) -> Result<InfoWithSubjects, InfoError> {
1174 let info = stream_info_with_details(
1175 self.context.clone(),
1176 self.name.clone(),
1177 0,
1178 self.deleted,
1179 self.subject.clone(),
1180 )
1181 .await?;
1182
1183 Ok(InfoWithSubjects::new(self.context, info, self.subject))
1184 }
1185}
1186
1187#[derive(Debug, Default, Serialize, Deserialize, Clone, PartialEq, Eq)]
1191pub struct Config {
1192 pub name: String,
1194 #[serde(default)]
1196 pub max_bytes: i64,
1197 #[serde(default, rename = "max_msgs")]
1199 pub max_messages: i64,
1200 #[serde(default, rename = "max_msgs_per_subject")]
1202 pub max_messages_per_subject: i64,
1203 pub discard: DiscardPolicy,
1206 #[serde(default, skip_serializing_if = "is_default")]
1208 pub discard_new_per_subject: bool,
1209 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1212 pub subjects: Vec<String>,
1213 pub retention: RetentionPolicy,
1215 #[serde(default)]
1217 pub max_consumers: i32,
1218 #[serde(default, with = "serde_nanos")]
1220 pub max_age: Duration,
1221 #[serde(default, skip_serializing_if = "is_default", rename = "max_msg_size")]
1223 pub max_message_size: i32,
1224 pub storage: StorageType,
1226 pub num_replicas: usize,
1228 #[serde(default, skip_serializing_if = "is_default")]
1230 pub no_ack: bool,
1231 #[serde(default, skip_serializing_if = "is_default", with = "serde_nanos")]
1233 pub duplicate_window: Duration,
1234 #[serde(default, skip_serializing_if = "is_default")]
1236 pub template_owner: String,
1237 #[serde(default, skip_serializing_if = "is_default")]
1239 pub sealed: bool,
1240 #[serde(default, skip_serializing_if = "is_default")]
1242 pub description: Option<String>,
1243 #[serde(
1244 default,
1245 rename = "allow_rollup_hdrs",
1246 skip_serializing_if = "is_default"
1247 )]
1248 pub allow_rollup: bool,
1250 #[serde(default, skip_serializing_if = "is_default")]
1251 pub deny_delete: bool,
1253 #[serde(default, skip_serializing_if = "is_default")]
1255 pub deny_purge: bool,
1256
1257 #[serde(default, skip_serializing_if = "is_default")]
1259 pub republish: Option<Republish>,
1260
1261 #[serde(default, skip_serializing_if = "is_default")]
1264 pub allow_direct: bool,
1265
1266 #[serde(default, skip_serializing_if = "is_default")]
1268 pub mirror_direct: bool,
1269
1270 #[serde(default, skip_serializing_if = "Option::is_none")]
1272 pub mirror: Option<Source>,
1273
1274 #[serde(default, skip_serializing_if = "Option::is_none")]
1276 pub sources: Option<Vec<Source>>,
1277
1278 #[cfg(feature = "server_2_10")]
1279 #[serde(default, skip_serializing_if = "is_default")]
1281 pub metadata: HashMap<String, String>,
1282
1283 #[cfg(feature = "server_2_10")]
1284 #[serde(default, skip_serializing_if = "Option::is_none")]
1286 pub subject_transform: Option<SubjectTransform>,
1287
1288 #[cfg(feature = "server_2_10")]
1289 #[serde(default, skip_serializing_if = "Option::is_none")]
1294 pub compression: Option<Compression>,
1295 #[cfg(feature = "server_2_10")]
1296 #[serde(default, deserialize_with = "default_consumer_limits_as_none")]
1298 pub consumer_limits: Option<ConsumerLimits>,
1299
1300 #[cfg(feature = "server_2_10")]
1301 #[serde(default, skip_serializing_if = "Option::is_none", rename = "first_seq")]
1303 pub first_sequence: Option<u64>,
1304
1305 #[serde(default, skip_serializing_if = "Option::is_none")]
1307 pub placement: Option<Placement>,
1308
1309 #[serde(default, skip_serializing_if = "Option::is_none")]
1311 pub persist_mode: Option<PersistenceMode>,
1312
1313 #[cfg(feature = "server_2_11")]
1315 #[serde(
1316 default,
1317 with = "rfc3339::option",
1318 skip_serializing_if = "Option::is_none"
1319 )]
1320 pub pause_until: Option<DateTime>,
1321
1322 #[cfg(feature = "server_2_11")]
1324 #[serde(default, skip_serializing_if = "is_default", rename = "allow_msg_ttl")]
1325 pub allow_message_ttl: bool,
1326
1327 #[cfg(feature = "server_2_11")]
1330 #[serde(default, skip_serializing_if = "Option::is_none", with = "serde_nanos")]
1331 pub subject_delete_marker_ttl: Option<Duration>,
1332
1333 #[cfg(feature = "server_2_12")]
1335 #[serde(default, skip_serializing_if = "is_default", rename = "allow_atomic")]
1336 pub allow_atomic_publish: bool,
1337
1338 #[cfg(feature = "server_2_12")]
1340 #[serde(
1341 default,
1342 skip_serializing_if = "is_default",
1343 rename = "allow_msg_schedules"
1344 )]
1345 pub allow_message_schedules: bool,
1346
1347 #[cfg(feature = "server_2_12")]
1349 #[serde(
1350 default,
1351 skip_serializing_if = "is_default",
1352 rename = "allow_msg_counter"
1353 )]
1354 pub allow_message_counter: bool,
1355
1356 #[cfg(feature = "server_2_14")]
1358 #[cfg_attr(docsrs, doc(cfg(feature = "server_2_14")))]
1359 #[serde(default, skip_serializing_if = "is_default", rename = "allow_batched")]
1360 pub allow_batch_publish: bool,
1361}
1362
1363impl From<&Config> for Config {
1364 fn from(sc: &Config) -> Config {
1365 sc.clone()
1366 }
1367}
1368
1369impl From<&str> for Config {
1370 fn from(s: &str) -> Config {
1371 Config {
1372 name: s.to_string(),
1373 ..Default::default()
1374 }
1375 }
1376}
1377
1378#[cfg(feature = "server_2_10")]
1379fn default_consumer_limits_as_none<'de, D>(
1380 deserializer: D,
1381) -> Result<Option<ConsumerLimits>, D::Error>
1382where
1383 D: Deserializer<'de>,
1384{
1385 let consumer_limits = Option::<ConsumerLimits>::deserialize(deserializer)?;
1386 if let Some(cl) = consumer_limits {
1387 if cl == ConsumerLimits::default() {
1388 Ok(None)
1389 } else {
1390 Ok(Some(cl))
1391 }
1392 } else {
1393 Ok(None)
1394 }
1395}
1396#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq, Default)]
1397pub struct ConsumerLimits {
1398 #[serde(default, with = "serde_nanos")]
1400 pub inactive_threshold: std::time::Duration,
1401 #[serde(default)]
1403 pub max_ack_pending: i64,
1404}
1405
1406#[derive(Serialize, Deserialize, Debug, Clone, Eq, PartialEq)]
1407pub enum Compression {
1408 #[serde(rename = "s2")]
1409 S2,
1410 #[serde(rename = "none")]
1411 None,
1412}
1413
1414#[derive(Serialize, Deserialize, Debug, Clone, Eq, PartialEq)]
1416pub struct SubjectTransform {
1417 #[serde(rename = "src")]
1418 pub source: String,
1419
1420 #[serde(rename = "dest")]
1421 pub destination: String,
1422}
1423
1424#[derive(Serialize, Deserialize, Debug, Clone, Eq, PartialEq)]
1426pub struct Republish {
1427 #[serde(rename = "src")]
1429 pub source: String,
1430 #[serde(rename = "dest")]
1432 pub destination: String,
1433 #[serde(default)]
1435 pub headers_only: bool,
1436}
1437
1438#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
1440pub struct Placement {
1441 #[serde(default, skip_serializing_if = "is_default")]
1443 pub cluster: Option<String>,
1444 #[serde(default, skip_serializing_if = "is_default")]
1446 pub tags: Vec<String>,
1447}
1448
1449#[derive(Default, Debug, Serialize, Deserialize, Clone, Copy, PartialEq, Eq)]
1452#[repr(u8)]
1453pub enum DiscardPolicy {
1454 #[default]
1456 #[serde(rename = "old")]
1457 Old = 0,
1458 #[serde(rename = "new")]
1460 New = 1,
1461}
1462
1463#[derive(Default, Debug, Serialize, Deserialize, Clone, Copy, PartialEq, Eq)]
1465#[repr(u8)]
1466pub enum RetentionPolicy {
1467 #[default]
1470 #[serde(rename = "limits")]
1471 Limits = 0,
1472 #[serde(rename = "interest")]
1474 Interest = 1,
1475 #[serde(rename = "workqueue")]
1477 WorkQueue = 2,
1478}
1479
1480#[derive(Default, Debug, Serialize, Deserialize, Clone, Copy, PartialEq, Eq)]
1482#[repr(u8)]
1483pub enum StorageType {
1484 #[default]
1486 #[serde(rename = "file")]
1487 File = 0,
1488 #[serde(rename = "memory")]
1490 Memory = 1,
1491}
1492
1493#[derive(Default, Debug, Serialize, Deserialize, Clone, Copy, PartialEq, Eq)]
1495#[repr(u8)]
1496pub enum PersistenceMode {
1497 #[default]
1499 #[serde(rename = "default")]
1500 Default = 0,
1501 #[serde(rename = "async")]
1503 Async = 1,
1504}
1505
1506async fn stream_info_with_details(
1507 context: Context,
1508 stream: String,
1509 offset: usize,
1510 deleted_details: bool,
1511 subjects_filter: String,
1512) -> Result<Info, InfoError> {
1513 let subject = format!("STREAM.INFO.{stream}");
1514
1515 let payload = StreamInfoRequest {
1516 offset,
1517 deleted_details,
1518 subjects_filter,
1519 };
1520
1521 let response: Response<Info> = context.request(subject, &payload).await?;
1522
1523 match response {
1524 Response::Ok(info) => Ok(info),
1525 Response::Err { error } => Err(error.into()),
1526 }
1527}
1528
1529type InfoRequest = BoxFuture<'static, Result<Info, InfoError>>;
1530
1531#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
1532pub struct StreamInfoRequest {
1533 offset: usize,
1534 deleted_details: bool,
1535 subjects_filter: String,
1536}
1537
1538pub struct InfoWithSubjects {
1539 stream: String,
1540 context: Context,
1541 pub info: Info,
1542 offset: usize,
1543 subjects: collections::hash_map::IntoIter<String, usize>,
1544 info_request: Option<InfoRequest>,
1545 subjects_filter: String,
1546 pages_done: bool,
1547}
1548
1549impl InfoWithSubjects {
1550 pub fn new(context: Context, mut info: Info, subject: String) -> Self {
1551 let subjects = info.state.subjects.take().unwrap_or_default();
1552 let name = info.config.name.clone();
1553 InfoWithSubjects {
1554 context,
1555 info,
1556 pages_done: subjects.is_empty(),
1557 offset: subjects.len(),
1558 subjects: subjects.into_iter(),
1559 subjects_filter: subject,
1560 stream: name,
1561 info_request: None,
1562 }
1563 }
1564}
1565
1566impl futures_util::Stream for InfoWithSubjects {
1567 type Item = Result<(String, usize), InfoError>;
1568
1569 fn poll_next(
1570 mut self: Pin<&mut Self>,
1571 cx: &mut std::task::Context<'_>,
1572 ) -> Poll<Option<Self::Item>> {
1573 match self.subjects.next() {
1574 Some((subject, count)) => Poll::Ready(Some(Ok((subject, count)))),
1575 None => {
1576 if self.pages_done {
1578 return Poll::Ready(None);
1579 }
1580 let stream = self.stream.clone();
1581 let context = self.context.clone();
1582 let subjects_filter = self.subjects_filter.clone();
1583 let offset = self.offset;
1584 match self
1585 .info_request
1586 .get_or_insert_with(|| {
1587 Box::pin(stream_info_with_details(
1588 context,
1589 stream,
1590 offset,
1591 false,
1592 subjects_filter,
1593 ))
1594 })
1595 .poll_unpin(cx)
1596 {
1597 Poll::Ready(resp) => match resp {
1598 Ok(info) => {
1599 let subjects = info.state.subjects.clone();
1600 self.offset += subjects.as_ref().map_or_else(|| 0, |s| s.len());
1601 self.info_request = None;
1602 let subjects = subjects.unwrap_or_default();
1603 self.subjects = info.state.subjects.unwrap_or_default().into_iter();
1604 let total = info.paged_info.map(|info| info.total).unwrap_or(0);
1605 if total <= self.offset || subjects.is_empty() {
1606 self.pages_done = true;
1607 }
1608 match self.subjects.next() {
1609 Some((subject, count)) => Poll::Ready(Some(Ok((subject, count)))),
1610 None => Poll::Ready(None),
1611 }
1612 }
1613 Err(err) => Poll::Ready(Some(Err(err))),
1614 },
1615 Poll::Pending => Poll::Pending,
1616 }
1617 }
1618 }
1619 }
1620}
1621
1622#[derive(Debug, Deserialize, Clone, PartialEq, Eq)]
1624pub struct Info {
1625 pub config: Config,
1627 #[serde(with = "rfc3339")]
1629 pub created: DateTime,
1630 pub state: State,
1632 pub cluster: Option<ClusterInfo>,
1634 #[serde(default)]
1636 pub mirror: Option<SourceInfo>,
1637 #[serde(default)]
1639 pub sources: Vec<SourceInfo>,
1640 #[serde(flatten)]
1641 paged_info: Option<PagedInfo>,
1642}
1643
1644#[derive(Debug, Deserialize, Clone, PartialEq, Eq)]
1645pub struct PagedInfo {
1646 offset: usize,
1647 total: usize,
1648 limit: usize,
1649}
1650
1651#[derive(Deserialize)]
1652pub struct DeleteStatus {
1653 pub success: bool,
1654}
1655
1656#[cfg(feature = "server_2_11")]
1657#[derive(Deserialize)]
1658pub struct PauseResponse {
1659 pub paused: bool,
1660 #[serde(with = "rfc3339")]
1661 pub pause_until: DateTime,
1662 #[serde(default, with = "serde_nanos")]
1663 pub pause_remaining: Option<Duration>,
1664}
1665
1666#[cfg(feature = "server_2_11")]
1667#[derive(Serialize, Debug)]
1668struct PauseResumeConsumerRequest {
1669 #[serde(with = "rfc3339::option", skip_serializing_if = "Option::is_none")]
1670 pause_until: Option<DateTime>,
1671}
1672
1673#[cfg(feature = "server_2_14")]
1674#[derive(Serialize, Debug)]
1675pub(crate) struct ConsumerResetRequest {
1676 #[serde(default, skip_serializing_if = "is_default")]
1677 pub(crate) seq: u64,
1678}
1679
1680#[cfg(feature = "server_2_14")]
1682#[cfg_attr(docsrs, doc(cfg(feature = "server_2_14")))]
1683#[derive(Debug, Deserialize, Clone)]
1684pub struct ConsumerResetResponse {
1685 #[serde(flatten)]
1688 pub info: super::consumer::Info,
1689 pub reset_seq: u64,
1693}
1694
1695#[derive(Debug, Deserialize, Clone, PartialEq, Eq)]
1697pub struct State {
1698 pub messages: u64,
1700 pub bytes: u64,
1702 #[serde(rename = "first_seq")]
1704 pub first_sequence: u64,
1705 #[serde(with = "rfc3339", rename = "first_ts")]
1707 pub first_timestamp: DateTime,
1708 #[serde(rename = "last_seq")]
1710 pub last_sequence: u64,
1711 #[serde(with = "rfc3339", rename = "last_ts")]
1713 pub last_timestamp: DateTime,
1714 pub consumer_count: usize,
1716 #[serde(default, rename = "num_subjects")]
1718 pub subjects_count: u64,
1719 #[serde(default, rename = "num_deleted")]
1721 pub deleted_count: Option<u64>,
1722 #[serde(default)]
1725 pub deleted: Option<Vec<u64>>,
1726
1727 pub(crate) subjects: Option<HashMap<String, usize>>,
1728}
1729
1730#[derive(Debug, Serialize, Deserialize, Clone)]
1732pub struct RawMessage {
1733 #[serde(rename = "subject")]
1735 pub subject: String,
1736
1737 #[serde(rename = "seq")]
1739 pub sequence: u64,
1740
1741 #[serde(default, rename = "data")]
1743 pub payload: String,
1744
1745 #[serde(default, rename = "hdrs")]
1747 pub headers: Option<String>,
1748
1749 #[serde(rename = "time", with = "rfc3339")]
1751 pub time: DateTime,
1752}
1753
1754impl TryFrom<RawMessage> for StreamMessage {
1755 type Error = crate::Error;
1756
1757 fn try_from(value: RawMessage) -> Result<Self, Self::Error> {
1758 let decoded_payload = STANDARD
1759 .decode(value.payload)
1760 .map_err(|err| Box::new(std::io::Error::other(err)))?;
1761 let decoded_headers = value
1762 .headers
1763 .map(|header| STANDARD.decode(header))
1764 .map_or(Ok(None), |v| v.map(Some))?;
1765
1766 let (headers, _, _) = decoded_headers
1767 .map_or_else(|| Ok((HeaderMap::new(), None, None)), |h| parse_headers(&h))?;
1768
1769 Ok(StreamMessage {
1770 subject: value.subject.into(),
1771 payload: decoded_payload.into(),
1772 headers,
1773 sequence: value.sequence,
1774 time: value.time,
1775 })
1776 }
1777}
1778
1779fn is_continuation(c: char) -> bool {
1780 c == ' ' || c == '\t'
1781}
1782const HEADER_LINE: &str = "NATS/1.0";
1783
1784#[allow(clippy::type_complexity)]
1785fn parse_headers(
1786 buf: &[u8],
1787) -> Result<(HeaderMap, Option<StatusCode>, Option<String>), crate::Error> {
1788 let mut headers = HeaderMap::new();
1789 let mut maybe_status: Option<StatusCode> = None;
1790 let mut maybe_description: Option<String> = None;
1791 let mut lines = if let Ok(line) = std::str::from_utf8(buf) {
1792 line.lines().peekable()
1793 } else {
1794 return Err(Box::new(std::io::Error::other("invalid header")));
1795 };
1796
1797 if let Some(line) = lines.next() {
1798 let line = line
1799 .strip_prefix(HEADER_LINE)
1800 .ok_or_else(|| {
1801 Box::new(std::io::Error::other(
1802 "version line does not start with NATS/1.0",
1803 ))
1804 })?
1805 .trim();
1806
1807 match line.split_once(' ') {
1808 Some((status, description)) => {
1809 if !status.is_empty() {
1810 maybe_status = Some(status.parse()?);
1811 }
1812
1813 if !description.is_empty() {
1814 maybe_description = Some(description.trim().to_string());
1815 }
1816 }
1817 None => {
1818 if !line.is_empty() {
1819 maybe_status = Some(line.parse()?);
1820 }
1821 }
1822 }
1823 } else {
1824 return Err(Box::new(std::io::Error::other(
1825 "expected header information not found",
1826 )));
1827 };
1828
1829 while let Some(line) = lines.next() {
1830 if line.is_empty() {
1831 continue;
1832 }
1833
1834 if let Some((k, v)) = line.split_once(':').to_owned() {
1835 let mut s = String::from(v.trim());
1836 while let Some(v) = lines.next_if(|s| s.starts_with(is_continuation)).to_owned() {
1837 s.push(' ');
1838 s.push_str(v.trim());
1839 }
1840
1841 headers.insert(
1842 HeaderName::from_str(k)?,
1843 HeaderValue::from_str(&s).map_err(|err| Box::new(io::Error::other(err)))?,
1844 );
1845 } else {
1846 return Err(Box::new(std::io::Error::other("malformed header line")));
1847 }
1848 }
1849
1850 if headers.is_empty() {
1851 Ok((HeaderMap::new(), maybe_status, maybe_description))
1852 } else {
1853 Ok((headers, maybe_status, maybe_description))
1854 }
1855}
1856
1857#[derive(Debug, Serialize, Deserialize, Clone)]
1858struct GetRawMessage {
1859 pub(crate) message: RawMessage,
1860}
1861
1862fn is_default<T: Default + Eq>(t: &T) -> bool {
1863 t == &T::default()
1864}
1865#[derive(Debug, Default, Deserialize, Clone, PartialEq, Eq)]
1867pub struct ClusterInfo {
1868 #[serde(default)]
1870 pub name: Option<String>,
1871 #[serde(default)]
1873 pub raft_group: Option<String>,
1874 #[serde(default)]
1876 pub leader: Option<String>,
1877 #[serde(default, with = "rfc3339::option")]
1879 pub leader_since: Option<DateTime>,
1880 #[cfg(feature = "server_2_12")]
1882 #[serde(default)]
1883 pub system_account: bool,
1885 #[cfg(feature = "server_2_12")]
1886 #[serde(default)]
1888 pub traffic_account: Option<String>,
1889 #[serde(default)]
1891 pub replicas: Vec<PeerInfo>,
1892}
1893
1894#[derive(Debug, Default, Deserialize, Clone, PartialEq, Eq)]
1896pub struct PeerInfo {
1897 pub name: String,
1899 pub current: bool,
1901 #[serde(with = "serde_nanos")]
1903 pub active: Duration,
1904 #[serde(default)]
1906 pub offline: bool,
1907 pub lag: Option<u64>,
1909}
1910
1911#[derive(Debug, Clone, Deserialize, PartialEq, Eq)]
1912pub struct SourceInfo {
1913 pub name: String,
1915 pub lag: u64,
1917 #[serde(deserialize_with = "negative_duration_as_none")]
1919 pub active: Option<std::time::Duration>,
1920 #[serde(default)]
1922 pub filter_subject: Option<String>,
1923 #[serde(default)]
1925 pub subject_transform_dest: Option<String>,
1926 #[serde(default)]
1928 pub subject_transforms: Vec<SubjectTransform>,
1929}
1930
1931fn negative_duration_as_none<'de, D>(
1932 deserializer: D,
1933) -> Result<Option<std::time::Duration>, D::Error>
1934where
1935 D: Deserializer<'de>,
1936{
1937 let n = i64::deserialize(deserializer)?;
1938 if n.is_negative() {
1939 Ok(None)
1940 } else {
1941 Ok(Some(std::time::Duration::from_nanos(n as u64)))
1942 }
1943}
1944
1945#[derive(Debug, Deserialize, Clone, Copy)]
1947pub struct PurgeResponse {
1948 pub success: bool,
1950 pub purged: u64,
1952}
1953#[derive(Default, Debug, Serialize, Clone)]
1955pub struct PurgeRequest {
1956 #[serde(default, rename = "seq", skip_serializing_if = "is_default")]
1958 pub sequence: Option<u64>,
1959
1960 #[serde(default, skip_serializing_if = "is_default")]
1962 pub filter: Option<String>,
1963
1964 #[serde(default, skip_serializing_if = "is_default")]
1966 pub keep: Option<u64>,
1967}
1968
1969#[derive(Clone, Debug, Serialize, Deserialize, Eq, PartialEq, Default)]
1970pub struct Source {
1971 pub name: String,
1973 #[serde(default, rename = "opt_start_seq", skip_serializing_if = "is_default")]
1975 pub start_sequence: Option<u64>,
1976 #[serde(
1977 default,
1978 rename = "opt_start_time",
1979 skip_serializing_if = "is_default",
1980 with = "rfc3339::option"
1981 )]
1982 pub start_time: Option<DateTime>,
1984 #[serde(default, skip_serializing_if = "is_default")]
1986 pub filter_subject: Option<String>,
1987 #[serde(default, skip_serializing_if = "Option::is_none")]
1989 pub external: Option<External>,
1990 #[serde(default, skip_serializing_if = "is_default")]
1992 pub domain: Option<String>,
1993 #[cfg(feature = "server_2_10")]
1995 #[serde(default, skip_serializing_if = "is_default")]
1996 pub subject_transforms: Vec<SubjectTransform>,
1997
1998 #[cfg(feature = "server_2_14")]
2005 #[cfg_attr(docsrs, doc(cfg(feature = "server_2_14")))]
2006 #[serde(default, skip_serializing_if = "Option::is_none")]
2007 pub consumer: Option<StreamConsumerSource>,
2008}
2009
2010#[cfg(feature = "server_2_14")]
2013#[cfg_attr(docsrs, doc(cfg(feature = "server_2_14")))]
2014#[derive(Clone, Debug, Serialize, Deserialize, Eq, PartialEq, Default)]
2015pub struct StreamConsumerSource {
2016 #[serde(default, skip_serializing_if = "is_default")]
2018 pub name: String,
2019 #[serde(default, skip_serializing_if = "is_default")]
2021 pub deliver_subject: String,
2022}
2023
2024#[cfg(feature = "server_2_14")]
2025impl StreamConsumerSource {
2026 pub fn new(name: impl Into<String>, deliver_subject: impl Into<String>) -> Self {
2029 Self {
2030 name: name.into(),
2031 deliver_subject: deliver_subject.into(),
2032 }
2033 }
2034}
2035
2036#[derive(Clone, Debug, Serialize, Deserialize, Eq, PartialEq, Default)]
2037pub struct External {
2038 #[serde(rename = "api")]
2040 pub api_prefix: String,
2041 #[serde(rename = "deliver", skip_serializing_if = "is_default")]
2043 pub delivery_prefix: Option<String>,
2044}
2045
2046use std::marker::PhantomData;
2047
2048#[derive(Debug, Default)]
2049pub struct Yes;
2050#[derive(Debug, Default)]
2051pub struct No;
2052
2053pub trait ToAssign: Debug {}
2054
2055impl ToAssign for Yes {}
2056impl ToAssign for No {}
2057
2058#[derive(Debug)]
2059pub struct Purge<SEQUENCE, KEEP>
2060where
2061 SEQUENCE: ToAssign,
2062 KEEP: ToAssign,
2063{
2064 inner: PurgeRequest,
2065 sequence_set: PhantomData<SEQUENCE>,
2066 keep_set: PhantomData<KEEP>,
2067 context: Context,
2068 stream_name: String,
2069}
2070
2071impl<SEQUENCE, KEEP> Purge<SEQUENCE, KEEP>
2072where
2073 SEQUENCE: ToAssign,
2074 KEEP: ToAssign,
2075{
2076 pub fn filter<T: Into<String>>(mut self, filter: T) -> Purge<SEQUENCE, KEEP> {
2078 self.inner.filter = Some(filter.into());
2079 self
2080 }
2081}
2082
2083impl Purge<No, No> {
2084 pub(crate) fn build<I>(stream: &Stream<I>) -> Purge<No, No> {
2085 Purge {
2086 context: stream.context.clone(),
2087 stream_name: stream.name.clone(),
2088 inner: Default::default(),
2089 sequence_set: PhantomData {},
2090 keep_set: PhantomData {},
2091 }
2092 }
2093}
2094
2095impl<KEEP> Purge<No, KEEP>
2096where
2097 KEEP: ToAssign,
2098{
2099 pub fn keep(self, keep: u64) -> Purge<No, Yes> {
2102 Purge {
2103 context: self.context.clone(),
2104 stream_name: self.stream_name.clone(),
2105 sequence_set: PhantomData {},
2106 keep_set: PhantomData {},
2107 inner: PurgeRequest {
2108 keep: Some(keep),
2109 ..self.inner
2110 },
2111 }
2112 }
2113}
2114impl<SEQUENCE> Purge<SEQUENCE, No>
2115where
2116 SEQUENCE: ToAssign,
2117{
2118 pub fn sequence(self, sequence: u64) -> Purge<Yes, No> {
2121 Purge {
2122 context: self.context.clone(),
2123 stream_name: self.stream_name.clone(),
2124 sequence_set: PhantomData {},
2125 keep_set: PhantomData {},
2126 inner: PurgeRequest {
2127 sequence: Some(sequence),
2128 ..self.inner
2129 },
2130 }
2131 }
2132}
2133
2134#[derive(Clone, Debug, PartialEq)]
2135pub enum PurgeErrorKind {
2136 Request,
2137 TimedOut,
2138 JetStream(super::errors::Error),
2139}
2140
2141impl Display for PurgeErrorKind {
2142 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2143 match self {
2144 Self::Request => write!(f, "request failed"),
2145 Self::TimedOut => write!(f, "timed out"),
2146 Self::JetStream(err) => write!(f, "JetStream error: {err}"),
2147 }
2148 }
2149}
2150
2151pub type PurgeError = Error<PurgeErrorKind>;
2152
2153impl<S, K> IntoFuture for Purge<S, K>
2154where
2155 S: ToAssign + std::marker::Send,
2156 K: ToAssign + std::marker::Send,
2157{
2158 type Output = Result<PurgeResponse, PurgeError>;
2159
2160 type IntoFuture = BoxFuture<'static, Result<PurgeResponse, PurgeError>>;
2161
2162 fn into_future(self) -> Self::IntoFuture {
2163 Box::pin(std::future::IntoFuture::into_future(async move {
2164 let request_subject = format!("STREAM.PURGE.{}", self.stream_name);
2165 let response: Response<PurgeResponse> = self
2166 .context
2167 .request(request_subject, &self.inner)
2168 .map_err(|err| match err.kind() {
2169 RequestErrorKind::TimedOut => PurgeError::new(PurgeErrorKind::TimedOut),
2170 _ => PurgeError::with_source(PurgeErrorKind::Request, err),
2171 })
2172 .await?;
2173
2174 match response {
2175 Response::Err { error } => Err(PurgeError::new(PurgeErrorKind::JetStream(error))),
2176 Response::Ok(response) => Ok(response),
2177 }
2178 }))
2179 }
2180}
2181
2182#[derive(Deserialize, Debug)]
2183struct ConsumerPage {
2184 total: usize,
2185 consumers: Option<Vec<String>>,
2186}
2187
2188#[derive(Deserialize, Debug)]
2189struct ConsumerInfoPage {
2190 total: usize,
2191 consumers: Option<Vec<super::consumer::Info>>,
2192}
2193
2194type ConsumerNamesErrorKind = StreamsErrorKind;
2195type ConsumerNamesError = StreamsError;
2196type PageRequest = BoxFuture<'static, Result<ConsumerPage, RequestError>>;
2197
2198pub struct ConsumerNames {
2199 context: Context,
2200 stream: String,
2201 offset: usize,
2202 page_request: Option<PageRequest>,
2203 consumers: Vec<String>,
2204 done: bool,
2205}
2206
2207impl futures_util::Stream for ConsumerNames {
2208 type Item = Result<String, ConsumerNamesError>;
2209
2210 fn poll_next(
2211 mut self: Pin<&mut Self>,
2212 cx: &mut std::task::Context<'_>,
2213 ) -> std::task::Poll<Option<Self::Item>> {
2214 match self.page_request.as_mut() {
2215 Some(page) => match page.try_poll_unpin(cx) {
2216 std::task::Poll::Ready(page) => {
2217 self.page_request = None;
2218 let page = page.map_err(|err| {
2219 ConsumerNamesError::with_source(ConsumerNamesErrorKind::Other, err)
2220 })?;
2221
2222 if let Some(consumers) = page.consumers {
2223 self.offset += consumers.len();
2224 self.consumers = consumers;
2225 if self.offset >= page.total {
2226 self.done = true;
2227 }
2228 match self.consumers.pop() {
2229 Some(stream) => Poll::Ready(Some(Ok(stream))),
2230 None => Poll::Ready(None),
2231 }
2232 } else {
2233 Poll::Ready(None)
2234 }
2235 }
2236 std::task::Poll::Pending => std::task::Poll::Pending,
2237 },
2238 None => {
2239 if let Some(stream) = self.consumers.pop() {
2240 Poll::Ready(Some(Ok(stream)))
2241 } else {
2242 if self.done {
2243 return Poll::Ready(None);
2244 }
2245 let context = self.context.clone();
2246 let offset = self.offset;
2247 let stream = self.stream.clone();
2248 self.page_request = Some(Box::pin(async move {
2249 match context
2250 .request(
2251 format!("CONSUMER.NAMES.{stream}"),
2252 &json!({
2253 "offset": offset,
2254 }),
2255 )
2256 .await?
2257 {
2258 Response::Err { error } => Err(RequestError::with_source(
2259 super::context::RequestErrorKind::Other,
2260 error,
2261 )),
2262 Response::Ok(page) => Ok(page),
2263 }
2264 }));
2265 self.poll_next(cx)
2266 }
2267 }
2268 }
2269 }
2270}
2271
2272pub type ConsumersErrorKind = StreamsErrorKind;
2273pub type ConsumersError = StreamsError;
2274type PageInfoRequest = BoxFuture<'static, Result<ConsumerInfoPage, RequestError>>;
2275
2276pub struct Consumers {
2277 context: Context,
2278 stream: String,
2279 offset: usize,
2280 page_request: Option<PageInfoRequest>,
2281 consumers: Vec<super::consumer::Info>,
2282 done: bool,
2283}
2284
2285impl futures_util::Stream for Consumers {
2286 type Item = Result<super::consumer::Info, ConsumersError>;
2287
2288 fn poll_next(
2289 mut self: Pin<&mut Self>,
2290 cx: &mut std::task::Context<'_>,
2291 ) -> std::task::Poll<Option<Self::Item>> {
2292 match self.page_request.as_mut() {
2293 Some(page) => match page.try_poll_unpin(cx) {
2294 std::task::Poll::Ready(page) => {
2295 self.page_request = None;
2296 let page = page.map_err(|err| {
2297 ConsumersError::with_source(ConsumersErrorKind::Other, err)
2298 })?;
2299 if let Some(consumers) = page.consumers {
2300 self.offset += consumers.len();
2301 self.consumers = consumers;
2302 if self.offset >= page.total {
2303 self.done = true;
2304 }
2305 match self.consumers.pop() {
2306 Some(consumer) => Poll::Ready(Some(Ok(consumer))),
2307 None => Poll::Ready(None),
2308 }
2309 } else {
2310 Poll::Ready(None)
2311 }
2312 }
2313 std::task::Poll::Pending => std::task::Poll::Pending,
2314 },
2315 None => {
2316 if let Some(stream) = self.consumers.pop() {
2317 Poll::Ready(Some(Ok(stream)))
2318 } else {
2319 if self.done {
2320 return Poll::Ready(None);
2321 }
2322 let context = self.context.clone();
2323 let offset = self.offset;
2324 let stream = self.stream.clone();
2325 self.page_request = Some(Box::pin(async move {
2326 match context
2327 .request(
2328 format!("CONSUMER.LIST.{stream}"),
2329 &json!({
2330 "offset": offset,
2331 }),
2332 )
2333 .await?
2334 {
2335 Response::Err { error } => Err(RequestError::with_source(
2336 super::context::RequestErrorKind::Other,
2337 error,
2338 )),
2339 Response::Ok(page) => Ok(page),
2340 }
2341 }));
2342 self.poll_next(cx)
2343 }
2344 }
2345 }
2346 }
2347}
2348
2349#[derive(Clone, Debug, PartialEq)]
2350pub enum LastRawMessageErrorKind {
2351 NoMessageFound,
2352 InvalidSubject,
2353 JetStream(super::errors::Error),
2354 Other,
2355}
2356
2357impl Display for LastRawMessageErrorKind {
2358 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2359 match self {
2360 Self::NoMessageFound => write!(f, "no message found"),
2361 Self::InvalidSubject => write!(f, "invalid subject"),
2362 Self::Other => write!(f, "failed to get last raw message"),
2363 Self::JetStream(err) => write!(f, "JetStream error: {err}"),
2364 }
2365 }
2366}
2367
2368pub type LastRawMessageError = Error<LastRawMessageErrorKind>;
2369pub type RawMessageErrorKind = LastRawMessageErrorKind;
2370pub type RawMessageError = LastRawMessageError;
2371
2372#[derive(Clone, Debug, PartialEq)]
2373pub enum ConsumerErrorKind {
2374 TimedOut,
2376 Request,
2377 InvalidConsumerType,
2378 InvalidName,
2379 JetStream(super::errors::Error),
2380 Other,
2381}
2382
2383impl Display for ConsumerErrorKind {
2384 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2385 match self {
2386 Self::TimedOut => write!(f, "timed out"),
2387 Self::Request => write!(f, "request failed"),
2388 Self::JetStream(err) => write!(f, "JetStream error: {err}"),
2389 Self::Other => write!(f, "consumer error"),
2390 Self::InvalidConsumerType => write!(f, "invalid consumer type"),
2391 Self::InvalidName => write!(f, "invalid consumer name"),
2392 }
2393 }
2394}
2395
2396pub type ConsumerError = Error<ConsumerErrorKind>;
2397
2398#[cfg(feature = "server_2_14")]
2399#[cfg_attr(docsrs, doc(cfg(feature = "server_2_14")))]
2400#[derive(Clone, Debug, PartialEq)]
2401pub enum ConsumerResetErrorKind {
2402 TimedOut,
2403 Request,
2404 NotFound,
2410 InvalidReset,
2414 JetStream(super::errors::Error),
2415}
2416
2417#[cfg(feature = "server_2_14")]
2418#[cfg_attr(docsrs, doc(cfg(feature = "server_2_14")))]
2419impl Display for ConsumerResetErrorKind {
2420 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2421 match self {
2422 Self::TimedOut => write!(f, "timed out"),
2423 Self::Request => write!(f, "request failed"),
2424 Self::NotFound => write!(f, "stream or consumer not found"),
2425 Self::InvalidReset => write!(f, "invalid reset"),
2426 Self::JetStream(err) => write!(f, "JetStream error: {err}"),
2427 }
2428 }
2429}
2430
2431#[cfg(feature = "server_2_14")]
2432pub type ConsumerResetError = Error<ConsumerResetErrorKind>;
2433
2434#[cfg(feature = "server_2_14")]
2435impl From<super::errors::Error> for ConsumerResetError {
2436 fn from(err: super::errors::Error) -> Self {
2437 match err.error_code() {
2438 super::errors::ErrorCode::CONSUMER_INVALID_RESET => {
2439 ConsumerResetError::new(ConsumerResetErrorKind::InvalidReset)
2440 }
2441 super::errors::ErrorCode::CONSUMER_NOT_FOUND
2442 | super::errors::ErrorCode::STREAM_NOT_FOUND => {
2443 ConsumerResetError::new(ConsumerResetErrorKind::NotFound)
2444 }
2445 _ => ConsumerResetError::new(ConsumerResetErrorKind::JetStream(err)),
2446 }
2447 }
2448}
2449
2450#[cfg(feature = "server_2_14")]
2451impl From<super::context::RequestError> for ConsumerResetError {
2452 fn from(err: super::context::RequestError) -> Self {
2453 match err.kind() {
2454 RequestErrorKind::TimedOut => ConsumerResetError::new(ConsumerResetErrorKind::TimedOut),
2455 RequestErrorKind::NoResponders => {
2456 ConsumerResetError::new(ConsumerResetErrorKind::NotFound)
2457 }
2458 _ => ConsumerResetError::with_source(ConsumerResetErrorKind::Request, err),
2459 }
2460 }
2461}
2462
2463#[derive(Clone, Debug, PartialEq)]
2464pub enum ConsumerCreateStrictErrorKind {
2465 TimedOut,
2467 Request,
2468 InvalidConsumerType,
2469 InvalidName,
2470 AlreadyExists,
2471 JetStream(super::errors::Error),
2472 Other,
2473}
2474
2475impl Display for ConsumerCreateStrictErrorKind {
2476 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2477 match self {
2478 Self::TimedOut => write!(f, "timed out"),
2479 Self::Request => write!(f, "request failed"),
2480 Self::JetStream(err) => write!(f, "JetStream error: {err}"),
2481 Self::Other => write!(f, "consumer error"),
2482 Self::InvalidConsumerType => write!(f, "invalid consumer type"),
2483 Self::InvalidName => write!(f, "invalid consumer name"),
2484 Self::AlreadyExists => write!(f, "consumer already exists"),
2485 }
2486 }
2487}
2488
2489pub type ConsumerCreateStrictError = Error<ConsumerCreateStrictErrorKind>;
2490
2491#[derive(Clone, Debug, PartialEq)]
2492pub enum ConsumerUpdateErrorKind {
2493 TimedOut,
2495 Request,
2496 InvalidConsumerType,
2497 InvalidName,
2498 DoesNotExist,
2499 JetStream(super::errors::Error),
2500 Other,
2501}
2502
2503impl Display for ConsumerUpdateErrorKind {
2504 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2505 match self {
2506 Self::TimedOut => write!(f, "timed out"),
2507 Self::Request => write!(f, "request failed"),
2508 Self::JetStream(err) => write!(f, "JetStream error: {err}"),
2509 Self::Other => write!(f, "consumer error"),
2510 Self::InvalidConsumerType => write!(f, "invalid consumer type"),
2511 Self::InvalidName => write!(f, "invalid consumer name"),
2512 Self::DoesNotExist => write!(f, "consumer does not exist"),
2513 }
2514 }
2515}
2516
2517pub type ConsumerUpdateError = Error<ConsumerUpdateErrorKind>;
2518
2519impl From<super::errors::Error> for ConsumerError {
2520 fn from(err: super::errors::Error) -> Self {
2521 ConsumerError::new(ConsumerErrorKind::JetStream(err))
2522 }
2523}
2524impl From<super::errors::Error> for ConsumerCreateStrictError {
2525 fn from(err: super::errors::Error) -> Self {
2526 if err.error_code() == super::errors::ErrorCode::CONSUMER_ALREADY_EXISTS {
2527 ConsumerCreateStrictError::new(ConsumerCreateStrictErrorKind::AlreadyExists)
2528 } else {
2529 ConsumerCreateStrictError::new(ConsumerCreateStrictErrorKind::JetStream(err))
2530 }
2531 }
2532}
2533impl From<super::errors::Error> for ConsumerUpdateError {
2534 fn from(err: super::errors::Error) -> Self {
2535 if err.error_code() == super::errors::ErrorCode::CONSUMER_DOES_NOT_EXIST {
2536 ConsumerUpdateError::new(ConsumerUpdateErrorKind::DoesNotExist)
2537 } else {
2538 ConsumerUpdateError::new(ConsumerUpdateErrorKind::JetStream(err))
2539 }
2540 }
2541}
2542impl From<ConsumerError> for ConsumerUpdateError {
2543 fn from(err: ConsumerError) -> Self {
2544 match err.kind() {
2545 ConsumerErrorKind::JetStream(err) => {
2546 if err.error_code() == super::errors::ErrorCode::CONSUMER_DOES_NOT_EXIST {
2547 ConsumerUpdateError::new(ConsumerUpdateErrorKind::DoesNotExist)
2548 } else {
2549 ConsumerUpdateError::new(ConsumerUpdateErrorKind::JetStream(err))
2550 }
2551 }
2552 ConsumerErrorKind::Request => {
2553 ConsumerUpdateError::new(ConsumerUpdateErrorKind::Request)
2554 }
2555 ConsumerErrorKind::TimedOut => {
2556 ConsumerUpdateError::new(ConsumerUpdateErrorKind::TimedOut)
2557 }
2558 ConsumerErrorKind::InvalidConsumerType => {
2559 ConsumerUpdateError::new(ConsumerUpdateErrorKind::InvalidConsumerType)
2560 }
2561 ConsumerErrorKind::InvalidName => {
2562 ConsumerUpdateError::new(ConsumerUpdateErrorKind::InvalidName)
2563 }
2564 ConsumerErrorKind::Other => ConsumerUpdateError::new(ConsumerUpdateErrorKind::Other),
2565 }
2566 }
2567}
2568
2569impl From<ConsumerError> for ConsumerCreateStrictError {
2570 fn from(err: ConsumerError) -> Self {
2571 match err.kind() {
2572 ConsumerErrorKind::JetStream(err) => {
2573 if err.error_code() == super::errors::ErrorCode::CONSUMER_ALREADY_EXISTS {
2574 ConsumerCreateStrictError::new(ConsumerCreateStrictErrorKind::AlreadyExists)
2575 } else {
2576 ConsumerCreateStrictError::new(ConsumerCreateStrictErrorKind::JetStream(err))
2577 }
2578 }
2579 ConsumerErrorKind::Request => {
2580 ConsumerCreateStrictError::new(ConsumerCreateStrictErrorKind::Request)
2581 }
2582 ConsumerErrorKind::TimedOut => {
2583 ConsumerCreateStrictError::new(ConsumerCreateStrictErrorKind::TimedOut)
2584 }
2585 ConsumerErrorKind::InvalidConsumerType => {
2586 ConsumerCreateStrictError::new(ConsumerCreateStrictErrorKind::InvalidConsumerType)
2587 }
2588 ConsumerErrorKind::InvalidName => {
2589 ConsumerCreateStrictError::new(ConsumerCreateStrictErrorKind::InvalidName)
2590 }
2591 ConsumerErrorKind::Other => {
2592 ConsumerCreateStrictError::new(ConsumerCreateStrictErrorKind::Other)
2593 }
2594 }
2595 }
2596}
2597
2598impl From<super::context::RequestError> for ConsumerError {
2599 fn from(err: super::context::RequestError) -> Self {
2600 match err.kind() {
2601 RequestErrorKind::TimedOut => ConsumerError::new(ConsumerErrorKind::TimedOut),
2602 _ => ConsumerError::with_source(ConsumerErrorKind::Request, err),
2603 }
2604 }
2605}
2606impl From<super::context::RequestError> for ConsumerUpdateError {
2607 fn from(err: super::context::RequestError) -> Self {
2608 match err.kind() {
2609 RequestErrorKind::TimedOut => {
2610 ConsumerUpdateError::new(ConsumerUpdateErrorKind::TimedOut)
2611 }
2612 _ => ConsumerUpdateError::with_source(ConsumerUpdateErrorKind::Request, err),
2613 }
2614 }
2615}
2616impl From<super::context::RequestError> for ConsumerCreateStrictError {
2617 fn from(err: super::context::RequestError) -> Self {
2618 match err.kind() {
2619 RequestErrorKind::TimedOut => {
2620 ConsumerCreateStrictError::new(ConsumerCreateStrictErrorKind::TimedOut)
2621 }
2622 _ => {
2623 ConsumerCreateStrictError::with_source(ConsumerCreateStrictErrorKind::Request, err)
2624 }
2625 }
2626 }
2627}
2628
2629#[derive(Debug, Serialize, Default)]
2630pub struct DirectGetRequest {
2631 #[serde(rename = "seq", skip_serializing_if = "Option::is_none")]
2632 sequence: Option<u64>,
2633 #[serde(rename = "last_by_subj", skip_serializing)]
2634 last_by_subject: Option<String>,
2635 #[serde(rename = "next_by_subj", skip_serializing_if = "Option::is_none")]
2636 next_by_subject: Option<String>,
2637}
2638
2639pub struct WithHeaders;
2641
2642pub struct WithoutHeaders;
2644
2645trait DirectGetResponse: Sized {
2647 fn from_message(message: crate::Message) -> Result<Self, DirectGetError>;
2648}
2649
2650impl DirectGetResponse for StreamMessage {
2651 fn from_message(message: crate::Message) -> Result<Self, DirectGetError> {
2652 StreamMessage::try_from(message).map_err(Into::into)
2653 }
2654}
2655
2656impl DirectGetResponse for StreamValue {
2657 fn from_message(message: crate::Message) -> Result<Self, DirectGetError> {
2658 Ok(StreamValue {
2659 data: message.payload,
2660 })
2661 }
2662}
2663
2664pub struct DirectGetBuilder<T = WithHeaders> {
2665 context: Context,
2666 stream_name: String,
2667 request: DirectGetRequest,
2668 _phantom: std::marker::PhantomData<T>,
2669}
2670
2671impl DirectGetBuilder<WithHeaders> {
2672 fn new(context: Context, stream_name: String) -> DirectGetBuilder<WithHeaders> {
2673 DirectGetBuilder {
2674 context,
2675 stream_name,
2676 request: DirectGetRequest::default(),
2677 _phantom: std::marker::PhantomData,
2678 }
2679 }
2680}
2681
2682impl<T> DirectGetBuilder<T> {
2683 async fn send_internal<R: DirectGetResponse>(&self) -> Result<R, DirectGetError> {
2685 let payload = if self.request.last_by_subject.is_some() {
2688 Bytes::new()
2689 } else {
2690 serde_json::to_vec(&self.request).map(Bytes::from)?
2691 };
2692
2693 let request_subject = if let Some(ref subject) = self.request.last_by_subject {
2694 format!(
2695 "{}.DIRECT.GET.{}.{}",
2696 self.context.prefix, self.stream_name, subject
2697 )
2698 } else {
2699 format!("{}.DIRECT.GET.{}", self.context.prefix, self.stream_name)
2700 };
2701
2702 let response = self
2703 .context
2704 .client
2705 .request(request_subject, payload)
2706 .await?;
2707
2708 if let Some(status) = response.status {
2710 if let Some(ref description) = response.description {
2711 match status {
2712 StatusCode::NOT_FOUND => {
2713 return Err(DirectGetError::new(DirectGetErrorKind::NotFound))
2714 }
2715 StatusCode::TIMEOUT => {
2717 return Err(DirectGetError::new(DirectGetErrorKind::InvalidSubject))
2718 }
2719 _ => {
2720 return Err(DirectGetError::new(DirectGetErrorKind::ErrorResponse(
2721 status,
2722 description.to_string(),
2723 )));
2724 }
2725 }
2726 }
2727 }
2728
2729 R::from_message(response)
2730 }
2731
2732 pub fn sequence(mut self, seq: u64) -> Self {
2734 self.request.sequence = Some(seq);
2735 self
2736 }
2737
2738 pub fn last_by_subject<S: Into<String>>(mut self, subject: S) -> Self {
2740 self.request.last_by_subject = Some(subject.into());
2741 self
2742 }
2743
2744 pub fn next_by_subject<S: Into<String>>(mut self, subject: S) -> Self {
2746 self.request.next_by_subject = Some(subject.into());
2747 self
2748 }
2749}
2750
2751impl DirectGetBuilder<WithHeaders> {
2752 pub async fn send(self) -> Result<StreamMessage, DirectGetError> {
2754 self.send_internal::<StreamMessage>().await
2755 }
2756}
2757
2758impl DirectGetBuilder<WithoutHeaders> {
2759 pub async fn send(self) -> Result<StreamValue, DirectGetError> {
2761 self.send_internal::<StreamValue>().await
2762 }
2763}
2764
2765pub struct StreamValue {
2766 pub data: Bytes,
2767}
2768
2769#[derive(Debug, Serialize, Default)]
2770pub struct RawMessageRequest {
2771 #[serde(rename = "seq", skip_serializing_if = "Option::is_none")]
2772 sequence: Option<u64>,
2773 #[serde(rename = "last_by_subj", skip_serializing_if = "Option::is_none")]
2774 last_by_subject: Option<String>,
2775 #[serde(rename = "next_by_subj", skip_serializing_if = "Option::is_none")]
2776 next_by_subject: Option<String>,
2777}
2778
2779trait RawMessageResponse: Sized {
2781 fn from_raw_message(message: RawMessage) -> Result<Self, RawMessageError>;
2782}
2783
2784impl RawMessageResponse for StreamMessage {
2785 fn from_raw_message(message: RawMessage) -> Result<Self, RawMessageError> {
2786 StreamMessage::try_from(message)
2787 .map_err(|err| RawMessageError::with_source(RawMessageErrorKind::Other, err))
2788 }
2789}
2790
2791impl RawMessageResponse for StreamValue {
2792 fn from_raw_message(message: RawMessage) -> Result<Self, RawMessageError> {
2793 use base64::engine::general_purpose::STANDARD;
2794 use base64::Engine;
2795
2796 let decoded_payload = STANDARD.decode(message.payload).map_err(|err| {
2797 RawMessageError::with_source(
2798 RawMessageErrorKind::Other,
2799 Box::new(std::io::Error::other(err)),
2800 )
2801 })?;
2802
2803 Ok(StreamValue {
2804 data: decoded_payload.into(),
2805 })
2806 }
2807}
2808
2809pub struct RawMessageBuilder<T = WithHeaders> {
2810 context: Context,
2811 stream_name: String,
2812 request: RawMessageRequest,
2813 _phantom: std::marker::PhantomData<T>,
2814}
2815
2816impl RawMessageBuilder<WithHeaders> {
2817 fn new(context: Context, stream_name: String) -> Self {
2818 RawMessageBuilder {
2819 context,
2820 stream_name,
2821 request: RawMessageRequest::default(),
2822 _phantom: std::marker::PhantomData,
2823 }
2824 }
2825}
2826
2827impl<T> RawMessageBuilder<T> {
2828 async fn send_internal<R: RawMessageResponse>(&self) -> Result<R, RawMessageError> {
2830 for subject in [&self.request.last_by_subject, &self.request.next_by_subject]
2832 .into_iter()
2833 .flatten()
2834 {
2835 if !is_valid_subject(subject) {
2836 return Err(RawMessageError::new(RawMessageErrorKind::InvalidSubject));
2837 }
2838 }
2839
2840 let subject = format!("STREAM.MSG.GET.{}", self.stream_name);
2841
2842 let response: Response<GetRawMessage> = self
2843 .context
2844 .request(subject, &self.request)
2845 .map_err(|err| RawMessageError::with_source(RawMessageErrorKind::Other, err))
2846 .await?;
2847
2848 match response {
2849 Response::Err { error } => {
2850 if error.error_code() == ErrorCode::NO_MESSAGE_FOUND {
2851 Err(RawMessageError::new(RawMessageErrorKind::NoMessageFound))
2852 } else {
2853 Err(RawMessageError::new(RawMessageErrorKind::JetStream(error)))
2854 }
2855 }
2856 Response::Ok(value) => R::from_raw_message(value.message),
2857 }
2858 }
2859
2860 pub fn sequence(mut self, seq: u64) -> Self {
2862 self.request.sequence = Some(seq);
2863 self
2864 }
2865
2866 pub fn last_by_subject<S: Into<String>>(mut self, subject: S) -> Self {
2868 self.request.last_by_subject = Some(subject.into());
2869 self
2870 }
2871
2872 pub fn next_by_subject<S: Into<String>>(mut self, subject: S) -> Self {
2874 self.request.next_by_subject = Some(subject.into());
2875 self
2876 }
2877}
2878
2879impl RawMessageBuilder<WithHeaders> {
2880 pub async fn send(self) -> Result<StreamMessage, RawMessageError> {
2882 self.send_internal::<StreamMessage>().await
2883 }
2884}
2885
2886impl RawMessageBuilder<WithoutHeaders> {
2887 pub async fn send(self) -> Result<StreamValue, RawMessageError> {
2889 self.send_internal::<StreamValue>().await
2890 }
2891}
2892
2893#[cfg(test)]
2894mod tests {
2895 use super::*;
2896
2897 #[test]
2898 fn consumer_limits_de() {
2899 let config = Config {
2900 ..Default::default()
2901 };
2902
2903 let roundtrip: Config = {
2904 let ser = serde_json::to_string(&config).unwrap();
2905 serde_json::from_str(&ser).unwrap()
2906 };
2907 assert_eq!(config, roundtrip);
2908 }
2909}