Skip to content

Commit d0d7145

Browse files
committed
fix(client,pgpool): read idle socket messages and isolate connection failures
1 parent d32ac99 commit d0d7145

12 files changed

Lines changed: 315 additions & 42 deletions

‎client/INTERNAL.md‎

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -53,9 +53,12 @@ temporary prepared-statement cleanup.
5353

5454
## Driver And Queues
5555

56-
The driver processes one `Request` at a time:
56+
The driver keeps one reader task on the socket for the entire established
57+
connection, including idle periods. It routes notices, notifications, and
58+
parameter updates immediately. Ordinary backend responses pass through a
59+
one-message queue to the request loop, which processes one `Request` at a time:
5760

58-
1. receive and retain the next accepted operation request before any socket I/O
61+
1. receive and retain the next accepted operation request
5962
2. write its frontend bytes
6063
3. route every ordinary response to that request's bounded response queue
6164
4. continue until `ReadyForQuery` or the request's protocol-specific terminal
@@ -65,15 +68,15 @@ The driver processes one `Request` at a time:
6568
There is no pending-request list, opportunistic pipeline decision, or
6669
oldest-request response routing.
6770

68-
Each request response queue remains bounded (currently eight messages).
69-
Backpressure therefore stops socket reads for the active request rather than
70-
growing memory without limit. Because execution is single-flight, a stalled
71-
consumer cannot be bypassed by a later request.
71+
Each request response queue remains bounded (currently eight messages). The
72+
one-message reader handoff keeps socket reads bounded when an active response
73+
consumer stalls. Because execution is single-flight, a later request cannot
74+
bypass that consumer.
7275

7376
COPY IN is the protocol exception that requires bidirectional coordination.
7477
Its producer queue is bounded to eight actions. After PostgreSQL enters COPY
7578
mode, the driver starts exactly one request-local writer for data, finish, or
76-
abort actions while the driver itself remains the sole backend reader. An early
79+
abort actions while its reader task remains the sole backend reader. An early
7780
`ErrorResponse` is forwarded to the sink, the input queue is cleared and closed,
7881
and the writer stops after its current complete frame. The driver joins that
7982
writer before forwarding the final `ReadyForQuery` or starting another request.
@@ -192,7 +195,8 @@ Out-of-band backend messages never enter an operation response queue:
192195

193196
They update shared state and are published through the Client async-message
194197
queue. `Client::next_message()` has one-consumer semantics and returns `None`
195-
when the driver terminates.
198+
when the driver terminates. The reader continues delivering these messages
199+
without a query in flight.
196200

197201
Before starting a physical connection, `pgpool` installs the internal async
198202
message handler. That handler forwards messages to the pool's buffer instead

‎client/INTERNAL_CN.md‎

Lines changed: 10 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -44,23 +44,25 @@
4444

4545
## 驱动器与队列
4646

47-
驱动器每次处理一个 `Request`:
47+
驱动器在连接建立后的整个生命周期维持一个套接字读取任务,包括空闲时段。
48+
提示、通知和参数更新会立即发布;普通响应经容量为一的队列交给请求循环。
49+
请求循环每次处理一个 `Request`:
4850

49-
1. 在进行任何套接字 I/O 前,先接收并保留下一个已接受操作的请求
51+
1. 接收并保留下一个已接受操作的请求
5052
2. 写入该请求的前端协议字节
5153
3. 将每条普通响应路由到该请求的有界响应队列
5254
4. 持续处理,直到收到 `ReadyForQuery` 或到达该请求协议规定的终止状态
5355
5. 转而处理下一个请求
5456

5557
驱动器不维护待处理请求列表,不会择机启用流水线,也不会按最早请求来路由响应。
5658

57-
每个请求的响应队列始终有容量上限(目前为八条消息)。因此,背压会暂停活动请求的
58-
套接字读取,而不会让内存无限增长。由于同一时刻只执行一个操作,后续请求无法绕过
59-
停滞的消费者。
59+
每个请求的响应队列始终有容量上限(目前为八条消息)。读取任务到请求循环的队列
60+
最多缓存一条消息,活动响应的消费者停滞时,背压仍会限制套接字读取和内存占用。
61+
由于同一时刻只执行一个操作,后续请求无法绕过停滞的消费者。
6062

