-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathqueries.h
More file actions
816 lines (736 loc) · 25.1 KB
/
Copy pathqueries.h
File metadata and controls
816 lines (736 loc) · 25.1 KB
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
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
/**
* @file queries.h
* @brief PostgreSQL query representation and management
*
* This file defines the data structures and classes for representing
* and managing SQL queries for the PostgreSQL client, including:
*
* - Storage for prepared queries
* - Parameter binding for prepared statements
* - SQL query execution with callbacks
* - Various query types (BEGIN, COMMIT, ROLLBACK, etc.)
*
* The implementation follows the PostgreSQL protocol for preparing
* and executing queries, supporting both simple and prepared statements.
*
* @see qb::pg::detail::Transaction
* @see qb::pg::detail::Database
*
* @author qb - C++ Actor Framework
* @copyright Copyright (c) 2011-2026 qb - isndev (cpp.actor)
* Licensed under the Apache License, Version 2.0 (http://www.apache.org/licenses/LICENSE-2.0)
* @ingroup Pgsql
*/
#pragma once
#include <iomanip>
#include <iostream>
#include <list> // P2-1: LRU cache
#include <string_view>
#include <type_traits>
#include <vector>
#include <qb/io.h>
#include <qb/system/container/unordered_map.h>
#include <qb/utility/branch_hints.h>
#include "./common.h"
#include "./error.h"
#include "./param_serializer.h"
#include "./protocol.h"
#include "./type_mapping.h"
namespace qb::pg::detail {
using namespace qb::pg;
/**
* @brief Structure for storing a prepared query definition
*
* Contains all the information needed to execute a prepared statement,
* including its name, SQL expression, parameter types, and result description.
*/
struct PreparedQuery {
std::string name; ///< Name of the prepared query
std::string expression; ///< SQL expression
std::vector<oid> param_types; ///< Types of parameters (was type_oid_sequence)
row_description_type row_description; ///< Description of result columns
};
/**
* @brief Storage for prepared queries with LRU eviction (P2-1)
*
* Provides a central repository for all prepared statements in the
* database session, allowing them to be referenced by name.
* Implements bounded capacity with LRU eviction to prevent unbounded growth.
*/
class PreparedStorage {
// LRU cache implementation
struct LruEntry {
std::string name; ///< Query name (key)
PreparedQuery query; ///< The prepared query
mutable std::list<std::string>::iterator lru_iter; ///< Iterator in LRU list (mutable for get())
};
qb::unordered_map<std::string, LruEntry> _prepared_queries; ///< Map of queries
std::list<std::string> _lru_list; ///< LRU order list
size_t _max_size{100}; ///< Max capacity
size_t _evicted_count{0}; ///< Stats: evicted
public:
/**
* @brief Default constructor with default max size
*/
PreparedStorage() = default;
/**
* @brief Construct with custom max size
* @param max_size Maximum number of prepared queries to keep
*/
explicit PreparedStorage(size_t max_size)
: _max_size(max_size > 0 ? max_size : 100) {}
/**
* @brief Set maximum cache size (applies to future insertions)
* @param max_size New maximum size
*/
void set_max_size(size_t max_size);
/**
* @brief Get current maximum cache size
* @return size_t Max size
*/
size_t
max_size() const {
return _max_size;
}
/**
* @brief Get current number of cached queries
* @return size_t Current size
*/
size_t
size() const {
return _prepared_queries.size();
}
/**
* @brief Get total number of evicted queries (stats)
* @return size_t Eviction count
*/
size_t
evicted_count() const {
return _evicted_count;
}
/**
* @brief Checks if a prepared query exists
*
* @param name Name of the prepared query
* @return bool True if the query exists, false otherwise
*/
bool
has(std::string_view name) const {
return _prepared_queries.find(std::string(name)) != _prepared_queries.cend();
}
/**
* @brief Adds a prepared query to storage
*
* Implements LRU eviction if over capacity.
* Updates LRU order on access.
*
* @param query Prepared query to add
* @return const PreparedQuery& Reference to the stored query
*/
const PreparedQuery &push(PreparedQuery &&query);
/**
* @brief Retrieves a prepared query by name
*
* Updates LRU order on access (marks as recently used).
*
* @param name Name of the prepared query
* @return PreparedQuery const& Reference to the prepared query
* @throws std::out_of_range If the query doesn't exist
*/
PreparedQuery const &get(std::string_view name) const;
/**
* @brief Clear all prepared queries
*/
void
clear() {
_prepared_queries.clear();
_lru_list.clear();
}
private:
/**
* @brief Evict least recently used items if over capacity
*/
void evict_if_needed();
};
// Maintain backward compatibility
using PreparedQueryStorage = PreparedStorage;
/**
* @brief Class for managing query parameters
*
* Encapsulates parameters for prepared statements, handling
* type conversion and binary encoding according to PostgreSQL protocol.
*/
class QueryParams {
std::vector<byte> _params; ///< Serialized parameters
std::vector<integer> _param_types; ///< OIDs for parameter types
public:
/**
* @brief Constructs an empty parameter set
*/
QueryParams() = default;
/**
* @brief Constructs parameter set from variadic arguments
*
* Uses template argument deduction to convert various parameter types
* to their PostgreSQL binary representation.
*
* The single-argument case is constrained to exclude QueryParams itself so
* this forwarding constructor does not hijack the copy/move constructors
* (a `QueryParams b(a)` from a non-const lvalue would otherwise bind here and
* try to *serialize* `a` instead of copying it).
*
* @tparam T Parameter types
* @param args Parameter values
*/
template <typename... T,
std::enable_if_t<!(sizeof...(T) == 1 && std::conjunction_v<std::is_same<std::decay_t<T>, QueryParams>...>), int> = 0>
QueryParams(T &&...args) {
if constexpr (sizeof...(T) > 0) {
// Do not use format_codes_buffer, no longer used
// Serialize parameters directly
ParamSerializer serializer;
serializer.serialize_params(std::forward<T>(args)...);
// GCC -O2 emits a spurious -Warray-bounds on these small-vector copies
// (a known GCC-14 middle-end false positive; clang/MSVC are clean).
#if defined(__GNUC__) && !defined(__clang__)
#pragma GCC diagnostic push
#pragma GCC diagnostic ignored "-Warray-bounds"
#endif
_params = serializer.params_buffer();
_param_types = serializer.param_types();
#if defined(__GNUC__) && !defined(__clang__)
#pragma GCC diagnostic pop
#endif
}
}
/**
* @brief Gets the serialized parameters
*
* @return std::vector<byte>& Reference to the serialized parameters
*/
std::vector<byte> &
get() {
return _params;
}
/**
* @brief Gets the serialized parameters (const version)
*
* @return const std::vector<byte>& Const reference to the serialized parameters
*/
const std::vector<byte> &
get() const {
return _params;
}
/**
* @brief Gets the parameter types
*
* @return const std::vector<integer>& Const reference to parameter OIDs
*/
const std::vector<integer> &
param_types() const {
return _param_types;
}
/**
* @brief Gets the number of parameters
*
* @return smallint The number of parameters
*/
smallint param_count() const;
/**
* @brief Checks if the parameter set is empty
*
* @return bool True if there are no parameters, false otherwise
*/
bool
empty() const {
return _params.empty();
}
};
/**
* @brief Interface for SQL queries
*
* Base class for all SQL query implementations, providing a common
* interface for getting the query message and handling callbacks.
*/
class ISqlQuery {
public:
ISqlQuery() = default;
virtual ~ISqlQuery() = default;
/**
* @brief Checks if the query is valid
*
* @return bool True if the query is valid, false otherwise
*/
virtual bool
is_valid() const {
return true;
}
/**
* @brief Gets the PostgreSQL protocol message for the query
*
* @return message Message to send to the server
*/
virtual message get() const = 0;
/**
* @brief Called when the query succeeds
*/
virtual void on_success() const = 0;
/**
* @brief Called when the query fails
*
* @param err Error information
*/
virtual void on_error(error::db_error const &err) const = 0;
};
/**
* @brief Base implementation of SQL query with callbacks
*
* Provides a base implementation for SQL queries with success and error callbacks.
*
* @tparam CB_SUCCESS Type of success callback
* @tparam CB_ERROR Type of error callback
*/
template <typename CB_SUCCESS, typename CB_ERROR>
class SqlQuery : public ISqlQuery {
CB_SUCCESS _on_success; ///< Success callback
CB_ERROR _on_error; ///< Error callback
public:
/**
* @brief Constructs a SQL query with callbacks
*
* @param success Success callback
* @param error Error callback
*/
SqlQuery(CB_SUCCESS &&success, CB_ERROR &&error)
: _on_success(std::forward<CB_SUCCESS>(success))
, _on_error(std::forward<CB_ERROR>(error)) {}
virtual ~SqlQuery() = default;
/**
* @brief Calls the success callback
*/
void
on_success() const final {
_on_success();
}
/**
* @brief Calls the error callback
*
* @param err Error information
*/
void
on_error(error::db_error const &err) const final {
_on_error(err);
}
};
/**
* @brief Query for beginning a transaction
*
* Creates a BEGIN statement with optional transaction mode (isolation level, etc.).
*
* @tparam CB_SUCCESS Type of success callback
* @tparam CB_ERROR Type of error callback
*/
template <typename CB_SUCCESS, typename CB_ERROR>
class BeginQuery final : public SqlQuery<CB_SUCCESS, CB_ERROR> {
transaction_mode _mode; ///< Transaction mode
int _statement_timeout_ms{}; ///< If > 0, append SET LOCAL statement_timeout (ms)
public:
/**
* @brief Constructs a BEGIN query
*
* @param mode Transaction mode
* @param statement_timeout If positive, same round-trip runs
* `SET LOCAL statement_timeout = N` (milliseconds) after BEGIN
* @param success Success callback
* @param error Error callback
*/
BeginQuery(transaction_mode mode, qb::duration statement_timeout, CB_SUCCESS &&success, CB_ERROR &&error)
: SqlQuery<CB_SUCCESS, CB_ERROR>(std::forward<CB_SUCCESS>(success), std::forward<CB_ERROR>(error))
, _mode(mode)
, _statement_timeout_ms(statement_timeout > qb::duration::zero()
? static_cast<int>(std::chrono::duration_cast<std::chrono::milliseconds>(statement_timeout).count())
: 0) {}
/**
* @brief Creates the BEGIN message
*
* @return message Message to send to the server
*/
message
get() const final {
std::string sql = "BEGIN ";
sql += to_string(_mode);
if (_statement_timeout_ms > 0) {
sql += "; SET LOCAL statement_timeout = ";
sql += std::to_string(_statement_timeout_ms);
}
QB_LOG_DEBUG("[pgsql] Send BEGIN: \"" << sql << "\"");
message m(query_tag);
m.write(sql);
return m;
}
};
/**
* @brief Query for committing a transaction
*
* Creates a COMMIT statement to finalize a transaction.
*
* @tparam CB_SUCCESS Type of success callback
* @tparam CB_ERROR Type of error callback
*/
template <typename CB_SUCCESS, typename CB_ERROR>
class CommitQuery final : public SqlQuery<CB_SUCCESS, CB_ERROR> {
public:
/**
* @brief Constructs a COMMIT query
*
* @param success Success callback
* @param error Error callback
*/
CommitQuery(CB_SUCCESS &&success, CB_ERROR &&error)
: SqlQuery<CB_SUCCESS, CB_ERROR>(std::forward<CB_SUCCESS>(success), std::forward<CB_ERROR>(error)) {}
/**
* @brief Creates the COMMIT message
*
* @return message Message to send to the server
*/
message
get() const final {
QB_LOG_DEBUG("[pgsql] Send COMMIT");
message m(query_tag);
m.write("commit");
return m;
}
};
/**
* @brief Query for rolling back a transaction
*
* Creates a ROLLBACK statement to abort a transaction.
*
* @tparam CB_SUCCESS Type of success callback
* @tparam CB_ERROR Type of error callback
*/
template <typename CB_SUCCESS, typename CB_ERROR>
class RollbackQuery final : public SqlQuery<CB_SUCCESS, CB_ERROR> {
public:
/**
* @brief Constructs a ROLLBACK query
*
* @param success Success callback
* @param error Error callback
*/
RollbackQuery(CB_SUCCESS &&success, CB_ERROR &&error)
: SqlQuery<CB_SUCCESS, CB_ERROR>(std::forward<CB_SUCCESS>(success), std::forward<CB_ERROR>(error)) {}
/**
* @brief Creates the ROLLBACK message
*
* @return message Message to send to the server
*/
message
get() const final {
QB_LOG_DEBUG("[pgsql] Send ROLLBACK");
message m(query_tag);
m.write("rollback");
return m;
}
};
/**
* @brief Quote a value as a PostgreSQL SQL identifier: wrap in double quotes and double any
* embedded double quote (per SQL identifier syntax, matching libpq's PQescapeIdentifier).
*
* SECURITY: savepoint names are the ONE user-supplied value the module injects into a
* simple-query string (`savepoint <name>`) rather than binding out-of-band. Without quoting, a
* name like `s; DROP TABLE users; --` executes a second statement. Quoting turns any name into a
* single literal identifier, so injection is impossible regardless of which API (callback or
* co_await) supplied it. The co_await path additionally rejects non-alnum names early; this is
* the belt-and-suspenders that also covers the callback path.
*/
[[nodiscard]] inline std::string
pg_quote_identifier(std::string const &ident) {
std::string out;
out.reserve(ident.size() + 2);
out.push_back('"');
for (char const c : ident) {
if (c == '"')
out.push_back('"'); // double an embedded quote
out.push_back(c);
}
out.push_back('"');
return out;
}
/**
* @brief Query for creating a savepoint
*
* Creates a SAVEPOINT statement within a transaction.
*
* @tparam CB_SUCCESS Type of success callback
* @tparam CB_ERROR Type of error callback
*/
template <typename CB_SUCCESS, typename CB_ERROR>
class SavePointQuery final : public SqlQuery<CB_SUCCESS, CB_ERROR> {
std::string _name; ///< Savepoint name (owned copy — never a dangling reference)
public:
/**
* @brief Constructs a SAVEPOINT query
*
* @param name Savepoint name
* @param success Success callback
* @param error Error callback
*/
SavePointQuery(std::string const &name, CB_SUCCESS &&success, CB_ERROR &&error)
: SqlQuery<CB_SUCCESS, CB_ERROR>(std::forward<CB_SUCCESS>(success), std::forward<CB_ERROR>(error))
, _name(name) {}
/**
* @brief Creates the SAVEPOINT message
*
* @return message Message to send to the server
*/
message
get() const final {
QB_LOG_DEBUG("[pgsql] Send SAVEPOINT " << _name);
message m(query_tag);
m.write("savepoint " + pg_quote_identifier(_name));
return m;
}
};
/**
* @brief Query for releasing a savepoint
*
* Creates a RELEASE SAVEPOINT statement within a transaction.
*
* @tparam CB_SUCCESS Type of success callback
* @tparam CB_ERROR Type of error callback
*/
template <typename CB_SUCCESS, typename CB_ERROR>
class ReleaseSavePointQuery final : public SqlQuery<CB_SUCCESS, CB_ERROR> {
std::string _name; ///< Savepoint name (owned copy — never a dangling reference)
public:
/**
* @brief Constructs a RELEASE SAVEPOINT query
*
* @param name Savepoint name
* @param success Success callback
* @param error Error callback
*/
ReleaseSavePointQuery(std::string const &name, CB_SUCCESS &&success, CB_ERROR &&error)
: SqlQuery<CB_SUCCESS, CB_ERROR>(std::forward<CB_SUCCESS>(success), std::forward<CB_ERROR>(error))
, _name(name) {}
/**
* @brief Creates the RELEASE SAVEPOINT message
*
* @return message Message to send to the server
*/
message
get() const final {
QB_LOG_DEBUG("[pgsql] Send RELEASE SAVEPOINT " << _name);
message m(query_tag);
m.write("release savepoint " + pg_quote_identifier(_name));
return m;
}
};
/**
* @brief Query for rolling back to a savepoint
*
* Creates a ROLLBACK TO SAVEPOINT statement within a transaction.
*
* @tparam CB_SUCCESS Type of success callback
* @tparam CB_ERROR Type of error callback
*/
template <typename CB_SUCCESS, typename CB_ERROR>
class RollbackSavePointQuery final : public SqlQuery<CB_SUCCESS, CB_ERROR> {
std::string _name; ///< Savepoint name (owned copy — never a dangling reference)
public:
/**
* @brief Constructs a ROLLBACK TO SAVEPOINT query
*
* @param name Savepoint name
* @param success Success callback
* @param error Error callback
*/
RollbackSavePointQuery(std::string const &name, CB_SUCCESS &&success, CB_ERROR &&error)
: SqlQuery<CB_SUCCESS, CB_ERROR>(std::forward<CB_SUCCESS>(success), std::forward<CB_ERROR>(error))
, _name(name) {}
/**
* @brief Creates the ROLLBACK TO SAVEPOINT message
*
* @return message Message to send to the server
*/
message
get() const final {
QB_LOG_DEBUG("[pgsql] Send ROLLBACK TO SAVEPOINT " << _name);
message m(query_tag);
m.write("rollback to savepoint " + pg_quote_identifier(_name));
return m;
}
};
/**
* @brief Query for executing a simple SQL statement
*
* Creates a query for direct execution of SQL expressions.
*
* @tparam CB_SUCCESS Type of success callback
* @tparam CB_ERROR Type of error callback
*/
template <typename CB_SUCCESS, typename CB_ERROR>
class SimpleQuery final : public SqlQuery<CB_SUCCESS, CB_ERROR> {
const std::string _expression; ///< SQL expression
public:
/**
* @brief Constructs a simple query
*
* @param expr SQL expression
* @param success Success callback
* @param error Error callback
*/
SimpleQuery(std::string &&expr, CB_SUCCESS &&success, CB_ERROR &&error)
: SqlQuery<CB_SUCCESS, CB_ERROR>(std::forward<CB_SUCCESS>(success), std::forward<CB_ERROR>(error))
, _expression(std::move(expr)) {}
/**
* @brief Creates the query message
*
* @return message Message to send to the server
*/
message
get() const final {
QB_LOG_DEBUG("[pgsql] Send QUERY \"" << _expression << "\"");
message m(query_tag);
m.write(_expression);
return m;
}
};
/**
* @brief Query for preparing a statement
*
* Creates a query for preparing a named statement according to the PostgreSQL protocol.
*
* @tparam CB_SUCCESS Type of success callback
* @tparam CB_ERROR Type of error callback
*/
template <typename CB_SUCCESS, typename CB_ERROR>
class ParseQuery final : public SqlQuery<CB_SUCCESS, CB_ERROR> {
PreparedQuery const &_query; ///< Prepared query definition
public:
ParseQuery(PreparedQuery const &query, CB_SUCCESS &&success, CB_ERROR &&error)
: SqlQuery<CB_SUCCESS, CB_ERROR>(std::forward<CB_SUCCESS>(success), std::forward<CB_ERROR>(error))
, _query(query) {}
bool
is_valid() const final {
// The Parse message encodes the parameter-type count as an int16, but get()
// writes every OID entry. More than 32767 declared types would truncate the
// count while still emitting all entries, desynchronizing the wire stream.
// Reject here (like ExecuteQuery's missing-statement check) so the failure goes
// through the normal on_error path instead of corrupting the connection. This is
// the Parse-side twin of the Bind guard ParamSerializer::ensure_param_count_fits().
if (qb::likely(_query.param_types.size() <= static_cast<std::size_t>(std::numeric_limits<smallint>::max())))
return true;
QB_LOG_CRIT("[pgsql] PARSE rejected: " << _query.param_types.size() << " parameter types exceed protocol max 32767");
return false;
}
[[nodiscard]] message
get() const final {
QB_LOG_DEBUG("[pgsql] Send PARSE QUERY \"" << _query.expression << "\"");
message cmd(parse_tag);
cmd.write(_query.name);
cmd.write(_query.expression);
cmd.write((smallint) _query.param_types.size());
for (auto oid_val : _query.param_types) {
cmd.write(static_cast<integer>(oid_val));
}
message describe(describe_tag);
describe.write('S');
describe.write(_query.name);
cmd.pack(describe);
cmd.pack(message(sync_tag));
return cmd;
}
};
/**
* @brief Prepared statement execution
*
* @tparam CB_SUCCESS Type of success callback
* @tparam CB_ERROR Type of error callback
*/
template <typename CB_SUCCESS, typename CB_ERROR>
class ExecuteQuery final : public SqlQuery<CB_SUCCESS, CB_ERROR> {
const PreparedStorage &_storage; ///< Prepared statement storage
std::string _query_name; ///< Query name to execute
QueryParams _params; ///< Query parameters
public:
/**
* @brief Constructs an execute query
*
* @param storage Prepared statement storage
* @param query_name Query name to execute
* @param params Query parameters
* @param success Success callback
* @param error Error callback
*/
ExecuteQuery(const PreparedStorage &storage, std::string_view query_name, QueryParams &¶ms, CB_SUCCESS &&success, CB_ERROR &&error)
: SqlQuery<CB_SUCCESS, CB_ERROR>(std::forward<CB_SUCCESS>(success), std::forward<CB_ERROR>(error))
, _storage(storage)
, _query_name(query_name)
, _params(std::move(params)) {}
bool
is_valid() const final {
if (qb::likely(_storage.has(_query_name)))
return true;
QB_LOG_CRIT("[pgsql] Error prepared query " << _query_name << " not registered");
return false;
}
message
get() const final {
const auto &query = _storage.get(_query_name);
message cmd(bind_tag);
// Exact format expected by PostgreSQL for a Bind message:
// 1. Portal name (empty = unnamed)
cmd.write("");
// 2. Prepared statement name
cmd.write(query.name);
// 3. Parameter format codes (PostgreSQL Bind message)
// 0 = no parameters or all default (text); 1 = single code applies to every
// parameter. We send binary (1) for all parameters when there are any.
const smallint param_count = _params.param_count();
if (param_count == 0) {
cmd.write(static_cast<smallint>(0));
} else {
cmd.write(static_cast<smallint>(1));
cmd.write(static_cast<smallint>(1)); // binary
}
// 4. Total number of parameters
cmd.write(param_count);
// 5. Parameter values
if (!_params.empty() && param_count > 0) {
// Skip the count in the parameters buffer
const std::vector<byte> ¶m_buffer = _params.get();
if (param_buffer.size() > sizeof(smallint)) {
const byte *data = param_buffer.data() + sizeof(smallint);
size_t data_size = param_buffer.size() - sizeof(smallint);
if (data_size == 0) {
QB_LOG_WARN("[pgsql] Bind: param_count=" << param_count << " but serialized payload is empty");
}
// Copy the raw data
auto out = cmd.output();
std::copy(data, data + data_size, out);
} else {
QB_LOG_WARN("[pgsql] Bind: param_count=" << param_count << " but parameter buffer missing payload");
}
}
// 6. Result-column format codes (one per column, or 0 = default all text).
// Binary for scalars; text for string-like OIDs so DataRow matches decoders.
const auto &rd = query.row_description;
const smallint ncol = static_cast<smallint>(rd.size());
cmd.write(ncol);
for (auto const &fd : rd) {
const smallint fmt = type_oid_prefers_binary_result_format(fd.type_oid) ? 1 : 0;
cmd.write(fmt);
}
// 7. Execute message (empty portal, no row limit)
message execute(execute_tag);
execute.write("");
execute.write(0);
cmd.pack(execute);
// 8. Sync message
cmd.pack(message(sync_tag));
return cmd;
}
};
} // namespace qb::pg::detail