1use crate::error::{KitError, Result};
4use crate::internal::{ensure_internal_tables, internal_tables_core};
5use crate::schema::{to_core_indexes, to_core_schema};
6use mongreldb_core::epoch::Snapshot;
7use mongreldb_core::memtable::Row as CoreRow;
8use mongreldb_core::memtable::Value as CoreValue;
9use mongreldb_core::schema::Schema as CoreSchema;
10use mongreldb_core::Database as CoreDatabase;
11use mongreldb_core::{AggState, ApproxAgg, NativeAgg, NativeAggResult, RowId};
12use mongreldb_kit_core::schema::Index as KitIndex;
13use mongreldb_kit_core::schema::IndexKind as KitIndexKind;
14use mongreldb_kit_core::schema::Schema as KitSchema;
15use mongreldb_kit_core::schema::Table as KitTable;
16use mongreldb_kit_core::{ProcedureSpec, TriggerSpec, ViewSpec};
17use serde_json::Value;
18
19use std::collections::HashMap;
20use std::path::{Path, PathBuf};
21use std::sync::Arc;
22
23const SCHEMA_FILE: &str = "kit_schema.json";
24
25#[derive(Clone, Copy, Debug, Default)]
30pub struct OpenOptions {
31 pub lock_timeout_ms: u32,
37}
38
39impl OpenOptions {
40 pub fn new() -> Self {
42 Self::default()
43 }
44
45 pub fn with_lock_timeout_ms(mut self, ms: u32) -> Self {
48 self.lock_timeout_ms = ms;
49 self
50 }
51}
52
53#[derive(Debug, Clone, Default)]
54pub struct SqlOptions {
55 pub query_id: Option<mongreldb_query::QueryId>,
56 pub timeout: Option<std::time::Duration>,
57}
58
59#[derive(Debug, Clone, Copy, PartialEq, Eq)]
60pub struct SqlOutputLimits {
61 pub max_rows: usize,
62 pub max_bytes: usize,
63}
64
65impl Default for SqlOutputLimits {
66 fn default() -> Self {
67 Self {
68 max_rows: 1_000_000,
69 max_bytes: 64 * 1024 * 1024,
70 }
71 }
72}
73
74pub struct SqlQueryHandle {
75 query_id: mongreldb_query::QueryId,
76 session: Arc<mongreldb_query::MongrelSession>,
77 worker: Option<
78 std::thread::JoinHandle<mongreldb_query::Result<mongreldb_query::ManagedQueryBatches>>,
79 >,
80}
81
82impl SqlQueryHandle {
83 pub fn id(&self) -> mongreldb_query::QueryId {
84 self.query_id
85 }
86
87 pub fn cancel(&self) -> mongreldb_query::CancelOutcome {
88 self.session.cancel_query(self.query_id)
89 }
90
91 pub fn status(&self) -> Option<mongreldb_query::QueryStatus> {
92 self.session.query_registry().status(self.query_id)
93 }
94
95 pub fn wait(self) -> Result<Vec<arrow::record_batch::RecordBatch>> {
96 let output = self.wait_for_serialization()?;
97 let batches = output.batches().to_vec();
98 complete_sql_output(output)?;
99 Ok(batches)
100 }
101
102 pub fn wait_arrow(self) -> Result<Vec<u8>> {
103 self.wait_arrow_with_limits(SqlOutputLimits::default())
104 }
105
106 pub fn wait_arrow_with_limits(self, limits: SqlOutputLimits) -> Result<Vec<u8>> {
107 let output = self.wait_for_serialization()?;
108 let result = crate::arrow_util::batches_to_ipc_controlled_with_limits(
109 output.batches(),
110 output.query(),
111 limits,
112 );
113 match result {
114 Ok(bytes) => {
115 complete_sql_output(output)?;
116 Ok(bytes)
117 }
118 Err(error) => {
119 fail_sql_output(output, &error);
120 Err(error)
121 }
122 }
123 }
124
125 pub fn wait_rows(self) -> Result<Vec<serde_json::Map<String, Value>>> {
126 self.wait_rows_with_limits(SqlOutputLimits::default())
127 }
128
129 pub fn wait_rows_with_limits(
130 self,
131 limits: SqlOutputLimits,
132 ) -> Result<Vec<serde_json::Map<String, Value>>> {
133 let output = self.wait_for_serialization()?;
134 let result = crate::arrow_util::batches_to_rows_controlled_with_limits(
135 output.batches(),
136 output.query(),
137 limits,
138 );
139 match result {
140 Ok(rows) => {
141 complete_sql_output(output)?;
142 Ok(rows)
143 }
144 Err(error) => {
145 fail_sql_output(output, &error);
146 Err(error)
147 }
148 }
149 }
150
151 pub fn wait_for_serialization(mut self) -> Result<mongreldb_query::ManagedQueryBatches> {
152 let result = self
153 .worker
154 .take()
155 .expect("SQL worker is present")
156 .join()
157 .map_err(|_| KitError::Storage("SQL worker panicked".into()))?;
158 let output = result.map_err(|error| {
159 let status = self.session.query_registry().status(self.query_id);
160 crate::error::query_error_with_status(error, status.as_ref())
161 })?;
162 self.session
163 .fire_test_hook(mongreldb_query::SqlTestHookPoint::BeforeSerializationBatch);
164 Ok(output)
165 }
166}
167
168#[doc(hidden)]
170pub fn complete_sql_output(output: mongreldb_query::ManagedQueryBatches) -> Result<()> {
171 let query = output.query().clone();
172 loop {
173 if let Err(error) = query.checkpoint() {
174 let status = query.status();
175 output.fail();
176 return Err(crate::error::query_error_with_status(error, Some(&status)));
177 }
178 let phase = query.phase();
179 match phase {
180 mongreldb_query::SqlQueryPhase::Serializing
181 | mongreldb_query::SqlQueryPhase::CommitCritical => {
182 if query
183 .transition(phase, mongreldb_query::SqlQueryPhase::Completed)
184 .is_ok()
185 {
186 break;
187 }
188 }
189 mongreldb_query::SqlQueryPhase::Completed => break,
190 mongreldb_query::SqlQueryPhase::Cancelling => std::thread::yield_now(),
191 phase => {
192 let error = mongreldb_query::MongrelQueryError::InvalidQueryState(format!(
193 "query {} cannot complete output conversion from {phase:?}",
194 query.id()
195 ));
196 let status = query.status();
197 output.fail_with_error(
198 error.code(),
199 mongreldb_query::QueryTerminalErrorCategory::Execution,
200 );
201 return Err(crate::error::query_error_with_status(error, Some(&status)));
202 }
203 }
204 }
205 output.complete().map_err(|error| {
206 let status = query.status();
207 crate::error::query_error_with_status(error, Some(&status))
208 })
209}
210
211#[doc(hidden)]
213pub fn fail_sql_output(output: mongreldb_query::ManagedQueryBatches, error: &KitError) {
214 match error {
215 KitError::ResultLimitExceeded { .. } => output.fail_result_limit(),
216 KitError::SerializationFailed { .. } => output.fail_serialization(),
217 _ => output.fail(),
218 }
219}
220
221impl Drop for SqlQueryHandle {
222 fn drop(&mut self) {
223 if self.worker.is_some() {
224 let _ = self.session.cancel_query(self.query_id);
225 }
226 }
227}
228
229pub type DefaultProvider = Box<dyn Fn() -> Value + Send + Sync>;
231
232#[derive(Debug, Clone)]
235pub struct ExplainPlan {
236 pub index_accelerated: bool,
238 pub exact: bool,
241 pub pushed_conditions: Vec<String>,
243}
244
245#[derive(Debug, Clone)]
247pub struct SimilarRow {
248 pub row: crate::schema::Row,
249 pub similarity: f64,
250}
251
252fn parse_string_set(value: Option<&Value>) -> std::collections::HashSet<String> {
256 let arr = match value {
257 Some(Value::Array(a)) => Some(a.clone()),
258 Some(Value::String(s)) => serde_json::from_str::<Value>(s)
259 .ok()
260 .and_then(|v| v.as_array().cloned()),
261 _ => None,
262 };
263 arr.into_iter()
264 .flatten()
265 .filter_map(|v| match v {
266 Value::String(s) => Some(s),
267 Value::Number(n) => Some(n.to_string()),
268 Value::Bool(b) => Some(b.to_string()),
269 _ => None,
270 })
271 .collect()
272}
273
274#[derive(Debug, Clone, Copy, PartialEq, Eq)]
276pub enum IncrementalAggKind {
277 Count,
278 Sum,
279 Min,
280 Max,
281 Avg,
282}
283
284#[derive(Debug, Clone)]
286pub struct IncrementalAggregate {
287 pub value: Value,
290 pub incremental: bool,
294 pub delta_rows: u64,
296}
297
298fn incremental_cache_key(
302 table_id: u32,
303 column: Option<u16>,
304 agg: IncrementalAggKind,
305 conditions: &[mongreldb_core::query::Condition],
306) -> u64 {
307 use std::hash::{Hash, Hasher};
308 let mut h = std::collections::hash_map::DefaultHasher::new();
309 table_id.hash(&mut h);
310 column.hash(&mut h);
311 (agg as u8).hash(&mut h);
312 format!("{conditions:?}").hash(&mut h);
314 h.finish()
315}
316
317fn agg_state_value(s: &AggState) -> Value {
321 let num_f64 = |x: f64| {
322 serde_json::Number::from_f64(x)
323 .map(Value::Number)
324 .unwrap_or(Value::Null)
325 };
326 match s {
327 AggState::Count(n) => Value::from(*n),
328 AggState::SumI { sum, .. } => i64::try_from(*sum)
329 .map(Value::from)
330 .unwrap_or_else(|_| num_f64(*sum as f64)),
331 AggState::SumF { sum, .. } => num_f64(*sum),
332 AggState::AvgI { sum, count } if *count > 0 => num_f64(*sum as f64 / *count as f64),
333 AggState::AvgF { sum, count } if *count > 0 => num_f64(*sum / *count as f64),
334 AggState::AvgI { .. } | AggState::AvgF { .. } => Value::Null,
335 AggState::MinI(n) | AggState::MaxI(n) => Value::from(*n),
336 AggState::MinF(f) | AggState::MaxF(f) => num_f64(*f),
337 AggState::Empty => Value::Null,
338 }
339}
340
341#[derive(Debug, Clone, Copy, PartialEq, Eq)]
343pub enum ApproxAggKind {
344 Count,
345 Sum,
346 Avg,
347}
348
349#[derive(Debug, Clone)]
353pub struct ApproxAggregate {
354 pub point: f64,
355 pub ci_low: f64,
356 pub ci_high: f64,
357 pub n_population: u64,
358 pub n_sample_live: usize,
359 pub n_passing: usize,
360}
361
362fn condition_label(c: &mongreldb_core::query::Condition) -> String {
365 let dbg = format!("{c:?}");
366 dbg.split(['(', '{', ' ']).next().unwrap_or("").to_string()
367}
368
369fn open_core_with_retry<T>(
370 timeout_ms: u32,
371 mut open: impl FnMut() -> mongreldb_core::Result<T>,
372) -> mongreldb_core::Result<T> {
373 if timeout_ms == 0 {
374 return open();
375 }
376 let deadline = std::time::Instant::now() + std::time::Duration::from_millis(timeout_ms as u64);
377 let mut next_sleep = std::time::Duration::from_millis(1);
378 loop {
379 match open() {
380 Ok(db) => return Ok(db),
381 Err(err) if is_lock_contention(&err) => {
382 let now = std::time::Instant::now();
383 if now >= deadline {
384 return Err(err);
385 }
386 let sleep = next_sleep.min(deadline - now);
387 std::thread::sleep(sleep);
388 next_sleep = next_sleep
389 .saturating_mul(10)
390 .min(std::time::Duration::from_millis(50));
391 }
392 Err(err) => return Err(err),
393 }
394 }
395}
396
397fn is_lock_contention(err: &mongreldb_core::MongrelError) -> bool {
398 matches!(err, mongreldb_core::MongrelError::DatabaseLocked { .. })
399}
400
401pub struct Database {
406 pub(crate) inner: Arc<CoreDatabase>,
407 pub(crate) schema: KitSchema,
408 pub(crate) root: PathBuf,
409 pub(crate) default_providers: HashMap<String, DefaultProvider>,
411 pub(crate) session: parking_lot::RwLock<Option<Arc<mongreldb_query::MongrelSession>>>,
418 sequence_lock: parking_lot::Mutex<()>,
419}
420
421impl Database {
422 pub fn open(path: &Path) -> Result<Self> {
424 let inner = Arc::new(CoreDatabase::open(path)?);
425 let schema = load_schema(path)?;
426 ensure_internal_tables(&inner)?;
428 reap_rotated_wal_segments(&inner);
429 Ok(Self {
430 inner,
431 schema,
432 root: path.to_path_buf(),
433 default_providers: HashMap::new(),
434 session: parking_lot::RwLock::new(None),
435 sequence_lock: parking_lot::Mutex::new(()),
436 })
437 }
438
439 pub fn open_with_options(path: &Path, opts: OpenOptions) -> Result<Self> {
447 let inner = Arc::new(open_core_with_retry(opts.lock_timeout_ms, || {
448 CoreDatabase::open(path)
449 })?);
450 let schema = load_schema(path)?;
451 ensure_internal_tables(&inner)?;
452 reap_rotated_wal_segments(&inner);
453 Ok(Self {
454 inner,
455 schema,
456 root: path.to_path_buf(),
457 default_providers: HashMap::new(),
458 session: parking_lot::RwLock::new(None),
459 sequence_lock: parking_lot::Mutex::new(()),
460 })
461 }
462
463 pub fn open_encrypted(path: &Path, passphrase: &str) -> Result<Self> {
465 let inner = Arc::new(CoreDatabase::open_encrypted(path, passphrase)?);
466 let schema = load_schema(path)?;
467 ensure_internal_tables(&inner)?;
468 reap_rotated_wal_segments(&inner);
469 Ok(Self {
470 inner,
471 schema,
472 root: path.to_path_buf(),
473 default_providers: HashMap::new(),
474 session: parking_lot::RwLock::new(None),
475 sequence_lock: parking_lot::Mutex::new(()),
476 })
477 }
478
479 pub fn open_encrypted_with_options(
483 path: &Path,
484 passphrase: &str,
485 opts: OpenOptions,
486 ) -> Result<Self> {
487 let inner = Arc::new(open_core_with_retry(opts.lock_timeout_ms, || {
488 CoreDatabase::open_encrypted(path, passphrase)
489 })?);
490 let schema = load_schema(path)?;
491 ensure_internal_tables(&inner)?;
492 reap_rotated_wal_segments(&inner);
493 Ok(Self {
494 inner,
495 schema,
496 root: path.to_path_buf(),
497 default_providers: HashMap::new(),
498 session: parking_lot::RwLock::new(None),
499 sequence_lock: parking_lot::Mutex::new(()),
500 })
501 }
502
503 pub fn create_encrypted(path: &Path, schema: KitSchema, passphrase: &str) -> Result<Self> {
507 std::fs::create_dir_all(path)?;
508 let inner = Arc::new(CoreDatabase::create_encrypted(path, passphrase)?);
509 ensure_internal_tables(&inner)?;
510 store_schema(path, &schema)?;
511 for table in &schema.tables {
512 create_core_table(&inner, &table.name, to_core_schema(table)?)?;
513 }
514 Ok(Self {
515 inner,
516 schema,
517 root: path.to_path_buf(),
518 default_providers: HashMap::new(),
519 session: parking_lot::RwLock::new(None),
520 sequence_lock: parking_lot::Mutex::new(()),
521 })
522 }
523
524 pub fn create(path: &Path, schema: KitSchema) -> Result<Self> {
526 std::fs::create_dir_all(path)?;
527 let inner = Arc::new(CoreDatabase::create(path)?);
528
529 ensure_internal_tables(&inner)?;
532
533 store_schema(path, &schema)?;
536
537 for table in &schema.tables {
539 create_core_table(&inner, &table.name, to_core_schema(table)?)?;
540 }
541
542 Ok(Self {
543 inner,
544 schema,
545 root: path.to_path_buf(),
546 default_providers: HashMap::new(),
547 session: parking_lot::RwLock::new(None),
548 sequence_lock: parking_lot::Mutex::new(()),
549 })
550 }
551
552 pub fn open_with_credentials(path: &Path, username: &str, password: &str) -> Result<Self> {
562 let inner = Arc::new(CoreDatabase::open_with_credentials(
563 path, username, password,
564 )?);
565 let schema = load_schema(path)?;
566 ensure_internal_tables(&inner)?;
567 reap_rotated_wal_segments(&inner);
568 Ok(Self {
569 inner,
570 schema,
571 root: path.to_path_buf(),
572 default_providers: HashMap::new(),
573 session: parking_lot::RwLock::new(None),
574 sequence_lock: parking_lot::Mutex::new(()),
575 })
576 }
577
578 pub fn open_with_credentials_and_options(
582 path: &Path,
583 username: &str,
584 password: &str,
585 opts: OpenOptions,
586 ) -> Result<Self> {
587 let inner = Arc::new(open_core_with_retry(opts.lock_timeout_ms, || {
588 CoreDatabase::open_with_credentials(path, username, password)
589 })?);
590 let schema = load_schema(path)?;
591 ensure_internal_tables(&inner)?;
592 reap_rotated_wal_segments(&inner);
593 Ok(Self {
594 inner,
595 schema,
596 root: path.to_path_buf(),
597 default_providers: HashMap::new(),
598 session: parking_lot::RwLock::new(None),
599 sequence_lock: parking_lot::Mutex::new(()),
600 })
601 }
602
603 pub fn create_with_credentials(
609 path: &Path,
610 schema: KitSchema,
611 admin_username: &str,
612 admin_password: &str,
613 ) -> Result<Self> {
614 std::fs::create_dir_all(path)?;
615 let inner = Arc::new(CoreDatabase::create_with_credentials(
616 path,
617 admin_username,
618 admin_password,
619 )?);
620 ensure_internal_tables(&inner)?;
621 store_schema(path, &schema)?;
622 for table in &schema.tables {
623 create_core_table(&inner, &table.name, to_core_schema(table)?)?;
624 }
625 Ok(Self {
626 inner,
627 schema,
628 root: path.to_path_buf(),
629 default_providers: HashMap::new(),
630 session: parking_lot::RwLock::new(None),
631 sequence_lock: parking_lot::Mutex::new(()),
632 })
633 }
634
635 pub fn open_encrypted_with_credentials(
638 path: &Path,
639 passphrase: &str,
640 username: &str,
641 password: &str,
642 ) -> Result<Self> {
643 let inner = Arc::new(CoreDatabase::open_encrypted_with_credentials(
644 path, passphrase, username, password,
645 )?);
646 let schema = load_schema(path)?;
647 ensure_internal_tables(&inner)?;
648 reap_rotated_wal_segments(&inner);
649 Ok(Self {
650 inner,
651 schema,
652 root: path.to_path_buf(),
653 default_providers: HashMap::new(),
654 session: parking_lot::RwLock::new(None),
655 sequence_lock: parking_lot::Mutex::new(()),
656 })
657 }
658
659 pub fn open_encrypted_with_credentials_and_options(
663 path: &Path,
664 passphrase: &str,
665 username: &str,
666 password: &str,
667 opts: OpenOptions,
668 ) -> Result<Self> {
669 let inner = Arc::new(open_core_with_retry(opts.lock_timeout_ms, || {
670 CoreDatabase::open_encrypted_with_credentials(path, passphrase, username, password)
671 })?);
672 let schema = load_schema(path)?;
673 ensure_internal_tables(&inner)?;
674 reap_rotated_wal_segments(&inner);
675 Ok(Self {
676 inner,
677 schema,
678 root: path.to_path_buf(),
679 default_providers: HashMap::new(),
680 session: parking_lot::RwLock::new(None),
681 sequence_lock: parking_lot::Mutex::new(()),
682 })
683 }
684
685 pub fn create_encrypted_with_credentials(
689 path: &Path,
690 schema: KitSchema,
691 passphrase: &str,
692 admin_username: &str,
693 admin_password: &str,
694 ) -> Result<Self> {
695 std::fs::create_dir_all(path)?;
696 let inner = Arc::new(CoreDatabase::create_encrypted_with_credentials(
697 path,
698 passphrase,
699 admin_username,
700 admin_password,
701 )?);
702 ensure_internal_tables(&inner)?;
703 store_schema(path, &schema)?;
704 for table in &schema.tables {
705 create_core_table(&inner, &table.name, to_core_schema(table)?)?;
706 }
707 Ok(Self {
708 inner,
709 schema,
710 root: path.to_path_buf(),
711 default_providers: HashMap::new(),
712 session: parking_lot::RwLock::new(None),
713 sequence_lock: parking_lot::Mutex::new(()),
714 })
715 }
716
717 pub fn enable_auth(&self, admin_username: &str, admin_password: &str) -> Result<()> {
721 self.inner
722 .enable_auth(admin_username, admin_password)
723 .map_err(KitError::from)
724 }
725
726 pub fn disable_auth(&self) -> Result<()> {
730 self.inner.disable_auth().map_err(KitError::from)
731 }
732
733 pub fn require_auth_enabled(&self) -> bool {
735 self.inner.require_auth_enabled()
736 }
737
738 pub fn refresh_principal(&self) -> Result<()> {
742 self.inner.refresh_principal().map_err(KitError::from)?;
743 *self.session.write() = None;
746 Ok(())
747 }
748
749 pub fn register_default(
752 &mut self,
753 name: impl Into<String>,
754 provider: impl Fn() -> Value + Send + Sync + 'static,
755 ) {
756 self.default_providers
757 .insert(name.into(), Box::new(provider));
758 }
759
760 pub fn embedding_providers(&self) -> &mongreldb_core::EmbeddingProviderRegistry {
767 self.inner.embedding_providers()
768 }
769
770 pub fn register_embedding_provider(
772 &self,
773 provider: Arc<dyn mongreldb_core::EmbeddingProvider>,
774 ) {
775 let _ = self.embedding_providers().register_new(provider);
776 }
777
778 pub fn embed_texts(
786 &self,
787 source: &mongreldb_kit_core::schema::EmbeddingSource,
788 texts: &[&str],
789 expected_dim: u32,
790 ) -> Result<Vec<Vec<f32>>> {
791 let core_source = crate::schema::to_core_embedding_source(source);
792 self.inner
793 .embedding_providers()
794 .embed(&core_source, texts, expected_dim)
795 .map_err(|e| KitError::Validation(e.to_string()))
796 }
797
798 pub fn raw(&self) -> &CoreDatabase {
802 &self.inner
803 }
804
805 pub fn start_create_index(&self, table: &str, index: &KitIndex) -> Result<u64> {
810 let table_schema = self
811 .schema
812 .table(table)
813 .ok_or_else(|| KitError::Validation(format!("table {table:?} not found")))?;
814 let definition = one_core_index(table_schema, index)?;
815 self.inner
816 .start_create_index(table, definition)
817 .map_err(KitError::from)
818 }
819
820 pub fn start_replace_index(
825 &self,
826 table: &str,
827 expected_old_name: &str,
828 replacement: &KitIndex,
829 ) -> Result<u64> {
830 let table_schema = self
831 .schema
832 .table(table)
833 .ok_or_else(|| KitError::Validation(format!("table {table:?} not found")))?;
834 let old = table_schema
835 .indexes
836 .iter()
837 .find(|index| index.name == expected_old_name)
838 .ok_or_else(|| {
839 KitError::Validation(format!(
840 "index {expected_old_name:?} not found on table {table:?}"
841 ))
842 })?;
843 let old_definition = one_core_index(table_schema, old)?;
844 let definition = one_core_index(table_schema, replacement)?;
845 self.inner
846 .start_replace_index(table, &old_definition.name, definition)
847 .map_err(KitError::from)
848 }
849
850 pub fn resume_index_build(&self, job_id: u64) -> Result<()> {
852 self.inner
853 .resume_index_build(job_id)
854 .map_err(KitError::from)
855 }
856
857 pub fn cancel_index_build(&self, job_id: u64) -> Result<()> {
859 self.inner
860 .job_registry()
861 .cancel(job_id)
862 .map_err(mongreldb_core::MongrelError::from)
863 .map_err(KitError::from)
864 }
865
866 pub fn index_build(&self, job_id: u64) -> Result<mongreldb_core::JobRecord> {
868 let record = self
869 .inner
870 .job_registry()
871 .get(job_id)
872 .ok_or_else(|| KitError::Validation(format!("job {job_id} not found")))?;
873 if record.kind != mongreldb_core::JobKind::IndexBuild {
874 return Err(KitError::Validation(format!(
875 "job {job_id} is not an index build"
876 )));
877 }
878 Ok(record)
879 }
880
881 pub fn wait_index_build(
883 &self,
884 job_id: u64,
885 timeout: std::time::Duration,
886 ) -> Result<mongreldb_core::JobRecord> {
887 self.inner
888 .job_registry()
889 .wait_terminal(job_id, timeout)
890 .map_err(mongreldb_core::MongrelError::from)
891 .map_err(KitError::from)
892 }
893
894 pub fn table_names(&self) -> Vec<String> {
896 self.schema
897 .tables
898 .iter()
899 .map(|t| t.name.clone())
900 .filter(|n| !n.starts_with("__kit_"))
901 .collect()
902 }
903
904 pub fn create_procedure(
905 &self,
906 spec: &ProcedureSpec,
907 ) -> Result<mongreldb_core::StoredProcedure> {
908 let procedure = core_procedure(spec)?;
909 self.inner
910 .create_procedure(procedure)
911 .map_err(KitError::from)
912 }
913
914 pub fn replace_procedure(
915 &self,
916 spec: &ProcedureSpec,
917 ) -> Result<mongreldb_core::StoredProcedure> {
918 let procedure = core_procedure(spec)?;
919 self.inner
920 .create_or_replace_procedure(procedure)
921 .map_err(KitError::from)
922 }
923
924 pub fn drop_procedure(&self, name: &str) -> Result<()> {
925 self.inner.drop_procedure(name).map_err(KitError::from)
926 }
927
928 pub fn call_procedure(
929 &self,
930 name: &str,
931 args: serde_json::Map<String, Value>,
932 ) -> Result<mongreldb_core::ProcedureCallResult> {
933 let args = args
934 .iter()
935 .map(|(key, value)| Ok((key.clone(), json_to_core_value(value)?)))
936 .collect::<Result<HashMap<_, _>>>()?;
937 self.inner
938 .call_procedure(name, args)
939 .map_err(KitError::from)
940 }
941
942 pub fn create_trigger(&self, spec: &TriggerSpec) -> Result<mongreldb_core::StoredTrigger> {
943 let trigger = core_trigger(spec)?;
944 self.inner.create_trigger(trigger).map_err(KitError::from)
945 }
946
947 pub fn replace_trigger(&self, spec: &TriggerSpec) -> Result<mongreldb_core::StoredTrigger> {
948 let trigger = core_trigger(spec)?;
949 self.inner
950 .create_or_replace_trigger(trigger)
951 .map_err(KitError::from)
952 }
953
954 pub fn drop_trigger(&self, name: &str) -> Result<()> {
955 self.inner.drop_trigger(name).map_err(KitError::from)
956 }
957
958 pub fn triggers(&self) -> Vec<mongreldb_core::StoredTrigger> {
959 self.inner.triggers()
960 }
961
962 pub fn trigger(&self, name: &str) -> Option<mongreldb_core::StoredTrigger> {
963 self.inner.trigger(name)
964 }
965
966 pub fn allocate_sequence(&self, name: &str, count: i64) -> Result<i64> {
972 use crate::internal::cols;
973 let _guard = self.sequence_lock.lock();
974 let mut attempt = 0;
975 loop {
976 let mut txn = self.inner.begin();
977 let snapshot = txn.read_snapshot();
978 let existing = self
979 .visible_core_rows_at(crate::internal::SEQUENCES, snapshot)?
980 .into_iter()
981 .find(|r| internal_bytes(r, cols::SEQ_NAME) == Some(name.to_string()));
982
983 let now = crate::internal::iso_now();
984 let (start, next, old_row_id) = match &existing {
988 Some(row) => {
989 let current = match row.columns.get(&cols::SEQ_NEXT) {
990 Some(CoreValue::Int64(i)) => *i,
991 _ => 1,
992 };
993 (current, current + count, Some(row.row_id))
994 }
995 None => (1, 1 + count, None),
996 };
997
998 if let Some(rid) = old_row_id {
999 txn.delete(crate::internal::SEQUENCES, rid)
1000 .map_err(KitError::from)?;
1001 }
1002 txn.put(
1003 crate::internal::SEQUENCES,
1004 vec![
1005 (cols::SEQ_NAME, CoreValue::Bytes(name.as_bytes().to_vec())),
1006 (cols::SEQ_NEXT, CoreValue::Int64(next)),
1007 (cols::SEQ_UPDATED, CoreValue::Bytes(now.into_bytes())),
1008 ],
1009 )
1010 .map_err(KitError::from)?;
1011 match txn.commit() {
1012 Ok(_) => return Ok(start),
1013 Err(mongreldb_core::MongrelError::Conflict(_)) if attempt < 10_000 => {
1014 attempt += 1;
1015 std::thread::yield_now();
1016 continue;
1017 }
1018 Err(e) => return Err(KitError::from(e)),
1019 }
1020 }
1021 }
1022
1023 pub fn transaction<T, F>(&self, max_retries: usize, mut f: F) -> Result<T>
1026 where
1027 F: FnMut(&mut crate::txn::Transaction<'_>) -> Result<T>,
1028 {
1029 let mut attempt = 0;
1030 loop {
1031 let mut txn = self.begin()?;
1032 match f(&mut txn) {
1033 Ok(value) => match txn.commit() {
1034 Ok(()) => return Ok(value),
1035 Err(KitError::Conflict(_)) if attempt < max_retries => {
1036 attempt += 1;
1037 continue;
1038 }
1039 Err(e) => return Err(e),
1040 },
1041 Err(KitError::Conflict(_)) if attempt < max_retries => {
1042 txn.rollback();
1043 attempt += 1;
1044 continue;
1045 }
1046 Err(e) => {
1047 txn.rollback();
1048 return Err(e);
1049 }
1050 }
1051 }
1052 }
1053
1054 pub fn table(&self, name: &str) -> Option<&KitTable> {
1056 self.schema.table(name)
1057 }
1058
1059 pub fn schema(&self) -> &KitSchema {
1061 &self.schema
1062 }
1063
1064 pub fn begin(&self) -> Result<crate::txn::Transaction<'_>> {
1066 let core_txn = self.inner.begin();
1067 Ok(crate::txn::Transaction::new(self, core_txn))
1068 }
1069
1070 pub fn set_schema(&mut self, schema: KitSchema) {
1072 self.schema = schema;
1073 }
1074
1075 pub fn check_internal_tables(&self) -> Result<()> {
1078 let schema_file = self.root.join(SCHEMA_FILE);
1079 if !schema_file.exists() {
1080 return Err(KitError::Integrity(format!(
1081 "schema file {} is missing",
1082 schema_file.display()
1083 )));
1084 }
1085 for (name, _) in internal_tables_core() {
1086 if self.inner.table_id(name).is_err() {
1087 return Err(KitError::Integrity(format!(
1088 "internal table {name} is missing"
1089 )));
1090 }
1091 }
1092 Ok(())
1093 }
1094
1095 pub fn gc(&self) -> Result<usize> {
1098 self.inner.gc().map_err(KitError::from)
1099 }
1100
1101 pub fn check(&self) -> Vec<serde_json::Value> {
1104 self.inner
1105 .check()
1106 .into_iter()
1107 .map(|i| {
1108 serde_json::json!({
1109 "table_id": i.table_id,
1110 "table_name": i.table_name,
1111 "severity": i.severity,
1112 "description": i.description,
1113 })
1114 })
1115 .collect()
1116 }
1117
1118 pub fn doctor(&self) -> Result<Vec<u64>> {
1120 self.inner.doctor().map_err(KitError::from)
1121 }
1122
1123 pub fn snapshot_epoch(&self) -> u64 {
1127 self.inner.snapshot().0.epoch.0
1128 }
1129
1130 pub fn set_history_retention_epochs(&self, epochs: u64) -> Result<()> {
1131 self.inner
1132 .set_history_retention_epochs(epochs)
1133 .map_err(KitError::from)
1134 }
1135
1136 pub fn history_retention_epochs(&self) -> u64 {
1137 self.inner.history_retention_epochs()
1138 }
1139
1140 pub fn earliest_retained_epoch(&self) -> u64 {
1141 self.inner.earliest_retained_epoch().0
1142 }
1143
1144 pub fn export_tsv(&self, table: &str) -> Result<String> {
1148 let t = self
1149 .schema
1150 .tables
1151 .iter()
1152 .find(|t| t.name == table)
1153 .ok_or_else(|| KitError::Validation(format!("unknown table '{table}'")))?
1154 .clone();
1155 let tx = self.begin()?;
1156 let rows = tx.all_rows(table)?;
1157 Ok(crate::tsv::rows_to_tsv(&t, &rows))
1158 }
1159
1160 pub fn import_tsv(&self, table: &str, text: &str) -> Result<usize> {
1164 let t = self
1165 .schema
1166 .tables
1167 .iter()
1168 .find(|t| t.name == table)
1169 .ok_or_else(|| KitError::Validation(format!("unknown table '{table}'")))?
1170 .clone();
1171 let rows = crate::tsv::tsv_to_rows(&t, text)?;
1172 let n = rows.len();
1173 self.transaction(1, |tx| {
1174 tx.insert_many(table, rows.clone())?;
1175 Ok(())
1176 })?;
1177 Ok(n)
1178 }
1179
1180 pub fn explain(
1185 &self,
1186 table: &str,
1187 predicate: &mongreldb_kit_core::query::Expr,
1188 ) -> Result<ExplainPlan> {
1189 let t = self
1190 .schema
1191 .tables
1192 .iter()
1193 .find(|t| t.name == table)
1194 .ok_or_else(|| KitError::Validation(format!("unknown table '{table}'")))?;
1195 Ok(match crate::pushdown::translate_predicate(t, predicate) {
1196 Some(p) => ExplainPlan {
1197 index_accelerated: p.can_push(),
1198 exact: p.fully_translated,
1199 pushed_conditions: p.conditions.iter().map(condition_label).collect(),
1200 },
1201 None => ExplainPlan {
1202 index_accelerated: false,
1203 exact: false,
1204 pushed_conditions: Vec::new(),
1205 },
1206 })
1207 }
1208
1209 pub fn rows_at_epoch(&self, table: &str, epoch: u64) -> Result<Vec<crate::schema::Row>> {
1215 let t = self
1216 .schema
1217 .tables
1218 .iter()
1219 .find(|t| t.name == table)
1220 .ok_or_else(|| KitError::Validation(format!("unknown table '{table}'")))?;
1221 let current = self.snapshot_epoch();
1222 if epoch > current {
1223 return Err(KitError::Validation(format!(
1224 "epoch {epoch} is in the future (current committed epoch is {current})"
1225 )));
1226 }
1227 let snap = Snapshot::at(mongreldb_core::epoch::Epoch(epoch));
1228 let rows = self.visible_core_rows_at(table, snap)?;
1229 rows.iter()
1230 .map(|r| crate::schema::core_row_to_json(r, t))
1231 .collect()
1232 }
1233
1234 pub fn approx_aggregate(
1240 &self,
1241 table: &str,
1242 column: Option<&str>,
1243 agg: ApproxAggKind,
1244 z: f64,
1245 ) -> Result<Option<ApproxAggregate>> {
1246 let t = self
1247 .schema
1248 .tables
1249 .iter()
1250 .find(|t| t.name == table)
1251 .ok_or_else(|| KitError::Validation(format!("unknown table '{table}'")))?;
1252 if matches!(agg, ApproxAggKind::Sum | ApproxAggKind::Avg) && column.is_none() {
1253 return Err(KitError::Validation(
1254 "approx sum/avg requires a column".into(),
1255 ));
1256 }
1257 let cid = match column {
1258 Some(name) => Some(
1259 t.columns
1260 .iter()
1261 .find(|c| c.name == name)
1262 .ok_or_else(|| KitError::Validation(format!("unknown column '{name}'")))?
1263 .id as u16,
1264 ),
1265 None => None,
1266 };
1267 let core_agg = match agg {
1268 ApproxAggKind::Count => ApproxAgg::Count,
1269 ApproxAggKind::Sum => ApproxAgg::Sum,
1270 ApproxAggKind::Avg => ApproxAgg::Avg,
1271 };
1272 let handle = self.inner.table(table).map_err(KitError::from)?;
1273 let mut guard = handle.lock();
1274 let res = guard
1275 .approx_aggregate(&[], cid, core_agg, z)
1276 .map_err(KitError::from)?;
1277 Ok(res.map(|r| ApproxAggregate {
1278 point: r.point,
1279 ci_low: r.ci_low,
1280 ci_high: r.ci_high,
1281 n_population: r.n_population,
1282 n_sample_live: r.n_sample_live,
1283 n_passing: r.n_passing,
1284 }))
1285 }
1286
1287 pub fn scan_batched<F>(&self, table: &str, batch_size: usize, mut f: F) -> Result<()>
1293 where
1294 F: FnMut(&[serde_json::Map<String, Value>]) -> Result<()>,
1295 {
1296 let kit_t = self
1297 .schema
1298 .tables
1299 .iter()
1300 .find(|t| t.name == table)
1301 .ok_or_else(|| KitError::Validation(format!("unknown table '{table}'")))?;
1302 let batch_size = batch_size.max(1);
1303 let (snapshot, _pin) = self.inner.snapshot();
1306 let handle = self.inner.table(table).map_err(KitError::from)?;
1307 let guard = handle.lock();
1308
1309 let mut projection: Vec<(u16, mongreldb_core::schema::TypeId)> = Vec::new();
1311 let mut meta: Vec<(String, mongreldb_kit_core::schema::ColumnType)> = Vec::new();
1312 for c in &guard.schema().columns {
1313 if let Some(kc) = kit_t.columns.iter().find(|kc| kc.id as u16 == c.id) {
1314 projection.push((c.id, c.ty.clone()));
1315 meta.push((kc.name.clone(), kc.storage_type));
1316 }
1317 }
1318
1319 match guard
1320 .scan_cursor(snapshot, projection, &[])
1321 .map_err(KitError::from)?
1322 {
1323 Some(mut cursor) => {
1324 let mut buf: Vec<serde_json::Map<String, Value>> = Vec::with_capacity(batch_size);
1325 while let Some(batch) = cursor.next_batch().map_err(KitError::from)? {
1326 let nrows = batch.first().map(|c| c.len()).unwrap_or(0);
1327 for j in 0..nrows {
1328 let mut m = serde_json::Map::new();
1329 for (ci, (name, ty)) in meta.iter().enumerate() {
1330 let cv = batch
1331 .get(ci)
1332 .and_then(|col| col.value_at(j))
1333 .unwrap_or(CoreValue::Null);
1334 m.insert(name.clone(), crate::schema::core_to_json(&cv, *ty)?);
1335 }
1336 buf.push(m);
1337 if buf.len() >= batch_size {
1338 f(&buf)?;
1339 buf.clear();
1340 }
1341 }
1342 }
1343 if !buf.is_empty() {
1344 f(&buf)?;
1345 }
1346 Ok(())
1347 }
1348 None => {
1349 drop(guard);
1350 let rows = self.visible_core_rows_at(table, snapshot)?;
1351 let maps: Vec<serde_json::Map<String, Value>> = rows
1352 .iter()
1353 .map(|r| crate::schema::core_row_to_json(r, kit_t).map(|row| row.values))
1354 .collect::<Result<Vec<_>>>()?;
1355 for chunk in maps.chunks(batch_size) {
1356 f(chunk)?;
1357 }
1358 Ok(())
1359 }
1360 }
1361 }
1362
1363 pub fn set_similarity(
1372 &self,
1373 table: &str,
1374 column: &str,
1375 query: &[String],
1376 k: usize,
1377 ) -> Result<Vec<SimilarRow>> {
1378 let t = self
1379 .schema
1380 .tables
1381 .iter()
1382 .find(|t| t.name == table)
1383 .ok_or_else(|| KitError::Validation(format!("unknown table '{table}'")))?;
1384 let col = t.columns.iter().find(|c| c.name == column).ok_or_else(|| {
1385 KitError::Validation(format!("unknown column '{column}' on table '{table}'"))
1386 })?;
1387 let query_set: std::collections::HashSet<String> = query.iter().cloned().collect();
1388
1389 let has_minhash = t.indexes.iter().any(|idx| {
1390 idx.kind == KitIndexKind::MinHash && idx.columns.iter().any(|c| c == column)
1391 });
1392 let rows = if has_minhash {
1393 let query_hashes: Vec<u64> = query
1395 .iter()
1396 .map(|s| mongreldb_core::index::minhash_token_hash(s))
1397 .collect();
1398 let cand_k = k.saturating_mul(8).max(k + 64);
1400 let cond = mongreldb_core::query::Condition::MinHashSimilar {
1401 column_id: col.id as u16,
1402 query: query_hashes,
1403 k: cand_k,
1404 };
1405 let (snapshot, _pin) = self.inner.snapshot();
1406 let core_rows = self.query_core_rows_at(table, &[cond], snapshot)?;
1407 core_rows
1408 .iter()
1409 .map(|r| crate::schema::core_row_to_json(r, t))
1410 .collect::<Result<Vec<_>>>()?
1411 } else {
1412 let tx = self.begin()?;
1413 tx.all_rows(table)?
1414 };
1415
1416 let mut scored: Vec<SimilarRow> = Vec::new();
1417 for row in rows {
1418 let set = parse_string_set(row.values.get(column));
1419 let inter = set.iter().filter(|x| query_set.contains(*x)).count();
1420 let union = set.len() + query_set.len() - inter;
1421 let sim = if union == 0 {
1422 0.0
1423 } else {
1424 inter as f64 / union as f64
1425 };
1426 if sim > 0.0 {
1427 scored.push(SimilarRow {
1428 row,
1429 similarity: sim,
1430 });
1431 }
1432 }
1433 scored.sort_by(|a, b| {
1434 b.similarity
1435 .partial_cmp(&a.similarity)
1436 .unwrap_or(std::cmp::Ordering::Equal)
1437 });
1438 scored.truncate(k);
1439 Ok(scored)
1440 }
1441
1442 pub fn flush(&self) -> Result<()> {
1446 for name in self.inner.table_names() {
1447 let handle = self.inner.table(&name).map_err(KitError::from)?;
1448 let mut guard = handle.lock();
1449 guard.flush().map_err(KitError::from)?;
1450 }
1451 Ok(())
1452 }
1453
1454 pub fn incremental_aggregate(
1466 &self,
1467 table: &str,
1468 column: Option<&str>,
1469 agg: IncrementalAggKind,
1470 filter: Option<&mongreldb_kit_core::query::Expr>,
1471 ) -> Result<IncrementalAggregate> {
1472 let t = self
1473 .schema
1474 .tables
1475 .iter()
1476 .find(|t| t.name == table)
1477 .ok_or_else(|| KitError::Validation(format!("unknown table '{table}'")))?;
1478 if !matches!(agg, IncrementalAggKind::Count) && column.is_none() {
1479 return Err(KitError::Validation(
1480 "sum/min/max/avg incremental aggregate requires a column".into(),
1481 ));
1482 }
1483 let cid = match column {
1484 Some(name) => Some(
1485 t.columns
1486 .iter()
1487 .find(|c| c.name == name)
1488 .ok_or_else(|| KitError::Validation(format!("unknown column '{name}'")))?
1489 .id as u16,
1490 ),
1491 None => None,
1492 };
1493 let conditions = match filter {
1494 Some(expr) => {
1495 let plan = crate::pushdown::translate_predicate(t, expr).ok_or_else(|| {
1496 KitError::Validation(
1497 "filter is not index-translatable for an incremental aggregate".into(),
1498 )
1499 })?;
1500 if !plan.fully_translated {
1501 return Err(KitError::Validation(
1502 "filter has a residual that an incremental aggregate cannot apply exactly"
1503 .into(),
1504 ));
1505 }
1506 plan.conditions
1507 }
1508 None => Vec::new(),
1509 };
1510 let core_agg = match agg {
1511 IncrementalAggKind::Count => NativeAgg::Count,
1512 IncrementalAggKind::Sum => NativeAgg::Sum,
1513 IncrementalAggKind::Min => NativeAgg::Min,
1514 IncrementalAggKind::Max => NativeAgg::Max,
1515 IncrementalAggKind::Avg => NativeAgg::Avg,
1516 };
1517 let cache_key = incremental_cache_key(t.id, cid, agg, &conditions);
1518 let handle = self.inner.table(table).map_err(KitError::from)?;
1519 let mut guard = handle.lock();
1520 let res = guard
1521 .aggregate_incremental(cache_key, &conditions, cid, core_agg)
1522 .map_err(KitError::from)?;
1523 Ok(IncrementalAggregate {
1524 value: agg_state_value(&res.state),
1525 incremental: res.incremental,
1526 delta_rows: res.delta_rows,
1527 })
1528 }
1529
1530 pub fn applied_migrations(&self) -> Result<Vec<mongreldb_kit_core::migrations::Migration>> {
1532 crate::migrate::load_applied_migrations(&self.inner)
1533 }
1534
1535 pub(crate) fn core_db(&self) -> &CoreDatabase {
1536 &self.inner
1537 }
1538
1539 pub(crate) fn core_arc(&self) -> Arc<CoreDatabase> {
1542 Arc::clone(&self.inner)
1543 }
1544
1545 pub fn close(&self) -> Result<()> {
1550 self.inner.close().map_err(KitError::from)
1551 }
1552
1553 pub fn compact_all(&self) -> Result<(usize, usize)> {
1558 self.inner.compact().map_err(KitError::from)
1559 }
1560
1561 pub fn compact_table(&self, name: &str) -> Result<bool> {
1564 self.inner.compact_table(name).map_err(KitError::from)
1565 }
1566
1567 pub fn rename_table(&mut self, from: &str, to: &str) -> Result<()> {
1577 if from.starts_with("__kit_") || to.starts_with("__kit_") {
1578 return Err(KitError::Validation(
1579 "rename_table: names beginning with '__kit_' are reserved for internal tables"
1580 .into(),
1581 ));
1582 }
1583 self.inner.rename_table(from, to).map_err(KitError::from)?;
1584 if !self.schema.rename_table(from, to) {
1587 return Err(KitError::Integrity(format!(
1590 "rename_table: kit schema has no table '{from}' (or '{to}' already exists)"
1591 )));
1592 }
1593 for table in &mut self.schema.tables {
1594 for fk in &mut table.foreign_keys {
1595 if fk.references_table == from {
1596 fk.references_table = to.to_string();
1597 }
1598 }
1599 }
1600 store_schema(&self.root, &self.schema)?;
1601 Ok(())
1602 }
1603
1604 pub fn analyze(&self) -> Result<()> {
1609 for name in self.inner.table_names() {
1610 let handle = self.inner.table(&name).map_err(KitError::from)?;
1611 handle.lock().ensure_indexes_complete()?;
1612 }
1613 Ok(())
1614 }
1615
1616 pub fn vacuum(&self) -> Result<usize> {
1620 self.inner.compact().map_err(KitError::from)?;
1621 self.inner.gc().map_err(KitError::from)
1622 }
1623
1624 pub fn create_view(&self, spec: &ViewSpec) -> Result<()> {
1630 self.sql(&spec.create_sql())?;
1631 Ok(())
1632 }
1633
1634 pub fn drop_view(&self, name: &str) -> Result<()> {
1636 self.sql(&format!("DROP VIEW IF EXISTS {name}"))?;
1637 Ok(())
1638 }
1639
1640 pub fn reserve_auto_inc(&self, table: &str) -> Result<Option<i64>> {
1648 let handle = self.inner.table(table).map_err(KitError::from)?;
1649 let mut guard = handle.lock();
1650 guard.reserve_auto_inc().map_err(KitError::from)
1651 }
1652
1653 pub fn create_user(&self, username: &str, password: &str) -> Result<()> {
1657 self.inner
1658 .create_user(username, password)
1659 .map_err(KitError::from)?;
1660 Ok(())
1661 }
1662
1663 pub fn drop_user(&self, username: &str) -> Result<()> {
1665 self.inner.drop_user(username).map_err(KitError::from)
1666 }
1667
1668 pub fn alter_user_password(&self, username: &str, new_password: &str) -> Result<()> {
1670 self.inner
1671 .alter_user_password(username, new_password)
1672 .map_err(KitError::from)
1673 }
1674
1675 pub fn verify_user(
1677 &self,
1678 username: &str,
1679 password: &str,
1680 ) -> Result<Option<mongreldb_core::auth::UserEntry>> {
1681 self.inner
1682 .verify_user(username, password)
1683 .map_err(KitError::from)
1684 }
1685
1686 pub fn set_user_admin(&self, username: &str, is_admin: bool) -> Result<()> {
1688 self.inner
1689 .set_user_admin(username, is_admin)
1690 .map_err(KitError::from)
1691 }
1692
1693 pub fn users(&self) -> Vec<String> {
1695 self.inner.users().into_iter().map(|u| u.username).collect()
1696 }
1697
1698 pub fn create_role(&self, name: &str) -> Result<()> {
1700 self.inner.create_role(name).map_err(KitError::from)?;
1701 Ok(())
1702 }
1703
1704 pub fn drop_role(&self, name: &str) -> Result<()> {
1706 self.inner.drop_role(name).map_err(KitError::from)
1707 }
1708
1709 pub fn roles(&self) -> Vec<String> {
1711 self.inner.roles().into_iter().map(|r| r.name).collect()
1712 }
1713
1714 pub fn grant_role(&self, username: &str, role_name: &str) -> Result<()> {
1716 self.inner
1717 .grant_role(username, role_name)
1718 .map_err(KitError::from)
1719 }
1720
1721 pub fn revoke_role(&self, username: &str, role_name: &str) -> Result<()> {
1723 self.inner
1724 .revoke_role(username, role_name)
1725 .map_err(KitError::from)
1726 }
1727
1728 pub fn grant_permission(
1730 &self,
1731 role_name: &str,
1732 permission: mongreldb_core::auth::Permission,
1733 ) -> Result<()> {
1734 self.inner
1735 .grant_permission(role_name, permission)
1736 .map_err(KitError::from)
1737 }
1738
1739 pub fn revoke_permission(
1741 &self,
1742 role_name: &str,
1743 permission: mongreldb_core::auth::Permission,
1744 ) -> Result<()> {
1745 self.inner
1746 .revoke_permission(role_name, permission)
1747 .map_err(KitError::from)
1748 }
1749
1750 pub fn set_spill_threshold(&self, bytes: u64) {
1756 self.inner.set_spill_threshold(bytes);
1757 }
1758
1759 pub fn set_recursive_triggers(&self, enabled: bool) {
1761 self.inner.set_recursive_triggers(enabled);
1762 }
1763
1764 pub fn trigger_config(&self) -> mongreldb_core::TriggerConfig {
1766 self.inner.trigger_config()
1767 }
1768
1769 pub fn set_trigger_config(&self, config: mongreldb_core::TriggerConfig) -> Result<()> {
1771 self.inner
1772 .set_trigger_config(config)
1773 .map_err(KitError::from)
1774 }
1775
1776 pub fn set_table_compaction_zstd_level(&self, table: &str, level: i32) -> Result<()> {
1778 let handle = self.inner.table(table).map_err(KitError::from)?;
1779 handle.lock().set_compaction_zstd_level(level);
1780 Ok(())
1781 }
1782
1783 pub fn set_table_result_cache_max_bytes(&self, table: &str, max_bytes: u64) -> Result<()> {
1785 let handle = self.inner.table(table).map_err(KitError::from)?;
1786 handle.lock().set_result_cache_max_bytes(max_bytes);
1787 Ok(())
1788 }
1789
1790 pub fn set_table_mutable_run_spill_bytes(&self, table: &str, bytes: u64) -> Result<()> {
1792 let handle = self.inner.table(table).map_err(KitError::from)?;
1793 handle.lock().set_mutable_run_spill_bytes(bytes);
1794 Ok(())
1795 }
1796
1797 pub fn set_table_sync_byte_threshold(&self, table: &str, threshold: u64) -> Result<()> {
1799 let handle = self.inner.table(table).map_err(KitError::from)?;
1800 handle.lock().set_sync_byte_threshold(threshold);
1801 Ok(())
1802 }
1803
1804 pub fn set_table_index_build_policy(
1807 &self,
1808 table: &str,
1809 policy: mongreldb_core::IndexBuildPolicy,
1810 ) -> Result<()> {
1811 let handle = self.inner.table(table).map_err(KitError::from)?;
1812 handle.lock().set_index_build_policy(policy);
1813 Ok(())
1814 }
1815
1816 pub fn table_page_cache_stats(&self, table: &str) -> Result<mongreldb_core::cache::CacheStats> {
1818 let handle = self.inner.table(table).map_err(KitError::from)?;
1819 let stats = handle.lock().page_cache_stats();
1820 Ok(stats)
1821 }
1822
1823 pub fn table_run_count(&self, table: &str) -> Result<usize> {
1825 let handle = self.inner.table(table).map_err(KitError::from)?;
1826 let n = handle.lock().run_count();
1827 Ok(n)
1828 }
1829
1830 pub fn table_memtable_len(&self, table: &str) -> Result<usize> {
1832 let handle = self.inner.table(table).map_err(KitError::from)?;
1833 let n = handle.lock().memtable_len();
1834 Ok(n)
1835 }
1836
1837 pub fn table_mutable_run_len(&self, table: &str) -> Result<usize> {
1839 let handle = self.inner.table(table).map_err(KitError::from)?;
1840 let n = handle.lock().mutable_run_len();
1841 Ok(n)
1842 }
1843
1844 pub fn table_page_cache_len(&self, table: &str) -> Result<usize> {
1846 let handle = self.inner.table(table).map_err(KitError::from)?;
1847 let n = handle.lock().page_cache_len();
1848 Ok(n)
1849 }
1850
1851 pub fn table_decoded_cache_len(&self, table: &str) -> Result<usize> {
1853 let handle = self.inner.table(table).map_err(KitError::from)?;
1854 let n = handle.lock().decoded_cache_len();
1855 Ok(n)
1856 }
1857
1858 pub fn sql(&self, statement: &str) -> Result<Vec<arrow::record_batch::RecordBatch>> {
1876 self.sql_with_options(statement, SqlOptions::default())
1877 }
1878
1879 fn sql_session(&self) -> Result<Arc<mongreldb_query::MongrelSession>> {
1880 if let Some(session) = self.session.read().as_ref() {
1881 return Ok(Arc::clone(session));
1882 }
1883 let session = Arc::new(
1884 mongreldb_query::MongrelSession::open(self.core_arc()).map_err(KitError::from)?,
1885 );
1886 let mut cached = self.session.write();
1887 Ok(Arc::clone(cached.get_or_insert(session)))
1888 }
1889
1890 #[doc(hidden)]
1891 pub fn set_sql_test_hook(&self, hook: Option<mongreldb_query::SqlTestHook>) -> Result<()> {
1892 self.sql_session()?.set_test_hook(hook);
1893 Ok(())
1894 }
1895
1896 pub fn start_sql(
1897 &self,
1898 statement: impl Into<String>,
1899 options: SqlOptions,
1900 ) -> Result<SqlQueryHandle> {
1901 let session = self.sql_session()?;
1902 let query = session
1903 .register_query(mongreldb_query::SqlQueryOptions {
1904 query_id: options.query_id,
1905 timeout: options.timeout,
1906 ..mongreldb_query::SqlQueryOptions::default()
1907 })
1908 .map_err(KitError::from)?;
1909 let query_id = query.id();
1910 let registration = mongreldb_query::RegisteredQueryGuard::new(query);
1911 let worker_session = Arc::clone(&session);
1912 let statement = statement.into();
1913 let worker = std::thread::Builder::new()
1914 .name(format!("mongreldb-kit-sql-{query_id}"))
1915 .spawn(move || {
1916 sql_runtime().block_on(
1917 worker_session
1918 .run_with_query_for_serialization(&statement, registration.into_query()),
1919 )
1920 })
1921 .map_err(|error| KitError::Storage(error.to_string()))?;
1922 Ok(SqlQueryHandle {
1923 query_id,
1924 session,
1925 worker: Some(worker),
1926 })
1927 }
1928
1929 pub fn sql_with_options(
1930 &self,
1931 statement: &str,
1932 options: SqlOptions,
1933 ) -> Result<Vec<arrow::record_batch::RecordBatch>> {
1934 self.start_sql(statement, options)?.wait()
1935 }
1936
1937 pub fn cancel_sql(&self, query_id: mongreldb_query::QueryId) -> mongreldb_query::CancelOutcome {
1938 self.session
1939 .read()
1940 .as_ref()
1941 .map_or(mongreldb_query::CancelOutcome::NotFound, |session| {
1942 session.cancel_query(query_id)
1943 })
1944 }
1945
1946 pub fn sql_query_status(
1947 &self,
1948 query_id: mongreldb_query::QueryId,
1949 ) -> Result<Option<mongreldb_query::QueryStatus>> {
1950 Ok(self.sql_session()?.query_registry().status(query_id))
1951 }
1952
1953 pub fn refresh_sql_session(&self) -> Result<()> {
1958 let session =
1959 mongreldb_query::MongrelSession::open(self.core_arc()).map_err(KitError::from)?;
1960 *self.session.write() = Some(Arc::new(session));
1961 Ok(())
1962 }
1963
1964 pub fn sql_arrow(&self, statement: &str) -> Result<Vec<u8>> {
1970 self.sql_arrow_with_options(statement, SqlOptions::default())
1971 }
1972
1973 pub fn sql_arrow_with_options(&self, statement: &str, options: SqlOptions) -> Result<Vec<u8>> {
1974 self.sql_serialized_with_options(statement, options, |output| {
1975 crate::arrow_util::batches_to_ipc_controlled(output.batches(), output.query())
1976 })
1977 }
1978
1979 pub fn sql_rows(&self, statement: &str) -> Result<Vec<serde_json::Map<String, Value>>> {
1983 self.sql_rows_with_options(statement, SqlOptions::default())
1984 }
1985
1986 pub fn sql_rows_with_options(
1987 &self,
1988 statement: &str,
1989 options: SqlOptions,
1990 ) -> Result<Vec<serde_json::Map<String, Value>>> {
1991 self.sql_serialized_with_options(statement, options, |output| {
1992 crate::arrow_util::batches_to_rows_controlled(output.batches(), output.query())
1993 })
1994 }
1995
1996 fn sql_serialized_with_options<T>(
1997 &self,
1998 statement: &str,
1999 options: SqlOptions,
2000 serialize: impl FnOnce(&mongreldb_query::ManagedQueryBatches) -> Result<T>,
2001 ) -> Result<T> {
2002 let session = self.sql_session()?;
2003 let query = session
2004 .register_query(mongreldb_query::SqlQueryOptions {
2005 query_id: options.query_id,
2006 timeout: options.timeout,
2007 ..mongreldb_query::SqlQueryOptions::default()
2008 })
2009 .map_err(KitError::from)?;
2010 let query_id = query.id();
2011 let output = sql_runtime()
2012 .block_on(session.run_with_query_for_serialization(statement, query))
2013 .map_err(|error| {
2014 let status = session.query_registry().status(query_id);
2015 crate::error::query_error_with_status(error, status.as_ref())
2016 })?;
2017 session.fire_test_hook(mongreldb_query::SqlTestHookPoint::BeforeSerializationBatch);
2018 match serialize(&output) {
2019 Ok(value) => {
2020 complete_sql_output(output)?;
2021 Ok(value)
2022 }
2023 Err(error) => {
2024 fail_sql_output(output, &error);
2025 Err(error)
2026 }
2027 }
2028 }
2029
2030 pub(crate) fn lookup_row_id(&self, table: &str, key: &[u8]) -> Result<Option<RowId>> {
2034 let handle = self.inner.table(table).map_err(KitError::from)?;
2035 let mut guard = handle.lock();
2036 guard.ensure_indexes_complete()?;
2037 Ok(guard.lookup_pk(key))
2038 }
2039
2040 pub(crate) fn root(&self) -> &Path {
2041 &self.root
2042 }
2043
2044 pub(crate) fn visible_core_rows_at(
2048 &self,
2049 table_name: &str,
2050 snapshot: Snapshot,
2051 ) -> Result<Vec<CoreRow>> {
2052 let handle = self.inner.table(table_name).map_err(KitError::from)?;
2053 let guard = handle.lock();
2054 guard.visible_rows(snapshot).map_err(KitError::from)
2055 }
2056
2057 pub(crate) fn query_core_rows_at(
2064 &self,
2065 table_name: &str,
2066 conditions: &[mongreldb_core::query::Condition],
2067 snapshot: Snapshot,
2068 ) -> Result<Vec<CoreRow>> {
2069 if conditions.is_empty() {
2070 return self.visible_core_rows_at(table_name, snapshot);
2071 }
2072 let handle = self.inner.table(table_name).map_err(KitError::from)?;
2073 let mut guard = handle.lock();
2074 let q = conditions
2075 .iter()
2076 .cloned()
2077 .fold(mongreldb_core::query::Query::new(), |query, condition| {
2078 query.and(condition)
2079 });
2080 guard.query(&q).map_err(KitError::from)
2081 }
2082
2083 pub(crate) fn flush_table(&self, table_name: &str) -> Result<()> {
2091 let handle = self.inner.table(table_name).map_err(KitError::from)?;
2092 handle.lock().flush().map_err(KitError::from)?;
2093 Ok(())
2094 }
2095
2096 pub(crate) fn count_core_rows_at(
2106 &self,
2107 table_name: &str,
2108 conditions: &[mongreldb_core::query::Condition],
2109 snapshot: Snapshot,
2110 ) -> Result<Option<u64>> {
2111 let handle = self.inner.table(table_name).map_err(KitError::from)?;
2112 let mut guard = handle.lock();
2113 if guard.snapshot().epoch != snapshot.epoch {
2114 return Ok(None); }
2116 guard
2117 .count_conditions(conditions, snapshot)
2118 .map_err(KitError::from)
2119 }
2120
2121 pub(crate) fn aggregate_core_at(
2130 &self,
2131 table_name: &str,
2132 column: Option<u16>,
2133 conditions: &[mongreldb_core::query::Condition],
2134 agg: NativeAgg,
2135 snapshot: Snapshot,
2136 ) -> Result<Option<NativeAggResult>> {
2137 let handle = self.inner.table(table_name).map_err(KitError::from)?;
2138 let guard = handle.lock();
2139 if guard.snapshot().epoch != snapshot.epoch {
2140 return Ok(None); }
2142 guard
2143 .aggregate_native(snapshot, column, conditions, agg)
2144 .map_err(KitError::from)
2145 }
2146
2147 pub(crate) fn count_distinct_core_at(
2156 &self,
2157 table_name: &str,
2158 column_id: u16,
2159 snapshot: Snapshot,
2160 ) -> Result<Option<u64>> {
2161 let handle = self.inner.table(table_name).map_err(KitError::from)?;
2162 let mut guard = handle.lock();
2163 if guard.snapshot().epoch != snapshot.epoch {
2164 return Ok(None); }
2166 guard
2167 .count_distinct_from_bitmap(column_id)
2168 .map_err(KitError::from)
2169 }
2170
2171 #[allow(dead_code)]
2173 pub(crate) fn get_core_row(&self, table_name: &str, row_id: u64) -> Result<Option<CoreRow>> {
2174 let handle = self.inner.table(table_name).map_err(KitError::from)?;
2175 let guard = handle.lock();
2176 let snapshot = guard.snapshot();
2177 Ok(guard.get(mongreldb_core::RowId(row_id), snapshot))
2178 }
2179}
2180
2181fn one_core_index(table: &KitTable, index: &KitIndex) -> Result<mongreldb_core::IndexDef> {
2182 let mut definitions = to_core_indexes(table, index)?;
2183 if definitions.len() != 1 {
2184 return Err(KitError::Validation(format!(
2185 "online index jobs require exactly one column; index {:?} has {}",
2186 index.name,
2187 definitions.len()
2188 )));
2189 }
2190 Ok(definitions.remove(0))
2191}
2192
2193pub(crate) fn create_core_table(db: &CoreDatabase, name: &str, schema: CoreSchema) -> Result<()> {
2194 if db.table_id(name).is_ok() {
2195 return Ok(());
2196 }
2197 db.create_table(name, schema).map_err(KitError::from)?;
2198 Ok(())
2199}
2200
2201fn sql_runtime() -> &'static tokio::runtime::Runtime {
2205 use std::sync::OnceLock;
2206 static RT: OnceLock<tokio::runtime::Runtime> = OnceLock::new();
2207 RT.get_or_init(|| {
2208 tokio::runtime::Builder::new_multi_thread()
2209 .worker_threads(4)
2210 .enable_all()
2211 .build()
2212 .expect("failed to build kit SQL tokio runtime")
2213 })
2214}
2215
2216fn core_procedure(spec: &ProcedureSpec) -> Result<mongreldb_core::StoredProcedure> {
2217 let parsed: mongreldb_core::StoredProcedure =
2218 serde_json::from_value(spec.json.clone()).map_err(KitError::from)?;
2219 mongreldb_core::StoredProcedure::new(parsed.name, parsed.mode, parsed.params, parsed.body, 0)
2220 .map_err(KitError::from)
2221}
2222
2223fn core_trigger(spec: &TriggerSpec) -> Result<mongreldb_core::StoredTrigger> {
2224 let parsed: mongreldb_core::StoredTrigger =
2225 serde_json::from_value(spec.json.clone()).map_err(KitError::from)?;
2226 mongreldb_core::StoredTrigger::new(
2227 parsed.name,
2228 mongreldb_core::TriggerDefinition {
2229 target: parsed.target,
2230 timing: parsed.timing,
2231 event: parsed.event,
2232 update_of: parsed.update_of,
2233 target_columns: parsed.target_columns,
2234 when: parsed.when,
2235 program: parsed.program,
2236 },
2237 0,
2238 )
2239 .map_err(KitError::from)
2240}
2241
2242fn json_to_core_value(value: &Value) -> Result<CoreValue> {
2243 match value {
2244 Value::Null => Ok(CoreValue::Null),
2245 Value::Bool(value) => Ok(CoreValue::Bool(*value)),
2246 Value::Number(value) => {
2247 if let Some(value) = value.as_i64() {
2248 Ok(CoreValue::Int64(value))
2249 } else if let Some(value) = value.as_f64() {
2250 Ok(CoreValue::Float64(value))
2251 } else {
2252 Err(KitError::Validation("unsupported JSON number".into()))
2253 }
2254 }
2255 Value::String(value) => Ok(CoreValue::Bytes(value.as_bytes().to_vec())),
2256 Value::Array(_) | Value::Object(_) => Err(KitError::Validation(
2257 "procedure args only support scalar JSON values".into(),
2258 )),
2259 }
2260}
2261
2262pub(crate) fn internal_bytes(row: &CoreRow, col_id: u16) -> Option<String> {
2264 match row.columns.get(&col_id) {
2265 Some(CoreValue::Bytes(b)) => String::from_utf8(b.clone()).ok(),
2266 _ => None,
2267 }
2268}
2269
2270fn reap_rotated_wal_segments(db: &CoreDatabase) {
2285 let _ = db.gc();
2286}
2287
2288pub(crate) fn load_schema(path: &Path) -> Result<KitSchema> {
2289 let file = path.join(SCHEMA_FILE);
2290 let json = std::fs::read_to_string(&file)
2291 .map_err(|e| KitError::Migration(format!("cannot read schema file: {e}")))?;
2292 let schema: KitSchema = serde_json::from_str(&json)?;
2293 Ok(schema)
2294}
2295
2296pub(crate) fn store_schema(path: &Path, schema: &KitSchema) -> Result<()> {
2297 let file = path.join(SCHEMA_FILE);
2298 let json = serde_json::to_string_pretty(schema)?;
2299 std::fs::write(&file, json)?;
2300 Ok(())
2301}
2302
2303pub(crate) fn persist_schema(db: &Database, schema: &KitSchema) -> Result<()> {
2305 store_schema(&db.root, schema)
2306}
2307
2308#[cfg(test)]
2309mod tests {
2310 use super::open_core_with_retry;
2311
2312 fn lock_error() -> mongreldb_core::MongrelError {
2313 mongreldb_core::MongrelError::DatabaseLocked {
2314 path: "/tmp/db".into(),
2315 message: "would block".into(),
2316 }
2317 }
2318
2319 #[test]
2320 fn open_retry_waits_for_lock_contention_only() {
2321 let mut calls = 0;
2322 let value = open_core_with_retry(50, || {
2323 calls += 1;
2324 if calls < 3 {
2325 Err(lock_error())
2326 } else {
2327 Ok(7)
2328 }
2329 })
2330 .unwrap();
2331 assert_eq!(value, 7);
2332 assert_eq!(calls, 3);
2333
2334 let mut non_lock_calls = 0;
2335 let err: mongreldb_core::Result<()> = open_core_with_retry(50, || {
2336 non_lock_calls += 1;
2337 Err(mongreldb_core::MongrelError::Other("nope".into()))
2338 });
2339 let err = err.unwrap_err();
2340 assert_eq!(non_lock_calls, 1);
2341 assert!(matches!(err, mongreldb_core::MongrelError::Other(_)));
2342 }
2343}