Coverage for gws-app/gws/base/job/manager.py: 68%

131 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-24 12:46 +0200

1"""Job manager.""" 

2 

3from typing import Optional 

4import sys 

5 

6import gws 

7import gws.lib.dynimport 

8import gws.lib.jsonx 

9import gws.lib.sqlitex 

10import gws.lib.datetimex as dtx 

11import gws.server.spool 

12 

13 

14class Object(gws.JobManager): 

15 TABLE = 'jobs' 

16 DDL = f""" 

17 CREATE TABLE IF NOT EXISTS {TABLE} ( 

18 uid TEXT NOT NULL PRIMARY KEY, 

19 userUid TEXT DEFAULT '', 

20 userStr TEXT DEFAULT '', 

21 worker TEXT DEFAULT '', 

22 state TEXT DEFAULT '', 

23 error TEXT DEFAULT '', 

24 numSteps INTEGER DEFAULT 0, 

25 step INTEGER DEFAULT 0, 

26 stepName TEXT DEFAULT '', 

27 payload TEXT DEFAULT '', 

28 result TEXT DEFAULT '', 

29 created INTEGER DEFAULT 0, 

30 updated INTEGER DEFAULT 0 

31 ) 

32 """ 

33 

34 dbPath: str 

35 

36 def configure(self): 

37 ver = self.root.specs.version.rpartition('.')[0] 

38 self.dbPath = self.cfg('path', default=f'{gws.c.MISC_DIR}/jobs.{ver}.sqlite') 

39 

40 def create_job(self, worker, user, payload=None): 

41 job_uid = gws.u.random_string(64) 

42 gws.log.debug(f'JOB {job_uid}: creating: {worker=} {user.uid=}') 

43 

44 self._db().insert(self.TABLE, dict(uid=job_uid)) 

45 

46 mod_path = gws.u.require(sys.modules.get(worker.__module__)).__file__ 

47 self._write( 

48 job_uid, 

49 dict( 

50 userUid=user.uid, 

51 userStr=self.root.app.authMgr.serialize_user(user), 

52 worker=f'{mod_path}:{worker.__name__}', 

53 state=gws.JobState.open, 

54 payload=payload, 

55 result={}, 

56 created=gws.u.stime(), 

57 updated=gws.u.stime(), 

58 ), 

59 ) 

60 

61 return self._get_job_or_fail(job_uid) 

62 

63 def get_job(self, job_uid: str, user=None, state=None): 

64 job, msg = self._get_job(job_uid, user, state) 

65 if not job: 

66 gws.log.error(msg) 

67 return job 

68 

69 def _get_job_or_fail(self, job_uid: str, user=None, state=None): 

70 job, msg = self._get_job(job_uid, user, state) 

71 if not job: 

72 raise gws.Error(msg) 

73 return job 

74 

75 def _get_job(self, job_uid: str, user=None, state=None): 

76 rs = self._db().select(f'SELECT * FROM {self.TABLE} WHERE uid=:uid', uid=job_uid) 

77 if not rs: 

78 return None, f'JOB {job_uid}: not found' 

79 

80 rec = rs[0] 

81 

82 job_user = None 

83 us = rec.get('userStr') 

84 if us: 

85 job_user = self.root.app.authMgr.unserialize_user(us) 

86 if not job_user: 

87 job_user = self.root.app.authMgr.guestUser 

88 

89 if user and job_user.uid != user.uid: 

90 return None, f'JOB {job_uid}: wrong user {job_user.uid=} {user.uid=}' 

91 

92 if state and rec.get('state') != state: 

93 return None, f'JOB {job_uid}: wrong state {rec.get("state")=} {state=}' 

94 

95 job = gws.Job( 

96 uid=rec['uid'], 

97 user=job_user, 

98 worker=rec['worker'], 

99 state=rec['state'], 

100 error=rec['error'], 

101 numSteps=rec['numSteps'] or 0, 

102 step=rec['step'] or 0, 

103 stepName=rec['stepName'] or '', 

104 payload=gws.lib.jsonx.from_string(rec['payload'] or '{}'), 

105 result=gws.lib.jsonx.from_string(rec['result'] or '{}'), 

106 timeCreated=dtx.from_timestamp(rec['created'] or 0), 

107 timeUpdated=dtx.from_timestamp(rec['updated'] or 0), 

108 ) 

109 return job, '' 

110 

111 def update_job(self, job, **kwargs): 

112 job = self.get_job(job.uid) 

113 if not job: 

114 return 

115 self._write(job.uid, kwargs) 

116 return self._get_job_or_fail(job.uid) 

117 

