stream.go 9.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329
  1. // Copyright 2021 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. "github.com/lf-edge/ekuiper/internal/conf"
  20. "github.com/lf-edge/ekuiper/internal/xsql"
  21. "github.com/lf-edge/ekuiper/pkg/ast"
  22. "github.com/lf-edge/ekuiper/pkg/errorx"
  23. "github.com/lf-edge/ekuiper/pkg/kv"
  24. "strings"
  25. )
  26. var (
  27. log = conf.Log
  28. )
  29. type StreamProcessor struct {
  30. db kv.KeyValue
  31. }
  32. //@params d : the directory of the DB to save the stream info
  33. func NewStreamProcessor(d string) *StreamProcessor {
  34. processor := &StreamProcessor{
  35. db: kv.GetDefaultKVStore(d),
  36. }
  37. return processor
  38. }
  39. func (p *StreamProcessor) ExecStmt(statement string) (result []string, err error) {
  40. parser := xsql.NewParser(strings.NewReader(statement))
  41. stmt, err := xsql.Language.Parse(parser)
  42. if err != nil {
  43. return nil, err
  44. }
  45. switch s := stmt.(type) {
  46. case *ast.StreamStmt: //Table is also StreamStmt
  47. var r string
  48. err = p.execSave(s, statement, false)
  49. stt := ast.StreamTypeMap[s.StreamType]
  50. if err != nil {
  51. err = fmt.Errorf("Create %s fails: %v.", stt, err)
  52. } else {
  53. r = fmt.Sprintf("%s %s is created.", strings.Title(stt), s.Name)
  54. log.Printf("%s", r)
  55. }
  56. result = append(result, r)
  57. case *ast.ShowStreamsStatement:
  58. result, err = p.execShow(ast.TypeStream)
  59. case *ast.ShowTablesStatement:
  60. result, err = p.execShow(ast.TypeTable)
  61. case *ast.DescribeStreamStatement:
  62. var r string
  63. r, err = p.execDescribe(s, ast.TypeStream)
  64. result = append(result, r)
  65. case *ast.DescribeTableStatement:
  66. var r string
  67. r, err = p.execDescribe(s, ast.TypeTable)
  68. result = append(result, r)
  69. case *ast.ExplainStreamStatement:
  70. var r string
  71. r, err = p.execExplain(s, ast.TypeStream)
  72. result = append(result, r)
  73. case *ast.ExplainTableStatement:
  74. var r string
  75. r, err = p.execExplain(s, ast.TypeTable)
  76. result = append(result, r)
  77. case *ast.DropStreamStatement:
  78. var r string
  79. r, err = p.execDrop(s, ast.TypeStream)
  80. result = append(result, r)
  81. case *ast.DropTableStatement:
  82. var r string
  83. r, err = p.execDrop(s, ast.TypeTable)
  84. result = append(result, r)
  85. default:
  86. return nil, fmt.Errorf("Invalid stream statement: %s", statement)
  87. }
  88. return
  89. }
  90. func (p *StreamProcessor) execSave(stmt *ast.StreamStmt, statement string, replace bool) error {
  91. err := p.db.Open()
  92. if err != nil {
  93. return fmt.Errorf("error when opening db: %v.", err)
  94. }
  95. defer p.db.Close()
  96. s, err := json.Marshal(xsql.StreamInfo{
  97. StreamType: stmt.StreamType,
  98. Statement: statement,
  99. })
  100. if err != nil {
  101. return fmt.Errorf("error when saving to db: %v.", err)
  102. }
  103. if replace {
  104. err = p.db.Set(string(stmt.Name), string(s))
  105. } else {
  106. err = p.db.Setnx(string(stmt.Name), string(s))
  107. }
  108. return err
  109. }
  110. func (p *StreamProcessor) ExecReplaceStream(statement string, st ast.StreamType) (string, error) {
  111. parser := xsql.NewParser(strings.NewReader(statement))
  112. stmt, err := xsql.Language.Parse(parser)
  113. if err != nil {
  114. return "", err
  115. }
  116. stt := ast.StreamTypeMap[st]
  117. switch s := stmt.(type) {
  118. case *ast.StreamStmt:
  119. if s.StreamType != st {
  120. return "", errorx.NewWithCode(errorx.NOT_FOUND, fmt.Sprintf("%s %s is not found", ast.StreamTypeMap[st], s.Name))
  121. }
  122. err = p.execSave(s, statement, true)
  123. if err != nil {
  124. return "", fmt.Errorf("Replace %s fails: %v.", stt, err)
  125. } else {
  126. info := fmt.Sprintf("%s %s is replaced.", strings.Title(stt), s.Name)
  127. log.Printf("%s", info)
  128. return info, nil
  129. }
  130. default:
  131. return "", fmt.Errorf("Invalid %s statement: %s", stt, statement)
  132. }
  133. }
  134. func (p *StreamProcessor) ExecStreamSql(statement string) (string, error) {
  135. r, err := p.ExecStmt(statement)
  136. if err != nil {
  137. return "", err
  138. } else {
  139. return strings.Join(r, "\n"), err
  140. }
  141. }
  142. func (p *StreamProcessor) execShow(st ast.StreamType) ([]string, error) {
  143. keys, err := p.ShowStream(st)
  144. if len(keys) == 0 {
  145. keys = append(keys, fmt.Sprintf("No %s definitions are found.", ast.StreamTypeMap[st]))
  146. }
  147. return keys, err
  148. }
  149. func (p *StreamProcessor) ShowStream(st ast.StreamType) ([]string, error) {
  150. stt := ast.StreamTypeMap[st]
  151. err := p.db.Open()
  152. if err != nil {
  153. return nil, fmt.Errorf("Show %ss fails, error when opening db: %v.", stt, err)
  154. }
  155. defer p.db.Close()
  156. keys, err := p.db.Keys()
  157. if err != nil {
  158. return nil, fmt.Errorf("Show %ss fails, error when loading data from db: %v.", stt, err)
  159. }
  160. var (
  161. v string
  162. vs = &xsql.StreamInfo{}
  163. result = make([]string, 0)
  164. )
  165. for _, k := range keys {
  166. if ok, _ := p.db.Get(k, &v); ok {
  167. if err := json.Unmarshal([]byte(v), vs); err == nil && vs.StreamType == st {
  168. result = append(result, k)
  169. }
  170. }
  171. }
  172. return result, nil
  173. }
  174. func (p *StreamProcessor) getStream(name string, st ast.StreamType) (string, error) {
  175. vs, err := xsql.GetDataSourceStatement(p.db, name)
  176. if vs != nil && vs.StreamType == st {
  177. return vs.Statement, nil
  178. }
  179. if err != nil {
  180. return "", err
  181. }
  182. return "", errorx.NewWithCode(errorx.NOT_FOUND, fmt.Sprintf("%s %s is not found", ast.StreamTypeMap[st], name))
  183. }
  184. func (p *StreamProcessor) execDescribe(stmt ast.NameNode, st ast.StreamType) (string, error) {
  185. streamStmt, err := p.DescStream(stmt.GetName(), st)
  186. if err != nil {
  187. return "", err
  188. }
  189. switch s := streamStmt.(type) {
  190. case *ast.StreamStmt:
  191. var buff bytes.Buffer
  192. buff.WriteString("Fields\n--------------------------------------------------------------------------------\n")
  193. for _, f := range s.StreamFields {
  194. buff.WriteString(f.Name + "\t")
  195. buff.WriteString(printFieldType(f.FieldType))
  196. buff.WriteString("\n")
  197. }
  198. buff.WriteString("\n")
  199. printOptions(s.Options, &buff)
  200. return buff.String(), err
  201. default:
  202. return "%s", fmt.Errorf("Error resolving the %s %s, the data in db may be corrupted.", ast.StreamTypeMap[st], stmt.GetName())
  203. }
  204. }
  205. func printOptions(opts *ast.Options, buff *bytes.Buffer) {
  206. if opts.CONF_KEY != "" {
  207. buff.WriteString(fmt.Sprintf("CONF_KEY: %s\n", opts.CONF_KEY))
  208. }
  209. if opts.DATASOURCE != "" {
  210. buff.WriteString(fmt.Sprintf("DATASOURCE: %s\n", opts.DATASOURCE))
  211. }
  212. if opts.FORMAT != "" {
  213. buff.WriteString(fmt.Sprintf("FORMAT: %s\n", opts.FORMAT))
  214. }
  215. if opts.KEY != "" {
  216. buff.WriteString(fmt.Sprintf("KEY: %s\n", opts.KEY))
  217. }
  218. if opts.RETAIN_SIZE != 0 {
  219. buff.WriteString(fmt.Sprintf("RETAIN_SIZE: %d\n", opts.RETAIN_SIZE))
  220. }
  221. if opts.SHARED {
  222. buff.WriteString(fmt.Sprintf("SHARED: %v\n", opts.SHARED))
  223. }
  224. if opts.STRICT_VALIDATION {
  225. buff.WriteString(fmt.Sprintf("STRICT_VALIDATION: %v\n", opts.STRICT_VALIDATION))
  226. }
  227. if opts.TIMESTAMP != "" {
  228. buff.WriteString(fmt.Sprintf("TIMESTAMP: %s\n", opts.TIMESTAMP))
  229. }
  230. if opts.TIMESTAMP_FORMAT != "" {
  231. buff.WriteString(fmt.Sprintf("TIMESTAMP_FORMAT: %s\n", opts.TIMESTAMP_FORMAT))
  232. }
  233. if opts.TYPE != "" {
  234. buff.WriteString(fmt.Sprintf("TYPE: %s\n", opts.TYPE))
  235. }
  236. }
  237. func (p *StreamProcessor) DescStream(name string, st ast.StreamType) (ast.Statement, error) {
  238. statement, err := p.getStream(name, st)
  239. if err != nil {
  240. return nil, fmt.Errorf("Describe %s fails, %s.", ast.StreamTypeMap[st], err)
  241. }
  242. parser := xsql.NewParser(strings.NewReader(statement))
  243. stream, err := xsql.Language.Parse(parser)
  244. if err != nil {
  245. return nil, err
  246. }
  247. return stream, nil
  248. }
  249. func (p *StreamProcessor) execExplain(stmt ast.NameNode, st ast.StreamType) (string, error) {
  250. _, err := p.getStream(stmt.GetName(), st)
  251. if err != nil {
  252. return "", fmt.Errorf("Explain %s fails, %s.", ast.StreamTypeMap[st], err)
  253. }
  254. return "TO BE SUPPORTED", nil
  255. }
  256. func (p *StreamProcessor) execDrop(stmt ast.NameNode, st ast.StreamType) (string, error) {
  257. s, err := p.DropStream(stmt.GetName(), st)
  258. if err != nil {
  259. return s, fmt.Errorf("Drop %s fails: %s.", ast.StreamTypeMap[st], err)
  260. }
  261. return s, nil
  262. }
  263. func (p *StreamProcessor) DropStream(name string, st ast.StreamType) (string, error) {
  264. defer p.db.Close()
  265. _, err := p.getStream(name, st)
  266. if err != nil {
  267. return "", err
  268. }
  269. err = p.db.Open()
  270. if err != nil {
  271. return "", fmt.Errorf("error when opening db: %v", err)
  272. }
  273. defer p.db.Close()
  274. err = p.db.Delete(name)
  275. if err != nil {
  276. return "", err
  277. } else {
  278. return fmt.Sprintf("%s %s is dropped.", strings.Title(ast.StreamTypeMap[st]), name), nil
  279. }
  280. }
  281. func printFieldType(ft ast.FieldType) (result string) {
  282. switch t := ft.(type) {
  283. case *ast.BasicType:
  284. result = t.Type.String()
  285. case *ast.ArrayType:
  286. result = "array("
  287. if t.FieldType != nil {
  288. result += printFieldType(t.FieldType)
  289. } else {
  290. result += t.Type.String()
  291. }
  292. result += ")"
  293. case *ast.RecType:
  294. result = "struct("
  295. isFirst := true
  296. for _, f := range t.StreamFields {
  297. if isFirst {
  298. isFirst = false
  299. } else {
  300. result += ", "
  301. }
  302. result = result + f.Name + " " + printFieldType(f.FieldType)
  303. }
  304. result += ")"
  305. }
  306. return
  307. }