read.go 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639
  1. package kafka
  2. import (
  3. "bufio"
  4. "errors"
  5. "fmt"
  6. "io"
  7. "reflect"
  8. )
  9. type readable interface {
  10. readFrom(*bufio.Reader, int) (int, error)
  11. }
  12. var errShortRead = errors.New("not enough bytes available to load the response")
  13. func peekRead(r *bufio.Reader, sz int, n int, f func([]byte)) (int, error) {
  14. if n > sz {
  15. return sz, errShortRead
  16. }
  17. b, err := r.Peek(n)
  18. if err != nil {
  19. return sz, err
  20. }
  21. f(b)
  22. return discardN(r, sz, n)
  23. }
  24. func readInt8(r *bufio.Reader, sz int, v *int8) (int, error) {
  25. return peekRead(r, sz, 1, func(b []byte) { *v = makeInt8(b) })
  26. }
  27. func readInt16(r *bufio.Reader, sz int, v *int16) (int, error) {
  28. return peekRead(r, sz, 2, func(b []byte) { *v = makeInt16(b) })
  29. }
  30. func readInt32(r *bufio.Reader, sz int, v *int32) (int, error) {
  31. return peekRead(r, sz, 4, func(b []byte) { *v = makeInt32(b) })
  32. }
  33. func readInt64(r *bufio.Reader, sz int, v *int64) (int, error) {
  34. return peekRead(r, sz, 8, func(b []byte) { *v = makeInt64(b) })
  35. }
  36. func readVarInt(r *bufio.Reader, sz int, v *int64) (remain int, err error) {
  37. // Optimistically assume that most of the time, there will be data buffered
  38. // in the reader. If this is not the case, the buffer will be refilled after
  39. // consuming zero bytes from the input.
  40. input, _ := r.Peek(r.Buffered())
  41. x := uint64(0)
  42. s := uint(0)
  43. for {
  44. if len(input) > sz {
  45. input = input[:sz]
  46. }
  47. for i, b := range input {
  48. if b < 0x80 {
  49. x |= uint64(b) << s
  50. *v = int64(x>>1) ^ -(int64(x) & 1)
  51. n, err := r.Discard(i + 1)
  52. return sz - n, err
  53. }
  54. x |= uint64(b&0x7f) << s
  55. s += 7
  56. }
  57. // Make room in the input buffer to load more data from the underlying
  58. // stream. The x and s variables are left untouched, ensuring that the
  59. // varint decoding can continue on the next loop iteration.
  60. n, _ := r.Discard(len(input))
  61. sz -= n
  62. if sz == 0 {
  63. return 0, errShortRead
  64. }
  65. // Fill the buffer: ask for one more byte, but in practice the reader
  66. // will load way more from the underlying stream.
  67. if _, err := r.Peek(1); err != nil {
  68. if err == io.EOF {
  69. err = errShortRead
  70. }
  71. return sz, err
  72. }
  73. // Grab as many bytes as possible from the buffer, then go on to the
  74. // next loop iteration which is going to consume it.
  75. input, _ = r.Peek(r.Buffered())
  76. }
  77. }
  78. func readBool(r *bufio.Reader, sz int, v *bool) (int, error) {
  79. return peekRead(r, sz, 1, func(b []byte) { *v = b[0] != 0 })
  80. }
  81. func readString(r *bufio.Reader, sz int, v *string) (int, error) {
  82. return readStringWith(r, sz, func(r *bufio.Reader, sz int, n int) (remain int, err error) {
  83. *v, remain, err = readNewString(r, sz, n)
  84. return
  85. })
  86. }
  87. func readStringWith(r *bufio.Reader, sz int, cb func(*bufio.Reader, int, int) (int, error)) (int, error) {
  88. var err error
  89. var len int16
  90. if sz, err = readInt16(r, sz, &len); err != nil {
  91. return sz, err
  92. }
  93. n := int(len)
  94. if n > sz {
  95. return sz, errShortRead
  96. }
  97. return cb(r, sz, n)
  98. }
  99. func readNewString(r *bufio.Reader, sz int, n int) (string, int, error) {
  100. b, sz, err := readNewBytes(r, sz, n)
  101. return string(b), sz, err
  102. }
  103. func readBytes(r *bufio.Reader, sz int, v *[]byte) (int, error) {
  104. return readBytesWith(r, sz, func(r *bufio.Reader, sz int, n int) (remain int, err error) {
  105. *v, remain, err = readNewBytes(r, sz, n)
  106. return
  107. })
  108. }
  109. func readBytesWith(r *bufio.Reader, sz int, cb func(*bufio.Reader, int, int) (int, error)) (int, error) {
  110. var err error
  111. var n int
  112. if sz, err = readArrayLen(r, sz, &n); err != nil {
  113. return sz, err
  114. }
  115. if n > sz {
  116. return sz, errShortRead
  117. }
  118. return cb(r, sz, n)
  119. }
  120. func readNewBytes(r *bufio.Reader, sz int, n int) ([]byte, int, error) {
  121. var err error
  122. var b []byte
  123. var shortRead bool
  124. if n > 0 {
  125. if sz < n {
  126. n = sz
  127. shortRead = true
  128. }
  129. b = make([]byte, n)
  130. n, err = io.ReadFull(r, b)
  131. b = b[:n]
  132. sz -= n
  133. if err == nil && shortRead {
  134. err = errShortRead
  135. }
  136. }
  137. return b, sz, err
  138. }
  139. func readArrayLen(r *bufio.Reader, sz int, n *int) (int, error) {
  140. var err error
  141. var len int32
  142. if sz, err = readInt32(r, sz, &len); err != nil {
  143. return sz, err
  144. }
  145. *n = int(len)
  146. return sz, nil
  147. }
  148. func readArrayWith(r *bufio.Reader, sz int, cb func(*bufio.Reader, int) (int, error)) (int, error) {
  149. var err error
  150. var len int32
  151. if sz, err = readInt32(r, sz, &len); err != nil {
  152. return sz, err
  153. }
  154. for n := int(len); n > 0; n-- {
  155. if sz, err = cb(r, sz); err != nil {
  156. break
  157. }
  158. }
  159. return sz, err
  160. }
  161. func readStringArray(r *bufio.Reader, sz int, v *[]string) (remain int, err error) {
  162. var content []string
  163. fn := func(r *bufio.Reader, size int) (fnRemain int, fnErr error) {
  164. var value string
  165. if fnRemain, fnErr = readString(r, size, &value); fnErr != nil {
  166. return
  167. }
  168. content = append(content, value)
  169. return
  170. }
  171. if remain, err = readArrayWith(r, sz, fn); err != nil {
  172. return
  173. }
  174. *v = content
  175. return
  176. }
  177. func readMapStringInt32(r *bufio.Reader, sz int, v *map[string][]int32) (remain int, err error) {
  178. var len int32
  179. if remain, err = readInt32(r, sz, &len); err != nil {
  180. return
  181. }
  182. content := make(map[string][]int32, len)
  183. for i := 0; i < int(len); i++ {
  184. var key string
  185. var values []int32
  186. if remain, err = readString(r, remain, &key); err != nil {
  187. return
  188. }
  189. fn := func(r *bufio.Reader, size int) (fnRemain int, fnErr error) {
  190. var value int32
  191. if fnRemain, fnErr = readInt32(r, size, &value); fnErr != nil {
  192. return
  193. }
  194. values = append(values, value)
  195. return
  196. }
  197. if remain, err = readArrayWith(r, remain, fn); err != nil {
  198. return
  199. }
  200. content[key] = values
  201. }
  202. *v = content
  203. return
  204. }
  205. func read(r *bufio.Reader, sz int, a interface{}) (int, error) {
  206. switch v := a.(type) {
  207. case *int8:
  208. return readInt8(r, sz, v)
  209. case *int16:
  210. return readInt16(r, sz, v)
  211. case *int32:
  212. return readInt32(r, sz, v)
  213. case *int64:
  214. return readInt64(r, sz, v)
  215. case *bool:
  216. return readBool(r, sz, v)
  217. case *string:
  218. return readString(r, sz, v)
  219. case *[]byte:
  220. return readBytes(r, sz, v)
  221. }
  222. switch v := reflect.ValueOf(a).Elem(); v.Kind() {
  223. case reflect.Struct:
  224. return readStruct(r, sz, v)
  225. case reflect.Slice:
  226. return readSlice(r, sz, v)
  227. default:
  228. panic(fmt.Sprintf("unsupported type: %T", a))
  229. }
  230. }
  231. func readAll(r *bufio.Reader, sz int, ptrs ...interface{}) (int, error) {
  232. var err error
  233. for _, ptr := range ptrs {
  234. if sz, err = readPtr(r, sz, ptr); err != nil {
  235. break
  236. }
  237. }
  238. return sz, err
  239. }
  240. func readPtr(r *bufio.Reader, sz int, ptr interface{}) (int, error) {
  241. switch v := ptr.(type) {
  242. case *int8:
  243. return readInt8(r, sz, v)
  244. case *int16:
  245. return readInt16(r, sz, v)
  246. case *int32:
  247. return readInt32(r, sz, v)
  248. case *int64:
  249. return readInt64(r, sz, v)
  250. case *string:
  251. return readString(r, sz, v)
  252. case *[]byte:
  253. return readBytes(r, sz, v)
  254. case readable:
  255. return v.readFrom(r, sz)
  256. default:
  257. panic(fmt.Sprintf("unsupported type: %T", v))
  258. }
  259. }
  260. func readStruct(r *bufio.Reader, sz int, v reflect.Value) (int, error) {
  261. var err error
  262. for i, n := 0, v.NumField(); i != n; i++ {
  263. if sz, err = read(r, sz, v.Field(i).Addr().Interface()); err != nil {
  264. return sz, err
  265. }
  266. }
  267. return sz, nil
  268. }
  269. func readSlice(r *bufio.Reader, sz int, v reflect.Value) (int, error) {
  270. var err error
  271. var len int32
  272. if sz, err = readInt32(r, sz, &len); err != nil {
  273. return sz, err
  274. }
  275. if n := int(len); n < 0 {
  276. v.Set(reflect.Zero(v.Type()))
  277. } else {
  278. v.Set(reflect.MakeSlice(v.Type(), n, n))
  279. for i := 0; i != n; i++ {
  280. if sz, err = read(r, sz, v.Index(i).Addr().Interface()); err != nil {
  281. return sz, err
  282. }
  283. }
  284. }
  285. return sz, nil
  286. }
  287. func readFetchResponseHeaderV2(r *bufio.Reader, size int) (throttle int32, watermark int64, remain int, err error) {
  288. var n int32
  289. var p struct {
  290. Partition int32
  291. ErrorCode int16
  292. HighwaterMarkOffset int64
  293. MessageSetSize int32
  294. }
  295. if remain, err = readInt32(r, size, &throttle); err != nil {
  296. return
  297. }
  298. if remain, err = readInt32(r, remain, &n); err != nil {
  299. return
  300. }
  301. // This error should never trigger, unless there's a bug in the kafka client
  302. // or server.
  303. if n != 1 {
  304. err = fmt.Errorf("1 kafka topic was expected in the fetch response but the client received %d", n)
  305. return
  306. }
  307. // We ignore the topic name because we've requests messages for a single
  308. // topic, unless there's a bug in the kafka server we will have received
  309. // the name of the topic that we requested.
  310. if remain, err = discardString(r, remain); err != nil {
  311. return
  312. }
  313. if remain, err = readInt32(r, remain, &n); err != nil {
  314. return
  315. }
  316. // This error should never trigger, unless there's a bug in the kafka client
  317. // or server.
  318. if n != 1 {
  319. err = fmt.Errorf("1 kafka partition was expected in the fetch response but the client received %d", n)
  320. return
  321. }
  322. if remain, err = read(r, remain, &p); err != nil {
  323. return
  324. }
  325. if p.ErrorCode != 0 {
  326. err = Error(p.ErrorCode)
  327. return
  328. }
  329. // This error should never trigger, unless there's a bug in the kafka client
  330. // or server.
  331. if remain != int(p.MessageSetSize) {
  332. err = fmt.Errorf("the size of the message set in a fetch response doesn't match the number of remaining bytes (message set size = %d, remaining bytes = %d)", p.MessageSetSize, remain)
  333. return
  334. }
  335. watermark = p.HighwaterMarkOffset
  336. return
  337. }
  338. func readFetchResponseHeaderV5(r *bufio.Reader, size int) (throttle int32, watermark int64, remain int, err error) {
  339. var n int32
  340. type AbortedTransaction struct {
  341. ProducerId int64
  342. FirstOffset int64
  343. }
  344. var p struct {
  345. Partition int32
  346. ErrorCode int16
  347. HighwaterMarkOffset int64
  348. LastStableOffset int64
  349. LogStartOffset int64
  350. }
  351. var messageSetSize int32
  352. var abortedTransactions []AbortedTransaction
  353. if remain, err = readInt32(r, size, &throttle); err != nil {
  354. return
  355. }
  356. if remain, err = readInt32(r, remain, &n); err != nil {
  357. return
  358. }
  359. // This error should never trigger, unless there's a bug in the kafka client
  360. // or server.
  361. if n != 1 {
  362. err = fmt.Errorf("1 kafka topic was expected in the fetch response but the client received %d", n)
  363. return
  364. }
  365. // We ignore the topic name because we've requests messages for a single
  366. // topic, unless there's a bug in the kafka server we will have received
  367. // the name of the topic that we requested.
  368. if remain, err = discardString(r, remain); err != nil {
  369. return
  370. }
  371. if remain, err = readInt32(r, remain, &n); err != nil {
  372. return
  373. }
  374. // This error should never trigger, unless there's a bug in the kafka client
  375. // or server.
  376. if n != 1 {
  377. err = fmt.Errorf("1 kafka partition was expected in the fetch response but the client received %d", n)
  378. return
  379. }
  380. if remain, err = read(r, remain, &p); err != nil {
  381. return
  382. }
  383. var abortedTransactionLen int
  384. if remain, err = readArrayLen(r, remain, &abortedTransactionLen); err != nil {
  385. return
  386. }
  387. if abortedTransactionLen == -1 {
  388. abortedTransactions = nil
  389. } else {
  390. abortedTransactions = make([]AbortedTransaction, abortedTransactionLen)
  391. for i := 0; i < abortedTransactionLen; i++ {
  392. if remain, err = read(r, remain, &abortedTransactions[i]); err != nil {
  393. return
  394. }
  395. }
  396. }
  397. if p.ErrorCode != 0 {
  398. err = Error(p.ErrorCode)
  399. return
  400. }
  401. remain, err = readInt32(r, remain, &messageSetSize)
  402. if err != nil {
  403. return
  404. }
  405. // This error should never trigger, unless there's a bug in the kafka client
  406. // or server.
  407. if remain != int(messageSetSize) {
  408. err = fmt.Errorf("the size of the message set in a fetch response doesn't match the number of remaining bytes (message set size = %d, remaining bytes = %d)", messageSetSize, remain)
  409. return
  410. }
  411. watermark = p.HighwaterMarkOffset
  412. return
  413. }
  414. func readFetchResponseHeaderV10(r *bufio.Reader, size int) (throttle int32, watermark int64, remain int, err error) {
  415. var n int32
  416. var errorCode int16
  417. type AbortedTransaction struct {
  418. ProducerId int64
  419. FirstOffset int64
  420. }
  421. var p struct {
  422. Partition int32
  423. ErrorCode int16
  424. HighwaterMarkOffset int64
  425. LastStableOffset int64
  426. LogStartOffset int64
  427. }
  428. var messageSetSize int32
  429. var abortedTransactions []AbortedTransaction
  430. if remain, err = readInt32(r, size, &throttle); err != nil {
  431. return
  432. }
  433. if remain, err = readInt16(r, remain, &errorCode); err != nil {
  434. return
  435. }
  436. if errorCode != 0 {
  437. err = Error(errorCode)
  438. return
  439. }
  440. if remain, err = discardInt32(r, remain); err != nil {
  441. return
  442. }
  443. if remain, err = readInt32(r, remain, &n); err != nil {
  444. return
  445. }
  446. // This error should never trigger, unless there's a bug in the kafka client
  447. // or server.
  448. if n != 1 {
  449. err = fmt.Errorf("1 kafka topic was expected in the fetch response but the client received %d", n)
  450. return
  451. }
  452. // We ignore the topic name because we've requests messages for a single
  453. // topic, unless there's a bug in the kafka server we will have received
  454. // the name of the topic that we requested.
  455. if remain, err = discardString(r, remain); err != nil {
  456. return
  457. }
  458. if remain, err = readInt32(r, remain, &n); err != nil {
  459. return
  460. }
  461. // This error should never trigger, unless there's a bug in the kafka client
  462. // or server.
  463. if n != 1 {
  464. err = fmt.Errorf("1 kafka partition was expected in the fetch response but the client received %d", n)
  465. return
  466. }
  467. if remain, err = read(r, remain, &p); err != nil {
  468. return
  469. }
  470. var abortedTransactionLen int
  471. if remain, err = readArrayLen(r, remain, &abortedTransactionLen); err != nil {
  472. return
  473. }
  474. if abortedTransactionLen == -1 {
  475. abortedTransactions = nil
  476. } else {
  477. abortedTransactions = make([]AbortedTransaction, abortedTransactionLen)
  478. for i := 0; i < abortedTransactionLen; i++ {
  479. if remain, err = read(r, remain, &abortedTransactions[i]); err != nil {
  480. return
  481. }
  482. }
  483. }
  484. if p.ErrorCode != 0 {
  485. err = Error(p.ErrorCode)
  486. return
  487. }
  488. remain, err = readInt32(r, remain, &messageSetSize)
  489. if err != nil {
  490. return
  491. }
  492. // This error should never trigger, unless there's a bug in the kafka client
  493. // or server.
  494. if remain != int(messageSetSize) {
  495. err = fmt.Errorf("the size of the message set in a fetch response doesn't match the number of remaining bytes (message set size = %d, remaining bytes = %d)", messageSetSize, remain)
  496. return
  497. }
  498. watermark = p.HighwaterMarkOffset
  499. return
  500. }
  501. func readMessageHeader(r *bufio.Reader, sz int) (offset int64, attributes int8, timestamp int64, remain int, err error) {
  502. var version int8
  503. if remain, err = readInt64(r, sz, &offset); err != nil {
  504. return
  505. }
  506. // On discarding the message size and CRC:
  507. // ---------------------------------------
  508. //
  509. // - Not sure why kafka gives the message size here, we already have the
  510. // number of remaining bytes in the response and kafka should only truncate
  511. // the trailing message.
  512. //
  513. // - TCP is already taking care of ensuring data integrity, no need to
  514. // waste resources doing it a second time so we just skip the message CRC.
  515. //
  516. if remain, err = discardN(r, remain, 8); err != nil {
  517. return
  518. }
  519. if remain, err = readInt8(r, remain, &version); err != nil {
  520. return
  521. }
  522. if remain, err = readInt8(r, remain, &attributes); err != nil {
  523. return
  524. }
  525. switch version {
  526. case 0:
  527. case 1:
  528. remain, err = readInt64(r, remain, &timestamp)
  529. default:
  530. err = fmt.Errorf("unsupported message version %d found in fetch response", version)
  531. }
  532. return
  533. }