surrealdb-core 3.3.1

A scalable, distributed, collaborative, document-graph database, for the realtime web
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
//! LET operator - binds a value to a parameter name.
//!
//! LET is a context-mutating operator that adds a new parameter binding
//! to the execution context.

use std::sync::Arc;

use futures::stream;
use surrealdb_strand::Strand;
use surrealdb_types::{SqlFormat, ToSql};

use crate::err::Error;
use crate::exec::context::{ContextLevel, ExecutionContext};
use crate::exec::plan_or_compute::collect_stream;
use crate::exec::{
	AccessMode, BoxFut, CardinalityHint, Error as ExecError, ExecOperator, FlowResult,
	OperatorMetrics, OutputShape, ValueBatchStream, buffer_stream,
};
use crate::expr::{ControlFlow, Kind};
use crate::val::{Array, Value};

/// LET operator - binds a value to a parameter.
///
/// Implements `OperatorPlan` with `mutates_context() = true`.
/// The `output_context()` method evaluates the value and adds it to the
/// context parameters.
///
/// The value can be:
/// - A scalar expression (wrapped in `ExprPlan`) - evaluates to a single value
/// - A query - results are collected into an array
#[derive(Debug)]
pub struct LetPlan {
	/// Parameter name to bind (without $)
	pub name: Strand,
	/// Optional declared type for the binding — when present, the computed
	/// value is coerced to this kind before binding (mirrors
	/// `SetStatement::compute`).
	pub kind: Option<Kind>,
	/// Metrics for EXPLAIN ANALYZE
	pub(crate) metrics: Arc<OperatorMetrics>,
	/// Value to bind - either an ExprPlan for scalars or a query plan
	pub value: Arc<dyn ExecOperator>,
}

impl LetPlan {
	pub(crate) fn new(name: Strand, kind: Option<Kind>, value: Arc<dyn ExecOperator>) -> Self {
		Self {
			name,
			kind,
			value,
			metrics: Arc::new(OperatorMetrics::new()),
		}
	}

	fn coerce(&self, value: Value) -> Result<Value, Error> {
		match &self.kind {
			Some(kind) => value.coerce_to_kind(kind).map_err(|e| {
				ExecError::SetCoerce {
					name: self.name.to_string(),
					error: Box::new(e),
				}
				.into()
			}),
			None => Ok(value),
		}
	}

