@@ -1376,6 +1376,125 @@ def get_video_job_status(api_key, job_id):
13761376 return response .json ()
13771377
13781378
1379+ # ---------------------------------------------------------------------------
1380+ # Batch Processing (Asset Library orchestration)
1381+ # ---------------------------------------------------------------------------
1382+
1383+
1384+ def _batch_processing_url (workspace_url , suffix = "" ):
1385+ return f"{ API_URL } /batch-processing/v1/external/{ workspace_url } /asset-library/jobs{ suffix } "
1386+
1387+
1388+ def _batch_processing_headers (api_key ):
1389+ # Keep credentials out of URLs, proxy logs, and shell history. validateToken supports Bearer.
1390+ return {"Authorization" : f"Bearer { api_key } " }
1391+
1392+
1393+ def _raise_for_batch_processing_response (response ):
1394+ message = response .text
1395+ try :
1396+ body = response .json ()
1397+ if isinstance (body , dict ):
1398+ error = body .get ("error" )
1399+ if isinstance (error , dict ):
1400+ message = error .get ("message" ) or error .get ("hint" ) or message
1401+ elif error :
1402+ message = str (error )
1403+ else :
1404+ message = body .get ("message" ) or message
1405+ except (TypeError , ValueError ):
1406+ pass
1407+ raise RoboflowError (message , status_code = response .status_code )
1408+
1409+
1410+ def create_asset_library_batch_job (
1411+ api_key ,
1412+ workspace_url ,
1413+ * ,
1414+ workflow_id ,
1415+ idempotency_key ,
1416+ image_ids = None ,
1417+ query = None ,
1418+ machine_type = "cpu" ,
1419+ display_name = None ,
1420+ ):
1421+ """Queue a published Workflow over an exact Asset Library selection."""
1422+ payload = {
1423+ "workflowId" : workflow_id ,
1424+ "idempotencyKey" : idempotency_key ,
1425+ "machineType" : machine_type ,
1426+ }
1427+ if image_ids is not None :
1428+ payload ["imageIds" ] = image_ids
1429+ if query is not None :
1430+ payload ["query" ] = query
1431+ if display_name :
1432+ payload ["displayName" ] = display_name
1433+ response = requests .post (
1434+ _batch_processing_url (workspace_url ),
1435+ headers = _batch_processing_headers (api_key ),
1436+ json = payload ,
1437+ )
1438+ if response .status_code != 202 :
1439+ _raise_for_batch_processing_response (response )
1440+ return response .json ()
1441+
1442+
1443+ def list_batch_processing_jobs (api_key , workspace_url , * , page_size = 10 , next_page_token = None , search = None ):
1444+ """List durable Batch Processing jobs in a workspace."""
1445+ params = {"pageSize" : page_size }
1446+ if next_page_token :
1447+ params ["nextPageToken" ] = next_page_token
1448+ if search :
1449+ params ["search" ] = search
1450+ response = requests .get (
1451+ _batch_processing_url (workspace_url ),
1452+ headers = _batch_processing_headers (api_key ),
1453+ params = params ,
1454+ )
1455+ if response .status_code != 200 :
1456+ _raise_for_batch_processing_response (response )
1457+ return response .json ()
1458+
1459+
1460+ def get_batch_processing_job (api_key , workspace_url , job_id ):
1461+ """Get current metadata for one Batch Processing job."""
1462+ encoded = quote (job_id , safe = "" )
1463+ response = requests .get (
1464+ _batch_processing_url (workspace_url , f"/{ encoded } " ),
1465+ headers = _batch_processing_headers (api_key ),
1466+ )
1467+ if response .status_code != 200 :
1468+ _raise_for_batch_processing_response (response )
1469+ return response .json ()
1470+
1471+
1472+ def abort_batch_processing_job (api_key , workspace_url , job_id ):
1473+ """Abort one Batch Processing job."""
1474+ encoded = quote (job_id , safe = "" )
1475+ response = requests .post (
1476+ _batch_processing_url (workspace_url , f"/{ encoded } /abort" ),
1477+ headers = _batch_processing_headers (api_key ),
1478+ json = {},
1479+ )
1480+ if response .status_code != 200 :
1481+ _raise_for_batch_processing_response (response )
1482+ return response .json ()
1483+
1484+
1485+ def restart_batch_processing_job (api_key , workspace_url , job_id ):
1486+ """Restart one Batch Processing job with its existing configuration."""
1487+ encoded = quote (job_id , safe = "" )
1488+ response = requests .post (
1489+ _batch_processing_url (workspace_url , f"/{ encoded } /restart" ),
1490+ headers = _batch_processing_headers (api_key ),
1491+ json = {},
1492+ )
1493+ if response .status_code != 200 :
1494+ _raise_for_batch_processing_response (response )
1495+ return response .json ()
1496+
1497+
13791498# ---------------------------------------------------------------------------
13801499# Phase 2: Universe search
13811500# ---------------------------------------------------------------------------
0 commit comments