orion-server 1.11.1

Turn business logic into live REST/Kafka services, declared as JSON
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
use std::sync::Arc;

use async_trait::async_trait;
use dataflow_rs::engine::error::DataflowError;
use dataflow_rs::engine::task_context::TaskContext;
use serde_json::Value;

use super::connector_handler::{ConnectorHandler, Produced};
use super::connector_helpers::{
    ConnectorCall, QueryBudget, QueryFailure, acquire_conn, decode_failure, encode_failure,
    reject_mongo_connector, require_op_allowed, resolve_bind_params, resolve_row_format,
    to_connect_error,
};
use super::schema::{FieldKind, FieldSchema};
use super::templated_input::TemplatedInput;
use crate::connector::ConnectorRegistry;
use crate::connector::pool_cache::SqlPoolCache;
use crate::engine::HandlerError;

/// This handler's name, for the row-conversion helpers below — a reference to
/// the one place it is written (F48), not a second spelling of it.
const NAME: &str = <DbReadHandler as ConnectorHandler>::NAME;

/// Executes SQL SELECT queries against external databases configured via connectors.
pub struct DbReadHandler {
    pub pool_cache: Arc<SqlPoolCache>,
    pub registry: Arc<ConnectorRegistry>,
    /// Hard row cap, from `query.max_limit` (F10). Raw SQL can't have a
    /// LIMIT injected reliably, so rows are streamed and counted — one
    /// `SELECT * FROM big_table` must not OOM the process.
    pub max_rows: usize,
}

/// The statement and its bind values.
///
/// The SQL text is a literal read from the task; only the parameters come from
/// the message, which is what keeps them the sole request-controlled part of
/// the statement.
pub struct DbRead {
    query: String,
    params: Vec<Value>,
    format: crate::connector::sql_decode::RowFormat,
}

impl DbRead {
    /// The parse both raw-SQL handlers do. Shared rather than copied because
    /// `db_read` and `db_write` differ in what the database does with the
    /// statement, not in what the task says.
    /// The half `db_read` and `db_write` genuinely share: a literal statement
    /// and message-derived binds.
    ///
    /// The rendering choices (`numeric_as`, `binary_as`) are **not** read here.
    /// They govern how a decoded row renders,
    /// and `db_write` decodes no rows — it answers `rows_affected`. Reading it
    /// on both paths meant a wrong value produced `db_write: 'numeric_as' must
    /// be one of number/string`, an error naming a field `DB_WRITE_FIELDS` does
    /// not declare and the schema validator already reports as `UNKNOWN_FIELD`.
    /// One of the two had to go, and the runtime is the one that was wrong.
    pub(super) fn parse_statement(
        call: &ConnectorCall<'_>,
        input: &TemplatedInput,
        ctx: &TaskContext<'_>,
    ) -> Result<Self, HandlerError> {
        Ok(Self {
            query: call.require_str(input, "query")?.to_string(),
            params: resolve_bind_params(input, call.name, ctx)?,
            format: crate::connector::sql_decode::RowFormat::default(),
        })
    }

    /// [`parse_statement`](Self::parse_statement) plus the read-only rendering
    /// choice — and the check that the statement is actually a read.
    pub(super) fn parse_read(
        call: &ConnectorCall<'_>,
        input: &TemplatedInput,
        ctx: &TaskContext<'_>,
    ) -> Result<Self, HandlerError> {
        let read = Self {
            format: resolve_row_format(input, call.name, ctx)?,
            ..Self::parse_statement(call, input, ctx)?
        };
        require_read_only(&read.query, call.name)?;
        Ok(read)
    }

    pub(super) fn query(&self) -> &str {
        &self.query
    }

    pub(super) fn params(&self) -> &[Value] {
        &self.params
    }

    pub(super) fn format(&self) -> crate::connector::sql_decode::RowFormat {
        self.format
    }
}

#[async_trait]
impl ConnectorHandler for DbReadHandler {
    const NAME: &'static str = "db_read";
    type Kind = crate::connector::kind::Db;
    type Input = TemplatedInput;
    type Parsed = DbRead;

    fn registry(&self) -> &Arc<ConnectorRegistry> {
        &self.registry
    }

