xsql_processor_test.go 27 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295
  1. package processors
  2. import (
  3. "encoding/json"
  4. "engine/common"
  5. "engine/xsql"
  6. "engine/xstream"
  7. "engine/xstream/test"
  8. "fmt"
  9. "path"
  10. "reflect"
  11. "strings"
  12. "testing"
  13. "time"
  14. )
  15. var BadgerDir string
  16. func init(){
  17. dataDir, err := common.GetDataLoc()
  18. if err != nil {
  19. log.Panic(err)
  20. }else{
  21. log.Infof("db location is %s", dataDir)
  22. }
  23. BadgerDir = path.Join(path.Dir(dataDir), "test")
  24. log.Infof("badge location is %s", BadgerDir)
  25. }
  26. func TestStreamCreateProcessor(t *testing.T) {
  27. var tests = []struct {
  28. s string
  29. r []string
  30. err string
  31. }{
  32. {
  33. s: `SHOW STREAMS;`,
  34. r: []string{"no stream definition found"},
  35. },
  36. {
  37. s: `EXPLAIN STREAM topic1;`,
  38. err: "stream topic1 not found",
  39. },
  40. {
  41. s: `CREATE STREAM topic1 (
  42. USERID BIGINT,
  43. FIRST_NAME STRING,
  44. LAST_NAME STRING,
  45. NICKNAMES ARRAY(STRING),
  46. Gender BOOLEAN,
  47. ADDRESS STRUCT(STREET_NAME STRING, NUMBER BIGINT),
  48. ) WITH (DATASOURCE="users", FORMAT="AVRO", KEY="USERID");`,
  49. r: []string{"stream topic1 created"},
  50. },
  51. {
  52. s: `CREATE STREAM topic1 (
  53. USERID BIGINT,
  54. ) WITH (DATASOURCE="users", FORMAT="AVRO", KEY="USERID");`,
  55. err: "key topic1 already exist, delete it before creating a new one",
  56. },
  57. {
  58. s: `EXPLAIN STREAM topic1;`,
  59. r: []string{"TO BE SUPPORTED"},
  60. },
  61. {
  62. s: `DESCRIBE STREAM topic1;`,
  63. r: []string{"Fields\n--------------------------------------------------------------------------------\nUSERID\tbigint\nFIRST_NAME\tstring\nLAST_NAME\tstring\nNICKNAMES\t" +
  64. "array(string)\nGender\tboolean\nADDRESS\tstruct(STREET_NAME string, NUMBER bigint)\n\n" +
  65. "DATASOURCE: users\nFORMAT: AVRO\nKEY: USERID\n"},
  66. },
  67. {
  68. s: `SHOW STREAMS;`,
  69. r: []string{"topic1"},
  70. },
  71. {
  72. s: `DROP STREAM topic1;`,
  73. r: []string{"stream topic1 dropped"},
  74. },
  75. {
  76. s: `DESCRIBE STREAM topic1;`,
  77. err: "stream topic1 not found",
  78. },
  79. {
  80. s: `DROP STREAM topic1;`,
  81. err: "Key not found",
  82. },
  83. }
  84. fmt.Printf("The test bucket size is %d.\n\n", len(tests))
  85. streamDB := path.Join(BadgerDir, "streamTest")
  86. for i, tt := range tests {
  87. results, err := NewStreamProcessor(tt.s, streamDB).Exec()
  88. if !reflect.DeepEqual(tt.err, errstring(err)) {
  89. t.Errorf("%d. %q: error mismatch:\n exp=%s\n got=%s\n\n", i, tt.s, tt.err, err)
  90. } else if tt.err == "" {
  91. if !reflect.DeepEqual(tt.r, results) {
  92. t.Errorf("%d. %q\n\nstmt mismatch:\n\ngot=%#v\n\n", i, tt.s, results)
  93. }
  94. }
  95. }
  96. }
  97. func createStreams(t *testing.T){
  98. demo := `CREATE STREAM demo (
  99. color STRING,
  100. size BIGINT,
  101. ts BIGINT
  102. ) WITH (DATASOURCE="demo", FORMAT="json", KEY="ts");`
  103. _, err := NewStreamProcessor(demo, path.Join(BadgerDir, "stream")).Exec()
  104. if err != nil{
  105. t.Log(err)
  106. }
  107. demo1 := `CREATE STREAM demo1 (
  108. temp FLOAT,
  109. hum BIGINT,
  110. ts BIGINT
  111. ) WITH (DATASOURCE="demo1", FORMAT="json", KEY="ts");`
  112. _, err = NewStreamProcessor(demo1, path.Join(BadgerDir, "stream")).Exec()
  113. if err != nil{
  114. t.Log(err)
  115. }
  116. sessionDemo := `CREATE STREAM sessionDemo (
  117. temp FLOAT,
  118. hum BIGINT,
  119. ts BIGINT
  120. ) WITH (DATASOURCE="sessionDemo", FORMAT="json", KEY="ts");`
  121. _, err = NewStreamProcessor(sessionDemo, path.Join(BadgerDir, "stream")).Exec()
  122. if err != nil{
  123. t.Log(err)
  124. }
  125. }
  126. func dropStreams(t *testing.T){
  127. demo := `DROP STREAM demo`
  128. _, err := NewStreamProcessor(demo, path.Join(BadgerDir, "stream")).Exec()
  129. if err != nil{
  130. t.Log(err)
  131. }
  132. demo1 := `DROP STREAM demo1`
  133. _, err = NewStreamProcessor(demo1, path.Join(BadgerDir, "stream")).Exec()
  134. if err != nil{
  135. t.Log(err)
  136. }
  137. sessionDemo := `DROP STREAM sessionDemo`
  138. _, err = NewStreamProcessor(sessionDemo, path.Join(BadgerDir, "stream")).Exec()
  139. if err != nil{
  140. t.Log(err)
  141. }
  142. }
  143. func getMockSource(name string, done chan<- struct{}, size int) *test.MockSource{
  144. var data []*xsql.Tuple
  145. switch name{
  146. case "demo":
  147. data = []*xsql.Tuple{
  148. {
  149. Emitter: name,
  150. Message: map[string]interface{}{
  151. "color": "red",
  152. "size": 3,
  153. "ts": 1541152486013,
  154. },
  155. Timestamp: 1541152486013,
  156. },
  157. {
  158. Emitter: name,
  159. Message: map[string]interface{}{
  160. "color": "blue",
  161. "size": 6,
  162. "ts": 1541152486822,
  163. },
  164. Timestamp: 1541152486822,
  165. },
  166. {
  167. Emitter: name,
  168. Message: map[string]interface{}{
  169. "color": "blue",
  170. "size": 2,
  171. "ts": 1541152487632,
  172. },
  173. Timestamp: 1541152487632,
  174. },
  175. {
  176. Emitter: name,
  177. Message: map[string]interface{}{
  178. "color": "yellow",
  179. "size": 4,
  180. "ts": 1541152488442,
  181. },
  182. Timestamp: 1541152488442,
  183. },
  184. {
  185. Emitter: name,
  186. Message: map[string]interface{}{
  187. "color": "red",
  188. "size": 1,
  189. "ts": 1541152489252,
  190. },
  191. Timestamp: 1541152489252,
  192. },
  193. }
  194. case "demo1":
  195. data = []*xsql.Tuple{
  196. {
  197. Emitter: name,
  198. Message: map[string]interface{}{
  199. "temp": 25.5,
  200. "hum": 65,
  201. "ts": 1541152486013,
  202. },
  203. Timestamp: 1541152486013,
  204. },
  205. {
  206. Emitter: name,
  207. Message: map[string]interface{}{
  208. "temp": 27.5,
  209. "hum": 59,
  210. "ts": 1541152486823,
  211. },
  212. Timestamp: 1541152486823,
  213. },
  214. {
  215. Emitter: name,
  216. Message: map[string]interface{}{
  217. "temp": 28.1,
  218. "hum": 75,
  219. "ts": 1541152487632,
  220. },
  221. Timestamp: 1541152487632,
  222. },
  223. {
  224. Emitter: name,
  225. Message: map[string]interface{}{
  226. "temp": 27.4,
  227. "hum": 80,
  228. "ts": 1541152488442,
  229. },
  230. Timestamp: 1541152488442,
  231. },
  232. {
  233. Emitter: name,
  234. Message: map[string]interface{}{
  235. "temp": 25.5,
  236. "hum": 62,
  237. "ts": 1541152489252,
  238. },
  239. Timestamp: 1541152489252,
  240. },
  241. }
  242. case "sessionDemo":
  243. data = []*xsql.Tuple{
  244. {
  245. Emitter: name,
  246. Message: map[string]interface{}{
  247. "temp": 25.5,
  248. "hum": 65,
  249. "ts": 1541152486013,
  250. },
  251. Timestamp: 1541152486013,
  252. },
  253. {
  254. Emitter: name,
  255. Message: map[string]interface{}{
  256. "temp": 27.5,
  257. "hum": 59,
  258. "ts": 1541152486823,
  259. },
  260. Timestamp: 1541152486823,
  261. },
  262. {
  263. Emitter: name,
  264. Message: map[string]interface{}{
  265. "temp": 28.1,
  266. "hum": 75,
  267. "ts": 1541152487932,
  268. },
  269. Timestamp: 1541152487932,
  270. },
  271. {
  272. Emitter: name,
  273. Message: map[string]interface{}{
  274. "temp": 27.4,
  275. "hum": 80,
  276. "ts": 1541152488442,
  277. },
  278. Timestamp: 1541152488442,
  279. },
  280. {
  281. Emitter: name,
  282. Message: map[string]interface{}{
  283. "temp": 25.5,
  284. "hum": 62,
  285. "ts": 1541152489252,
  286. },
  287. Timestamp: 1541152489252,
  288. },
  289. {
  290. Emitter: name,
  291. Message: map[string]interface{}{
  292. "temp": 26.2,
  293. "hum": 63,
  294. "ts": 1541152490062,
  295. },
  296. Timestamp: 1541152490062,
  297. },
  298. {
  299. Emitter: name,
  300. Message: map[string]interface{}{
  301. "temp": 26.8,
  302. "hum": 71,
  303. "ts": 1541152490872,
  304. },
  305. Timestamp: 1541152490872,
  306. },
  307. {
  308. Emitter: name,
  309. Message: map[string]interface{}{
  310. "temp": 28.9,
  311. "hum": 85,
  312. "ts": 1541152491682,
  313. },
  314. Timestamp: 1541152491682,
  315. },
  316. {
  317. Emitter: name,
  318. Message: map[string]interface{}{
  319. "temp": 29.1,
  320. "hum": 92,
  321. "ts": 1541152492492,
  322. },
  323. Timestamp: 1541152492492,
  324. },
  325. {
  326. Emitter: name,
  327. Message: map[string]interface{}{
  328. "temp": 32.2,
  329. "hum": 99,
  330. "ts": 1541152493202,
  331. },
  332. Timestamp: 1541152493202,
  333. },
  334. {
  335. Emitter: name,
  336. Message: map[string]interface{}{
  337. "temp": 30.9,
  338. "hum": 87,
  339. "ts": 1541152494112,
  340. },
  341. Timestamp: 1541152494112,
  342. },
  343. }
  344. }
  345. return test.NewMockSource(data[:size], name, done, false)
  346. }
  347. func TestSingleSQL(t *testing.T) {
  348. var tests = []struct {
  349. name string
  350. sql string
  351. r [][]map[string]interface{}
  352. }{
  353. {
  354. name: `rule1`,
  355. sql: `SELECT * FROM demo`,
  356. r: [][]map[string]interface{}{
  357. {{
  358. "color": "red",
  359. "size": float64(3),
  360. "ts": float64(1541152486013),
  361. }},
  362. {{
  363. "color": "blue",
  364. "size": float64(6),
  365. "ts": float64(1541152486822),
  366. }},
  367. {{
  368. "color": "blue",
  369. "size": float64(2),
  370. "ts": float64(1541152487632),
  371. }},
  372. {{
  373. "color": "yellow",
  374. "size": float64(4),
  375. "ts": float64(1541152488442),
  376. }},
  377. {{
  378. "color": "red",
  379. "size": float64(1),
  380. "ts": float64(1541152489252),
  381. }},
  382. },
  383. }, {
  384. name: `rule2`,
  385. sql: `SELECT color, ts FROM demo where size > 3`,
  386. r: [][]map[string]interface{}{
  387. {{
  388. "color": "blue",
  389. "ts": float64(1541152486822),
  390. }},
  391. {{
  392. "color": "yellow",
  393. "ts": float64(1541152488442),
  394. }},
  395. },
  396. },
  397. }
  398. fmt.Printf("The test bucket size is %d.\n\n", len(tests))
  399. createStreams(t)
  400. defer dropStreams(t)
  401. done := make(chan struct{})
  402. defer close(done)
  403. for i, tt := range tests {
  404. p := NewRuleProcessor(BadgerDir)
  405. parser := xsql.NewParser(strings.NewReader(tt.sql))
  406. var sources []xstream.Source
  407. if stmt, err := xsql.Language.Parse(parser); err != nil{
  408. t.Errorf("parse sql %s error: %s", tt.sql , err)
  409. }else {
  410. if selectStmt, ok := stmt.(*xsql.SelectStatement); !ok {
  411. t.Errorf("sql %s is not a select statement", tt.sql)
  412. } else {
  413. streams := xsql.GetStreams(selectStmt)
  414. for _, stream := range streams{
  415. source := getMockSource(stream, done, 5)
  416. sources = append(sources, source)
  417. }
  418. }
  419. }
  420. tp, inputs, err := p.createTopoWithSources(&xstream.Rule{Id:tt.name, Sql: tt.sql}, sources)
  421. if err != nil{
  422. t.Error(err)
  423. }
  424. sink := test.NewMockSink("mockSink", tt.name)
  425. tp.AddSink(inputs, sink)
  426. count := len(sources)
  427. errCh := tp.Open()
  428. func(){
  429. for{
  430. select{
  431. case err = <- errCh:
  432. t.Log(err)
  433. tp.Cancel()
  434. return
  435. case <- done:
  436. count--
  437. log.Infof("%d sources remaining", count)
  438. if count <= 0{
  439. log.Info("stream stopping")
  440. time.Sleep(1 * time.Second)
  441. tp.Cancel()
  442. return
  443. }
  444. default:
  445. }
  446. }
  447. }()
  448. results := sink.GetResults()
  449. var maps [][]map[string]interface{}
  450. for _, v := range results{
  451. var mapRes []map[string]interface{}
  452. err := json.Unmarshal(v, &mapRes)
  453. if err != nil {
  454. t.Errorf("Failed to parse the input into map")
  455. continue
  456. }
  457. maps = append(maps, mapRes)
  458. }
  459. if !reflect.DeepEqual(tt.r, maps) {
  460. t.Errorf("%d. %q\n\nresult mismatch:\n\nexp=%#v\n\ngot=%#v\n\n", i, tt.sql, tt.r, maps)
  461. }
  462. }
  463. }
  464. func TestWindow(t *testing.T) {
  465. common.IsTesting = true
  466. var tests = []struct {
  467. name string
  468. sql string
  469. size int
  470. r [][]map[string]interface{}
  471. }{
  472. {
  473. name: `rule1`,
  474. sql: `SELECT * FROM demo GROUP BY HOPPINGWINDOW(ss, 2, 1)`,
  475. size: 5,
  476. r: [][]map[string]interface{}{
  477. {{
  478. "color": "red",
  479. "size": float64(3),
  480. "ts": float64(1541152486013),
  481. },{
  482. "color": "blue",
  483. "size": float64(6),
  484. "ts": float64(1541152486822),
  485. }},
  486. {{
  487. "color": "red",
  488. "size": float64(3),
  489. "ts": float64(1541152486013),
  490. },{
  491. "color": "blue",
  492. "size": float64(6),
  493. "ts": float64(1541152486822),
  494. },{
  495. "color": "blue",
  496. "size": float64(2),
  497. "ts": float64(1541152487632),
  498. }},
  499. {{
  500. "color": "blue",
  501. "size": float64(2),
  502. "ts": float64(1541152487632),
  503. },{
  504. "color": "yellow",
  505. "size": float64(4),
  506. "ts": float64(1541152488442),
  507. }},
  508. },
  509. }, {
  510. name: `rule2`,
  511. sql: `SELECT color, ts FROM demo where size > 2 GROUP BY tumblingwindow(ss, 1)`,
  512. size: 5,
  513. r: [][]map[string]interface{}{
  514. {{
  515. "color": "red",
  516. "ts": float64(1541152486013),
  517. },{
  518. "color": "blue",
  519. "ts": float64(1541152486822),
  520. }},
  521. {{
  522. "color": "yellow",
  523. "ts": float64(1541152488442),
  524. }},
  525. },
  526. }, {
  527. name: `rule3`,
  528. sql: `SELECT color, temp, ts FROM demo INNER JOIN demo1 ON demo.ts = demo1.ts GROUP BY SlidingWindow(ss, 1)`,
  529. size: 5,
  530. r: [][]map[string]interface{}{
  531. {{
  532. "color": "red",
  533. "temp": 25.5,
  534. "ts": float64(1541152486013),
  535. }},{{
  536. "color": "red",
  537. "temp": 25.5,
  538. "ts": float64(1541152486013),
  539. }},{{
  540. "color": "red",
  541. "temp": 25.5,
  542. "ts": float64(1541152486013),
  543. }},{{
  544. "color": "blue",
  545. "temp": 28.1,
  546. "ts": float64(1541152487632),
  547. }},{{
  548. "color": "blue",
  549. "temp": 28.1,
  550. "ts": float64(1541152487632),
  551. }},{{
  552. "color": "blue",
  553. "temp": 28.1,
  554. "ts": float64(1541152487632),
  555. },{
  556. "color": "yellow",
  557. "temp": 27.4,
  558. "ts": float64(1541152488442),
  559. }},{{
  560. "color": "yellow",
  561. "temp": 27.4,
  562. "ts": float64(1541152488442),
  563. }},{{
  564. "color": "yellow",
  565. "temp": 27.4,
  566. "ts": float64(1541152488442),
  567. },{
  568. "color": "red",
  569. "temp": 25.5,
  570. "ts": float64(1541152489252),
  571. }},
  572. },
  573. }, {
  574. name: `rule4`,
  575. sql: `SELECT color FROM demo GROUP BY SlidingWindow(ss, 2), color ORDER BY color`,
  576. size: 5,
  577. r: [][]map[string]interface{}{
  578. {{
  579. "color": "red",
  580. }},{{
  581. "color": "blue",
  582. },{
  583. "color": "red",
  584. }},{{
  585. "color": "blue",
  586. },{
  587. "color": "red",
  588. }},{{
  589. "color": "blue",
  590. },{
  591. "color": "yellow",
  592. }},{{
  593. "color": "blue",
  594. }, {
  595. "color": "red",
  596. },{
  597. "color": "yellow",
  598. }},
  599. },
  600. },{
  601. name: `rule5`,
  602. sql: `SELECT temp FROM sessionDemo GROUP BY SessionWindow(ss, 2, 1) `,
  603. size: 11,
  604. r: [][]map[string]interface{}{
  605. {{
  606. "temp": 25.5,
  607. },{
  608. "temp": 27.5,
  609. }},{{
  610. "temp": 28.1,
  611. },{
  612. "temp": 27.4,
  613. },{
  614. "temp": 25.5,
  615. }},{{
  616. "temp": 26.2,
  617. },{
  618. "temp": 26.8,
  619. },{
  620. "temp": 28.9,
  621. },{
  622. "temp": 29.1,
  623. },{
  624. "temp": 32.2,
  625. }},
  626. },
  627. },{
  628. name: `rule6`,
  629. sql: `SELECT max(temp) as m, count(color) as c FROM demo INNER JOIN demo1 ON demo.ts = demo1.ts GROUP BY SlidingWindow(ss, 1)`,
  630. size: 5,
  631. r: [][]map[string]interface{}{
  632. {{
  633. "m": 25.5,
  634. "c": float64(1),
  635. }},{{
  636. "m": 25.5,
  637. "c": float64(1),
  638. }},{{
  639. "m": 25.5,
  640. "c": float64(1),
  641. }},{{
  642. "m": 28.1,
  643. "c": float64(1),
  644. }},{{
  645. "m": 28.1,
  646. "c": float64(1),
  647. }},{{
  648. "m": 28.1,
  649. "c": float64(2),
  650. }},{{
  651. "m": 27.4,
  652. "c": float64(1),
  653. }},{{
  654. "m": 27.4,
  655. "c": float64(2),
  656. }},
  657. },
  658. },
  659. }
  660. fmt.Printf("The test bucket size is %d.\n\n", len(tests))
  661. createStreams(t)
  662. defer dropStreams(t)
  663. done := make(chan struct{})
  664. defer close(done)
  665. common.ResetMockTicker()
  666. for i, tt := range tests {
  667. p := NewRuleProcessor(BadgerDir)
  668. parser := xsql.NewParser(strings.NewReader(tt.sql))
  669. var sources []xstream.Source
  670. if stmt, err := xsql.Language.Parse(parser); err != nil{
  671. t.Errorf("parse sql %s error: %s", tt.sql , err)
  672. }else {
  673. if selectStmt, ok := stmt.(*xsql.SelectStatement); !ok {
  674. t.Errorf("sql %s is not a select statement", tt.sql)
  675. } else {
  676. streams := xsql.GetStreams(selectStmt)
  677. for _, stream := range streams{
  678. source := getMockSource(stream, done, tt.size)
  679. sources = append(sources, source)
  680. }
  681. }
  682. }
  683. tp, inputs, err := p.createTopoWithSources(&xstream.Rule{Id:tt.name, Sql: tt.sql}, sources)
  684. if err != nil{
  685. t.Error(err)
  686. }
  687. sink := test.NewMockSink("mockSink", tt.name)
  688. tp.AddSink(inputs, sink)
  689. count := len(sources)
  690. errCh := tp.Open()
  691. func(){
  692. for{
  693. select{
  694. case err = <- errCh:
  695. t.Log(err)
  696. tp.Cancel()
  697. return
  698. case <- done:
  699. count--
  700. log.Infof("%d sources remaining", count)
  701. if count <= 0{
  702. log.Info("stream stopping")
  703. time.Sleep(1 * time.Second)
  704. tp.Cancel()
  705. return
  706. }
  707. default:
  708. }
  709. }
  710. }()
  711. results := sink.GetResults()
  712. var maps [][]map[string]interface{}
  713. for _, v := range results{
  714. var mapRes []map[string]interface{}
  715. err := json.Unmarshal(v, &mapRes)
  716. if err != nil {
  717. t.Errorf("Failed to parse the input into map")
  718. continue
  719. }
  720. maps = append(maps, mapRes)
  721. }
  722. if !reflect.DeepEqual(tt.r, maps) {
  723. t.Errorf("%d. %q\n\nresult mismatch:\n\nexp=%#v\n\ngot=%#v\n\n", i, tt.sql, tt.r, maps)
  724. }
  725. }
  726. }
  727. func createEventStreams(t *testing.T){
  728. demo := `CREATE STREAM demoE (
  729. color STRING,
  730. size BIGINT,
  731. ts BIGINT
  732. ) WITH (DATASOURCE="demoE", FORMAT="json", KEY="ts", TIMESTAMP="ts");`
  733. _, err := NewStreamProcessor(demo, path.Join(BadgerDir, "stream")).Exec()
  734. if err != nil{
  735. t.Log(err)
  736. }
  737. demo1 := `CREATE STREAM demo1E (
  738. temp FLOAT,
  739. hum BIGINT,
  740. ts BIGINT
  741. ) WITH (DATASOURCE="demo1E", FORMAT="json", KEY="ts", TIMESTAMP="ts");`
  742. _, err = NewStreamProcessor(demo1, path.Join(BadgerDir, "stream")).Exec()
  743. if err != nil{
  744. t.Log(err)
  745. }
  746. sessionDemo := `CREATE STREAM sessionDemoE (
  747. temp FLOAT,
  748. hum BIGINT,
  749. ts BIGINT
  750. ) WITH (DATASOURCE="sessionDemoE", FORMAT="json", KEY="ts", TIMESTAMP="ts");`
  751. _, err = NewStreamProcessor(sessionDemo, path.Join(BadgerDir, "stream")).Exec()
  752. if err != nil{
  753. t.Log(err)
  754. }
  755. }
  756. func dropEventStreams(t *testing.T){
  757. demo := `DROP STREAM demoE`
  758. _, err := NewStreamProcessor(demo, path.Join(BadgerDir, "stream")).Exec()
  759. if err != nil{
  760. t.Log(err)
  761. }
  762. demo1 := `DROP STREAM demo1E`
  763. _, err = NewStreamProcessor(demo1, path.Join(BadgerDir, "stream")).Exec()
  764. if err != nil{
  765. t.Log(err)
  766. }
  767. sessionDemo := `DROP STREAM sessionDemoE`
  768. _, err = NewStreamProcessor(sessionDemo, path.Join(BadgerDir, "stream")).Exec()
  769. if err != nil{
  770. t.Log(err)
  771. }
  772. }
  773. func getEventMockSource(name string, done chan<- struct{}, size int) *test.MockSource{
  774. var data []*xsql.Tuple
  775. switch name{
  776. case "demoE":
  777. data = []*xsql.Tuple{
  778. {
  779. Emitter: name,
  780. Message: map[string]interface{}{
  781. "color": "red",
  782. "size": 3,
  783. "ts": 1541152486013,
  784. },
  785. Timestamp: 1541152486013,
  786. },
  787. {
  788. Emitter: name,
  789. Message: map[string]interface{}{
  790. "color": "blue",
  791. "size": 2,
  792. "ts": 1541152487632,
  793. },
  794. Timestamp: 1541152487632,
  795. },
  796. {
  797. Emitter: name,
  798. Message: map[string]interface{}{
  799. "color": "red",
  800. "size": 1,
  801. "ts": 1541152489252,
  802. },
  803. Timestamp: 1541152489252,
  804. },
  805. { //dropped item
  806. Emitter: name,
  807. Message: map[string]interface{}{
  808. "color": "blue",
  809. "size": 6,
  810. "ts": 1541152486822,
  811. },
  812. Timestamp: 1541152486822,
  813. },
  814. {
  815. Emitter: name,
  816. Message: map[string]interface{}{
  817. "color": "yellow",
  818. "size": 4,
  819. "ts": 1541152488442,
  820. },
  821. Timestamp: 1541152488442,
  822. },
  823. { //To lift the watermark and issue all windows
  824. Emitter: name,
  825. Message: map[string]interface{}{
  826. "color": "yellow",
  827. "size": 4,
  828. "ts": 1541152492342,
  829. },
  830. Timestamp: 1541152488442,
  831. },
  832. }
  833. case "demo1E":
  834. data = []*xsql.Tuple{
  835. {
  836. Emitter: name,
  837. Message: map[string]interface{}{
  838. "temp": 27.5,
  839. "hum": 59,
  840. "ts": 1541152486823,
  841. },
  842. Timestamp: 1541152486823,
  843. },
  844. {
  845. Emitter: name,
  846. Message: map[string]interface{}{
  847. "temp": 25.5,
  848. "hum": 65,
  849. "ts": 1541152486013,
  850. },
  851. Timestamp: 1541152486013,
  852. },
  853. {
  854. Emitter: name,
  855. Message: map[string]interface{}{
  856. "temp": 27.4,
  857. "hum": 80,
  858. "ts": 1541152488442,
  859. },
  860. Timestamp: 1541152488442,
  861. },
  862. {
  863. Emitter: name,
  864. Message: map[string]interface{}{
  865. "temp": 28.1,
  866. "hum": 75,
  867. "ts": 1541152487632,
  868. },
  869. Timestamp: 1541152487632,
  870. },
  871. {
  872. Emitter: name,
  873. Message: map[string]interface{}{
  874. "temp": 25.5,
  875. "hum": 62,
  876. "ts": 1541152489252,
  877. },
  878. Timestamp: 1541152489252,
  879. },
  880. {
  881. Emitter: name,
  882. Message: map[string]interface{}{
  883. "temp": 25.5,
  884. "hum": 62,
  885. "ts": 1541152499252,
  886. },
  887. Timestamp: 1541152499252,
  888. },
  889. }
  890. case "sessionDemoE":
  891. data = []*xsql.Tuple{
  892. {
  893. Emitter: name,
  894. Message: map[string]interface{}{
  895. "temp": 25.5,
  896. "hum": 65,
  897. "ts": 1541152486013,
  898. },
  899. Timestamp: 1541152486013,
  900. },
  901. {
  902. Emitter: name,
  903. Message: map[string]interface{}{
  904. "temp": 28.1,
  905. "hum": 75,
  906. "ts": 1541152487932,
  907. },
  908. Timestamp: 1541152487932,
  909. },
  910. {
  911. Emitter: name,
  912. Message: map[string]interface{}{
  913. "temp": 27.5,
  914. "hum": 59,
  915. "ts": 1541152486823,
  916. },
  917. Timestamp: 1541152486823,
  918. },
  919. {
  920. Emitter: name,
  921. Message: map[string]interface{}{
  922. "temp": 25.5,
  923. "hum": 62,
  924. "ts": 1541152489252,
  925. },
  926. Timestamp: 1541152489252,
  927. },
  928. {
  929. Emitter: name,
  930. Message: map[string]interface{}{
  931. "temp": 27.4,
  932. "hum": 80,
  933. "ts": 1541152488442,
  934. },
  935. Timestamp: 1541152488442,
  936. },
  937. {
  938. Emitter: name,
  939. Message: map[string]interface{}{
  940. "temp": 26.2,
  941. "hum": 63,
  942. "ts": 1541152490062,
  943. },
  944. Timestamp: 1541152490062,
  945. },
  946. {
  947. Emitter: name,
  948. Message: map[string]interface{}{
  949. "temp": 28.9,
  950. "hum": 85,
  951. "ts": 1541152491682,
  952. },
  953. Timestamp: 1541152491682,
  954. },
  955. {
  956. Emitter: name,
  957. Message: map[string]interface{}{
  958. "temp": 26.8,
  959. "hum": 71,
  960. "ts": 1541152490872,
  961. },
  962. Timestamp: 1541152490872,
  963. },
  964. {
  965. Emitter: name,
  966. Message: map[string]interface{}{
  967. "temp": 29.1,
  968. "hum": 92,
  969. "ts": 1541152492492,
  970. },
  971. Timestamp: 1541152492492,
  972. },
  973. {
  974. Emitter: name,
  975. Message: map[string]interface{}{
  976. "temp": 30.9,
  977. "hum": 87,
  978. "ts": 1541152494112,
  979. },
  980. Timestamp: 1541152494112,
  981. },
  982. {
  983. Emitter: name,
  984. Message: map[string]interface{}{
  985. "temp": 32.2,
  986. "hum": 99,
  987. "ts": 1541152493202,
  988. },
  989. Timestamp: 1541152493202,
  990. },
  991. {
  992. Emitter: name,
  993. Message: map[string]interface{}{
  994. "temp": 32.2,
  995. "hum": 99,
  996. "ts": 1541152499202,
  997. },
  998. Timestamp: 1541152499202,
  999. },
  1000. }
  1001. }
  1002. return test.NewMockSource(data[:size], name, done, true)
  1003. }
  1004. func TestEventWindow(t *testing.T) {
  1005. common.IsTesting = true
  1006. var tests = []struct {
  1007. name string
  1008. sql string
  1009. size int
  1010. r [][]map[string]interface{}
  1011. }{
  1012. {
  1013. name: `rule1`,
  1014. sql: `SELECT * FROM demoE GROUP BY HOPPINGWINDOW(ss, 2, 1)`,
  1015. size: 6,
  1016. r: [][]map[string]interface{}{
  1017. {{
  1018. "color": "red",
  1019. "size": float64(3),
  1020. "ts": float64(1541152486013),
  1021. }},
  1022. {{
  1023. "color": "red",
  1024. "size": float64(3),
  1025. "ts": float64(1541152486013),
  1026. },{
  1027. "color": "blue",
  1028. "size": float64(2),
  1029. "ts": float64(1541152487632),
  1030. }},
  1031. {{
  1032. "color": "blue",
  1033. "size": float64(2),
  1034. "ts": float64(1541152487632),
  1035. },{
  1036. "color": "yellow",
  1037. "size": float64(4),
  1038. "ts": float64(1541152488442),
  1039. }},{{
  1040. "color": "yellow",
  1041. "size": float64(4),
  1042. "ts": float64(1541152488442),
  1043. },{
  1044. "color": "red",
  1045. "size": float64(1),
  1046. "ts": float64(1541152489252),
  1047. }},{{
  1048. "color": "red",
  1049. "size": float64(1),
  1050. "ts": float64(1541152489252),
  1051. }},
  1052. },
  1053. }, {
  1054. name: `rule2`,
  1055. sql: `SELECT color, ts FROM demoE where size > 2 GROUP BY tumblingwindow(ss, 1)`,
  1056. size: 6,
  1057. r: [][]map[string]interface{}{
  1058. {{
  1059. "color": "red",
  1060. "ts": float64(1541152486013),
  1061. }},
  1062. {{
  1063. "color": "yellow",
  1064. "ts": float64(1541152488442),
  1065. }},
  1066. },
  1067. }, {
  1068. name: `rule3`,
  1069. sql: `SELECT color, temp, ts FROM demoE INNER JOIN demo1E ON demoE.ts = demo1E.ts GROUP BY SlidingWindow(ss, 1)`,
  1070. size: 6,
  1071. r: [][]map[string]interface{}{
  1072. {{
  1073. "color": "red",
  1074. "temp": 25.5,
  1075. "ts": float64(1541152486013),
  1076. }},{{
  1077. "color": "red",
  1078. "temp": 25.5,
  1079. "ts": float64(1541152486013),
  1080. }},{{
  1081. "color": "blue",
  1082. "temp": 28.1,
  1083. "ts": float64(1541152487632),
  1084. }},{{
  1085. "color": "blue",
  1086. "temp": 28.1,
  1087. "ts": float64(1541152487632),
  1088. },{
  1089. "color": "yellow",
  1090. "temp": 27.4,
  1091. "ts": float64(1541152488442),
  1092. }},{{
  1093. "color": "yellow",
  1094. "temp": 27.4,
  1095. "ts": float64(1541152488442),
  1096. },{
  1097. "color": "red",
  1098. "temp": 25.5,
  1099. "ts": float64(1541152489252),
  1100. }},
  1101. },
  1102. }, {
  1103. name: `rule4`,
  1104. sql: `SELECT color FROM demoE GROUP BY SlidingWindow(ss, 2), color ORDER BY color`,
  1105. size: 6,
  1106. r: [][]map[string]interface{}{
  1107. {{
  1108. "color": "red",
  1109. }},{{
  1110. "color": "blue",
  1111. },{
  1112. "color": "red",
  1113. }},{{
  1114. "color": "blue",
  1115. },{
  1116. "color": "yellow",
  1117. }},{{
  1118. "color": "blue",
  1119. }, {
  1120. "color": "red",
  1121. },{
  1122. "color": "yellow",
  1123. }},
  1124. },
  1125. },{
  1126. name: `rule5`,
  1127. sql: `SELECT temp FROM sessionDemoE GROUP BY SessionWindow(ss, 2, 1) `,
  1128. size: 12,
  1129. r: [][]map[string]interface{}{
  1130. {{
  1131. "temp": 25.5,
  1132. }},{{
  1133. "temp": 28.1,
  1134. },{
  1135. "temp": 27.4,
  1136. },{
  1137. "temp": 25.5,
  1138. }},{{
  1139. "temp": 26.2,
  1140. },{
  1141. "temp": 26.8,
  1142. },{
  1143. "temp": 28.9,
  1144. },{
  1145. "temp": 29.1,
  1146. },{
  1147. "temp": 32.2,
  1148. }},{{
  1149. "temp": 30.9,
  1150. }},
  1151. },
  1152. },{
  1153. name: `rule6`,
  1154. sql: `SELECT max(temp) as m, count(color) as c FROM demoE INNER JOIN demo1E ON demoE.ts = demo1E.ts GROUP BY SlidingWindow(ss, 1)`,
  1155. size: 6,
  1156. r: [][]map[string]interface{}{
  1157. {{
  1158. "m": 25.5,
  1159. "c": float64(1),
  1160. }},{{
  1161. "m": 25.5,
  1162. "c": float64(1),
  1163. }},{{
  1164. "m": 28.1,
  1165. "c": float64(1),
  1166. }},{{
  1167. "m": 28.1,
  1168. "c": float64(2),
  1169. }},{{
  1170. "m": 27.4,
  1171. "c": float64(2),
  1172. }},
  1173. },
  1174. },
  1175. }
  1176. fmt.Printf("The test bucket size is %d.\n\n", len(tests))
  1177. createEventStreams(t)
  1178. defer dropEventStreams(t)
  1179. done := make(chan struct{})
  1180. defer close(done)
  1181. common.ResetMockTicker()
  1182. //mock ticker
  1183. realTicker := time.NewTicker(500 * time.Millisecond)
  1184. tickerDone := make(chan bool)
  1185. go func(){
  1186. ticker := common.GetTicker(1000).(*common.MockTicker)
  1187. timer := common.GetTimer(1000).(*common.MockTimer)
  1188. for {
  1189. select {
  1190. case <-tickerDone:
  1191. log.Infof("real ticker exiting...")
  1192. return
  1193. case t := <-realTicker.C:
  1194. ts := common.TimeToUnixMilli(t)
  1195. if ticker != nil {
  1196. go ticker.DoTick(ts)
  1197. }
  1198. if timer != nil {
  1199. go timer.DoTick(ts)
  1200. }
  1201. }
  1202. }
  1203. }()
  1204. for i, tt := range tests {
  1205. p := NewRuleProcessor(BadgerDir)
  1206. parser := xsql.NewParser(strings.NewReader(tt.sql))
  1207. var sources []xstream.Source
  1208. if stmt, err := xsql.Language.Parse(parser); err != nil{
  1209. t.Errorf("parse sql %s error: %s", tt.sql , err)
  1210. }else {
  1211. if selectStmt, ok := stmt.(*xsql.SelectStatement); !ok {
  1212. t.Errorf("sql %s is not a select statement", tt.sql)
  1213. } else {
  1214. streams := xsql.GetStreams(selectStmt)
  1215. for _, stream := range streams{
  1216. source := getEventMockSource(stream, done, tt.size)
  1217. sources = append(sources, source)
  1218. }
  1219. }
  1220. }
  1221. tp, inputs, err := p.createTopoWithSources(&xstream.Rule{
  1222. Id:tt.name, Sql: tt.sql,
  1223. Options: map[string]interface{}{
  1224. "isEventTime": true,
  1225. "lateTolerance": float64(1000),
  1226. },
  1227. }, sources)
  1228. if err != nil{
  1229. t.Error(err)
  1230. }
  1231. sink := test.NewMockSink("mockSink", tt.name)
  1232. tp.AddSink(inputs, sink)
  1233. count := len(sources)
  1234. errCh := tp.Open()
  1235. func(){
  1236. for{
  1237. select{
  1238. case err = <- errCh:
  1239. t.Log(err)
  1240. tp.Cancel()
  1241. return
  1242. case <- done:
  1243. count--
  1244. log.Infof("%d sources remaining", count)
  1245. if count <= 0{
  1246. log.Info("stream stopping")
  1247. time.Sleep(1 * time.Second)
  1248. tp.Cancel()
  1249. return
  1250. }
  1251. default:
  1252. }
  1253. }
  1254. }()
  1255. results := sink.GetResults()
  1256. var maps [][]map[string]interface{}
  1257. for _, v := range results{
  1258. var mapRes []map[string]interface{}
  1259. err := json.Unmarshal(v, &mapRes)
  1260. if err != nil {
  1261. t.Errorf("Failed to parse the input into map")
  1262. continue
  1263. }
  1264. maps = append(maps, mapRes)
  1265. }
  1266. if !reflect.DeepEqual(tt.r, maps) {
  1267. t.Errorf("%d. %q\n\nresult mismatch:\n\nexp=%#v\n\ngot=%#v\n\n", i, tt.sql, tt.r, maps)
  1268. }
  1269. }
  1270. realTicker.Stop()
  1271. tickerDone <- true
  1272. close(tickerDone)
  1273. }
  1274. func errstring(err error) string {
  1275. if err != nil {
  1276. return err.Error()
  1277. }
  1278. return ""
  1279. }