2
0

protocol.go 6.1 KB

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