conn.go 7.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367
  1. // Copyright (c) 2012 The gocql Authors. All rights reserved.
  2. // Use of this source code is governed by a BSD-style
  3. // license that can be found in the LICENSE file.
  4. package gocql
  5. import (
  6. "net"
  7. "sync"
  8. "sync/atomic"
  9. "time"
  10. )
  11. const defaultFrameSize = 4096
  12. // Conn is a single connection to a Cassandra node. It can be used to execute
  13. // queries, but users are usually advised to use a more reliable, higher
  14. // level API.
  15. type Conn struct {
  16. conn net.Conn
  17. timeout time.Duration
  18. uniq chan uint8
  19. calls []callReq
  20. nwait int32
  21. prepMu sync.Mutex
  22. prep map[string]*queryInfo
  23. keyspace string
  24. }
  25. // Connect establishes a connection to a Cassandra node.
  26. // You must also call the Serve method before you can execute any queries.
  27. func Connect(addr, version string, timeout time.Duration) (*Conn, error) {
  28. conn, err := net.DialTimeout("tcp", addr, timeout)
  29. if err != nil {
  30. return nil, err
  31. }
  32. c := &Conn{
  33. conn: conn,
  34. uniq: make(chan uint8, 128),
  35. calls: make([]callReq, 128),
  36. prep: make(map[string]*queryInfo),
  37. timeout: timeout,
  38. }
  39. for i := 0; i < cap(c.uniq); i++ {
  40. c.uniq <- uint8(i)
  41. }
  42. if err := c.init(version); err != nil {
  43. return nil, err
  44. }
  45. return c, nil
  46. }
  47. func (c *Conn) init(version string) error {
  48. req := make(frame, headerSize, defaultFrameSize)
  49. req.setHeader(protoRequest, 0, 0, opStartup)
  50. req.writeStringMap(map[string]string{
  51. "CQL_VERSION": version,
  52. })
  53. resp, err := c.callSimple(req)
  54. if err != nil {
  55. return err
  56. } else if resp[3] == opError {
  57. return resp.readErrorFrame()
  58. } else if resp[3] != opReady {
  59. return ErrProtocol
  60. }
  61. /* if cfg.Keyspace != "" {
  62. qry := &Query{stmt: "USE " + cfg.Keyspace}
  63. frame, err = c.executeQuery(qry)
  64. if err != nil {
  65. return err
  66. }
  67. } */
  68. return nil
  69. }
  70. // Serve starts the stream multiplexer for this connection, which is required
  71. // to execute any queries. This method runs as long as the connection is
  72. // open and is therefore usually called in a separate goroutine.
  73. func (c *Conn) Serve() error {
  74. var err error
  75. for {
  76. var frame frame
  77. frame, err = c.recv()
  78. if err != nil {
  79. break
  80. }
  81. c.dispatch(frame)
  82. }
  83. c.conn.Close()
  84. for id := 0; id < len(c.calls); id++ {
  85. req := &c.calls[id]
  86. if atomic.LoadInt32(&req.active) == 1 {
  87. req.resp <- callResp{nil, err}
  88. }
  89. }
  90. return err
  91. }
  92. func (c *Conn) recv() (frame, error) {
  93. resp := make(frame, headerSize, headerSize+512)
  94. c.conn.SetReadDeadline(time.Now().Add(c.timeout))
  95. n, last, pinged := 0, 0, false
  96. for n < len(resp) {
  97. nn, err := c.conn.Read(resp[n:])
  98. n += nn
  99. if err != nil {
  100. if err, ok := err.(net.Error); ok && err.Timeout() {
  101. if n > last {
  102. // we hit the deadline but we made progress.
  103. // simply extend the deadline
  104. c.conn.SetReadDeadline(time.Now().Add(c.timeout))
  105. last = n
  106. } else if n == 0 && !pinged {
  107. c.conn.SetReadDeadline(time.Now().Add(c.timeout))
  108. if atomic.LoadInt32(&c.nwait) > 0 {
  109. go c.ping()
  110. pinged = true
  111. }
  112. } else {
  113. return nil, err
  114. }
  115. } else {
  116. return nil, err
  117. }
  118. }
  119. if n == headerSize && len(resp) == headerSize {
  120. if resp[0] != protoResponse {
  121. return nil, ErrProtocol
  122. }
  123. resp.grow(resp.Length())
  124. }
  125. }
  126. return resp, nil
  127. }
  128. func (c *Conn) callSimple(req frame) (frame, error) {
  129. req.setLength(len(req) - headerSize)
  130. if _, err := c.conn.Write(req); err != nil {
  131. c.conn.Close()
  132. return nil, err
  133. }
  134. return c.recv()
  135. }
  136. func (c *Conn) call(req frame) (frame, error) {
  137. id := <-c.uniq
  138. req[2] = id
  139. call := &c.calls[id]
  140. call.resp = make(chan callResp, 1)
  141. atomic.AddInt32(&c.nwait, 1)
  142. atomic.StoreInt32(&call.active, 1)
  143. req.setLength(len(req) - headerSize)
  144. if _, err := c.conn.Write(req); err != nil {
  145. c.conn.Close()
  146. return nil, err
  147. }
  148. reply := <-call.resp
  149. call.resp = nil
  150. c.uniq <- id
  151. return reply.buf, reply.err
  152. }
  153. func (c *Conn) dispatch(resp frame) {
  154. id := int(resp[2])
  155. if id >= len(c.calls) {
  156. return
  157. }
  158. call := &c.calls[id]
  159. if !atomic.CompareAndSwapInt32(&call.active, 1, 0) {
  160. return
  161. }
  162. atomic.AddInt32(&c.nwait, -1)
  163. call.resp <- callResp{resp, nil}
  164. }
  165. func (c *Conn) ping() error {
  166. req := make(frame, headerSize)
  167. req.setHeader(protoRequest, 0, 0, opOptions)
  168. _, err := c.call(req)
  169. return err
  170. }
  171. func (c *Conn) prepareStatement(stmt string) (*queryInfo, error) {
  172. c.prepMu.Lock()
  173. info := c.prep[stmt]
  174. if info != nil {
  175. c.prepMu.Unlock()
  176. info.wg.Wait()
  177. return info, nil
  178. }
  179. info = new(queryInfo)
  180. info.wg.Add(1)
  181. c.prep[stmt] = info
  182. c.prepMu.Unlock()
  183. frame := make(frame, headerSize, defaultFrameSize)
  184. frame.setHeader(protoRequest, 0, 0, opPrepare)
  185. frame.writeLongString(stmt)
  186. frame.setLength(len(frame) - headerSize)
  187. frame, err := c.call(frame)
  188. if err != nil {
  189. return nil, err
  190. }
  191. if frame[3] == opError {
  192. return nil, frame.readErrorFrame()
  193. }
  194. frame.skipHeader()
  195. frame.readInt() // kind
  196. info.id = frame.readShortBytes()
  197. info.args = frame.readMetaData()
  198. info.rval = frame.readMetaData()
  199. info.wg.Done()
  200. return info, nil
  201. }
  202. func (c *Conn) switchKeyspace(keyspace string) error {
  203. if keyspace == "" || c.keyspace == keyspace {
  204. return nil
  205. }
  206. if _, err := c.ExecuteQuery(&Query{Stmt: "USE " + keyspace}); err != nil {
  207. return err
  208. }
  209. return nil
  210. }
  211. func (c *Conn) ExecuteQuery(qry *Query) (*Iter, error) {
  212. if err := c.switchKeyspace(qry.Keyspace); err != nil {
  213. return nil, err
  214. }
  215. frame, err := c.executeQuery(qry)
  216. if err != nil {
  217. return nil, err
  218. }
  219. if frame[3] == opError {
  220. return nil, frame.readErrorFrame()
  221. } else if frame[3] == opResult {
  222. iter := new(Iter)
  223. iter.readFrame(frame)
  224. return iter, nil
  225. }
  226. return nil, nil
  227. }
  228. func (c *Conn) ExecuteBatch(batch *Batch) error {
  229. frame := make(frame, headerSize, defaultFrameSize)
  230. frame.setHeader(protoRequest, 0, 0, opBatch)
  231. frame.writeByte(byte(batch.Type))
  232. frame.writeShort(uint16(len(batch.Entries)))
  233. for i := 0; i < len(batch.Entries); i++ {
  234. entry := &batch.Entries[i]
  235. var info *queryInfo
  236. if len(entry.Args) > 0 {
  237. info, err := c.prepareStatement(entry.Stmt)
  238. if err != nil {
  239. return err
  240. }
  241. frame.writeByte(1)
  242. frame.writeShortBytes(info.id)
  243. } else {
  244. frame.writeByte(0)
  245. frame.writeLongString(entry.Stmt)
  246. }
  247. frame.writeShort(uint16(len(entry.Args)))
  248. for j := 0; j < len(entry.Args); j++ {
  249. val, err := Marshal(info.args[j].TypeInfo, entry.Args[i])
  250. if err != nil {
  251. return err
  252. }
  253. frame.writeBytes(val)
  254. }
  255. }
  256. frame.writeConsistency(batch.Cons)
  257. frame, err := c.call(frame)
  258. if err != nil {
  259. return err
  260. }
  261. if frame[3] == opError {
  262. return frame.readErrorFrame()
  263. }
  264. return nil
  265. }
  266. func (c *Conn) Close() {
  267. c.conn.Close()
  268. }
  269. func (c *Conn) executeQuery(query *Query) (frame, error) {
  270. var info *queryInfo
  271. if len(query.Args) > 0 {
  272. var err error
  273. info, err = c.prepareStatement(query.Stmt)
  274. if err != nil {
  275. return nil, err
  276. }
  277. }
  278. frame := make(frame, headerSize, defaultFrameSize)
  279. if info == nil {
  280. frame.setHeader(protoRequest, 0, 0, opQuery)
  281. frame.writeLongString(query.Stmt)
  282. } else {
  283. frame.setHeader(protoRequest, 0, 0, opExecute)
  284. frame.writeShortBytes(info.id)
  285. }
  286. frame.writeConsistency(query.Cons)
  287. flags := uint8(0)
  288. if len(query.Args) > 0 {
  289. flags |= flagQueryValues
  290. }
  291. frame.writeByte(flags)
  292. if len(query.Args) > 0 {
  293. frame.writeShort(uint16(len(query.Args)))
  294. for i := 0; i < len(query.Args); i++ {
  295. val, err := Marshal(info.args[i].TypeInfo, query.Args[i])
  296. if err != nil {
  297. return nil, err
  298. }
  299. frame.writeBytes(val)
  300. }
  301. }
  302. frame, err := c.call(frame)
  303. if err != nil {
  304. return nil, err
  305. }
  306. if frame[3] == opError {
  307. frame.skipHeader()
  308. code := frame.readInt()
  309. desc := frame.readString()
  310. return nil, Error{code, desc}
  311. }
  312. return frame, nil
  313. }
  314. type queryInfo struct {
  315. id []byte
  316. args []ColumnInfo
  317. rval []ColumnInfo
  318. wg sync.WaitGroup
  319. }
  320. type callReq struct {
  321. active int32
  322. resp chan callResp
  323. }
  324. type callResp struct {
  325. buf frame
  326. err error
  327. }