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
use anyhow::anyhow;
use reblessive::tree::Stk;
use super::IgnoreError;
use crate::catalog::providers::TableProvider;
use crate::ctx::FrozenContext;
use crate::dbs::{Options, Statement};
use crate::doc::{Document, Error as DocError};
use crate::err::Error;
use crate::val::Value;
impl Document {
pub(crate) async fn upsert(
&mut self,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
stm: &Statement<'_>,
) -> Result<Value, IgnoreError> {
// SECURITY (GHSA-2v9j): confine record users to their own tenant. The
// read path enforces this at the scan operators, but writes perform no
// scan, so the gate must be applied explicitly here.
self.check_record_user_access(opt)?;
// Error for tracking initial failures
let mut error: Option<anyhow::Error> = None;
// Skip the create attempt when we already have the document
if !self.is_iteration_initial() {
return self.upsert_update(stk, ctx, opt, stm).await;
}
// Save point so a failed create attempt can be rolled back
ctx.tx().new_save_point().await?;
// Try create first; on a recoverable conflict, fall through to update
let retry = match self.upsert_create(stk, ctx, opt, stm).await {
// Record created successfully
Ok(x) => {
ctx.tx().release_last_save_point().await?;
return Ok(x);
}
// We should ignore this record
Err(IgnoreError::Ignore) => {
ctx.tx().release_last_save_point().await?;
return Err(IgnoreError::Ignore);
}
// There was an error creating the record
Err(IgnoreError::Error(e)) => {
// A unique-index conflict is raised by the index layer, so it is
// recovered ahead of the core-error ladder below, which cannot
// see it. The recovery only borrows, so `e` stays whole for the
// arms that re-raise it.
let index_conflict = crate::idx::index_exists_record(&e);
match index_conflict {
// We got an index exists error
Some(record) if !self.is_specific_record_id() => record,
// An index conflict against a record id the statement named
// itself cannot be retried, so it is reported like any other
// create conflict
Some(_) => {
ctx.tx().rollback_to_save_point().await?;
self.mutated = false;
return Err(IgnoreError::Error(e));
}
None => match e.downcast::<DocError>() {
// This record already exists
Ok(DocError::RecordExists {
record,
}) => record,
// There was a possible schema error
Ok(e) if e.is_schema_related() && stm.is_repeatable() => {
error = Some(e.into());
self.inner_id()?
}
// Any other document failure is a conflict
Ok(e) => {
ctx.tx().rollback_to_save_point().await?;
self.mutated = false;
return Err(IgnoreError::Error(anyhow!(e)));
}
Err(e) => match e.downcast::<Error>() {
// There was a conflict error
Ok(e) => {
ctx.tx().rollback_to_save_point().await?;
self.mutated = false;
return Err(IgnoreError::Error(anyhow!(e)));
}
// Unrelated error — always surface
Err(e) => {
ctx.tx().rollback_to_save_point().await?;
self.mutated = false;
return Err(IgnoreError::Error(e));
}
},
},
}
}
};
// Roll back the create attempt before falling through to update
ctx.tx().rollback_to_save_point().await?;
// Reset any mutation tracking, for retry
self.mutated = false;
// Check if the request is finished
if ctx.is_done(None).await? {
return Err(IgnoreError::Ignore);
}
// Get the namespace id
let ns = self.doc_ctx.ns().namespace_id;
// Get the database id
let db = self.doc_ctx.db().database_id;
// Get the already stored record
let val = ctx.tx().get_record(ns, db, &retry.table, &retry.key, opt.version).await?;
// Reset the document for the retry
self.modify_for_update_retry(retry, val);
// Update the document with `UPDATE`
let res = self.upsert_update(stk, ctx, opt, stm).await;
// Return any carried over error if set
match error {
Some(e) => match res {
Err(_) => Err(IgnoreError::Error(e)),
Ok(v) => Ok(v),
},
None => res,
}
}
/// Attempt to run an UPSERT statement to
/// create a record which does not exist
async fn upsert_create(
&mut self,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
stm: &Statement<'_>,
) -> Result<Value, IgnoreError> {
// Ensure we can store this type of record
self.check_table_type_upsert()?;
// Ensure we can write to the table at all
self.check_permissions_quick_create(ctx, opt)?;
// Reject writes to read-only view tables (after the permission gate)
self.check_table_not_view(opt)?;
// Ensure any input data is computed
self.compute_input_data(stk, ctx, opt, stm).await?;
// Set the specified record content
self.process_record_data(stk, ctx, opt).await?;
// Generate a new record id if necessary
self.generate_record_id(stk, ctx, opt).await?;
// Ensure all special fields are valid
self.check_data_fields()?;
// Set the default record field values
self.default_record_data()?;
// Process the field schema for the table
self.process_table_fields(stk, ctx, opt, stm).await?;
// Clean up table fields and NONE values
self.cleanup_table_fields()?;
// Check table permissions after create
self.check_create_permissions(stk, ctx, opt, &self.current).await?;
// Store the document and index data
self.store_record_data(ctx, stm).await?;
self.store_index_data(stk, ctx, opt).await?;
// Materialise the computed fields the record's observers read: they are
// stripped before storage, so events, live queries, changefeeds and the
// output projection would otherwise see the record without them.
self.materialise_observed_fields(stk, ctx, opt).await?;
// Process additional table operations
self.process_table_references(stk, ctx, opt).await?;
self.process_table_views(stk, ctx, opt, super::Action::Create).await?;
self.process_table_events(stk, ctx, opt, super::Action::Create).await?;
self.process_table_lives(stk, ctx, opt, super::Action::Create).await?;
self.process_changefeeds(ctx, opt).await?;
// Check table permissions for output
self.check_select_permissions(stk, ctx, opt, &self.current).await?;
// Process the projected output document
self.output_write(stk, ctx, opt, stm.output(), stm).await
}
/// Attempt to run an UPSERT statement to
/// update a record which already exists
async fn upsert_update(
&mut self,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
stm: &Statement<'_>,
) -> Result<Value, IgnoreError> {
// Ensure the record actually exists
self.check_record_exists()?;
// Ensure we can store this type of record
self.check_table_type_upsert()?;
// SECURITY: evaluate the table-level update permission BEFORE any
// user-supplied expression in the WHERE clause or data clause.
// Otherwise a `WHERE THROW ...` / `SET x = THROW ...` could exfiltrate
// field values before the permission check rejects the operation.
self.check_update_permissions(stk, ctx, opt, &self.current).await?;
// Reject writes to read-only view tables (after the permission gate)
self.check_table_not_view(opt)?;
// Check if the WHERE condition is truthy BEFORE evaluating the data
// clause, so a data clause with side effects runs only for records the
// condition accepts. This also matches the index-backed plan, where
// records rejected by the condition never enter this pipeline at all.
self.check_where_condition(stk, ctx, opt, stm.cond()).await?;
// Ensure any input data is computed
self.compute_input_data(stk, ctx, opt, stm).await?;
// Ensure all special fields are valid
self.check_data_fields()?;
// Set the specified record content
self.process_record_data(stk, ctx, opt).await?;
// Set the default record field values
self.default_record_data()?;
// Process the field schema for the table
self.process_table_fields(stk, ctx, opt, stm).await?;
// Clean up table fields and NONE values
self.cleanup_table_fields()?;
// Check table permissions after update
self.recheck_update_permissions(stk, ctx, opt, &self.current).await?;
// Store the document and index data
self.store_record_data(ctx, stm).await?;
self.store_inline_adjacency_data(ctx, opt).await?;
self.store_index_data(stk, ctx, opt).await?;
// Materialise the computed fields the record's observers read: they are
// stripped before storage, so events, live queries, changefeeds and the
// output projection would otherwise see the record without them.
self.materialise_observed_fields(stk, ctx, opt).await?;
// Process additional table operations
self.process_table_references(stk, ctx, opt).await?;
self.process_table_views(stk, ctx, opt, super::Action::Update).await?;
self.process_table_events(stk, ctx, opt, super::Action::Update).await?;
self.process_table_lives(stk, ctx, opt, super::Action::Update).await?;
self.process_changefeeds(ctx, opt).await?;
// Check table permissions for output
self.check_select_permissions(stk, ctx, opt, &self.current).await?;
// Process the projected output document
self.output_write(stk, ctx, opt, stm.output(), stm).await
}
}