	/// Run the value plan and reduce its output to the single value this LET
	/// binds, before coercion.
	///
	/// A scalar value plan emits exactly one row, which becomes the bound value;
	/// an empty stream binds `NONE`. A query value plan binds all of its rows as
	/// an array.
	///
	/// Control-flow signals raised by the value plan are returned to the caller
	/// rather than resolved here, and reach it identically whether the plan
	/// raised them while building its stream or while draining it.
	async fn compute_value(&self, input: &ExecutionContext) -> crate::expr::FlowResult<Value> {
		let stream = buffer_stream(
			self.value.execute(input)?,
			self.value.access_mode(),
			self.value.cardinality_hint(),
			input.root().ctx.config.exec.operator_buffer_size,
		);
		let results = collect_stream(stream).await?;

		Ok(if self.value.output_shape().is_scalar() {
			results.into_iter().next().unwrap_or(Value::None)
		} else {
			Value::Array(Array(results))
		})
	}
}
impl ExecOperator for LetPlan {
	fn name(&self) -> &'static str {
		"Let"
	}

	fn attrs(&self) -> Vec<(String, String)> {
		vec![("name".to_string(), format!("${}", self.name.as_str()))]
	}

	fn required_context(&self) -> ContextLevel {
		self.value.required_context()
	}

	fn access_mode(&self) -> AccessMode {
		self.value.access_mode()
	}

	fn cardinality_hint(&self) -> CardinalityHint {
		CardinalityHint::AtMostOne
	}

	fn execute(&self, _ctx: &ExecutionContext) -> FlowResult<ValueBatchStream> {
		// LET returns NONE as its result (the binding happens in output_context)
		Ok(Box::pin(stream::once(async { Ok(crate::exec::ValueBatch::new(vec![Value::None])) })))
	}

	fn mutates_context(&self) -> bool {
		true
	}

	fn output_context<'a>(
		&'a self,
		input: &'a ExecutionContext,
	) -> BoxFut<'a, crate::expr::FlowResult<ExecutionContext>> {
		Box::pin(async move {
			// A RETURN from the value plan supplies the value to bind, and does so
			// whether it was raised as the plan started or partway through
			// producing rows.
			//
			// Everything else travels onwards unchanged: BREAK and CONTINUE belong
			// to an enclosing loop, so a LET in a FOR body breaks or continues that
			// loop rather than failing at the binding, and a stray one is rejected
			// by the statement boundary that has no loop to offer it. Errors reach
			// the boundary with their type intact; see the trait's
			// `output_context`.
			let computed_value = match self.compute_value(input).await {
				Ok(v) => v,
				Err(ControlFlow::Return(v)) => v,
				Err(ctrl) => return Err(ctrl),
			};

			// Apply declared type coercion (mirrors `SetStatement::compute`).
			let coerced = self.coerce(computed_value).map_err(|e| ControlFlow::Err(e.into()))?;
			Ok(input.with_param(self.name.clone(), coerced))
		})
	}

	fn children(&self) -> Vec<&Arc<dyn ExecOperator>> {
		vec![&self.value]
	}

	fn metrics(&self) -> Option<&OperatorMetrics> {
		Some(&self.metrics)
	}

	fn output_shape(&self) -> OutputShape {
		// `execute()` always emits exactly one row (`NONE`), so a caller reading
		// this plan's own result as an expression value — e.g. a block reduced to
		// a single LET statement — must see that one value unwrapped, not an
		// array containing it.
		OutputShape::Scalar
	}
}

impl ToSql for LetPlan {
	fn fmt_sql(&self, f: &mut String, _fmt: SqlFormat) {
		f.push_str("LET $");
		f.push_str(self.name.as_str());
		f.push_str(" = ");
		if self.value.output_shape().is_scalar() {
			f.push_str("<expr>");
		} else {
			f.push_str("(<query>)");
		}
	}
}

#[cfg(test)]
mod tests {
	use super::*;
	use crate::exec::operators::test_util::root_ctx;
	use crate::exec::{OutputOrdering, OutputShape, ValueBatch};

	/// Which signal a [`StubValue`] raises once its rows are exhausted.
	///
	/// A discriminant rather than a [`ControlFlow`]: the error arm owns an
	/// `anyhow::Error`, so `ControlFlow` is neither `Copy` nor `Clone` and cannot
	/// be handed to a stream that may be polled more than once.
	#[derive(Debug, Clone, Copy)]
	enum Signal {
		Break,
		Continue,
		Return(i64),
		/// Raises a fresh `ExecError::Thrown` each time, so the arm stays `Copy`
		/// even though `ControlFlow::Err` owns an `anyhow::Error`.
		Throw(&'static str),
	}

	impl Signal {
		fn build(self) -> ControlFlow {
			match self {
				Signal::Break => ControlFlow::Break,
				Signal::Continue => ControlFlow::Continue,
				Signal::Return(v) => ControlFlow::Return(Value::from(v)),
				Signal::Throw(msg) => {
					ControlFlow::Err(anyhow::Error::new(crate::exec::Error::Thrown(msg.to_owned())))
				}
			}
		}
	}

	/// A stub value plan for [`LetPlan`], covering the two places a value plan
	/// can raise a control-flow signal.
	///
	/// `rows` are emitted as one batch, then `signal` is raised. With `eager`
	/// set, `execute()` returns the signal instead of a stream and `rows` never
	/// appear; otherwise the signal arrives as a later stream item. `scalar`
	/// drives `output_shape()`, which decides whether `LetPlan` binds one value
	/// or an array.
	#[derive(Debug)]
	struct StubValue {
		rows: Vec<Value>,
		signal: Option<Signal>,
		eager: bool,
		scalar: bool,
		access_mode: AccessMode,
		required_context: ContextLevel,
	}

