stream_test.go 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382
  1. // Copyright 2021-2023 EMQ Technologies Co., Ltd.
  2. //
  3. // Licensed under the Apache License, Version 2.0 (the "License");
  4. // you may not use this file except in compliance with the License.
  5. // You may obtain a copy of the License at
  6. //
  7. // http://www.apache.org/licenses/LICENSE-2.0
  8. //
  9. // Unless required by applicable law or agreed to in writing, software
  10. // distributed under the License is distributed on an "AS IS" BASIS,
  11. // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  12. // See the License for the specific language governing permissions and
  13. // limitations under the License.
  14. package processor
  15. import (
  16. "fmt"
  17. "os"
  18. "path/filepath"
  19. "reflect"
  20. "strconv"
  21. "testing"
  22. "github.com/gdexlab/go-render/render"
  23. "github.com/lf-edge/ekuiper/internal/conf"
  24. "github.com/lf-edge/ekuiper/internal/schema"
  25. "github.com/lf-edge/ekuiper/internal/testx"
  26. "github.com/lf-edge/ekuiper/pkg/ast"
  27. )
  28. func init() {
  29. testx.InitEnv()
  30. }
  31. func TestStreamCreateProcessor(t *testing.T) {
  32. tests := []struct {
  33. s string
  34. r []string
  35. err string
  36. }{
  37. {
  38. s: `SHOW STREAMS;`,
  39. r: []string{"No stream definitions are found."},
  40. },
  41. {
  42. s: `EXPLAIN STREAM topic1;`,
  43. err: "Explain stream fails, topic1 is not found.",
  44. },
  45. {
  46. s: `CREATE STREAM topic1 (
  47. USERID BIGINT,
  48. FIRST_NAME STRING,
  49. LAST_NAME STRING,
  50. NICKNAMES ARRAY(STRING),
  51. Gender BOOLEAN,
  52. ADDRESS STRUCT(STREET_NAME STRING, NUMBER BIGINT, BUILDING STRUCT(NAME STRING, ROOM BIGINT)),
  53. ) WITH (DATASOURCE="users", FORMAT="JSON", KEY="USERID");`,
  54. r: []string{"Stream topic1 is created."},
  55. },
  56. {
  57. s: `CREATE STREAM ` + "`stream`" + ` (
  58. USERID BIGINT,
  59. FIRST_NAME STRING,
  60. LAST_NAME STRING,
  61. NICKNAMES ARRAY(STRING),
  62. Gender BOOLEAN,
  63. ` + "`地址`" + ` STRUCT(STREET_NAME STRING, NUMBER BIGINT),
  64. ) WITH (DATASOURCE="users", FORMAT="JSON", KEY="USERID");`,
  65. r: []string{"Stream stream is created."},
  66. },
  67. {
  68. s: `CREATE STREAM topic1 (
  69. USERID BIGINT,
  70. ) WITH (DATASOURCE="users", FORMAT="JSON", KEY="USERID");`,
  71. err: "Create stream fails: Item topic1 already exists.",
  72. },
  73. {
  74. s: `EXPLAIN STREAM topic1;`,
  75. r: []string{"TO BE SUPPORTED"},
  76. },
  77. {
  78. s: `DESCRIBE STREAM topic1;`,
  79. r: []string{"Fields\n--------------------------------------------------------------------------------\nUSERID\tbigint\nFIRST_NAME\tstring\nLAST_NAME\tstring\nNICKNAMES\t" +
  80. "array(string)\nGender\tboolean\nADDRESS\tstruct(STREET_NAME string, NUMBER bigint, BUILDING struct(NAME string, ROOM bigint))\n\n" +
  81. "DATASOURCE: users\nFORMAT: JSON\nKEY: USERID\n"},
  82. },
  83. {
  84. s: `DROP STREAM topic1;`,
  85. r: []string{"Stream topic1 is dropped."},
  86. },
  87. {
  88. s: `SHOW STREAMS;`,
  89. r: []string{"stream"},
  90. },
  91. {
  92. s: `DESCRIBE STREAM topic1;`,
  93. err: "Describe stream fails, topic1 is not found.",
  94. },
  95. {
  96. s: `DROP STREAM topic1;`,
  97. err: "Drop stream fails: topic1 is not found.",
  98. },
  99. {
  100. s: "DROP STREAM `stream`;",
  101. r: []string{"Stream stream is dropped."},
  102. },
  103. }
  104. fmt.Printf("The test bucket size is %d.\n\n", len(tests))
  105. p := NewStreamProcessor()
  106. for i, tt := range tests {
  107. results, err := p.ExecStmt(tt.s)
  108. if !reflect.DeepEqual(tt.err, testx.Errstring(err)) {
  109. t.Errorf("%d. %q: error mismatch:\n exp=%s\n got=%s\n\n", i, tt.s, tt.err, err)
  110. } else if tt.err == "" {
  111. if !reflect.DeepEqual(tt.r, results) {
  112. t.Errorf("%d. %q\n\nstmt mismatch:\nexp=%s\ngot=%#v\n\n", i, tt.s, tt.r, results)
  113. }
  114. }
  115. }
  116. }
  117. func TestTableProcessor(t *testing.T) {
  118. tests := []struct {
  119. s string
  120. r []string
  121. err string
  122. }{
  123. {
  124. s: `SHOW TABLES;`,
  125. r: []string{"No table definitions are found."},
  126. },
  127. {
  128. s: `EXPLAIN TABLE topic1;`,
  129. err: "Explain table fails, topic1 is not found.",
  130. },
  131. {
  132. s: `CREATE TABLE topic1 (
  133. USERID BIGINT,
  134. FIRST_NAME STRING,
  135. LAST_NAME STRING,
  136. NICKNAMES ARRAY(STRING),
  137. Gender BOOLEAN,
  138. ADDRESS STRUCT(STREET_NAME STRING, NUMBER BIGINT),
  139. ) WITH (DATASOURCE="users", FORMAT="JSON", KEY="USERID");`,
  140. r: []string{"Table topic1 is created."},
  141. },
  142. {
  143. s: `CREATE TABLE ` + "`stream`" + ` (
  144. USERID BIGINT,
  145. FIRST_NAME STRING,
  146. LAST_NAME STRING,
  147. NICKNAMES ARRAY(STRING),
  148. Gender BOOLEAN,
  149. ` + "`地址`" + ` STRUCT(STREET_NAME STRING, NUMBER BIGINT),
  150. ) WITH (DATASOURCE="users", FORMAT="JSON", KEY="USERID");`,
  151. r: []string{"Table stream is created."},
  152. },
  153. {
  154. s: `CREATE TABLE topic1 (
  155. USERID BIGINT,
  156. ) WITH (DATASOURCE="users", FORMAT="JSON", KEY="USERID");`,
  157. err: "Create table fails: Item topic1 already exists.",
  158. },
  159. {
  160. s: `EXPLAIN TABLE topic1;`,
  161. r: []string{"TO BE SUPPORTED"},
  162. },
  163. {
  164. s: `DESCRIBE TABLE topic1;`,
  165. r: []string{"Fields\n--------------------------------------------------------------------------------\nUSERID\tbigint\nFIRST_NAME\tstring\nLAST_NAME\tstring\nNICKNAMES\t" +
  166. "array(string)\nGender\tboolean\nADDRESS\tstruct(STREET_NAME string, NUMBER bigint)\n\n" +
  167. "DATASOURCE: users\nFORMAT: JSON\nKEY: USERID\n"},
  168. },
  169. {
  170. s: `DROP TABLE topic1;`,
  171. r: []string{"Table topic1 is dropped."},
  172. },
  173. {
  174. s: `SHOW TABLES;`,
  175. r: []string{"stream"},
  176. },
  177. {
  178. s: `DESCRIBE TABLE topic1;`,
  179. err: "Describe table fails, topic1 is not found.",
  180. },
  181. {
  182. s: `DROP TABLE topic1;`,
  183. err: "Drop table fails: topic1 is not found.",
  184. },
  185. {
  186. s: "DROP TABLE `stream`;",
  187. r: []string{"Table stream is dropped."},
  188. },
  189. }
  190. fmt.Printf("The test bucket size is %d.\n\n", len(tests))
  191. p := NewStreamProcessor()
  192. for i, tt := range tests {
  193. results, err := p.ExecStmt(tt.s)
  194. if !reflect.DeepEqual(tt.err, testx.Errstring(err)) {
  195. t.Errorf("%d. %q: error mismatch:\n exp=%s\n got=%s\n\n", i, tt.s, tt.err, err)
  196. } else if tt.err == "" {
  197. if !reflect.DeepEqual(tt.r, results) {
  198. t.Errorf("%d. %q\n\nstmt mismatch:\nexp=%s\ngot=%#v\n\n", i, tt.s, tt.r, results)
  199. }
  200. }
  201. }
  202. }
  203. func TestTableList(t *testing.T) {
  204. p := NewStreamProcessor()
  205. p.ExecStmt(`CREATE TABLE tt1 () WITH (DATASOURCE="users", FORMAT="JSON", KIND="scan")`)
  206. p.ExecStmt(`CREATE TABLE tt2 () WITH (DATASOURCE="users", TYPE="memory", FORMAT="JSON", KEY="id", KIND="lookup")`)
  207. p.ExecStmt(`CREATE TABLE tt3 () WITH (DATASOURCE="users", TYPE="memory", FORMAT="JSON", KEY="id", KIND="lookup")`)
  208. p.ExecStmt(`CREATE TABLE tt4 () WITH (DATASOURCE="users", FORMAT="JSON")`)
  209. defer func() {
  210. p.ExecStmt(`DROP TABLE tt1`)
  211. p.ExecStmt(`DROP TABLE tt2`)
  212. p.ExecStmt(`DROP TABLE tt3`)
  213. p.ExecStmt(`DROP TABLE tt4`)
  214. }()
  215. la, err := p.ShowTable("lookup")
  216. if err != nil {
  217. t.Errorf("Show lookup table fails: %s", err)
  218. return
  219. }
  220. le := []string{"tt2", "tt3"}
  221. if !reflect.DeepEqual(le, la) {
  222. t.Errorf("Show lookup table mismatch:\nexp=%s\ngot=%s", le, la)
  223. return
  224. }
  225. ls, err := p.ShowTable("scan")
  226. if err != nil {
  227. t.Errorf("Show scan table fails: %s", err)
  228. return
  229. }
  230. lse := []string{"tt1", "tt4"}
  231. if !reflect.DeepEqual(lse, ls) {
  232. t.Errorf("Show scan table mismatch:\nexp=%s\ngot=%s", lse, ls)
  233. return
  234. }
  235. }
  236. func TestAll(t *testing.T) {
  237. expected := map[string]map[string]string{
  238. "streams": {
  239. "demo": "create stream demo () WITH (FORMAT=\"JSON\", DATASOURCE=\"demo\", SHARED=\"TRUE\")",
  240. "demo1": "create stream demo1 () WITH (FORMAT=\"JSON\", DATASOURCE=\"demo\")",
  241. "demo2": "create stream demo2 () WITH (FORMAT=\"JSON\", DATASOURCE=\"demo\", SHARED=\"TRUE\")",
  242. "demo3": "create stream demo3 () WITH (FORMAT=\"JSON\", DATASOURCE=\"demo\", SHARED=\"TRUE\")",
  243. },
  244. "tables": {
  245. "tt1": `CREATE TABLE tt1 () WITH (DATASOURCE="users", FORMAT="JSON", KIND="scan")`,
  246. "tt3": `CREATE TABLE tt3 () WITH (DATASOURCE="users", TYPE="memory", FORMAT="JSON", KEY="id", KIND="lookup")`,
  247. },
  248. }
  249. p := NewStreamProcessor()
  250. p.db.Clean()
  251. defer p.db.Clean()
  252. for st, m := range expected {
  253. for k, v := range m {
  254. p.ExecStmt(v)
  255. defer p.ExecStmt("Drop " + st + " " + k)
  256. }
  257. }
  258. all, err := p.GetAll()
  259. if err != nil {
  260. t.Error(err)
  261. return
  262. }
  263. if !reflect.DeepEqual(all, expected) {
  264. t.Errorf("Expect\t %v\nBut got\t%v", expected, all)
  265. }
  266. }
  267. func TestInferredStream(t *testing.T) {
  268. // init schema
  269. // Prepare test schema file
  270. conf.IsTesting = false
  271. dataDir, err := conf.GetDataLoc()
  272. if err != nil {
  273. t.Fatal(err)
  274. }
  275. etcDir := filepath.Join(dataDir, "schemas", "custom")
  276. err = os.MkdirAll(etcDir, os.ModePerm)
  277. if err != nil {
  278. t.Fatal(err)
  279. }
  280. defer func() {
  281. err = os.RemoveAll(etcDir)
  282. if err != nil {
  283. t.Fatal(err)
  284. }
  285. }()
  286. // build the so file into data/test prior to running the test
  287. bytesRead, err := os.ReadFile(filepath.Join(dataDir, "myFormat.so"))
  288. if err != nil {
  289. t.Fatal(err)
  290. }
  291. err = os.WriteFile(filepath.Join(etcDir, "myFormat.so"), bytesRead, 0o755)
  292. if err != nil {
  293. t.Fatal(err)
  294. }
  295. petcDir := filepath.Join(dataDir, "schemas", "protobuf")
  296. err = os.MkdirAll(petcDir, os.ModePerm)
  297. if err != nil {
  298. t.Fatal(err)
  299. }
  300. // Copy test2.proto
  301. bytesRead, err = os.ReadFile("../schema/test/test2.proto")
  302. if err != nil {
  303. t.Fatal(err)
  304. }
  305. err = os.WriteFile(filepath.Join(petcDir, "test2.proto"), bytesRead, 0o755)
  306. if err != nil {
  307. t.Fatal(err)
  308. }
  309. defer func() {
  310. err = os.RemoveAll(petcDir)
  311. if err != nil {
  312. t.Fatal(err)
  313. }
  314. }()
  315. schema.InitRegistry()
  316. tests := []struct {
  317. s string
  318. r map[string]*ast.JsonStreamField
  319. err string
  320. }{
  321. {
  322. s: `CREATE STREAM demo0 (USERID bigint, NAME string) WITH (FORMAT="JSON", DATASOURCE="demo", SHARED="TRUE")`,
  323. r: map[string]*ast.JsonStreamField{
  324. "USERID": {Type: "bigint"},
  325. "NAME": {Type: "string"},
  326. },
  327. }, {
  328. s: `CREATE STREAM demo1 (USERID bigint, NAME string) WITH (FORMAT="protobuf", DATASOURCE="demo", SCHEMAID="test2.Book")`,
  329. r: map[string]*ast.JsonStreamField{
  330. "name": {Type: "string"},
  331. "author": {Type: "string"},
  332. },
  333. }, {
  334. s: `CREATE STREAM demo2 () WITH (FORMAT="custom", DATASOURCE="demo", SCHEMAID="myFormat.Sample")`,
  335. r: map[string]*ast.JsonStreamField{
  336. "id": {Type: "bigint"},
  337. "name": {Type: "string"},
  338. "age": {Type: "bigint"},
  339. "hobbies": {
  340. Type: "struct",
  341. Properties: map[string]*ast.JsonStreamField{
  342. "indoor": {Type: "array", Items: &ast.JsonStreamField{Type: "string"}},
  343. "outdoor": {Type: "array", Items: &ast.JsonStreamField{Type: "string"}},
  344. },
  345. },
  346. },
  347. },
  348. }
  349. p := NewStreamProcessor()
  350. p.db.Clean()
  351. defer p.db.Clean()
  352. for i, tt := range tests {
  353. _, err := p.ExecStmt(tt.s)
  354. if err != nil {
  355. t.Errorf("%d. ExecStmt(%q) error: %v", i, tt.s, err)
  356. continue
  357. }
  358. sf, err := p.GetInferredJsonSchema("demo"+strconv.Itoa(i), ast.TypeStream)
  359. if err != nil {
  360. t.Errorf("GetInferredJsonSchema fails: %s", err)
  361. continue
  362. }
  363. if !reflect.DeepEqual(sf, tt.r) {
  364. t.Errorf("GetInferredJsonSchema mismatch:\nexp=%v\ngot=%v", render.AsCode(tt.r), render.AsCode(sf))
  365. }
  366. }
  367. }