rpc.go 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458
  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 server
  15. import (
  16. "bytes"
  17. "encoding/json"
  18. "fmt"
  19. "github.com/lf-edge/ekuiper/internal/pkg/model"
  20. "github.com/lf-edge/ekuiper/internal/plugin"
  21. "github.com/lf-edge/ekuiper/internal/service"
  22. "github.com/lf-edge/ekuiper/internal/topo/sink"
  23. "strings"
  24. "time"
  25. )
  26. const QueryRuleId = "internal-ekuiper_query_rule"
  27. type Server int
  28. func (t *Server) CreateQuery(sql string, reply *string) error {
  29. if _, ok := registry.Load(QueryRuleId); ok {
  30. stopQuery()
  31. }
  32. tp, err := ruleProcessor.ExecQuery(QueryRuleId, sql)
  33. if err != nil {
  34. return err
  35. } else {
  36. rs := &RuleState{Name: QueryRuleId, Topology: tp, Triggered: true}
  37. registry.Store(QueryRuleId, rs)
  38. msg := fmt.Sprintf("Query was submit successfully.")
  39. logger.Println(msg)
  40. *reply = fmt.Sprintf(msg)
  41. }
  42. return nil
  43. }
  44. func stopQuery() {
  45. if rs, ok := registry.Load(QueryRuleId); ok {
  46. logger.Printf("stop the query.")
  47. (*rs.Topology).Cancel()
  48. registry.Delete(QueryRuleId)
  49. }
  50. }
  51. /**
  52. * qid is not currently used.
  53. */
  54. func (t *Server) GetQueryResult(qid string, reply *string) error {
  55. if rs, ok := registry.Load(QueryRuleId); ok {
  56. c := (*rs.Topology).GetContext()
  57. if c != nil && c.Err() != nil {
  58. return c.Err()
  59. }
  60. }
  61. sink.QR.LastFetch = time.Now()
  62. sink.QR.Mux.Lock()
  63. if len(sink.QR.Results) > 0 {
  64. *reply = strings.Join(sink.QR.Results, "")
  65. sink.QR.Results = make([]string, 10)
  66. } else {
  67. *reply = ""
  68. }
  69. sink.QR.Mux.Unlock()
  70. return nil
  71. }
  72. func (t *Server) Stream(stream string, reply *string) error {
  73. content, err := streamProcessor.ExecStmt(stream)
  74. if err != nil {
  75. return fmt.Errorf("Stream command error: %s", err)
  76. } else {
  77. for _, c := range content {
  78. *reply = *reply + fmt.Sprintln(c)
  79. }
  80. }
  81. return nil
  82. }
  83. func (t *Server) CreateRule(rule *model.RPCArgDesc, reply *string) error {
  84. r, err := ruleProcessor.ExecCreate(rule.Name, rule.Json)
  85. if err != nil {
  86. return fmt.Errorf("Create rule error : %s.", err)
  87. } else {
  88. *reply = fmt.Sprintf("Rule %s was created successfully, please use 'bin/kuiper getstatus rule %s' command to get rule status.", rule.Name, rule.Name)
  89. }
  90. //Start the rule
  91. rs, err := createRuleState(r)
  92. if err != nil {
  93. return err
  94. }
  95. err = doStartRule(rs)
  96. if err != nil {
  97. return err
  98. }
  99. return nil
  100. }
  101. func (t *Server) GetStatusRule(name string, reply *string) error {
  102. if r, err := getRuleStatus(name); err != nil {
  103. return err
  104. } else {
  105. *reply = r
  106. }
  107. return nil
  108. }
  109. func (t *Server) GetTopoRule(name string, reply *string) error {
  110. if r, err := getRuleTopo(name); err != nil {
  111. return err
  112. } else {
  113. dst := &bytes.Buffer{}
  114. if err = json.Indent(dst, []byte(r), "", " "); err != nil {
  115. *reply = r
  116. } else {
  117. *reply = dst.String()
  118. }
  119. }
  120. return nil
  121. }
  122. func (t *Server) StartRule(name string, reply *string) error {
  123. if err := startRule(name); err != nil {
  124. return err
  125. } else {
  126. *reply = fmt.Sprintf("Rule %s was started", name)
  127. }
  128. return nil
  129. }
  130. func (t *Server) StopRule(name string, reply *string) error {
  131. *reply = stopRule(name)
  132. return nil
  133. }
  134. func (t *Server) RestartRule(name string, reply *string) error {
  135. err := restartRule(name)
  136. if err != nil {
  137. return err
  138. }
  139. *reply = fmt.Sprintf("Rule %s was restarted.", name)
  140. return nil
  141. }
  142. func (t *Server) DescRule(name string, reply *string) error {
  143. r, err := ruleProcessor.ExecDesc(name)
  144. if err != nil {
  145. return fmt.Errorf("Desc rule error : %s.", err)
  146. } else {
  147. *reply = r
  148. }
  149. return nil
  150. }
  151. func (t *Server) ShowRules(_ int, reply *string) error {
  152. r, err := getAllRulesWithStatus()
  153. if err != nil {
  154. return fmt.Errorf("Show rule error : %s.", err)
  155. }
  156. if len(r) == 0 {
  157. *reply = "No rule definitions are found."
  158. } else {
  159. result, err := json.Marshal(r)
  160. if err != nil {
  161. return fmt.Errorf("Show rule error : %s.", err)
  162. }
  163. dst := &bytes.Buffer{}
  164. if err := json.Indent(dst, result, "", " "); err != nil {
  165. return fmt.Errorf("Show rule error : %s.", err)
  166. }
  167. *reply = dst.String()
  168. }
  169. return nil
  170. }
  171. func (t *Server) DropRule(name string, reply *string) error {
  172. deleteRule(name)
  173. r, err := ruleProcessor.ExecDrop(name)
  174. if err != nil {
  175. return fmt.Errorf("Drop rule error : %s.", err)
  176. } else {
  177. err := t.StopRule(name, reply)
  178. if err != nil {
  179. return err
  180. }
  181. }
  182. *reply = r
  183. return nil
  184. }
  185. func (t *Server) CreatePlugin(arg *model.PluginDesc, reply *string) error {
  186. pt := plugin.PluginType(arg.Type)
  187. p, err := getPluginByJson(arg, pt)
  188. if err != nil {
  189. return fmt.Errorf("Create plugin error: %s", err)
  190. }
  191. if p.GetFile() == "" {
  192. return fmt.Errorf("Create plugin error: Missing plugin file url.")
  193. }
  194. err = pluginManager.Register(pt, p)
  195. if err != nil {
  196. return fmt.Errorf("Create plugin error: %s", err)
  197. } else {
  198. *reply = fmt.Sprintf("Plugin %s is created.", p.GetName())
  199. }
  200. return nil
  201. }
  202. func (t *Server) RegisterPlugin(arg *model.PluginDesc, reply *string) error {
  203. p, err := getPluginByJson(arg, plugin.FUNCTION)
  204. if err != nil {
  205. return fmt.Errorf("Register plugin functions error: %s", err)
  206. }
  207. if len(p.GetSymbols()) == 0 {
  208. return fmt.Errorf("Register plugin functions error: Missing function list.")
  209. }
  210. err = pluginManager.RegisterFuncs(p.GetName(), p.GetSymbols())
  211. if err != nil {
  212. return fmt.Errorf("Create plugin error: %s", err)
  213. } else {
  214. *reply = fmt.Sprintf("Plugin %s is created.", p.GetName())
  215. }
  216. return nil
  217. }
  218. func (t *Server) DropPlugin(arg *model.PluginDesc, reply *string) error {
  219. pt := plugin.PluginType(arg.Type)
  220. p, err := getPluginByJson(arg, pt)
  221. if err != nil {
  222. return fmt.Errorf("Drop plugin error: %s", err)
  223. }
  224. err = pluginManager.Delete(pt, p.GetName(), arg.Stop)
  225. if err != nil {
  226. return fmt.Errorf("Drop plugin error: %s", err)
  227. } else {
  228. if arg.Stop {
  229. *reply = fmt.Sprintf("Plugin %s is dropped and Kuiper will be stopped.", p.GetName())
  230. } else {
  231. *reply = fmt.Sprintf("Plugin %s is dropped and Kuiper must restart for the change to take effect.", p.GetName())
  232. }
  233. }
  234. return nil
  235. }
  236. func (t *Server) ShowPlugins(arg int, reply *string) error {
  237. pt := plugin.PluginType(arg)
  238. l, err := pluginManager.List(pt)
  239. if err != nil {
  240. return fmt.Errorf("Show plugin error: %s", err)
  241. } else {
  242. if len(l) == 0 {
  243. l = append(l, "No plugin is found.")
  244. }
  245. *reply = strings.Join(l, "\n")
  246. }
  247. return nil
  248. }
  249. func (t *Server) ShowUdfs(_ int, reply *string) error {
  250. l, err := pluginManager.ListSymbols()
  251. if err != nil {
  252. return fmt.Errorf("Show UDFs error: %s", err)
  253. } else {
  254. if len(l) == 0 {
  255. l = append(l, "No udf is found.")
  256. }
  257. *reply = strings.Join(l, "\n")
  258. }
  259. return nil
  260. }
  261. func (t *Server) DescPlugin(arg *model.PluginDesc, reply *string) error {
  262. pt := plugin.PluginType(arg.Type)
  263. p, err := getPluginByJson(arg, pt)
  264. if err != nil {
  265. return fmt.Errorf("Describe plugin error: %s", err)
  266. }
  267. m, ok := pluginManager.Get(pt, p.GetName())
  268. if !ok {
  269. return fmt.Errorf("Describe plugin error: not found")
  270. } else {
  271. r, err := marshalDesc(m)
  272. if err != nil {
  273. return fmt.Errorf("Describe plugin error: %v", err)
  274. }
  275. *reply = r
  276. }
  277. return nil
  278. }
  279. func (t *Server) DescUdf(arg string, reply *string) error {
  280. m, ok := pluginManager.GetSymbol(arg)
  281. if !ok {
  282. return fmt.Errorf("Describe udf error: not found")
  283. } else {
  284. j := map[string]string{
  285. "name": arg,
  286. "plugin": m,
  287. }
  288. r, err := marshalDesc(j)
  289. if err != nil {
  290. return fmt.Errorf("Describe udf error: %v", err)
  291. }
  292. *reply = r
  293. }
  294. return nil
  295. }
  296. func (t *Server) CreateService(arg *model.RPCArgDesc, reply *string) error {
  297. sd := &service.ServiceCreationRequest{}
  298. if arg.Json != "" {
  299. if err := json.Unmarshal([]byte(arg.Json), sd); err != nil {
  300. return fmt.Errorf("Parse service %s error : %s.", arg.Json, err)
  301. }
  302. }
  303. if sd.Name != arg.Name {
  304. return fmt.Errorf("Create service error: name mismatch.")
  305. }
  306. if sd.File == "" {
  307. return fmt.Errorf("Create service error: Missing service file url.")
  308. }
  309. err := serviceManager.Create(sd)
  310. if err != nil {
  311. return fmt.Errorf("Create service error: %s", err)
  312. } else {
  313. *reply = fmt.Sprintf("Service %s is created.", arg.Name)
  314. }
  315. return nil
  316. }
  317. func (t *Server) DescService(name string, reply *string) error {
  318. s, err := serviceManager.Get(name)
  319. if err != nil {
  320. return fmt.Errorf("Desc service error : %s.", err)
  321. } else {
  322. r, err := marshalDesc(s)
  323. if err != nil {
  324. return fmt.Errorf("Describe service error: %v", err)
  325. }
  326. *reply = r
  327. }
  328. return nil
  329. }
  330. func (t *Server) DescServiceFunc(name string, reply *string) error {
  331. s, err := serviceManager.GetFunction(name)
  332. if err != nil {
  333. return fmt.Errorf("Desc service func error : %s.", err)
  334. } else {
  335. r, err := marshalDesc(s)
  336. if err != nil {
  337. return fmt.Errorf("Describe service func error: %v", err)
  338. }
  339. *reply = r
  340. }
  341. return nil
  342. }
  343. func (t *Server) DropService(name string, reply *string) error {
  344. err := serviceManager.Delete(name)
  345. if err != nil {
  346. return fmt.Errorf("Drop service error : %s.", err)
  347. }
  348. *reply = fmt.Sprintf("Service %s is dropped", name)
  349. return nil
  350. }
  351. func (t *Server) ShowServices(_ int, reply *string) error {
  352. s, err := serviceManager.List()
  353. if err != nil {
  354. return fmt.Errorf("Show service error: %s.", err)
  355. }
  356. if len(s) == 0 {
  357. *reply = "No service definitions are found."
  358. } else {
  359. r, err := marshalDesc(s)
  360. if err != nil {
  361. return fmt.Errorf("Show service error: %v", err)
  362. }
  363. *reply = r
  364. }
  365. return nil
  366. }
  367. func (t *Server) ShowServiceFuncs(_ int, reply *string) error {
  368. s, err := serviceManager.ListFunctions()
  369. if err != nil {
  370. return fmt.Errorf("Show service funcs error: %s.", err)
  371. }
  372. if len(s) == 0 {
  373. *reply = "No service definitions are found."
  374. } else {
  375. r, err := marshalDesc(s)
  376. if err != nil {
  377. return fmt.Errorf("Show service funcs error: %v", err)
  378. }
  379. *reply = r
  380. }
  381. return nil
  382. }
  383. func marshalDesc(m interface{}) (string, error) {
  384. s, err := json.Marshal(m)
  385. if err != nil {
  386. return "", fmt.Errorf("invalid json %v", m)
  387. }
  388. dst := &bytes.Buffer{}
  389. if err := json.Indent(dst, s, "", " "); err != nil {
  390. return "", fmt.Errorf("indent json error %v", err)
  391. }
  392. return dst.String(), nil
  393. }
  394. func getPluginByJson(arg *model.PluginDesc, pt plugin.PluginType) (plugin.Plugin, error) {
  395. p := plugin.NewPluginByType(pt)
  396. if arg.Json != "" {
  397. if err := json.Unmarshal([]byte(arg.Json), p); err != nil {
  398. return nil, fmt.Errorf("Parse plugin %s error : %s.", arg.Json, err)
  399. }
  400. }
  401. p.SetName(arg.Name)
  402. return p, nil
  403. }
  404. func init() {
  405. ticker := time.NewTicker(time.Second * 5)
  406. go func() {
  407. for {
  408. <-ticker.C
  409. if _, ok := registry.Load(QueryRuleId); !ok {
  410. continue
  411. }
  412. n := time.Now()
  413. w := 10 * time.Second
  414. if v := n.Sub(sink.QR.LastFetch); v >= w {
  415. logger.Printf("The client seems no longer fetch the query result, stop the query now.")
  416. stopQuery()
  417. ticker.Stop()
  418. return
  419. }
  420. }
  421. }()
  422. }