    fn parse(
        &self,
        call: &ConnectorCall<'_>,
        input: &TemplatedInput,
        ctx: &TaskContext<'_>,
    ) -> Result<Self::Parsed, HandlerError> {
        DbRead::parse_read(call, input, ctx)
    }

    fn gate(
        _parsed: &Self::Parsed,
        conn: &crate::connector::DbConnectorConfig,
        connector: &str,
    ) -> Result<(), HandlerError> {
        require_op_allowed(&conn.operations, "read", connector)?;
        // A Mongo connection string in a `db` connector is the right type and
        // the wrong backend: both are `ConnectorConfig::Db`, and only the
        // string tells them apart.
        Ok(reject_mongo_connector(
            <Self as ConnectorHandler>::NAME,
            connector,
            conn,
        )?)
    }

    async fn run(
        &self,
        read: Self::Parsed,
        db_config: &crate::connector::DbConnectorConfig,
        call: &ConnectorCall<'_>,
        _input: &TemplatedInput,
        _ctx: &mut TaskContext<'_>,
    ) -> Result<Produced, HandlerError> {
        let pool = self
            .pool_cache
            .get_pool(call.connector, db_config)
            .await
            .map_err(to_connect_error)?;

        let max_rows = self.max_rows;
        let params = read.params();
        let format = read.format();
        let query = read.query();

        // One budget over both legs. `acquire` is a round trip of its own, and
        // giving it a timeout of its own would silently make the bound the
        // connector's owner set into `connect_timeout_ms` plus
        // `query_timeout_ms`.
        let budget = QueryBudget::start(db_config.query_timeout_ms);

        // Bound once, outside the dispatch, because the vocabulary is the same
        // whichever driver runs it.
        let scalars: Vec<crate::connector::sql_encode::Scalar> =
            params.iter().map(Into::into).collect();

        // One body, three drivers: the macro binds the concrete pool, its
        // decoder and its binders, and each arm is type-checked on its own.
        let json = crate::connector::pool_cache::dispatch_sql_pool!(
            &pool, p, rows_to_json, bind, typed_args, _write_result => {
                // The connection is named rather than left to sqlx to acquire
                // per statement, because PostgreSQL caches a prepared
                // statement's parameter types per connection — asking the
                // server what it declared and then binding against it only
                // means anything if both happen on the same one.
                let mut conn = acquire_conn(&budget, call.name, p).await?;
                // Another leg of the same budget: on a connection that has
                // already seen this SQL it is a cache hit and no round trip at
                // all, but the first time through it is a real one.
                let bound = budget
                    .run(call.name, async {
                        typed_args(&mut conn, query, Some(&scalars))
                            .await
                            .map_err(|e| QueryFailure::Classified(encode_failure(NAME, e)))
                    })
                    .await?;
                let rows = budget.run(call.name, async {
                    use futures::TryStreamExt;
                    // `AssertSqlSafe` states what this handler is: the
                    // raw-SQL escape hatch, whose statement is authored in the
                    // workflow rather than assembled here. sqlx 0.9 asks the
                    // caller to own that, and the answer is the same as it was
                    // before it asked — the text comes from the definition, the
                    // author-supplied *values* travel as bind parameters beside
                    // it, and a definition reaches the runtime only through the
                    // admin API. Authors who want values checked use the
                    // portable dialect (`data_query`/`data_write`) instead.
                    let sqlx_query = match bound {
                        crate::connector::sql_encode::Bound::Typed(args) => {
                            sqlx::query_with(sqlx::AssertSqlSafe(query), args)
                        }
                        crate::connector::sql_encode::Bound::Fallback { cache } => {
                            bind(sqlx::query(sqlx::AssertSqlSafe(query)), params).persistent(cache)
                        }
                    };
                    let mut stream = sqlx_query.fetch(&mut *conn);
                    let mut rows = Vec::new();
                    // Not `.map_err(|e| e.to_string())`: stringifying here
                    // converted through `From<String>`, which is
                    // unconditionally a backend failure, so a constraint the
                    // driver had already classified was thrown away before
                    // `QueryFailure` could see it.
                    while let Some(row) = stream.try_next().await? {
                        if rows.len() >= max_rows {
                            // F42: classified so `timed_query` reports it as a
                            // 400 with the text intact rather than a 500 with
                            // the guidance sanitised away. The guidance *is*
                            // the message, so losing it loses the point.
                            return Err(
                                crate::engine::functions::connector_helpers::QueryFailure::Limit(
                                    format!(
                                        "{} result exceeds query.max_limit ({max_rows} rows) \
                                         — add a LIMIT to the query or raise the cap",
                                        call.name
                                    ),
                                ),
                            );
                        }
                        rows.push(row);
                    }
                    Ok(rows)
                })
                .await?;
                rows_to_json(&rows, format).map_err(|e| decode_failure(NAME, e))?
            }
        );

        Ok(Value::Array(json).into())
    }
}

