From 5490a31909de7cb6bdc5a8cafad29f2281b0b459 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Mi=C5=82osz=20Sm=C3=B3=C5=82ka?= Date: Wed, 13 May 2026 20:12:57 +0200 Subject: [PATCH 1/2] Follow up to ba82ee1ad17497e0e269465beed2910ed69393be Fix WHERE condition priority for PostgreSQL --- pkg/sql/queue_schema_adapter_postgresql.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/sql/queue_schema_adapter_postgresql.go b/pkg/sql/queue_schema_adapter_postgresql.go index 5faa769..eaa2ce0 100644 --- a/pkg/sql/queue_schema_adapter_postgresql.go +++ b/pkg/sql/queue_schema_adapter_postgresql.go @@ -104,7 +104,7 @@ func (s PostgreSQLQueueSchema) SelectQuery(params SelectQueryParams) (Query, err if s.GenerateWhereClause != nil { where, args = s.GenerateWhereClause(whereParams) if where != "" { - where = "AND " + where + where = "AND (" + where + ")" } } From afc6fb34ca722e950a1c47d68ad29fde7754e737 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Mi=C5=82osz=20Sm=C3=B3=C5=82ka?= Date: Wed, 13 May 2026 20:21:51 +0200 Subject: [PATCH 2/2] Fix lint errors --- pkg/sql/adapters_pgx.go | 8 ++++---- pkg/sql/queue_schema_adapter_mysql_test.go | 3 ++- pkg/sql/queue_schema_adapter_postgresql.go | 2 +- pkg/sql/queue_schema_adapter_postgresql_test.go | 3 ++- pkg/sql/schema_adapter_postgresql.go | 2 +- 5 files changed, 10 insertions(+), 8 deletions(-) diff --git a/pkg/sql/adapters_pgx.go b/pkg/sql/adapters_pgx.go index 81db775..67b5422 100644 --- a/pkg/sql/adapters_pgx.go +++ b/pkg/sql/adapters_pgx.go @@ -66,25 +66,25 @@ func (c PgxBeginner) BeginTx(ctx context.Context, options *stdSQL.TxOptions) (Tx } func (c PgxBeginner) ExecContext(ctx context.Context, query string, args ...any) (Result, error) { - res, err := c.Conn.Exec(ctx, query, args...) + res, err := c.Exec(ctx, query, args...) return PgxResult{res}, err } func (c PgxBeginner) QueryContext(ctx context.Context, query string, args ...any) (Rows, error) { - rows, err := c.Conn.Query(ctx, query, args...) + rows, err := c.Query(ctx, query, args...) return PgxRows{rows}, err } func (t PgxTx) ExecContext(ctx context.Context, query string, args ...any) (Result, error) { - res, err := t.Tx.Exec(ctx, query, args...) + res, err := t.Exec(ctx, query, args...) return PgxResult{res}, err } func (t PgxTx) QueryContext(ctx context.Context, query string, args ...any) (Rows, error) { - rows, err := t.Tx.Query(ctx, query, args...) + rows, err := t.Query(ctx, query, args...) return PgxRows{rows}, err } diff --git a/pkg/sql/queue_schema_adapter_mysql_test.go b/pkg/sql/queue_schema_adapter_mysql_test.go index a1b0787..e0c6c8a 100644 --- a/pkg/sql/queue_schema_adapter_mysql_test.go +++ b/pkg/sql/queue_schema_adapter_mysql_test.go @@ -56,6 +56,7 @@ func TestMySQLQueueSchemaAdapter(t *testing.T) { } var receivedMessages []*message.Message +ReceiveLoop: for i := 0; i < 5; i++ { select { case msg := <-messages: @@ -63,7 +64,7 @@ func TestMySQLQueueSchemaAdapter(t *testing.T) { msg.Ack() case <-time.After(5 * time.Second): t.Errorf("expected to receive message") - break + break ReceiveLoop } } diff --git a/pkg/sql/queue_schema_adapter_postgresql.go b/pkg/sql/queue_schema_adapter_postgresql.go index eaa2ce0..974a3ab 100644 --- a/pkg/sql/queue_schema_adapter_postgresql.go +++ b/pkg/sql/queue_schema_adapter_postgresql.go @@ -74,7 +74,7 @@ func queueInsertMarkers(count int) string { index := 1 for i := 0; i < count; i++ { - result.WriteString(fmt.Sprintf("($%d,$%d,$%d),", index, index+1, index+2)) + fmt.Fprintf(&result, "($%d,$%d,$%d),", index, index+1, index+2) index += 3 } diff --git a/pkg/sql/queue_schema_adapter_postgresql_test.go b/pkg/sql/queue_schema_adapter_postgresql_test.go index 20f151b..73a0a0e 100644 --- a/pkg/sql/queue_schema_adapter_postgresql_test.go +++ b/pkg/sql/queue_schema_adapter_postgresql_test.go @@ -56,6 +56,7 @@ func TestPostgreSQLQueueSchemaAdapter(t *testing.T) { } var receivedMessages []*message.Message +ReceiveLoop: for i := 0; i < 5; i++ { select { case msg := <-messages: @@ -63,7 +64,7 @@ func TestPostgreSQLQueueSchemaAdapter(t *testing.T) { msg.Ack() case <-time.After(5 * time.Second): t.Errorf("expected to receive message") - break + break ReceiveLoop } } diff --git a/pkg/sql/schema_adapter_postgresql.go b/pkg/sql/schema_adapter_postgresql.go index 4088cd0..1197e8b 100644 --- a/pkg/sql/schema_adapter_postgresql.go +++ b/pkg/sql/schema_adapter_postgresql.go @@ -95,7 +95,7 @@ func defaultInsertMarkers(count int) string { index := 1 for i := 0; i < count; i++ { - result.WriteString(fmt.Sprintf("($%d,$%d,$%d,pg_current_xact_id()),", index, index+1, index+2)) + fmt.Fprintf(&result, "($%d,$%d,$%d,pg_current_xact_id()),", index, index+1, index+2) index += 3 }