118 def _write(self, job_uid, rec): 

119 rec['updated'] = gws.u.stime() 

120 if 'payload' in rec: 

121 rec['payload'] = gws.lib.jsonx.to_string(rec['payload'] or {}) 

122 if 'result' in rec: 

123 rec['result'] = gws.lib.jsonx.to_string(rec['result'] or {}) 

124 

125 gws.log.debug(f'JOB {job_uid}: save {rec=}') 

126 self._db().update(self.TABLE, rec, job_uid) 

127 

128 def schedule_job(self, job: gws.Job): 

129 job = self._get_job_or_fail(job.uid, state=gws.JobState.open) 

130 if gws.server.spool.is_active(): 

131 gws.server.spool.add(job) 

132 return job 

133 return self.run_job(job) 

134 

135 def run_job(self, job: gws.Job): 

136 job_uid = job.uid 

137 

138 # atomically mark an 'open' job as 'running' 

139 

140 tmp = gws.u.random_string(64) 

141 self._db().execute( 

142 f'UPDATE {self.TABLE} SET state=:tmp WHERE uid=:uid AND state=:state', 

143 uid=job_uid, 

144 tmp=tmp, 

145 state=gws.JobState.open, 

146 ) 

147 job = self._get_job_or_fail(job_uid, state=tmp) 

148 self._db().execute( 

149 f'UPDATE {self.TABLE} SET state=:s WHERE uid=:uid', 

150 uid=job.uid, 

151 s=gws.JobState.running, 

152 ) 

153 job = self._get_job_or_fail(job_uid, state=gws.JobState.running) 

154 

155 # now it's ours, let's run it 

156 

157 try: 

158 mod_path, _, fn_name = job.worker.rpartition(':') 

159 mod = gws.lib.dynimport.import_from_path(mod_path) 

160 worker_cls = getattr(mod, fn_name) 

161 worker_cls.run(self.root, job) 

162 except gws.JobTerminated as exc: 

163 gws.log.error(f'JOB {job_uid}: JobTerminated: {exc.args[0]!r}') 

164 self.update_job(job, state=gws.JobState.error) 

165 except Exception as exc: 

166 gws.log.error(f'JOB {job_uid}: FAILED {exc=}') 

167 gws.log.exception() 

168 self.update_job(job, state=gws.JobState.error, error=repr(exc)) 

169 

170 return self.get_job(job_uid) 

171 

172 def cancel_job(self, job: gws.Job) -> Optional[gws.Job]: 

173 return self.update_job(job, state=gws.JobState.cancel) 

174 

175 def remove_job(self, job: gws.Job): 

176 self._db().delete(self.TABLE, job.uid) 

177 

178 def require_job(self, req, p): 

179 job = self.root.app.jobMgr.get_job(p.jobUid, req.user) 

180 if not job: 

181 raise gws.NotFoundError(f'JOB {p.jobUid}: not found') 

182 return job 

183 

184 def require_result(self, req, p): 

185 job = self.require_job(req, p) 

186 if job.state != gws.JobState.complete: 

187 raise gws.NotFoundError(f'JOB {p.jobUid}: wrong state {job.state!r}') 

188 if not job.result: 

189 raise gws.NotFoundError(f'JOB {p.jobUid}: no result') 

190 return job.result 

191 

192 def handle_status_request(self, req, p): 

193 job = self.require_job(req, p) 

194 return self.job_status_response(job) 

195 

196 def handle_cancel_request(self, req, p): 

197 job = self.require_job(req, p) 

198 job = self.cancel_job(job) 

199 if not job: 

200 raise gws.NotFoundError(f'JOB {p.jobUid}: not found') 

201 return self.job_status_response(job) 

202 

203 def job_status_response(self, job, **kwargs): 

204 d = dict( 

205 jobUid=job.uid, 

206 state=job.state, 

207 stepName=job.stepName or '', 

208 progress=self._get_progress(job), 

209 output={}, 

210 ) 

211 d.update(kwargs) 

212 return gws.JobStatusResponse(d) 

213 

214 def _get_progress(self, job): 

215 if job.state == gws.JobState.complete: 

216 return 100 

217 if job.state != gws.JobState.running: 

218 return 0 

219 if not job.numSteps: 

220 return 0 

221 return int(min(100.0, job.step * 100.0 / job.numSteps)) 

222 

223 ## 

224 

225 _sqlitex: gws.lib.sqlitex.Object 

226 

227 def _db(self): 

228 if getattr(self, '_sqlitex', None) is None: 

229 self._sqlitex = gws.lib.sqlitex.Object(self.dbPath, self.DDL) 

230 return self._sqlitex