Skip to content

Commit ad880cc

Browse files
authored
Merge pull request #42 from Bibek200619/feature/web-ui-live-camera-ugv-simulation
Build minimal operator UI, live camera relay, and UGV simulation
2 parents 8e3ffd0 + fabb0e8 commit ad880cc

75 files changed

Lines changed: 7723 additions & 5049 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎.idea/workspace.xml‎

Lines changed: 26 additions & 27 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎web_app/README.md‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -66,3 +66,16 @@ Before adding a database field, API name, WebSocket event, or ROS-to-app mapping
6666
- [`docs/DATABASE_SCHEMA.md`](docs/DATABASE_SCHEMA.md)
6767

6868
Database/API-facing names use `snake_case`; frontend models may use `camelCase` through explicit adapters.
69+
70+
## Live camera and operator workspace
71+
72+
The redesigned workspace includes an authenticated live MJPEG camera viewer with snapshots,
73+
pause/resume, full screen, and explicit connection states. See
74+
[Live camera setup](docs/LIVE_CAMERA.md) for vehicle configuration and verification.
75+
76+
## First product demo — simulated UGV
77+
78+
Run `./simulation/start.sh` to launch a local 3D off-road patrol with mountain, rocky, forest, and custom terrain, plus a dedicated
79+
operator dashboard with simulated video and telemetry. Open <http://127.0.0.1:8010>
80+
and click **Run guided demo**. The demo dashboard runs at <http://127.0.0.1:5174>.
81+
See [simulation/README.md](simulation/README.md) for the presentation script and controls.

‎web_app/backend/.env.example‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,3 +20,10 @@ TELEMETRY_PERSISTENCE_RATE_HZ=2
2020
WEBSOCKET_QUEUE_SIZE=100
2121
WEBSOCKET_MAX_DROPPED_MESSAGES=25
2222
WEBSOCKET_AUTH_TIMEOUT_SECONDS=5
23+
24+
# Read-only HTTP MJPEG camera source. The URL stays on the backend.
25+
# For ROS web_video_server, e.g. http://<vehicle-host>:8080/stream?topic=/camera/image_raw&type=mjpeg
26+
CAMERA_STREAM_URL=
27+
CAMERA_TIMEOUT_SECONDS=8
28+
CAMERA_MAX_VIEWERS=4
29+
CAMERA_MAX_FRAME_BYTES=4194304

‎web_app/backend/app/api/v1/router.py‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,9 @@
11
from fastapi import APIRouter
22

3-
from app.api.v1.routes import admin, commands, logs, missions, robots, system
3+
from app.api.v1.routes import admin, cameras, commands, logs, missions, robots, system
44

55
router = APIRouter()
6+
router.include_router(cameras.router, prefix="/cameras", tags=["cameras"])
67
router.include_router(system.router, tags=["system"])
78
router.include_router(robots.router, prefix="/robots", tags=["robots"])
89
router.include_router(missions.router, prefix="/missions", tags=["missions"])
Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,44 @@
1+
from typing import Annotated
2+
3+
from fastapi import APIRouter, Depends, Request
4+
from fastapi.responses import StreamingResponse
5+
from pydantic import BaseModel
6+
from starlette.background import BackgroundTask
7+
8+
from app.api.dependencies import require_authenticated
9+
from app.core.security import CurrentUser
10+
from app.services.camera_service import CameraService
11+
12+
router = APIRouter()
13+
14+
15+
class CameraResponse(BaseModel):
16+
camera_id: str = "primary"
17+
name: str = "Front camera"
18+
configured: bool
19+
transport: str = "mjpeg"
20+
21+
22+
@router.get("/primary", response_model=CameraResponse)
23+
async def camera_status(
24+
request: Request, _: Annotated[CurrentUser, Depends(require_authenticated)]
25+
) -> CameraResponse:
26+
return CameraResponse(configured=bool(request.app.state.camera_service.url))
27+
28+
29+
@router.get("/primary/stream")
30+
async def camera_stream(
31+
request: Request, _: Annotated[CurrentUser, Depends(require_authenticated)]
32+
) -> StreamingResponse:
33+
service: CameraService = request.app.state.camera_service
34+
upstream = await service.open()
35+
return StreamingResponse(
36+
service.frames(upstream),
37+
background=BackgroundTask(service.release, upstream),
38+
media_type="multipart/x-mixed-replace; boundary=frame",
39+
headers={
40+
"Cache-Control": "no-store",
41+
"X-Accel-Buffering": "no",
42+
"X-Content-Type-Options": "nosniff",
43+
},
44+
)

‎web_app/backend/app/core/config.py‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
from functools import lru_cache
22
from typing import Annotated, Literal
3+
from urllib.parse import urlsplit
34
from uuid import UUID
45

56
from pydantic import Field, SecretStr, field_validator, model_validator
@@ -27,6 +28,19 @@ class Settings(BaseSettings):
2728
websocket_queue_size: int = Field(default=100, ge=5, le=1000)
2829
websocket_max_dropped_messages: int = Field(default=25, ge=1, le=1000)
2930
websocket_auth_timeout_seconds: float = Field(default=5.0, gt=0, le=30)
31+
camera_stream_url: str = ""
32+
camera_timeout_seconds: float = Field(default=8.0, ge=1, le=30)
33+
camera_max_viewers: int = Field(default=4, ge=1, le=20)
34+
camera_max_frame_bytes: int = Field(default=4_194_304, ge=65_536, le=16_777_216)
35+
36+
@field_validator("camera_stream_url")
37+
@classmethod
38+
def validate_camera_url(cls, value: str) -> str:
39+
if value:
40+
url = urlsplit(value)
41+
if url.scheme not in {"http", "https"} or not url.hostname or url.fragment:
42+
raise ValueError("CAMERA_STREAM_URL must be an HTTP(S) MJPEG endpoint")
43+
return value
3044

