typedb_driver/database/
database.rs1use std::{
21 io::{BufWriter, Write},
22 path::Path,
23 sync::Arc,
24};
25
26use prost::Message;
27
28use crate::{
29 common::{Error, Result, info::DatabaseInfo},
30 connection::server::{server_manager::ServerManager, server_routing::ServerRouting},
31 database::migration::{DatabaseExportAnswer, try_create_export_file, try_open_existing_export_file},
32 error::MigrationError,
33 resolve,
34};
35
36#[derive(Debug, Clone)]
38pub struct Database {
39 name: String,
40 server_manager: Arc<ServerManager>,
41}
42
43impl Database {
44 pub(super) fn new(database_info: DatabaseInfo, server_manager: Arc<ServerManager>) -> Result<Self> {
45 Ok(Self { name: database_info.name, server_manager })
46 }
47
48 #[cfg_attr(feature = "sync", maybe_async::must_be_sync)]
49 pub(super) async fn get(name: String, server_manager: Arc<ServerManager>) -> Result<Self> {
50 Ok(Self { name, server_manager })
51 }
52
53 pub fn name(&self) -> &str {
55 self.name.as_str()
56 }
57
58 #[cfg_attr(feature = "sync", doc = "database.delete();")]
64 #[cfg_attr(not(feature = "sync"), doc = "database.delete().await;")]
65 #[cfg_attr(feature = "sync", maybe_async::must_be_sync)]
67 pub async fn delete(self: Arc<Self>) -> Result {
68 self.server_manager
69 .execute(ServerRouting::Auto, |server_connection| {
70 let name = self.name.clone();
71 async move { server_connection.delete_database(name).await }
72 })
73 .await
74 }
75
76 #[cfg_attr(feature = "sync", doc = "database.schema();")]
82 #[cfg_attr(not(feature = "sync"), doc = "database.schema().await;")]
83 #[cfg_attr(feature = "sync", maybe_async::must_be_sync)]
85 pub async fn schema(&self) -> Result<String> {
86 self.server_manager
87 .execute(ServerRouting::Auto, |server_connection| {
88 let name = self.name.clone();
89 async move { server_connection.database_schema(name).await }
90 })
91 .await
92 }
93
94 #[cfg_attr(feature = "sync", doc = "database.type_schema();")]
100 #[cfg_attr(not(feature = "sync"), doc = "database.type_schema().await;")]
101 #[cfg_attr(feature = "sync", maybe_async::must_be_sync)]
103 pub async fn type_schema(&self) -> Result<String> {
104 self.server_manager
105 .execute(ServerRouting::Auto, |server_connection| {
106 let name = self.name.clone();
107 async move { server_connection.database_type_schema(name).await }
108 })
109 .await
110 }
111
112 #[cfg_attr(feature = "sync", doc = "database.export_to_file(schema_path, data_path);")]
124 #[cfg_attr(not(feature = "sync"), doc = "database.export_to_file(schema_path, data_path).await;")]
125 #[cfg_attr(feature = "sync", maybe_async::must_be_sync)]
127 pub async fn export_to_file(&self, schema_file_path: impl AsRef<Path>, data_file_path: impl AsRef<Path>) -> Result {
128 let schema_file_path = schema_file_path.as_ref();
129 let data_file_path = data_file_path.as_ref();
130 if schema_file_path == data_file_path {
131 return Err(Error::Migration(MigrationError::CannotExportToTheSameFile));
132 }
133
134 let _ = try_create_export_file(schema_file_path)?;
135 if let Err(err) = try_create_export_file(data_file_path) {
136 let _ = std::fs::remove_file(schema_file_path);
137 return Err(err);
138 }
139
140 let result = self
141 .server_manager
142 .execute(ServerRouting::Auto, |server_connection| {
143 let name = self.name.clone();
144 async move {
145 let mut schema_file = try_open_existing_export_file(schema_file_path)?;
147 let data_file = try_open_existing_export_file(data_file_path)?;
148 let mut export_stream = server_connection.database_export(name).await?;
149 let mut data_writer = BufWriter::new(data_file);
150
151 loop {
152 match resolve!(export_stream.next())? {
153 DatabaseExportAnswer::Done => break,
154 DatabaseExportAnswer::Schema(schema) => {
155 schema_file.write_all(schema.as_bytes())?;
156 schema_file.flush()?;
157 }
158 DatabaseExportAnswer::Items(items) => {
159 for item in items {
160 let mut buf = Vec::new();
161 item.encode_length_delimited(&mut buf)
162 .map_err(|_| Error::Migration(MigrationError::CannotEncodeExportedConcept))?;
163 data_writer.write_all(&buf)?;
164 }
165 }
166 }
167 }
168
169 data_writer.flush()?;
170 Ok(())
171 }
172 })
173 .await;
174
175 if result.is_err() {
176 let _ = std::fs::remove_file(schema_file_path);
177 let _ = std::fs::remove_file(data_file_path);
178 }
179 result
180 }
181}