Jelajahi Sumber

postgres wire

Danilo Fragoso 4 bulan lalu
induk
melakukan
97fc1897c1
9 mengubah file dengan 1224 tambahan dan 3 penghapusan
  1. TEMPAT SAMPAH
      bin/pizzasql
  2. TEMPAT SAMPAH
      bin/sqllogictest
  3. 126 0
      cmd/sqllogictest/main.go
  4. 4 1
      go.mod
  5. 60 0
      go.sum
  6. 78 2
      main.go
  7. 558 0
      pkg/pgserver/connection.go
  8. 251 0
      pkg/pgserver/protocol.go
  9. 147 0
      pkg/pgserver/server.go

TEMPAT SAMPAH
bin/pizzasql


TEMPAT SAMPAH
bin/sqllogictest


+ 126 - 0
cmd/sqllogictest/main.go

@@ -11,7 +11,9 @@ package main
 import (
 	"bufio"
 	"bytes"
+	"context"
 	"crypto/md5"
+	"database/sql"
 	"flag"
 	"fmt"
 	"math"
@@ -24,6 +26,7 @@ import (
 	"time"
 
 	"github.com/goccy/go-json"
+	_ "github.com/lib/pq" // PostgreSQL driver
 )
 
 const engineName = "pizzasql"
@@ -81,6 +84,8 @@ type queryResponse struct {
 type runner struct {
 	baseURL    string
 	client     *http.Client
+	pgDB       *sql.DB // PostgreSQL connection (if using -pg flag)
+	usePG      bool    // Use PostgreSQL wire protocol instead of HTTP
 	verbose    bool
 	stopOnFail bool
 	passed     int
@@ -94,6 +99,10 @@ type runner struct {
 
 func main() {
 	urlFlag := flag.String("url", "http://localhost:8080", "PizzaSQL server URL")
+	pgFlag := flag.Bool("pg", false, "Use PostgreSQL wire protocol instead of HTTP")
+	pgHostFlag := flag.String("pg-host", "localhost", "PostgreSQL server host")
+	pgPortFlag := flag.Int("pg-port", 5432, "PostgreSQL server port")
+	pgDBFlag := flag.String("pg-db", "pizzasql", "PostgreSQL database name")
 	dirFlag := flag.String("dir", "testdata/sqllogictest", "Directory containing .test files")
 	fileFlag := flag.String("file", "", "Single .test file to run (overrides -dir)")
 	verboseFlag := flag.Bool("v", false, "Print each passing record")
@@ -104,11 +113,35 @@ func main() {
 	r := &runner{
 		baseURL:    strings.TrimRight(*urlFlag, "/"),
 		client:     &http.Client{Timeout: 120 * time.Second},
+		usePG:      *pgFlag,
 		verbose:    *verboseFlag,
 		stopOnFail: *stopFlag,
 		logPath:    *logFlag,
 	}
 
+	// If using PostgreSQL wire protocol, establish connection
+	if *pgFlag {
+		connStr := fmt.Sprintf("host=%s port=%d dbname=%s sslmode=disable",
+			*pgHostFlag, *pgPortFlag, *pgDBFlag)
+		db, err := sql.Open("postgres", connStr)
+		if err != nil {
+			fmt.Fprintf(os.Stderr, "Failed to connect to PostgreSQL: %v\n", err)
+			os.Exit(1)
+		}
+		defer db.Close()
+
+		// Test connection
+		ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
+		defer cancel()
+		if err := db.PingContext(ctx); err != nil {
+			fmt.Fprintf(os.Stderr, "Failed to ping PostgreSQL server: %v\n", err)
+			os.Exit(1)
+		}
+
+		r.pgDB = db
+		fmt.Printf("Connected to PostgreSQL at %s:%d (database: %s)\n", *pgHostFlag, *pgPortFlag, *pgDBFlag)
+	}
+
 	if *logFlag != "" {
 		lf, err := os.Create(*logFlag)
 		if err != nil {
@@ -583,6 +616,13 @@ func equalSlices(a, b []string) bool {
 }
 
 func (r *runner) execQuery(sql string) (*queryResponse, error) {
+	if r.usePG {
+		return r.execQueryPG(sql)
+	}
+	return r.execQueryHTTP(sql)
+}
+
+func (r *runner) execQueryHTTP(sql string) (*queryResponse, error) {
 	body, _ := json.Marshal(queryRequest{SQL: sql})
 	resp, err := r.client.Post(r.baseURL+"/query", "application/json", bytes.NewReader(body))
 	if err != nil {
@@ -596,6 +636,92 @@ func (r *runner) execQuery(sql string) (*queryResponse, error) {
 	return &qr, nil
 }
 
+func (r *runner) execQueryPG(sql string) (*queryResponse, error) {
+	ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
+	defer cancel()
+
+	// Check if it's a query or statement
+	sqlUpper := strings.TrimSpace(strings.ToUpper(sql))
+	isSelect := strings.HasPrefix(sqlUpper, "SELECT") ||
+		strings.HasPrefix(sqlUpper, "PRAGMA") ||
+		strings.HasPrefix(sqlUpper, "EXPLAIN")
+
+	var qr queryResponse
+
+	if isSelect {
+		// Execute query and get results
+		rows, err := r.pgDB.QueryContext(ctx, sql)
+		if err != nil {
+			qr.Error = &struct {
+				Code    string `json:"code"`
+				Message string `json:"message"`
+			}{
+				Code:    "QUERY_ERROR",
+				Message: err.Error(),
+			}
+			return &qr, nil
+		}
+		defer rows.Close()
+
+		// Get column information
+		colTypes, err := rows.ColumnTypes()
+		if err != nil {
+			return nil, fmt.Errorf("get column types: %w", err)
+		}
+
+		for _, ct := range colTypes {
+			qr.Columns = append(qr.Columns, struct {
+				Name string `json:"name"`
+				Type string `json:"type"`
+			}{
+				Name: ct.Name(),
+				Type: ct.DatabaseTypeName(),
+			})
+		}
+
+		// Read all rows
+		for rows.Next() {
+			values := make([]interface{}, len(colTypes))
+			valuePtrs := make([]interface{}, len(colTypes))
+			for i := range values {
+				valuePtrs[i] = &values[i]
+			}
+
+			if err := rows.Scan(valuePtrs...); err != nil {
+				return nil, fmt.Errorf("scan row: %w", err)
+			}
+
+			// Convert byte arrays to strings (PostgreSQL returns some types as []byte)
+			for i, v := range values {
+				if b, ok := v.([]byte); ok {
+					values[i] = string(b)
+				}
+			}
+
+			qr.Rows = append(qr.Rows, values)
+		}
+
+		if err := rows.Err(); err != nil {
+			return nil, fmt.Errorf("rows error: %w", err)
+		}
+	} else {
+		// Execute statement (INSERT, UPDATE, DELETE, CREATE, etc.)
+		_, err := r.pgDB.ExecContext(ctx, sql)
+		if err != nil {
+			qr.Error = &struct {
+				Code    string `json:"code"`
+				Message string `json:"message"`
+			}{
+				Code:    "EXEC_ERROR",
+				Message: err.Error(),
+			}
+			return &qr, nil
+		}
+	}
+
+	return &qr, nil
+}
+
 func (r *runner) pass(rec *record) {
 	r.passed++
 	if r.verbose && r.logW != nil {

+ 4 - 1
go.mod

@@ -2,4 +2,7 @@ module github.com/danfragoso/pizzasql-next
 
 go 1.24
 
-require github.com/goccy/go-json v0.10.6 // indirect
+require (
+	github.com/goccy/go-json v0.10.6 // indirect
+	github.com/lib/pq v1.12.3 // indirect
+)

+ 60 - 0
go.sum

@@ -0,0 +1,60 @@
+github.com/bytedance/gopkg v0.1.3 h1:TPBSwH8RsouGCBcMBktLt1AymVo2TVsBVCY4b6TnZ/M=
+github.com/bytedance/gopkg v0.1.3/go.mod h1:576VvJ+eJgyCzdjS+c4+77QF3p7ubbtiKARP3TxducM=
+github.com/bytedance/sonic v1.15.1 h1:nJD5PmM0vY7J8CT6MxoqbVAAMhkSmV2HgRAUrrpLoOw=
+github.com/bytedance/sonic v1.15.1/go.mod h1:mT2NbXunuaEbnZ+mRIX/vYqKISmgEuHFDI4UzmKx2SA=
+github.com/bytedance/sonic/loader v0.5.1 h1:Ygpfa9zwRCCKSlrp5bBP/b/Xzc3VxsAW+5NIYXrOOpI=
+github.com/bytedance/sonic/loader v0.5.1/go.mod h1:AR4NYCk5DdzZizZ5djGqQ92eEhCCcdf5x77udYiSJRo=
+github.com/cloudwego/base64x v0.1.6 h1:t11wG9AECkCDk5fMSoxmufanudBtJ+/HemLstXDLI2M=
+github.com/cloudwego/base64x v0.1.6/go.mod h1:OFcloc187FXDaYHvrNIjxSe8ncn0OOM8gEHfghB2IPU=
+github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
+github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
+github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY=
+github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
+github.com/goccy/go-json v0.10.6 h1:p8HrPJzOakx/mn/bQtjgNjdTcN+/S6FcG2CTtQOrHVU=
+github.com/goccy/go-json v0.10.6/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M=
+github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
+github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
+github.com/klauspost/compress v1.15.15 h1:EF27CXIuDsYJ6mmvtBRlEuB2UVOqHG1tAXgZ7yIO+lw=
+github.com/klauspost/compress v1.15.15/go.mod h1:ZcK2JAFqKOpnBlxcLsJzYfrS9X1akm9fHZNnD9+Vo/4=
+github.com/klauspost/cpuid/v2 v2.2.3 h1:sxCkb+qR91z4vsqw4vGGZlDgPz3G7gjaLyK3V8y70BU=
+github.com/klauspost/cpuid/v2 v2.2.3/go.mod h1:RVVoqg1df56z8g3pUjL/3lE5UfnlrJX8tyFgg4nqhuY=
+github.com/klauspost/cpuid/v2 v2.2.9 h1:66ze0taIn2H33fBvCkXuv9BmCwDfafmiIVpKV9kKGuY=
+github.com/klauspost/cpuid/v2 v2.2.9/go.mod h1:rqkxqrZ1EhYM9G+hXH7YdowN5R5RGN6NK4QwQ3WMXF8=
+github.com/lib/pq v1.12.3 h1:tTWxr2YLKwIvK90ZXEw8GP7UFHtcbTtty8zsI+YjrfQ=
+github.com/lib/pq v1.12.3/go.mod h1:/p+8NSbOcwzAEI7wiMXFlgydTwcgTr3OSKMsD2BitpA=
+github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY=
+github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y=
+github.com/minio/simdjson-go v0.4.5 h1:r4IQwjRGmWCQ2VeMc7fGiilu1z5du0gJ/I/FsKwgo5A=
+github.com/minio/simdjson-go v0.4.5/go.mod h1:eoNz0DcLQRyEDeaPr4Ru6JpjlZPzbA0IodxVJk8lO8E=
+github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w=
+github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls=
+github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
+github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE=
+github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo=
+github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
+github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw=
+github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo=
+github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA=
+github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
+github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU=
+github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo=
+github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY=
+github.com/twitchyliquid64/golang-asm v0.15.1 h1:SU5vSMR7hnwNxj24w34ZyCi/FmDZTkS4MhqMhdFk5YI=
+github.com/twitchyliquid64/golang-asm v0.15.1/go.mod h1:a1lVb/DtPvCB8fslRZhAngC2+aY1QWCk3Cedj/Gdt08=
+golang.org/x/arch v0.0.0-20210923205945-b76863e36670 h1:18EFjUmQOcUvxNYSkA6jO9VAiXCnxFY6NyDX0bHDmkU=
+golang.org/x/arch v0.0.0-20210923205945-b76863e36670/go.mod h1:5om86z9Hs0C8fWVUuoMHwpExlXzs5Tkyp9hOrfG7pp8=
+golang.org/x/sys v0.0.0-20220704084225-05e143d24a9e/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
+golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
+golang.org/x/sys v0.42.0 h1:omrd2nAlyT5ESRdCLYdm3+fMfNFE/+Rf4bDIQImRJeo=
+golang.org/x/sys v0.42.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
+gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
+gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+modernc.org/libc v1.72.0 h1:IEu559v9a0XWjw0DPoVKtXpO2qt5NVLAnFaBbjq+n8c=
+modernc.org/libc v1.72.0/go.mod h1:tTU8DL8A+XLVkEY3x5E/tO7s2Q/q42EtnNWda/L5QhQ=
+modernc.org/mathutil v1.7.1 h1:GCZVGXdaN8gTqB1Mf/usp1Y/hSqgI2vAGGP4jZMCxOU=
+modernc.org/mathutil v1.7.1/go.mod h1:4p5IwJITfppl0G4sUEDtCr4DthTaT47/N3aT6MhfgJg=
+modernc.org/memory v1.11.0 h1:o4QC8aMQzmcwCK3t3Ux/ZHmwFPzE6hf2Y5LbkRs+hbI=
+modernc.org/memory v1.11.0/go.mod h1:/JP4VbVC+K5sU2wZi9bHoq2MAkCnrt2r98UGeSK7Mjw=
+modernc.org/sqlite v1.50.0 h1:eMowQSWLK0MeiQTdmz3lqoF5dqclujdlIKeJA11+7oM=
+modernc.org/sqlite v1.50.0/go.mod h1:m0w8xhwYUVY3H6pSDwc3gkJ/irZT/0YEXwBlhaxQEew=

+ 78 - 2
main.go

@@ -19,6 +19,7 @@ import (
 	"github.com/danfragoso/pizzasql-next/pkg/kvmanager"
 	"github.com/danfragoso/pizzasql-next/pkg/lexer"
 	"github.com/danfragoso/pizzasql-next/pkg/parser"
+	"github.com/danfragoso/pizzasql-next/pkg/pgserver"
 	"github.com/danfragoso/pizzasql-next/pkg/sqlexport"
 	"github.com/danfragoso/pizzasql-next/pkg/sqlimport"
 	"github.com/danfragoso/pizzasql-next/pkg/storage"
@@ -38,9 +39,14 @@ var (
 	httpCORS        = flag.Bool("http-cors", true, "Enable CORS")
 	httpAuth        = flag.Bool("http-auth", false, "Enable authentication")
 	httpCompression = flag.Bool("http-compression", true, "Enable HTTP response compression")
-	httpQuiet       = flag.Bool("quiet", false, "Disable request logging")
+	quiet           = flag.Bool("quiet", false, "Disable request/query logging")
 	apiKeys         = flag.String("api-keys", "", "Comma-separated API keys")
 
+	// PostgreSQL wire protocol server flags
+	pgEnable = flag.Bool("pg", false, "Enable PostgreSQL wire protocol server")
+	pgHost   = flag.String("pg-host", "localhost", "PostgreSQL server host")
+	pgPort   = flag.Int("pg-port", 5432, "PostgreSQL server port")
+
 	// Export/Import flags
 	exportFile   = flag.String("o", "", "Output file for export")
 	importFile   = flag.String("i", "", "Input file for import")
@@ -82,6 +88,12 @@ func main() {
 		return
 	}
 
+	// Check if PostgreSQL server mode is enabled
+	if *pgEnable {
+		runPGServer()
+		return
+	}
+
 	// Check for export command
 	if *exportFile != "" {
 		runExport()
@@ -752,7 +764,7 @@ func runHTTPServer() {
 	config.EnableCORS = *httpCORS
 	config.EnableAuth = *httpAuth
 	config.EnableCompression = *httpCompression
-	config.EnableLogging = !*httpQuiet
+	config.EnableLogging = !*quiet
 
 	if *apiKeys != "" {
 		config.APIKeys = strings.Split(*apiKeys, ",")
@@ -826,6 +838,70 @@ func runHTTPServer() {
 	fmt.Println("Server stopped")
 }
 
+func runPGServer() {
+	// Connect to PizzaKV
+	pool, err := storage.NewKVPool(*kvAddr, *poolSize, *timeout)
+	if err != nil {
+		fmt.Fprintf(os.Stderr, "Failed to connect to PizzaKV at %s: %v\n", *kvAddr, err)
+		fmt.Fprintf(os.Stderr, "Make sure PizzaKV is running: pizzakv\n")
+		os.Exit(1)
+	}
+	defer pool.Close()
+
+	// Create database manager for multi-database support
+	dbManagerConfig := &storage.DatabaseManagerConfig{
+		DefaultDatabase: *database,
+		AutoCreate:      true,
+	}
+	dbManager := storage.NewDatabaseManager(pool, dbManagerConfig)
+
+	// Configure PostgreSQL server
+	config := pgserver.DefaultConfig()
+	config.Host = *pgHost
+	config.Port = *pgPort
+	config.DefaultDatabase = *database
+	config.Quiet = *quiet
+
+	// Create and start server
+	server := pgserver.New(config, dbManager)
+
+	// Handle graceful shutdown
+	stop := make(chan os.Signal, 1)
+	signal.Notify(stop, os.Interrupt, syscall.SIGTERM)
+
+	// Start server in goroutine
+	go func() {
+		if err := server.Start(); err != nil {
+			fmt.Fprintf(os.Stderr, "PostgreSQL server error: %v\n", err)
+			os.Exit(1)
+		}
+	}()
+
+	fmt.Printf("PizzaSQL PostgreSQL Server started on %s:%d\n", *pgHost, *pgPort)
+	fmt.Printf("Default database: %s\n", *database)
+	fmt.Printf("PizzaKV: %s\n", *kvAddr)
+	fmt.Println()
+	fmt.Println("Connect using psql:")
+	fmt.Printf("  psql -h %s -p %d -d %s\n", *pgHost, *pgPort, *database)
+	fmt.Println()
+	fmt.Println("Or any PostgreSQL client library:")
+	fmt.Printf("  postgresql://%s:%d/%s\n", *pgHost, *pgPort, *database)
+	fmt.Println()
+	fmt.Println("Press Ctrl+C to stop")
+
+	<-stop
+	fmt.Println("\nShutting down server...")
+
+	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
+	defer cancel()
+
+	if err := server.Shutdown(ctx); err != nil {
+		fmt.Fprintf(os.Stderr, "Error during shutdown: %v\n", err)
+	}
+
+	fmt.Println("Server stopped")
+}
+
 // launchPizzaKV starts a PizzaKV instance and updates kvAddr
 func launchPizzaKV() error {
 	kvManager = kvmanager.NewManager()

+ 558 - 0
pkg/pgserver/connection.go

@@ -0,0 +1,558 @@
+package pgserver
+
+import (
+	"bufio"
+	"bytes"
+	"encoding/binary"
+	"fmt"
+	"io"
+	"log"
+	"net"
+	"strings"
+
+	"github.com/danfragoso/pizzasql-next/pkg/executor"
+	"github.com/danfragoso/pizzasql-next/pkg/lexer"
+	"github.com/danfragoso/pizzasql-next/pkg/parser"
+	"github.com/danfragoso/pizzasql-next/pkg/storage"
+)
+
+// Connection represents a client connection
+type Connection struct {
+	conn      net.Conn
+	reader    *bufio.Reader
+	writer    *bufio.Writer
+	executor  *executor.Executor
+	schema    *storage.SchemaManager
+	dbManager *storage.DatabaseManager
+	database  string
+	params    map[string]string
+	txStatus  byte
+	quiet     bool // Disable query logging
+}
+
+// NewConnection creates a new connection handler
+func NewConnection(conn net.Conn, dbManager *storage.DatabaseManager, quiet bool) *Connection {
+	return &Connection{
+		conn:      conn,
+		reader:    bufio.NewReader(conn),
+		writer:    bufio.NewWriter(conn),
+		dbManager: dbManager,
+		params:    make(map[string]string),
+		txStatus:  TxStatusIdle,
+		quiet:     quiet,
+	}
+}
+
+// Handle processes the connection
+func (c *Connection) Handle() error {
+	defer c.conn.Close()
+
+	// First, check for SSL request (sent before startup message)
+	// SSL request is 8 bytes: length(4) + code(4) where code = 80877103
+	firstBytes := make([]byte, 8)
+	n, err := io.ReadFull(c.reader, firstBytes)
+	if err != nil {
+		return fmt.Errorf("failed to read initial bytes: %w", err)
+	}
+
+	// Check if it's an SSL request (code 80877103 = 0x04D2162F)
+	if n == 8 {
+		length := binary.BigEndian.Uint32(firstBytes[0:4])
+		code := binary.BigEndian.Uint32(firstBytes[4:8])
+
+		if length == 8 && code == 80877103 {
+			// SSL request - we don't support SSL, send 'N'
+			if !c.quiet {
+				log.Printf("Client requested SSL, sending rejection")
+			}
+			if _, err := c.conn.Write([]byte{'N'}); err != nil {
+				return fmt.Errorf("failed to send SSL rejection: %w", err)
+			}
+			// Now read the actual startup message
+		} else {
+			// Not SSL request, this is part of startup message
+			// We need to prepend these bytes back for ReadStartupMessage
+			// Create a multi-reader that first reads our buffered bytes, then continues with the reader
+			c.reader = bufio.NewReader(io.MultiReader(bytes.NewReader(firstBytes), c.reader))
+		}
+	}
+
+	// Read startup message
+	if !c.quiet {
+		log.Printf("Reading startup message...")
+	}
+	params, err := ReadStartupMessage(c.reader)
+	if err != nil {
+		return fmt.Errorf("failed to read startup message: %w", err)
+	}
+
+	c.params = params
+	if !c.quiet {
+		log.Printf("Startup params: %+v", params)
+	}
+
+	// Get database name from params (default to "pizzasql")
+	dbName := params["database"]
+	if dbName == "" {
+		dbName = "pizzasql"
+	}
+	c.database = dbName
+
+	if !c.quiet {
+		log.Printf("New connection: user=%s database=%s", params["user"], dbName)
+	}
+
+	// Initialize database
+	if err := c.initDatabase(dbName); err != nil {
+		c.sendError("FATAL", ErrCodeConnectionFailure, fmt.Sprintf("Failed to initialize database: %v", err))
+		return err
+	}
+
+	// Send authentication OK (no auth for now)
+	if err := c.sendAuthenticationOk(); err != nil {
+		return err
+	}
+
+	// Send parameter status messages
+	if err := c.sendParameterStatus("server_version", "14.0 (PizzaSQL)"); err != nil {
+		return err
+	}
+	if err := c.sendParameterStatus("server_encoding", "UTF8"); err != nil {
+		return err
+	}
+	if err := c.sendParameterStatus("client_encoding", "UTF8"); err != nil {
+		return err
+	}
+	if err := c.sendParameterStatus("DateStyle", "ISO, MDY"); err != nil {
+		return err
+	}
+	if err := c.sendParameterStatus("TimeZone", "UTC"); err != nil {
+		return err
+	}
+
+	// Send backend key data (for cancellation - we don't implement this yet)
+	if err := c.sendBackendKeyData(12345, 67890); err != nil {
+		return err
+	}
+
+	// Send ready for query
+	if err := c.sendReadyForQuery(); err != nil {
+		return err
+	}
+
+	// Message loop
+	for {
+		msg, err := ReadMessage(c.reader)
+		if err != nil {
+			if err == io.EOF {
+				log.Printf("Connection closed by client")
+				return nil
+			}
+			return fmt.Errorf("failed to read message: %w", err)
+		}
+
+		if err := c.handleMessage(msg); err != nil {
+			if err == io.EOF {
+				// Normal termination
+				log.Printf("Connection closed normally")
+				return nil
+			}
+			log.Printf("Error handling message: %v", err)
+			return err
+		}
+	}
+}
+
+// initDatabase initializes the database connection
+func (c *Connection) initDatabase(dbName string) error {
+	db, err := c.dbManager.GetDatabase(dbName)
+	if err != nil {
+		return err
+	}
+
+	c.schema = db.Schema
+	c.executor = executor.New(db.Schema, db.Table)
+	c.executor.SyncCatalog()
+
+	return nil
+}
+
+// handleMessage processes a client message
+func (c *Connection) handleMessage(msg *Message) error {
+	switch msg.Type {
+	case MsgQuery:
+		return c.handleQuery(msg)
+
+	case MsgTerminate:
+		log.Printf("Client requested termination")
+		return io.EOF
+
+	case MsgParse, MsgBind, MsgDescribe, MsgExecute, MsgSync, MsgClose:
+		// Extended query protocol - not implemented yet
+		c.sendError("ERROR", ErrCodeFeatureNotSupported, "Extended query protocol not yet supported")
+		return c.sendReadyForQuery()
+
+	default:
+		log.Printf("Unknown message type: %c (%d)", msg.Type, msg.Type)
+		c.sendError("ERROR", ErrCodeProtocolViolation, fmt.Sprintf("Unknown message type: %c", msg.Type))
+		return c.sendReadyForQuery()
+	}
+}
+
+// handleQuery processes a simple query
+func (c *Connection) handleQuery(msg *Message) error {
+	// Parse query string (null-terminated)
+	sql := string(msg.Data[:len(msg.Data)-1])
+
+	if !c.quiet {
+		log.Printf("Query: %s", sql)
+	}
+
+	// Handle empty query
+	sqlTrimmed := strings.TrimSpace(sql)
+	if sqlTrimmed == "" || sqlTrimmed == ";" {
+		if err := c.sendEmptyQueryResponse(); err != nil {
+			return err
+		}
+		return c.sendReadyForQuery()
+	}
+
+	// Handle special PostgreSQL system queries that drivers send
+	sqlUpper := strings.ToUpper(strings.TrimSpace(sql))
+
+	// lib/pq and other drivers query these for connection validation
+	if strings.Contains(sqlUpper, "SELECT VERSION()") {
+		// Return a fake PostgreSQL version
+		return c.handleVersionQuery()
+	}
+
+	if strings.Contains(sqlUpper, "SELECT CURRENT_USER") {
+		// Return the current user
+		return c.handleCurrentUserQuery()
+	}
+
+	if strings.Contains(sqlUpper, "SHOW") && (strings.Contains(sqlUpper, "SERVER_VERSION") ||
+		strings.Contains(sqlUpper, "SERVER_ENCODING") ||
+		strings.Contains(sqlUpper, "CLIENT_ENCODING")) {
+		// Handle SHOW commands
+		return c.handleShowCommand(sqlUpper)
+	}
+
+	// Execute query
+	l := lexer.New(sql)
+	p := parser.New(l)
+	stmt, err := p.Parse()
+	if err != nil {
+		c.sendError("ERROR", ErrCodeSyntaxError, fmt.Sprintf("Syntax error: %v", err))
+		return c.sendReadyForQuery()
+	}
+
+	result, err := c.executor.Execute(stmt)
+	if err != nil {
+		c.sendError("ERROR", ErrCodeInternalError, fmt.Sprintf("Execution error: %v", err))
+		return c.sendReadyForQuery()
+	}
+
+	// Send result based on statement type
+	if err := c.sendResult(result, stmt); err != nil {
+		return err
+	}
+
+	return c.sendReadyForQuery()
+}
+
+// sendResult sends query results
+func (c *Connection) sendResult(result *executor.Result, stmt parser.Statement) error {
+	// For SELECT statements, send row description and data rows
+	if _, isSelect := stmt.(*parser.SelectStmt); isSelect && len(result.Columns) > 0 {
+		// Send row description
+		if err := c.sendRowDescription(result.Columns, result.ColumnTypes); err != nil {
+			return err
+		}
+
+		// Send data rows
+		for _, row := range result.Rows {
+			if err := c.sendDataRow(row, result.Columns); err != nil {
+				return err
+			}
+		}
+
+		// Send command complete
+		tag := fmt.Sprintf("SELECT %d", len(result.Rows))
+		return c.sendCommandComplete(tag)
+	}
+
+	// For other statements, just send command complete
+	tag := c.getCommandTag(stmt, result)
+	return c.sendCommandComplete(tag)
+}
+
+// getCommandTag returns the command completion tag
+func (c *Connection) getCommandTag(stmt parser.Statement, result *executor.Result) string {
+	switch stmt.(type) {
+	case *parser.CreateTableStmt:
+		return "CREATE TABLE"
+	case *parser.DropTableStmt:
+		return "DROP TABLE"
+	case *parser.CreateIndexStmt:
+		return "CREATE INDEX"
+	case *parser.DropIndexStmt:
+		return "DROP INDEX"
+	case *parser.InsertStmt:
+		return fmt.Sprintf("INSERT 0 %d", result.RowsAffected)
+	case *parser.UpdateStmt:
+		return fmt.Sprintf("UPDATE %d", result.RowsAffected)
+	case *parser.DeleteStmt:
+		return fmt.Sprintf("DELETE %d", result.RowsAffected)
+	case *parser.BeginStmt:
+		c.txStatus = TxStatusInBlock
+		return "BEGIN"
+	case *parser.CommitStmt:
+		c.txStatus = TxStatusIdle
+		return "COMMIT"
+	case *parser.RollbackStmt:
+		c.txStatus = TxStatusIdle
+		return "ROLLBACK"
+	default:
+		return "OK"
+	}
+}
+
+// sendAuthenticationOk sends authentication OK message
+func (c *Connection) sendAuthenticationOk() error {
+	mb := NewMessageBuilder()
+	mb.WriteInt32(0) // Auth OK
+	return c.writeMessage(MsgAuthenticationOk, mb.Bytes())
+}
+
+// sendParameterStatus sends a parameter status message
+func (c *Connection) sendParameterStatus(name, value string) error {
+	mb := NewMessageBuilder()
+	mb.WriteString(name)
+	mb.WriteString(value)
+	return c.writeMessage(MsgParameterStatus, mb.Bytes())
+}
+
+// sendBackendKeyData sends backend key data
+func (c *Connection) sendBackendKeyData(processID, secretKey int32) error {
+	mb := NewMessageBuilder()
+	mb.WriteInt32(processID)
+	mb.WriteInt32(secretKey)
+	return c.writeMessage(MsgBackendKeyData, mb.Bytes())
+}
+
+// sendReadyForQuery sends ready for query message
+func (c *Connection) sendReadyForQuery() error {
+	mb := NewMessageBuilder()
+	mb.WriteByte(c.txStatus)
+	return c.writeMessage(MsgReadyForQuery, mb.Bytes())
+}
+
+// sendEmptyQueryResponse sends empty query response
+func (c *Connection) sendEmptyQueryResponse() error {
+	return c.writeMessage(MsgEmptyQueryResponse, []byte{})
+}
+
+// sendRowDescription sends row description (column metadata)
+func (c *Connection) sendRowDescription(columns []string, columnTypes []string) error {
+	mb := NewMessageBuilder()
+	mb.WriteInt16(int16(len(columns)))
+
+	for i, col := range columns {
+		colType := ""
+		if i < len(columnTypes) {
+			colType = columnTypes[i]
+		}
+		mb.WriteString(col)
+		mb.WriteInt32(0)                             // table OID
+		mb.WriteInt16(0)                             // column attribute number
+		mb.WriteInt32(c.getOIDForType(colType))      // type OID
+		mb.WriteInt16(c.getTypeSizeForType(colType)) // type size
+		mb.WriteInt32(-1)                            // type modifier
+		mb.WriteInt16(0)                             // format code (text)
+	}
+
+	return c.writeMessage(MsgRowDescription, mb.Bytes())
+}
+
+// sendDataRow sends a data row
+func (c *Connection) sendDataRow(row []interface{}, columns []string) error {
+	mb := NewMessageBuilder()
+	mb.WriteInt16(int16(len(row)))
+
+	for _, value := range row {
+		if value == nil {
+			mb.WriteInt32(-1) // NULL indicator
+			continue
+		}
+
+		// Convert value to string
+		strValue := c.valueToString(value)
+		mb.WriteInt32(int32(len(strValue)))
+		mb.WriteBytes([]byte(strValue))
+	}
+
+	return c.writeMessage(MsgDataRow, mb.Bytes())
+}
+
+// sendCommandComplete sends command complete message
+func (c *Connection) sendCommandComplete(tag string) error {
+	mb := NewMessageBuilder()
+	mb.WriteString(tag)
+	return c.writeMessage(MsgCommandComplete, mb.Bytes())
+}
+
+// sendError sends an error response
+func (c *Connection) sendError(severity, code, message string) error {
+	mb := NewMessageBuilder()
+	mb.WriteByte(ErrorFieldSeverity)
+	mb.WriteString(severity)
+	mb.WriteByte(ErrorFieldCode)
+	mb.WriteString(code)
+	mb.WriteByte(ErrorFieldMessage)
+	mb.WriteString(message)
+	mb.WriteByte(0) // Terminator
+
+	return c.writeMessage(MsgErrorResponse, mb.Bytes())
+}
+
+// writeMessage writes a message to the connection
+func (c *Connection) writeMessage(msgType byte, data []byte) error {
+	if !c.quiet {
+		log.Printf("Sending message type=%c length=%d", msgType, len(data)+4)
+	}
+	if err := WriteMessage(c.writer, msgType, data); err != nil {
+		return err
+	}
+	return c.writer.Flush()
+}
+
+// getOIDForType returns PostgreSQL OID for type
+func (c *Connection) getOIDForType(typeName string) int32 {
+	switch strings.ToUpper(typeName) {
+	case "INTEGER", "INT":
+		return 23 // INT4OID
+	case "TEXT", "VARCHAR", "CHAR":
+		return 25 // TEXTOID
+	case "REAL", "FLOAT":
+		return 700 // FLOAT4OID
+	case "DOUBLE":
+		return 701 // FLOAT8OID
+	case "BOOLEAN", "BOOL":
+		return 16 // BOOLOID
+	case "BLOB":
+		return 17 // BYTEAOID
+	default:
+		return 25 // Default to TEXT
+	}
+}
+
+// getTypeSizeForType returns type size
+func (c *Connection) getTypeSizeForType(typeName string) int16 {
+	switch strings.ToUpper(typeName) {
+	case "INTEGER", "INT":
+		return 4
+	case "REAL", "FLOAT":
+		return 4
+	case "DOUBLE":
+		return 8
+	case "BOOLEAN", "BOOL":
+		return 1
+	default:
+		return -1 // Variable length
+	}
+}
+
+// valueToString converts a value to string
+func (c *Connection) valueToString(value interface{}) string {
+	if value == nil {
+		return ""
+	}
+	return fmt.Sprintf("%v", value)
+}
+
+// handleVersionQuery handles SELECT version()
+func (c *Connection) handleVersionQuery() error {
+	columns := []string{"version"}
+	columnTypes := []string{"TEXT"}
+
+	if err := c.sendRowDescription(columns, columnTypes); err != nil {
+		return err
+	}
+
+	row := []interface{}{"PostgreSQL 14.0 (PizzaSQL)"}
+	if err := c.sendDataRow(row, columns); err != nil {
+		return err
+	}
+
+	if err := c.sendCommandComplete("SELECT 1"); err != nil {
+		return err
+	}
+
+	return c.sendReadyForQuery()
+}
+
+// handleCurrentUserQuery handles SELECT current_user
+func (c *Connection) handleCurrentUserQuery() error {
+	columns := []string{"current_user"}
+	columnTypes := []string{"TEXT"}
+
+	if err := c.sendRowDescription(columns, columnTypes); err != nil {
+		return err
+	}
+
+	user := c.params["user"]
+	if user == "" {
+		user = "pizzasql"
+	}
+
+	row := []interface{}{user}
+	if err := c.sendDataRow(row, columns); err != nil {
+		return err
+	}
+
+	if err := c.sendCommandComplete("SELECT 1"); err != nil {
+		return err
+	}
+
+	return c.sendReadyForQuery()
+}
+
+// handleShowCommand handles SHOW commands
+func (c *Connection) handleShowCommand(sqlUpper string) error {
+	var value string
+	var name string
+
+	if strings.Contains(sqlUpper, "SERVER_VERSION") {
+		name = "server_version"
+		value = "14.0"
+	} else if strings.Contains(sqlUpper, "SERVER_ENCODING") {
+		name = "server_encoding"
+		value = "UTF8"
+	} else if strings.Contains(sqlUpper, "CLIENT_ENCODING") {
+		name = "client_encoding"
+		value = "UTF8"
+	} else {
+		// Unknown SHOW command
+		c.sendError("ERROR", ErrCodeFeatureNotSupported, "SHOW command not supported")
+		return c.sendReadyForQuery()
+	}
+
+	columns := []string{name}
+	columnTypes := []string{"TEXT"}
+
+	if err := c.sendRowDescription(columns, columnTypes); err != nil {
+		return err
+	}
+
+	row := []interface{}{value}
+	if err := c.sendDataRow(row, columns); err != nil {
+		return err
+	}
+
+	if err := c.sendCommandComplete("SHOW"); err != nil {
+		return err
+	}
+
+	return c.sendReadyForQuery()
+}

+ 251 - 0
pkg/pgserver/protocol.go

@@ -0,0 +1,251 @@
+package pgserver
+
+import (
+	"encoding/binary"
+	"fmt"
+	"io"
+)
+
+// Message type constants (first byte of message)
+const (
+	// Frontend (client) messages
+	MsgStartup   = 0   // Startup message (no type byte)
+	MsgQuery     = 'Q' // Simple query
+	MsgTerminate = 'X' // Terminate
+	MsgPassword  = 'p' // Password message
+	MsgParse     = 'P' // Parse (prepared statement)
+	MsgBind      = 'B' // Bind
+	MsgDescribe  = 'D' // Describe
+	MsgExecute   = 'E' // Execute
+	MsgSync      = 'S' // Sync
+	MsgFlush     = 'H' // Flush
+	MsgClose     = 'C' // Close
+
+	// Backend (server) messages
+	MsgAuthenticationOk     = 'R' // Authentication request
+	MsgBackendKeyData       = 'K' // Backend key data
+	MsgBindComplete         = '2' // Bind complete
+	MsgCloseComplete        = '3' // Close complete
+	MsgCommandComplete      = 'C' // Command complete
+	MsgDataRow              = 'D' // Data row
+	MsgEmptyQueryResponse   = 'I' // Empty query response
+	MsgErrorResponse        = 'E' // Error response
+	MsgNoData               = 'n' // No data
+	MsgNoticeResponse       = 'N' // Notice response
+	MsgParameterDescription = 't' // Parameter description
+	MsgParameterStatus      = 'S' // Parameter status
+	MsgParseComplete        = '1' // Parse complete
+	MsgReadyForQuery        = 'Z' // Ready for query
+	MsgRowDescription       = 'T' // Row description
+	MsgNotificationResponse = 'A' // Notification response
+)
+
+// Transaction status
+const (
+	TxStatusIdle    = 'I' // Idle (not in transaction)
+	TxStatusInBlock = 'T' // In transaction block
+	TxStatusFailed  = 'E' // In failed transaction block
+)
+
+// Error field types
+const (
+	ErrorFieldSeverity         = 'S'
+	ErrorFieldCode             = 'C'
+	ErrorFieldMessage          = 'M'
+	ErrorFieldDetail           = 'D'
+	ErrorFieldHint             = 'H'
+	ErrorFieldPosition         = 'P'
+	ErrorFieldInternalPosition = 'p'
+	ErrorFieldInternalQuery    = 'q'
+	ErrorFieldWhere            = 'W'
+	ErrorFieldSchemaName       = 's'
+	ErrorFieldTableName        = 't'
+	ErrorFieldColumnName       = 'c'
+	ErrorFieldDataTypeName     = 'd'
+	ErrorFieldConstraintName   = 'n'
+	ErrorFieldFile             = 'F'
+	ErrorFieldLine             = 'L'
+	ErrorFieldRoutine          = 'R'
+)
+
+// PostgreSQL error codes (subset)
+const (
+	ErrCodeSuccess             = "00000"
+	ErrCodeSyntaxError         = "42601"
+	ErrCodeUndefinedTable      = "42P01"
+	ErrCodeUndefinedColumn     = "42703"
+	ErrCodeDuplicateTable      = "42P07"
+	ErrCodeDuplicateColumn     = "42701"
+	ErrCodeInvalidParameter    = "22023"
+	ErrCodeInternalError       = "XX000"
+	ErrCodeConnectionFailure   = "08006"
+	ErrCodeProtocolViolation   = "08P01"
+	ErrCodeFeatureNotSupported = "0A000"
+)
+
+// Message represents a PostgreSQL protocol message
+type Message struct {
+	Type byte
+	Data []byte
+}
+
+// WriteMessage writes a message to the writer
+func WriteMessage(w io.Writer, msgType byte, data []byte) error {
+	// Write message type
+	if _, err := w.Write([]byte{msgType}); err != nil {
+		return err
+	}
+
+	// Write message length (includes itself, 4 bytes)
+	length := uint32(len(data) + 4)
+	if err := binary.Write(w, binary.BigEndian, length); err != nil {
+		return err
+	}
+
+	// Write message data
+	if _, err := w.Write(data); err != nil {
+		return err
+	}
+
+	return nil
+}
+
+// ReadMessage reads a message from the reader
+func ReadMessage(r io.Reader) (*Message, error) {
+	// Read message type
+	typeBuf := make([]byte, 1)
+	if _, err := io.ReadFull(r, typeBuf); err != nil {
+		return nil, err
+	}
+
+	// Read message length
+	var length uint32
+	if err := binary.Read(r, binary.BigEndian, &length); err != nil {
+		return nil, err
+	}
+
+	if length < 4 {
+		return nil, fmt.Errorf("invalid message length: %d", length)
+	}
+
+	// Read message data
+	data := make([]byte, length-4)
+	if _, err := io.ReadFull(r, data); err != nil {
+		return nil, err
+	}
+
+	return &Message{
+		Type: typeBuf[0],
+		Data: data,
+	}, nil
+}
+
+// ReadStartupMessage reads the initial startup message (no type byte)
+func ReadStartupMessage(r io.Reader) (map[string]string, error) {
+	// Read message length
+	var length uint32
+	if err := binary.Read(r, binary.BigEndian, &length); err != nil {
+		return nil, err
+	}
+
+	if length < 8 {
+		return nil, fmt.Errorf("invalid startup message length: %d", length)
+	}
+
+	// Read protocol version
+	var version uint32
+	if err := binary.Read(r, binary.BigEndian, &version); err != nil {
+		return nil, err
+	}
+
+	// Read parameters
+	data := make([]byte, length-8)
+	if _, err := io.ReadFull(r, data); err != nil {
+		return nil, err
+	}
+
+	params := make(map[string]string)
+	params["protocol_version"] = fmt.Sprintf("%d", version)
+
+	// Parse null-terminated key-value pairs
+	i := 0
+	for i < len(data) {
+		if data[i] == 0 {
+			break
+		}
+
+		// Read key
+		keyStart := i
+		for i < len(data) && data[i] != 0 {
+			i++
+		}
+		if i >= len(data) {
+			break
+		}
+		key := string(data[keyStart:i])
+		i++ // skip null
+
+		// Read value
+		valueStart := i
+		for i < len(data) && data[i] != 0 {
+			i++
+		}
+		if i > len(data) {
+			break
+		}
+		value := string(data[valueStart:i])
+		i++ // skip null
+
+		params[key] = value
+	}
+
+	return params, nil
+}
+
+// MessageBuilder helps build protocol messages
+type MessageBuilder struct {
+	data []byte
+}
+
+// NewMessageBuilder creates a new message builder
+func NewMessageBuilder() *MessageBuilder {
+	return &MessageBuilder{
+		data: make([]byte, 0, 1024),
+	}
+}
+
+// WriteByte writes a single byte
+func (mb *MessageBuilder) WriteByte(b byte) {
+	mb.data = append(mb.data, b)
+}
+
+// WriteInt16 writes a 16-bit integer
+func (mb *MessageBuilder) WriteInt16(n int16) {
+	mb.data = append(mb.data, byte(n>>8), byte(n))
+}
+
+// WriteInt32 writes a 32-bit integer
+func (mb *MessageBuilder) WriteInt32(n int32) {
+	mb.data = append(mb.data, byte(n>>24), byte(n>>16), byte(n>>8), byte(n))
+}
+
+// WriteString writes a null-terminated string
+func (mb *MessageBuilder) WriteString(s string) {
+	mb.data = append(mb.data, []byte(s)...)
+	mb.data = append(mb.data, 0)
+}
+
+// WriteBytes writes raw bytes
+func (mb *MessageBuilder) WriteBytes(b []byte) {
+	mb.data = append(mb.data, b...)
+}
+
+// Bytes returns the built message data
+func (mb *MessageBuilder) Bytes() []byte {
+	return mb.data
+}
+
+// Reset resets the builder for reuse
+func (mb *MessageBuilder) Reset() {
+	mb.data = mb.data[:0]
+}

+ 147 - 0
pkg/pgserver/server.go

@@ -0,0 +1,147 @@
+package pgserver
+
+import (
+	"context"
+	"fmt"
+	"log"
+	"net"
+	"sync"
+	"time"
+
+	"github.com/danfragoso/pizzasql-next/pkg/storage"
+)
+
+// Config holds server configuration
+type Config struct {
+	Host            string
+	Port            int
+	MaxConnections  int
+	ReadTimeout     time.Duration
+	WriteTimeout    time.Duration
+	DefaultDatabase string
+	Quiet           bool // Disable query logging
+}
+
+// DefaultConfig returns default configuration
+func DefaultConfig() *Config {
+	return &Config{
+		Host:            "localhost",
+		Port:            5432,
+		MaxConnections:  100,
+		ReadTimeout:     30 * time.Second,
+		WriteTimeout:    30 * time.Second,
+		DefaultDatabase: "pizzasql",
+	}
+}
+
+// Server represents the PostgreSQL wire protocol server
+type Server struct {
+	config    *Config
+	dbManager *storage.DatabaseManager
+	listener  net.Listener
+	ctx       context.Context
+	cancel    context.CancelFunc
+	wg        sync.WaitGroup
+}
+
+// New creates a new PostgreSQL wire protocol server
+func New(config *Config, dbManager *storage.DatabaseManager) *Server {
+	if config == nil {
+		config = DefaultConfig()
+	}
+
+	ctx, cancel := context.WithCancel(context.Background())
+
+	return &Server{
+		config:    config,
+		dbManager: dbManager,
+		ctx:       ctx,
+		cancel:    cancel,
+	}
+}
+
+// Start starts the server
+func (s *Server) Start() error {
+	addr := fmt.Sprintf("%s:%d", s.config.Host, s.config.Port)
+
+	listener, err := net.Listen("tcp", addr)
+	if err != nil {
+		return fmt.Errorf("failed to start PostgreSQL server: %w", err)
+	}
+
+	s.listener = listener
+	log.Printf("PostgreSQL wire protocol server listening on %s", addr)
+
+	// Accept connections
+	for {
+		conn, err := listener.Accept()
+		if err != nil {
+			select {
+			case <-s.ctx.Done():
+				// Server is shutting down
+				return nil
+			default:
+				log.Printf("Error accepting connection: %v", err)
+				continue
+			}
+		}
+
+		// Handle connection in a goroutine
+		s.wg.Add(1)
+		go func() {
+			defer s.wg.Done()
+			s.handleConnection(conn)
+		}()
+	}
+}
+
+// handleConnection handles a client connection
+func (s *Server) handleConnection(conn net.Conn) {
+	// Don't set static deadlines - let the connection be persistent
+	// Timeouts will be handled by context cancellation if needed
+
+	c := NewConnection(conn, s.dbManager, s.config.Quiet)
+
+	if err := c.Handle(); err != nil {
+		log.Printf("Connection error: %v", err)
+	}
+}
+
+// Shutdown gracefully shuts down the server
+func (s *Server) Shutdown(ctx context.Context) error {
+	log.Println("Shutting down PostgreSQL server...")
+
+	// Cancel context to stop accepting new connections
+	s.cancel()
+
+	// Close listener
+	if s.listener != nil {
+		if err := s.listener.Close(); err != nil {
+			log.Printf("Error closing listener: %v", err)
+		}
+	}
+
+	// Wait for connections to finish with timeout
+	done := make(chan struct{})
+	go func() {
+		s.wg.Wait()
+		close(done)
+	}()
+
+	select {
+	case <-done:
+		log.Println("PostgreSQL server stopped gracefully")
+		return nil
+	case <-ctx.Done():
+		log.Println("PostgreSQL server shutdown timeout")
+		return ctx.Err()
+	}
+}
+
+// Addr returns the server address
+func (s *Server) Addr() string {
+	if s.listener != nil {
+		return s.listener.Addr().String()
+	}
+	return fmt.Sprintf("%s:%d", s.config.Host, s.config.Port)
+}