	impl StubValue {
		/// A plan that only yields `rows`, with no signal.
		fn rows(rows: Vec<Value>) -> Self {
			Self {
				rows,
				signal: None,
				eager: false,
				scalar: true,
				access_mode: AccessMode::ReadOnly,
				required_context: ContextLevel::Root,
			}
		}

		/// A plan whose `execute()` raises `signal` before building a stream.
		fn eager(signal: Signal) -> Self {
			Self {
				rows: Vec::new(),
				signal: Some(signal),
				eager: true,
				scalar: true,
				access_mode: AccessMode::ReadOnly,
				required_context: ContextLevel::Root,
			}
		}

		/// A plan that emits `rows`, then raises `signal` from the stream.
		fn rows_then(rows: Vec<Value>, signal: Signal) -> Self {
			Self {
				rows,
				signal: Some(signal),
				eager: false,
				scalar: true,
				access_mode: AccessMode::ReadOnly,
				required_context: ContextLevel::Root,
			}
		}

		fn non_scalar(mut self) -> Self {
			self.scalar = false;
			self
		}

		/// Declare the metadata `LetPlan` is expected to inherit.
		fn metadata(mut self, access_mode: AccessMode, required_context: ContextLevel) -> Self {
			self.access_mode = access_mode;
			self.required_context = required_context;
			self
		}

		fn into_operator(self) -> Arc<dyn ExecOperator> {
			Arc::new(self)
		}
	}

	impl ExecOperator for StubValue {
		fn name(&self) -> &'static str {
			"StubValue"
		}

		fn required_context(&self) -> ContextLevel {
			self.required_context
		}

		fn access_mode(&self) -> AccessMode {
			self.access_mode
		}

		fn cardinality_hint(&self) -> CardinalityHint {
			CardinalityHint::Unbounded
		}

		fn output_ordering(&self) -> OutputOrdering {
			OutputOrdering::Unordered
		}

		fn output_shape(&self) -> OutputShape {
			if self.scalar {
				OutputShape::Scalar
			} else {
				OutputShape::Rows
			}
		}

		fn execute(&self, _ctx: &ExecutionContext) -> FlowResult<ValueBatchStream> {
			if let Some(signal) = self.signal.filter(|_| self.eager) {
				return Err(signal.build());
			}

			let mut items: Vec<FlowResult<ValueBatch>> = Vec::new();
			if !self.rows.is_empty() {
				items.push(Ok(ValueBatch::new(self.rows.clone())));
			}
			if let Some(signal) = self.signal {
				items.push(Err(signal.build()));
			}

			Ok(Box::pin(stream::iter(items)))
		}
	}

	/// Build a `LET $x = <value>` plan with no declared type.
	fn let_plan(value: Arc<dyn ExecOperator>) -> LetPlan {
		LetPlan::new(Strand::new("x"), None, value)
	}

	/// Run `output_context` and read back the value bound to `$x`.
	async fn bound_value(plan: &LetPlan) -> crate::expr::FlowResult<Value> {
		let ctx = root_ctx();
		let out = plan.output_context(&ctx).await?;
		Ok(out.value("x").cloned().unwrap_or(Value::None))
	}

	/// Run `output_context` and return the signal it passed on, panicking if it
	/// bound a value instead.
	async fn propagated_signal(plan: &LetPlan) -> ControlFlow {
		match bound_value(plan).await {
			Ok(v) => panic!("expected a control-flow signal, got the binding: {v:?}"),
			Err(ctrl) => ctrl,
		}
	}

