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

133 statements  

« prev     ^ index     » next       coverage.py v7.16.2, created at 2026-10-05 13:35 +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 """Job manager.""" 

16 

17 TABLE = 'jobs' 

18 """Name of the jobs table.""" 

19 DDL = f""" 

20 CREATE TABLE IF NOT EXISTS {TABLE} ( 

21 uid TEXT NOT NULL PRIMARY KEY, 

22 userUid TEXT DEFAULT '', 

23 userStr TEXT DEFAULT '', 

24 worker TEXT DEFAULT '', 

25 state TEXT DEFAULT '', 

26 error TEXT DEFAULT '', 

27 numSteps INTEGER DEFAULT 0, 

28 step INTEGER DEFAULT 0, 

29 stepName TEXT DEFAULT '', 

30 payload TEXT DEFAULT '', 

31 result TEXT DEFAULT '', 

32 created INTEGER DEFAULT 0, 

33 updated INTEGER DEFAULT 0 

34 ) 

35 """ 

36 """DDL statement for the jobs table.""" 

37 

38 dbPath: str 

39 """Path of the SQLite database.""" 

40 

41 def configure(self): 

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

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

44 

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

46 job_uid = gws.u.random_string(64) 

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

48 

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

50 

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

52 self._write( 

53 job_uid, 

54 dict( 

55 userUid=user.uid, 

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

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

58 state=gws.JobState.open, 

59 payload=payload, 

60 result={}, 

61 created=gws.u.stime(), 

62 updated=gws.u.stime(), 

63 ), 

64 ) 

65 

66 return self._get_job_or_fail(job_uid) 

67 

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

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

70 if not job: 

71 gws.log.error(msg) 

72 return job 

73 

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

75 """Return a job like ``get_job``, but raise ``gws.Error`` if it is not found.""" 

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

77 if not job: 

78 raise gws.Error(msg) 

79 return job 

80 

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

82 """Read a job, checking the user and state; return the job and an error message.""" 

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

84 if not rs: 

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

86 

87 rec = rs[0] 

88 

89 job_user = None 

90 us = rec.get('userStr') 

91 if us: 

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

93 if not job_user: 

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

95 

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

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

98 

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

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

101 

102 job = gws.Job( 

103 uid=rec['uid'], 

104 user=job_user, 

105 worker=rec['worker'], 

106 state=rec['state'], 

107 error=rec['error'], 

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

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

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

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

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

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

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

115 ) 

116 return job, '' 

117 

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

119 job = self.get_job(job.uid) 

120 if not job: 

121 return 

122 self._write(job.uid, kwargs) 

123 return self._get_job_or_fail(job.uid) 

124 

125 def _write(self, job_uid, rec): 

126 """Update a job record, serializing the payload and the result.""" 

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

128 if 'payload' in rec: 

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

130 if 'result' in rec: 

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

132 

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

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

135 

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

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

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

139 gws.server.spool.add(job) 

140 return job 

141 return self.run_job(job) 

142 

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

144 job_uid = job.uid 

145 

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

147 

148 tmp = gws.u.random_string(64) 

149 self._db().execute( 

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

151 uid=job_uid, 

152 tmp=tmp, 

153 state=gws.JobState.open, 

154 ) 

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

156 self._db().execute( 

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

158 uid=job.uid, 

159 s=gws.JobState.running, 

160 ) 

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

162 

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

164 

165 try: 

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

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

168 worker_cls = getattr(mod, fn_name) 

169 worker_cls.run(self.root, job) 

170 except gws.JobTerminated as exc: 

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

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

173 except Exception as exc: 

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

175 gws.log.exception() 

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

177 

178 return self.get_job(job_uid) 

179 

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

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

182 

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

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

185 

186 def require_job(self, req, p): 

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

188 if not job: 

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

190 return job 

191 

192 def require_result(self, req, p): 

193 job = self.require_job(req, p) 

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

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

196 if not job.result: 

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

198 return job.result 

199 

200 def handle_status_request(self, req, p): 

201 job = self.require_job(req, p) 

202 return self.job_status_response(job) 

203 

204 def handle_cancel_request(self, req, p): 

205 job = self.require_job(req, p) 

206 job = self.cancel_job(job) 

207 if not job: 

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

209 return self.job_status_response(job) 

210 

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

212 d = dict( 

213 jobUid=job.uid, 

214 state=job.state, 

215 stepName=job.stepName or '', 

216 progress=self._get_progress(job), 

217 output={}, 

218 ) 

219 d.update(kwargs) 

220 return gws.JobStatusResponse(d) 

221 

222 def _get_progress(self, job): 

223 """Return the progress of a job in percent.""" 

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

225 return 100 

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

227 return 0 

228 if not job.numSteps: 

229 return 0 

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

231 

232 ## 

233 

234 _sqlitex: gws.lib.sqlitex.Object 

235 """Database object, created on first use.""" 

236 

237 def _db(self): 

238 """Return the database object, creating it if needed.""" 

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

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

241 return self._sqlitex