feat(chalice): fixed jobs execution
This commit is contained in:
parent
16e7be5e99
commit
048a9767ac
6 changed files with 61 additions and 114 deletions
1
.github/workflows/api-ee.yaml
vendored
1
.github/workflows/api-ee.yaml
vendored
|
|
@ -10,6 +10,7 @@ on:
|
|||
branches:
|
||||
- dev
|
||||
- api-*
|
||||
- v1.11.0-patch
|
||||
paths:
|
||||
- "ee/api/**"
|
||||
- "api/**"
|
||||
|
|
|
|||
1
.github/workflows/api.yaml
vendored
1
.github/workflows/api.yaml
vendored
|
|
@ -10,6 +10,7 @@ on:
|
|||
branches:
|
||||
- dev
|
||||
- api-*
|
||||
- v1.11.0-patch
|
||||
paths:
|
||||
- "api/**"
|
||||
- "!api/.gitignore"
|
||||
|
|
|
|||
|
|
@ -17,11 +17,9 @@ class JobStatus:
|
|||
def get(job_id):
|
||||
with pg_client.PostgresClient() as cur:
|
||||
query = cur.mogrify(
|
||||
"""\
|
||||
SELECT
|
||||
*
|
||||
FROM public.jobs
|
||||
WHERE job_id = %(job_id)s;""",
|
||||
"""SELECT *
|
||||
FROM public.jobs
|
||||
WHERE job_id = %(job_id)s;""",
|
||||
{"job_id": job_id}
|
||||
)
|
||||
cur.execute(query=query)
|
||||
|
|
@ -37,11 +35,9 @@ def get(job_id):
|
|||
def get_all(project_id):
|
||||
with pg_client.PostgresClient() as cur:
|
||||
query = cur.mogrify(
|
||||
"""\
|
||||
SELECT
|
||||
*
|
||||
FROM public.jobs
|
||||
WHERE project_id = %(project_id)s;""",
|
||||
"""SELECT *
|
||||
FROM public.jobs
|
||||
WHERE project_id = %(project_id)s;""",
|
||||
{"project_id": project_id}
|
||||
)
|
||||
cur.execute(query=query)
|
||||
|
|
@ -59,15 +55,10 @@ def create(project_id, data):
|
|||
**data
|
||||
}
|
||||
|
||||
query = cur.mogrify("""\
|
||||
INSERT INTO public.jobs(
|
||||
project_id, description, status, action,
|
||||
reference_id, start_at
|
||||
)
|
||||
VALUES (
|
||||
%(project_id)s, %(description)s, %(status)s, %(action)s,
|
||||
%(reference_id)s, %(start_at)s
|
||||
) RETURNING *;""", job)
|
||||
query = cur.mogrify(
|
||||
"""INSERT INTO public.jobs(project_id, description, status, action,reference_id, start_at)
|
||||
VALUES (%(project_id)s, %(description)s, %(status)s, %(action)s,%(reference_id)s, %(start_at)s)
|
||||
RETURNING *;""", job)
|
||||
|
||||
cur.execute(query=query)
|
||||
|
||||
|
|
@ -90,14 +81,13 @@ def update(job_id, job):
|
|||
**job
|
||||
}
|
||||
|
||||
query = cur.mogrify("""\
|
||||
UPDATE public.jobs
|
||||
SET
|
||||
updated_at = timezone('utc'::text, now()),
|
||||
status = %(status)s,
|
||||
errors = %(errors)s
|
||||
WHERE
|
||||
job_id = %(job_id)s RETURNING *;""", job_data)
|
||||
query = cur.mogrify(
|
||||
"""UPDATE public.jobs
|
||||
SET updated_at = timezone('utc'::text, now()),
|
||||
status = %(status)s,
|
||||
errors = %(errors)s
|
||||
WHERE job_id = %(job_id)s
|
||||
RETURNING *;""", job_data)
|
||||
|
||||
cur.execute(query=query)
|
||||
|
||||
|
|
@ -116,11 +106,12 @@ def format_datetime(r):
|
|||
def get_scheduled_jobs():
|
||||
with pg_client.PostgresClient() as cur:
|
||||
query = cur.mogrify(
|
||||
"""\
|
||||
SELECT * FROM public.jobs
|
||||
WHERE status = %(status)s AND start_at <= (now() at time zone 'utc');""",
|
||||
{"status": JobStatus.SCHEDULED}
|
||||
)
|
||||
"""SELECT * FROM public.jobs
|
||||
WHERE status = %(status)s AND start_at <= (now() at time zone 'utc');""",
|
||||
{"status": JobStatus.SCHEDULED})
|
||||
print("------------------")
|
||||
print(query)
|
||||
print("------------------")
|
||||
cur.execute(query=query)
|
||||
data = cur.fetchall()
|
||||
for record in data:
|
||||
|
|
@ -131,43 +122,28 @@ def get_scheduled_jobs():
|
|||
|
||||
def execute_jobs():
|
||||
jobs = get_scheduled_jobs()
|
||||
if len(jobs) == 0:
|
||||
# No jobs to execute
|
||||
return
|
||||
|
||||
for job in jobs:
|
||||
print(f"job can be executed {job['id']}")
|
||||
print(f"Executing jobId:{job['jobId']}")
|
||||
try:
|
||||
if job["action"] == Actions.DELETE_USER_DATA:
|
||||
session_ids = sessions.get_session_ids_by_user_ids(
|
||||
project_id=job["projectId"], user_ids=job["referenceId"]
|
||||
)
|
||||
|
||||
sessions.delete_sessions_by_session_ids(session_ids)
|
||||
sessions_mobs.delete_mobs(session_ids=session_ids, project_id=job["projectId"])
|
||||
session_ids = sessions.get_session_ids_by_user_ids(project_id=job["projectId"],
|
||||
user_ids=[job["referenceId"]])
|
||||
if len(session_ids) > 0:
|
||||
print(f"Deleting {len(session_ids)} sessions")
|
||||
sessions.delete_sessions_by_session_ids(session_ids)
|
||||
sessions_mobs.delete_mobs(session_ids=session_ids, project_id=job["projectId"])
|
||||
else:
|
||||
raise Exception(f"The action {job['action']} not supported.")
|
||||
raise Exception(f"The action '{job['action']}' not supported.")
|
||||
|
||||
job["status"] = JobStatus.COMPLETED
|
||||
print(f"job completed {job['id']}")
|
||||
print(f"Job completed {job['jobId']}")
|
||||
except Exception as e:
|
||||
print("-----")
|
||||
print(e)
|
||||
print("-----")
|
||||
job["status"] = JobStatus.FAILED
|
||||
job["error"] = str(e)
|
||||
print(f"job failed {job['id']}")
|
||||
print(f"Job failed {job['jobId']}")
|
||||
|
||||
update(job["job_id"], job)
|
||||
|
||||
|
||||
def group_user_ids_by_project_id(jobs, now):
|
||||
project_id_user_ids = {}
|
||||
for job in jobs:
|
||||
if job["startAt"] > now:
|
||||
continue
|
||||
|
||||
project_id = job["projectId"]
|
||||
if project_id not in project_id_user_ids:
|
||||
project_id_user_ids[project_id] = []
|
||||
|
||||
project_id_user_ids[project_id].append(job)
|
||||
|
||||
return project_id_user_ids
|
||||
update(job["jobId"], job)
|
||||
|
|
|
|||
|
|
@ -1068,23 +1068,21 @@ def get_session_user(project_id, user_id):
|
|||
def get_session_ids_by_user_ids(project_id, user_ids):
|
||||
with pg_client.PostgresClient() as cur:
|
||||
query = cur.mogrify(
|
||||
"""\
|
||||
SELECT session_id FROM public.sessions
|
||||
WHERE
|
||||
project_id = %(project_id)s AND user_id IN %(userId)s;""",
|
||||
{"project_id": project_id, "userId": tuple(user_ids)}
|
||||
)
|
||||
ids = cur.execute(query=query)
|
||||
return ids
|
||||
"""SELECT session_id
|
||||
FROM public.sessions
|
||||
WHERE project_id = %(project_id)s
|
||||
AND user_id IN %(userId)s;""",
|
||||
{"project_id": project_id, "userId": tuple(user_ids)})
|
||||
cur.execute(query=query)
|
||||
ids = cur.fetchall()
|
||||
return [s["session_id"] for s in ids]
|
||||
|
||||
|
||||
def delete_sessions_by_session_ids(session_ids):
|
||||
with pg_client.PostgresClient(unlimited_query=True) as cur:
|
||||
query = cur.mogrify(
|
||||
"""\
|
||||
DELETE FROM public.sessions
|
||||
WHERE
|
||||
session_id IN %(session_ids)s;""",
|
||||
"""DELETE FROM public.sessions
|
||||
WHERE session_id IN %(session_ids)s;""",
|
||||
{"session_ids": tuple(session_ids)}
|
||||
)
|
||||
cur.execute(query=query)
|
||||
|
|
@ -1092,20 +1090,6 @@ def delete_sessions_by_session_ids(session_ids):
|
|||
return True
|
||||
|
||||
|
||||
def delete_sessions_by_user_ids(project_id, user_ids):
|
||||
with pg_client.PostgresClient(unlimited_query=True) as cur:
|
||||
query = cur.mogrify(
|
||||
"""\
|
||||
DELETE FROM public.sessions
|
||||
WHERE
|
||||
project_id = %(project_id)s AND user_id IN %(userId)s;""",
|
||||
{"project_id": project_id, "userId": tuple(user_ids)}
|
||||
)
|
||||
cur.execute(query=query)
|
||||
|
||||
return True
|
||||
|
||||
|
||||
def count_all():
|
||||
with pg_client.PostgresClient(unlimited_query=True) as cur:
|
||||
cur.execute(query="SELECT COUNT(session_id) AS count FROM public.sessions")
|
||||
|
|
|
|||
|
|
@ -57,5 +57,6 @@ def get_ios(session_id):
|
|||
|
||||
def delete_mobs(project_id, session_ids):
|
||||
for session_id in session_ids:
|
||||
for k in __get_mob_keys(project_id=project_id, session_id=session_id):
|
||||
for k in __get_mob_keys(project_id=project_id, session_id=session_id) \
|
||||
+ __get_mob_keys_deprecated(session_id=session_id):
|
||||
s3.schedule_for_deletion(config("sessions_bucket"), k)
|
||||
|
|
|
|||
|
|
@ -1399,23 +1399,21 @@ def get_session_user(project_id, user_id):
|
|||
def get_session_ids_by_user_ids(project_id, user_ids):
|
||||
with pg_client.PostgresClient() as cur:
|
||||
query = cur.mogrify(
|
||||
"""\
|
||||
SELECT session_id FROM public.sessions
|
||||
WHERE
|
||||
project_id = %(project_id)s AND user_id IN %(userId)s;""",
|
||||
{"project_id": project_id, "userId": tuple(user_ids)}
|
||||
)
|
||||
ids = cur.execute(query=query)
|
||||
return ids
|
||||
"""SELECT session_id
|
||||
FROM public.sessions
|
||||
WHERE project_id = %(project_id)s
|
||||
AND user_id IN %(userId)s;""",
|
||||
{"project_id": project_id, "userId": tuple(user_ids)})
|
||||
cur.execute(query=query)
|
||||
ids = cur.fetchall()
|
||||
return [s["session_id"] for s in ids]
|
||||
|
||||
|
||||
def delete_sessions_by_session_ids(session_ids):
|
||||
with pg_client.PostgresClient(unlimited_query=True) as cur:
|
||||
query = cur.mogrify(
|
||||
"""\
|
||||
DELETE FROM public.sessions
|
||||
WHERE
|
||||
session_id IN %(session_ids)s;""",
|
||||
"""DELETE FROM public.sessions
|
||||
WHERE session_id IN %(session_ids)s;""",
|
||||
{"session_ids": tuple(session_ids)}
|
||||
)
|
||||
cur.execute(query=query)
|
||||
|
|
@ -1423,20 +1421,6 @@ def delete_sessions_by_session_ids(session_ids):
|
|||
return True
|
||||
|
||||
|
||||
def delete_sessions_by_user_ids(project_id, user_ids):
|
||||
with pg_client.PostgresClient(unlimited_query=True) as cur:
|
||||
query = cur.mogrify(
|
||||
"""\
|
||||
DELETE FROM public.sessions
|
||||
WHERE
|
||||
project_id = %(project_id)s AND user_id IN %(userId)s;""",
|
||||
{"project_id": project_id, "userId": tuple(user_ids)}
|
||||
)
|
||||
cur.execute(query=query)
|
||||
|
||||
return True
|
||||
|
||||
|
||||
def count_all():
|
||||
with ch_client.ClickHouseClient() as cur:
|
||||
row = cur.execute(query=f"SELECT COUNT(session_id) AS count FROM {exp_ch_helper.get_main_sessions_table()}")
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue