conn.go 8.5 KB

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