Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 47 additions & 0 deletions pgduck_server/src/pgsession/pgsession.c
Original file line number Diff line number Diff line change
Expand Up @@ -878,6 +878,46 @@ process_execute_message(PGSession * pgSession, StringInfo inputMessage)
}


/*
* PGDUCK_ENGINE_ERROR_PREFIX marks log lines that carry only a canned error
* class, so a log collector can match on the prefix alone and collect nothing
* else from this process. The text after the prefix must always come from the
* fixed literal set in error_class_for_sqlstate -- never interpolate DuckDB
* text, a statement, a URL, or client identity into it. Every other log line
* here may carry all of those.
*/
#define PGDUCK_ENGINE_ERROR_PREFIX "pgduck_engine_error: "

/*
* error_class_for_sqlstate maps a SQLSTATE this server reports to a PII-free
* class. Codes we do not map are "other", and NULL is one of them.
*
* NULL means no error type was recorded; pgsession_send_postgres_error reports
* it to the client as feature_not_supported, which is not a category worth a
* class of its own.
*/
static const char *
error_class_for_sqlstate(const char *sqlState)
{
if (sqlState == NULL)
return "other";

if (strcmp(sqlState, PGDUCK_SQLSTATE_OUT_OF_MEMORY) == 0)
return "out_of_memory";

if (strcmp(sqlState, PGDUCK_SQLSTATE_IO_ERROR) == 0)
return "io_error";

if (strcmp(sqlState, PGDUCK_SQLSTATE_INVALID_PARAMETER) == 0)
return "invalid_input";

if (strcmp(sqlState, PGDUCK_SQLSTATE_INTERNAL_ERROR) == 0)
return "internal_error";

return "other";
}


/*
* sqlstate_for_status maps the statuses that are raised without a DuckDB error
* type, so no SQLSTATE was recorded on the session. NULL leaves the default.
Expand Down Expand Up @@ -913,6 +953,13 @@ handle_pgsession_error_message(DuckDBStatus status, PGSession * pgSession, char

pgSession->duckSession.errorSqlState = NULL;

/*
* Every reportable status passes through here, including the fatal ones
* the caller exits on, so one line here covers all of them.
*/
PGDUCK_SERVER_LOG(PGDUCK_ENGINE_ERROR_PREFIX "%s",
Comment thread
sfc-gh-rsirma marked this conversation as resolved.
error_class_for_sqlstate(sqlState));