// -- Input schema (F53) --
//
// The table describing this handler's `function.input` lives next to the
// handler it describes. It used to sit in `schema.rs` with the other nine,
// which is how every schema/handler divergence in the 1.0 audit happened:
// a field was added, renamed or made conditional here and the table saying
// so was in a different file.

pub(super) const DB_READ_FIELDS: &[FieldSchema] = &[
    FieldSchema {
        name: "connector",
        description: "Name of the SQL connector to query.",
        kind: FieldKind::String,
        required: true,
        ..FieldSchema::DEFAULT
    },
    FieldSchema {
        name: "query",
        description: "Read statement — SELECT, WITH, VALUES or TABLE; a write belongs in db_write, which has its own 'raw_write' connector gate. Bind placeholders are the backend's own spelling: ? for SQLite and MySQL, $1, $2, ... for PostgreSQL.",
        kind: FieldKind::String,
        required: true,
        ..FieldSchema::DEFAULT
    },
    FieldSchema {
        name: "params",
        description: "Array of values to bind to query placeholders, in order. Accepts {\"var\": \"path\"} to read the value from the message.",
        kind: FieldKind::Array,
        resolvable: true,
        ..FieldSchema::DEFAULT
    },
    FieldSchema {
        name: "numeric_as",
        description: "How an arbitrary-precision decimal column is rendered: \"number\" (default) or \"string\". A number is computable in JSONLogic and rounds beyond 2^53 or on most decimal fractions; a string keeps every digit, which is what a money column needs.",
        kind: FieldKind::String,
        resolvable: true,
        ..FieldSchema::DEFAULT
    },
    FieldSchema {
        name: "binary_as",
        description: "How a binary column is rendered: \"auto\" (default), \"hex\", \"base64\" or \"text\". Auto reads the bytes as text when they are valid UTF-8 and as hex when they are not, so its result shape depends on the data; name an encoding for a column that is genuinely binary.",
        kind: FieldKind::String,
        resolvable: true,
        ..FieldSchema::DEFAULT
    },
    FieldSchema {
        name: "output",
        description: "Dotted path in the message where rows are written. Defaults to \"data\".",
        kind: FieldKind::String,
        template_at: &[""],
        ..FieldSchema::DEFAULT
    },
];

// -- Read-only statement check --
//
// `db_read` gates on the connector's `read` operation and then runs whatever
// statement the task carries. `fetch` executes any statement — it merely
// streams whatever rows come back — so `DELETE FROM t RETURNING id` ran here on
// PostgreSQL and SQLite, and a bare `DELETE`/`UPDATE`/`INSERT` ran (returning no
// rows) on all three, while `raw_write: false` was set on the connector.
//
// That made the operation gates advertise more than they enforced. The gate
// table's own claim — "SQL writes are fully bounded by allowed_entities once
// `raw_write: false` leaves `data_write` as the only write path" — is only true
// with this check in place, because `db_read` was the second write path.
//
// The statement is a workflow-authored literal, never caller-supplied, so this
// is not an injection guard; it is what makes "delete-proof connector" a
// property an operator can rely on rather than a convention authors are asked
// to keep.

