scheduler.go 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405
  1. package scheduler
  2. import (
  3. "encoding/json"
  4. "errors"
  5. "os"
  6. "path/filepath"
  7. "time"
  8. "imuslab.com/arozos/mod/agi"
  9. "imuslab.com/arozos/mod/filesystem"
  10. "imuslab.com/arozos/mod/info/logger"
  11. "imuslab.com/arozos/mod/user"
  12. "imuslab.com/arozos/mod/utils"
  13. )
  14. // WebRootFshID is stored in Job.FshID for scripts that live inside the webapp
  15. // folder (./web/<AppName>/…) rather than in a user's virtual filesystem.
  16. // At execution time the scheduler resolves the script via the OS filesystem
  17. // instead of the user's storage pool.
  18. const WebRootFshID = "__webroot__"
  19. // WebRootBase is the directory from which app-relative script paths are resolved.
  20. const WebRootBase = "./web"
  21. /*
  22. ArozOS System Scheduler
  23. author: tobychui
  24. This module provide scheduling executable feature for ArozOS
  25. Some feature was migrated from the v1.113 aecron module
  26. */
  27. type Job struct {
  28. Name string //The name of this job
  29. Creator string //The creator of this job. When execute, this user permission will be used
  30. Description string //Job description, can be empty
  31. ExecutionInterval int64 //Execuation interval in seconds
  32. BaseTime int64 //Exeuction basetime. The next interval is calculated using (current time - base time ) % execution interval
  33. FshID string //The target FSH ID that this script file is stored
  34. ScriptVpath string //The agi script file being called, require Vpath
  35. AppName string //The webapp/module that registered this job (empty for manual jobs)
  36. lastExecutionTime int64 //Last time this job being executed
  37. lastExecutionOutput string //The output of last execution
  38. }
  39. type ScheudlerOption struct {
  40. UserHandler *user.UserHandler
  41. Gateway *agi.Gateway
  42. Logger *logger.Logger
  43. CronFile string //The location of the cronfile which store the jobs registry in file format
  44. }
  45. type Scheduler struct {
  46. jobs []*Job
  47. options *ScheudlerOption
  48. ticker chan bool
  49. }
  50. var ()
  51. func NewScheduler(option *ScheudlerOption) (*Scheduler, error) {
  52. if !utils.FileExists(option.CronFile) {
  53. //Cronfile not exists. Create it
  54. emptyJobList := []*Job{}
  55. ls, _ := json.Marshal(emptyJobList)
  56. err := os.WriteFile(option.CronFile, ls, 0755)
  57. if err != nil {
  58. return nil, err
  59. }
  60. }
  61. //Load previous jobs from file
  62. jobs, err := loadJobsFromFile(option.CronFile)
  63. if err != nil {
  64. return nil, err
  65. }
  66. //Create the ArOZ Emulated Crontask
  67. thisScheduler := Scheduler{
  68. jobs: jobs,
  69. options: option,
  70. }
  71. option.Logger.PrintAndLog("Scheduler", "Scheduler started", nil)
  72. //Start the cronjob at 1 minute ticker interval
  73. go func() {
  74. //Delay start: Wait until seconds = 0
  75. for time.Now().Unix()%60 > 0 {
  76. time.Sleep(500 * time.Millisecond)
  77. }
  78. stopChannel := thisScheduler.createTicker(1 * time.Minute)
  79. thisScheduler.ticker = stopChannel
  80. option.Logger.PrintAndLog("Scheduler", "ArozOS System Scheduler Started", nil)
  81. }()
  82. //Return the crontask
  83. return &thisScheduler, nil
  84. }
  85. func (a *Scheduler) createTicker(duration time.Duration) chan bool {
  86. ticker := time.NewTicker(duration)
  87. stop := make(chan bool, 1)
  88. go func() {
  89. defer logger.PrintAndLog("Scheduler", "Scheduler Stopped", nil)
  90. for {
  91. select {
  92. case <-ticker.C:
  93. //Run jobs
  94. for _, thisJob := range a.jobs {
  95. if (time.Now().Unix()-thisJob.BaseTime)%thisJob.ExecutionInterval == 0 {
  96. a.executeJob(thisJob)
  97. }
  98. }
  99. case <-stop:
  100. return
  101. }
  102. }
  103. }()
  104. return stop
  105. }
  106. // executeJob resolves the script path and runs it in a goroutine.
  107. // It supports two modes:
  108. //
  109. // 1. App-root scripts (FshID == WebRootFshID):
  110. // ScriptVpath is relative to ./web/, e.g. "MyApp/cron.agi".
  111. // The file is read directly from the OS; no user FSH is involved.
  112. //
  113. // 2. User virtual-path scripts (any other FshID):
  114. // ScriptVpath is a virtual path like "user:/path/to/script.agi".
  115. // The file is resolved through the creator's storage pool.
  116. func (a *Scheduler) executeJob(thisJob *Job) {
  117. targetUser, err := a.options.UserHandler.GetUserInfoFromUsername(thisJob.Creator)
  118. if err != nil {
  119. a.cronlogError("User "+thisJob.Creator+" no longer exists", err)
  120. return
  121. }
  122. cloned := *thisJob
  123. if thisJob.FshID == WebRootFshID {
  124. // ── App-root script ───────────────────────────────────────────────
  125. rpath := filepath.Join(WebRootBase, filepath.FromSlash(thisJob.ScriptVpath))
  126. if _, statErr := os.Stat(rpath); os.IsNotExist(statErr) {
  127. a.cronlog("Removing job " + thisJob.Name + " by " + thisJob.Creator + " as app script no longer exists: " + rpath)
  128. a.RemoveJobFromScheduleList(thisJob.Name)
  129. a.saveJobsToCronFile()
  130. return
  131. }
  132. ext := filepath.Ext(rpath)
  133. if ext != ".js" && ext != ".agi" {
  134. a.cronlogError("Unsupported app script extension: "+ext, errors.New("unsupported extension"))
  135. return
  136. }
  137. go func(job Job, realPath string, u *user.User) {
  138. job.lastExecutionTime = time.Now().Unix()
  139. // fsh == nil --> ExecuteAGIScriptAsUser reads via os.ReadFile
  140. execID, resp, execErr := a.options.Gateway.ExecuteAGIScriptAsUser(nil, realPath, u, nil, nil)
  141. if execErr != nil {
  142. a.cronlogError("["+execID+"] "+job.Name+" execution error: "+execErr.Error(), execErr)
  143. job.lastExecutionOutput = execErr.Error()
  144. } else {
  145. a.cronlog("[" + execID + "] " + job.Name + " executed: " + resp)
  146. job.lastExecutionOutput = resp
  147. }
  148. }(cloned, rpath, targetUser)
  149. return
  150. }
  151. // ── User virtual-path script ──────────────────────────────────────────
  152. fsh, err := targetUser.GetFileSystemHandlerFromVirtualPath(thisJob.ScriptVpath)
  153. if err != nil {
  154. a.cronlogError("Unable to resolve vpath for job: "+thisJob.Name+" (user: "+thisJob.Creator+")", err)
  155. return
  156. }
  157. rpath, err := fsh.FileSystemAbstraction.VirtualPathToRealPath(thisJob.ScriptVpath, targetUser.Username)
  158. if err != nil {
  159. a.cronlogError("Unable to get real path for job: "+thisJob.Name, err)
  160. return
  161. }
  162. if !fsh.FileSystemAbstraction.FileExists(rpath) {
  163. a.cronlog("Removing job " + thisJob.Name + " by " + thisJob.Creator + " as script no longer exists")
  164. a.RemoveJobFromScheduleList(thisJob.Name)
  165. a.saveJobsToCronFile()
  166. return
  167. }
  168. ext := filepath.Ext(rpath)
  169. if ext != ".js" && ext != ".agi" {
  170. a.cronlogError("Unsupported script extension: "+ext, errors.New("unsupported extension"))
  171. return
  172. }
  173. go func(job Job, f *filesystem.FileSystemHandler, realPath string, u *user.User) {
  174. job.lastExecutionTime = time.Now().Unix()
  175. execID, resp, execErr := a.options.Gateway.ExecuteAGIScriptAsUser(f, realPath, u, nil, nil)
  176. if execErr != nil {
  177. a.cronlogError("["+execID+"] "+job.Name+" execution error: "+execErr.Error(), execErr)
  178. job.lastExecutionOutput = execErr.Error()
  179. } else {
  180. a.cronlog("[" + execID + "] " + job.Name + " executed: " + resp)
  181. job.lastExecutionOutput = resp
  182. }
  183. }(cloned, fsh, rpath, targetUser)
  184. }
  185. func (a *Scheduler) Close() {
  186. if a.ticker != nil {
  187. //Stop the ticker
  188. a.ticker <- true
  189. }
  190. }
  191. // Add an job object to system scheduler
  192. func (a *Scheduler) AddJobToScheduler(job *Job) error {
  193. a.jobs = append(a.jobs, job)
  194. return nil
  195. }
  196. func (a *Scheduler) GetScheduledJobByName(name string) *Job {
  197. for _, thisJob := range a.jobs {
  198. if thisJob.Name == name {
  199. return thisJob
  200. }
  201. }
  202. return nil
  203. }
  204. func (a *Scheduler) RemoveJobFromScheduleList(taskName string) {
  205. newJobSlice := []*Job{}
  206. for _, j := range a.jobs {
  207. if j.Name != taskName {
  208. thisJob := j
  209. newJobSlice = append(newJobSlice, thisJob)
  210. }
  211. }
  212. a.jobs = newJobSlice
  213. }
  214. func (a *Scheduler) JobExists(name string) bool {
  215. targetJob := a.GetScheduledJobByName(name)
  216. if targetJob == nil {
  217. return false
  218. } else {
  219. return true
  220. }
  221. }
  222. // GetJobsByApp returns all jobs registered by the given app name
  223. func (a *Scheduler) GetJobsByApp(appName string) []*Job {
  224. result := []*Job{}
  225. for _, j := range a.jobs {
  226. if j.AppName == appName {
  227. result = append(result, j)
  228. }
  229. }
  230. return result
  231. }
  232. // RemoveJobsByApp removes all scheduler jobs associated with a given app name and saves to disk
  233. func (a *Scheduler) RemoveJobsByApp(appName string) {
  234. newJobSlice := []*Job{}
  235. for _, j := range a.jobs {
  236. if j.AppName != appName {
  237. newJobSlice = append(newJobSlice, j)
  238. }
  239. }
  240. a.jobs = newJobSlice
  241. a.saveJobsToCronFile()
  242. }
  243. // AppJobExists checks whether a job with the given name was registered by the given app and creator
  244. func (a *Scheduler) AppJobExists(appName, creator, taskName string) bool {
  245. for _, j := range a.jobs {
  246. if j.AppName == appName && j.Creator == creator && j.Name == taskName {
  247. return true
  248. }
  249. }
  250. return false
  251. }
  252. // RegisterJobFromAGI creates and saves a new job on behalf of a user/app from AGI scripts.
  253. //
  254. // When appName is non-empty and scriptVpath does not contain ":" (i.e. it is not
  255. // a user virtual path), the script is treated as app-root-relative:
  256. //
  257. // scriptVpath = "cron.agi" --> stored as "appName/cron.agi", FshID = WebRootFshID
  258. //
  259. // Otherwise scriptVpath must be a full virtual path (e.g. "user:/path/to/script.agi")
  260. // and is resolved through the creator's storage pool as usual.
  261. func (a *Scheduler) RegisterJobFromAGI(creator, appName, taskName, scriptVpath, description string, interval, baseTime int64) error {
  262. // Validate name uniqueness
  263. for _, j := range a.jobs {
  264. if j.Name == taskName {
  265. return errors.New("task name already occupied: " + taskName)
  266. }
  267. }
  268. var fshID string
  269. var storedVpath string
  270. isAppScript := appName != "" && !containsVpathSeparator(scriptVpath)
  271. if isAppScript {
  272. // App-root script: resolve relative to ./web/<appName>/
  273. relPath := appName + "/" + scriptVpath
  274. realPath := filepath.Join(WebRootBase, filepath.FromSlash(relPath))
  275. if _, err := os.Stat(realPath); os.IsNotExist(err) {
  276. return errors.New("app script not found: " + realPath)
  277. }
  278. fshID = WebRootFshID
  279. storedVpath = relPath
  280. } else {
  281. // User virtual-path script
  282. targetUser, err := a.options.UserHandler.GetUserInfoFromUsername(creator)
  283. if err != nil {
  284. return err
  285. }
  286. fsh, err := targetUser.GetFileSystemHandlerFromVirtualPath(scriptVpath)
  287. if err != nil {
  288. return err
  289. }
  290. fshID = fsh.UUID
  291. storedVpath = scriptVpath
  292. }
  293. newJob := &Job{
  294. Name: taskName,
  295. Creator: creator,
  296. AppName: appName,
  297. Description: description,
  298. ExecutionInterval: interval,
  299. BaseTime: alignBaseTime(baseTime),
  300. ScriptVpath: storedVpath,
  301. FshID: fshID,
  302. }
  303. a.jobs = append(a.jobs, newJob)
  304. return a.saveJobsToCronFile()
  305. }
  306. // alignBaseTime floors t to the nearest whole minute so that the scheduler's
  307. // per-minute ticker (which fires at unix timestamps divisible by 60) can
  308. // satisfy the condition (ticker - baseTime) % interval == 0.
  309. // Without this alignment, any job with interval < 86400 that was registered
  310. // at a non-minute boundary will never fire.
  311. func alignBaseTime(t int64) int64 {
  312. return (t / 60) * 60
  313. }
  314. // containsVpathSeparator returns true when s contains the ":" that marks a
  315. // virtual-path root (e.g. "user:/…" or "tmp:/…").
  316. func containsVpathSeparator(s string) bool {
  317. for _, c := range s {
  318. if c == ':' {
  319. return true
  320. }
  321. }
  322. return false
  323. }
  324. // UnregisterJobFromAGI removes a job by task name for a given creator (or admin)
  325. func (a *Scheduler) UnregisterJobFromAGI(creator, taskName string) error {
  326. targetJob := a.GetScheduledJobByName(taskName)
  327. if targetJob == nil {
  328. return errors.New("job not found: " + taskName)
  329. }
  330. if targetJob.Creator != creator {
  331. // Check if creator is admin (requires a UserHandler; deny if unavailable)
  332. if a.options.UserHandler == nil {
  333. return errors.New("permission denied")
  334. }
  335. targetUser, err := a.options.UserHandler.GetUserInfoFromUsername(creator)
  336. if err != nil || !targetUser.IsAdmin() {
  337. return errors.New("permission denied")
  338. }
  339. }
  340. a.RemoveJobFromScheduleList(taskName)
  341. return a.saveJobsToCronFile()
  342. }
  343. //Write the output to log file. Default to ./system/aecron/{date}.log
  344. /*
  345. func cronlog(message string) {
  346. currentTime := time.Now()
  347. timestamp := currentTime.Format("2006-01-02 15:04:05")
  348. message = timestamp + " " + message
  349. currentLogFile := filepath.ToSlash(filepath.Clean(logFolder)) + "/" + time.Now().Format("2006-02-01") + ".log"
  350. f, err := os.OpenFile(currentLogFile, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644)
  351. if err != nil {
  352. //Unable to write to file. Log to STDOUT instead
  353. logger.PrintAndLog("Scheduler", fmt.Sprint(message), nil)
  354. return
  355. }
  356. if _, err := f.WriteString(message + "\n"); err != nil {
  357. logger.PrintAndLog("Scheduler", fmt.Sprint(message), nil)
  358. return
  359. }
  360. defer f.Close()
  361. }
  362. */