cluster.go 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772
  1. package main
  2. import (
  3. "io"
  4. "net/http"
  5. "net/url"
  6. "path/filepath"
  7. "sync"
  8. "time"
  9. "encoding/json"
  10. "errors"
  11. "strings"
  12. "imuslab.com/arozos/mod/cluster/capability"
  13. "imuslab.com/arozos/mod/cluster/events"
  14. "imuslab.com/arozos/mod/cluster/identity"
  15. "imuslab.com/arozos/mod/cluster/jobs"
  16. "imuslab.com/arozos/mod/cluster/membership"
  17. "imuslab.com/arozos/mod/cluster/metadata"
  18. "imuslab.com/arozos/mod/cluster/replication"
  19. "imuslab.com/arozos/mod/cluster/scheduling"
  20. "imuslab.com/arozos/mod/cluster/storage"
  21. fs "imuslab.com/arozos/mod/filesystem"
  22. "imuslab.com/arozos/mod/filesystem/abstractions/clusterfs"
  23. "imuslab.com/arozos/mod/filesystem/arozfs"
  24. "imuslab.com/arozos/mod/info/usageinfo"
  25. "imuslab.com/arozos/mod/network/neighbour"
  26. prout "imuslab.com/arozos/mod/prouter"
  27. "imuslab.com/arozos/mod/time/nightly"
  28. )
  29. /*
  30. Functions related to ArozOS clusters
  31. Author: tobychui
  32. This is a section of the arozos core that handle cluster
  33. related function endpoints:
  34. - Neighbourhood: mDNS discovery of nearby ArozOS hosts (LAN only)
  35. - Cluster membership: the cluster agent of this node (mod/cluster/membership)
  36. reachable by other nodes under /cluster/acn/* (see main.router.go)
  37. */
  38. var (
  39. NeighbourDiscoverer *neighbour.Discoverer
  40. clusterManager *membership.Manager
  41. clusterIdentity *identity.Manager
  42. clusterMetadata *metadata.Manager
  43. clusterStorage *storage.Service
  44. clusterReplication *replication.Manager
  45. clusterEvents *events.Bus
  46. clusterJobs *jobs.Manager
  47. clusterScheduling *scheduling.Manager
  48. clusterMountMu sync.Mutex
  49. )
  50. // clusterRunHook executes an event hook script as its owner with the event
  51. // readable through postPara("event") (JSON).
  52. func clusterRunHook(h events.Hook, ev events.Event) error {
  53. if AGIGateway == nil {
  54. return errors.New("AGI gateway not ready")
  55. }
  56. u, err := userHandler.GetUserInfoFromUsername(h.Owner)
  57. if err != nil {
  58. return err
  59. }
  60. fsh, err := u.GetFileSystemHandlerFromVirtualPath(h.Script)
  61. if err != nil {
  62. return err
  63. }
  64. rpath, err := fsh.FileSystemAbstraction.VirtualPathToRealPath(h.Script, u.Username)
  65. if err != nil {
  66. return err
  67. }
  68. if !fsh.FileSystemAbstraction.FileExists(rpath) {
  69. return errors.New("hook script not found: " + h.Script)
  70. }
  71. js, _ := json.Marshal(ev)
  72. form := "event=" + url.QueryEscape(string(js))
  73. req, _ := http.NewRequest(http.MethodPost, "/system/cluster/events/hook", strings.NewReader(form))
  74. req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
  75. _, _, err = AGIGateway.ExecuteAGIScriptAsUser(fsh, rpath, u, nil, req)
  76. return err
  77. }
  78. // clusterProvider adapts the cluster services to the AGI cluster library.
  79. type clusterProvider struct{}
  80. func (c *clusterProvider) InCluster() bool {
  81. return clusterManager != nil && clusterManager.InCluster()
  82. }
  83. func (c *clusterProvider) Self() interface{} {
  84. for _, n := range clusterManager.NodeViews() {
  85. if n.Local {
  86. return n
  87. }
  88. }
  89. return nil
  90. }
  91. func (c *clusterProvider) Nodes() interface{} { return clusterManager.NodeViews() }
  92. func (c *clusterProvider) Status() interface{} {
  93. out := map[string]interface{}{"cluster": clusterManager.Cluster(), "identityOrigin": clusterManager.IdentityOrigin()}
  94. if clusterMetadata != nil {
  95. out["metadata"] = clusterMetadata.Status()
  96. out["volumes"] = clusterMetadata.Volumes()
  97. }
  98. if clusterStorage != nil {
  99. out["storage"] = clusterStorage.Status()
  100. }
  101. if clusterReplication != nil {
  102. out["replication"] = clusterReplication.Status()
  103. }
  104. return out
  105. }
  106. func (c *clusterProvider) Stat(path string) (interface{}, error) {
  107. if clusterMetadata == nil {
  108. return nil, errors.New("metadata store not available")
  109. }
  110. return clusterMetadata.Stat(path)
  111. }
  112. func (c *clusterProvider) List(path string) (interface{}, error) {
  113. if clusterMetadata == nil {
  114. return nil, errors.New("metadata store not available")
  115. }
  116. return clusterMetadata.ListDir(path)
  117. }
  118. func (c *clusterProvider) SetReplicas(path string, n int) error {
  119. if clusterMetadata == nil {
  120. return errors.New("metadata store not available")
  121. }
  122. if n < 0 || n > 16 {
  123. return errors.New("replicas must be between 0 and 16")
  124. }
  125. rec, err := clusterMetadata.Stat(path)
  126. if err != nil {
  127. return err
  128. }
  129. rec.Replicas = n
  130. return clusterMetadata.Submit(metadata.KindFile, rec)
  131. }
  132. func (c *clusterProvider) SetPolicy(folder string, n int) error {
  133. if clusterMetadata == nil {
  134. return errors.New("metadata store not available")
  135. }
  136. if n < 1 || n > 16 {
  137. return errors.New("replicas must be between 1 and 16")
  138. }
  139. return clusterMetadata.Submit(metadata.KindPolicy, &metadata.Policy{Folder: folder, Replicas: n})
  140. }
  141. func (c *clusterProvider) AddHook(owner string, types []string, script string) (interface{}, error) {
  142. if clusterEvents == nil {
  143. return nil, errors.New("event bus not available")
  144. }
  145. return clusterEvents.AddHook(owner, types, script)
  146. }
  147. func (c *clusterProvider) RemoveHook(id string, owner string) error {
  148. if clusterEvents == nil {
  149. return errors.New("event bus not available")
  150. }
  151. return clusterEvents.RemoveHook(id, owner)
  152. }
  153. func (c *clusterProvider) Hooks(owner string) interface{} {
  154. if clusterEvents == nil {
  155. return []events.Hook{}
  156. }
  157. return clusterEvents.Hooks(owner)
  158. }
  159. func (c *clusterProvider) SubmitJob(owner string, name string, scriptVpath string, args []byte, inputs []string, features []string, nodes []string, timeoutSec int, dataset string, partitionMax int) (interface{}, error) {
  160. if clusterJobs == nil {
  161. return nil, errors.New("cluster jobs not available")
  162. }
  163. source, err := clusterReadScript(owner, scriptVpath)
  164. if err != nil {
  165. return nil, err
  166. }
  167. req := jobs.SubmitRequest{
  168. Name: name, ScriptVpath: scriptVpath, Script: source,
  169. Inputs: inputs, Features: features, Nodes: nodes, TimeoutSec: timeoutSec,
  170. Dataset: dataset, PartitionMax: partitionMax,
  171. }
  172. if len(args) > 0 && json.Valid(args) {
  173. req.Args = json.RawMessage(args)
  174. }
  175. return clusterJobs.Submit(jobs.SpecFrom(req, owner))
  176. }
  177. func (c *clusterProvider) JobStatus(id string, requester string, isAdmin bool) (interface{}, error) {
  178. if clusterJobs == nil {
  179. return nil, errors.New("cluster jobs not available")
  180. }
  181. rec, ok := clusterJobs.Get(id)
  182. if !ok || (!isAdmin && rec.Spec.Owner != requester) {
  183. return nil, errors.New("job not found")
  184. }
  185. return rec, nil
  186. }
  187. func (c *clusterProvider) JobList(owner string) interface{} {
  188. if clusterJobs == nil {
  189. return []jobs.Record{}
  190. }
  191. return clusterJobs.List(owner)
  192. }
  193. func (c *clusterProvider) CancelJob(id string, requester string, isAdmin bool) error {
  194. if clusterJobs == nil {
  195. return errors.New("cluster jobs not available")
  196. }
  197. return clusterJobs.Cancel(id, requester, isAdmin)
  198. }
  199. func (c *clusterProvider) WaitJob(id string, timeoutSec int, requester string, isAdmin bool) (interface{}, error) {
  200. if clusterJobs == nil {
  201. return nil, errors.New("cluster jobs not available")
  202. }
  203. rec, ok := clusterJobs.Get(id)
  204. if !ok || (!isAdmin && rec.Spec.Owner != requester) {
  205. return nil, errors.New("job not found")
  206. }
  207. return clusterJobs.WaitFor(id, timeoutSec)
  208. }
  209. func (c *clusterProvider) Emit(user string, evType string, data []byte) error {
  210. if clusterEvents == nil {
  211. return errors.New("event bus not available")
  212. }
  213. evType = strings.ToLower(strings.TrimSpace(evType))
  214. if evType == "" {
  215. return errors.New("event type required")
  216. }
  217. if !strings.HasPrefix(evType, "app.") {
  218. evType = "app." + evType
  219. }
  220. if len(data) > 64<<10 {
  221. return errors.New("event data too large (64 KB max)")
  222. }
  223. clusterEvents.Publish(events.Event{Type: evType, User: user, Data: json.RawMessage(data)})
  224. return nil
  225. }
  226. // clusterLocalRoots lists the local, non-buffered drives a volume may live on.
  227. func clusterLocalRoots() map[string]string {
  228. roots := map[string]string{}
  229. for _, fsh := range GetAllLoadedFsh() {
  230. if fsh == nil || fsh.Closed || fsh.RequireBuffer || fsh.UUID == "tmp" || fsh.UUID == "cluster" {
  231. continue
  232. }
  233. if arozfs.IsNetworkDrive(fsh.Filesystem) {
  234. continue
  235. }
  236. roots[fsh.UUID] = filepath.Clean(fsh.Path)
  237. }
  238. return roots
  239. }
  240. // clusterBackend adapts the storage service to the clusterfs backend interface.
  241. type clusterBackend struct{ s *storage.Service }
  242. func toInfo(i storage.Info) clusterfs.Info {
  243. return clusterfs.Info{Path: i.Path, IsDir: i.IsDir, Size: i.Size, ModTime: i.ModTime}
  244. }
  245. func (b *clusterBackend) Stat(p string) (clusterfs.Info, error) {
  246. i, err := b.s.Stat(p)
  247. return toInfo(i), err
  248. }
  249. func (b *clusterBackend) List(p string) ([]clusterfs.Info, error) {
  250. items, err := b.s.List(p)
  251. if err != nil {
  252. return nil, err
  253. }
  254. out := make([]clusterfs.Info, 0, len(items))
  255. for _, i := range items {
  256. out = append(out, toInfo(i))
  257. }
  258. return out, nil
  259. }
  260. func (b *clusterBackend) Mkdir(p string, owner string) error { return b.s.Mkdir(p, owner) }
  261. func (b *clusterBackend) Remove(p string, recursive bool) error { return b.s.Remove(p, recursive) }
  262. func (b *clusterBackend) Rename(o string, n string) error { return b.s.Rename(o, n) }
  263. func (b *clusterBackend) OpenRead(p string) (io.ReadCloser, error) { return b.s.OpenRead(p) }
  264. func (b *clusterBackend) Write(p string, r io.Reader, owner string) error {
  265. return b.s.Write(p, r, owner)
  266. }
  267. func (b *clusterBackend) Ready() bool { return b.s.Ready() }
  268. // clusterMountDrive attaches the cluster:/ drive to the base storage pool so
  269. // every user sees it, and clusterUnmountDrive removes it again.
  270. // clusterDriveWanted reports whether cluster:/ should be visible: this node
  271. // is in a cluster and the cluster has at least one volume to keep files in.
  272. // Without either, the drive is left out of File Manager and every other part
  273. // of the system, nightly tasks included.
  274. func clusterDriveWanted() bool {
  275. if clusterManager == nil || clusterStorage == nil || clusterMetadata == nil || !clusterManager.InCluster() {
  276. return false
  277. }
  278. return clusterHasVolume(clusterMetadata.Volumes())
  279. }
  280. // clusterHasVolume reports whether any volume is still part of the cluster.
  281. func clusterHasVolume(vols []metadata.Volume) bool {
  282. for _, v := range vols {
  283. if !v.Removed {
  284. return true
  285. }
  286. }
  287. return false
  288. }
  289. /*
  290. Master node
  291. Nightly maintenance of cluster:/ (expired trash, old version history)
  292. acts on files every member can see, so letting every node run it means
  293. doing the same scan, and the same deletions, once per node. The master
  294. node is the node holding the metadata leader lease, which is also the
  295. node that hands out placement and replication decisions.
  296. */
  297. // clusterIsMasterNode reports whether this node maintains the storage shared
  298. // by the cluster. A node with no cluster is the only node there is, so it is
  299. // its own master.
  300. func clusterIsMasterNode() bool {
  301. if *disable_cluster || clusterManager == nil || clusterMetadata == nil || !clusterManager.InCluster() {
  302. return true
  303. }
  304. return clusterMetadata.IsLeader()
  305. }
  306. // nightlyFshOption returns how nightly maintenance of one file system handler
  307. // should be run: a drive shared by the whole cluster is maintained by the
  308. // master node only, every other drive by the host it belongs to.
  309. func nightlyFshOption(fsh *fs.FileSystemHandler) nightly.TaskOption {
  310. if fsh != nil && fsh.Filesystem == "cluster" {
  311. return nightly.TaskOption{Name: "cluster:/ maintenance", MasterNodeOnly: true}
  312. }
  313. return nightly.TaskOption{}
  314. }
  315. // nightlyShouldMaintainFsh reports whether tonight's maintenance of this file
  316. // system handler belongs to this host.
  317. func nightlyShouldMaintainFsh(fsh *fs.FileSystemHandler) bool {
  318. return nightlyManager.ShouldRun(nightlyFshOption(fsh))
  319. }
  320. // clusterSyncDrive mounts or unmounts cluster:/ to match clusterDriveWanted.
  321. func clusterSyncDrive() {
  322. if clusterDriveWanted() {
  323. clusterMountDrive()
  324. } else {
  325. clusterUnmountDrive()
  326. }
  327. }
  328. func clusterMountDrive() {
  329. clusterMountMu.Lock()
  330. defer clusterMountMu.Unlock()
  331. if clusterStorage == nil || baseStoragePool == nil || !clusterManager.InCluster() {
  332. return
  333. }
  334. if baseStoragePool.ContainDiskID("cluster") {
  335. return
  336. }
  337. handler := &fs.FileSystemHandler{
  338. Name: "Cluster",
  339. UUID: "cluster",
  340. Path: "cluster:/",
  341. Hierarchy: "public",
  342. HierarchyConfig: fs.DefaultEmptyHierarchySpecificConfig,
  343. ReadOnly: false,
  344. RequireBuffer: true,
  345. InitiationTime: time.Now().Unix(),
  346. FileSystemAbstraction: clusterfs.New("cluster", &clusterBackend{s: clusterStorage}),
  347. Filesystem: "cluster",
  348. StartOptions: fs.FileSystemOption{Name: "Cluster", Uuid: "cluster", Path: "cluster:/", Hierarchy: "public", Filesystem: "cluster"},
  349. RuntimePersistenceConfig: fs.RuntimePersistenceConfig{
  350. LocalBufferPath: *tmp_directory,
  351. },
  352. }
  353. if err := baseStoragePool.AttachFsHandler(handler); err != nil {
  354. systemWideLogger.PrintAndLog("Cluster", "Unable to mount cluster drive: "+err.Error(), err)
  355. return
  356. }
  357. systemWideLogger.PrintAndLog("Cluster", "cluster:/ mounted for all users", nil)
  358. go clusterRemoveThumbnailFolders()
  359. }
  360. func clusterUnmountDrive() {
  361. clusterMountMu.Lock()
  362. defer clusterMountMu.Unlock()
  363. if baseStoragePool == nil || !baseStoragePool.ContainDiskID("cluster") {
  364. return
  365. }
  366. baseStoragePool.DetachFsHandler("cluster")
  367. systemWideLogger.PrintAndLog("Cluster", "cluster:/ unmounted", nil)
  368. }
  369. // clusterAccountStore adapts the auth agent and permission handler to the
  370. // identity service's AccountStore interface.
  371. type clusterAccountStore struct{}
  372. func (s *clusterAccountStore) UserExists(username string) bool { return authAgent.UserExists(username) }
  373. func (s *clusterAccountStore) PasswordHash(username string) (string, error) {
  374. return authAgent.GetPasswordHash(username)
  375. }
  376. func (s *clusterAccountStore) SetPasswordHash(username string, hash string) error {
  377. return authAgent.SetPasswordHash(username, hash)
  378. }
  379. func (s *clusterAccountStore) Groups(username string) ([]string, error) {
  380. return authAgent.GetUserGroups(username)
  381. }
  382. func (s *clusterAccountStore) SetGroups(username string, groups []string) error {
  383. return authAgent.SetUserGroups(username, groups)
  384. }
  385. func (s *clusterAccountStore) ListUsers() []string { return authAgent.ListUsers() }
  386. func (s *clusterAccountStore) DeleteUser(username string) error {
  387. return authAgent.UnregisterUser(username)
  388. }
  389. func (s *clusterAccountStore) GroupExists(group string) bool {
  390. return permissionHandler.GroupExists(group)
  391. }
  392. // clusterNotifyPasswordChanged forwards a password change of a replicated
  393. // account to the identity origin. Safe to call when clustering is off.
  394. func clusterNotifyPasswordChanged(username string, passwordHash string) {
  395. if clusterIdentity != nil {
  396. clusterIdentity.NotifyPasswordChanged(username, passwordHash)
  397. }
  398. }
  399. // clusterHealthProvider builds the load snapshot shipped with each heartbeat.
  400. // RAM figures are cached because reading them shells out on some platforms.
  401. type clusterHealthProvider struct {
  402. mu sync.Mutex
  403. ramUsed int64
  404. ramTotal int64
  405. ramSample time.Time
  406. }
  407. func (p *clusterHealthProvider) snapshot() membership.Health {
  408. h := membership.Health{}
  409. if cpu, _, _, _, ready := usageinfo.GetCachedStats(); ready {
  410. h.CPUUsage = cpu
  411. }
  412. p.mu.Lock()
  413. if time.Since(p.ramSample) > time.Minute {
  414. used, total := usageinfo.GetNumericRAMUsage()
  415. if total > 0 {
  416. p.ramUsed, p.ramTotal = used, total
  417. }
  418. p.ramSample = time.Now()
  419. }
  420. h.RAMUsed, h.RAMTotal = p.ramUsed, p.ramTotal
  421. p.mu.Unlock()
  422. if free, total, err := capability.DiskUsage(filepath.Clean(*root_directory)); err == nil {
  423. h.DiskFree, h.DiskTotal = free, total
  424. }
  425. return h
  426. }
  427. func ClusterInit() {
  428. if *disable_cluster {
  429. systemWideLogger.PrintAndLog("Cluster", "Cluster features are disabled by the -disable_cluster flag", nil)
  430. } else {
  431. clusterStartAgent()
  432. }
  433. clusterStartNeighbourhood()
  434. //Let the nightly tasks know which node maintains the shared storage
  435. nightlyManager.SetMasterNodeResolver(clusterIsMasterNode)
  436. }
  437. // clusterStartAgent starts the cluster agent and every cluster service on
  438. // top of it. It is always available unless -disable_cluster is set, so a node
  439. // can create or join a cluster from System Settings regardless of the LAN
  440. // discovery features.
  441. func clusterStartAgent() {
  442. health := &clusterHealthProvider{}
  443. manager, err := membership.NewManager(membership.Option{
  444. NodeID: deviceUUID,
  445. DBFile: filepath.Join("system", "cluster.db"),
  446. KeyFile: filepath.Join("system", "cluster", "node.key"),
  447. Version: build_version + " " + internal_version,
  448. DefaultName: *host_name,
  449. Health: health.snapshot,
  450. })
  451. if err != nil {
  452. systemWideLogger.PrintAndLog("Cluster", "Unable to start cluster agent: "+err.Error(), err)
  453. } else {
  454. clusterManager = manager
  455. //Settings first, then Info; Cluster Jobs is added once the job runtime starts
  456. registerSetting(settingModule{
  457. Name: "Cluster Settings",
  458. Desc: "Create or join a cluster and set up its storage and scheduling",
  459. IconPath: "SystemAO/cluster/img/small_icon.png",
  460. Group: "Cluster",
  461. StartDir: "SystemAO/cluster/cluster.html",
  462. RequireAdmin: true,
  463. })
  464. registerSetting(settingModule{
  465. Name: "Cluster Info",
  466. Desc: "Health of this node, its peers, replication and scheduling",
  467. IconPath: "SystemAO/cluster/img/small_icon.png",
  468. Group: "Cluster",
  469. StartDir: "SystemAO/cluster/clusterinfo.html",
  470. RequireAdmin: true,
  471. })
  472. adminRouter := prout.NewModuleRouter(prout.RouterOption{
  473. ModuleName: "System Setting",
  474. AdminOnly: true,
  475. UserHandler: userHandler,
  476. DeniedHandler: func(w http.ResponseWriter, r *http.Request) {
  477. errorHandlePermissionDenied(w, r)
  478. },
  479. })
  480. registerAdmin := func(pattern string, handler func(http.ResponseWriter, *http.Request)) {
  481. adminRouter.HandleFunc(pattern, handler)
  482. }
  483. clusterManager.RegisterAdminRoutes(registerAdmin)
  484. //Identity service: forward auth to the cluster's identity origin with
  485. //replicated accounts as fallback, plus signed user assertions
  486. idm, err := identity.New(identity.Option{
  487. Membership: clusterManager,
  488. Accounts: &clusterAccountStore{},
  489. })
  490. if err != nil {
  491. systemWideLogger.PrintAndLog("Cluster", "Unable to start cluster identity service: "+err.Error(), err)
  492. } else {
  493. clusterIdentity = idm
  494. authAgent.ForwardAuth = clusterIdentity.ForwardAuth
  495. clusterIdentity.RegisterAdminRoutes(registerAdmin)
  496. }
  497. //Metadata store: replicated namespace index + leader lease
  498. mdm, err := metadata.New(metadata.Option{Membership: clusterManager})
  499. if err != nil {
  500. systemWideLogger.PrintAndLog("Cluster", "Unable to start cluster metadata store: "+err.Error(), err)
  501. } else {
  502. clusterMetadata = mdm
  503. clusterMetadata.RegisterAdminRoutes(registerAdmin)
  504. //Placement scoring shared by jobs, writes and replication
  505. if sch, err := scheduling.New(clusterManager, clusterMetadata); err != nil {
  506. systemWideLogger.PrintAndLog("Cluster", "Unable to start cluster scheduling: "+err.Error(), err)
  507. } else {
  508. clusterScheduling = sch
  509. clusterScheduling.RegisterAdminRoutes(registerAdmin, clusterSchedExplain)
  510. }
  511. //Storage: volumes, chunked transfer and the cluster:/ drive
  512. sto, err := storage.New(storage.Option{
  513. Membership: clusterManager,
  514. Metadata: clusterMetadata,
  515. Scheduler: clusterScheduling,
  516. TmpDir: *tmp_directory,
  517. LocalRoots: clusterLocalRoots,
  518. Thumbnailer: clusterRenderThumbnail,
  519. })
  520. if err != nil {
  521. systemWideLogger.PrintAndLog("Cluster", "Unable to start cluster storage: "+err.Error(), err)
  522. } else {
  523. clusterStorage = sto
  524. clusterStorage.RegisterAdminRoutes(registerAdmin)
  525. //Empty thumbnail folders earlier builds left in cluster:/
  526. nightlyManager.RegisterNightlyTask(func() { clusterRemoveThumbnailFolders() })
  527. prevChange := clusterManager.OnMembershipChange
  528. clusterManager.OnMembershipChange = func(in bool) {
  529. if prevChange != nil {
  530. prevChange(in)
  531. }
  532. clusterSyncDrive()
  533. }
  534. //Show or hide cluster:/ as volumes come and go
  535. clusterMetadata.OnChange = func(kind string, payload []byte) {
  536. if kind == metadata.KindVolume {
  537. go clusterSyncDrive()
  538. }
  539. }
  540. clusterSyncDrive()
  541. //Event bus: file / replica / node events, script hooks, web feed
  542. bus, err := events.New(clusterManager, clusterRunHook)
  543. if err != nil {
  544. systemWideLogger.PrintAndLog("Cluster", "Unable to start cluster event bus: "+err.Error(), err)
  545. } else {
  546. clusterEvents = bus
  547. clusterStorage.OnFileWritten = func(rec metadata.FileRecord) {
  548. bus.Publish(events.Event{Type: "file.created", Path: rec.Path, FileID: rec.ID, User: rec.Owner})
  549. }
  550. clusterStorage.OnFileRemoved = func(rec metadata.FileRecord) {
  551. bus.Publish(events.Event{Type: "file.removed", Path: rec.Path, FileID: rec.ID, User: rec.Owner})
  552. }
  553. clusterStorage.OnFileRenamed = func(oldPath string, rec metadata.FileRecord) {
  554. data, _ := json.Marshal(map[string]string{"from": oldPath})
  555. bus.Publish(events.Event{Type: "file.renamed", Path: rec.Path, FileID: rec.ID, User: rec.Owner, Data: data})
  556. }
  557. clusterStorage.OnReplicaVerified = func(rec metadata.FileRecord, volumeID string) {
  558. data, _ := json.Marshal(map[string]string{"volume": volumeID})
  559. bus.Publish(events.Event{Type: "replica.verified", Path: rec.Path, FileID: rec.ID, Data: data})
  560. }
  561. clusterStorage.OnReplicaStale = func(rec metadata.FileRecord, volumeID string) {
  562. data, _ := json.Marshal(map[string]string{"volume": volumeID})
  563. bus.Publish(events.Event{Type: "replica.stale", Path: rec.Path, FileID: rec.ID, Data: data})
  564. }
  565. clusterStorage.OnDiskFull = func(vol metadata.Volume) {
  566. data, _ := json.Marshal(map[string]interface{}{"volume": vol.ID, "name": vol.Name, "free": vol.Free, "capacity": vol.Capacity})
  567. bus.Publish(events.Event{Type: "node.diskfull", Node: vol.NodeID, Data: data})
  568. }
  569. userRouter := prout.NewModuleRouter(prout.RouterOption{
  570. ModuleName: "System Setting",
  571. AdminOnly: false,
  572. UserHandler: userHandler,
  573. DeniedHandler: func(w http.ResponseWriter, r *http.Request) {
  574. errorHandlePermissionDenied(w, r)
  575. },
  576. })
  577. userRouter.HandleFunc("/system/cluster/events/ws", bus.HandleWebSocket)
  578. userRouter.HandleFunc("/system/cluster/events/hooks", func(w http.ResponseWriter, r *http.Request) {
  579. u, err := userHandler.GetUserInfoFromRequest(w, r)
  580. if err != nil {
  581. errorHandlePermissionDenied(w, r)
  582. return
  583. }
  584. owner := u.Username
  585. if u.IsAdmin() {
  586. owner = ""
  587. }
  588. js, _ := json.Marshal(bus.Hooks(owner))
  589. w.Header().Set("Content-Type", "application/json")
  590. w.Write(js)
  591. })
  592. }
  593. //Job runtime and scheduler
  594. jm, err := jobs.New(jobs.Option{
  595. Membership: clusterManager,
  596. Metadata: clusterMetadata,
  597. Executor: &clusterJobExecutor{},
  598. Scheduler: clusterScheduling,
  599. LocalityBytes: clusterJobLocality,
  600. })
  601. if err != nil {
  602. systemWideLogger.PrintAndLog("Cluster", "Unable to start cluster jobs: "+err.Error(), err)
  603. } else {
  604. clusterJobs = jm
  605. if clusterEvents != nil {
  606. clusterJobs.OnFinished = func(rec jobs.Record) {
  607. evType := "job.completed"
  608. if rec.State.Status != jobs.StatusSucceeded {
  609. evType = "job.failed"
  610. }
  611. data, _ := json.Marshal(map[string]string{"name": rec.Spec.Name, "status": rec.State.Status, "error": rec.State.Error})
  612. clusterEvents.Publish(events.Event{Type: evType, FileID: rec.Spec.ID, User: rec.Spec.Owner, Node: rec.State.Node, Data: data})
  613. }
  614. }
  615. jobRouter := prout.NewModuleRouter(prout.RouterOption{
  616. ModuleName: "Tasks Scheduler",
  617. AdminOnly: false,
  618. UserHandler: userHandler,
  619. DeniedHandler: func(w http.ResponseWriter, r *http.Request) {
  620. errorHandlePermissionDenied(w, r)
  621. },
  622. })
  623. clusterJobs.RegisterUserRoutes(func(pattern string, handler func(http.ResponseWriter, *http.Request)) {
  624. jobRouter.HandleFunc(pattern, handler)
  625. }, clusterJobCaller, clusterJobSubmit)
  626. registerSetting(settingModule{
  627. Name: "Cluster Jobs",
  628. Desc: "Run scripts on any node of the cluster",
  629. IconPath: "SystemAO/cluster/img/small_icon.png",
  630. Group: "Cluster",
  631. StartDir: "SystemAO/cluster/jobs.html",
  632. RequireAdmin: false,
  633. })
  634. }
  635. //AGI "cluster" library
  636. if AGIGateway != nil {
  637. AGIGateway.Option.ClusterProvider = &clusterProvider{}
  638. AGIGateway.ClusterLibRegister()
  639. }
  640. //Replication: keep every file at its policy's copy count
  641. rep, err := replication.New(replication.Option{
  642. Membership: clusterManager,
  643. Metadata: clusterMetadata,
  644. Storage: clusterStorage,
  645. Scheduler: clusterScheduling,
  646. })
  647. if err != nil {
  648. systemWideLogger.PrintAndLog("Cluster", "Unable to start cluster replication: "+err.Error(), err)
  649. } else {
  650. clusterReplication = rep
  651. clusterReplication.RegisterAdminRoutes(registerAdmin)
  652. //Nightly integrity pass over this node's copies (1 GB budget)
  653. nightlyManager.RegisterNightlyTask(func() {
  654. if clusterStorage != nil && clusterManager.InCluster() {
  655. clusterStorage.VerifyLocalCopies(1 << 30)
  656. }
  657. })
  658. }
  659. }
  660. }
  661. }
  662. }
  663. // clusterStartNeighbourhood runs the older mDNS neighbour discovery, which is
  664. // separate from the cluster and governed by -allow_mdns.
  665. func clusterStartNeighbourhood() {
  666. //Only enable neighbourhood scanning on mdns enabled mode
  667. if *allow_mdns && MDNS != nil {
  668. //Start the network discovery
  669. thisDiscoverer := neighbour.NewDiscoverer(MDNS, sysdb)
  670. //Start a scan immediately (in go routine for non blocking)
  671. go func() {
  672. thisDiscoverer.UpdateScan(10)
  673. }()
  674. //Setup the scanning timer
  675. thisDiscoverer.StartScanning(300, 15)
  676. NeighbourDiscoverer = &thisDiscoverer
  677. //Register the settings
  678. registerSetting(settingModule{
  679. Name: "Neighbourhood",
  680. Desc: "Nearby ArOZ Host for Clustering",
  681. IconPath: "SystemAO/cluster/img/small_icon.png",
  682. Group: "Cluster",
  683. StartDir: "SystemAO/cluster/neighbour.html",
  684. RequireAdmin: false,
  685. })
  686. //Register cluster scanning endpoints
  687. router := prout.NewModuleRouter(prout.RouterOption{
  688. ModuleName: "System Setting",
  689. UserHandler: userHandler,
  690. DeniedHandler: func(w http.ResponseWriter, r *http.Request) {
  691. errorHandlePermissionDenied(w, r)
  692. },
  693. })
  694. router.HandleFunc("/system/cluster/scan", NeighbourDiscoverer.HandleScanningRequest)
  695. router.HandleFunc("/system/cluster/record", NeighbourDiscoverer.HandleScanRecord)
  696. router.HandleFunc("/system/cluster/wol", NeighbourDiscoverer.HandleWakeOnLan)
  697. } else {
  698. systemWideLogger.PrintAndLog("Cluster", "MDNS not enabled or startup failed. Skipping Cluster Scanner initiation.", nil)
  699. }
  700. }
  701. // ClusterShutdown stops heartbeats and tunnels before the process exits.
  702. func ClusterShutdown() {
  703. if clusterJobs != nil {
  704. clusterJobs.Close()
  705. }
  706. if clusterEvents != nil {
  707. clusterEvents.Close()
  708. }
  709. if clusterReplication != nil {
  710. clusterReplication.Close()
  711. }
  712. if clusterStorage != nil {
  713. clusterStorage.Close()
  714. }
  715. if clusterMetadata != nil {
  716. clusterMetadata.Close()
  717. }
  718. if clusterIdentity != nil {
  719. clusterIdentity.Close()
  720. }
  721. if clusterManager != nil {
  722. clusterManager.Close()
  723. }
  724. }