From b28fe9d58f991533f171cd4aa4eb545b6e47f16d Mon Sep 17 00:00:00 2001 From: henryruhs Date: Mon, 24 Aug 2026 07:59:58 +0200 Subject: [PATCH] first piece for /jobs endpoint part5 --- facefusion/apis/core.py | 6 +- facefusion/apis/endpoints/jobs.py | 95 ++++++++++++++--- facefusion/apis/locales.py | 7 +- tests/test_api_jobs.py | 168 ++++++++++++++++++++++++++++-- 4 files changed, 254 insertions(+), 22 deletions(-) diff --git a/facefusion/apis/core.py b/facefusion/apis/core.py index 352378f0..4a1ddc0a 100644 --- a/facefusion/apis/core.py +++ b/facefusion/apis/core.py @@ -8,7 +8,7 @@ from starlette.routing import Route, WebSocketRoute from facefusion.apis.endpoints.assets import delete_assets, get_asset, get_assets, upload_asset from facefusion.apis.endpoints.capabilities import get_capabilities -from facefusion.apis.endpoints.jobs import create_job, delete_job, delete_jobs, get_job, get_jobs, submit_job, submit_jobs +from facefusion.apis.endpoints.jobs import create_job, delete_job, delete_jobs, get_job, get_jobs, update_job, update_jobs from facefusion.apis.endpoints.metrics import get_metrics, websocket_metrics from facefusion.apis.endpoints.ping import websocket_ping from facefusion.apis.endpoints.session import create_session, destroy_session, get_session, refresh_session @@ -50,10 +50,10 @@ def create_api() -> Starlette: Route('/stream', delete_stream, methods = [ 'DELETE' ], name = 'delete_stream', middleware = [ session_guard ]), Route('/jobs', get_jobs, methods = [ 'GET' ], middleware = [ session_guard ]), Route('/jobs', create_job, methods = [ 'POST' ], middleware = [ session_guard ]), - Route('/jobs/submit', submit_jobs, methods = [ 'PUT' ], middleware = [ session_guard ]), + Route('/jobs', update_jobs, methods = [ 'PATCH' ], middleware = [ session_guard ]), Route('/jobs', delete_jobs, methods = [ 'DELETE' ], middleware = [ session_guard ]), Route('/jobs/{job_id}', get_job, methods = [ 'GET' ], middleware = [ session_guard ]), - Route('/jobs/{job_id}/submit', submit_job, methods = [ 'PUT' ], middleware = [ session_guard ]), + Route('/jobs/{job_id}', update_job, methods = [ 'PATCH' ], middleware = [ session_guard ]), Route('/jobs/{job_id}', delete_job, methods = [ 'DELETE' ], middleware = [ session_guard ]), WebSocketRoute('/metrics', websocket_metrics, middleware = [ session_guard ]), WebSocketRoute('/ping', websocket_ping, middleware = [ session_guard ]), diff --git a/facefusion/apis/endpoints/jobs.py b/facefusion/apis/endpoints/jobs.py index 2941fbc6..88d57a02 100644 --- a/facefusion/apis/endpoints/jobs.py +++ b/facefusion/apis/endpoints/jobs.py @@ -1,10 +1,14 @@ +from functools import partial + +from starlette.background import BackgroundTask from starlette.requests import Request from starlette.responses import JSONResponse -from starlette.status import HTTP_200_OK, HTTP_201_CREATED, HTTP_400_BAD_REQUEST, HTTP_404_NOT_FOUND +from starlette.status import HTTP_200_OK, HTTP_201_CREATED, HTTP_202_ACCEPTED, HTTP_400_BAD_REQUEST, HTTP_404_NOT_FOUND import facefusion.choices +import facefusion.core from facefusion import state_manager, translator -from facefusion.jobs import job_helper, job_manager +from facefusion.jobs import job_helper, job_manager, job_runner async def get_jobs(request : Request) -> JSONResponse: @@ -58,31 +62,98 @@ async def create_job(request : Request) -> JSONResponse: }, status_code = HTTP_400_BAD_REQUEST) -async def submit_jobs(request : Request) -> JSONResponse: - if job_manager.submit_jobs(state_manager.get_item('halt_on_error')): +async def update_jobs(request : Request) -> JSONResponse: + action = request.query_params.get('action') + + if action == 'submit': + if job_manager.submit_jobs(state_manager.get_item('halt_on_error')): + return JSONResponse( + { + 'message': translator.get('ok', 'facefusion.apis') + }, status_code = HTTP_200_OK) + return JSONResponse( { - 'message': translator.get('ok', 'facefusion.apis') - }, status_code = HTTP_200_OK) + 'message': translator.get('job_all_not_submitted', 'facefusion.apis') + }, status_code = HTTP_400_BAD_REQUEST) + + if action == 'run': + if job_manager.find_job_ids('queued'): + run_jobs_task = BackgroundTask(partial(job_runner.run_jobs, facefusion.core.process_step, state_manager.get_item('halt_on_error'))) + return JSONResponse( + { + 'message': translator.get('ok', 'facefusion.apis') + }, status_code = HTTP_202_ACCEPTED, background = run_jobs_task) + + return JSONResponse( + { + 'message': translator.get('job_all_not_run', 'facefusion.apis') + }, status_code = HTTP_400_BAD_REQUEST) + + if action == 'retry': + if job_manager.find_job_ids('failed'): + retry_jobs_task = BackgroundTask(partial(job_runner.retry_jobs, facefusion.core.process_step, state_manager.get_item('halt_on_error'))) + return JSONResponse( + { + 'message': translator.get('ok', 'facefusion.apis') + }, status_code = HTTP_202_ACCEPTED, background = retry_jobs_task) + + return JSONResponse( + { + 'message': translator.get('job_all_not_retried', 'facefusion.apis') + }, status_code = HTTP_400_BAD_REQUEST) return JSONResponse( { - 'message': translator.get('job_all_not_submitted', 'facefusion.apis') + 'message': translator.get('invalid_job_action', 'facefusion.apis') }, status_code = HTTP_400_BAD_REQUEST) -async def submit_job(request : Request) -> JSONResponse: +async def update_job(request : Request) -> JSONResponse: job_id = request.path_params.get('job_id') + action = request.query_params.get('action') + + if action == 'submit': + if job_manager.submit_job(job_id): + return JSONResponse( + { + 'message': translator.get('ok', 'facefusion.apis') + }, status_code = HTTP_200_OK) - if job_manager.submit_job(job_id): return JSONResponse( { - 'message': translator.get('ok', 'facefusion.apis') - }, status_code = HTTP_200_OK) + 'message': translator.get('job_not_submitted', 'facefusion.apis') + }, status_code = HTTP_400_BAD_REQUEST) + + if action == 'run': + if job_id in job_manager.find_job_ids('queued'): + run_job_task = BackgroundTask(partial(job_runner.run_job, job_id, facefusion.core.process_step)) + return JSONResponse( + { + 'message': translator.get('ok', 'facefusion.apis') + }, status_code = HTTP_202_ACCEPTED, background = run_job_task) + + return JSONResponse( + { + 'message': translator.get('job_not_run', 'facefusion.apis') + }, status_code = HTTP_400_BAD_REQUEST) + + if action == 'retry': + if job_id in job_manager.find_job_ids('failed'): + retry_job_task = BackgroundTask(partial(job_runner.retry_job, job_id, facefusion.core.process_step)) + return JSONResponse( + { + 'message': translator.get('ok', 'facefusion.apis') + }, status_code = HTTP_202_ACCEPTED, background = retry_job_task) + + return JSONResponse( + { + 'message': translator.get('job_not_retried', 'facefusion.apis') + }, status_code = HTTP_400_BAD_REQUEST) return JSONResponse( { - 'message': translator.get('job_not_submitted', 'facefusion.apis') + 'message': translator.get('invalid_job_action', 'facefusion.apis') }, status_code = HTTP_400_BAD_REQUEST) diff --git a/facefusion/apis/locales.py b/facefusion/apis/locales.py index d4ff2496..f0d90315 100644 --- a/facefusion/apis/locales.py +++ b/facefusion/apis/locales.py @@ -12,11 +12,16 @@ LOCALES : Locales =\ 'target_asset_not_found': 'target asset not found', 'invalid_state_key': 'invalid state key', 'invalid_job_status': 'invalid job status', + 'invalid_job_action': 'invalid job action', 'job_not_found': 'job not found', 'job_not_created': 'job not created', 'job_not_submitted': 'job not submitted', 'job_not_deleted': 'job not deleted', + 'job_not_run': 'job not run', + 'job_not_retried': 'job not retried', 'job_all_not_submitted': 'jobs not submitted', - 'job_all_not_deleted': 'jobs not deleted' + 'job_all_not_deleted': 'jobs not deleted', + 'job_all_not_run': 'jobs not run', + 'job_all_not_retried': 'jobs not retried' } } diff --git a/tests/test_api_jobs.py b/tests/test_api_jobs.py index 25b5b0ff..216dd543 100644 --- a/tests/test_api_jobs.py +++ b/tests/test_api_jobs.py @@ -129,7 +129,7 @@ def test_create_job(test_client : TestClient) -> None: def test_submit_jobs(test_client : TestClient) -> None: - submit_jobs_response = test_client.put('/jobs/submit') + submit_jobs_response = test_client.patch('/jobs?action=submit') assert submit_jobs_response.status_code == 401 @@ -140,9 +140,18 @@ def test_submit_jobs(test_client : TestClient) -> None: create_session_body = create_session_response.json() access_token = create_session_body.get('access_token') + submit_jobs_response = test_client.patch('/jobs?action=invalid', headers = + { + 'Authorization': 'Bearer ' + access_token + }) + submit_jobs_body = submit_jobs_response.json() + + assert submit_jobs_body.get('message') == 'invalid job action' + assert submit_jobs_response.status_code == 400 + create_job('job-test-submit-jobs') - submit_jobs_response = test_client.put('/jobs/submit', headers = + submit_jobs_response = test_client.patch('/jobs?action=submit', headers = { 'Authorization': 'Bearer ' + access_token }) @@ -152,7 +161,7 @@ def test_submit_jobs(test_client : TestClient) -> None: assert submit_jobs_response.status_code == 400 with patch('facefusion.jobs.job_manager.submit_jobs', return_value = True): - submit_jobs_response = test_client.put('/jobs/submit', headers = + submit_jobs_response = test_client.patch('/jobs?action=submit', headers = { 'Authorization': 'Bearer ' + access_token }) @@ -163,7 +172,7 @@ def test_submit_jobs(test_client : TestClient) -> None: def test_submit_job(test_client : TestClient) -> None: - submit_job_response = test_client.put('/jobs/job-test-submit-job/submit') + submit_job_response = test_client.patch('/jobs/job-test-submit-job?action=submit') assert submit_job_response.status_code == 401 @@ -174,9 +183,18 @@ def test_submit_job(test_client : TestClient) -> None: create_session_body = create_session_response.json() access_token = create_session_body.get('access_token') + submit_job_response = test_client.patch('/jobs/job-test-submit-job?action=invalid', headers = + { + 'Authorization': 'Bearer ' + access_token + }) + submit_job_body = submit_job_response.json() + + assert submit_job_body.get('message') == 'invalid job action' + assert submit_job_response.status_code == 400 + create_job('job-test-submit-job') - submit_job_response = test_client.put('/jobs/job-test-submit-job/submit', headers = + submit_job_response = test_client.patch('/jobs/job-test-submit-job?action=submit', headers = { 'Authorization': 'Bearer ' + access_token }) @@ -186,7 +204,7 @@ def test_submit_job(test_client : TestClient) -> None: assert submit_job_response.status_code == 400 with patch('facefusion.jobs.job_manager.submit_job', return_value = True): - submit_job_response = test_client.put('/jobs/job-test-submit-job/submit', headers = + submit_job_response = test_client.patch('/jobs/job-test-submit-job?action=submit', headers = { 'Authorization': 'Bearer ' + access_token }) @@ -196,6 +214,144 @@ def test_submit_job(test_client : TestClient) -> None: assert submit_job_response.status_code == 200 +def test_run_jobs(test_client : TestClient) -> None: + run_jobs_response = test_client.patch('/jobs?action=run') + + assert run_jobs_response.status_code == 401 + + create_session_response = test_client.post('/session', json = + { + 'client_version': metadata.get('version') + }) + create_session_body = create_session_response.json() + access_token = create_session_body.get('access_token') + + run_jobs_response = test_client.patch('/jobs?action=run', headers = + { + 'Authorization': 'Bearer ' + access_token + }) + run_jobs_body = run_jobs_response.json() + + assert run_jobs_body.get('message') == 'jobs not run' + assert run_jobs_response.status_code == 400 + + with patch('facefusion.jobs.job_manager.find_job_ids', return_value = [ 'job-test-run-jobs' ]): + with patch('facefusion.jobs.job_runner.run_jobs', return_value = True) as run_jobs_mock: + run_jobs_response = test_client.patch('/jobs?action=run', headers = + { + 'Authorization': 'Bearer ' + access_token + }) + run_jobs_body = run_jobs_response.json() + + assert run_jobs_body.get('message') == 'ok' + assert run_jobs_response.status_code == 202 + assert run_jobs_mock.called is True + + +def test_run_job(test_client : TestClient) -> None: + run_job_response = test_client.patch('/jobs/job-test-run-job?action=run') + + assert run_job_response.status_code == 401 + + create_session_response = test_client.post('/session', json = + { + 'client_version': metadata.get('version') + }) + create_session_body = create_session_response.json() + access_token = create_session_body.get('access_token') + + create_job('job-test-run-job') + + run_job_response = test_client.patch('/jobs/job-test-run-job?action=run', headers = + { + 'Authorization': 'Bearer ' + access_token + }) + run_job_body = run_job_response.json() + + assert run_job_body.get('message') == 'job not run' + assert run_job_response.status_code == 400 + + with patch('facefusion.jobs.job_manager.find_job_ids', return_value = [ 'job-test-run-job' ]): + with patch('facefusion.jobs.job_runner.run_job', return_value = True) as run_job_mock: + run_job_response = test_client.patch('/jobs/job-test-run-job?action=run', headers = + { + 'Authorization': 'Bearer ' + access_token + }) + run_job_body = run_job_response.json() + + assert run_job_body.get('message') == 'ok' + assert run_job_response.status_code == 202 + assert run_job_mock.called is True + + +def test_retry_jobs(test_client : TestClient) -> None: + retry_jobs_response = test_client.patch('/jobs?action=retry') + + assert retry_jobs_response.status_code == 401 + + create_session_response = test_client.post('/session', json = + { + 'client_version': metadata.get('version') + }) + create_session_body = create_session_response.json() + access_token = create_session_body.get('access_token') + + retry_jobs_response = test_client.patch('/jobs?action=retry', headers = + { + 'Authorization': 'Bearer ' + access_token + }) + retry_jobs_body = retry_jobs_response.json() + + assert retry_jobs_body.get('message') == 'jobs not retried' + assert retry_jobs_response.status_code == 400 + + with patch('facefusion.jobs.job_manager.find_job_ids', return_value = [ 'job-test-retry-jobs' ]): + with patch('facefusion.jobs.job_runner.retry_jobs', return_value = True) as retry_jobs_mock: + retry_jobs_response = test_client.patch('/jobs?action=retry', headers = + { + 'Authorization': 'Bearer ' + access_token + }) + retry_jobs_body = retry_jobs_response.json() + + assert retry_jobs_body.get('message') == 'ok' + assert retry_jobs_response.status_code == 202 + assert retry_jobs_mock.called is True + + +def test_retry_job(test_client : TestClient) -> None: + retry_job_response = test_client.patch('/jobs/job-test-retry-job?action=retry') + + assert retry_job_response.status_code == 401 + + create_session_response = test_client.post('/session', json = + { + 'client_version': metadata.get('version') + }) + create_session_body = create_session_response.json() + access_token = create_session_body.get('access_token') + + retry_job_response = test_client.patch('/jobs/job-test-retry-job?action=retry', headers = + { + 'Authorization': 'Bearer ' + access_token + }) + retry_job_body = retry_job_response.json() + + assert retry_job_body.get('message') == 'job not retried' + assert retry_job_response.status_code == 400 + + with patch('facefusion.jobs.job_manager.find_job_ids', return_value = [ 'job-test-retry-job' ]): + with patch('facefusion.jobs.job_runner.retry_job', return_value = True) as retry_job_mock: + retry_job_response = test_client.patch('/jobs/job-test-retry-job?action=retry', headers = + { + 'Authorization': 'Bearer ' + access_token + }) + retry_job_body = retry_job_response.json() + + assert retry_job_body.get('message') == 'ok' + assert retry_job_response.status_code == 202 + assert retry_job_mock.called is True + + def test_delete_jobs(test_client : TestClient) -> None: delete_jobs_response = test_client.delete('/jobs')