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
« prev ^ index » next coverage.py v7.16.2, created at 2026-10-05 13:35 +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 """Job manager."""
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."""
38 dbPath: str
39 """Path of the SQLite database."""
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')
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=}')
49 self._db().insert(self.TABLE, dict(uid=job_uid))
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 )
66 return self._get_job_or_fail(job_uid)
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
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
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'
87 rec = rs[0]
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
96 if user and job_user.uid != user.uid:
97 return None, f'JOB {job_uid}: wrong user {job_user.uid=} {user.uid=}'
99 if state and rec.get('state') != state:
100 return None, f'JOB {job_uid}: wrong state {rec.get("state")=} {state=}'
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, ''
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)
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 {})
133 gws.log.debug(f'JOB {job_uid}: save {rec=}')
134 self._db().update(self.TABLE, rec, job_uid)
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)
143 def run_job(self, job: gws.Job):
144 job_uid = job.uid
146 # atomically mark an 'open' job as 'running'
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)
163 # now it's ours, let's run it
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))
178 return self.get_job(job_uid)
180 def cancel_job(self, job: gws.Job) -> Optional[gws.Job]:
181 return self.update_job(job, state=gws.JobState.cancel)
183 def remove_job(self, job: gws.Job):
184 self._db().delete(self.TABLE, job.uid)
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
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
200 def handle_status_request(self, req, p):
201 job = self.require_job(req, p)
202 return self.job_status_response(job)
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)
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)
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))
232 ##
234 _sqlitex: gws.lib.sqlitex.Object
235 """Database object, created on first use."""
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