pub struct Connection {
pub config: Arc<RwLock<ClientConfig>>,
pub command_tx: Sender<Message>,
/* private fields */
}Expand description
Управляет соединением с сервером CloudPub.
Структура Connection предоставляет высокоуровневый интерфейс для взаимодействия с
сервером CloudPub, обработки аутентификации, регистрации сервисов и
управления жизненным циклом.
§Пример
use cloudpub_sdk::Connection;
use cloudpub_common::protocol::{Protocol, Auth, Endpoint};
// Создание и настройка соединения
let mut conn = Connection::builder()
.credentials("user@example.com", "password")
.timeout_secs(30)
.build()
.await?;
// Регистрация сервиса
let endpoint = conn.publish(
Protocol::Http,
"localhost:8080".to_string(),
Some("My Service".to_string()),
Some(Auth::None),
None,
None,
None,
).await?;
println!("Сервис доступен по адресу: {}", endpoint.as_url());
// Список всех сервисов
let services = conn.ls().await?;
for service in services {
let name = service.client.as_ref()
.and_then(|c| c.description.clone())
.unwrap_or_else(|| "Безымянный".to_string());
println!("Сервис: {} - {}", name, service.as_url());
}§Потокобезопасность
Connection использует внутреннюю синхронизацию и может быть безопасно разделен между
потоками с использованием Arc<Mutex<Connection>> при необходимости.
Fields§
§config: Arc<RwLock<ClientConfig>>§command_tx: Sender<Message>Implementations§
Source§impl Connection
impl Connection
Sourcepub fn builder() -> ConnectionBuilder
pub fn builder() -> ConnectionBuilder
Создает новый построитель соединения для настройки параметров соединения.
Это рекомендуемый способ создания нового соединения.
§Пример
use cloudpub_sdk::Connection;
let conn = Connection::builder()
.config_path("/path/to/config.toml")
.credentials("user@example.com", "password")
.timeout_secs(30)
.verbose(true)
.build()
.await?;Sourcepub async fn wait_for_event(
&self,
target_event: impl Fn(&ConnectionEvent) -> bool + Send,
) -> Result<ConnectionEvent>
pub async fn wait_for_event( &self, target_event: impl Fn(&ConnectionEvent) -> bool + Send, ) -> Result<ConnectionEvent>
Ожидает наступления определенного события с таймаутом.
Этот метод блокируется до получения целевого события или истечения таймаута. Он используется внутренне другими методами, но также может использоваться напрямую для пользовательской обработки событий.
§Аргументы
target_event- Предикатная функция, которая возвращает true при получении желаемого события
§Возвращает
Возвращает соответствующее событие или ошибку, если:
- Истек таймаут
- Соединение закрыто
- Получено событие ошибки
§Пример
use cloudpub_sdk::ConnectionEvent;
// Ожидать установки соединения
let event = conn.wait_for_event(|e| matches!(e, ConnectionEvent::Connected)).await?;
// Ожидать любой операции с эндпоинтом
let event = conn.wait_for_event(|e| matches!(e, ConnectionEvent::Endpoint(_))).await?;
if let ConnectionEvent::Endpoint(endpoint) = event {
println!("Эндпоинт зарегистрирован: {}", endpoint.guid);
}Sourcepub fn logout(&self) -> Result<()>
pub fn logout(&self) -> Result<()>
Выходит из системы, очищая токен аутентификации.
Это удаляет сохраненный токен аутентификации как из памяти, так и из файла конфигурации. Соединение нужно будет повторно аутентифицировать для дальнейших операций.
§Пример
conn.logout()?;
println!("Выход выполнен успешно");Sourcepub async fn register(
&mut self,
protocol: Protocol,
address: String,
name: Option<String>,
auth: Option<Auth>,
acl: Option<Vec<Acl>>,
headers: Option<Vec<Header>>,
rules: Option<Vec<FilterRule>>,
) -> Result<ServerEndpoint>
pub async fn register( &mut self, protocol: Protocol, address: String, name: Option<String>, auth: Option<Auth>, acl: Option<Vec<Acl>>, headers: Option<Vec<Header>>, rules: Option<Vec<FilterRule>>, ) -> Result<ServerEndpoint>
Регистрирует сервис на сервере CloudPub.
Этот метод регистрирует локальный сервис для доступа через CloudPub. Сервису будет присвоен уникальный GUID и URL для удаленного доступа.
§Аргументы
protocol- Тип протокола (HTTP, HTTPS, TCP, UDP, WS, WSS, RTSP)address- Локальный адрес сервиса. Для RTSP с авторизацией используйте формат:rtsp://user:pass@host:port/pathname- Опциональное человекочитаемое имя для сервисаauth- Опциональный метод аутентификации для доступа к сервисуacl- Опциональный список контроля доступа для фильтрации IPheaders- Опциональные HTTP заголовки для добавления к ответамrules- Опциональные правила фильтрации для фильтрации запросов
§Возвращает
Возвращает ServerEndpoint, содержащий:
guid: Уникальный идентификатор для сервиса- URL информацию для доступа к сервису
- Статус сервиса и метаданные
§Пример
use cloudpub_common::protocol::{Protocol, Auth, Endpoint, Acl, Header, FilterRule};
// HTTP service
let endpoint = conn.register(
Protocol::Http,
"localhost:3000".to_string(),
Some("My Web App".to_string()),
Some(Auth::Basic),
None,
None,
None,
).await?;
// HTTP service with ACL and headers
let acl = vec![Acl { user: "admin".to_string(), role: cloudpub_common::protocol::Role::Admin as i32 }];
let headers = vec![Header { name: "X-Custom".to_string(), value: "test".to_string() }];
let endpoint = conn.register(
Protocol::Http,
"localhost:8080".to_string(),
Some("API Server".to_string()),
Some(Auth::None),
Some(acl),
Some(headers),
None,
).await?;
// RTSP service with credentials in URL
let rtsp_endpoint = conn.register(
Protocol::Rtsp,
"rtsp://camera:secret@localhost:554/stream".to_string(),
Some("Security Camera".to_string()),
Some(Auth::None),
None,
None,
None,
).await?;
println!("Сервис зарегистрирован по адресу: {}", endpoint.as_url());
println!("GUID сервиса: {}", endpoint.guid);Sourcepub async fn publish(
&mut self,
protocol: Protocol,
address: String,
name: Option<String>,
auth: Option<Auth>,
acl: Option<Vec<Acl>>,
headers: Option<Vec<Header>>,
rules: Option<Vec<FilterRule>>,
) -> Result<ServerEndpoint>
pub async fn publish( &mut self, protocol: Protocol, address: String, name: Option<String>, auth: Option<Auth>, acl: Option<Vec<Acl>>, headers: Option<Vec<Header>>, rules: Option<Vec<FilterRule>>, ) -> Result<ServerEndpoint>
Публикует сервис на сервере CloudPub.
Это псевдоним для register(), который запускает сервис сразу после регистрации
§Аргументы
protocol- Тип протокола (HTTP, HTTPS, TCP, UDP, WS, WSS, RTSP)address- Локальный адрес сервиса. Для RTSP с авторизацией используйте формат:rtsp://user:pass@host:port/pathname- Опциональное человекочитаемое имя для сервисаauth- Опциональный метод аутентификации для доступа к сервисуacl- Опциональный список контроля доступа для фильтрации IPheaders- Опциональные HTTP заголовки для добавления к ответамrules- Опциональные правила фильтрации для фильтрации запросов
§Возвращает
Возвращает ServerEndpoint с деталями сервиса и URL доступа.
§Пример
use cloudpub_common::protocol::{Protocol, Auth, Endpoint, Acl, FilterRule};
// Publish a TCP service
let endpoint = conn.publish(
Protocol::Tcp,
"localhost:8080".to_string(),
Some("TCP Server".to_string()),
Some(Auth::None),
None,
None,
None,
).await?;
// Publish HTTP service with ACL
let acl = vec![Acl { user: "reader".to_string(), role: cloudpub_common::protocol::Role::Reader as i32 }];
let endpoint = conn.publish(
Protocol::Http,
"localhost:3000".to_string(),
Some("Internal API".to_string()),
Some(Auth::Basic),
Some(acl),
None,
None,
).await?;
// Publish RTSP with embedded credentials
let rtsp = conn.publish(
Protocol::Rtsp,
"rtsp://admin:password@192.168.1.100:554/live/ch0".to_string(),
Some("IP Camera".to_string()),
Some(Auth::Basic),
None,
None,
None,
).await?;
println!("TCP сервис доступен по адресу: {}", endpoint.as_url());Sourcepub async fn ls(&mut self) -> Result<Vec<ServerEndpoint>>
pub async fn ls(&mut self) -> Result<Vec<ServerEndpoint>>
Выводит список всех зарегистрированных сервисов.
Возвращает список всех сервисов, в данный момент зарегистрированных на сервере, включая их статус, URL и метаданные.
§Возвращает
Вектор структур ServerEndpoint, содержащих информацию о сервисах.
§Пример
use cloudpub_common::protocol::Endpoint;
let services = conn.ls().await?;
for service in services {
let name = service.client.as_ref()
.and_then(|c| c.description.clone())
.unwrap_or_else(|| "Безымянный".to_string());
println!("Сервис: {} ({})" ,
name,
service.guid
);
println!(" URL: {}", service.as_url());
println!(" Статус: {}", service.status.unwrap_or_else(|| "Неизвестен".to_string()));
}Sourcepub async fn start(&mut self, guid: String) -> Result<()>
pub async fn start(&mut self, guid: String) -> Result<()>
Запускает сервис по его GUID.
Запускает ранее зарегистрированный сервис, который мог быть остановлен. Это делает сервис снова доступным через его CloudPub URL.
§Аргументы
guid- Уникальный идентификатор сервиса для запуска
§Пример
let services = conn.ls().await?;
if let Some(service) = services.first() {
conn.start(service.guid.clone()).await?;
println!("Запущен сервис: {}", service.guid);
}Sourcepub async fn stop(&mut self, guid: String) -> Result<()>
pub async fn stop(&mut self, guid: String) -> Result<()>
Останавливает сервис по его GUID.
Временно останавливает сервис, делая его недоступным через CloudPub. Регистрация сервиса сохраняется и может быть перезапущена позже.
§Arguments
guid- Уникальный идентификатор сервиса для остановки
§Пример
conn.stop("service-guid-123".to_string()).await?;
println!("Сервис остановлен");Sourcepub async fn unpublish(&mut self, guid: String) -> Result<()>
pub async fn unpublish(&mut self, guid: String) -> Result<()>
Отменяет публикацию (удаляет) сервис по его GUID.
Навсегда удаляет регистрацию сервиса с сервера. Сервис больше не будет доступен через CloudPub.
§Arguments
guid- Уникальный идентификатор сервиса для отмены публикации
§Пример
conn.unpublish("service-guid-123".to_string()).await?;
println!("Публикация сервиса отменена и удалена");Sourcepub async fn clean(&mut self) -> Result<()>
pub async fn clean(&mut self) -> Result<()>
Удаляет все зарегистрированные сервисы.
Этот метод удаляет все сервисы, зарегистрированные текущим пользователем. Используйте с осторожностью, так как эта операция не может быть отменена.
§Пример
conn.clean().await?;
println!("Все сервисы удалены");
let services = conn.ls().await?;
assert_eq!(services.len(), 0);Sourcepub async fn ping(&mut self) -> Result<u64>
pub async fn ping(&mut self) -> Result<u64>
Пингует сервер для измерения задержки.
Создает временную конечную точку пинга и измеряет время прохождения туда и обратно до сервера. Полезно для проверки работоспособности соединения и задержки.
§Возвращает
Задержку пинга в микросекундах.
§Пример
let latency_us = conn.ping().await?;
println!("Задержка сервера: {}μs ({:.2}мс)", latency_us, latency_us as f64 / 1000.0);
if latency_us > 100_000 { // 100ms
println!("Предупреждение: Обнаружена высокая задержка");
}Trait Implementations§
Source§impl Drop for Connection
impl Drop for Connection
Auto Trait Implementations§
impl !RefUnwindSafe for Connection
impl !UnwindSafe for Connection
impl Freeze for Connection
impl Send for Connection
impl Sync for Connection
impl Unpin for Connection
impl UnsafeUnpin for Connection
Blanket Implementations§
Source§impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
Source§impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
fn into_either(self, into_left: bool) -> Either<Self, Self>
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more