Skip to main content

surrealdb_engine_api/
lib.rs

1//! The service-provider interface between the SurrealDB Rust SDK and the
2//! engines it drives.
3//!
4//! An engine owns a datastore or a connection to one, and serves
5//! [`SurrealEngine`] to the SDK. Some run as a task the SDK hands a [`Route`]
6//! per request, answering on the route's response channel; session lifetime
7//! travels alongside on a separate [`SessionId`] channel either way. This crate
8//! holds exactly the types that cross that boundary, so an engine can live in
9//! its own crate without the SDK depending on it, or on anything it in turn
10//! depends on.
11//!
12//! # Stability
13//!
14//! This is an internal interface between crates released together. It carries
15//! no stability guarantee and may change in any release, including a patch
16//! release. Depend on it only if you implement an engine; application code
17//! should use the [`surrealdb`](https://docs.rs/surrealdb) crate.
18
19use std::borrow::Cow;
20use std::fmt::Debug;
21use std::future::Future;
22use std::path::PathBuf;
23use std::pin::Pin;
24
25use async_channel::Sender;
26pub use surrealdb_rpc::QUERY_STREAM_BUFFER;
27pub use surrealdb_rpc::export::Config as DbExportConfig;
28use surrealdb_rpc::{QueryResult, QueryStreamItem, Token, items_for_result};
29use surrealdb_types::{
30	Array, ConnectionError, Error, NotFoundError, Notification, Object, SurrealValue, Value,
31	Variables,
32};
33use uuid::Uuid;
34
35pub mod session;
36
37pub use session::{Established, SessionEntry, SessionRegistry};
38
39/// A future boxed for storage behind a trait object, as
40/// [`SurrealEngine`]'s methods require.
41///
42/// `Send` but deliberately not `Sync`: an engine may await a future from a
43/// client library that is not itself `Sync` (tonic's are not), and requiring
44/// `Sync` here would rule those engines out entirely. Nothing polls one of
45/// these from two threads at once, so `Sync` buys nothing.
46pub type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
47
48/// A request travelling from the SDK to an engine, tagged with the session it
49/// belongs to.
50#[derive(Debug)]
51pub struct RequestData {
52	/// The operation the engine should perform.
53	pub command: Command,
54	/// The session the command runs under.
55	pub session_id: Uuid,
56}
57
58/// A request paired with the channel its response must be sent on.
59#[derive(Debug)]
60pub struct Route {
61	/// The request to execute.
62	pub request: RequestData,
63	/// Where the engine sends the outcome. Exactly one message is expected.
64	pub response: Sender<Result<Vec<QueryResult>, Error>>,
65}
66
67/// A session-lifetime event. Engines keep per-session state (the authenticated
68/// session, its variables, its live queries), and these events tell them when
69/// to create, copy and discard it.
70#[derive(Debug, Clone, Copy)]
71pub enum SessionId {
72	/// A new session was opened.
73	Initial(Uuid),
74	/// A session was cloned; `new` starts as a copy of `old`.
75	Clone {
76		/// The session being copied.
77		old: Uuid,
78		/// The session receiving the copy.
79		new: Uuid,
80	},
81	/// A session was dropped and its state can be released.
82	Drop(Uuid),
83}
84
85/// Why a route could not be matched to a session.
86#[derive(Debug, Clone)]
87pub enum SessionError {
88	/// No state is registered for the session the route targets.
89	NotFound(Uuid),
90	/// The remote end reported a session failure.
91	Remote(String),
92}
93
94impl From<SessionError> for Error {
95	fn from(error: SessionError) -> Self {
96		session_error_to_error(error)
97	}
98}
99
100/// Convert a session error into the error type the SDK surfaces to callers.
101pub fn session_error_to_error(e: SessionError) -> Error {
102	match e {
103		SessionError::NotFound(id) => Error::not_found(
104			format!("Session not found: {id}"),
105			NotFoundError::Session {
106				id: Some(id.to_string()),
107			},
108		),
109		SessionError::Remote(msg) => Error::internal(msg),
110	}
111}
112
113/// Which machine learning model to export, by name and version.
114#[derive(Debug, Clone)]
115pub struct MlExportConfig {
116	/// The model name.
117	pub name: String,
118	/// The model version.
119	pub version: String,
120}
121
122/// The operations an engine can be asked to perform.
123#[derive(Debug, Clone)]
124pub enum Command {
125	/// Select the namespace and/or database for the session.
126	Use {
127		/// The namespace to use, or `None` to leave it unchanged.
128		namespace: Option<String>,
129		/// The database to use, or `None` to leave it unchanged.
130		database: Option<String>,
131	},
132	/// Sign up as a record user and authenticate the session.
133	Signup {
134		/// The signup credentials.
135		credentials: Object,
136	},
137	/// Sign in and authenticate the session.
138	Signin {
139		/// The signin credentials.
140		credentials: Object,
141	},
142	/// Authenticate the session with an existing token.
143	Authenticate {
144		/// The token to authenticate with.
145		token: Token,
146	},
147	/// Exchange a refresh token for a fresh token pair.
148	Refresh {
149		/// The token carrying the refresh token.
150		token: Token,
151	},
152	/// Drop the session's authentication.
153	Invalidate,
154	/// Begin a manual transaction.
155	Begin,
156	/// Cancel a manual transaction.
157	Rollback {
158		/// The transaction to cancel.
159		txn: Uuid,
160	},
161	/// Commit a manual transaction.
162	Commit {
163		/// The transaction to commit.
164		txn: Uuid,
165	},
166	/// Invalidate a refresh token so it can no longer be redeemed.
167	Revoke {
168		/// The token carrying the refresh token.
169		token: Token,
170	},
171	/// Execute a query, optionally inside a manual transaction.
172	Query {
173		/// The transaction to run in, or `None` for an implicit one.
174		txn: Option<Uuid>,
175		/// The query text.
176		query: Cow<'static, str>,
177		/// The variables bound for this query only.
178		variables: Variables,
179	},
180	/// Export the database to a file.
181	ExportFile {
182		/// The file to write.
183		path: PathBuf,
184		/// What to include in the export.
185		config: Option<DbExportConfig>,
186	},
187	/// Export a machine learning model to a file.
188	ExportMl {
189		/// The file to write.
190		path: PathBuf,
191		/// The model to export.
192		config: MlExportConfig,
193	},
194	/// Export the database to a channel.
195	ExportBytes {
196		/// Where the export is streamed.
197		bytes: Sender<Result<Vec<u8>, Error>>,
198		/// What to include in the export.
199		config: Option<DbExportConfig>,
200	},
201	/// Export a machine learning model to a channel.
202	ExportBytesMl {
203		/// Where the export is streamed.
204		bytes: Sender<Result<Vec<u8>, Error>>,
205		/// The model to export.
206		config: MlExportConfig,
207	},
208	/// Import a database export from a file.
209	ImportFile {
210		/// The file to read.
211		path: PathBuf,
212	},
213	/// Import a machine learning model from a file.
214	ImportMl {
215		/// The file to read.
216		path: PathBuf,
217	},
218	/// Check that the engine is reachable.
219	Health,
220	/// Report the database version.
221	Version,
222	/// Set a session variable.
223	Set {
224		/// The variable name.
225		key: String,
226		/// The value to bind, or `Value::None` to remove it.
227		value: Value,
228	},
229	/// Remove a session variable.
230	Unset {
231		/// The variable name.
232		key: String,
233	},
234	/// Register the channel a live query's notifications are delivered on.
235	SubscribeLive {
236		/// The live query.
237		uuid: Uuid,
238		/// Where notifications are delivered.
239		notification_sender: Sender<Result<Notification, Error>>,
240	},
241	/// Kill a live query.
242	Kill {
243		/// The live query.
244		uuid: Uuid,
245	},
246	/// Adopt an existing remote session.
247	Attach {
248		/// The session to adopt.
249		session_id: Uuid,
250	},
251	/// Release a remote session without ending it.
252	Detach {
253		/// The session to release.
254		session_id: Uuid,
255	},
256	/// Call a function or a machine learning model.
257	Run {
258		/// The function or model name.
259		name: String,
260		/// The model version, if any.
261		version: Option<String>,
262		/// The arguments to pass.
263		args: Array,
264	},
265}
266
267/// Which session, and which explicit transaction, a request applies to.
268///
269/// Mirrors the `RequestContext` every request carries in the SurrealDB
270/// network protocol, minus the fields the SDK does not populate: an engine
271/// applies its own configured timeouts rather than being told them per call.
272#[derive(Debug, Clone, Copy)]
273pub struct EngineContext {
274	/// The session the request runs under.
275	pub session: Uuid,
276	/// The explicit transaction to run in, or `None` for an implicit one.
277	///
278	/// Only meaningful on [`SurrealEngine::query`]; the transaction-lifecycle
279	/// methods take the transaction they act on as an explicit argument.
280	pub transaction: Option<Uuid>,
281}
282
283impl EngineContext {
284	/// A context for a request outside any explicit transaction.
285	pub fn new(session: Uuid) -> Self {
286		Self {
287			session,
288			transaction: None,
289		}
290	}
291
292	/// A context for a request inside the given explicit transaction.
293	pub fn with_transaction(session: Uuid, transaction: Option<Uuid>) -> Self {
294		Self {
295			session,
296			transaction,
297		}
298	}
299}
300
301/// A boxed engine result.
302pub type EngineFuture<'a, T> = BoxFuture<'a, Result<T, Error>>;
303
304/// Answers a streaming caller from an engine's buffered `query`.
305///
306/// The caller sees the same items either way, just all at once. This is the
307/// right answer for any engine whose transport cannot carry results before the
308/// query ends -- which is most of them.
309fn buffered_query_stream<'a, E>(
310	engine: &'a E,
311	ctx: EngineContext,
312	query: Cow<'static, str>,
313	variables: Variables,
314	items: Sender<QueryStreamItem>,
315) -> EngineFuture<'a, ()>
316where
317	E: SurrealEngine + ?Sized,
318{
319	Box::pin(async move {
320		for (index, result) in engine.query(ctx, query, variables).await?.into_iter().enumerate() {
321			for item in items_for_result(index, result) {
322				// A receiver that has gone away wants no more items, and there
323				// is nothing else to do with them.
324				if items.send(item).await.is_err() {
325					return Ok(());
326				}
327			}
328		}
329		Ok(())
330	})
331}
332
333/// The interface every SurrealDB engine implements, and the only thing the
334/// Rust SDK calls to reach a database.
335///
336/// One method per operation, each taking and returning the types that
337/// operation actually deals in, so no engine has to encode a result into a
338/// generic [`Value`] purely for the SDK to take it apart again. An embedded
339/// engine hands its own values straight back; a remote engine converts once,
340/// from its wire format.
341///
342/// Methods for capabilities an engine may not have -- live queries, export
343/// and import -- default to reporting that they are unsupported, so an engine
344/// implements only what it serves. The SDK gates most of these on
345/// `ExtraFeatures` before calling, so the default is a backstop rather than
346/// the usual path.
347///
348/// # Stability
349///
350/// This is an internal interface between crates released together. It carries
351/// no stability guarantee and may change in any release, including a patch
352/// release.
353pub trait SurrealEngine: Debug + Send + Sync + 'static {
354	// ------------------------------------------------------------------
355	// Queries
356	// ------------------------------------------------------------------
357
358	/// Executes SurrealQL, returning one result per statement, in order.
359	fn query(
360		&self,
361		ctx: EngineContext,
362		query: Cow<'static, str>,
363		variables: Variables,
364	) -> EngineFuture<'_, Vec<QueryResult>>;
365
366	/// Executes SurrealQL, sending results into `items` as they are produced.
367	///
368	/// The returned future is the execution: drive it while draining `items`,
369	/// and treat the channel closing as "no more results" rather than as
370	/// success, since a failure that belongs to no single statement is reported
371	/// by the future.
372	///
373	/// The default answers from [`Self::query`] and replays the finished
374	/// results, which is what an engine whose transport cannot carry
375	/// incremental results should do — the caller sees the same items either
376	/// way, just all at once. Overriding it is worthwhile only where results
377	/// can actually reach the caller before the query ends.
378	///
379	/// Rows are provisional until their statement's
380	/// [`Finished`](QueryStreamItem::Finished) item arrives; see
381	/// [`QueryStreamItem`].
382	fn query_stream(
383		&self,
384		ctx: EngineContext,
385		query: Cow<'static, str>,
386		variables: Variables,
387		items: Sender<QueryStreamItem>,
388	) -> EngineFuture<'_, ()> {
389		buffered_query_stream(self, ctx, query, variables, items)
390	}
391
392	/// Calls a function, or a machine learning model when `version` is set.
393	fn run(
394		&self,
395		ctx: EngineContext,
396		name: String,
397		version: Option<String>,
398		args: Array,
399	) -> EngineFuture<'_, Value>;
400
401	// ------------------------------------------------------------------
402	// Session state
403	// ------------------------------------------------------------------
404
405	/// Selects the namespace and/or database, returning the resulting
406	/// selection. `None` for either argument leaves that one unchanged.
407	fn use_ns_db(
408		&self,
409		ctx: EngineContext,
410		namespace: Option<String>,
411		database: Option<String>,
412	) -> EngineFuture<'_, (Option<String>, Option<String>)>;
413
414	/// Binds a session variable.
415	fn set(&self, ctx: EngineContext, key: String, value: Value) -> EngineFuture<'_, ()>;
416
417	/// Removes a session variable.
418	fn unset(&self, ctx: EngineContext, key: String) -> EngineFuture<'_, ()>;
419
420	// ------------------------------------------------------------------
421	// Authentication
422	// ------------------------------------------------------------------
423
424	/// Registers a record user and authenticates the session as them.
425	fn signup(&self, ctx: EngineContext, credentials: Object) -> EngineFuture<'_, Token>;
426
427	/// Authenticates the session with credentials.
428	fn signin(&self, ctx: EngineContext, credentials: Object) -> EngineFuture<'_, Token>;
429
430	/// Authenticates the session with an existing token, returning the token
431	/// now in effect.
432	///
433	/// A server may hand back a token of its own rather than the one it was
434	/// given, so the result is what the session is authenticated with -- not
435	/// necessarily the argument.
436	fn authenticate(&self, ctx: EngineContext, token: Token) -> EngineFuture<'_, Token>;
437
438	/// Exchanges a refresh token for a fresh token pair.
439	fn refresh(&self, ctx: EngineContext, token: Token) -> EngineFuture<'_, Token>;
440
441	/// Invalidates a refresh token so it can no longer be redeemed.
442	fn revoke(&self, ctx: EngineContext, token: Token) -> EngineFuture<'_, ()>;
443
444	/// Drops the session's authentication.
445	fn invalidate(&self, ctx: EngineContext) -> EngineFuture<'_, ()>;
446
447	// ------------------------------------------------------------------
448	// Transactions
449	// ------------------------------------------------------------------
450
451	/// Opens an explicit transaction, returning its id.
452	fn begin(&self, ctx: EngineContext) -> EngineFuture<'_, Uuid>;
453
454	/// Commits an explicit transaction.
455	fn commit(&self, ctx: EngineContext, txn: Uuid) -> EngineFuture<'_, ()>;
456
457	/// Cancels an explicit transaction.
458	fn rollback(&self, ctx: EngineContext, txn: Uuid) -> EngineFuture<'_, ()>;
459
460	// ------------------------------------------------------------------
461	// Connection
462	// ------------------------------------------------------------------
463
464	/// Checks that the engine is reachable.
465	fn health(&self, ctx: EngineContext) -> EngineFuture<'_, ()>;
466
467	/// Reports the database version, as the server spells it (for example
468	/// `surrealdb-3.0.0`).
469	fn version(&self, ctx: EngineContext) -> EngineFuture<'_, String>;
470
471	// ------------------------------------------------------------------
472	// Live queries
473	// ------------------------------------------------------------------
474
475	/// Registers the channel a live query's notifications are delivered on.
476	fn subscribe_live(
477		&self,
478		_ctx: EngineContext,
479		_uuid: Uuid,
480		_notifications: Sender<Result<Notification, Error>>,
481	) -> EngineFuture<'_, ()> {
482		Box::pin(async { Err(unsupported("Live queries")) })
483	}
484
485	/// Kills a live query.
486	fn kill(&self, _ctx: EngineContext, _uuid: Uuid) -> EngineFuture<'_, ()> {
487		Box::pin(async { Err(unsupported("Live queries")) })
488	}
489
490	// ------------------------------------------------------------------
491	// Export and import
492	// ------------------------------------------------------------------
493
494	/// Exports the database to a file.
495	fn export_file(
496		&self,
497		_ctx: EngineContext,
498		_path: PathBuf,
499		_config: Option<DbExportConfig>,
500	) -> EngineFuture<'_, ()> {
501		Box::pin(async { Err(unsupported("Export")) })
502	}
503
504	/// Exports the database, streaming it to a channel.
505	fn export_bytes(
506		&self,
507		_ctx: EngineContext,
508		_bytes: Sender<Result<Vec<u8>, Error>>,
509		_config: Option<DbExportConfig>,
510	) -> EngineFuture<'_, ()> {
511		Box::pin(async { Err(unsupported("Export")) })
512	}
513
514	/// Exports a machine learning model to a file.
515	fn export_ml_file(
516		&self,
517		_ctx: EngineContext,
518		_path: PathBuf,
519		_config: MlExportConfig,
520	) -> EngineFuture<'_, ()> {
521		Box::pin(async { Err(unsupported("Machine learning model export")) })
522	}
523
524	/// Exports a machine learning model, streaming it to a channel.
525	fn export_ml_bytes(
526		&self,
527		_ctx: EngineContext,
528		_bytes: Sender<Result<Vec<u8>, Error>>,
529		_config: MlExportConfig,
530	) -> EngineFuture<'_, ()> {
531		Box::pin(async { Err(unsupported("Machine learning model export")) })
532	}
533
534	/// Imports a database export from a file.
535	fn import_file(&self, _ctx: EngineContext, _path: PathBuf) -> EngineFuture<'_, ()> {
536		Box::pin(async { Err(unsupported("Import")) })
537	}
538
539	/// Imports a machine learning model from a file.
540	fn import_ml_file(&self, _ctx: EngineContext, _path: PathBuf) -> EngineFuture<'_, ()> {
541		Box::pin(async { Err(unsupported("Machine learning model import")) })
542	}
543}
544
545/// The error an engine reports for an operation it does not serve.
546fn unsupported(what: &str) -> Error {
547	Error::configuration(format!("{what} is not supported by this engine"), None)
548}
549
550/// Flattens the results of an operation that runs a single statement into the
551/// one value it produced.
552///
553/// An empty reply reads as [`Value::None`]: an operation with no result may
554/// answer with either, and both mean the same thing. Anything longer than one
555/// result is a bug in the engine rather than something to report to the user.
556pub fn single_result(mut results: Vec<QueryResult>) -> Result<Value, Error> {
557	match results.len() {
558		0 => Ok(Value::None),
559		1 => results.remove(0).result,
560		_ => Err(Error::internal("expected the database to return one or no results".to_string())),
561	}
562}
563
564/// A [`SurrealEngine`] that drives an engine which consumes [`Route`]s.
565///
566/// The WebSocket and HTTP engines each run a task that reads `Route`s off a
567/// channel and answers on the response channel a `Route` carries -- the
568/// channel is also their session replay log, so a reconnection can rebuild
569/// what the connection had established. This adapter is the whole of what it
570/// takes to expose one of them through [`SurrealEngine`]: it turns each typed
571/// call back into the [`Command`] that task already understands, and unwraps
572/// the single response into the type the method promises.
573///
574/// [`Command`] is therefore an implementation detail of those two engines, not
575/// part of the interface: an engine with no route channel -- the embedded and
576/// gRPC ones -- never constructs a `Command` at all.
577#[derive(Debug, Clone)]
578pub struct RouteChannelEngine {
579	sender: Sender<Route>,
580}
581
582impl RouteChannelEngine {
583	/// Wraps a route sender as a [`SurrealEngine`].
584	pub fn new(sender: Sender<Route>) -> Self {
585		Self {
586			sender,
587		}
588	}
589
590	/// Sends one command and awaits its single response, flattening the
591	/// engine's `Vec<QueryResult>` reply into the one value these
592	/// non-`query` operations return.
593	///
594	/// An empty reply reads as [`Value::None`]: the route protocol lets an
595	/// engine answer a no-result operation with either an empty vector or a
596	/// single `Value::None`, and both mean the same thing.
597	async fn value(&self, command: Command, session: Uuid) -> Result<Value, Error> {
598		single_result(self.results(command, session).await?)
599	}
600
601	/// Sends one command and awaits its single response.
602	async fn results(&self, command: Command, session: Uuid) -> Result<Vec<QueryResult>, Error> {
603		let (response, receiver) = async_channel::bounded(1);
604		let route = Route {
605			request: RequestData {
606				command,
607				session_id: session,
608			},
609			response,
610		};
611		// Both failure modes mean the engine task is gone, which callers
612		// distinguish from a database error with `Error::is_connection()` to
613		// decide whether reconnecting is worth trying.
614		self.sender.send(route).await.map_err(|e| {
615			Error::connection(
616				format!("Failed to send command: {e}"),
617				ConnectionError::ConnectionFailed,
618			)
619		})?;
620		receiver.recv().await.map_err(|_| {
621			Error::connection(
622				"The engine dropped the request without answering".to_string(),
623				ConnectionError::ConnectionFailed,
624			)
625		})?
626	}
627
628	/// Sends one command whose response carries nothing of interest.
629	async fn unit(&self, command: Command, session: Uuid) -> Result<(), Error> {
630		match self.value(command, session).await? {
631			Value::None | Value::Null => Ok(()),
632			Value::Array(array) if array.is_empty() => Ok(()),
633			_ => Err(Error::internal("expected the database to return nothing".to_string())),
634		}
635	}
636}
637
638/// Converts the value an engine returns for signin/signup/refresh into a
639/// [`Token`].
640///
641/// These engines answer with the token's wire form (a bare string, or an
642/// object carrying `token` and `refresh`), which is exactly what `Token`
643/// deserialises from.
644fn value_to_token(value: Value) -> Result<Token, Error> {
645	// signin/signup answers historically arrive wrapped in a single-element
646	// array from some engines; unwrap that before converting.
647	let value = match value {
648		Value::Array(array) if array.len() == 1 => {
649			array.into_iter().next().expect("array has exactly one element")
650		}
651		value => value,
652	};
653	Token::from_value(value)
654}
655
656impl SurrealEngine for RouteChannelEngine {
657	fn query(
658		&self,
659		ctx: EngineContext,
660		query: Cow<'static, str>,
661		variables: Variables,
662	) -> EngineFuture<'_, Vec<QueryResult>> {
663		Box::pin(self.results(
664			Command::Query {
665				txn: ctx.transaction,
666				query,
667				variables,
668			},
669			ctx.session,
670		))
671	}
672
673	fn run(
674		&self,
675		ctx: EngineContext,
676		name: String,
677		version: Option<String>,
678		args: Array,
679	) -> EngineFuture<'_, Value> {
680		Box::pin(self.value(
681			Command::Run {
682				name,
683				version,
684				args,
685			},
686			ctx.session,
687		))
688	}
689
690	fn use_ns_db(
691		&self,
692		ctx: EngineContext,
693		namespace: Option<String>,
694		database: Option<String>,
695	) -> EngineFuture<'_, (Option<String>, Option<String>)> {
696		Box::pin(async move {
697			let value = self
698				.value(
699					Command::Use {
700						namespace,
701						database,
702					},
703					ctx.session,
704				)
705				.await?;
706			// Engines that predate reporting the resulting selection answer
707			// with something other than an object; report "unknown" rather
708			// than failing, as the SDK has always done.
709			let Value::Object(object) = value else {
710				return Ok((None, None));
711			};
712			let read = |key: &str| object.get(key).and_then(|v| v.as_string()).map(String::from);
713			Ok((read("namespace"), read("database")))
714		})
715	}
716
717	fn set(&self, ctx: EngineContext, key: String, value: Value) -> EngineFuture<'_, ()> {
718		Box::pin(self.unit(
719			Command::Set {
720				key,
721				value,
722			},
723			ctx.session,
724		))
725	}
726
727	fn unset(&self, ctx: EngineContext, key: String) -> EngineFuture<'_, ()> {
728		Box::pin(self.unit(
729			Command::Unset {
730				key,
731			},
732			ctx.session,
733		))
734	}
735
736	fn signup(&self, ctx: EngineContext, credentials: Object) -> EngineFuture<'_, Token> {
737		Box::pin(async move {
738			let value = self
739				.value(
740					Command::Signup {
741						credentials,
742					},
743					ctx.session,
744				)
745				.await?;
746			value_to_token(value)
747		})
748	}
749
750	fn signin(&self, ctx: EngineContext, credentials: Object) -> EngineFuture<'_, Token> {
751		Box::pin(async move {
752			let value = self
753				.value(
754					Command::Signin {
755						credentials,
756					},
757					ctx.session,
758				)
759				.await?;
760			value_to_token(value)
761		})
762	}
763
764	fn authenticate(&self, ctx: EngineContext, token: Token) -> EngineFuture<'_, Token> {
765		Box::pin(async move {
766			let value = self
767				.value(
768					Command::Authenticate {
769						token,
770					},
771					ctx.session,
772				)
773				.await?;
774			value_to_token(value)
775		})
776	}
777
778	fn refresh(&self, ctx: EngineContext, token: Token) -> EngineFuture<'_, Token> {
779		Box::pin(async move {
780			let value = self
781				.value(
782					Command::Refresh {
783						token,
784					},
785					ctx.session,
786				)
787				.await?;
788			value_to_token(value)
789		})
790	}
791
792	fn revoke(&self, ctx: EngineContext, token: Token) -> EngineFuture<'_, ()> {
793		Box::pin(self.unit(
794			Command::Revoke {
795				token,
796			},
797			ctx.session,
798		))
799	}
800
801	fn invalidate(&self, ctx: EngineContext) -> EngineFuture<'_, ()> {
802		Box::pin(self.unit(Command::Invalidate, ctx.session))
803	}
804
805	fn begin(&self, ctx: EngineContext) -> EngineFuture<'_, Uuid> {
806		Box::pin(async move {
807			let value = self.value(Command::Begin, ctx.session).await?;
808			let uuid = value.into_uuid().map_err(|e| Error::internal(e.to_string()))?;
809			Ok(uuid.into_inner())
810		})
811	}
812
813	fn commit(&self, ctx: EngineContext, txn: Uuid) -> EngineFuture<'_, ()> {
814		Box::pin(async move {
815			self.value(
816				Command::Commit {
817					txn,
818				},
819				ctx.session,
820			)
821			.await?;
822			Ok(())
823		})
824	}
825
826	fn rollback(&self, ctx: EngineContext, txn: Uuid) -> EngineFuture<'_, ()> {
827		Box::pin(async move {
828			self.value(
829				Command::Rollback {
830					txn,
831				},
832				ctx.session,
833			)
834			.await?;
835			Ok(())
836		})
837	}
838
839	fn health(&self, ctx: EngineContext) -> EngineFuture<'_, ()> {
840		Box::pin(self.unit(Command::Health, ctx.session))
841	}
842
843	fn version(&self, ctx: EngineContext) -> EngineFuture<'_, String> {
844		Box::pin(async move {
845			let value = self.value(Command::Version, ctx.session).await?;
846			value.into_string().map_err(|e| Error::internal(e.to_string()))
847		})
848	}
849
850	fn subscribe_live(
851		&self,
852		ctx: EngineContext,
853		uuid: Uuid,
854		notifications: Sender<Result<Notification, Error>>,
855	) -> EngineFuture<'_, ()> {
856		Box::pin(self.unit(
857			Command::SubscribeLive {
858				uuid,
859				notification_sender: notifications,
860			},
861			ctx.session,
862		))
863	}
864
865	fn kill(&self, ctx: EngineContext, uuid: Uuid) -> EngineFuture<'_, ()> {
866		Box::pin(self.unit(
867			Command::Kill {
868				uuid,
869			},
870			ctx.session,
871		))
872	}
873
874	fn export_file(
875		&self,
876		ctx: EngineContext,
877		path: PathBuf,
878		config: Option<DbExportConfig>,
879	) -> EngineFuture<'_, ()> {
880		Box::pin(self.unit(
881			Command::ExportFile {
882				path,
883				config,
884			},
885			ctx.session,
886		))
887	}
888
889	fn export_bytes(
890		&self,
891		ctx: EngineContext,
892		bytes: Sender<Result<Vec<u8>, Error>>,
893		config: Option<DbExportConfig>,
894	) -> EngineFuture<'_, ()> {
895		Box::pin(self.unit(
896			Command::ExportBytes {
897				bytes,
898				config,
899			},
900			ctx.session,
901		))
902	}
903
904	fn export_ml_file(
905		&self,
906		ctx: EngineContext,
907		path: PathBuf,
908		config: MlExportConfig,
909	) -> EngineFuture<'_, ()> {
910		Box::pin(self.unit(
911			Command::ExportMl {
912				path,
913				config,
914			},
915			ctx.session,
916		))
917	}
918
919	fn export_ml_bytes(
920		&self,
921		ctx: EngineContext,
922		bytes: Sender<Result<Vec<u8>, Error>>,
923		config: MlExportConfig,
924	) -> EngineFuture<'_, ()> {
925		Box::pin(self.unit(
926			Command::ExportBytesMl {
927				bytes,
928				config,
929			},
930			ctx.session,
931		))
932	}
933
934	fn import_file(&self, ctx: EngineContext, path: PathBuf) -> EngineFuture<'_, ()> {
935		Box::pin(self.unit(
936			Command::ImportFile {
937				path,
938			},
939			ctx.session,
940		))
941	}
942
943	fn import_ml_file(&self, ctx: EngineContext, path: PathBuf) -> EngineFuture<'_, ()> {
944		Box::pin(self.unit(
945			Command::ImportMl {
946				path,
947			},
948			ctx.session,
949		))
950	}
951}
952
953#[cfg(test)]
954mod tests {
955	use surrealdb_types::Value;
956
957	use super::*;
958
959	/// The buffered adaptation produces the same items a streaming engine
960	/// would, so a caller cannot tell which one answered.
961	#[tokio::test]
962	async fn the_buffered_adaptation_produces_the_same_items() {
963		let (sender, routes) = async_channel::bounded(1);
964		let engine = RouteChannelEngine::new(sender);
965		let (items, received) = async_channel::bounded(8);
966		let stream = engine.query_stream(
967			EngineContext::new(Uuid::nil()),
968			Cow::Borrowed("SELECT * FROM thing"),
969			Variables::default(),
970			items,
971		);
972		let serve = async {
973			let route = routes.recv().await.expect("a route");
974			let _ = route
975				.response
976				.send(Ok(vec![QueryResult {
977					time: std::time::Duration::ZERO,
978					result: Ok(Value::Array(vec![Value::Bool(true)].into())),
979					query_type: surrealdb_rpc::QueryType::Other,
980				}]))
981				.await;
982		};
983		let (outcome, ()) = futures::future::join(stream, serve).await;
984		outcome.expect("the engine answered");
985
986		let mut items = Vec::new();
987		while let Ok(item) = received.try_recv() {
988			items.push(item);
989		}
990		assert!(matches!(items[0], QueryStreamItem::Rows { .. }), "a list becomes rows");
991		assert!(
992			matches!(
993				items[1],
994				QueryStreamItem::Finished {
995					error: None,
996					..
997				}
998			),
999			"and the statement is terminated"
1000		);
1001		assert_eq!(items.len(), 2);
1002	}
1003}