use std::time::Duration;
use tokio::sync::mpsc::UnboundedSender;
use tonic::Streaming;
use typedb_protocol::{database, database_manager, migration::Item, transaction};
use uuid::Uuid;
use crate::{
Credentials, QueryOptions, TransactionOptions, TransactionType,
analyze::AnalyzedQuery,
answer::{
QueryType,
concept_document::{ConceptDocumentHeader, Node},
concept_row::ConceptRowHeader,
},
common::{RequestID, info::DatabaseInfo},
concept::Concept,
connection::server::{Server, server_version::ServerVersion},
error::ServerError,
info::UserInfo,
};
#[derive(Debug)]
pub(super) enum Request {
ConnectionOpen { driver_lang: String, driver_version: String, credentials: Credentials },
ServersAll,
ServersGet,
ServerVersion,
DatabasesAll,
DatabaseGet { database_name: String },
DatabasesContains { database_name: String },
DatabaseCreate { database_name: String },
DatabaseImport(DatabaseImportRequest),
DatabaseDelete { database_name: String },
DatabaseSchema { database_name: String },
DatabaseTypeSchema { database_name: String },
DatabaseExport { database_name: String },
Transaction(TransactionRequest),
UsersAll,
UsersGet { name: String },
UsersContains { name: String },
UsersCreate { user: UserInfo },
UsersUpdate { username: String, user: UserInfo },
UsersDelete { name: String },
}
#[derive(Debug)]
pub(super) enum Response {
ConnectionOpen {
connection_id: Uuid,
server_duration_millis: u64,
servers: Vec<Server>,
},
ServersAll {
servers: Vec<Server>,
},
ServersGet {
server: Server,
},
ServerVersion {
server_version: ServerVersion,
},
DatabasesContains {
contains: bool,
},
DatabaseCreate {
database: DatabaseInfo,
},
DatabaseImport {
request_sink: UnboundedSender<database_manager::import::Client>,
response_source: Streaming<database_manager::import::Server>,
},
DatabaseGet {
database: DatabaseInfo,
},
DatabasesAll {
databases: Vec<DatabaseInfo>,
},
DatabaseDelete,
DatabaseSchema {
schema: String,
},
DatabaseTypeSchema {
schema: String,
},
DatabaseExportStream {
response_source: Streaming<database::export::Server>,
},
TransactionStream {
open_request_id: RequestID,
request_sink: UnboundedSender<transaction::Client>,
response_source: Streaming<transaction::Server>,
server_duration_millis: u64,
},
UsersAll {
users: Vec<UserInfo>,
},
UsersContain {
contains: bool,
},
UsersCreate,
UsersUpdate,
UsersDelete,
UsersGet {
user: Option<UserInfo>,
},
}
#[derive(Debug)]
pub(super) enum DatabaseImportRequest {
Initial { name: String, schema: String },
ItemPart { items: Vec<Item> },
Done,
}
#[derive(Debug)]
pub(super) enum DatabaseExportResponse {
Schema(String),
Items(Vec<Item>),
Done,
}
#[derive(Debug)]
pub(super) enum TransactionRequest {
Open { database: String, transaction_type: TransactionType, options: TransactionOptions, network_latency: Duration },
Commit,
Rollback,
Analyze { query: String },
Query(QueryRequest),
Stream { request_id: RequestID },
}
#[derive(Debug)]
pub(super) enum TransactionResponse {
Open { server_duration_millis: u64 },
Commit,
Rollback,
Query(QueryResponse),
Analyze(AnalyzeResponse),
Close,
}
#[derive(Debug)]
pub(super) enum AnalyzeResponse {
Ok(AnalyzedQuery),
Err(ServerError),
}
#[derive(Debug)]
pub(super) enum QueryRequest {
Query { query: String, options: QueryOptions },
}
#[derive(Debug)]
pub(super) enum QueryResponse {
Ok(QueryType),
ConceptRowsHeader(ConceptRowHeader),
ConceptDocumentsHeader(ConceptDocumentHeader),
StreamConceptRows(Vec<(Vec<Option<Concept>>, Option<Vec<u8>>)>),
StreamConceptDocuments(Vec<Option<Node>>),
Error(ServerError),
}