	#[tokio::test]
	async fn a_scalar_value_plan_binds_its_single_row() {
		let plan = let_plan(StubValue::rows(vec![Value::from(7)]).into_operator());
		assert_eq!(bound_value(&plan).await.unwrap(), Value::from(7));
	}

	#[tokio::test]
	async fn an_empty_scalar_value_plan_binds_none() {
		let plan = let_plan(StubValue::rows(Vec::new()).into_operator());
		assert_eq!(bound_value(&plan).await.unwrap(), Value::None);
	}

	#[tokio::test]
	async fn a_non_scalar_value_plan_binds_all_rows_as_an_array() {
		let plan = let_plan(
			StubValue::rows(vec![Value::from(1), Value::from(2)]).non_scalar().into_operator(),
		);
		assert_eq!(
			bound_value(&plan).await.unwrap(),
			Value::Array(Array(vec![Value::from(1), Value::from(2)]))
		);
	}

	#[tokio::test]
	async fn a_loop_signal_raised_by_the_value_plans_execute_travels_onwards() {
		let plan = let_plan(StubValue::eager(Signal::Break).into_operator());
		assert!(matches!(propagated_signal(&plan).await, ControlFlow::Break));

		let plan = let_plan(StubValue::eager(Signal::Continue).into_operator());
		assert!(matches!(propagated_signal(&plan).await, ControlFlow::Continue));
	}

	#[tokio::test]
	async fn a_loop_signal_from_inside_the_value_stream_travels_onwards_too() {
		// The signal arrives after a row has already been collected. Dropping it
		// here is what let a `BREAK` inside a LET value leave the enclosing loop
		// running, with a plausible-looking binding in place of the signal.
		let rows = vec![Value::from(1)];

		let plan = let_plan(StubValue::rows_then(rows.clone(), Signal::Break).into_operator());
		assert!(matches!(propagated_signal(&plan).await, ControlFlow::Break));

		let plan = let_plan(StubValue::rows_then(rows, Signal::Continue).into_operator());
		assert!(matches!(propagated_signal(&plan).await, ControlFlow::Continue));
	}

	#[tokio::test]
	async fn a_return_raised_by_the_value_plans_execute_supplies_the_bound_value() {
		let plan = let_plan(StubValue::eager(Signal::Return(9)).into_operator());
		assert_eq!(bound_value(&plan).await.unwrap(), Value::from(9));
	}

	#[tokio::test]
	async fn a_return_from_inside_the_value_stream_wins_over_the_rows_before_it() {
		let plan = let_plan(
			StubValue::rows_then(vec![Value::from(1), Value::from(2)], Signal::Return(9))
				.into_operator(),
		);
		assert_eq!(bound_value(&plan).await.unwrap(), Value::from(9));
	}

	#[tokio::test]
	async fn a_declared_type_coerces_the_bound_value() {
		let plan = LetPlan::new(
			Strand::new("x"),
			Some(Kind::String),
			StubValue::rows(vec![Value::from(7)]).into_operator(),
		);
		let err = match bound_value(&plan).await {
			Err(ControlFlow::Err(e)) => e,
			other => panic!("expected a coercion failure, got: {other:?}"),
		};
		assert!(
			format!("{err}").contains("$x"),
			"coercion failure should name the parameter, got: {err}"
		);
	}

	// =========================================================================
	// Binding publication, shadowing, and inherited metadata
	// =========================================================================

	#[tokio::test]
	async fn the_binding_is_published_through_output_context_not_through_execute() {
		let ctx = root_ctx();
		let plan = let_plan(StubValue::rows(vec![Value::from(3i64)]).into_operator());

		// The executor only asks for the modified context when this is true.
		assert!(plan.mutates_context());

		assert_eq!(bound_value(&plan).await.expect("binding should succeed"), Value::from(3i64));

		// The input context is left alone; the binding travels only forward.
		assert!(ctx.value("x").is_none());
	}

