ruleState_test.go 9.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465
  1. // Copyright 2022-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 rule
  15. import (
  16. "errors"
  17. "fmt"
  18. "reflect"
  19. "sync"
  20. "testing"
  21. "time"
  22. "github.com/lf-edge/ekuiper/internal/conf"
  23. "github.com/lf-edge/ekuiper/internal/processor"
  24. "github.com/lf-edge/ekuiper/internal/testx"
  25. "github.com/lf-edge/ekuiper/pkg/api"
  26. )
  27. var defaultOption = &api.RuleOption{
  28. IsEventTime: false,
  29. LateTol: 1000,
  30. Concurrency: 1,
  31. BufferLength: 1024,
  32. SendMetaToSink: false,
  33. SendError: true,
  34. Qos: api.AtMostOnce,
  35. CheckpointInterval: 300000,
  36. Restart: &api.RestartStrategy{
  37. Attempts: 0,
  38. Delay: 1000,
  39. Multiplier: 2,
  40. MaxDelay: 30000,
  41. JitterFactor: 0.1,
  42. },
  43. }
  44. func init() {
  45. testx.InitEnv()
  46. }
  47. func TestCreate(t *testing.T) {
  48. sp := processor.NewStreamProcessor()
  49. sp.ExecStmt(`CREATE STREAM demo () WITH (DATASOURCE="users", FORMAT="JSON")`)
  50. defer sp.ExecStmt(`DROP STREAM demo`)
  51. tests := []struct {
  52. r *api.Rule
  53. e error
  54. }{
  55. {
  56. r: &api.Rule{
  57. Triggered: false,
  58. Id: "test",
  59. Sql: "SELECT ts FROM demo",
  60. Actions: []map[string]interface{}{
  61. {
  62. "log": map[string]interface{}{},
  63. },
  64. },
  65. Options: defaultOption,
  66. },
  67. e: nil,
  68. },
  69. {
  70. r: &api.Rule{
  71. Triggered: false,
  72. Id: "test",
  73. Sql: "SELECT FROM demo",
  74. Actions: []map[string]interface{}{
  75. {
  76. "log": map[string]interface{}{},
  77. },
  78. },
  79. Options: defaultOption,
  80. },
  81. e: errors.New("Parse SQL SELECT FROM demo error: found \"FROM\", expected expression.."),
  82. },
  83. {
  84. r: &api.Rule{
  85. Triggered: false,
  86. Id: "test",
  87. Sql: "SELECT * FROM demo1",
  88. Actions: []map[string]interface{}{
  89. {
  90. "log": map[string]interface{}{},
  91. },
  92. },
  93. Options: defaultOption,
  94. },
  95. e: errors.New("fail to get stream demo1, please check if stream is created"),
  96. },
  97. }
  98. for i, tt := range tests {
  99. _, err := NewRuleState(tt.r)
  100. if !reflect.DeepEqual(err, tt.e) {
  101. t.Errorf("%d.\n\nerror mismatch:\n\nexp=%#v\n\ngot=%#v\n\n", i, tt.e, err)
  102. }
  103. }
  104. }
  105. func TestUpdate(t *testing.T) {
  106. sp := processor.NewStreamProcessor()
  107. sp.ExecStmt(`CREATE STREAM demo () WITH (DATASOURCE="users", FORMAT="JSON")`)
  108. defer sp.ExecStmt(`DROP STREAM demo`)
  109. rs, err := NewRuleState(&api.Rule{
  110. Triggered: false,
  111. Id: "test",
  112. Sql: "SELECT ts FROM demo",
  113. Actions: []map[string]interface{}{
  114. {
  115. "log": map[string]interface{}{},
  116. },
  117. },
  118. Options: defaultOption,
  119. })
  120. if err != nil {
  121. t.Error(err)
  122. return
  123. }
  124. defer rs.Close()
  125. err = rs.Start()
  126. if err != nil {
  127. t.Error(err)
  128. return
  129. }
  130. tests := []struct {
  131. r *api.Rule
  132. e error
  133. }{
  134. {
  135. r: &api.Rule{
  136. Triggered: false,
  137. Id: "test",
  138. Sql: "SELECT FROM demo",
  139. Actions: []map[string]interface{}{
  140. {
  141. "log": map[string]interface{}{},
  142. },
  143. },
  144. Options: defaultOption,
  145. },
  146. e: errors.New("Parse SQL SELECT FROM demo error: found \"FROM\", expected expression.."),
  147. },
  148. {
  149. r: &api.Rule{
  150. Triggered: false,
  151. Id: "test",
  152. Sql: "SELECT * FROM demo1",
  153. Actions: []map[string]interface{}{
  154. {
  155. "log": map[string]interface{}{},
  156. },
  157. },
  158. Options: defaultOption,
  159. },
  160. e: errors.New("fail to get stream demo1, please check if stream is created"),
  161. },
  162. {
  163. r: &api.Rule{
  164. Triggered: false,
  165. Id: "test",
  166. Sql: "SELECT * FROM demo",
  167. Actions: []map[string]interface{}{
  168. {
  169. "log": map[string]interface{}{},
  170. },
  171. },
  172. Options: defaultOption,
  173. },
  174. e: nil,
  175. },
  176. }
  177. for i, tt := range tests {
  178. err = rs.UpdateTopo(tt.r)
  179. if !reflect.DeepEqual(err, tt.e) {
  180. t.Errorf("%d.\n\nerror mismatch:\n\nexp=%#v\n\ngot=%#v\n\n", i, tt.e, err)
  181. }
  182. }
  183. }
  184. func TestMultipleAccess(t *testing.T) {
  185. sp := processor.NewStreamProcessor()
  186. sp.ExecStmt(`CREATE STREAM demo () WITH (DATASOURCE="users", FORMAT="JSON")`)
  187. defer sp.ExecStmt(`DROP STREAM demo`)
  188. rs, err := NewRuleState(&api.Rule{
  189. Triggered: false,
  190. Id: "test",
  191. Sql: "SELECT ts FROM demo",
  192. Actions: []map[string]interface{}{
  193. {
  194. "log": map[string]interface{}{},
  195. },
  196. },
  197. Options: defaultOption,
  198. })
  199. if err != nil {
  200. t.Error(err)
  201. return
  202. }
  203. defer rs.Close()
  204. err = rs.Start()
  205. if err != nil {
  206. t.Error(err)
  207. return
  208. }
  209. var wg sync.WaitGroup
  210. wg.Add(10)
  211. for i := 0; i < 10; i++ {
  212. if i%3 == 0 {
  213. go func(i int) {
  214. rs.Stop()
  215. fmt.Printf("%d:%d\n", i, rs.triggered)
  216. wg.Done()
  217. }(i)
  218. } else {
  219. go func(i int) {
  220. rs.Start()
  221. fmt.Printf("%d:%d\n", i, rs.triggered)
  222. wg.Done()
  223. }(i)
  224. }
  225. }
  226. wg.Wait()
  227. rs.Start()
  228. fmt.Printf("%d:%d\n", 10, rs.triggered)
  229. if rs.triggered != 1 {
  230. t.Errorf("triggered mismatch:\n\nexp=%#v\n\ngot=%#v\n\n", 1, rs.triggered)
  231. }
  232. }
  233. // Test rule state message
  234. func TestRuleState_Start(t *testing.T) {
  235. sp := processor.NewStreamProcessor()
  236. sp.ExecStmt(`CREATE STREAM demo () WITH (TYPE="neuron", FORMAT="JSON")`)
  237. defer sp.ExecStmt(`DROP STREAM demo`)
  238. // Test rule not triggered
  239. r := &api.Rule{
  240. Triggered: false,
  241. Id: "test",
  242. Sql: "SELECT ts FROM demo",
  243. Actions: []map[string]interface{}{
  244. {
  245. "log": map[string]interface{}{},
  246. },
  247. },
  248. Options: defaultOption,
  249. }
  250. const ruleStopped = "Stopped: canceled manually."
  251. const ruleStarted = "Running"
  252. t.Run("test rule loaded but not started", func(t *testing.T) {
  253. rs, err := NewRuleState(r)
  254. if err != nil {
  255. t.Error(err)
  256. return
  257. }
  258. state, err := rs.GetState()
  259. if err != nil {
  260. t.Errorf("get rule state error: %v", err)
  261. return
  262. }
  263. if state != ruleStopped {
  264. t.Errorf("rule state mismatch: exp=%v, got=%v", ruleStopped, state)
  265. return
  266. }
  267. })
  268. t.Run("test rule started", func(t *testing.T) {
  269. rs, err := NewRuleState(r)
  270. if err != nil {
  271. t.Error(err)
  272. return
  273. }
  274. err = rs.Start()
  275. if err != nil {
  276. t.Error(err)
  277. return
  278. }
  279. time.Sleep(100 * time.Millisecond)
  280. state, err := rs.GetState()
  281. if err != nil {
  282. t.Errorf("get rule state error: %v", err)
  283. return
  284. }
  285. if state != ruleStarted {
  286. t.Errorf("rule state mismatch: exp=%v, got=%v", ruleStopped, state)
  287. return
  288. }
  289. })
  290. t.Run("test rule loaded and stopped", func(t *testing.T) {
  291. rs, err := NewRuleState(r)
  292. if err != nil {
  293. t.Error(err)
  294. return
  295. }
  296. err = rs.Start()
  297. if err != nil {
  298. t.Error(err)
  299. return
  300. }
  301. err = rs.Close()
  302. if err != nil {
  303. t.Error(err)
  304. return
  305. }
  306. state, err := rs.GetState()
  307. if err != nil {
  308. t.Errorf("get rule state error: %v", err)
  309. return
  310. }
  311. if state != ruleStopped {
  312. t.Errorf("rule state mismatch: exp=%v, got=%v", ruleStopped, state)
  313. return
  314. }
  315. })
  316. }
  317. func TestScheduleRule(t *testing.T) {
  318. conf.IsTesting = true
  319. sp := processor.NewStreamProcessor()
  320. sp.ExecStmt(`CREATE STREAM demo () WITH (TYPE="neuron", FORMAT="JSON")`)
  321. defer sp.ExecStmt(`DROP STREAM demo`)
  322. // Test rule not triggered
  323. r := &api.Rule{
  324. Triggered: false,
  325. Id: "test",
  326. Sql: "SELECT ts FROM demo",
  327. Actions: []map[string]interface{}{
  328. {
  329. "log": map[string]interface{}{},
  330. },
  331. },
  332. Options: defaultOption,
  333. }
  334. r.Options.Cron = "mockCron"
  335. r.Options.Duration = "1s"
  336. const ruleStarted = "Running"
  337. const ruleStopped = "Stopped: waiting for next schedule."
  338. func() {
  339. rs, err := NewRuleState(r)
  340. if err != nil {
  341. t.Error(err)
  342. return
  343. }
  344. if err := rs.startScheduleRule(); err != nil {
  345. t.Error(err)
  346. return
  347. }
  348. time.Sleep(500 * time.Millisecond)
  349. state, err := rs.GetState()
  350. if err != nil {
  351. t.Errorf("get rule state error: %v", err)
  352. return
  353. }
  354. if state != ruleStarted {
  355. t.Errorf("rule state mismatch: exp=%v, got=%v", ruleStarted, state)
  356. return
  357. }
  358. if !rs.cronState.isInSchedule {
  359. t.Error("cron state should be in schedule")
  360. return
  361. }
  362. }()
  363. func() {
  364. rs, err := NewRuleState(r)
  365. if err != nil {
  366. t.Error(err)
  367. return
  368. }
  369. if err := rs.startScheduleRule(); err != nil {
  370. t.Error(err)
  371. return
  372. }
  373. time.Sleep(1500 * time.Millisecond)
  374. state, err := rs.GetState()
  375. if err != nil {
  376. t.Errorf("get rule state error: %v", err)
  377. return
  378. }
  379. if state != ruleStopped {
  380. t.Errorf("rule state mismatch: exp=%v, got=%v", ruleStopped, state)
  381. return
  382. }
  383. if !rs.cronState.isInSchedule {
  384. t.Error("cron state should be in schedule")
  385. return
  386. }
  387. }()
  388. func() {
  389. rs, err := NewRuleState(r)
  390. if err != nil {
  391. t.Error(err)
  392. return
  393. }
  394. if err := rs.startScheduleRule(); err != nil {
  395. t.Error(err)
  396. return
  397. }
  398. if err := rs.startScheduleRule(); err == nil {
  399. t.Error("rule can't be register in cron twice")
  400. return
  401. } else {
  402. if err.Error() != "rule test is already in schedule" {
  403. t.Error("error message wrong")
  404. return
  405. }
  406. }
  407. }()
  408. func() {
  409. rs, err := NewRuleState(r)
  410. if err != nil {
  411. t.Error(err)
  412. return
  413. }
  414. if err := rs.startScheduleRule(); err != nil {
  415. t.Error(err)
  416. return
  417. }
  418. if err := rs.Stop(); err != nil {
  419. t.Error(err)
  420. return
  421. }
  422. state, err := rs.GetState()
  423. if err != nil {
  424. t.Errorf("get rule state error: %v", err)
  425. return
  426. }
  427. if state != "Stopped: canceled manually." {
  428. t.Errorf("rule state mismatch: exp=%v, got=%v", "Stopped: canceled manually.", state)
  429. return
  430. }
  431. if rs.cronState.isInSchedule {
  432. t.Error("cron state shouldn't be in schedule")
  433. return
  434. }
  435. }()
  436. func() {
  437. rs, err := NewRuleState(r)
  438. if err != nil {
  439. t.Error(err)
  440. return
  441. }
  442. if err := rs.Stop(); err != nil {
  443. t.Error(err)
  444. return
  445. }
  446. if err := rs.Close(); err != nil {
  447. t.Error(err)
  448. return
  449. }
  450. }()
  451. }