3145
model_config = SettingsConfigDict(
3246
env_file=".env", env_file_encoding="utf-8", extra="ignore", case_sensitive=False

‎web_app/backend/app/main.py‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
SensorStatusRepository,
2929
)
3030
from app.services.audit_service import AuditService
31+
from app.services.camera_service import CameraService
3132
from app.services.telemetry_service import TelemetryService
3233
from app.ugv_integration.bridge import RosbridgeUGVBridge
3334
from app.ugv_integration.ingestion import UGVIngestionCoordinator
@@ -81,6 +82,11 @@ async def lifespan(application: FastAPI) -> AsyncIterator[None]:
8182
timeout_seconds=app_settings.ugv_connection_timeout_seconds,
8283
)
8384
application.state.settings = app_settings
85+
camera = CameraService(
86+
app_settings.camera_stream_url, app_settings.camera_timeout_seconds,
87+
app_settings.camera_max_viewers, app_settings.camera_max_frame_bytes,
88+
)
89+
application.state.camera_service = camera
8490
application.state.db_client = database
8591
application.state.auth_client = auth
8692
application.state.websocket_manager = manager
@@ -111,6 +117,7 @@ async def lifespan(application: FastAPI) -> AsyncIterator[None]:
111117
try:
112118
yield
113119
finally:
120+
await camera.close()
114121
if ingestion is not None:
115122
await ingestion.stop()
116123
else:
Lines changed: 105 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,105 @@
1+
import asyncio
2+
from collections.abc import AsyncIterator
3+
4+
import httpx
5+
6+
from app.core.errors import AppError
7+
8+
9+
class CameraService:
10+
"""Bounded, read-only relay of the operator-configured MJPEG source."""
11+
12+
def __init__(self, url: str, timeout: float, max_viewers: int, max_frame_bytes: int) -> None:
13+
self.url = url
14+
self.max_frame_bytes = max_frame_bytes
15+
self.slots = asyncio.Semaphore(max_viewers)
16+
self._responses: set[httpx.Response] = set()
17+
self.client = httpx.AsyncClient(
18+
timeout=httpx.Timeout(timeout), follow_redirects=False, trust_env=False
19+
)
20+
21+
async def open(self) -> httpx.Response:
22+
if not self.url:
23+
raise AppError(
24+
"CAMERA_NOT_CONFIGURED", "The camera source is not configured.", status_code=503
25+
)
26+
if self.slots.locked():
27+
raise AppError("CAMERA_BUSY", "All camera viewing slots are in use.", status_code=503)
28+
await self.slots.acquire()
29+
response = None
30+
try:
31+
response = await self.client.send(
32+
self.client.build_request(
33+
"GET", self.url, headers={"Accept": "multipart/x-mixed-replace"}
34+
),
35+
stream=True,
36+
)
37+
response.raise_for_status()
38+
if (
39+
not response.headers.get("content-type", "")
40+
.lower()
41+
.startswith("multipart/x-mixed-replace")
42+
):
43+
raise ValueError("Expected MJPEG")
44+
self._responses.add(response)
45+
return response
46+
except BaseException as exc:
47+
if response is not None:
48+
await response.aclose()
49+
self.slots.release()
50+
if isinstance(exc, (httpx.HTTPError, ValueError)):
51+
raise AppError(
52+
"CAMERA_UNAVAILABLE", "The camera source is unavailable.", status_code=502
53+
) from None
54+
raise
55+
56+
async def frames(self, response: httpx.Response) -> AsyncIterator[bytes]:
57+
# Normalize multipart boundaries and discard headers from the private upstream.
58+
buffer = bytearray()
59+
try:
60+
async for chunk in response.aiter_bytes():
61+
buffer.extend(chunk)
62+
while True:
63+
start = buffer.find(b"\xff\xd8")
64+
if start < 0:
65+
if len(buffer) > self.max_frame_bytes:
66+
return
67+
break
68+
if start:
69+
del buffer[:start]
70+
end = buffer.find(b"\xff\xd9", 2)
71+
if end < 0:
72+
if len(buffer) > self.max_frame_bytes:
73+
return
74+
break
75+
if end + 2 > self.max_frame_bytes:
76+
return
77+
frame = bytes(buffer[: end + 2])
78+
del buffer[: end + 2]
79+
yield (
80+
b"--frame\r\nContent-Type: image/jpeg\r\nContent-Length: "
81+
+ str(len(frame)).encode()
82+
+ b"\r\n\r\n"
83+
+ frame
84+
+ b"\r\n"
85+
)
86+
except httpx.HTTPError:
87+
# The response has started; closing tells the viewer to reconnect.
88+
return
89+
finally:
90+
await self.release(response)
91+
92+
async def release(self, response: httpx.Response) -> None:
93+
"""Idempotent cleanup, including responses cancelled before iteration starts."""
94+
if response not in self._responses:
95+
return
96+
self._responses.remove(response)
97+
try:
98+
await response.aclose()
99+
finally:
100+
self.slots.release()
101+
102+
async def close(self) -> None:
103+
for response in tuple(self._responses):
104+
await self.release(response)
105+
await self.client.aclose()

0 commit comments

Comments
 (0)