6163
COPY IN 是协议中的例外,需要双向协调。其生产者队列最多容纳八个动作。PostgreSQL
6264
进入 COPY 模式后,驱动器会为该请求创建且仅创建一个写入任务,用来处理数据、完成或
63-
中止动作,而驱动器自身仍是唯一的后端消息读取者。提前到达的 `ErrorResponse` 会被
65+
中止动作,而读取任务仍是唯一的后端消息读取者。提前到达的 `ErrorResponse` 会被
6466
转发给写入端(sink),输入队列随即被清空并关闭,写入任务则在完成当前整个帧的写入后
6567
停止。驱动器会等待该写入任务结束,然后才转发最终的 `ReadyForQuery` 或开始另一个请求。
6668
整个过程仍然属于同一个逻辑操作。
@@ -142,7 +144,8 @@ async 0.22 在离开受保护代码块时,不会自动传播待处理的取消
142144
- `ParameterStatus`
143145

144146
它们会更新共享状态,并通过 Client 的异步消息队列发布。`Client::next_message()`
145-
采用单消费者语义,在驱动器终止时返回 `None`。
147+
采用单消费者语义,在驱动器终止时返回 `None`。即使没有查询在执行,读取任务也会
148+
持续交付这些消息。
146149

147150
`pgpool` 在启动物理连接前安装内部异步消息处理器,将消息转发到连接池的缓冲队列,
148151
而不是 Client 队列;连接池通过 `Pool::next_message()` 消费这些消息。

