-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathpgsql.h
More file actions
2893 lines (2671 loc) · 125 KB
/
Copy pathpgsql.h
File metadata and controls
2893 lines (2671 loc) · 125 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
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
/**
* @file pgsql.h
* @brief PostgreSQL client for the QB Actor Framework
*
* This file implements an asynchronous PostgreSQL client integrated with the QB Actor
* Framework. It provides a non-blocking interface for database operations such as:
*
* - Connection management to PostgreSQL databases
* - Transaction management (begin, commit, rollback)
* - Support for savepoints within transactions
* - Simple and prepared statement execution with parameter binding
* - Efficient query result retrieval and processing
* - Support for multiple authentication methods (MD5, SCRAM-SHA-256, etc.)
*
* The implementation is designed to work with the actor model, allowing
* database operations to be performed without blocking actor threads. The client
* fully implements the PostgreSQL wire protocol for efficient communication.
*
* Connection / query API (single-threaded qb-io: one event loop + coroutine scheduler
* per thread; never block the loop inside a callback or coroutine except via
* explicit suspension).
*
* Two orthogonal styles — same method names, different completion model:
*
* - **Coroutine completion:** overloads **without** user callbacks return
* `pg_reply_awaiter<T>`. Use **only** inside a coroutine: `auto r = co_await db.query("…");`
* or `co_await db.execute("…")`, `co_await db.prepare(…)`, `co_await db.begin()`, etc.
* From synchronous code, drive the coroutine with `qb::io::async::run_sync` (see
* `qb/io/async/coroutine/utils.h`) or spawn a `task` on `coro_scheduler()`. **Do not** call a
* blocking wait on the awaiter itself — there is no
* `.await()` on `pg_reply_awaiter`.
*
* - **Callback + synchronous drain:** pass success/error callbacks; overloads return
* `Transaction&` for fluent chaining. To run queued work to completion on the current
* thread, call **`Transaction::await()`** (or `qb::pg::await(db)`). For SQL/prepared
* ops when you have no real handler, use the constexpr discards:
* `execute(sql, discard_query, discard_error)`,
* `prepare(name, sql, types, discard_prepare, discard_error)`,
* `execute(name, params, discard_query, discard_error)`.
*
* - **Connection:** `co_await db.connect()` or `run_sync(db.connect())` — see
* `qb/io/async/coroutine/utils.h`.
*
* - **Coroutine transaction scope:** `co_await with_transaction(db, [](Transaction &tr) -> task<int>
* { ... })` runs `BEGIN`, awaits your `task` body (use `tr` / `db` for `co_await tr.execute` /
* `query`), then `COMMIT` on success or `ROLLBACK` on `transaction_abort`, `commit` failure, or a
* C++ exception. Overload `with_transaction(db, transaction_mode{...}, f)` sets isolation /
* read-only / deferrable. When `!reply.ok()` after an operation, throw
* `transaction_abort{reply.error()}` so the helper rolls back and returns `Reply::failure` instead
* of calling `COMMIT` on an aborted transaction.
*
* **Large-project conventions**
*
* - **One style per call stack:** In a `begin` success callback, use only callback overloads
* (`execute(..., cb, err)` or discards) and `Transaction::await()` — do not mix with discarded
* `execute("…")` coroutine awaiters (they are not driven there). In coroutines, use only
* `co_await` overloads and `with_transaction` / manual `begin` / `commit` / `rollback`.
* - **`Transaction&` vs `database&`:** `tcp::database` *is-a* `Transaction`; `with_transaction` and
* `co_await tr.query` use the same connection. Prefer passing `Transaction&` in helpers so code
* works with any concrete client type.
* - **Avoid nesting `with_transaction`:** It issues a second `BEGIN` on the same connection;
* behavior depends on the server (some configurations reject it and abort the block — use
* `transaction_abort{inner.error()}` and never `COMMIT` an aborted transaction; others may accept
* the pattern). Prefer a single scope plus `savepoint` / `release_savepoint` /
* `rollback_savepoint` for nested units of work.
* - **READ ONLY:** PostgreSQL still allows writes to **temporary** tables in a read-only
* transaction; only non-temporary relations are restricted.
* - **Errors:** For coroutine bodies, treat `Reply::ok()` as mandatory; on failure either throw
* `transaction_abort` (handled scope) or let exceptions propagate (rollback + rethrow).
*
* Key features:
* - Asynchronous I/O using the QB Actor Framework
* - Support for both plain TCP and SSL/TLS connections
* - Comprehensive transaction management
* - Prepared statement caching for performance
* - Detailed error reporting and handling
*
* @see qb::pg::detail::Database
* @see qb::pg::detail::Transaction
*
* @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 <chrono>
#include <coroutine>
#include <cstdint>
#include <functional>
#include <memory>
#include <optional>
#include <string>
#include <string_view>
#include <type_traits>
#include <qb/io/async.h>
#include <qb/io/async/tcp/connector.h>
#include <qb/io/crypto.h>
#ifdef QB_HAS_SSL
#include <qb/io/tcp/ssl/socket.h>
#endif
#include <qb/system/allocator/pipe.h>
#include <qb/system/cpu.h> // qb::scope_guard
#include <qb/system/parse.h> // qb::to_number (locale-free, non-throwing)
// P1-1: Socket includes for keepalive support
#ifdef _WIN32
#include <winsock2.h>
#include <ws2tcpip.h>
#else
#include <netinet/in.h>
#include <netinet/tcp.h>
#include <sys/socket.h>
#endif
#include "./commands.h"
#include "./pg_reply.h"
#include "./transaction.h"
/**
* @brief Maximum length for attribute names in PostgreSQL protocol
*
* Defines the maximum length in bytes for attribute names when parsing
* PostgreSQL protocol messages. This limit helps prevent buffer overflow
* attacks and ensures efficient memory usage.
*/
constexpr const uint32_t ATTRIBUTE_NAME_MAX = 1024; // 1 KB
/**
* @brief Maximum length for attribute values in PostgreSQL protocol
*
* Defines the maximum length in bytes for attribute values when parsing
* PostgreSQL protocol messages. This larger limit accommodates typical
* PostgreSQL data values while preventing excessively large allocations.
*/
constexpr const uint32_t ATTRIBUTE_VALUE_MAX = 1024 * 1024; // 1 MB
/**
* @brief Checks if a character is a control character
*
* Used during attribute parsing to validate input and ensure security.
* Control characters are generally not allowed in attribute names or values
* as they may indicate malformed or malicious input.
*
* @param c Character to check
* @return true if the character is a control character (ASCII 0-31 or 127)
*/
inline bool
is_control(int c) {
return ((c >= 0 && c <= 31) || c == 127);
}
/**
* @brief Parses header attributes from a PostgreSQL protocol message
*
* Parses a buffer of header attributes into a case-insensitive map.
* This function is primarily used during SCRAM authentication to process
* challenge-response data between client and server.
*
* Features:
* - Supports both quoted and unquoted attribute values
* - Handles attribute separators (comma and semicolon)
* - Enforces size limits to prevent buffer overflows
* - Validates input to reject control characters
* - Properly handles whitespace according to the PostgreSQL protocol
*
* @param ptr Pointer to the buffer containing attributes
* @param len Length of the buffer
* @return Case-insensitive map of attribute names to values
* @throws std::runtime_error If parsing fails due to control characters or exceeding
* size limits
*/
qb::icase_unordered_map<std::string> parse_header_attributes(const char *ptr, const size_t len);
namespace qb::protocol {
/**
* @brief PostgreSQL protocol implementation for the QB actor framework
*
* Handles the message framing and parsing according to the PostgreSQL
* wire protocol specification. This class is responsible for:
*
* - Extracting complete messages from the input stream
* - Managing protocol state between messages
* - Forwarding complete messages to the appropriate handlers
* - Implementing the PostgreSQL message format requirements
*
* The protocol handler processes the incoming byte stream and constructs
* well-formed PostgreSQL protocol messages. It maintains internal state
* to handle partial messages that arrive in multiple network packets.
*
* @tparam IO_ I/O handler type that provides input/output stream access
*/
template <typename IO_>
class pgsql final : public qb::io::async::AProtocol<IO_> {
public:
/**
* @brief PostgreSQL protocol message type
*
* Represents a complete PostgreSQL protocol message including
* message type, length, and payload data.
*/
using message = std::unique_ptr<pg::detail::message>;
private:
message message_; ///< Current message being processed
std::size_t offset_ = 0; ///< Current offset in the input buffer
public:
pgsql() = delete;
/**
* @brief Constructs a PostgreSQL protocol handler
*
* Initializes the protocol handler with a reference to the I/O
* subsystem that provides access to input and output streams.
*
* @param io Reference to the I/O handler
*/
explicit pgsql(IO_ &io) noexcept
: qb::io::async::AProtocol<IO_>(io) {}
/**
* @brief Copy data from input iterator to output iterator
*
* Helper method to copy data between iterators with a maximum limit.
* Used internally for buffer management and message construction.
*
* @tparam InputIter Input iterator type
* @tparam OutputIter Output iterator type
* @param in Start of input range
* @param end End of input range
* @param max Maximum number of items to copy
* @param out Output iterator
* @return InputIter Iterator after the last copied element
*/
template <typename InputIter, typename OutputIter>
InputIter
copy(InputIter in, InputIter end, size_t max, OutputIter out) {
for (size_t i = 0; i < max && in != end; ++i) {
*out++ = *in++;
}
return in;
}
/**
* @brief Calculate the size of a complete PostgreSQL message
*
* Inspects the input buffer to determine if a complete message is available.
* This method implements the PostgreSQL message framing protocol by:
*
* 1. Reading the message type byte and length field
* 2. Creating a new message object if needed
* 3. Reading message payload data up to the expected length
* 4. Determining if the message is complete
*
* If a message is complete, returns its size in bytes. If incomplete,
* returns 0 to indicate more data is needed from the network.
*
* @return std::size_t Size of the complete message, or 0 if incomplete
*/
std::size_t
getMessageSize() noexcept final {
constexpr const size_t header_size = sizeof(qb::pg::integer) + sizeof(qb::pg::byte);
const auto &in = this->_io.in();
if (in.size() < offset_ + header_size)
return 0; // read more
auto max_bytes = in.size() - offset_;
if (!message_) {
message_ = std::make_unique<pg::detail::message>();
// OPTIMIZED: Use std::copy_n for batch copy instead of byte-by-byte
// This provides ~10x performance improvement for large messages
auto header_begin = in.begin();
auto out = message_->output();
std::copy_n(header_begin, header_size, out);
offset_ += header_size;
max_bytes -= header_size;
const qb::pg::uinteger wire_len = static_cast<qb::pg::uinteger>(message_->length());
if (wire_len < 4u || wire_len > qb::pg::PG_PROTOCOL_MAX_MESSAGE_BYTES) {
QB_LOG_CRIT("[pgsql] Invalid wire message length " << wire_len << " (must be 4.." << qb::pg::PG_PROTOCOL_MAX_MESSAGE_BYTES
<< "); dropping connection");
message_.reset();
offset_ = 0;
// Mark the protocol invalid (the documented contract for a malformed size
// field): the I/O layer then disposes and fires event::disconnected, whose
// handler fails every pending query and resumes any pending connect awaiter.
// The previous prepare_reconnect() tore the transport down synchronously from
// inside the read handler WITHOUT firing on(disconnected) — so queued query
// awaiters (and a pending co_await connect(), whose handle it cleared without
// resuming) hung forever, and it left the read watcher on a closed fd.
this->not_ok();
return 0;
}
}
if (message_->length() > message_->size()) {
// Read the message body
auto out = message_->output();
const std::size_t to_copy = std::min(message_->length() - message_->size(), max_bytes);
// OPTIMIZED: Use std::copy_n for batch copy instead of byte-by-byte
auto data_begin = in.begin() + offset_;
std::copy_n(data_begin, to_copy, out);
offset_ += to_copy;
}
if (message_->length() == message_->size()) {
return message_->buffer_size();
}
return 0;
}
/**
* @brief Handle a complete PostgreSQL message
*
* Called by the protocol framework when a complete message has been
* received and parsed according to the PostgreSQL protocol rules.
*
* This method:
* 1. Validates the connection state
* 2. Resets the message read pointer
* 3. Forwards the complete message to the I/O handler for processing
* 4. Resets the protocol state to prepare for the next message
*
* @param size Size of the message (unused in this implementation)
*/
void
onMessage(std::size_t) noexcept final {
if (!this->ok())
return;
message_->reset_read();
// This is a noexcept boundary: a message handler that throws would call
// std::terminate. A hostile or misconfigured server can make on_authentication
// throw before the connection is established (an unsupported auth method falls
// into its `default:` throw; a malformed SCRAM server message makes
// parse_header_attributes / the iteration-count parse / PBKDF2 throw). Contain any
// handler exception and mark the protocol invalid so the I/O layer disposes and
// fires event::disconnected, whose handler fails pending queries and resumes a
// pending connect awaiter with an error instead of crashing the process.
try {
this->_io.on(std::move(message_));
} catch (std::exception const &e) {
QB_LOG_CRIT("[pgsql] exception in message handler, dropping connection: " << e.what());
this->not_ok();
} catch (...) {
QB_LOG_CRIT("[pgsql] unknown exception in message handler, dropping connection");
this->not_ok();
}
reset();
}
/**
* @brief Reset the protocol state
*
* Prepares the protocol handler for the next message by resetting
* internal state variables. This ensures that each new message is
* processed from a clean initial state.
*/
void
reset() noexcept final {
offset_ = 0;
}
};
} // namespace qb::protocol
namespace qb::pg {
/**
* @brief Asynchronous payload from PostgreSQL NOTIFY (after LISTEN on the same or another session).
*
* Delivered on the I/O thread when a `NotificationResponse` is received; see
* `tcp::notify_cb_consumer` / `tcp::notify_co_consumer`, or `database::on_incoming_notify`.
*/
struct notification {
int server_backend_pid{};
std::string channel;
std::string payload;
};
namespace detail {
using namespace qb::io;
using namespace qb::pg;
/**
* @brief Validate the SCRAM server nonce against the client nonce (RFC 5802 §5.1).
*
* In SCRAM-SHA-256 the server echoes the client's nonce and appends its own, so the
* combined nonce in the server-first message MUST begin with the exact nonce the
* client sent AND be strictly longer (the server has to contribute entropy). A
* server — or a man-in-the-middle — that fails this is not replaying our own first
* message faithfully, so the exchange must be aborted before deriving any proof.
*
* @param client_nonce The nonce this client generated and sent (`r=` in client-first).
* @param server_nonce The combined nonce returned by the server (`r=` in server-first).
* @return true iff server_nonce starts with client_nonce and is longer.
*/
[[nodiscard]] inline bool
scram_server_nonce_extends_client(std::string_view client_nonce, std::string_view server_nonce) noexcept {
return !client_nonce.empty() && server_nonce.size() > client_nonce.size() && server_nonce.starts_with(client_nonce);
}
/// Hard upper bound on the SCRAM-SHA-256 iteration count this client will honour.
/// PostgreSQL's server default is 4096; this is ~244x that. The count is
/// server-controlled and feeds PBKDF2 SYNCHRONOUSLY on the I/O event-loop thread,
/// so an unbounded value (e.g. i=2147483647 => ~minutes of HMAC) would stall every
/// actor on the core — a trivial denial of service from a hostile or MITM'd server.
inline constexpr int kMaxScramIterations = 1'000'000;
/**
* @brief Parse and bound-check a SCRAM server-first `i=` iteration count.
*
* Accepts only a canonical base-10 integer in [1, ::qb::pg::kMaxScramIterations].
* Rejects a missing/empty, non-numeric, overflowing, non-positive, or absurdly
* large value — the last being the denial-of-service guard documented on
* ::qb::pg::kMaxScramIterations.
*
* @param text The raw `i=` attribute value from the SCRAM server-first message.
* @return The validated iteration count.
* @throws error::connection_error if @p text is malformed or out of range.
*/
[[nodiscard]] inline int
scram_validate_iteration_count(std::string_view text) {
const auto parsed = qb::to_number<int>(text);
if (!parsed || *parsed < 1 || *parsed > kMaxScramIterations) {
throw error::connection_error("SCRAM iteration count missing, malformed, or out of range");
}
return *parsed;
}
/**
* @brief Escape a SCRAM `saslname` (RFC 5802): `=` -> `=3D`, `,` -> `=2C`.
*
* The `n=<username>` field in the SCRAM client-first message uses the `saslname`
* production, where `,` and `=` must be percent-style escaped (and `=3D` must be
* applied before `=2C` so the inserted `=` is not re-escaped). PostgreSQL ignores
* the SCRAM `n=` (it authenticates the startup-packet user), but a role name
* containing a comma/equals would otherwise emit a malformed client-first message
* that a strict RFC-5802 parser (proxy/pooler/non-PG server) rejects.
*/
[[nodiscard]] std::string scram_escape_saslname(std::string_view name);
/**
* @brief Opportunistic-TLS (STARTTLS) negotiator for the PostgreSQL protocol.
*
* Plugs into `qb::io::async::tcp::starttls_connect`. After the cleartext TCP connect,
* it sends the 8-byte **SSLRequest** packet (int32 length = 8, int32 request code =
* 80877103 / `0x04D2162F`, big-endian) and reads the single-byte server reply:
* - `'S'` → upgrade: the connector then performs the TLS handshake asynchronously.
* - anything else (`'N'`, EOF, error) → fail: a secure database **requires** TLS
* (use the plain `tcp::database` for cleartext). This drops the old, broken `'N'`
* fallback that produced an unusable handle-less `ssl::socket`.
*
* All I/O is non-blocking; the connector drives the readiness events.
*/
struct postgres_ssl_negotiator {
static constexpr bool enabled = true;
/**
* @brief Build the 8-byte PostgreSQL SSLRequest packet (big-endian).
*
* Layout: int32 length = 8, int32 request code = 80877103 (`0x04D2162F`).
*
* @return The serialized SSLRequest bytes ready to write on the cleartext socket.
*/
static std::array<std::uint8_t, 8> make_request() noexcept;
std::array<std::uint8_t, 8> request_{make_request()};
std::size_t written_{0};
std::uint8_t verdict_{0};
bool got_verdict_{false};
/**
* @brief Drive one step of the STARTTLS negotiation on a ready socket.
*
* Non-blocking state machine called by the connector on each readiness event:
* 1. write the remaining SSLRequest bytes,
* 2. read the single-byte server verdict,
* 3. decide: `'S'` -> upgrade to TLS, anything else / EOF / error -> fail
* (a secure database requires TLS).
*
* @param sock The cleartext socket connected to the server (not yet upgraded).
* @param revents Readiness flags from the event loop (unused; the phase is tracked
* internally).
* @return The next `starttls_action`: `want_write` / `want_read` to be polled again,
* `upgrade` to start the TLS handshake, or `fail` to abort the connect.
*/
qb::io::async::tcp::starttls_action advance(qb::io::tcp::socket &sock, int revents) noexcept;
};
/**
* @brief PostgreSQL database client implementation
*
* Core implementation of the PostgreSQL client that handles connection
* establishment, authentication, and query execution. This class provides
* the foundation for asynchronous database operations with PostgreSQL.
*
* Key features:
* - Asynchronous TCP/IP connection management
* - Multiple authentication methods support (Cleartext, MD5, SCRAM-SHA-256)
* - Transaction management (inherited from Transaction class)
* - Query execution and result processing
* - Prepared statement caching and execution
* - Event-driven message handling
*
* The Database class inherits from both the TCP client base class for network
* connectivity and the Transaction class for query and transaction management.
*
* @tparam QB_IO_ I/O handler type that provides networking capabilities
* @tparam NotifyDerived CRTP notify consumer type (`void` for plain `database`); receives NOTIFY via
* `consume_pg_notify` / `deliver_pg_notify`. See `notify_consumer` / `notify_co_consumer`
* (`notify_cb_consumer` is an alias for the same class).
*/
template <typename QB_IO_, typename NotifyDerived = void>
class Database
: public qb::io::async::tcp::client<Database<QB_IO_, NotifyDerived>, QB_IO_, void>
, public Transaction {
public:
/**
* @brief PostgreSQL protocol handler type
*
* Type alias for the protocol handler used by this database client.
*/
using pg_protocol = qb::protocol::pgsql<Database<QB_IO_, NotifyDerived>>;
private:
connection_options conn_opts_; ///< Database connection options
client_options_type client_opts_; ///< Client-supplied startup options (sent in the StartupMessage)
client_options_type server_params_; ///< Server-reported ParameterStatus cache (NOT echoed back at connect)
integer serverPid_{}; ///< Server process ID
integer serverSecret_{}; ///< Server secret for protocol operations
PreparedQueryStorage storage_; ///< Storage for prepared statements
bool is_connected_ = false; ///< Flag indicating if the connection is established
/// Outstanding `co_await connect()` handshake (coroutine resume + validity token)
bool connect_coroutine_pending_{false};
bool connect_handshake_failed_{false};
std::coroutine_handle<> connect_suspend_handle_{};
std::shared_ptr<bool> connect_suspend_valid_{};
/// Bumps on each new handshake / reconnect prep so stale `callback(timeout)` ignores
std::uint64_t connect_timer_generation_{0};
/// Owned handshake-deadline timer. MUST be a ScopedTimeout (cancelled on destruction),
/// not a fire-and-forget async::callback: the latter is a self-deleting heap Timeout that
/// outlives this Database and would dereference a freed `this` if the connection is dropped
/// within the timeout window (the common case — handshake finishes in ms, the deadline is
/// seconds). Destroying this member with the Database stops the watcher.
std::unique_ptr<qb::io::async::ScopedTimeout<std::function<void()>>> connect_deadline_timer_{};
/// When `NotifyDerived` is `void`, optional handler for `NotificationResponse` (plain
/// `database`).
std::function<void(::qb::pg::notification &&)> inbound_notify_handler_{};
void
try_resume_connect_wait() {
if (!connect_coroutine_pending_)
return;
const bool terminal = is_connected_ || connect_handshake_failed_ || (has_error() && !is_connected_);
if (!terminal)
return;
auto h = connect_suspend_handle_;
auto v = connect_suspend_valid_;
connect_coroutine_pending_ = false;
connect_suspend_handle_ = {};
connect_suspend_valid_.reset();
if (v && !*v)
return;
if (h)
qb::io::async::coro_scheduler().schedule_resume(h);
}
#ifndef QB_HAS_SSL
/**
* @brief Refuse a server-selected authentication method this build cannot perform.
*
* MD5 and SCRAM-SHA-256 both need real cryptography (MD5 / PBKDF2-HMAC-SHA256 /
* HMAC / CSPRNG nonce), which only OpenSSL provides. A `QB_WITH_SSL=OFF` build has
* none of it, so those code paths are not compiled at all — but the *server* picks
* the method, so the mismatch is only knowable at runtime, on receipt of the
* Authentication message. Fail here, loudly and specifically, rather than sending a
* bogus or empty password and letting the server report a generic auth failure.
*
* Uses the same terminal-failure shape as the SCRAM MITM refusal above (QB_LOG_CRIT +
* flags + resume): this runs inside the protocol message handler on the event-loop
* thread, so an escaping exception would tear down the loop rather than the
* connection.
*
* @param method Human-readable name of the method the server asked for.
*/
void
refuse_auth_without_ssl(const char *method) {
QB_LOG_CRIT("[pgsql] Server requested " << method
<< " authentication, which this build cannot perform: qb was built without OpenSSL "
"(QB_HAS_SSL undefined / QB_WITH_SSL=OFF). Rebuild qb with OpenSSL, or configure the "
"PostgreSQL role for 'trust' or 'password' (cleartext) authentication.");
connect_handshake_failed_ = true;
is_connected_ = false;
try_resume_connect_wait();
}
#endif // !QB_HAS_SSL
#ifdef QB_HAS_SSL
/**
* @brief The value-semantic client TLS context the connection options describe -- built the
* same way for the session and for an out-of-band `cancel_async()` (Huly QB-113):
* default -> Context::client() (TLS 1.2+, system trust store, verify chain + host);
* ssl_verify=none -> verification off (encrypt only);
* ssl_root_cert -> trust a private CA IN ADDITION to the system store (libpq sslrootcert);
* ssl_cert+ssl_key -> present a client certificate (mutual TLS; libpq sslcert/sslkey).
* A bad CA/cert/key path leaves `ok()` false: callers fail CLOSED on it.
*/
[[nodiscard]] qb::io::ssl::Context
make_client_tls_context() const {
auto tls = qb::io::ssl::Context::client();
if (conn_opts_.ssl_verify != qb::pg::ssl_verify_mode::full)
tls.verify(qb::io::ssl::VerifyMode::none);
if (!conn_opts_.ssl_root_cert.empty())
tls.trust(conn_opts_.ssl_root_cert);
if (!conn_opts_.ssl_cert.empty() && !conn_opts_.ssl_key.empty())
tls.identity(conn_opts_.ssl_cert, conn_opts_.ssl_key);
return tls;
}
#endif // QB_HAS_SSL
/**
* @brief The 16-byte PostgreSQL CancelRequest for this session's backend, all big-endian:
* int32 length = 16, int32 request code = 80877102 (0x04D2162E), int32 backend process
* id, int32 backend secret key -- the two captured from BackendKeyData at connect.
*/
[[nodiscard]] std::array<std::uint8_t, 16>
cancel_request_packet() const noexcept {
std::array<std::uint8_t, 16> pkt{};
const std::uint32_t len = htonl(16u);
const std::uint32_t code = htonl(80877102u);
const std::uint32_t pid = htonl(static_cast<std::uint32_t>(serverPid_));
const std::uint32_t key = htonl(static_cast<std::uint32_t>(serverSecret_));
std::memcpy(pkt.data() + 0, &len, 4);
std::memcpy(pkt.data() + 4, &code, 4);
std::memcpy(pkt.data() + 8, &pid, 4);
std::memcpy(pkt.data() + 12, &key, 4);
return pkt;
}
/**
* @brief The connect budget of an out-of-band cancel: `connect_timeout` (10 s when unset),
* capped at 2 s. The cancel targets the SAME already-reachable endpoint as the live
* session, so the connect is normally sub-millisecond; the cap only bounds the
* pathological unreachable case -- the whole loop for `cancel()`, one coroutine for
* `cancel_async()`.
*/
[[nodiscard]] qb::duration
cancel_connect_budget() const noexcept {
const qb::duration cfg = conn_opts_.connect_timeout > qb::duration::zero()
? conn_opts_.connect_timeout
: std::chrono::duration_cast<qb::duration>(std::chrono::seconds(10));
return std::min(cfg, std::chrono::duration_cast<qb::duration>(std::chrono::seconds(2)));
}
/**
* @brief Starts outbound TCP using the async framework (`qb::io::async::tcp::connect`).
*
* Uses the same connector path as Redis: non-blocking `n_connect`, `EV_WRITE` completion,
* and an optional deadline (`connect_timeout` or default 10s). The coroutine awaiter is
* resumed from `try_resume_connect_wait()` after TCP + PostgreSQL pre-startup steps or on
* failure.
*
* @param h Coroutine handle to resume when the handshake attempt finishes (success or failure)
* @param valid Shared flag cleared when the awaiter is destroyed (ignore stale callbacks)
* @param timeout_override If positive, overrides `conn_opts_.connect_timeout` for this attempt
*/
void
start_connect_from_awaiter(std::coroutine_handle<> h, std::shared_ptr<bool> valid, qb::duration timeout_override) {
++connect_timer_generation_;
const std::uint64_t timer_gen = connect_timer_generation_;
connect_suspend_handle_ = h;
connect_suspend_valid_ = std::move(valid);
connect_coroutine_pending_ = true;
connect_handshake_failed_ = false;
_error = error::db_error{"unknown error"};
if (is_connected_) {
try_resume_connect_wait();
return;
}
const double t_out =
timeout_override > qb::duration::zero()
? qb::detail::to_ev_seconds(timeout_override)
: (conn_opts_.connect_timeout > qb::duration::zero() ? qb::detail::to_ev_seconds(conn_opts_.connect_timeout) : 10.0);
const qb::io::uri connect_uri{conn_opts_.schema + "://" + conn_opts_.uri};
auto awaiter_valid = connect_suspend_valid_;
using transport_sock = std::remove_cvref_t<typename QB_IO_::transport_io_type>;
auto cb = [this, timer_gen, t_out, awaiter_valid](transport_sock &&sock) {
if (awaiter_valid && !*awaiter_valid)
return;
on_transport_ready(std::move(sock), timer_gen, t_out);
};
if constexpr (transport_sock::is_secure()) {
#ifdef QB_HAS_SSL
// PostgreSQL negotiates TLS in-band (the cleartext SSLRequest packet) BEFORE
// the handshake. Drive the whole connect → SSLRequest → TLS handshake through
// the connector's STARTTLS path so it runs fully asynchronously on the event
// loop (no blocking send/recv/handshake). A server that declines SSL fails the
// connect — a secure database requires TLS; use the plain tcp::database for
// cleartext.
// ssl_verify_mode::full -> verify the chain + host (the connector passes the remote host to
// ssl::socket::init_client); ::none -> encrypt only. The socket's ssl::Context (built below)
// governs verification, so the STARTTLS connector is handed a ready socket, not a verify bool.
const bool verify = (conn_opts_.ssl_verify == qb::pg::ssl_verify_mode::full);
if (!verify) {
// P2: encrypt-only (ssl_verify=none, the documented libpq-matching default) does NOT
// check the server certificate chain or hostname, leaving an active-MITM window on an
// untrusted network. SCRAM-SHA-256 mutual auth (enforced since the P1 fix) still
// authenticates the server, but non-SCRAM auth and the TLS channel itself are
// unprotected. Surface it ONCE so an unverified secure connection is never silent.
static bool warned_unverified_tls = false;
if (!warned_unverified_tls) {
warned_unverified_tls = true;
QB_LOG_WARN("[pgsql] TLS WITHOUT certificate verification (ssl_verify=none): the server "
"certificate chain/hostname is NOT verified. Set ssl_verify_mode::full to "
"authenticate the server (or rely on SCRAM-SHA-256 mutual auth).");
}
}
auto tls = make_client_tls_context();
if (!tls.ok()) {
// Fail CLOSED on a bad CA/cert/key path rather than silently connecting without it.
connect_handshake_failed_ = true;
_error = error::connection_error{"pgsql: TLS context configuration failed: " + tls.error()};
try_resume_connect_wait();
return;
}
transport_sock sock{std::move(tls)};
qb::io::async::tcp::starttls_connect<transport_sock, postgres_ssl_negotiator>(std::move(sock), connect_uri, std::move(cb),
qb::detail::from_ev_seconds(t_out));
#else
connect_handshake_failed_ = true;
_error = error::connection_error{"ssl transport requires QB_HAS_SSL"};
try_resume_connect_wait();
#endif
} else {
qb::io::async::tcp::connect<transport_sock>(connect_uri, std::move(cb), qb::detail::from_ev_seconds(t_out));
}
}
/**
* @brief Switches to the PostgreSQL protocol, starts read/write watchers, sends startup.
*
* Schedules the existing application-level handshake timeout (authentication / ReadyForQuery),
* distinct from the TCP connector deadline in `start_connect_from_awaiter`.
*
* @param timer_gen Generation counter; stale timers ignore the callback after reconnect
* @param t_out Handshake timeout in seconds passed to `async::callback`
*/
void
attach_pg_protocol_and_handshake_timer(std::uint64_t timer_gen, double t_out) {
this->template switch_protocol<pg_protocol>(*this);
this->start();
send_startup_message();
// Owned, cancellable deadline (see connect_deadline_timer_): reassigning here cancels
// any prior pending timer, and ~Database cancels this one — so the callback can never
// fire into a freed `this`. The timer_gen guard still covers an in-place reconnect.
connect_deadline_timer_ = qb::io::async::scoped_callback(std::function<void()>([this, t_out, timer_gen]() {
if (timer_gen != connect_timer_generation_)
return;
if (!connect_coroutine_pending_ || is_connected_)
return;
QB_LOG_WARN("[pgsql] Connection timed out after " << t_out << "s");
_error = error::db_error{"connection timeout"};
connect_handshake_failed_ = true;
try_resume_connect_wait();
}),
qb::detail::from_ev_seconds(t_out));
try_resume_connect_wait();
}
/**
* @brief Installs the ready transport and starts the PostgreSQL session.
*
* The connector delivers a fully-established socket: for a **secure** database it
* has already run the cleartext SSLRequest negotiation AND the TLS handshake
* asynchronously (via `starttls_connect` + `postgres_ssl_negotiator`); for a
* **plain** database it is the connected cleartext TCP socket. Either way this
* switches to the PostgreSQL protocol and sends the startup message — there is no
* blocking send/recv/handshake on the event loop anymore.
*
* @param sock Ready transport socket; empty/closed if the connect or TLS handshake failed.
* @param timer_gen Generation counter passed through to the handshake timer.
* @param t_out Authentication / ReadyForQuery timeout in seconds.
*/
template <typename Sock_>
void
on_transport_ready(Sock_ &&sock, std::uint64_t timer_gen, double t_out) {
if (!sock.is_open()) {
connect_handshake_failed_ = true;
_error = error::connection_error{"connection / TLS handshake failed"};
try_resume_connect_wait();
return;
}
this->clear_protocols(); // idempotent: drops any prior protocol, resets to the NoProtocol sentinel
this->transport() = std::forward<Sock_>(sock);
// Re-arm the SCRAM mutual-auth gate at the START of every handshake. These flags gate
// AuthenticationOk in on_authentication; they are also cleared in prepare_reconnect(), but a
// *bare* reconnect — connect() after disconnect() with no prepare_reconnect(), a supported and
// tested path (ReconnectWithoutPrepareReconnectIsUsable) — would otherwise carry a stale
// _scram_server_verified=true from a prior successful SCRAM login into the new handshake. An
// impersonating server on the second connection could then skip SASL entirely and send a bare
// AuthenticationOk to a client that never re-authenticated. on_transport_ready is the single
// choke point every connect path funnels through immediately before send_startup_message(), so
// resetting here (not only in prepare_reconnect) closes the reuse bypass on ALL connect paths.
_scram_pending = false;
_scram_server_verified = false;
_password_salt.clear();
_auth_message.clear();
attach_pg_protocol_and_handshake_timer(timer_gen, t_out);
}
/**
* @brief Creates a startup message for PostgreSQL connection
*
* Builds the startup message according to the PostgreSQL protocol specification.
* The message includes:
* - Protocol version
* - User authentication information
* - Target database name
* - Client parameters and options
*
* @param m Message object to populate with startup information
*/
void
create_startup_message(message &m) {
m.write(PROTOCOL_VERSION);
// Startup packet: null-terminated name=value pairs (PostgreSQL wire protocol).
// write(std::string) appends the required '\0' terminator; write_sv does not.
m.write(std::string(options::USER));
m.write(conn_opts_.user);
m.write(std::string(options::DATABASE));
m.write(conn_opts_.database);
for (auto &opt : client_opts_) {
m.write(opt.first);
m.write(opt.second);
}
// trailing terminator
m.write('\0');
}
/**
* @brief Sends the startup message to the PostgreSQL server
*/
void
send_startup_message() {
message m(empty_tag);
create_startup_message(m);
*this << m;
}
/**
* @brief Handles new command events in the transaction
*/
void
on_new_command() final {
process_if_query_ready();
}
/**
* @brief Handles sub-command status updates
*
* Propagates leaf `ResultQuery` / nested command outcomes to the root `Transaction`
* so `await()` and `status::operator bool` reflect failures (not only child `_result`).
*/
void
on_sub_command_status(bool status) final {
Transaction::on_sub_command_status(status);
}
/// Root `Transaction` for this connection (never null while `Database` lives).
Transaction *
root_transaction() noexcept {
return static_cast<Transaction *>(static_cast<Database<QB_IO_, NotifyDerived> *>(this));
}
Transaction *_current_command = this; ///< Current transaction being processed
ISqlQuery *_current_query = nullptr; ///< Current query being executed
bool _ready_for_query = false; ///< Flag indicating if ready for next query
/**
* @brief Finds the next transaction to execute
*
* Recursively traverses the transaction tree to find the
* deepest (leaf) transaction that should be executed next.
*
* @param cmd Current transaction
* @return Transaction* Next transaction to execute
*/
static Transaction *
next_transaction(Transaction *cmd) {
if (!cmd)
return nullptr;
auto sub = cmd->next_transaction();
if (!sub)
return cmd;
else
return next_transaction(sub);
}
/**
* @brief Processes a query in the transaction
*
* Fetches and executes the next query from the given transaction.
* If no more queries are in the current transaction, moves to parent.
*
* @param cmd Transaction containing the query
* @return bool true if a query was processed, false if no queries remain
*/
bool
process_query(Transaction *cmd) {
_ready_for_query = false;
if (!cmd)
cmd = root_transaction();
_current_command = next_transaction(cmd);
if (!_current_command)
_current_command = root_transaction();
_current_query = _current_command->next_query();
if (_current_query) {
if (qb::likely(_current_query->is_valid())) {
*this << _current_query->get();
return true;
} else {
QB_LOG_DEBUG("[pgsql] error processing query not valid");
_error = error::client_error{"query couldn't be processed check logs for more infos"};
on_error_query(error());
return process_query(_current_command) || (_ready_for_query = true);
}
} else if (_current_command->parent()) {
auto next_cmd = _current_command->parent();
do {
next_cmd->pop_transaction();
} while (!next_cmd->result() && (next_cmd = next_cmd->parent()));
return process_query(next_cmd);
}
return false;
}
/**
* @brief Processes queries if the client is ready
*/
void
process_if_query_ready() {
if (_ready_for_query) {
process_query(_current_command);
}
}
/**
* @brief Handles successful query completion
*/
void
on_success_query() {
if (!_current_query)
return;
if (!_current_command) {
_current_query = nullptr;
return;
}
auto query = _current_command->pop_query();
query->on_success();
_current_query = nullptr;
}
/**
* @brief Handles query error
*
* @param err Error information
*/
void
on_error_query(error::db_error const &err) {
_error = err;
if (!_current_query)
return;
if (!_current_command) {
_current_query = nullptr;
return;
}