Skip to content

Commit df84766

Browse files
committed
Merge branch 'fix/batch-cursor-field-descriptors-2626'
2 parents c91dd3b + 5862cf6 commit df84766

4 files changed

Lines changed: 244 additions & 3 deletions

File tree

batch_test.go

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -178,6 +178,57 @@ func TestConnSendBatchEmptyQuery(t *testing.T) {
178178
})
179179
}
180180

181+
// https://github.com/jackc/pgx/issues/2626
182+
func TestConnSendBatchCursorFetch(t *testing.T) {
183+
t.Parallel()
184+
185+
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
186+
defer cancel()
187+
188+
pgxtest.RunWithQueryExecModes(ctx, t, defaultConnTestRunner, nil, func(ctx context.Context, t testing.TB, conn *pgx.Conn) {
189+
pgxtest.SkipCockroachDB(t, conn, "Server does not support cursors in implicit transactions")
190+
191+
tx, err := conn.Begin(ctx)
192+
require.NoError(t, err)
193+
defer tx.Rollback(ctx)
194+
195+
// The FETCH is prepared before the DECLARE has executed, so the cursor does not exist and the server
196+
// describes the FETCH result as NoData. The actual fields are only known at execution time.
197+
batch := &pgx.Batch{}
198+
batch.Queue(`declare test_cursor cursor for select n, n::text as str from generate_series(1, 3) n`)
199+
batch.Queue(`fetch all in test_cursor`)
200+
201+
br := tx.SendBatch(ctx, batch)
202+
203+
_, err = br.Exec()
204+
require.NoError(t, err)
205+
206+
rows, err := br.Query()
207+
require.NoError(t, err)
208+
209+
fds := rows.FieldDescriptions()
210+
require.Len(t, fds, 2)
211+
require.Equal(t, "n", fds[0].Name)
212+
require.Equal(t, "str", fds[1].Name)
213+
214+
var ns []int32
215+
var strs []string
216+
for rows.Next() {
217+
var n int32
218+
var str string
219+
require.NoError(t, rows.Scan(&n, &str))
220+
ns = append(ns, n)
221+
strs = append(strs, str)
222+
}
223+
require.NoError(t, rows.Err())
224+
require.Equal(t, []int32{1, 2, 3}, ns)
225+
require.Equal(t, []string{"1", "2", "3"}, strs)
226+
227+
err = br.Close()
228+
require.NoError(t, err)
229+
})
230+
}
231+
181232
func TestConnSendBatchQueuedQuery(t *testing.T) {
182233
t.Parallel()
183234

pgconn/pgconn.go

Lines changed: 38 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1303,7 +1303,9 @@ func (pgConn *PgConn) ExecPrepared(ctx context.Context, stmtName string, paramVa
13031303
//
13041304
// This differs from [PgConn.ExecPrepared] in that it takes a [*StatementDescription] instead of the prepared statement name.
13051305
// Because it has the [*StatementDescription] it can avoid the Describe Portal message that [PgConn.ExecPrepared] must send to get
1306-
// the result column descriptions.
1306+
// the result column descriptions. However, if the statement description has no fields then a Describe is still sent, as
1307+
// an empty Fields may mean the results were not knowable at prepare time, e.g. a FETCH from a cursor that did not exist
1308+
// yet.
13071309
//
13081310
// paramValues are the parameter values. It must be encoded in the format given by paramFormats.
13091311
//
@@ -1364,7 +1366,10 @@ func (pgConn *PgConn) execExtendedPrefix(ctx context.Context, paramValues [][]by
13641366
}
13651367

13661368
func (pgConn *PgConn) execExtendedSuffix(result *ResultReader, statementDescription *StatementDescription, resultFormats []int16) {
1367-
if statementDescription == nil {
1369+
if statementDescription == nil || len(statementDescription.Fields) == 0 {
1370+
// The cached field descriptions are missing or empty. Empty field descriptions can occur when the statement's
1371+
// result set was not known at prepare time, e.g. a FETCH from a cursor that did not exist yet. Send a Describe
1372+
// so the server supplies the actual row description when the statement is executed.
13681373
pgConn.frontend.SendDescribe(&pgproto3.Describe{ObjectType: 'P'})
13691374
}
13701375
pgConn.frontend.SendExecute(&pgproto3.Execute{})
@@ -2523,6 +2528,12 @@ func (p *Pipeline) SendQueryStatement(statementDescription *StatementDescription
25232528
}
25242529

25252530
p.conn.frontend.SendBind(&pgproto3.Bind{PreparedStatement: statementDescription.Name, ParameterFormatCodes: paramFormats, Parameters: paramValues, ResultFormatCodes: resultFormats})
2531+
if len(statementDescription.Fields) == 0 {
2532+
// The cached field descriptions are empty. This can occur when the statement's result set is
2533+
// not known at prepare time, e.g. a FETCH from a cursor that did not exist yet. Send a
2534+
// Describe so the server supplies the actual row description when the statement is executed.
2535+
p.conn.frontend.SendDescribe(&pgproto3.Describe{ObjectType: 'P'})
2536+
}
25262537
p.conn.frontend.SendExecute(&pgproto3.Execute{})
25272538
p.state.PushBackRequestType(pipelineQueryStatement)
25282539
p.state.PushBackStatementData(statementDescription, resultFormats)
@@ -2729,12 +2740,36 @@ func (p *Pipeline) getResultsQueryStatement() (*ResultReader, error) {
27292740
return nil, err
27302741
}
27312742

2743+
sdFields := sd.Fields
2744+
if len(sdFields) == 0 {
2745+
// A Describe was sent for this statement (see SendQueryStatement). Read the server-provided
2746+
// row description which may include fields that were not known at prepare time.
2747+
msg, err := p.receiveMessage()
2748+
if err != nil {
2749+
return nil, err
2750+
}
2751+
2752+
switch msg := msg.(type) {
2753+
case *pgproto3.RowDescription:
2754+
sdFields = make([]FieldDescription, len(msg.Fields))
2755+
convertRowDescription(sdFields, msg)
2756+
case *pgproto3.NoData:
2757+
// Statement returns no rows.
2758+
case *pgproto3.ErrorResponse:
2759+
pgErr := ErrorResponseToPgError(msg)
2760+
p.state.HandleError(pgErr)
2761+
p.conn.resultReader.closed = true
2762+
return nil, pgErr
2763+
default:
2764+
return nil, p.handleUnexpectedMessage("QueryStatement RowDescription or NoData", msg)
2765+
}
2766+
}
2767+
27322768
msg, err := p.receiveMessage()
27332769
if err != nil {
27342770
return nil, err
27352771
}
27362772

2737-
sdFields := sd.Fields
27382773
fieldDescriptions := p.conn.getFieldDescriptionSlice(len(sdFields))
27392774
err = combineFieldDescriptionsAndResultFormats(fieldDescriptions, sdFields, resultFormats)
27402775
if err != nil {

pgconn/pgconn_test.go

Lines changed: 108 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1711,6 +1711,49 @@ func TestConnExecStatement(t *testing.T) {
17111711
ensureConnValid(t, pgConn)
17121712
}
17131713

1714+
// https://github.com/jackc/pgx/issues/2626
1715+
func TestConnExecStatementCursorFetch(t *testing.T) {
1716+
t.Parallel()
1717+
1718+
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
1719+
defer cancel()
1720+
1721+
pgConn, err := pgconn.Connect(ctx, os.Getenv("PGX_TEST_DATABASE"))
1722+
require.NoError(t, err)
1723+
defer closeConn(t, pgConn)
1724+
1725+
if pgConn.ParameterStatus("crdb_version") != "" {
1726+
t.Skip("Server does not support cursors in implicit transactions")
1727+
}
1728+
1729+
// Prepare the FETCH before the cursor exists. The server describes the result as NoData so the statement
1730+
// description has no fields. The actual fields are only known at execution time.
1731+
sd, err := pgConn.Prepare(ctx, "ps_fetch", `fetch all in "exec_statement_cursor"`, nil)
1732+
require.NoError(t, err)
1733+
require.Empty(t, sd.Fields)
1734+
1735+
// DECLARE CURSOR requires an explicit transaction block.
1736+
_, err = pgConn.Exec(ctx, "begin").ReadAll()
1737+
require.NoError(t, err)
1738+
1739+
_, err = pgConn.Exec(ctx, `declare "exec_statement_cursor" cursor for select n, n::text from generate_series(1, 3) n`).ReadAll()
1740+
require.NoError(t, err)
1741+
1742+
result := pgConn.ExecStatement(ctx, sd, nil, nil, nil).Read()
1743+
require.NoError(t, result.Err)
1744+
require.Len(t, result.FieldDescriptions, 2)
1745+
require.Equal(t, uint32(pgtype.Int4OID), result.FieldDescriptions[0].DataTypeOID)
1746+
require.Equal(t, uint32(pgtype.TextOID), result.FieldDescriptions[1].DataTypeOID)
1747+
require.Len(t, result.Rows, 3)
1748+
require.Equal(t, "1", string(result.Rows[0][0]))
1749+
require.Equal(t, "3", string(result.Rows[2][1]))
1750+
1751+
_, err = pgConn.Exec(ctx, "rollback").ReadAll()
1752+
require.NoError(t, err)
1753+
1754+
ensureConnValid(t, pgConn)
1755+
}
1756+
17141757
type byteCounterConn struct {
17151758
conn net.Conn
17161759
bytesRead int
@@ -3682,6 +3725,71 @@ func TestPipelineQueryStatementBindError(t *testing.T) {
36823725
ensureConnValid(t, pgConn)
36833726
}
36843727

3728+
// https://github.com/jackc/pgx/issues/2626
3729+
func TestPipelineQueryStatementCursorFetch(t *testing.T) {
3730+
t.Parallel()
3731+
3732+
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
3733+
defer cancel()
3734+
3735+
pgConn, err := pgconn.Connect(ctx, os.Getenv("PGX_TEST_DATABASE"))
3736+
require.NoError(t, err)
3737+
defer closeConn(t, pgConn)
3738+
3739+
if pgConn.ParameterStatus("crdb_version") != "" {
3740+
t.Skip("Server does not support cursors in implicit transactions")
3741+
}
3742+
3743+
// Prepare the FETCH before the cursor exists. The server describes the result as NoData so the statement
3744+
// description has no fields. The actual fields are only known at execution time.
3745+
sd, err := pgConn.Prepare(ctx, "ps_fetch", `fetch all in "pipeline_cursor"`, nil)
3746+
require.NoError(t, err)
3747+
require.Empty(t, sd.Fields)
3748+
3749+
// DECLARE CURSOR requires an explicit transaction block.
3750+
_, err = pgConn.Exec(ctx, "begin").ReadAll()
3751+
require.NoError(t, err)
3752+
3753+
pipeline := pgConn.StartPipeline(ctx)
3754+
pipeline.SendQueryParams(`declare "pipeline_cursor" cursor for select n, n::text from generate_series(1, 3) n`, nil, nil, nil, nil)
3755+
pipeline.SendQueryStatement(sd, nil, nil, nil)
3756+
err = pipeline.Sync()
3757+
require.NoError(t, err)
3758+
3759+
results, err := pipeline.GetResults()
3760+
require.NoError(t, err)
3761+
rr, ok := results.(*pgconn.ResultReader)
3762+
require.Truef(t, ok, "expected ResultReader, got: %#v", results)
3763+
readResult := rr.Read()
3764+
require.NoError(t, readResult.Err)
3765+
3766+
results, err = pipeline.GetResults()
3767+
require.NoError(t, err)
3768+
rr, ok = results.(*pgconn.ResultReader)
3769+
require.Truef(t, ok, "expected ResultReader, got: %#v", results)
3770+
readResult = rr.Read()
3771+
require.NoError(t, readResult.Err)
3772+
require.Len(t, readResult.FieldDescriptions, 2)
3773+
require.Equal(t, uint32(pgtype.Int4OID), readResult.FieldDescriptions[0].DataTypeOID)
3774+
require.Equal(t, uint32(pgtype.TextOID), readResult.FieldDescriptions[1].DataTypeOID)
3775+
require.Len(t, readResult.Rows, 3)
3776+
require.Equal(t, "1", string(readResult.Rows[0][0]))
3777+
require.Equal(t, "3", string(readResult.Rows[2][1]))
3778+
3779+
results, err = pipeline.GetResults()
3780+
require.NoError(t, err)
3781+
_, ok = results.(*pgconn.PipelineSync)
3782+
require.Truef(t, ok, "expected PipelineSync, got: %#v", results)
3783+
3784+
err = pipeline.Close()
3785+
require.NoError(t, err)
3786+
3787+
_, err = pgConn.Exec(ctx, "rollback").ReadAll()
3788+
require.NoError(t, err)
3789+
3790+
ensureConnValid(t, pgConn)
3791+
}
3792+
36853793
func TestPipelineGetResultsNilResultsOnError(t *testing.T) {
36863794
t.Parallel()
36873795

query_test.go

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,53 @@ func TestConnQueryRowsFieldDescriptionsBeforeNext(t *testing.T) {
7272
assert.Equal(t, "msg", rows.FieldDescriptions()[0].Name)
7373
}
7474

75+
// https://github.com/jackc/pgx/issues/2626
76+
func TestConnQueryPreparedCursorFetch(t *testing.T) {
77+
t.Parallel()
78+
79+
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
80+
defer cancel()
81+
82+
conn := mustConnectString(t, os.Getenv("PGX_TEST_DATABASE"))
83+
defer closeConn(t, conn)
84+
85+
pgxtest.SkipCockroachDB(t, conn, "Server does not support cursors in implicit transactions")
86+
87+
// Prepare the FETCH before the cursor exists. The server describes the result as NoData so the statement
88+
// description has no fields. The actual fields are only known at execution time.
89+
sd, err := conn.Prepare(ctx, "fetch_cursor", `fetch all in "query_cursor"`)
90+
require.NoError(t, err)
91+
require.Empty(t, sd.Fields)
92+
93+
tx, err := conn.Begin(ctx)
94+
require.NoError(t, err)
95+
defer tx.Rollback(ctx)
96+
97+
_, err = tx.Exec(ctx, `declare "query_cursor" cursor for select n, n::text as str from generate_series(1, 3) n`)
98+
require.NoError(t, err)
99+
100+
rows, err := tx.Query(ctx, "fetch_cursor")
101+
require.NoError(t, err)
102+
103+
fds := rows.FieldDescriptions()
104+
require.Len(t, fds, 2)
105+
require.Equal(t, "n", fds[0].Name)
106+
require.Equal(t, "str", fds[1].Name)
107+
108+
var ns []int32
109+
var strs []string
110+
for rows.Next() {
111+
var n int32
112+
var str string
113+
require.NoError(t, rows.Scan(&n, &str))
114+
ns = append(ns, n)
115+
strs = append(strs, str)
116+
}
117+
require.NoError(t, rows.Err())
118+
require.Equal(t, []int32{1, 2, 3}, ns)
119+
require.Equal(t, []string{"1", "2", "3"}, strs)
120+
}
121+
75122
func TestConnQueryWithoutResultSetCommandTag(t *testing.T) {
76123
t.Parallel()
77124

0 commit comments

Comments
 (0)