‎client/README.mbt.md‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -353,7 +353,8 @@ async fn _read_async_message(client : @client.Client) -> @client.AsyncMessage? {
353353
}
354354
```
355355

356-
Use one dedicated consumer for `next_message()`. It returns `None` after the
356+
Use one dedicated consumer for `next_message()`. Notifications arrive while the
357+
connection is idle, without sending another query. It returns `None` after the
357358
driver closes, after buffered messages are read. By default the client retains
358359
at most 256 messages. `Client::create_with_async_message_capacity(config,
359360
capacity)` accepts a positive capacity. When full, the driver discards the

‎client/connection_idle_wbtest.mbt‎

Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
1+
///|
2+
fn idle_notification_packet(
3+
channel : String,
4+
payload : String,
5+
) -> Bytes raise @proto.ProtocolError {
6+
let body = Buffer()
7+
body.write_int_be(123)
8+
body.write_bytes(@proto.utf8_encode(channel))
9+
body.write_byte(0)
10+
body.write_bytes(@proto.utf8_encode(payload))
11+
body.write_byte(0)
12+
connect_backend_packet(@backend.NOTIFICATION_RESPONSE_TAG, body)
13+
}
14+
15+
///|
16+
async test "idle connection delivers notifications without another request" {
17+
@async.with_task_group(group => {
18+
let (client, connection, backend) = connection_test_pair(group)
19+
defer backend.close()
20+
let driver = group.spawn(no_wait=true, allow_failure=true, () => {
21+
connection.run(group)
22+
})
23+
client.shared.driver_task.val = Some(driver)
24+
defer client.abort()
25+
@async.with_timeout(2000, () => {
26+
backend.write(idle_notification_packet("moon_idle", "first"))
27+
assert_true(
28+
client.next_message()
29+
is Some(Notification({ channel: "moon_idle", payload: "first", .. })),
30+
)
31+
let ping = group.spawn(() => client.check_connection())
32+
assert_eq(backend.read_exactly(5), sync_bytes())
33+
backend.write(ready_for_query_packet(b'I'))
34+
ping.wait()
35+
})
36+
})
37+
}
38+
39+
///|
40+
async test "partial idle notification stays aligned when a request arrives" {
41+
@async.with_task_group(group => {
42+
let (client, connection, backend) = connection_test_pair(group)
43+
defer backend.close()
44+
let driver = group.spawn(no_wait=true, allow_failure=true, () => {
45+
connection.run(group)
46+
})
47+
client.shared.driver_task.val = Some(driver)
48+
defer client.abort()
49+
@async.with_timeout(2000, () => {
50+
let notification = idle_notification_packet("moon_partial", "second")
51+
backend.write(notification[:3].to_owned())
52+
let ping = group.spawn(() => client.check_connection())
53+
assert_eq(backend.read_exactly(5), sync_bytes())
54+
backend.write(notification[3:].to_owned())
55+
backend.write(ready_for_query_packet(b'I'))
56+
assert_true(
57+
client.next_message()
58+
is Some(
59+
Notification({ channel: "moon_partial", payload: "second", .. })
60+
),
61+
)
62+
ping.wait()
63+
})
64+
})
65+
}
66+
67+
///|
68+
async test "idle socket EOF closes the driver without another request" {
69+
@async.with_task_group(group => {
70+
let (client, connection, backend) = connection_test_pair(group)
71+
let driver = group.spawn(no_wait=true, allow_failure=true, () => {
72+
connection.run(group)
73+
})
74+
@async.with_timeout(2000, () => {
75+
backend.close()
76+
driver.wait()
77+
assert_true(client.is_closed())
78+
assert_true(client.shared.runtime_error.val is Some(_))
79+
assert_true(client.next_message() is None)
80+
})
81+
})
82+
}

‎client/connection_loop.mbt‎

Lines changed: 62 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,10 @@
1212
async fn Driver::run(self : Driver, group : @async.TaskGroup[Unit]) -> Unit {
1313
let mut active : Request? = None
1414
let mut cleaned_up = false
15+
// The reader must survive the transition between idle time and a request.
16+
// Cancelling an in-progress frame read there would lose protocol alignment.
17+
let incoming : @async.Queue[@backend.Message] = Queue(kind=Blocking(1))
18+
let reading_request = @ref.new(false)
1519
fn cleanup(error : Error?) {
1620
if cleaned_up {
1721
return
@@ -43,6 +47,15 @@ async fn Driver::run(self : Driver, group : @async.TaskGroup[Unit]) -> Unit {
4347
self.shared.control.close()
4448
}
4549
}
50+
let reader = group.spawn(no_wait=true, allow_failure=true, () => {
51+
self.receive_messages(incoming, reading_request)
52+
})
53+
defer @async.protect_from_cancel(() => {
54+
reader.cancel()
55+
reader.wait() catch {
56+
_ => ()
57+
}
58+
})
4659
try {
4760
while true {
4861
let request = self.shared.requests.get()
@@ -55,8 +68,9 @@ async fn Driver::run(self : Driver, group : @async.TaskGroup[Unit]) -> Unit {
5568
return
5669
}
5770
Messages | CopyIn => {
71+
reading_request.val = true
5872
self.stream.write(request.bytes)
59-
self.drain_request(request, group)
73+
self.drain_request(request, incoming, group)
6074
self.close_request(request, None)
6175
active = None
6276
}
@@ -67,6 +81,48 @@ async fn Driver::run(self : Driver, group : @async.TaskGroup[Unit]) -> Unit {
6781
}
6882
}
6983

84+
///|
85+
/// Keep one owner of socket reads even when no request is in flight.
86+
async fn Driver::receive_messages(
87+
self : Driver,
88+
incoming : @async.Queue[@backend.Message],
89+
reading_request : @ref.Ref[Bool],
90+
) -> Unit {
91+
let mut failure : Error? = None
92+
try {
93+
while true {
94+
let message = self.read_message()
95+
if self.handle_async_message(message) {
96+
continue
97+
}
98+
if !reading_request.val {
99+
match message {
100+
ErrorResponse(body) =>
101+
raise ClientError::Database(parse_database_error(body.fields()))
102+
_ =>
103+
raise ClientError::UnexpectedMessage(
104+
"backend response without an active request",
105+
)
106+
}
107+
}
108+
if message is ReadyForQuery(_) {
109+
reading_request.val = false
110+
}
111+
incoming.put(message)
112+
}
113+
} catch {
114+
error => failure = Some(error)
115+
}
116+
match failure {
117+
Some(error) => {
118+
// Preserve queued requests for Driver::cleanup to close their own queues.
119+
incoming.close(error~, clear=true)
120+
self.shared.requests.close(error~)
121+
}
122+
None => ()
123+
}
124+
}
125+
70126
///|
71127
/// Receive the next asynchronous message buffered by the private driver.
72128
///
@@ -83,17 +139,15 @@ pub async fn Client::next_message(self : Client) -> AsyncMessage? {
83139
async fn Driver::drain_request(
84140
self : Driver,
85141
request : Request,
142+
incoming : @async.Queue[@backend.Message],
86143
group : @async.TaskGroup[Unit],
87144
) -> Unit {
88145
if request.kind is CopyIn {
89-
self.drain_copy_in_request(request, group)
146+
self.drain_copy_in_request(request, incoming, group)
90147
return
91148
}
92149
while true {
93-
let message = self.read_message()
94-
if self.handle_async_message(message) {
95-
continue
96-
}
150+
let message = incoming.get()
97151
match request.kind {
98152
Messages => self.forward_message_to_request(request, message)
99153
CopyIn => abort("COPY IN requests use drain_copy_in_request")
@@ -131,17 +185,15 @@ async fn Driver::forward_message_to_request(
131185
async fn Driver::drain_copy_in_request(
132186
self : Driver,
133187
request : Request,
188+
incoming : @async.Queue[@backend.Message],
134189
group : @async.TaskGroup[Unit],
135190
) -> Unit {
136191
let copy_input = request.copy_input.unwrap()
137192
let copy_error = request.copy_error.unwrap()
138193
let writer_stopped = @ref.new(false)
139194
let mut writer : @async.Task[Unit]? = None
140195
while true {
141-
let message = self.read_message()
142-
if self.handle_async_message(message) {
143-
continue
144-
}
196+
let message = incoming.get()
145197
match message {
146198
CopyInResponse(_) => {
147199
guard writer is None else {

‎pgpool/INTERNAL.md‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -36,5 +36,6 @@ only then may it allocate a request ID and mark a new request active.
3636
Streaming callbacks finish or abandon their raw handles before returning the
3737
connection. Detached drains save expected connection closure for `finish()`;
3838
database errors remain in stream state, and unexpected protocol errors fail the
39-
owning client executor and then the pool executor. Ordinary task cancellation
40-
does not imply a PostgreSQL cancel packet, automatic retry, or physical abort.
39+
owning client executor. The pool retires that physical connection on return or
40+
its next checkout. Ordinary task cancellation does not imply a PostgreSQL
41+
cancel packet, automatic retry, or physical abort.

‎pgpool/INTERNAL_CN.md‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,5 +28,6 @@
2828

2929
流式回调在归还连接前,会先结束或放弃其底层句柄。分离的排空任务会保存预期的连接
3030
关闭结果,供 `finish()` 读取;数据库错误保留在流状态中,
31-
意外的协议错误会使所属连接和连接池执行器失败。普通的任务取消并不意味着发送 PostgreSQL 取消报文、
31+
意外的协议错误会使所属客户端执行器失败;连接池在归还该连接或下一次借用时将其退役。
32+
普通的任务取消并不意味着发送 PostgreSQL 取消报文、
3233
自动重试或中止物理连接。

‎pgpool/README.mbt.md‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -289,6 +289,7 @@ fills, the oldest message is discarded and
289289
`Pool::dropped_async_messages()` increases. `PoolConfig::PoolConfig` accepts a
290290
positive `async_message_capacity`; the default is 256. After `Pool::close()`,
291291
buffered messages can still be read, then `next_message()` returns `None`.
292+
An idle listening connection delivers notifications without another checkout.
292293

293294
Create and recycle deadlines are hard for client I/O: when either expires, the
294295
pool first calls `Client::abort()` on the candidate physical connection and
@@ -312,8 +313,10 @@ and await `pool.ready()` before use. A repeated `run()` raises
312313
`PoolError::ExecutorAlreadyStarted`. `Pool::close()` rejects new operations and
313314
closes idle connections. Active callbacks are allowed to finish; their physical
314315
connections close instead of returning to the idle set. Await the executor task
315-
to know that all connections have stopped. Cancellation or an unexpected
316-
background connection failure closes the whole pool.
316+
to know that all connections have stopped. Pool executor cancellation or an
317+
unexpected pool maintenance failure closes the whole pool. An individual
318+
physical connection failure closes that client; the pool discards it on return
319+
or on its next checkout and can create a replacement.
317320

318321
`Pool::resize()` never interrupts checked-out sessions. When shrinking cannot
319322
remove enough currently available capacity tokens, the pool records resize

0 commit comments

Comments
 (0)