Websocket Support (#159)

This commit is contained in:
Gregory Schier
2025-01-31 09:00:11 -08:00
committed by GitHub
parent d411713502
commit c8be8082c5
122 changed files with 5090 additions and 616 deletions
+387 -6
View File
@@ -1,6 +1,15 @@
use crate::error::Error::ModelNotFound;
use crate::error::Result;
use crate::models::{AnyModel, CookieJar, CookieJarIden, Environment, EnvironmentIden, Folder, FolderIden, GrpcConnection, GrpcConnectionIden, GrpcConnectionState, GrpcEvent, GrpcEventIden, GrpcRequest, GrpcRequestIden, HttpRequest, HttpRequestIden, HttpResponse, HttpResponseHeader, HttpResponseIden, HttpResponseState, KeyValue, KeyValueIden, ModelType, Plugin, PluginIden, PluginKeyValue, PluginKeyValueIden, Settings, SettingsIden, SyncState, SyncStateIden, Workspace, WorkspaceIden, WorkspaceMeta, WorkspaceMetaIden};
use crate::models::{
AnyModel, CookieJar, CookieJarIden, Environment, EnvironmentIden, Folder, FolderIden,
GrpcConnection, GrpcConnectionIden, GrpcConnectionState, GrpcEvent, GrpcEventIden, GrpcRequest,
GrpcRequestIden, HttpRequest, HttpRequestIden, HttpResponse, HttpResponseHeader,
HttpResponseIden, HttpResponseState, KeyValue, KeyValueIden, ModelType, Plugin, PluginIden,
PluginKeyValue, PluginKeyValueIden, Settings, SettingsIden, SyncState, SyncStateIden,
WebsocketConnection, WebsocketConnectionIden, WebsocketEvent, WebsocketEventIden,
WebsocketRequest, WebsocketRequestIden, Workspace, WorkspaceIden, WorkspaceMeta,
WorkspaceMetaIden,
};
use crate::plugin::SqliteConnection;
use chrono::{NaiveDateTime, Utc};
use log::{debug, error, info, warn};
@@ -16,8 +25,7 @@ use std::path::Path;
use tauri::{AppHandle, Emitter, Listener, Manager, Runtime, WebviewWindow};
use ts_rs::TS;
const MAX_GRPC_CONNECTIONS_PER_REQUEST: usize = 20;
const MAX_HTTP_RESPONSES_PER_REQUEST: usize = MAX_GRPC_CONNECTIONS_PER_REQUEST;
const MAX_HISTORY_ITEMS: usize = 20;
pub async fn set_key_value_string<R: Runtime>(
mgr: &WebviewWindow<R>,
@@ -659,8 +667,8 @@ pub async fn upsert_grpc_connection<R: Runtime>(
update_source: &UpdateSource,
) -> Result<GrpcConnection> {
let connections =
list_http_responses_for_request(window, connection.request_id.as_str(), None).await?;
for c in connections.iter().skip(MAX_GRPC_CONNECTIONS_PER_REQUEST - 1) {
list_grpc_connections_for_request(window, connection.request_id.as_str()).await?;
for c in connections.iter().skip(MAX_HISTORY_ITEMS - 1) {
debug!("Deleting old grpc connection {}", c.id);
delete_grpc_connection(window, c.id.as_str(), update_source).await?;
}
@@ -911,6 +919,367 @@ pub async fn list_grpc_events<R: Runtime>(
Ok(items.map(|v| v.unwrap()).collect())
}
pub async fn delete_websocket_request<R: Runtime>(
window: &WebviewWindow<R>,
id: &str,
update_source: &UpdateSource,
) -> Result<WebsocketRequest> {
let request = match get_websocket_request(window, id).await? {
Some(r) => r,
None => {
return Err(ModelNotFound(id.to_string()));
}
};
let dbm = &*window.app_handle().state::<SqliteConnection>();
let db = dbm.0.lock().await.get().unwrap();
let (sql, params) = Query::delete()
.from_table(WebsocketRequestIden::Table)
.cond_where(Expr::col(WebsocketRequestIden::Id).eq(id))
.build_rusqlite(SqliteQueryBuilder);
db.execute(sql.as_str(), &*params.as_params())?;
emit_deleted_model(window, &AnyModel::WebsocketRequest(request.to_owned()), update_source);
Ok(request)
}
pub async fn delete_websocket_connection<R: Runtime>(
window: &WebviewWindow<R>,
id: &str,
update_source: &UpdateSource,
) -> Result<WebsocketConnection> {
let m = get_websocket_connection(window, id).await?;
let dbm = &*window.app_handle().state::<SqliteConnection>();
let db = dbm.0.lock().await.get().unwrap();
let (sql, params) = Query::delete()
.from_table(WebsocketConnectionIden::Table)
.cond_where(Expr::col(WebsocketConnectionIden::Id).eq(id))
.build_rusqlite(SqliteQueryBuilder);
db.execute(sql.as_str(), &*params.as_params())?;
emit_deleted_model(window, &AnyModel::WebsocketConnection(m.to_owned()), update_source);
Ok(m)
}
pub async fn delete_all_websocket_connections<R: Runtime>(
window: &WebviewWindow<R>,
request_id: &str,
update_source: &UpdateSource,
) -> Result<()> {
for c in list_websocket_connections_for_request(window, request_id).await? {
delete_websocket_connection(window, &c.id, update_source).await?;
}
Ok(())
}
pub async fn delete_all_websocket_connections_for_workspace<R: Runtime>(
window: &WebviewWindow<R>,
workspace_id: &str,
update_source: &UpdateSource,
) -> Result<()> {
for c in list_websocket_connections_for_workspace(window, workspace_id).await? {
delete_websocket_connection(window, &c.id, update_source).await?;
}
Ok(())
}
pub async fn get_websocket_connection<R: Runtime>(
mgr: &impl Manager<R>,
id: &str,
) -> Result<WebsocketConnection> {
let dbm = &*mgr.state::<SqliteConnection>();
let db = dbm.0.lock().await.get().unwrap();
let (sql, params) = Query::select()
.from(WebsocketConnectionIden::Table)
.column(Asterisk)
.cond_where(Expr::col(WebsocketConnectionIden::Id).eq(id))
.build_rusqlite(SqliteQueryBuilder);
let mut stmt = db.prepare(sql.as_str())?;
Ok(stmt.query_row(&*params.as_params(), |row| row.try_into())?)
}
pub async fn upsert_websocket_event<R: Runtime>(
window: &WebviewWindow<R>,
event: WebsocketEvent,
update_source: &UpdateSource,
) -> Result<WebsocketEvent> {
let id = match event.id.as_str() {
"" => generate_model_id(ModelType::TypeWebSocketEvent),
_ => event.id.to_string(),
};
let dbm = &*window.app_handle().state::<SqliteConnection>();
let db = dbm.0.lock().await.get().unwrap();
let (sql, params) = Query::insert()
.into_table(WebsocketEventIden::Table)
.columns([
WebsocketEventIden::Id,
WebsocketEventIden::CreatedAt,
WebsocketEventIden::UpdatedAt,
WebsocketEventIden::WorkspaceId,
WebsocketEventIden::ConnectionId,
WebsocketEventIden::RequestId,
WebsocketEventIden::MessageType,
WebsocketEventIden::IsServer,
WebsocketEventIden::Message,
])
.values_panic([
id.into(),
timestamp_for_upsert(update_source, event.created_at).into(),
timestamp_for_upsert(update_source, event.updated_at).into(),
event.workspace_id.into(),
event.connection_id.into(),
event.request_id.into(),
serde_json::to_string(&event.message_type)?.into(),
event.is_server.into(),
event.message.into(),
])
.on_conflict(
OnConflict::column(WebsocketEventIden::Id)
.update_columns([
WebsocketEventIden::UpdatedAt,
WebsocketEventIden::MessageType,
WebsocketEventIden::IsServer,
WebsocketEventIden::Message,
])
.to_owned(),
)
.returning_all()
.build_rusqlite(SqliteQueryBuilder);
let mut stmt = db.prepare(sql.as_str())?;
let m: WebsocketEvent = stmt.query_row(&*params.as_params(), |row| row.try_into())?;
emit_upserted_model(window, &AnyModel::WebsocketEvent(m.to_owned()), update_source);
Ok(m)
}
pub async fn upsert_websocket_request<R: Runtime>(
window: &WebviewWindow<R>,
request: WebsocketRequest,
update_source: &UpdateSource,
) -> Result<WebsocketRequest> {
let id = match request.id.as_str() {
"" => generate_model_id(ModelType::TypeWebsocketRequest),
_ => request.id.to_string(),
};
let trimmed_name = request.name.trim();
let dbm = &*window.app_handle().state::<SqliteConnection>();
let db = dbm.0.lock().await.get().unwrap();
let (sql, params) = Query::insert()
.into_table(WebsocketRequestIden::Table)
.columns([
WebsocketRequestIden::Id,
WebsocketRequestIden::CreatedAt,
WebsocketRequestIden::UpdatedAt,
WebsocketRequestIden::WorkspaceId,
WebsocketRequestIden::FolderId,
WebsocketRequestIden::Authentication,
WebsocketRequestIden::AuthenticationType,
WebsocketRequestIden::Description,
WebsocketRequestIden::Headers,
WebsocketRequestIden::Message,
WebsocketRequestIden::Name,
WebsocketRequestIden::SortPriority,
WebsocketRequestIden::Url,
WebsocketRequestIden::UrlParameters,
])
.values_panic([
id.into(),
timestamp_for_upsert(update_source, request.created_at).into(),
timestamp_for_upsert(update_source, request.updated_at).into(),
request.workspace_id.into(),
request.folder_id.as_ref().map(|s| s.as_str()).into(),
serde_json::to_string(&request.authentication)?.into(),
request.authentication_type.as_ref().map(|s| s.as_str()).into(),
request.description.into(),
serde_json::to_string(&request.headers)?.into(),
request.message.into(),
trimmed_name.into(),
request.sort_priority.into(),
request.url.into(),
serde_json::to_string(&request.url_parameters)?.into(),
])
.on_conflict(
OnConflict::column(WebsocketRequestIden::Id)
.update_columns([
WebsocketRequestIden::UpdatedAt,
WebsocketRequestIden::WorkspaceId,
WebsocketRequestIden::FolderId,
WebsocketRequestIden::Authentication,
WebsocketRequestIden::AuthenticationType,
WebsocketRequestIden::Description,
WebsocketRequestIden::Headers,
WebsocketRequestIden::Message,
WebsocketRequestIden::Name,
WebsocketRequestIden::SortPriority,
WebsocketRequestIden::Url,
WebsocketRequestIden::UrlParameters,
])
.to_owned(),
)
.returning_all()
.build_rusqlite(SqliteQueryBuilder);
let mut stmt = db.prepare(sql.as_str())?;
let m: WebsocketRequest = stmt.query_row(&*params.as_params(), |row| row.try_into())?;
emit_upserted_model(window, &AnyModel::WebsocketRequest(m.to_owned()), update_source);
Ok(m)
}
pub async fn list_websocket_connections_for_workspace<R: Runtime>(
mgr: &impl Manager<R>,
workspace_id: &str,
) -> Result<Vec<WebsocketConnection>> {
let dbm = &*mgr.state::<SqliteConnection>();
let db = dbm.0.lock().await.get().unwrap();
let (sql, params) = Query::select()
.from(WebsocketConnectionIden::Table)
.cond_where(Expr::col(WebsocketConnectionIden::WorkspaceId).eq(workspace_id))
.column(Asterisk)
.order_by(WebsocketConnectionIden::CreatedAt, Order::Desc)
.build_rusqlite(SqliteQueryBuilder);
let mut stmt = db.prepare(sql.as_str())?;
let items = stmt.query_map(&*params.as_params(), |row| row.try_into())?;
Ok(items.map(|v| v.unwrap()).collect())
}
pub async fn list_websocket_connections_for_request<R: Runtime>(
mgr: &impl Manager<R>,
request_id: &str,
) -> Result<Vec<WebsocketConnection>> {
let dbm = &*mgr.state::<SqliteConnection>();
let db = dbm.0.lock().await.get().unwrap();
let (sql, params) = Query::select()
.from(WebsocketConnectionIden::Table)
.cond_where(Expr::col(WebsocketConnectionIden::RequestId).eq(request_id))
.column(Asterisk)
.order_by(WebsocketConnectionIden::CreatedAt, Order::Desc)
.build_rusqlite(SqliteQueryBuilder);
let mut stmt = db.prepare(sql.as_str())?;
let items = stmt.query_map(&*params.as_params(), |row| row.try_into())?;
Ok(items.map(|v| v.unwrap()).collect())
}
pub async fn upsert_websocket_connection<R: Runtime>(
window: &WebviewWindow<R>,
connection: &WebsocketConnection,
update_source: &UpdateSource,
) -> Result<WebsocketConnection> {
let connections =
list_websocket_connections_for_request(window, connection.request_id.as_str()).await?;
for c in connections.iter().skip(MAX_HISTORY_ITEMS - 1) {
debug!("Deleting old websocket connection {}", c.id);
delete_websocket_connection(window, c.id.as_str(), update_source).await?;
}
let id = match connection.id.as_str() {
"" => generate_model_id(ModelType::TypeWebSocketConnection),
_ => connection.id.to_string(),
};
let dbm = &*window.app_handle().state::<SqliteConnection>();
let db = dbm.0.lock().await.get().unwrap();
let (sql, params) = Query::insert()
.into_table(WebsocketConnectionIden::Table)
.columns([
WebsocketConnectionIden::Id,
WebsocketConnectionIden::CreatedAt,
WebsocketConnectionIden::UpdatedAt,
WebsocketConnectionIden::WorkspaceId,
WebsocketConnectionIden::RequestId,
WebsocketConnectionIden::Elapsed,
WebsocketConnectionIden::Error,
WebsocketConnectionIden::Headers,
WebsocketConnectionIden::State,
WebsocketConnectionIden::Status,
WebsocketConnectionIden::Url,
])
.values_panic([
id.as_str().into(),
timestamp_for_upsert(update_source, connection.created_at).into(),
timestamp_for_upsert(update_source, connection.updated_at).into(),
connection.workspace_id.as_str().into(),
connection.request_id.as_str().into(),
connection.elapsed.into(),
connection.error.as_ref().map(|s| s.as_str()).into(),
serde_json::to_string(&connection.headers)?.into(),
serde_json::to_value(&connection.state)?.as_str().into(),
connection.status.into(),
connection.url.as_str().into(),
])
.on_conflict(
OnConflict::column(WebsocketConnectionIden::Id)
.update_columns([
WebsocketConnectionIden::UpdatedAt,
WebsocketConnectionIden::Elapsed,
WebsocketConnectionIden::Error,
WebsocketConnectionIden::Headers,
WebsocketConnectionIden::State,
WebsocketConnectionIden::Status,
WebsocketConnectionIden::Url,
])
.to_owned(),
)
.returning_all()
.build_rusqlite(SqliteQueryBuilder);
let mut stmt = db.prepare(sql.as_str())?;
let m: WebsocketConnection = stmt.query_row(&*params.as_params(), |row| row.try_into())?;
emit_upserted_model(window, &AnyModel::WebsocketConnection(m.to_owned()), update_source);
Ok(m)
}
pub async fn get_websocket_request<R: Runtime>(
mgr: &impl Manager<R>,
id: &str,
) -> Result<Option<WebsocketRequest>> {
let dbm = &*mgr.state::<SqliteConnection>();
let db = dbm.0.lock().await.get().unwrap();
let (sql, params) = Query::select()
.from(WebsocketRequestIden::Table)
.column(Asterisk)
.cond_where(Expr::col(WebsocketRequestIden::Id).eq(id))
.build_rusqlite(SqliteQueryBuilder);
let mut stmt = db.prepare(sql.as_str())?;
Ok(stmt.query_row(&*params.as_params(), |row| row.try_into()).optional()?)
}
pub async fn list_websocket_requests<R: Runtime>(
mgr: &impl Manager<R>,
workspace_id: &str,
) -> Result<Vec<WebsocketRequest>> {
let dbm = &*mgr.state::<SqliteConnection>();
let db = dbm.0.lock().await.get().unwrap();
let (sql, params) = Query::select()
.from(WebsocketRequestIden::Table)
.cond_where(Expr::col(WebsocketRequestIden::WorkspaceId).eq(workspace_id))
.column(Asterisk)
.build_rusqlite(SqliteQueryBuilder);
let mut stmt = db.prepare(sql.as_str())?;
let items = stmt.query_map(&*params.as_params(), |row| row.try_into())?;
Ok(items.map(|v| v.unwrap()).collect())
}
pub async fn list_websocket_events<R: Runtime>(
mgr: &impl Manager<R>,
connection_id: &str,
) -> Result<Vec<WebsocketEvent>> {
let dbm = &*mgr.state::<SqliteConnection>();
let db = dbm.0.lock().await.get().unwrap();
let (sql, params) = Query::select()
.from(WebsocketEventIden::Table)
.cond_where(Expr::col(WebsocketEventIden::ConnectionId).eq(connection_id))
.column(Asterisk)
.build_rusqlite(SqliteQueryBuilder);
let mut stmt = db.prepare(sql.as_str())?;
let items = stmt.query_map(&*params.as_params(), |row| row.try_into())?;
Ok(items.map(|v| v.unwrap()).collect())
}
pub async fn upsert_cookie_jar<R: Runtime>(
window: &WebviewWindow<R>,
cookie_jar: &CookieJar,
@@ -1676,7 +2045,7 @@ pub async fn create_http_response<R: Runtime>(
update_source: &UpdateSource,
) -> Result<HttpResponse> {
let responses = list_http_responses_for_request(window, request_id, None).await?;
for response in responses.iter().skip(MAX_HTTP_RESPONSES_PER_REQUEST - 1) {
for response in responses.iter().skip(MAX_HISTORY_ITEMS - 1) {
debug!("Deleting old response {}", response.id);
delete_http_response(window, response.id.as_str(), update_source).await?;
}
@@ -2170,6 +2539,7 @@ pub struct BatchUpsertResult {
pub folders: Vec<Folder>,
pub http_requests: Vec<HttpRequest>,
pub grpc_requests: Vec<GrpcRequest>,
pub websocket_requests: Vec<WebsocketRequest>,
}
pub async fn batch_upsert<R: Runtime>(
@@ -2179,6 +2549,7 @@ pub async fn batch_upsert<R: Runtime>(
folders: Vec<Folder>,
http_requests: Vec<HttpRequest>,
grpc_requests: Vec<GrpcRequest>,
websocket_requests: Vec<WebsocketRequest>,
update_source: &UpdateSource,
) -> Result<BatchUpsertResult> {
let mut imported_resources = BatchUpsertResult::default();
@@ -2246,6 +2617,14 @@ pub async fn batch_upsert<R: Runtime>(
info!("Imported {} grpc_requests", imported_resources.grpc_requests.len());
}
if websocket_requests.len() > 0 {
for v in websocket_requests {
let x = upsert_websocket_request(&window, v, update_source).await?;
imported_resources.websocket_requests.push(x.clone());
}
info!("Imported {} websocket_requests", imported_resources.websocket_requests.len());
}
Ok(imported_resources)
}
@@ -2264,6 +2643,7 @@ pub async fn get_workspace_export_resources<R: Runtime>(
folders: Vec::new(),
http_requests: Vec::new(),
grpc_requests: Vec::new(),
websocket_requests: Vec::new(),
},
};
@@ -2273,6 +2653,7 @@ pub async fn get_workspace_export_resources<R: Runtime>(
data.resources.folders.append(&mut list_folders(mgr, workspace_id).await?);
data.resources.http_requests.append(&mut list_http_requests(mgr, workspace_id).await?);
data.resources.grpc_requests.append(&mut list_grpc_requests(mgr, workspace_id).await?);
data.resources.websocket_requests.append(&mut list_websocket_requests(mgr, workspace_id).await?);
}
// Nuke environments if we don't want them