stream.go 14 KB

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