diff --git a/docs/core/project.md b/docs/core/project.md index 7b2e6d68..56a602e0 100644 --- a/docs/core/project.md +++ b/docs/core/project.md @@ -1 +1,37 @@ :::roboflow.core.project + +## Upload a native Action Recognition video + +`Project.upload_video` sends the original MP4 or MOV bytes to a signed upload +URL. It creates a video Source in the project; it does not extract frames or +run inference. The platform processes the upload asynchronously. + +```python +project = rf.workspace("my-workspace").project("my-actions") +status = project.upload_video( + "clip.mp4", + batch_name="session-1", + tag_names=["indoor"], + metadata={"camera": "front"}, + split="train", +) +if status["status"] == "pending": + status = project.wait_for_video_upload(status["videoId"], poll_timeout=300) + +if status["status"] == "failed": + raise RuntimeError(status["message"]) + +source_id = status["videoId"] # Use this Source ID for video annotations. +``` + +`upload_video(..., wait=True)` performs the bounded wait in one call. The +returned status is the API response: `pending`, `uploaded` (with +`resolvedBatch`), or `failed` (with `message`). Poll later with +`project.get_video_upload_status(video_id)`. Always use `videoId` from the +final `uploaded` response because ingestion can deduplicate onto another +Source. Batch, tags, metadata, and split follow the platform upload API; +the API validates their values. A timeout leaves the upload running, so +poll its original ID later. `poll_timeout=0` makes one status request and +returns a terminal result if available. Status requests use the remaining +polling budget as their connection and read inactivity timeout; this is not +a strict whole-response wall-clock limit for a slowly streaming server. diff --git a/roboflow/adapters/rfapi.py b/roboflow/adapters/rfapi.py index 31690e87..95793070 100644 --- a/roboflow/adapters/rfapi.py +++ b/roboflow/adapters/rfapi.py @@ -910,6 +910,62 @@ def _save_annotation_error(response): return AnnotationSaveError(str(responsejson), status_code=response.status_code) +# --------------------------------------------------------------------------- +# Native video upload endpoints +# --------------------------------------------------------------------------- + +VIDEO_UPLOAD_PREPARE_TIMEOUT = (5, 30) +VIDEO_UPLOAD_STATUS_TIMEOUT = 30 + + +def prepare_video_upload(api_key, workspace_url, project_url, body) -> dict: + """Prepare a native video Source upload and obtain its signed PUT URL.""" + try: + response = requests.post( + f"{API_URL}/{workspace_url}/upload/video", + params={"api_key": api_key}, + json={"project": project_url, **body}, + timeout=VIDEO_UPLOAD_PREPARE_TIMEOUT, + ) + except RequestException as error: + raise RoboflowError(f"Video upload preparation request failed: {type(error).__name__}") from None + if not response.ok: + raise RoboflowError(response.text, status_code=response.status_code) + return response.json() + + +def put_video_upload(signed_url, video_path, required_headers, content_type) -> None: + """Stream original video bytes to the API-issued signed URL.""" + with open(video_path, "rb") as video: + response = requests.put( + signed_url, + data=video, + headers={"Content-Type": content_type, **required_headers}, + timeout=(30, 3600), + ) + if not response.ok: + raise RoboflowError(response.text, status_code=response.status_code) + + +def get_video_upload_status(api_key, workspace_url, video_id, *, timeout=None) -> dict: + """Read processing state and the canonical Source ID after ingestion.""" + if timeout is None: + timeout = VIDEO_UPLOAD_STATUS_TIMEOUT + if timeout <= 0: + raise ValueError("Video upload status timeout must be positive") + try: + response = requests.get( + f"{API_URL}/{workspace_url}/upload/video/{video_id}", + params={"api_key": api_key}, + timeout=timeout, + ) + except RequestException as error: + raise RoboflowError(f"Video upload status request failed for {video_id}: {type(error).__name__}") from None + if not response.ok: + raise RoboflowError(response.text, status_code=response.status_code) + return response.json() + + # --------------------------------------------------------------------------- # Zip upload endpoints # --------------------------------------------------------------------------- diff --git a/roboflow/core/project.py b/roboflow/core/project.py index 03dc0ab6..8cfeb017 100644 --- a/roboflow/core/project.py +++ b/roboflow/core/project.py @@ -868,6 +868,89 @@ def __str__(self): return json.dumps(json_str, indent=2) + def upload_video( + self, + video_path: str, + *, + batch_name: Optional[str] = None, + tag_names: Optional[Union[str, List[str]]] = None, + metadata: Optional[Dict] = None, + split: Optional[str] = None, + wait: bool = False, + poll_interval: float = 2, + poll_timeout: float = 300, + ) -> Dict: + """Upload original MP4/MOV bytes as a native video Source. + + Returns the API processing status, including ``videoId``. Once the + status is ``uploaded``, that ID is the canonical Source ID to annotate. + The ID can change during ingestion if the video is deduplicated. + ``wait=False`` reads status once after the signed PUT; use + :meth:`wait_for_video_upload` to continue polling later. + """ + if not os.path.isfile(video_path): + raise ValueError(f"Video file not found: {video_path}") + content_type = {".mp4": "video/mp4", ".mov": "video/quicktime"}.get(os.path.splitext(video_path)[1].lower()) + if content_type is None: + raise ValueError("Native video upload accepts .mp4 and .mov files") + + body: Dict = {"name": os.path.basename(video_path), "contentType": content_type} + if batch_name is not None: + body["batch"] = batch_name + if tag_names is not None: + body["tag"] = tag_names + if metadata is not None: + body["metadata"] = metadata + if split is not None: + body["split"] = split + + prepared = rfapi.prepare_video_upload(self.__api_key, self.__workspace, self.__project_name, body) + video_id = prepared["videoId"] + if not prepared.get("signedUrl"): + raise rfapi.RoboflowError("Video upload API did not return a signedUrl") + rfapi.put_video_upload(prepared["signedUrl"], video_path, prepared.get("requiredHeaders", {}), content_type) + if wait: + return self.wait_for_video_upload(video_id, poll_interval=poll_interval, poll_timeout=poll_timeout) + return self.get_video_upload_status(video_id) + + def get_video_upload_status(self, video_id: str, *, timeout: Optional[float] = None) -> Dict: + """Get a native video's processing state and canonical Source ID. + + ``timeout`` limits connection and response-read inactivity. It is not + a strict total request-duration cap. + """ + return rfapi.get_video_upload_status(self.__api_key, self.__workspace, video_id, timeout=timeout) + + def wait_for_video_upload(self, video_id: str, *, poll_interval: float = 2, poll_timeout: float = 300) -> Dict: + """Poll until uploaded or failed, limiting each status read to the remaining budget. + + With ``poll_timeout=0``, perform one status read using the default + transport timeout and return a terminal result if it is already ready. + Requests' timeouts measure connection/read inactivity, so this is not + a strict wall-clock cap on a slowly streaming response. + """ + if poll_interval <= 0 or poll_timeout < 0: + raise ValueError("poll_interval must be positive and poll_timeout must be nonnegative") + deadline = time.monotonic() + poll_timeout + while True: + remaining = deadline - time.monotonic() + if poll_timeout > 0 and remaining <= 0: + raise rfapi.RoboflowError( + f"Video upload {video_id} did not finish within the {poll_timeout}s polling budget; " + "call get_video_upload_status to check later" + ) + request_timeout = min(rfapi.VIDEO_UPLOAD_STATUS_TIMEOUT, remaining) if poll_timeout > 0 else None + status = self.get_video_upload_status(video_id, timeout=request_timeout) + if status.get("status") in {"uploaded", "failed"}: + return status + remaining = deadline - time.monotonic() + if remaining <= 0: + raise rfapi.RoboflowError( + f"Video upload {video_id} is still {status.get('status')} after {poll_timeout}s; " + "call get_video_upload_status to check later" + ) + time.sleep(min(poll_interval, remaining)) + def image(self, image_id: str) -> Dict: """ Fetch the details of a specific image from the Roboflow API. diff --git a/tests/test_native_video_upload.py b/tests/test_native_video_upload.py new file mode 100644 index 00000000..95d8391f --- /dev/null +++ b/tests/test_native_video_upload.py @@ -0,0 +1,287 @@ +"""Public Project calls through the native video preparation, PUT and status API.""" + +import json +import os +import tempfile +import threading +import time +import unittest +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from unittest.mock import patch +from urllib.parse import parse_qs, urlparse + +import requests +import responses +from responses.matchers import json_params_matcher + +from roboflow.adapters.rfapi import RoboflowError +from roboflow.config import API_URL +from roboflow.core.project import Project +from tests import ROBOFLOW_API_KEY, WORKSPACE_NAME, RoboflowTest + + +class TestNativeVideoUpload(RoboflowTest): + def setUp(self): + super().setUp() + self.project.type = "action-recognition" + self.temp_dir = tempfile.TemporaryDirectory(dir=".") + self.video_path = os.path.join(self.temp_dir.name, "clip.mp4") + with open(self.video_path, "wb") as video: + video.write(b"original video bytes") + self.addCleanup(self.temp_dir.cleanup) + self.prepare_url = f"{API_URL}/{WORKSPACE_NAME}/upload/video?api_key={ROBOFLOW_API_KEY}" + self.status_url = f"{API_URL}/{WORKSPACE_NAME}/upload/video/upload-1?api_key={ROBOFLOW_API_KEY}" + + def test_upload_streams_original_bytes_and_returns_canonical_source(self): + headers = {"x-goog-content-length-range": "1,100", "x-goog-if-generation-match": "0"} + responses.add( + responses.POST, + self.prepare_url, + json={"videoId": "upload-1", "signedUrl": "https://signed.example/video", "requiredHeaders": headers}, + status=200, + match=[ + json_params_matcher( + { + "project": "test-project", + "name": "clip.mp4", + "contentType": "video/mp4", + "batch": "clips", + "tag": ["indoor"], + "metadata": {"camera": "one"}, + "split": "valid", + } + ) + ], + ) + uploaded = {} + + def receive_video(request): + uploaded["body"] = request.body + uploaded["headers"] = request.headers + uploaded["url"] = request.url + return 200, {}, "" + + responses.add_callback(responses.PUT, "https://signed.example/video", callback=receive_video) + responses.add( + responses.GET, + self.status_url, + json={"videoId": "source-2", "status": "uploaded", "duplicate": True, "resolvedBatch": None}, + status=200, + ) + + result = self.project.upload_video( + self.video_path, + batch_name="clips", + tag_names=["indoor"], + metadata={"camera": "one"}, + split="valid", + ) + + self.assertEqual(result["videoId"], "source-2") + self.assertEqual(result["status"], "uploaded") + self.assertEqual(uploaded["body"], b"original video bytes") + self.assertEqual(uploaded["headers"]["Content-Type"], "video/mp4") + for key, value in headers.items(): + self.assertEqual(uploaded["headers"][key], value) + self.assertNotIn("api_key", uploaded["url"]) + + def test_pending_then_wait_returns_uploaded_with_batch(self): + responses.add( + responses.POST, + self.prepare_url, + json={"videoId": "upload-1", "signedUrl": "https://signed.example/video", "requiredHeaders": {}}, + status=200, + ) + responses.add(responses.PUT, "https://signed.example/video", status=200) + responses.add(responses.GET, self.status_url, json={"videoId": "upload-1", "status": "pending"}, status=200) + responses.add(responses.GET, self.status_url, json={"videoId": "upload-1", "status": "pending"}, status=200) + responses.add( + responses.GET, + self.status_url, + json={"videoId": "upload-1", "status": "uploaded", "resolvedBatch": {"id": "b1", "name": "clips"}}, + status=200, + ) + + first = self.project.upload_video(self.video_path) + with patch("roboflow.core.project.time.sleep"): + final = self.project.wait_for_video_upload(first["videoId"], poll_interval=0.1, poll_timeout=10) + + self.assertEqual(first["status"], "pending") + self.assertEqual(final["resolvedBatch"], {"id": "b1", "name": "clips"}) + + def test_server_errors_and_bounded_wait(self): + responses.add(responses.POST, self.prepare_url, json={"error": "quota exceeded"}, status=403) + with self.assertRaises(RoboflowError) as error: + self.project.upload_video(self.video_path) + self.assertIn("quota exceeded", str(error.exception)) + self.assertEqual(error.exception.status_code, 403) + + responses.add(responses.GET, self.status_url, json={"videoId": "upload-1", "status": "pending"}, status=200) + with self.assertRaises(RoboflowError) as timeout: + self.project.wait_for_video_upload("upload-1", poll_timeout=0) + self.assertIn("upload-1", str(timeout.exception)) + + def test_signed_put_error_keeps_server_response(self): + responses.add( + responses.POST, + self.prepare_url, + json={"videoId": "upload-1", "signedUrl": "https://signed.example/video", "requiredHeaders": {}}, + status=200, + ) + responses.add(responses.PUT, "https://signed.example/video", body="object already exists", status=412) + + with self.assertRaises(RoboflowError) as error: + self.project.upload_video(self.video_path) + + self.assertIn("object already exists", str(error.exception)) + self.assertEqual(error.exception.status_code, 412) + self.assertFalse( + any(call.request.method == "GET" and "/upload/video/" in call.request.url for call in responses.calls) + ) + + def test_failed_processing_returns_api_message(self): + responses.add( + responses.GET, + self.status_url, + json={"videoId": "upload-1", "status": "failed", "message": "Invalid video media"}, + status=200, + ) + self.assertEqual( + self.project.wait_for_video_upload("upload-1", poll_timeout=0), + {"videoId": "upload-1", "status": "failed", "message": "Invalid video media"}, + ) + + def test_invalid_file_does_not_prepare_upload(self): + before = len(responses.calls) + with self.assertRaises(ValueError): + self.project.upload_video(os.path.join(self.temp_dir.name, "missing.mp4")) + self.assertEqual(len(responses.calls), before) + + def test_preparation_and_status_transport_errors_are_sanitized(self): + with patch("roboflow.adapters.rfapi.requests.post", side_effect=requests.exceptions.ConnectTimeout) as post: + with self.assertRaises(RoboflowError) as preparation: + self.project.upload_video(self.video_path) + self.assertIn("ConnectTimeout", str(preparation.exception)) + self.assertNotIn(ROBOFLOW_API_KEY, str(preparation.exception)) + self.assertEqual(post.call_args.kwargs["timeout"], (5, 30)) + + with patch("roboflow.adapters.rfapi.requests.get", side_effect=requests.exceptions.ConnectionError) as get: + with self.assertRaises(RoboflowError) as status: + self.project.get_video_upload_status("upload-1") + self.assertIn("ConnectionError", str(status.exception)) + self.assertNotIn(ROBOFLOW_API_KEY, str(status.exception)) + self.assertEqual(get.call_args.kwargs["timeout"], 30) + + +class TestNativeVideoUploadOverHttp(unittest.TestCase): + @staticmethod + def make_project(): + return Project( + "test-key", + { + "annotation": "", + "classes": {}, + "colors": {}, + "created": 0, + "id": "region-workspace/actions", + "images": 0, + "name": "Actions", + "public": False, + "splits": {}, + "type": "action-recognition", + "unannotated": 0, + "updated": 0, + }, + ) + + def test_public_call_uses_configured_host_and_streams_original_bytes(self): + received = {} + + class Handler(BaseHTTPRequestHandler): + def log_message(self, *_args): + pass + + def respond(self, body): + payload = json.dumps(body).encode() + self.send_response(200) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(payload))) + self.end_headers() + self.wfile.write(payload) + + def do_POST(self): + received["prepare_path"] = self.path + received["prepare_body"] = json.loads(self.rfile.read(int(self.headers["Content-Length"]))) + self.respond( + { + "videoId": "upload-1", + "signedUrl": f"http://127.0.0.1:{self.server.server_port}/signed", + "requiredHeaders": {"x-goog-if-generation-match": "0"}, + } + ) + + def do_PUT(self): + received["put_path"] = self.path + received["put_body"] = self.rfile.read(int(self.headers["Content-Length"])) + received["put_header"] = self.headers["x-goog-if-generation-match"] + self.respond({}) + + def do_GET(self): + received["status_path"] = self.path + self.respond({"videoId": "source-2", "status": "uploaded", "resolvedBatch": None}) + + server = ThreadingHTTPServer(("127.0.0.1", 0), Handler) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + self.addCleanup(thread.join) + self.addCleanup(server.server_close) + self.addCleanup(server.shutdown) + + with tempfile.TemporaryDirectory(dir=".") as temp_dir: + video_path = os.path.join(temp_dir, "native.mp4") + with open(video_path, "wb") as video: + video.write(b"original video bytes") + project = self.make_project() + with patch("roboflow.adapters.rfapi.API_URL", f"http://127.0.0.1:{server.server_port}"): + status = project.upload_video(video_path) + + self.assertEqual(status["videoId"], "source-2") + self.assertEqual( + received["prepare_body"], {"project": "actions", "name": "native.mp4", "contentType": "video/mp4"} + ) + self.assertEqual(urlparse(received["prepare_path"]).path, "/region-workspace/upload/video") + self.assertEqual(parse_qs(urlparse(received["prepare_path"]).query), {"api_key": ["test-key"]}) + self.assertEqual(received["put_path"], "/signed") + self.assertEqual(received["put_body"], b"original video bytes") + self.assertEqual(received["put_header"], "0") + self.assertEqual(urlparse(received["status_path"]).path, "/region-workspace/upload/video/upload-1") + + def test_wait_uses_remaining_budget_for_silent_status_server(self): + entered = threading.Event() + release = threading.Event() + + class Handler(BaseHTTPRequestHandler): + def log_message(self, *_args): + pass + + def do_GET(self): + entered.set() + release.wait(2) + + server = ThreadingHTTPServer(("127.0.0.1", 0), Handler) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + self.addCleanup(thread.join) + self.addCleanup(server.server_close) + self.addCleanup(server.shutdown) + self.addCleanup(release.set) + + with patch("roboflow.adapters.rfapi.API_URL", f"http://127.0.0.1:{server.server_port}"): + started = time.monotonic() + with self.assertRaises(RoboflowError) as error: + self.make_project().wait_for_video_upload("upload-1", poll_timeout=0.15) + elapsed = time.monotonic() - started + + self.assertTrue(entered.is_set()) + self.assertIn("ReadTimeout", str(error.exception)) + self.assertLess(elapsed, 1.5)