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
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-24 12:46 +0200
1"""Job manager."""
3from typing import Optional
4import sys
6import gws
7import gws.lib.dynimport
8import gws.lib.jsonx
9import gws.lib.sqlitex
10import gws.lib.datetimex as dtx
11import gws.server.spool
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 """
34 dbPath: str
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')
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=}')
44 self._db().insert(self.TABLE, dict(uid=job_uid))
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 )
61 return self._get_job_or_fail(job_uid)
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
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
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'
80 rec = rs[0]
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
89 if user and job_user.uid != user.uid:
90 return None, f'JOB {job_uid}: wrong user {job_user.uid=} {user.uid=}'
92 if state and rec.get('state') != state:
93 return None, f'JOB {job_uid}: wrong state {rec.get("state")=} {state=}'
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, ''
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)
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 {})
125 gws.log.debug(f'JOB {job_uid}: save {rec=}')
126 self._db().update(self.TABLE, rec, job_uid)
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)
135 def run_job(self, job: gws.Job):
136 job_uid = job.uid
138 # atomically mark an 'open' job as 'running'
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)
155 # now it's ours, let's run it
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))
170 return self.get_job(job_uid)
172 def cancel_job(self, job: gws.Job) -> Optional[gws.Job]:
173 return self.update_job(job, state=gws.JobState.cancel)
175 def remove_job(self, job: gws.Job):
176 self._db().delete(self.TABLE, job.uid)
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
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
192 def handle_status_request(self, req, p):
193 job = self.require_job(req, p)
194 return self.job_status_response(job)
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)
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)
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))
223 ##
225 _sqlitex: gws.lib.sqlitex.Object
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