	#[tokio::test]
	async fn execute_emits_a_single_none_row_whatever_the_bound_value_is() {
		let ctx = root_ctx();
		let plan: Arc<dyn ExecOperator> =
			Arc::new(let_plan(StubValue::rows(vec![Value::from(3i64)]).into_operator()));
		assert_eq!(
			crate::exec::operators::test_util::collect(&plan, &ctx).await,
			vec![Value::None],
			"LET is not an expression: its own output is always NONE"
		);
	}

	#[tokio::test]
	async fn output_shape_is_scalar_so_a_lone_let_statement_reduces_to_bare_none() {
		// A block reduced to this one statement (`{ LET $x = 1 }`) is read as an
		// expression value via `output_shape()`, not by draining `execute()`
		// itself. It must see the single NONE row unwrapped, matching every
		// other non-query expression, rather than an array containing it.
		let plan = let_plan(StubValue::rows(vec![Value::from(3i64)]).into_operator());
		assert_eq!(plan.output_shape(), OutputShape::Scalar);
	}

	#[tokio::test]
	async fn a_binding_shadows_an_outer_one_of_the_same_name_and_the_outer_stays_intact() {
		let outer = root_ctx().with_param("x", Value::from(1i64));
		let plan = let_plan(StubValue::rows(vec![Value::from(2i64)]).into_operator());

		let inner = plan.output_context(&outer).await.expect("binding should succeed");
		assert_eq!(inner.value("x").cloned(), Some(Value::from(2i64)));

		// `output_context` layers a child context; the context it was given still
		// resolves `$x` to the outer value.
		assert_eq!(outer.value("x").cloned(), Some(Value::from(1i64)));
	}

	#[tokio::test]
	async fn metadata_is_inherited_from_the_value_plan() {
		// `access_mode` selects the transaction type, so `LET $x = (CREATE …)` has
		// to report ReadWrite; `required_context` drives the executor's pre-flight
		// level check.
		let value = StubValue::rows(Vec::new())
			.metadata(AccessMode::ReadWrite, ContextLevel::Database)
			.into_operator();
		let plan = let_plan(value);
		assert_eq!(plan.access_mode(), AccessMode::ReadWrite);
		assert_eq!(plan.required_context(), ContextLevel::Database);
	}

	// =========================================================================
	// Coercion failure and error transparency
	// =========================================================================

	#[tokio::test]
	async fn a_value_that_does_not_satisfy_the_declared_type_fails_naming_the_parameter() {
		let ctx = root_ctx();
		let plan = LetPlan::new(
			Strand::new("x"),
			Some(Kind::Int),
			StubValue::rows(vec![Value::from("not a number")]).into_operator(),
		);
		let ctrl = plan.output_context(&ctx).await.expect_err("a string cannot coerce to int");
		let ControlFlow::Err(err) = ctrl else {
			panic!("a coercion failure is an error, not a control-flow signal");
		};
		assert!(
			matches!(
				err.downcast_ref::<crate::err::Error>(),
				Some(crate::err::Error::Exec(crate::exec::Error::SetCoerce { name, .. }))
					if name == "x"
			),
			"expected SetCoerce naming $x, got {err:?}"
		);
	}

	/// An error must reach the caller as itself, not flattened into a string: the
	/// transactor has to recognise a write conflict to retry it, and a cancelled
	/// or timed-out query has to keep reporting as such.
	#[tokio::test]
	async fn an_error_reaches_the_caller_downcastable_from_either_raise_site() {
		for value in [
			StubValue::eager(Signal::Throw("boom")).into_operator(),
			StubValue::rows_then(vec![Value::from(1i64)], Signal::Throw("boom")).into_operator(),
		] {
			let plan = let_plan(value);
			let ControlFlow::Err(err) = propagated_signal(&plan).await else {
				panic!("expected an error");
			};
			assert!(
				matches!(
					err.downcast_ref::<crate::exec::Error>(),
					Some(crate::exec::Error::Thrown(msg)) if msg == "boom"
				),
				"expected the original Thrown error, got {err:?}"
			);
		}
	}
}