protocol.go 6.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268
  1. package pgserver
  2. import (
  3. "encoding/binary"
  4. "fmt"
  5. "io"
  6. )
  7. // Maximum message sizes enforced before any allocation, to bound memory use
  8. // against malformed or hostile clients. Regular protocol messages (queries,
  9. // bind parameters, etc.) are capped at 16 MiB; startup messages are much
  10. // smaller and capped at 1 MiB.
  11. const (
  12. MaxMessageSize = 16 * 1024 * 1024 // 16 MiB
  13. MaxStartupMessageSize = 1 * 1024 * 1024 // 1 MiB
  14. )
  15. // Message type constants (first byte of message)
  16. const (
  17. // Frontend (client) messages
  18. MsgStartup = 0 // Startup message (no type byte)
  19. MsgQuery = 'Q' // Simple query
  20. MsgTerminate = 'X' // Terminate
  21. MsgPassword = 'p' // Password message
  22. MsgParse = 'P' // Parse (prepared statement)
  23. MsgBind = 'B' // Bind
  24. MsgDescribe = 'D' // Describe
  25. MsgExecute = 'E' // Execute
  26. MsgSync = 'S' // Sync
  27. MsgFlush = 'H' // Flush
  28. MsgClose = 'C' // Close
  29. // Backend (server) messages
  30. MsgAuthenticationOk = 'R' // Authentication request
  31. MsgBackendKeyData = 'K' // Backend key data
  32. MsgBindComplete = '2' // Bind complete
  33. MsgCloseComplete = '3' // Close complete
  34. MsgCommandComplete = 'C' // Command complete
  35. MsgDataRow = 'D' // Data row
  36. MsgEmptyQueryResponse = 'I' // Empty query response
  37. MsgErrorResponse = 'E' // Error response
  38. MsgNoData = 'n' // No data
  39. MsgNoticeResponse = 'N' // Notice response
  40. MsgParameterDescription = 't' // Parameter description
  41. MsgParameterStatus = 'S' // Parameter status
  42. MsgParseComplete = '1' // Parse complete
  43. MsgReadyForQuery = 'Z' // Ready for query
  44. MsgRowDescription = 'T' // Row description
  45. MsgNotificationResponse = 'A' // Notification response
  46. )
  47. // Transaction status
  48. const (
  49. TxStatusIdle = 'I' // Idle (not in transaction)
  50. TxStatusInBlock = 'T' // In transaction block
  51. TxStatusFailed = 'E' // In failed transaction block
  52. )
  53. // Error field types
  54. const (
  55. ErrorFieldSeverity = 'S'
  56. ErrorFieldCode = 'C'
  57. ErrorFieldMessage = 'M'
  58. ErrorFieldDetail = 'D'
  59. ErrorFieldHint = 'H'
  60. ErrorFieldPosition = 'P'
  61. ErrorFieldInternalPosition = 'p'
  62. ErrorFieldInternalQuery = 'q'
  63. ErrorFieldWhere = 'W'
  64. ErrorFieldSchemaName = 's'
  65. ErrorFieldTableName = 't'
  66. ErrorFieldColumnName = 'c'
  67. ErrorFieldDataTypeName = 'd'
  68. ErrorFieldConstraintName = 'n'
  69. ErrorFieldFile = 'F'
  70. ErrorFieldLine = 'L'
  71. ErrorFieldRoutine = 'R'
  72. )
  73. // PostgreSQL error codes (subset)
  74. const (
  75. ErrCodeSuccess = "00000"
  76. ErrCodeSyntaxError = "42601"
  77. ErrCodeUndefinedTable = "42P01"
  78. ErrCodeUndefinedColumn = "42703"
  79. ErrCodeDuplicateTable = "42P07"
  80. ErrCodeDuplicateColumn = "42701"
  81. ErrCodeInvalidParameter = "22023"
  82. ErrCodeInternalError = "XX000"
  83. ErrCodeConnectionFailure = "08006"
  84. ErrCodeProtocolViolation = "08P01"
  85. ErrCodeFeatureNotSupported = "0A000"
  86. ErrCodeTransactionAborted = "25P02"
  87. ErrCodeSerializationFailure = "40001"
  88. )
  89. // Message represents a PostgreSQL protocol message
  90. type Message struct {
  91. Type byte
  92. Data []byte
  93. }
  94. // WriteMessage writes a message to the writer
  95. func WriteMessage(w io.Writer, msgType byte, data []byte) error {
  96. // Write message type
  97. if _, err := w.Write([]byte{msgType}); err != nil {
  98. return err
  99. }
  100. // Write message length (includes itself, 4 bytes)
  101. length := uint32(len(data) + 4)
  102. if err := binary.Write(w, binary.BigEndian, length); err != nil {
  103. return err
  104. }
  105. // Write message data
  106. if _, err := w.Write(data); err != nil {
  107. return err
  108. }
  109. return nil
  110. }
  111. // ReadMessage reads a message from the reader
  112. func ReadMessage(r io.Reader) (*Message, error) {
  113. // Read message type
  114. typeBuf := make([]byte, 1)
  115. if _, err := io.ReadFull(r, typeBuf); err != nil {
  116. return nil, err
  117. }
  118. // Read message length
  119. var length uint32
  120. if err := binary.Read(r, binary.BigEndian, &length); err != nil {
  121. return nil, err
  122. }
  123. if length < 4 {
  124. return nil, fmt.Errorf("invalid message length: %d", length)
  125. }
  126. if length > MaxMessageSize {
  127. return nil, fmt.Errorf("message length %d exceeds maximum %d", length, MaxMessageSize)
  128. }
  129. // Read message data
  130. data := make([]byte, length-4)
  131. if _, err := io.ReadFull(r, data); err != nil {
  132. return nil, err
  133. }
  134. return &Message{
  135. Type: typeBuf[0],
  136. Data: data,
  137. }, nil
  138. }
  139. // ReadStartupMessage reads the initial startup message (no type byte)
  140. func ReadStartupMessage(r io.Reader) (map[string]string, error) {
  141. // Read message length
  142. var length uint32
  143. if err := binary.Read(r, binary.BigEndian, &length); err != nil {
  144. return nil, err
  145. }
  146. if length < 8 {
  147. return nil, fmt.Errorf("invalid startup message length: %d", length)
  148. }
  149. if length > MaxStartupMessageSize {
  150. return nil, fmt.Errorf("startup message length %d exceeds maximum %d", length, MaxStartupMessageSize)
  151. }
  152. // Read protocol version
  153. var version uint32
  154. if err := binary.Read(r, binary.BigEndian, &version); err != nil {
  155. return nil, err
  156. }
  157. // Read parameters
  158. data := make([]byte, length-8)
  159. if _, err := io.ReadFull(r, data); err != nil {
  160. return nil, err
  161. }
  162. params := make(map[string]string)
  163. params["protocol_version"] = fmt.Sprintf("%d", version)
  164. // Parse null-terminated key-value pairs
  165. i := 0
  166. for i < len(data) {
  167. if data[i] == 0 {
  168. break
  169. }
  170. // Read key
  171. keyStart := i
  172. for i < len(data) && data[i] != 0 {
  173. i++
  174. }
  175. if i >= len(data) {
  176. break
  177. }
  178. key := string(data[keyStart:i])
  179. i++ // skip null
  180. // Read value
  181. valueStart := i
  182. for i < len(data) && data[i] != 0 {
  183. i++
  184. }
  185. if i > len(data) {
  186. break
  187. }
  188. value := string(data[valueStart:i])
  189. i++ // skip null
  190. params[key] = value
  191. }
  192. return params, nil
  193. }
  194. // MessageBuilder helps build protocol messages
  195. type MessageBuilder struct {
  196. data []byte
  197. }
  198. // NewMessageBuilder creates a new message builder
  199. func NewMessageBuilder() *MessageBuilder {
  200. return &MessageBuilder{
  201. data: make([]byte, 0, 1024),
  202. }
  203. }
  204. // AppendByte appends a single byte.
  205. func (mb *MessageBuilder) AppendByte(b byte) {
  206. mb.data = append(mb.data, b)
  207. }
  208. // WriteInt16 writes a 16-bit integer
  209. func (mb *MessageBuilder) WriteInt16(n int16) {
  210. mb.data = append(mb.data, byte(n>>8), byte(n))
  211. }
  212. // WriteInt32 writes a 32-bit integer
  213. func (mb *MessageBuilder) WriteInt32(n int32) {
  214. mb.data = append(mb.data, byte(n>>24), byte(n>>16), byte(n>>8), byte(n))
  215. }
  216. // WriteString writes a null-terminated string
  217. func (mb *MessageBuilder) WriteString(s string) {
  218. mb.data = append(mb.data, []byte(s)...)
  219. mb.data = append(mb.data, 0)
  220. }
  221. // WriteBytes writes raw bytes
  222. func (mb *MessageBuilder) WriteBytes(b []byte) {
  223. mb.data = append(mb.data, b...)
  224. }
  225. // Bytes returns the built message data
  226. func (mb *MessageBuilder) Bytes() []byte {
  227. return mb.data
  228. }
  229. // Reset resets the builder for reuse
  230. func (mb *MessageBuilder) Reset() {
  231. mb.data = mb.data[:0]
  232. }