agi.cluster.go 9.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277
  1. package agi
  2. /*
  3. AGI Cluster Library
  4. Author: tobychui
  5. Exposes the ArozOS cluster (nodes, namespace metadata, replica policies,
  6. event hooks) to AGI scripts. Ordinary file access does not need this
  7. library: filelib already works on cluster:/ paths through the mounted
  8. drive. Only present when the host wired a ClusterProvider in.
  9. Usage in AGI:
  10. requirelib("cluster");
  11. cluster.inCluster() // bool
  12. cluster.self() // this node
  13. cluster.nodes() // [ {id, name, state, ...} ]
  14. cluster.status() // cluster, leader, identity origin, volumes
  15. cluster.stat("cluster:/photos/a.jpg") // record with copies
  16. cluster.list("cluster:/photos") // records
  17. cluster.setReplicas(path, n) // admin, per file
  18. cluster.policy("/photos", n) // admin, per top-level folder
  19. cluster.on("file.created", "user:/hook.agi") -> hook id
  20. cluster.off(hookId)
  21. cluster.hooks() // this user's hooks
  22. cluster.emit("app.custom", {any: "json"})
  23. */
  24. import (
  25. "encoding/json"
  26. "errors"
  27. "fmt"
  28. "os"
  29. "strings"
  30. "github.com/robertkrimen/otto"
  31. "imuslab.com/arozos/mod/agi/static"
  32. "imuslab.com/arozos/mod/info/logger"
  33. )
  34. // ClusterProvider is implemented by the core on top of the cluster packages.
  35. type ClusterProvider interface {
  36. InCluster() bool
  37. Self() interface{}
  38. Nodes() interface{}
  39. Status() interface{}
  40. Stat(path string) (interface{}, error)
  41. List(path string) (interface{}, error)
  42. SetReplicas(path string, n int) error
  43. SetPolicy(folder string, n int) error
  44. AddHook(owner string, types []string, script string) (interface{}, error)
  45. RemoveHook(id string, owner string) error
  46. Hooks(owner string) interface{}
  47. Emit(user string, evType string, data []byte) error
  48. SubmitJob(owner string, name string, scriptVpath string, args []byte, inputs []string, features []string, nodes []string, timeoutSec int, dataset string, partitionMax int) (interface{}, error)
  49. JobStatus(id string, requester string, isAdmin bool) (interface{}, error)
  50. JobList(owner string) interface{}
  51. CancelJob(id string, requester string, isAdmin bool) error
  52. WaitJob(id string, timeoutSec int, requester string, isAdmin bool) (interface{}, error)
  53. }
  54. func (g *Gateway) ClusterLibRegister() {
  55. err := g.RegisterLib("cluster", g.injectClusterLibFunctions)
  56. if err != nil {
  57. logger.PrintAndLog("Agi", fmt.Sprint(err), nil)
  58. os.Exit(1)
  59. }
  60. }
  61. func (g *Gateway) injectClusterLibFunctions(payload *static.AgiLibInjectionPayload) {
  62. vm := payload.VM
  63. u := payload.User
  64. p := g.Option.ClusterProvider
  65. if p == nil {
  66. return
  67. }
  68. username := ""
  69. if u != nil {
  70. username = u.Username
  71. }
  72. isAdmin := func() bool { return u != nil && u.IsAdmin() }
  73. fail := func(err error) otto.Value {
  74. panic(vm.MakeCustomError("ClusterError", err.Error()))
  75. }
  76. toJSON := func(v interface{}) otto.Value {
  77. js, err := json.Marshal(v)
  78. if err != nil {
  79. return fail(err)
  80. }
  81. r, _ := vm.ToValue(string(js))
  82. return r
  83. }
  84. // paths accept cluster:/x or /x; normalise to a vpath for permission checks
  85. vpathArg := func(call otto.FunctionCall, i int) string {
  86. raw, _ := call.Argument(i).ToString()
  87. raw = strings.TrimSpace(raw)
  88. if !strings.HasPrefix(strings.ToLower(raw), "cluster:") {
  89. raw = "cluster:" + strings.TrimPrefix(raw, "/")
  90. raw = strings.Replace(raw, "cluster:", "cluster:/", 1)
  91. }
  92. return raw
  93. }
  94. vm.Set("_cluster_inCluster", func(call otto.FunctionCall) otto.Value {
  95. r, _ := vm.ToValue(p.InCluster())
  96. return r
  97. })
  98. vm.Set("_cluster_self", func(call otto.FunctionCall) otto.Value { return toJSON(p.Self()) })
  99. vm.Set("_cluster_nodes", func(call otto.FunctionCall) otto.Value { return toJSON(p.Nodes()) })
  100. vm.Set("_cluster_status", func(call otto.FunctionCall) otto.Value { return toJSON(p.Status()) })
  101. vm.Set("_cluster_stat", func(call otto.FunctionCall) otto.Value {
  102. vpath := vpathArg(call, 0)
  103. if u != nil && !u.CanRead(vpath) {
  104. return fail(errors.New("path access denied: " + vpath))
  105. }
  106. rec, err := p.Stat(vpath)
  107. if err != nil {
  108. return fail(err)
  109. }
  110. return toJSON(rec)
  111. })
  112. vm.Set("_cluster_list", func(call otto.FunctionCall) otto.Value {
  113. vpath := vpathArg(call, 0)
  114. if u != nil && !u.CanRead(vpath) {
  115. return fail(errors.New("path access denied: " + vpath))
  116. }
  117. recs, err := p.List(vpath)
  118. if err != nil {
  119. return fail(err)
  120. }
  121. return toJSON(recs)
  122. })
  123. vm.Set("_cluster_setReplicas", func(call otto.FunctionCall) otto.Value {
  124. if !isAdmin() {
  125. return fail(errors.New("admin permission required"))
  126. }
  127. n, _ := call.Argument(1).ToInteger()
  128. if err := p.SetReplicas(vpathArg(call, 0), int(n)); err != nil {
  129. return fail(err)
  130. }
  131. return otto.TrueValue()
  132. })
  133. vm.Set("_cluster_policy", func(call otto.FunctionCall) otto.Value {
  134. if !isAdmin() {
  135. return fail(errors.New("admin permission required"))
  136. }
  137. folder, _ := call.Argument(0).ToString()
  138. n, _ := call.Argument(1).ToInteger()
  139. if err := p.SetPolicy(folder, int(n)); err != nil {
  140. return fail(err)
  141. }
  142. return otto.TrueValue()
  143. })
  144. vm.Set("_cluster_on", func(call otto.FunctionCall) otto.Value {
  145. types, _ := call.Argument(0).ToString()
  146. script, _ := call.Argument(1).ToString()
  147. if payload.ScriptFsh != nil && u != nil {
  148. script = static.RelativeVpathRewrite(payload.ScriptFsh, script, vm, u)
  149. }
  150. if u != nil && !u.CanRead(script) {
  151. return fail(errors.New("script access denied: " + script))
  152. }
  153. h, err := p.AddHook(username, strings.Split(types, ","), script)
  154. if err != nil {
  155. return fail(err)
  156. }
  157. return toJSON(h)
  158. })
  159. vm.Set("_cluster_off", func(call otto.FunctionCall) otto.Value {
  160. id, _ := call.Argument(0).ToString()
  161. owner := username
  162. if isAdmin() {
  163. owner = ""
  164. }
  165. if err := p.RemoveHook(id, owner); err != nil {
  166. return fail(err)
  167. }
  168. return otto.TrueValue()
  169. })
  170. vm.Set("_cluster_hooks", func(call otto.FunctionCall) otto.Value {
  171. owner := username
  172. if isAdmin() {
  173. owner = ""
  174. }
  175. return toJSON(p.Hooks(owner))
  176. })
  177. vm.Set("_cluster_emit", func(call otto.FunctionCall) otto.Value {
  178. evType, _ := call.Argument(0).ToString()
  179. data, _ := call.Argument(1).ToString()
  180. if err := p.Emit(username, evType, []byte(data)); err != nil {
  181. return fail(err)
  182. }
  183. return otto.TrueValue()
  184. })
  185. vm.Set("_cluster_jobSubmit", func(call otto.FunctionCall) otto.Value {
  186. raw, _ := call.Argument(0).ToString()
  187. var req struct {
  188. Name string `json:"name"`
  189. Script string `json:"script"`
  190. Args json.RawMessage `json:"args"`
  191. Inputs []string `json:"inputs"`
  192. Features []string `json:"features"`
  193. Nodes []string `json:"nodes"`
  194. Timeout int `json:"timeout"`
  195. Dataset string `json:"dataset"`
  196. PartitionMax int `json:"partitionMax"`
  197. }
  198. if err := json.Unmarshal([]byte(raw), &req); err != nil {
  199. return fail(err)
  200. }
  201. if payload.ScriptFsh != nil && u != nil {
  202. req.Script = static.RelativeVpathRewrite(payload.ScriptFsh, req.Script, vm, u)
  203. }
  204. if u != nil && !u.CanRead(req.Script) {
  205. return fail(errors.New("script access denied: " + req.Script))
  206. }
  207. rec, err := p.SubmitJob(username, req.Name, req.Script, req.Args, req.Inputs, req.Features, req.Nodes, req.Timeout, req.Dataset, req.PartitionMax)
  208. if err != nil {
  209. return fail(err)
  210. }
  211. return toJSON(rec)
  212. })
  213. vm.Set("_cluster_jobStatus", func(call otto.FunctionCall) otto.Value {
  214. id, _ := call.Argument(0).ToString()
  215. rec, err := p.JobStatus(id, username, isAdmin())
  216. if err != nil {
  217. return fail(err)
  218. }
  219. return toJSON(rec)
  220. })
  221. vm.Set("_cluster_jobList", func(call otto.FunctionCall) otto.Value {
  222. owner := username
  223. if isAdmin() {
  224. owner = ""
  225. }
  226. return toJSON(p.JobList(owner))
  227. })
  228. vm.Set("_cluster_jobCancel", func(call otto.FunctionCall) otto.Value {
  229. id, _ := call.Argument(0).ToString()
  230. if err := p.CancelJob(id, username, isAdmin()); err != nil {
  231. return fail(err)
  232. }
  233. return otto.TrueValue()
  234. })
  235. vm.Set("_cluster_jobWait", func(call otto.FunctionCall) otto.Value {
  236. id, _ := call.Argument(0).ToString()
  237. secs, _ := call.Argument(1).ToInteger()
  238. rec, err := p.WaitJob(id, int(secs), username, isAdmin())
  239. if err != nil {
  240. return fail(err)
  241. }
  242. return toJSON(rec)
  243. })
  244. vm.Run(`
  245. var cluster = {};
  246. cluster.inCluster = function() { return _cluster_inCluster(); };
  247. cluster.self = function() { return JSON.parse(_cluster_self()); };
  248. cluster.nodes = function() { return JSON.parse(_cluster_nodes()); };
  249. cluster.status = function() { return JSON.parse(_cluster_status()); };
  250. cluster.stat = function(p) { return JSON.parse(_cluster_stat(p)); };
  251. cluster.list = function(p) { return JSON.parse(_cluster_list(p)); };
  252. cluster.setReplicas = function(p, n) { return _cluster_setReplicas(p, n); };
  253. cluster.policy = function(folder, n) { return _cluster_policy(folder, n); };
  254. cluster.on = function(types, script) { var h = JSON.parse(_cluster_on(Array.isArray(types) ? types.join(",") : types, script)); return h.id; };
  255. cluster.off = function(id) { return _cluster_off(id); };
  256. cluster.hooks = function() { return JSON.parse(_cluster_hooks()); };
  257. cluster.emit = function(type, data) { return _cluster_emit(type, JSON.stringify(data === undefined ? {} : data)); };
  258. cluster.jobs = {};
  259. cluster.jobs.submit = function(spec) { return JSON.parse(_cluster_jobSubmit(JSON.stringify(spec))).spec.id; };
  260. cluster.jobs.status = function(id) { return JSON.parse(_cluster_jobStatus(id)); };
  261. cluster.jobs.list = function() { return JSON.parse(_cluster_jobList()); };
  262. cluster.jobs.cancel = function(id) { return _cluster_jobCancel(id); };
  263. cluster.jobs.wait = function(id, timeoutSec) { return JSON.parse(_cluster_jobWait(id, timeoutSec === undefined ? 300 : timeoutSec)); };
  264. cluster.jobs.mapreduce = function(spec) { spec = spec || {}; if (!spec.dataset) { throw new Error("a map/reduce job needs a dataset pattern"); } return cluster.jobs.submit(spec); };
  265. `)
  266. }