conn.go 8.5 KB

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