switch (status)
{
case DUCKDB_QUERY_ERROR:
Expand Down
203 changes: 203 additions & 0 deletions pgduck_server/tests/pytests/test_engine_error_class_log.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,203 @@
"""Coverage for the classified log line pgduck_server emits for engine errors.

The line carries a canned class and nothing else, so that a log collector can
match on its prefix and collect nothing else from this process. Every other line
pgduck_server writes may legitimately contain the DuckDB message, the failing
statement, or the client address, which is why the assertions here isolate the
classified line instead of searching the whole output.
"""

import os
import tempfile
from contextlib import contextmanager

import pytest
from utils_pytest import *


PGDUCK_UNIX_DOMAIN_PATH = "/tmp"
PGDUCK_PORT = 8259 # its own port, so a shared server cannot answer instead
DUCKDB_DATABASE_FILE_PATH = "/tmp/pgduck_engine_error_class.db"

ENGINE_ERROR_PREFIX = "pgduck_engine_error: "

# Spill setup copied from test_server_start.py: a 0-byte temp cap makes the
# overflow deterministic, while a memory_limit well above DuckDB's block size and
# threads=1 keep it a recoverable temp-cap OOM rather than a fatal memory OOM.
SPILL_MEMORY_LIMIT = "32MB"
SPILL_QUERY = "SELECT count(*) FROM (SELECT i FROM range(10000000) t(i) GROUP BY i) g"

# A path that does not exist, distinctive enough that finding it on the classified
# line can only mean the line leaked its query.
MISSING_FILE = "/tmp/pgduck_class_log_absent_9f3c1d.parquet"


def _connect():
conn = psycopg2.connect(host=PGDUCK_UNIX_DOMAIN_PATH, port=PGDUCK_PORT)
conn.autocommit = True
return conn


def _start_server(need_output=True, extra_args=None):
"""Start a server and wait for it to accept connections.

pgduck_server binds its socket only after DuckDB is initialized, so
connecting without this wait races that startup.
"""
server = PgDuckServer(
port=PGDUCK_PORT,
duckdb_database_file_path=DUCKDB_DATABASE_FILE_PATH,
need_output=need_output,
extra_args=extra_args,
)
assert is_server_listening(server.socket_path), "pgduck_server did not start"
return server


@contextmanager
def _spill_server(extra_args=None):
"""Start a server whose spill knobs come from an init file, as production does."""
with tempfile.TemporaryDirectory(dir="/tmp") as cfg_dir:
init_file = os.path.join(cfg_dir, "init.sql")
with open(init_file, "w") as f:
f.write(f"SET GLOBAL memory_limit='{SPILL_MEMORY_LIMIT}';\n")
f.write("SET GLOBAL threads='1';\n")
f.write("SET GLOBAL max_temp_directory_size='0KiB';\n")
yield _start_server(
extra_args=["--init_file_path", init_file] + list(extra_args or [])
)


def _class_lines(server):
output = get_server_output(server.output_queue)
return [line for line in output.splitlines() if ENGINE_ERROR_PREFIX in line]


def _assert_class(server, expected, *must_not_appear):
"""Assert the expected class was logged, and that the line carries nothing else.

*must_not_appear* holds fragments of the statement and of the DuckDB message.
They are checked against the classified lines only: the neighbouring WARNING
lines contain all of them by design and would mask a leak.
"""
lines = _class_lines(server)

assert lines, f"pgduck_server logged no {ENGINE_ERROR_PREFIX!r} line"

assert any(
ENGINE_ERROR_PREFIX + expected in line for line in lines
), f"expected class {expected!r}, got: {lines!r}"

for line in lines:
for fragment in must_not_appear:
assert fragment not in line, (
f"classified line leaked {fragment!r}, which must never reach a "
f"collected log line: {line!r}"
)


def test_io_error_class_is_logged_without_the_path():
"""A missing file is an IO error, and the path must not ride along."""
server = _start_server()

cur = _connect().cursor()
with pytest.raises(psycopg2.Error) as exc_info:
cur.execute(f"SELECT * FROM read_parquet('{MISSING_FILE}')")

assert (
exc_info.value.pgcode == "58030"
), f"missing file should report SQLSTATE 58030, got {exc_info.value.pgcode}"

_assert_class(server, "io_error", MISSING_FILE, "read_parquet")


def test_invalid_input_class_is_logged_without_the_value():
"""A bad format string is invalid input; neither it nor the value may appear."""
server = _start_server()
cur = _connect().cursor()
with pytest.raises(psycopg2.Error) as exc_info:
cur.execute("SELECT strptime('not-a-date', '%Y-%m-%d')")

assert (
exc_info.value.pgcode == "22023"
), f"invalid input should report SQLSTATE 22023, got {exc_info.value.pgcode}"

_assert_class(server, "invalid_input", "not-a-date", "strptime")


def test_invalid_input_class_is_logged_over_extended_protocol():
"""A bind parameter puts psycopg2 on Parse/Bind/Execute, the path pg_lake uses."""
server = _start_server()
cur = _connect().cursor()
with pytest.raises(psycopg2.Error) as exc_info:
cur.execute("SELECT strptime(%s, %s)", ("not-a-date", "%Y-%m-%d"))

assert (
exc_info.value.pgcode == "22023"
), f"invalid input should report SQLSTATE 22023, got {exc_info.value.pgcode}"

_assert_class(server, "invalid_input", "not-a-date", "strptime")


def test_unmapped_error_is_logged_as_other():
"""A catalog error is not a class we distinguish, so it lands in "other"."""
server = _start_server()
cur = _connect().cursor()
with pytest.raises(psycopg2.Error) as exc_info:
cur.execute("SELECT * FROM no_such_table_1a2b3c")

assert (
exc_info.value.pgcode == "0A000"
), f"unmapped error should stay SQLSTATE 0A000, got {exc_info.value.pgcode}"

_assert_class(server, "other", "no_such_table_1a2b3c")


def test_recoverable_out_of_memory_class_is_logged():
"""A temp-cap overflow is a recoverable OOM: classified, and the server lives."""
with _spill_server() as server:
assert is_server_listening(server.socket_path)

cur = _connect().cursor()
with pytest.raises(psycopg2.Error) as exc_info:
cur.execute(SPILL_QUERY)

assert (
exc_info.value.pgcode == "53200"
), f"recoverable OOM should report SQLSTATE 53200, got {exc_info.value.pgcode}"

_assert_class(server, "out_of_memory", "range(10000000)")

assert is_server_listening(
server.socket_path
), "pgduck_server stopped accepting connections after a recoverable OOM"
assert server.process.poll() is None, "pgduck_server process exited"


# Distinctive so a leak on the classified line cannot be a coincidence.
SECRET_KEY_ID = "leaky-key-id-9f3c1d"
SECRET_SECRET = "leaky-secret-9f3c1d"
SECRET_TOKEN = "leaky-session-token-9f3c1d"


def test_failed_create_secret_class_line_does_not_carry_credentials():
"""A bad CREATE SECRET still logs a class, without KEY_ID/SECRET/token."""
server = _start_server()
cur = _connect().cursor()
with pytest.raises(psycopg2.Error):
cur.execute(
"CREATE SECRET leak_probe_9f3c1d ("
"TYPE NOT_A_SECRET_TYPE, "
f"KEY_ID '{SECRET_KEY_ID}', "
f"SECRET '{SECRET_SECRET}', "
f"SESSION_TOKEN '{SECRET_TOKEN}'"
")"
)

_assert_class(
server,
"invalid_input",
SECRET_KEY_ID,
SECRET_SECRET,
SECRET_TOKEN,
)
18 changes: 18 additions & 0 deletions pgduck_server/tests/pytests/test_server_start.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import pytest
import subprocess
import os
import re
import signal
import time
import tempfile
Expand Down Expand Up @@ -692,6 +693,23 @@ def test_genuine_oom_over_extended_protocol_terminates_server():
"Out of Memory Error" in server_output
), f"expected the genuine OOM message to be surfaced, got: {server_output}"

# The class line must be written before exit(): after the process dies
# the client often only sees lost_connection, so this record is the
# one a collector can still attribute to OOM.
classified = [
line
for line in server_output.splitlines()
if re.match(r"^\S+ LOG pgduck_engine_error: out_of_memory$", line)
]
assert classified, (
"expected a classified out_of_memory record before the fatal exit, "
f"got: {server_output}"
)
for line in classified:
assert (
"range(100000000)" not in line
), f"classified line leaked the statement: {line!r}"


# ---------------------------------------------------------------------------
# Spill location via the init file
Expand Down
Loading