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
use anyhow::Result;
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::val::Value;
impl Document {
pub(crate) async fn insert(
&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)?;
// Skip the create attempt when we already have the document
if !self.is_iteration_initial() {
return self.insert_update(stk, ctx, opt, stm).await;
}
// Generate a record id up front
self.generate_record_id(stk, ctx, opt).await?;
// ON DUPLICATE KEY UPDATE makes a conflict retryable as an update
let retryable = stm.update().is_some();
// 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.insert_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 is no ON DUPLICATE KEY UPDATE clause
Err(IgnoreError::Error(e)) if !retryable => {
ctx.tx().rollback_to_save_point().await?;
self.mutated = false;
if stm.is_ignore() {
return Err(IgnoreError::Ignore);
} else {
return Err(IgnoreError::Error(e));
}
}
// There was an error creating the record
Err(IgnoreError::Error(e)) => {
// Only two failures name an existing record the update can be
// retried against. A unique-index conflict is raised by the
// index layer and a duplicate record id by the document layer,
// so each is recovered on its own. Both recoveries only borrow,
// so `e` stays whole for the arm that re-raises it.
let conflict = match crate::idx::index_exists_record(&e) {
// An index conflict against a record id the statement named
// itself cannot be retried, so it is reported like any other
// create conflict
Some(record) if !self.is_specific_record_id() => Some(record),
Some(_) => None,
None => match e.downcast_ref::<DocError>() {
Some(DocError::RecordExists {
record,
}) => Some(record.clone()),
_ => None,
},
};
match conflict {
Some(record) => record,
// Every other create failure belongs to this row whichever
// layer raised it, so `INSERT IGNORE` skips the row exactly
// as it does when there is no `ON DUPLICATE KEY UPDATE`
// clause. Classifying by error type instead would silently
// stop ignoring a failure the moment it moved layers, which
// `reproductions/insert_ignore_on_duplicate_key_leaf_error`
// pins: a `VALUE (THROW ...)` is raised outside the document
// layer and must still be skipped.
//
// This is wider than the pre-split behaviour, which surfaced
// anything that did not downcast to `err::Error`. That is a
// deliberate behaviour change, not an oversight.
None => {
ctx.tx().rollback_to_save_point().await?;
self.mutated = false;
if stm.is_ignore() {
return Err(IgnoreError::Ignore);
} else {
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 `ON DUPLICATE KEY UPDATE`
self.insert_update(stk, ctx, opt, stm).await
}
/// Attempt to run an INSERT statement to
/// create a record which does not exist
async fn insert_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_insert()?;
// 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?;
// Ensure all special fields are valid
self.check_data_fields()?;
// Set the specified record content
self.process_merge_data()?;
// 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_edges_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::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 INSERT statement to
/// update a record which already exists
async fn insert_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_insert()?;
// SECURITY: evaluate the table-level update permission BEFORE any
// user-supplied expression in the data clause. Otherwise a `ON
// DUPLICATE KEY SET x = THROW ...` could exfiltrate field values
// 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)?;
// 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_edges_data(ctx, opt).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
}
}