|
1 | 1 | use super::{BackupsManager, DiskBackupPolicy, DiskBackupsManager, InMemoryBackupsManager}; |
2 | | -use crate::SiftChannel; |
3 | | -use crate::TimeValue; |
4 | 2 | use crate::backup::disk::AsyncBackupsManager; |
5 | | -use hyper_util::rt::TokioIo; |
6 | | -use sift_connect::grpc::interceptor::AuthInterceptor; |
| 3 | +use crate::{TimeValue, backup::sanitize_name}; |
7 | 4 | use sift_error::ErrorKind; |
8 | 5 | use sift_rs::ingest::v1::{ |
9 | 6 | IngestWithConfigDataChannelValue, IngestWithConfigDataStreamRequest, |
10 | 7 | ingest_with_config_data_channel_value::Type, |
11 | 8 | }; |
12 | | -use sift_rs::ingest::v1::{ |
13 | | - IngestWithConfigDataStreamResponse, |
14 | | - ingest_service_server::{IngestService, IngestServiceServer}, |
15 | | -}; |
16 | 9 | use std::fs; |
17 | | -use std::sync::{Arc, Mutex}; |
18 | 10 | use tempdir::TempDir; |
19 | | -use tonic::transport::{Endpoint, Server, Uri}; |
20 | | -use tonic::{Request, Response, Status}; |
21 | | -use tower::{ServiceBuilder, service_fn}; |
22 | | - |
23 | | -#[derive(Debug, Clone)] |
24 | | -struct MockIngestService { |
25 | | - captured_data: Arc<Mutex<Vec<IngestWithConfigDataStreamRequest>>>, |
26 | | -} |
27 | | - |
28 | | -impl MockIngestService { |
29 | | - fn new() -> Self { |
30 | | - Self { |
31 | | - captured_data: Arc::new(Mutex::new(Vec::new())), |
32 | | - } |
33 | | - } |
34 | 11 |
|
35 | | - fn get_captured_data(&self) -> Vec<IngestWithConfigDataStreamRequest> { |
36 | | - self.captured_data.lock().unwrap().clone() |
| 12 | +#[test] |
| 13 | +fn test_sanitize_name_with_illegal_chars() { |
| 14 | + let illegal_chars = vec![ |
| 15 | + ':', '/', '\\', '*', '?', '"', '<', '>', '|', '.', ' ', '\t', '\n', '\r', |
| 16 | + ]; |
| 17 | + for char in illegal_chars { |
| 18 | + assert_eq!(sanitize_name(&format!("test{}test", char)), "test_test"); |
37 | 19 | } |
38 | 20 | } |
39 | 21 |
|
40 | | -#[tonic::async_trait] |
41 | | -impl IngestService for MockIngestService { |
42 | | - async fn ingest_with_config_data_stream( |
43 | | - &self, |
44 | | - request: Request<tonic::Streaming<IngestWithConfigDataStreamRequest>>, |
45 | | - ) -> Result<Response<IngestWithConfigDataStreamResponse>, Status> { |
46 | | - let mut stream = request.into_inner(); |
47 | | - |
48 | | - while let Some(data) = stream.message().await? { |
49 | | - self.captured_data.lock().unwrap().push(data); |
50 | | - } |
51 | | - |
52 | | - Ok(Response::new(IngestWithConfigDataStreamResponse {})) |
53 | | - } |
54 | | - |
55 | | - async fn ingest_arbitrary_protobuf_data_stream( |
56 | | - &self, |
57 | | - _request: Request< |
58 | | - tonic::Streaming<sift_rs::ingest::v1::IngestArbitraryProtobufDataStreamRequest>, |
59 | | - >, |
60 | | - ) -> Result<Response<sift_rs::ingest::v1::IngestArbitraryProtobufDataStreamResponse>, Status> |
61 | | - { |
62 | | - Err(Status::unimplemented("Not implemented for test")) |
63 | | - } |
64 | | -} |
65 | | - |
66 | | -async fn create_mock_grpc_channel_with_service() -> (SiftChannel, Arc<MockIngestService>) { |
67 | | - let mock_service = Arc::new(MockIngestService::new()); |
68 | | - let service_clone = mock_service.clone(); |
69 | | - |
70 | | - let (client, server) = tokio::io::duplex(1024); |
71 | | - |
72 | | - tokio::spawn(async move { |
73 | | - Server::builder() |
74 | | - .add_service(IngestServiceServer::new(service_clone.as_ref().clone())) |
75 | | - .serve_with_incoming(tokio_stream::once(Ok::<_, std::io::Error>(server))) |
76 | | - .await |
77 | | - .unwrap(); |
78 | | - }); |
79 | | - |
80 | | - let mut client = Some(client); |
81 | | - let channel = Endpoint::try_from("http://[::]:50051") |
82 | | - .unwrap() |
83 | | - .connect_with_connector(service_fn(move |_: Uri| { |
84 | | - let client = client.take(); |
85 | | - |
86 | | - async move { |
87 | | - if let Some(client) = client { |
88 | | - Ok(TokioIo::new(client)) |
89 | | - } else { |
90 | | - Err(std::io::Error::other("Client already taken")) |
91 | | - } |
92 | | - } |
93 | | - })) |
94 | | - .await |
95 | | - .unwrap(); |
96 | | - |
97 | | - let sift_channel = ServiceBuilder::new() |
98 | | - .layer(tonic::service::interceptor(AuthInterceptor { |
99 | | - apikey: "test-api-key".to_string(), |
100 | | - })) |
101 | | - .service(channel); |
102 | | - |
103 | | - (sift_channel, mock_service) |
| 22 | +#[test] |
| 23 | +fn test_sanitize_name_with_legal_chars() { |
| 24 | + assert_eq!(sanitize_name("test"), "test"); |
| 25 | + assert_eq!(sanitize_name("test_test"), "test_test"); |
| 26 | + assert_eq!(sanitize_name("test-test"), "test-test"); |
104 | 27 | } |
105 | 28 |
|
106 | 29 | #[tokio::test] |
@@ -229,7 +152,7 @@ async fn test_async_backups_manager_retrieve_data_with_graceful_termination() { |
229 | 152 | ..Default::default() |
230 | 153 | }; |
231 | 154 | let backup_retry_policy = crate::RetryPolicy::default(); |
232 | | - let (grpc_channel, mock_service) = create_mock_grpc_channel_with_service().await; |
| 155 | + let (grpc_channel, mock_service) = crate::test::create_mock_grpc_channel_with_service().await; |
233 | 156 |
|
234 | 157 | let mut backups_manager = AsyncBackupsManager::<IngestWithConfigDataStreamRequest>::new( |
235 | 158 | &backups_dir, |
@@ -308,7 +231,7 @@ async fn test_async_backups_manager_discard_data_with_graceful_termination() { |
308 | 231 | ..Default::default() |
309 | 232 | }; |
310 | 233 | let backup_retry_policy = crate::RetryPolicy::default(); |
311 | | - let (grpc_channel, mock_service) = create_mock_grpc_channel_with_service().await; |
| 234 | + let (grpc_channel, mock_service) = crate::test::create_mock_grpc_channel_with_service().await; |
312 | 235 |
|
313 | 236 | let mut backups_manager = AsyncBackupsManager::<IngestWithConfigDataStreamRequest>::new( |
314 | 237 | &backups_dir, |
|
0 commit comments