first piece for /jobs endpoint part5

This commit is contained in:
henryruhs
2026-08-24 07:59:58 +02:00
parent eb60333f19
commit b28fe9d58f
4 changed files with 254 additions and 22 deletions
+3 -3
View File
@@ -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 ]),
+83 -12
View File
@@ -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)
+6 -1
View File
@@ -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'
}
}
+162 -6
View File
@@ -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')