/// Refuse a `db_read` statement that is not a read, in the words the
/// handler has always used. The judgement is [`crate::sql_lex`]'s — the one
/// lexer every surface reads SQL with — read lossily, as the run-time check
/// always has been.
///
/// # Errors
///
/// [`DataflowError::Validation`] when the statement does not open with one of
/// [`crate::sql_lex::READ_STATEMENTS`], or when it carries a data-modifying
/// CTE.
fn require_read_only(query: &str, handler_name: &str) -> Result<(), HandlerError> {
    use crate::sql_lex::{READ_STATEMENTS, ReadOnlyViolation};
    let message = match crate::sql_lex::read_only_violation(query) {
        None => return Ok(()),
        Some(ReadOnlyViolation::Empty) => format!("{handler_name} 'query' has no statement to run"),
        Some(ReadOnlyViolation::NotARead { keyword }) => format!(
            "{handler_name} runs read statements only, but this one starts with \
             '{keyword}' — use db_write for INSERT/UPDATE/DELETE (it has its own \
             'raw_write' connector gate). Reads start with {}",
            READ_STATEMENTS.join(", ")
        ),
        Some(ReadOnlyViolation::ModifyingCte { keyword }) => format!(
            "{handler_name} runs read statements only, but this one carries a \
             data-modifying '{keyword}' common table expression — use db_write \
             (it has its own 'raw_write' connector gate)"
        ),
    };
    Err(DataflowError::Validation(message).into())
}

#[cfg(test)]
mod tests {
    use super::*;

    fn check(sql: &str) -> Result<(), String> {
        require_read_only(sql, "db_read").map_err(|e| {
            let e: DataflowError = e.into();
            e.to_string()
        })
    }

    #[test]
    fn reads_are_admitted() {
        for sql in [
            "SELECT id FROM users WHERE id = $1",
            "  \n select 1",
            "-- a comment\nSELECT 1",
            "/* block */ SELECT 1",
            "(SELECT 1) UNION (SELECT 2)",
            "WITH recent AS (SELECT * FROM orders) SELECT * FROM recent",
            "VALUES (1), (2)",
            "TABLE users",
            // A locking read is a read: `FOR UPDATE` must not be mistaken for
            // an UPDATE statement.
            "SELECT id FROM jobs ORDER BY id FOR UPDATE SKIP LOCKED",
            // The word only appears inside data.
            "SELECT id FROM notes WHERE body = 'delete from users'",
            "SELECT \"delete\" FROM t",
            "SELECT total AS deleted FROM t",
            "SELECT CAST(a AS text) FROM t",
        ] {
            assert!(
                check(sql).is_ok(),
                "must be admitted: {sql} — {:?}",
                check(sql)
            );
        }
    }

    #[test]
    fn writes_are_refused() {
        for sql in [
            "DELETE FROM audit_log WHERE id > 0 RETURNING id",
            "delete from audit_log",
            "INSERT INTO t (a) VALUES (1)",
            "UPDATE t SET a = 1",
            "TRUNCATE t",
            "DROP TABLE t",
            "PRAGMA journal_mode = WAL",
            // `EXPLAIN ANALYZE` executes the statement it explains.
            "EXPLAIN ANALYZE DELETE FROM t",
            "  -- lead in\n  DELETE FROM t",
        ] {
            let err = check(sql).expect_err(&format!("must be refused: {sql}"));
            assert!(err.contains("read statements only"), "{sql}: {err}");
        }
    }

    #[test]
    fn a_data_modifying_cte_is_refused() {
        for sql in [
            "WITH gone AS (DELETE FROM t RETURNING id) SELECT * FROM gone",
            "WITH added AS (INSERT INTO t (a) VALUES (1) RETURNING id) SELECT * FROM added",
            "with m as materialized (update t set a = 1 returning id) select * from m",
        ] {
            let err = check(sql).expect_err(&format!("must be refused: {sql}"));
            assert!(err.contains("data-modifying"), "{sql}: {err}");
        }
    }

    /// A statement whose text merely *mentions* a modifying keyword inside a
    /// literal, a comment or an identifier stays a read — the check reads
    /// syntax, not data.
    #[test]
    fn quoted_text_is_not_syntax() {
        assert!(check("SELECT 1 /* AS (DELETE */").is_ok());
        assert!(check("SELECT 'x AS (DELETE FROM t)' AS s").is_ok());
        assert!(check("SELECT $tag$ AS (DELETE FROM t) $tag$ AS s").is_ok());
        assert!(check("SELECT * FROM t WHERE a = $1 AND b = $2").is_ok());
    }

    #[test]
    fn an_empty_statement_is_refused() {
        let err = check("   -- nothing here\n").expect_err("empty");
        assert!(err.contains("no statement"), "{err}");
    }
}