stream.go 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489
  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. "bytes"
  17. "encoding/json"
  18. "fmt"
  19. "strings"
  20. "golang.org/x/text/cases"
  21. "golang.org/x/text/language"
  22. "github.com/lf-edge/ekuiper/internal/conf"
  23. "github.com/lf-edge/ekuiper/internal/pkg/store"
  24. "github.com/lf-edge/ekuiper/internal/schema"
  25. "github.com/lf-edge/ekuiper/internal/topo/lookup"
  26. "github.com/lf-edge/ekuiper/internal/xsql"
  27. "github.com/lf-edge/ekuiper/pkg/ast"
  28. "github.com/lf-edge/ekuiper/pkg/errorx"
  29. "github.com/lf-edge/ekuiper/pkg/kv"
  30. )
  31. var (
  32. log = conf.Log
  33. )
  34. type StreamProcessor struct {
  35. db kv.KeyValue
  36. streamStatusDb kv.KeyValue
  37. tableStatusDb kv.KeyValue
  38. }
  39. func NewStreamProcessor() *StreamProcessor {
  40. db, err := store.GetKV("stream")
  41. if err != nil {
  42. panic(fmt.Sprintf("Can not initialize store for the stream processor at path 'stream': %v", err))
  43. }
  44. streamDb, err := store.GetKV("streamStatus")
  45. if err != nil {
  46. panic(fmt.Sprintf("Can not initialize store for the stream processor at path 'stream': %v", err))
  47. }
  48. tableDb, err := store.GetKV("tableStatus")
  49. if err != nil {
  50. panic(fmt.Sprintf("Can not initialize store for the stream processor at path 'stream': %v", err))
  51. }
  52. processor := &StreamProcessor{
  53. db: db,
  54. streamStatusDb: streamDb,
  55. tableStatusDb: tableDb,
  56. }
  57. return processor
  58. }
  59. func (p *StreamProcessor) ExecStmt(statement string) (result []string, err error) {
  60. parser := xsql.NewParser(strings.NewReader(statement))
  61. stmt, err := xsql.Language.Parse(parser)
  62. if err != nil {
  63. return nil, err
  64. }
  65. switch s := stmt.(type) {
  66. case *ast.StreamStmt: //Table is also StreamStmt
  67. var r string
  68. err = p.execSave(s, statement, false)
  69. stt := ast.StreamTypeMap[s.StreamType]
  70. if err != nil {
  71. err = fmt.Errorf("Create %s fails: %v.", stt, err)
  72. } else {
  73. r = fmt.Sprintf("%s %s is created.", cases.Title(language.Und).String(stt), s.Name)
  74. log.Printf("%s", r)
  75. }
  76. result = append(result, r)
  77. case *ast.ShowStreamsStatement:
  78. result, err = p.execShow(ast.TypeStream)
  79. case *ast.ShowTablesStatement:
  80. result, err = p.execShow(ast.TypeTable)
  81. case *ast.DescribeStreamStatement:
  82. var r string
  83. r, err = p.execDescribe(s, ast.TypeStream)
  84. result = append(result, r)
  85. case *ast.DescribeTableStatement:
  86. var r string
  87. r, err = p.execDescribe(s, ast.TypeTable)
  88. result = append(result, r)
  89. case *ast.ExplainStreamStatement:
  90. var r string
  91. r, err = p.execExplain(s, ast.TypeStream)
  92. result = append(result, r)
  93. case *ast.ExplainTableStatement:
  94. var r string
  95. r, err = p.execExplain(s, ast.TypeTable)
  96. result = append(result, r)
  97. case *ast.DropStreamStatement:
  98. var r string
  99. r, err = p.execDrop(s, ast.TypeStream)
  100. result = append(result, r)
  101. case *ast.DropTableStatement:
  102. var r string
  103. r, err = p.execDrop(s, ast.TypeTable)
  104. result = append(result, r)
  105. default:
  106. return nil, fmt.Errorf("Invalid stream statement: %s", statement)
  107. }
  108. return
  109. }
  110. func (p *StreamProcessor) RecoverLookupTable() error {
  111. keys, err := p.db.Keys()
  112. if err != nil {
  113. return fmt.Errorf("error loading data from db: %v.", err)
  114. }
  115. var (
  116. v string
  117. vs = &xsql.StreamInfo{}
  118. )
  119. for _, k := range keys {
  120. if ok, _ := p.db.Get(k, &v); ok {
  121. if err := json.Unmarshal([]byte(v), vs); err == nil && vs.StreamType == ast.TypeTable {
  122. parser := xsql.NewParser(strings.NewReader(vs.Statement))
  123. stmt, e := xsql.Language.Parse(parser)
  124. if e != nil {
  125. log.Error(e)
  126. }
  127. switch s := stmt.(type) {
  128. case *ast.StreamStmt:
  129. log.Infof("Starting lookup table %s", s.Name)
  130. e = lookup.CreateInstance(string(s.Name), s.Options.TYPE, s.Options)
  131. if e != nil {
  132. log.Errorf("%s", e.Error())
  133. return e
  134. }
  135. default:
  136. log.Errorf("Invalid lookup table statement: %s", vs.Statement)
  137. }
  138. }
  139. }
  140. }
  141. return nil
  142. }
  143. func (p *StreamProcessor) execSave(stmt *ast.StreamStmt, statement string, replace bool) error {
  144. if stmt.StreamType == ast.TypeTable && stmt.Options.KIND == ast.StreamKindLookup {
  145. log.Infof("Creating lookup table %s", stmt.Name)
  146. err := lookup.CreateInstance(string(stmt.Name), stmt.Options.TYPE, stmt.Options)
  147. if err != nil {
  148. return err
  149. }
  150. }
  151. s, err := json.Marshal(xsql.StreamInfo{
  152. StreamType: stmt.StreamType,
  153. Statement: statement,
  154. StreamKind: stmt.Options.KIND,
  155. })
  156. if err != nil {
  157. return fmt.Errorf("error when saving to db: %v.", err)
  158. }
  159. if replace {
  160. err = p.db.Set(string(stmt.Name), string(s))
  161. } else {
  162. err = p.db.Setnx(string(stmt.Name), string(s))
  163. }
  164. return err
  165. }
  166. func (p *StreamProcessor) ExecReplaceStream(name string, statement string, st ast.StreamType) (string, error) {
  167. parser := xsql.NewParser(strings.NewReader(statement))
  168. stmt, err := xsql.Language.Parse(parser)
  169. if err != nil {
  170. return "", err
  171. }
  172. stt := ast.StreamTypeMap[st]
  173. switch s := stmt.(type) {
  174. case *ast.StreamStmt:
  175. if s.StreamType != st {
  176. return "", errorx.NewWithCode(errorx.NOT_FOUND, fmt.Sprintf("%s %s is not found", ast.StreamTypeMap[st], s.Name))
  177. }
  178. if string(s.Name) != name {
  179. return "", fmt.Errorf("Replace %s fails: the sql statement must update the %s source.", name, name)
  180. }
  181. err = p.execSave(s, statement, true)
  182. if err != nil {
  183. return "", fmt.Errorf("Replace %s fails: %v.", stt, err)
  184. } else {
  185. info := fmt.Sprintf("%s %s is replaced.", cases.Title(language.Und).String(stt), s.Name)
  186. log.Printf("%s", info)
  187. return info, nil
  188. }
  189. default:
  190. return "", fmt.Errorf("Invalid %s statement: %s", stt, statement)
  191. }
  192. }
  193. func (p *StreamProcessor) ExecStreamSql(statement string) (string, error) {
  194. r, err := p.ExecStmt(statement)
  195. if err != nil {
  196. return "", err
  197. } else {
  198. return strings.Join(r, "\n"), err
  199. }
  200. }
  201. func (p *StreamProcessor) execShow(st ast.StreamType) ([]string, error) {
  202. keys, err := p.ShowStream(st)
  203. if len(keys) == 0 {
  204. keys = append(keys, fmt.Sprintf("No %s definitions are found.", ast.StreamTypeMap[st]))
  205. }
  206. return keys, err
  207. }
  208. func (p *StreamProcessor) ShowStream(st ast.StreamType) ([]string, error) {
  209. stt := ast.StreamTypeMap[st]
  210. keys, err := p.db.Keys()
  211. if err != nil {
  212. return nil, fmt.Errorf("Show %ss fails, error when loading data from db: %v.", stt, err)
  213. }
  214. var (
  215. v string
  216. vs = &xsql.StreamInfo{}
  217. result = make([]string, 0)
  218. )
  219. for _, k := range keys {
  220. if ok, _ := p.db.Get(k, &v); ok {
  221. if err := json.Unmarshal([]byte(v), vs); err == nil && vs.StreamType == st {
  222. result = append(result, k)
  223. }
  224. }
  225. }
  226. return result, nil
  227. }
  228. func (p *StreamProcessor) ShowTable(kind string) ([]string, error) {
  229. if kind == "" {
  230. return p.ShowStream(ast.TypeTable)
  231. }
  232. keys, err := p.db.Keys()
  233. if err != nil {
  234. return nil, fmt.Errorf("Show tables fails, error when loading data from db: %v.", err)
  235. }
  236. var (
  237. v string
  238. vs = &xsql.StreamInfo{}
  239. result = make([]string, 0)
  240. )
  241. for _, k := range keys {
  242. if ok, _ := p.db.Get(k, &v); ok {
  243. if err := json.Unmarshal([]byte(v), vs); err == nil && vs.StreamType == ast.TypeTable {
  244. if kind == "scan" && (vs.StreamKind == ast.StreamKindScan || vs.StreamKind == "") {
  245. result = append(result, k)
  246. } else if kind == "lookup" && vs.StreamKind == ast.StreamKindLookup {
  247. result = append(result, k)
  248. }
  249. }
  250. }
  251. }
  252. return result, nil
  253. }
  254. func (p *StreamProcessor) GetStream(name string, st ast.StreamType) (string, error) {
  255. vs, err := xsql.GetDataSourceStatement(p.db, name)
  256. if vs != nil && vs.StreamType == st {
  257. return vs.Statement, nil
  258. }
  259. if err != nil {
  260. return "", err
  261. }
  262. return "", errorx.NewWithCode(errorx.NOT_FOUND, fmt.Sprintf("%s %s is not found", ast.StreamTypeMap[st], name))
  263. }
  264. func (p *StreamProcessor) execDescribe(stmt ast.NameNode, st ast.StreamType) (string, error) {
  265. streamStmt, err := p.DescStream(stmt.GetName(), st)
  266. if err != nil {
  267. return "", err
  268. }
  269. switch s := streamStmt.(type) {
  270. case *ast.StreamStmt:
  271. var buff bytes.Buffer
  272. buff.WriteString("Fields\n--------------------------------------------------------------------------------\n")
  273. for _, f := range s.StreamFields {
  274. buff.WriteString(f.Name + "\t")
  275. buff.WriteString(printFieldType(f.FieldType))
  276. buff.WriteString("\n")
  277. }
  278. buff.WriteString("\n")
  279. printOptions(s.Options, &buff)
  280. return buff.String(), err
  281. default:
  282. return "%s", fmt.Errorf("Error resolving the %s %s, the data in db may be corrupted.", ast.StreamTypeMap[st], stmt.GetName())
  283. }
  284. }
  285. func printOptions(opts *ast.Options, buff *bytes.Buffer) {
  286. if opts.CONF_KEY != "" {
  287. buff.WriteString(fmt.Sprintf("CONF_KEY: %s\n", opts.CONF_KEY))
  288. }
  289. if opts.DATASOURCE != "" {
  290. buff.WriteString(fmt.Sprintf("DATASOURCE: %s\n", opts.DATASOURCE))
  291. }
  292. if opts.FORMAT != "" {
  293. buff.WriteString(fmt.Sprintf("FORMAT: %s\n", opts.FORMAT))
  294. }
  295. if opts.SCHEMAID != "" {
  296. buff.WriteString(fmt.Sprintf("SCHEMAID: %s\n", opts.SCHEMAID))
  297. }
  298. if opts.KEY != "" {
  299. buff.WriteString(fmt.Sprintf("KEY: %s\n", opts.KEY))
  300. }
  301. if opts.RETAIN_SIZE != 0 {
  302. buff.WriteString(fmt.Sprintf("RETAIN_SIZE: %d\n", opts.RETAIN_SIZE))
  303. }
  304. if opts.SHARED {
  305. buff.WriteString(fmt.Sprintf("SHARED: %v\n", opts.SHARED))
  306. }
  307. if opts.STRICT_VALIDATION {
  308. buff.WriteString(fmt.Sprintf("STRICT_VALIDATION: %v\n", opts.STRICT_VALIDATION))
  309. }
  310. if opts.TIMESTAMP != "" {
  311. buff.WriteString(fmt.Sprintf("TIMESTAMP: %s\n", opts.TIMESTAMP))
  312. }
  313. if opts.TIMESTAMP_FORMAT != "" {
  314. buff.WriteString(fmt.Sprintf("TIMESTAMP_FORMAT: %s\n", opts.TIMESTAMP_FORMAT))
  315. }
  316. if opts.TYPE != "" {
  317. buff.WriteString(fmt.Sprintf("TYPE: %s\n", opts.TYPE))
  318. }
  319. }
  320. func (p *StreamProcessor) DescStream(name string, st ast.StreamType) (ast.Statement, error) {
  321. statement, err := p.GetStream(name, st)
  322. if err != nil {
  323. return nil, fmt.Errorf("Describe %s fails, %s.", ast.StreamTypeMap[st], err)
  324. }
  325. parser := xsql.NewParser(strings.NewReader(statement))
  326. stream, err := xsql.Language.Parse(parser)
  327. if err != nil {
  328. return nil, err
  329. }
  330. return stream, nil
  331. }
  332. func (p *StreamProcessor) GetInferredSchema(name string, st ast.StreamType) (ast.StreamFields, error) {
  333. statement, err := p.GetStream(name, st)
  334. if err != nil {
  335. return nil, fmt.Errorf("Describe %s fails, %s.", ast.StreamTypeMap[st], err)
  336. }
  337. parser := xsql.NewParser(strings.NewReader(statement))
  338. stream, err := xsql.Language.Parse(parser)
  339. if err != nil {
  340. return nil, err
  341. }
  342. stmt, ok := stream.(*ast.StreamStmt)
  343. if !ok {
  344. return nil, fmt.Errorf("Describe %s fails, cannot parse the data \"%s\" to a stream statement", ast.StreamTypeMap[st], statement)
  345. }
  346. if stmt.Options.SCHEMAID != "" {
  347. return schema.InferFromSchemaFile(stmt.Options.FORMAT, stmt.Options.SCHEMAID)
  348. }
  349. return nil, nil
  350. }
  351. // GetInferredJsonSchema return schema in json schema type
  352. func (p *StreamProcessor) GetInferredJsonSchema(name string, st ast.StreamType) (map[string]*ast.JsonStreamField, error) {
  353. statement, err := p.GetStream(name, st)
  354. if err != nil {
  355. return nil, fmt.Errorf("Describe %s fails, %s.", ast.StreamTypeMap[st], err)
  356. }
  357. parser := xsql.NewParser(strings.NewReader(statement))
  358. stream, err := xsql.Language.Parse(parser)
  359. if err != nil {
  360. return nil, err
  361. }
  362. stmt, ok := stream.(*ast.StreamStmt)
  363. if !ok {
  364. return nil, fmt.Errorf("Describe %s fails, cannot parse the data \"%s\" to a stream statement", ast.StreamTypeMap[st], statement)
  365. }
  366. sfs := stmt.StreamFields
  367. if stmt.Options.SCHEMAID != "" {
  368. sfs, err = schema.InferFromSchemaFile(stmt.Options.FORMAT, stmt.Options.SCHEMAID)
  369. if err != nil {
  370. return nil, err
  371. }
  372. }
  373. return sfs.ToJsonSchema(), nil
  374. }
  375. func (p *StreamProcessor) execExplain(stmt ast.NameNode, st ast.StreamType) (string, error) {
  376. _, err := p.GetStream(stmt.GetName(), st)
  377. if err != nil {
  378. return "", fmt.Errorf("Explain %s fails, %s.", ast.StreamTypeMap[st], err)
  379. }
  380. return "TO BE SUPPORTED", nil
  381. }
  382. func (p *StreamProcessor) execDrop(stmt ast.NameNode, st ast.StreamType) (string, error) {
  383. s, err := p.DropStream(stmt.GetName(), st)
  384. if err != nil {
  385. return s, fmt.Errorf("Drop %s fails: %s.", ast.StreamTypeMap[st], err)
  386. }
  387. return s, nil
  388. }
  389. func (p *StreamProcessor) DropStream(name string, st ast.StreamType) (string, error) {
  390. if st == ast.TypeTable {
  391. err := lookup.DropInstance(name)
  392. if err != nil {
  393. return "", err
  394. }
  395. }
  396. _, err := p.GetStream(name, st)
  397. if err != nil {
  398. return "", err
  399. }
  400. err = p.db.Delete(name)
  401. if err != nil {
  402. return "", err
  403. } else {
  404. return fmt.Sprintf("%s %s is dropped.", cases.Title(language.Und).String(ast.StreamTypeMap[st]), name), nil
  405. }
  406. }
  407. func printFieldType(ft ast.FieldType) (result string) {
  408. switch t := ft.(type) {
  409. case *ast.BasicType:
  410. result = t.Type.String()
  411. case *ast.ArrayType:
  412. result = "array("
  413. if t.FieldType != nil {
  414. result += printFieldType(t.FieldType)
  415. } else {
  416. result += t.Type.String()
  417. }
  418. result += ")"
  419. case *ast.RecType:
  420. result = "struct("
  421. isFirst := true
  422. for _, f := range t.StreamFields {
  423. if isFirst {
  424. isFirst = false
  425. } else {
  426. result += ", "
  427. }
  428. result = result + f.Name + " " + printFieldType(f.FieldType)
  429. }
  430. result += ")"
  431. }
  432. return
  433. }
  434. // GetAll return all streams and tables defined to export.
  435. func (p *StreamProcessor) GetAll() (result map[string]map[string]string, e error) {
  436. defs, err := p.db.All()
  437. if err != nil {
  438. e = err
  439. return
  440. }
  441. var (
  442. vs = &xsql.StreamInfo{}
  443. )
  444. result = map[string]map[string]string{
  445. "streams": make(map[string]string),
  446. "tables": make(map[string]string),
  447. }
  448. for k, v := range defs {
  449. if err := json.Unmarshal([]byte(v), vs); err == nil {
  450. switch vs.StreamType {
  451. case ast.TypeStream:
  452. result["streams"][k] = vs.Statement
  453. case ast.TypeTable:
  454. result["tables"][k] = vs.Statement
  455. }
  456. }
  457. }
  458. return
  459. }