common_test.go 23 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102
  1. package processors
  2. import (
  3. "encoding/json"
  4. "fmt"
  5. "github.com/emqx/kuiper/common"
  6. "github.com/emqx/kuiper/xsql"
  7. "github.com/emqx/kuiper/xstream"
  8. "github.com/emqx/kuiper/xstream/api"
  9. "github.com/emqx/kuiper/xstream/nodes"
  10. "github.com/emqx/kuiper/xstream/test"
  11. "os"
  12. "path"
  13. "reflect"
  14. "strings"
  15. "testing"
  16. "time"
  17. )
  18. const SOURCELEAP = 200 // Time change before sending a data
  19. const POSTLEAP = 1000 // Time change after all data sends out
  20. type ruleTest struct {
  21. name string
  22. sql string
  23. r interface{} // The result
  24. m map[string]interface{} // final metrics
  25. }
  26. var DbDir = getDbDir()
  27. func getDbDir() string {
  28. common.InitConf()
  29. dbDir, err := common.GetDataLoc()
  30. if err != nil {
  31. log.Panic(err)
  32. }
  33. log.Infof("db location is %s", dbDir)
  34. return dbDir
  35. }
  36. func cleanStateData() {
  37. dbDir, err := common.GetDataLoc()
  38. if err != nil {
  39. log.Panic(err)
  40. }
  41. c := path.Join(dbDir, "checkpoints")
  42. err = os.RemoveAll(c)
  43. if err != nil {
  44. log.Errorf("%s", err)
  45. }
  46. s := path.Join(dbDir, "sink", "cache")
  47. err = os.RemoveAll(s)
  48. if err != nil {
  49. log.Errorf("%s", err)
  50. }
  51. }
  52. func compareMetrics(tp *xstream.TopologyNew, m map[string]interface{}) (err error) {
  53. keys, values := tp.GetMetrics()
  54. if common.Config.Basic.Debug == true {
  55. for i, k := range keys {
  56. log.Printf("%s:%v", k, values[i])
  57. }
  58. }
  59. for k, v := range m {
  60. var (
  61. index int
  62. key string
  63. matched bool
  64. )
  65. for index, key = range keys {
  66. if k == key {
  67. if strings.HasSuffix(k, "process_latency_ms") {
  68. if values[index].(int64) >= v.(int64) {
  69. matched = true
  70. continue
  71. } else {
  72. break
  73. }
  74. }
  75. if values[index] == v {
  76. matched = true
  77. }
  78. break
  79. }
  80. }
  81. if matched {
  82. continue
  83. }
  84. //do not find
  85. if index < len(values) {
  86. return fmt.Errorf("metrics mismatch for %s:\n\nexp=%#v(%t)\n\ngot=%#v(%t)\n\n", k, v, v, values[index], values[index])
  87. } else {
  88. return fmt.Errorf("metrics mismatch for %s:\n\nexp=%#v\n\ngot=nil\n\n", k, v)
  89. }
  90. }
  91. return nil
  92. }
  93. func errstring(err error) string {
  94. if err != nil {
  95. return err.Error()
  96. }
  97. return ""
  98. }
  99. // The time diff must larger than timeleap
  100. var testData = map[string][]*xsql.Tuple{
  101. "demo": {
  102. {
  103. Emitter: "demo",
  104. Message: map[string]interface{}{
  105. "color": "red",
  106. "size": 3,
  107. "ts": 1541152486013,
  108. },
  109. Timestamp: 1541152486013,
  110. },
  111. {
  112. Emitter: "demo",
  113. Message: map[string]interface{}{
  114. "color": "blue",
  115. "size": 6,
  116. "ts": 1541152486822,
  117. },
  118. Timestamp: 1541152486822,
  119. },
  120. {
  121. Emitter: "demo",
  122. Message: map[string]interface{}{
  123. "color": "blue",
  124. "size": 2,
  125. "ts": 1541152487632,
  126. },
  127. Timestamp: 1541152487632,
  128. },
  129. {
  130. Emitter: "demo",
  131. Message: map[string]interface{}{
  132. "color": "yellow",
  133. "size": 4,
  134. "ts": 1541152488442,
  135. },
  136. Timestamp: 1541152488442,
  137. },
  138. {
  139. Emitter: "demo",
  140. Message: map[string]interface{}{
  141. "color": "red",
  142. "size": 1,
  143. "ts": 1541152489252,
  144. },
  145. Timestamp: 1541152489252,
  146. },
  147. },
  148. "demoError": {
  149. {
  150. Emitter: "demoError",
  151. Message: map[string]interface{}{
  152. "color": 3,
  153. "size": "red",
  154. "ts": 1541152486013,
  155. },
  156. Timestamp: 1541152486013,
  157. },
  158. {
  159. Emitter: "demoError",
  160. Message: map[string]interface{}{
  161. "color": "blue",
  162. "size": 6,
  163. "ts": "1541152486822",
  164. },
  165. Timestamp: 1541152486822,
  166. },
  167. {
  168. Emitter: "demoError",
  169. Message: map[string]interface{}{
  170. "color": "blue",
  171. "size": 2,
  172. "ts": 1541152487632,
  173. },
  174. Timestamp: 1541152487632,
  175. },
  176. {
  177. Emitter: "demoError",
  178. Message: map[string]interface{}{
  179. "color": 7,
  180. "size": 4,
  181. "ts": 1541152488442,
  182. },
  183. Timestamp: 1541152488442,
  184. },
  185. {
  186. Emitter: "demoError",
  187. Message: map[string]interface{}{
  188. "color": "red",
  189. "size": "blue",
  190. "ts": 1541152489252,
  191. },
  192. Timestamp: 1541152489252,
  193. },
  194. },
  195. "demo1": {
  196. {
  197. Emitter: "demo1",
  198. Message: map[string]interface{}{
  199. "temp": 25.5,
  200. "hum": 65,
  201. "from": "device1",
  202. "ts": 1541152486013,
  203. },
  204. Timestamp: 1541152486013,
  205. },
  206. {
  207. Emitter: "demo1",
  208. Message: map[string]interface{}{
  209. "temp": 27.5,
  210. "hum": 59,
  211. "from": "device2",
  212. "ts": 1541152486823,
  213. },
  214. Timestamp: 1541152486823,
  215. },
  216. {
  217. Emitter: "demo1",
  218. Message: map[string]interface{}{
  219. "temp": 28.1,
  220. "hum": 75,
  221. "from": "device3",
  222. "ts": 1541152487632,
  223. },
  224. Timestamp: 1541152487632,
  225. },
  226. {
  227. Emitter: "demo1",
  228. Message: map[string]interface{}{
  229. "temp": 27.4,
  230. "hum": 80,
  231. "from": "device1",
  232. "ts": 1541152488442,
  233. },
  234. Timestamp: 1541152488442,
  235. },
  236. {
  237. Emitter: "demo1",
  238. Message: map[string]interface{}{
  239. "temp": 25.5,
  240. "hum": 62,
  241. "from": "device3",
  242. "ts": 1541152489252,
  243. },
  244. Timestamp: 1541152489252,
  245. },
  246. },
  247. "sessionDemo": {
  248. {
  249. Emitter: "sessionDemo",
  250. Message: map[string]interface{}{
  251. "temp": 25.5,
  252. "hum": 65,
  253. "ts": 1541152486013,
  254. },
  255. Timestamp: 1541152486013,
  256. },
  257. {
  258. Emitter: "sessionDemo",
  259. Message: map[string]interface{}{
  260. "temp": 27.5,
  261. "hum": 59,
  262. "ts": 1541152486823,
  263. },
  264. Timestamp: 1541152486823,
  265. },
  266. {
  267. Emitter: "sessionDemo",
  268. Message: map[string]interface{}{
  269. "temp": 28.1,
  270. "hum": 75,
  271. "ts": 1541152487932,
  272. },
  273. Timestamp: 1541152487932,
  274. },
  275. {
  276. Emitter: "sessionDemo",
  277. Message: map[string]interface{}{
  278. "temp": 27.4,
  279. "hum": 80,
  280. "ts": 1541152488442,
  281. },
  282. Timestamp: 1541152488442,
  283. },
  284. {
  285. Emitter: "sessionDemo",
  286. Message: map[string]interface{}{
  287. "temp": 25.5,
  288. "hum": 62,
  289. "ts": 1541152489252,
  290. },
  291. Timestamp: 1541152489252,
  292. },
  293. {
  294. Emitter: "sessionDemo",
  295. Message: map[string]interface{}{
  296. "temp": 26.2,
  297. "hum": 63,
  298. "ts": 1541152490062,
  299. },
  300. Timestamp: 1541152490062,
  301. },
  302. {
  303. Emitter: "sessionDemo",
  304. Message: map[string]interface{}{
  305. "temp": 26.8,
  306. "hum": 71,
  307. "ts": 1541152490872,
  308. },
  309. Timestamp: 1541152490872,
  310. },
  311. {
  312. Emitter: "sessionDemo",
  313. Message: map[string]interface{}{
  314. "temp": 28.9,
  315. "hum": 85,
  316. "ts": 1541152491682,
  317. },
  318. Timestamp: 1541152491682,
  319. },
  320. {
  321. Emitter: "sessionDemo",
  322. Message: map[string]interface{}{
  323. "temp": 29.1,
  324. "hum": 92,
  325. "ts": 1541152492492,
  326. },
  327. Timestamp: 1541152492492,
  328. },
  329. {
  330. Emitter: "sessionDemo",
  331. Message: map[string]interface{}{
  332. "temp": 32.2,
  333. "hum": 99,
  334. "ts": 1541152493202,
  335. },
  336. Timestamp: 1541152493202,
  337. },
  338. {
  339. Emitter: "sessionDemo",
  340. Message: map[string]interface{}{
  341. "temp": 30.9,
  342. "hum": 87,
  343. "ts": 1541152494112,
  344. },
  345. Timestamp: 1541152494112,
  346. },
  347. },
  348. "demoE": {
  349. {
  350. Emitter: "demoE",
  351. Message: map[string]interface{}{
  352. "color": "red",
  353. "size": 3,
  354. "ts": 1541152486013,
  355. },
  356. Timestamp: 1541152486023,
  357. },
  358. {
  359. Emitter: "demoE",
  360. Message: map[string]interface{}{
  361. "color": "blue",
  362. "size": 2,
  363. "ts": 1541152487632,
  364. },
  365. Timestamp: 1541152487822,
  366. },
  367. {
  368. Emitter: "demoE",
  369. Message: map[string]interface{}{
  370. "color": "red",
  371. "size": 1,
  372. "ts": 1541152489252,
  373. },
  374. Timestamp: 1541152489632,
  375. },
  376. { //dropped item
  377. Emitter: "demoE",
  378. Message: map[string]interface{}{
  379. "color": "blue",
  380. "size": 6,
  381. "ts": 1541152486822,
  382. },
  383. Timestamp: 1541152489842,
  384. },
  385. {
  386. Emitter: "demoE",
  387. Message: map[string]interface{}{
  388. "color": "yellow",
  389. "size": 4,
  390. "ts": 1541152488442,
  391. },
  392. Timestamp: 1541152490052,
  393. },
  394. { //To lift the watermark and issue all windows
  395. Emitter: "demoE",
  396. Message: map[string]interface{}{
  397. "color": "yellow",
  398. "size": 4,
  399. "ts": 1541152492342,
  400. },
  401. Timestamp: 1541152498888,
  402. },
  403. },
  404. "demo1E": {
  405. {
  406. Emitter: "demo1E",
  407. Message: map[string]interface{}{
  408. "temp": 27.5,
  409. "hum": 59,
  410. "ts": 1541152486823,
  411. },
  412. Timestamp: 1541152487250,
  413. },
  414. {
  415. Emitter: "demo1E",
  416. Message: map[string]interface{}{
  417. "temp": 25.5,
  418. "hum": 65,
  419. "ts": 1541152486013,
  420. },
  421. Timestamp: 1541152487751,
  422. },
  423. {
  424. Emitter: "demo1E",
  425. Message: map[string]interface{}{
  426. "temp": 27.4,
  427. "hum": 80,
  428. "ts": 1541152488442,
  429. },
  430. Timestamp: 1541152489252,
  431. },
  432. {
  433. Emitter: "demo1E",
  434. Message: map[string]interface{}{
  435. "temp": 28.1,
  436. "hum": 75,
  437. "ts": 1541152487632,
  438. },
  439. Timestamp: 1541152489753,
  440. },
  441. {
  442. Emitter: "demo1E",
  443. Message: map[string]interface{}{
  444. "temp": 25.5,
  445. "hum": 62,
  446. "ts": 1541152489252,
  447. },
  448. Timestamp: 1541152489954,
  449. },
  450. {
  451. Emitter: "demo1E",
  452. Message: map[string]interface{}{
  453. "temp": 25.5,
  454. "hum": 62,
  455. "ts": 1541152499252,
  456. },
  457. Timestamp: 1541152499755,
  458. },
  459. },
  460. "sessionDemoE": {
  461. {
  462. Emitter: "sessionDemoE",
  463. Message: map[string]interface{}{
  464. "temp": 25.5,
  465. "hum": 65,
  466. "ts": 1541152486013,
  467. },
  468. Timestamp: 1541152486250,
  469. },
  470. {
  471. Emitter: "sessionDemoE",
  472. Message: map[string]interface{}{
  473. "temp": 28.1,
  474. "hum": 75,
  475. "ts": 1541152487932,
  476. },
  477. Timestamp: 1541152487951,
  478. },
  479. {
  480. Emitter: "sessionDemoE",
  481. Message: map[string]interface{}{
  482. "temp": 27.5,
  483. "hum": 59,
  484. "ts": 1541152486823,
  485. },
  486. Timestamp: 1541152488552,
  487. },
  488. {
  489. Emitter: "sessionDemoE",
  490. Message: map[string]interface{}{
  491. "temp": 25.5,
  492. "hum": 62,
  493. "ts": 1541152489252,
  494. },
  495. Timestamp: 1541152489353,
  496. },
  497. {
  498. Emitter: "sessionDemoE",
  499. Message: map[string]interface{}{
  500. "temp": 27.4,
  501. "hum": 80,
  502. "ts": 1541152488442,
  503. },
  504. Timestamp: 1541152489854,
  505. },
  506. {
  507. Emitter: "sessionDemoE",
  508. Message: map[string]interface{}{
  509. "temp": 26.2,
  510. "hum": 63,
  511. "ts": 1541152490062,
  512. },
  513. Timestamp: 1541152490155,
  514. },
  515. {
  516. Emitter: "sessionDemoE",
  517. Message: map[string]interface{}{
  518. "temp": 28.9,
  519. "hum": 85,
  520. "ts": 1541152491682,
  521. },
  522. Timestamp: 1541152491686,
  523. },
  524. {
  525. Emitter: "sessionDemoE",
  526. Message: map[string]interface{}{
  527. "temp": 26.8,
  528. "hum": 71,
  529. "ts": 1541152490872,
  530. },
  531. Timestamp: 1541152491972,
  532. },
  533. {
  534. Emitter: "sessionDemoE",
  535. Message: map[string]interface{}{
  536. "temp": 29.1,
  537. "hum": 92,
  538. "ts": 1541152492492,
  539. },
  540. Timestamp: 1541152492592,
  541. },
  542. {
  543. Emitter: "sessionDemoE",
  544. Message: map[string]interface{}{
  545. "temp": 30.9,
  546. "hum": 87,
  547. "ts": 1541152494112,
  548. },
  549. Timestamp: 1541152494212,
  550. },
  551. {
  552. Emitter: "sessionDemoE",
  553. Message: map[string]interface{}{
  554. "temp": 32.2,
  555. "hum": 99,
  556. "ts": 1541152493202,
  557. },
  558. Timestamp: 1541152495202,
  559. },
  560. {
  561. Emitter: "sessionDemoE",
  562. Message: map[string]interface{}{
  563. "temp": 32.2,
  564. "hum": 99,
  565. "ts": 1541152499202,
  566. },
  567. Timestamp: 1541152499402,
  568. },
  569. },
  570. "demoErr": {
  571. {
  572. Emitter: "demoErr",
  573. Message: map[string]interface{}{
  574. "color": "red",
  575. "size": 3,
  576. "ts": 1541152486013,
  577. },
  578. Timestamp: 1541152486221,
  579. },
  580. {
  581. Emitter: "demoErr",
  582. Message: map[string]interface{}{
  583. "color": 2,
  584. "size": "blue",
  585. "ts": 1541152487632,
  586. },
  587. Timestamp: 1541152487722,
  588. },
  589. {
  590. Emitter: "demoErr",
  591. Message: map[string]interface{}{
  592. "color": "red",
  593. "size": 1,
  594. "ts": 1541152489252,
  595. },
  596. Timestamp: 1541152489332,
  597. },
  598. { //dropped item
  599. Emitter: "demoErr",
  600. Message: map[string]interface{}{
  601. "color": "blue",
  602. "size": 6,
  603. "ts": 1541152486822,
  604. },
  605. Timestamp: 1541152489822,
  606. },
  607. {
  608. Emitter: "demoErr",
  609. Message: map[string]interface{}{
  610. "color": "yellow",
  611. "size": 4,
  612. "ts": 1541152488442,
  613. },
  614. Timestamp: 1541152490042,
  615. },
  616. { //To lift the watermark and issue all windows
  617. Emitter: "demoErr",
  618. Message: map[string]interface{}{
  619. "color": "yellow",
  620. "size": 4,
  621. "ts": 1541152492342,
  622. },
  623. Timestamp: 1541152493842,
  624. },
  625. },
  626. "ldemo": {
  627. {
  628. Emitter: "ldemo",
  629. Message: map[string]interface{}{
  630. "color": "red",
  631. "size": 3,
  632. "ts": 1541152486013,
  633. },
  634. Timestamp: 1541152486013,
  635. },
  636. {
  637. Emitter: "ldemo",
  638. Message: map[string]interface{}{
  639. "color": "blue",
  640. "size": "string",
  641. "ts": 1541152486822,
  642. },
  643. Timestamp: 1541152486822,
  644. },
  645. {
  646. Emitter: "ldemo",
  647. Message: map[string]interface{}{
  648. "size": 3,
  649. "ts": 1541152487632,
  650. },
  651. Timestamp: 1541152487632,
  652. },
  653. {
  654. Emitter: "ldemo",
  655. Message: map[string]interface{}{
  656. "color": 49,
  657. "size": 2,
  658. "ts": 1541152488442,
  659. },
  660. Timestamp: 1541152488442,
  661. },
  662. {
  663. Emitter: "ldemo",
  664. Message: map[string]interface{}{
  665. "color": "red",
  666. "ts": 1541152489252,
  667. },
  668. Timestamp: 1541152489252,
  669. },
  670. },
  671. "ldemo1": {
  672. {
  673. Emitter: "ldemo1",
  674. Message: map[string]interface{}{
  675. "temp": 25.5,
  676. "hum": 65,
  677. "ts": 1541152486013,
  678. },
  679. Timestamp: 1541152486013,
  680. },
  681. {
  682. Emitter: "ldemo1",
  683. Message: map[string]interface{}{
  684. "temp": 27.5,
  685. "hum": 59,
  686. "ts": 1541152486823,
  687. },
  688. Timestamp: 1541152486823,
  689. },
  690. {
  691. Emitter: "ldemo1",
  692. Message: map[string]interface{}{
  693. "temp": 28.1,
  694. "hum": 75,
  695. "ts": 1541152487632,
  696. },
  697. Timestamp: 1541152487632,
  698. },
  699. {
  700. Emitter: "ldemo1",
  701. Message: map[string]interface{}{
  702. "temp": 27.4,
  703. "hum": 80,
  704. "ts": "1541152488442",
  705. },
  706. Timestamp: 1541152488442,
  707. },
  708. {
  709. Emitter: "ldemo1",
  710. Message: map[string]interface{}{
  711. "temp": 25.5,
  712. "hum": 62,
  713. "ts": 1541152489252,
  714. },
  715. Timestamp: 1541152489252,
  716. },
  717. },
  718. "lsessionDemo": {
  719. {
  720. Emitter: "lsessionDemo",
  721. Message: map[string]interface{}{
  722. "temp": 25.5,
  723. "hum": 65,
  724. "ts": 1541152486013,
  725. },
  726. Timestamp: 1541152486013,
  727. },
  728. {
  729. Emitter: "lsessionDemo",
  730. Message: map[string]interface{}{
  731. "temp": 27.5,
  732. "hum": 59,
  733. "ts": 1541152486823,
  734. },
  735. Timestamp: 1541152486823,
  736. },
  737. {
  738. Emitter: "lsessionDemo",
  739. Message: map[string]interface{}{
  740. "temp": 28.1,
  741. "hum": 75,
  742. "ts": 1541152487932,
  743. },
  744. Timestamp: 1541152487932,
  745. },
  746. {
  747. Emitter: "lsessionDemo",
  748. Message: map[string]interface{}{
  749. "temp": 27.4,
  750. "hum": 80,
  751. "ts": 1541152488442,
  752. },
  753. Timestamp: 1541152488442,
  754. },
  755. {
  756. Emitter: "lsessionDemo",
  757. Message: map[string]interface{}{
  758. "temp": 25.5,
  759. "hum": 62,
  760. "ts": 1541152489252,
  761. },
  762. Timestamp: 1541152489252,
  763. },
  764. {
  765. Emitter: "lsessionDemo",
  766. Message: map[string]interface{}{
  767. "temp": 26.2,
  768. "hum": 63,
  769. "ts": 1541152490062,
  770. },
  771. Timestamp: 1541152490062,
  772. },
  773. {
  774. Emitter: "lsessionDemo",
  775. Message: map[string]interface{}{
  776. "temp": 26.8,
  777. "hum": 71,
  778. "ts": 1541152490872,
  779. },
  780. Timestamp: 1541152490872,
  781. },
  782. {
  783. Emitter: "lsessionDemo",
  784. Message: map[string]interface{}{
  785. "temp": 28.9,
  786. "hum": 85,
  787. "ts": 1541152491682,
  788. },
  789. Timestamp: 1541152491682,
  790. },
  791. {
  792. Emitter: "lsessionDemo",
  793. Message: map[string]interface{}{
  794. "temp": 29.1,
  795. "hum": 92,
  796. "ts": 1541152492492,
  797. },
  798. Timestamp: 1541152492492,
  799. },
  800. {
  801. Emitter: "lsessionDemo",
  802. Message: map[string]interface{}{
  803. "temp": 2.2,
  804. "hum": 99,
  805. "ts": 1541152493202,
  806. },
  807. Timestamp: 1541152493202,
  808. },
  809. {
  810. Emitter: "lsessionDemo",
  811. Message: map[string]interface{}{
  812. "temp": 30.9,
  813. "hum": 87,
  814. "ts": 1541152494112,
  815. },
  816. Timestamp: 1541152494112,
  817. },
  818. },
  819. "text": {
  820. {
  821. Emitter: "text",
  822. Message: map[string]interface{}{
  823. "slogan": "Impossible is nothing",
  824. "brand": "Adidas",
  825. },
  826. Timestamp: 1541152486500,
  827. },
  828. {
  829. Emitter: "text",
  830. Message: map[string]interface{}{
  831. "slogan": "Stronger than dirt",
  832. "brand": "Ajax",
  833. },
  834. Timestamp: 1541152487400,
  835. },
  836. {
  837. Emitter: "text",
  838. Message: map[string]interface{}{
  839. "slogan": "Belong anywhere",
  840. "brand": "Airbnb",
  841. },
  842. Timestamp: 1541152488300,
  843. },
  844. {
  845. Emitter: "text",
  846. Message: map[string]interface{}{
  847. "slogan": "I can't believe I ate the whole thing",
  848. "brand": "Alka Seltzer",
  849. },
  850. Timestamp: 1541152489200,
  851. },
  852. {
  853. Emitter: "text",
  854. Message: map[string]interface{}{
  855. "slogan": "You're in good hands",
  856. "brand": "Allstate",
  857. },
  858. Timestamp: 1541152490100,
  859. },
  860. {
  861. Emitter: "text",
  862. Message: map[string]interface{}{
  863. "slogan": "Don't leave home without it",
  864. "brand": "American Express",
  865. },
  866. Timestamp: 1541152491200,
  867. },
  868. {
  869. Emitter: "text",
  870. Message: map[string]interface{}{
  871. "slogan": "Think different",
  872. "brand": "Apple",
  873. },
  874. Timestamp: 1541152492300,
  875. },
  876. {
  877. Emitter: "text",
  878. Message: map[string]interface{}{
  879. "slogan": "We try harder",
  880. "brand": "Avis",
  881. },
  882. Timestamp: 1541152493400,
  883. },
  884. },
  885. }
  886. func commonResultFunc(result [][]byte) interface{} {
  887. var maps [][]map[string]interface{}
  888. for _, v := range result {
  889. var mapRes []map[string]interface{}
  890. err := json.Unmarshal(v, &mapRes)
  891. if err != nil {
  892. panic("Failed to parse the input into map")
  893. }
  894. maps = append(maps, mapRes)
  895. }
  896. return maps
  897. }
  898. func doRuleTest(t *testing.T, tests []ruleTest, j int, opt *api.RuleOption) {
  899. doRuleTestBySinkProps(t, tests, j, opt, nil, commonResultFunc)
  900. }
  901. func doRuleTestBySinkProps(t *testing.T, tests []ruleTest, j int, opt *api.RuleOption, sinkProps map[string]interface{}, resultFunc func(result [][]byte) interface{}) {
  902. fmt.Printf("The test bucket for option %d size is %d.\n\n", j, len(tests))
  903. for i, tt := range tests {
  904. datas, dataLength, tp, mockSink, errCh := createStream(t, tt, j, opt, sinkProps)
  905. if err := sendData(t, dataLength, tt.m, datas, errCh, tp, POSTLEAP); err != nil {
  906. t.Errorf("send data error %s", err)
  907. break
  908. }
  909. compareResult(t, mockSink, resultFunc, tt, i, tp)
  910. }
  911. }
  912. func compareResult(t *testing.T, mockSink *test.MockSink, resultFunc func(result [][]byte) interface{}, tt ruleTest, i int, tp *xstream.TopologyNew) {
  913. // Check results
  914. results := mockSink.GetResults()
  915. maps := resultFunc(results)
  916. if !reflect.DeepEqual(tt.r, maps) {
  917. t.Errorf("%d. %q\n\nresult mismatch:\n\nexp=%#v\n\ngot=%#v\n\n", i, tt.sql, tt.r, maps)
  918. }
  919. if err := compareMetrics(tp, tt.m); err != nil {
  920. t.Errorf("%d. %q\n\nmetrics mismatch:\n\n%s\n\n", i, tt.sql, err)
  921. }
  922. tp.Cancel()
  923. }
  924. func sendData(t *testing.T, dataLength int, metrics map[string]interface{}, datas [][]*xsql.Tuple, errCh <-chan error, tp *xstream.TopologyNew, postleap int) error {
  925. // Send data and move time
  926. mockClock := test.GetMockClock()
  927. // TODO assume multiple data source send the data in order and has the same length
  928. for i := 0; i < dataLength; i++ {
  929. mockClock.Add(SOURCELEAP * time.Millisecond)
  930. common.Log.Debugf("Clock add to %d", common.GetNowInMilli())
  931. time.Sleep(1)
  932. for _, d := range datas {
  933. // Make sure time is going forward only
  934. if d[i].Timestamp > common.GetNowInMilli() {
  935. mockClock.Set(common.TimeFromUnixMilli(d[i].Timestamp))
  936. common.Log.Debugf("Clock set to %d", common.GetNowInMilli())
  937. time.Sleep(1)
  938. }
  939. select {
  940. case err := <-errCh:
  941. t.Log(err)
  942. tp.Cancel()
  943. return err
  944. default:
  945. }
  946. }
  947. }
  948. mockClock.Add(time.Duration(postleap) * time.Millisecond)
  949. common.Log.Debugf("Clock add to %d", common.GetNowInMilli())
  950. time.Sleep(1)
  951. // Check if stream done. Poll for metrics,
  952. for retry := 100; retry > 0; retry-- {
  953. time.Sleep(time.Duration(retry) * time.Millisecond)
  954. if err := compareMetrics(tp, metrics); err == nil {
  955. break
  956. } else {
  957. common.Log.Debugf("check metrics error at %d: %s", retry, err)
  958. }
  959. }
  960. return nil
  961. }
  962. func createStream(t *testing.T, tt ruleTest, j int, opt *api.RuleOption, sinkProps map[string]interface{}) ([][]*xsql.Tuple, int, *xstream.TopologyNew, *test.MockSink, <-chan error) {
  963. // Rest for each test
  964. cleanStateData()
  965. test.ResetClock(1541152485800)
  966. // Create stream
  967. var (
  968. sources []*nodes.SourceNode
  969. datas [][]*xsql.Tuple
  970. dataLength int
  971. )
  972. p := NewRuleProcessor(DbDir)
  973. parser := xsql.NewParser(strings.NewReader(tt.sql))
  974. if stmt, err := xsql.Language.Parse(parser); err != nil {
  975. t.Errorf("parse sql %s error: %s", tt.sql, err)
  976. } else {
  977. if selectStmt, ok := stmt.(*xsql.SelectStatement); !ok {
  978. t.Errorf("sql %s is not a select statement", tt.sql)
  979. } else {
  980. streams := xsql.GetStreams(selectStmt)
  981. for _, stream := range streams {
  982. data := testData[stream]
  983. dataLength = len(data)
  984. datas = append(datas, data)
  985. source := nodes.NewSourceNodeWithSource(stream, test.NewMockSource(data), map[string]string{
  986. "DATASOURCE": stream,
  987. })
  988. sources = append(sources, source)
  989. }
  990. }
  991. }
  992. tp, inputs, err := p.createTopoWithSources(&api.Rule{Id: fmt.Sprintf("%s_%d", tt.name, j), Sql: tt.sql, Options: opt}, sources)
  993. if err != nil {
  994. t.Error(err)
  995. }
  996. mockSink := test.NewMockSink()
  997. sink := nodes.NewSinkNodeWithSink("mockSink", mockSink, sinkProps)
  998. tp.AddSink(inputs, sink)
  999. errCh := tp.Open()
  1000. return datas, dataLength, tp, mockSink, errCh
  1001. }
  1002. // Create or drop streams
  1003. func handleStream(createOrDrop bool, names []string, t *testing.T) {
  1004. p := NewStreamProcessor(path.Join(DbDir, "stream"))
  1005. for _, name := range names {
  1006. var sql string
  1007. if createOrDrop {
  1008. switch name {
  1009. case "demo":
  1010. sql = `CREATE STREAM demo (
  1011. color STRING,
  1012. size BIGINT,
  1013. ts BIGINT
  1014. ) WITH (DATASOURCE="demo", FORMAT="json", KEY="ts");`
  1015. case "demoError":
  1016. sql = `CREATE STREAM demoError (
  1017. color STRING,
  1018. size BIGINT,
  1019. ts BIGINT
  1020. ) WITH (DATASOURCE="demoError", FORMAT="json", KEY="ts");`
  1021. case "demo1":
  1022. sql = `CREATE STREAM demo1 (
  1023. temp FLOAT,
  1024. hum BIGINT,` +
  1025. "`from`" + ` STRING,
  1026. ts BIGINT
  1027. ) WITH (DATASOURCE="demo1", FORMAT="json", KEY="ts");`
  1028. case "sessionDemo":
  1029. sql = `CREATE STREAM sessionDemo (
  1030. temp FLOAT,
  1031. hum BIGINT,
  1032. ts BIGINT
  1033. ) WITH (DATASOURCE="sessionDemo", FORMAT="json", KEY="ts");`
  1034. case "demoE":
  1035. sql = `CREATE STREAM demoE (
  1036. color STRING,
  1037. size BIGINT,
  1038. ts BIGINT
  1039. ) WITH (DATASOURCE="demoE", FORMAT="json", KEY="ts", TIMESTAMP="ts");`
  1040. case "demo1E":
  1041. sql = `CREATE STREAM demo1E (
  1042. temp FLOAT,
  1043. hum BIGINT,
  1044. ts BIGINT
  1045. ) WITH (DATASOURCE="demo1E", FORMAT="json", KEY="ts", TIMESTAMP="ts");`
  1046. case "sessionDemoE":
  1047. sql = `CREATE STREAM sessionDemoE (
  1048. temp FLOAT,
  1049. hum BIGINT,
  1050. ts BIGINT
  1051. ) WITH (DATASOURCE="sessionDemoE", FORMAT="json", KEY="ts", TIMESTAMP="ts");`
  1052. case "demoErr":
  1053. sql = `CREATE STREAM demoErr (
  1054. color STRING,
  1055. size BIGINT,
  1056. ts BIGINT
  1057. ) WITH (DATASOURCE="demoErr", FORMAT="json", KEY="ts", TIMESTAMP="ts");`
  1058. case "ldemo":
  1059. sql = `CREATE STREAM ldemo (
  1060. ) WITH (DATASOURCE="ldemo", FORMAT="json");`
  1061. case "ldemo1":
  1062. sql = `CREATE STREAM ldemo1 (
  1063. ) WITH (DATASOURCE="ldemo1", FORMAT="json");`
  1064. case "lsessionDemo":
  1065. sql = `CREATE STREAM lsessionDemo (
  1066. ) WITH (DATASOURCE="lsessionDemo", FORMAT="json");`
  1067. case "ext":
  1068. sql = "CREATE STREAM ext (count bigint) WITH (DATASOURCE=\"users\", FORMAT=\"JSON\", TYPE=\"random\", CONF_KEY=\"ext\")"
  1069. case "ext2":
  1070. sql = "CREATE STREAM ext2 (count bigint) WITH (DATASOURCE=\"users\", FORMAT=\"JSON\", TYPE=\"random\", CONF_KEY=\"dedup\")"
  1071. case "text":
  1072. sql = "CREATE STREAM text (slogan string, brand string) WITH (DATASOURCE=\"users\", FORMAT=\"JSON\")"
  1073. default:
  1074. t.Errorf("create stream %s fail", name)
  1075. }
  1076. } else {
  1077. sql = `DROP STREAM ` + name
  1078. }
  1079. _, err := p.ExecStmt(sql)
  1080. if err != nil {
  1081. t.Log(err)
  1082. }
  1083. }
  1084. }