|
14 | 14 |
|
15 | 15 | from app import AUTHORIZATION, POSTGRES_PORT_WRITE, VERSIONING |
16 | 16 | from app.db.asyncpg_db import get_pool, get_pool_w |
17 | | -from app.utils.utils import require_json_content_type, validate_payload_keys |
| 17 | +from app.rbac_roles import check_create_permission |
| 18 | +from app.utils.utils import validate_payload_keys, require_json_content_type |
18 | 19 | from app.v1.endpoints.functions import set_role |
19 | 20 | from asyncpg.exceptions import InsufficientPrivilegeError |
20 | 21 | from fastapi import APIRouter, Body, Depends, Header, Request, status |
@@ -85,6 +86,18 @@ async def create_datastream( |
85 | 86 | current_user=user, |
86 | 87 | pool=Depends(get_pool_w) if POSTGRES_PORT_WRITE else Depends(get_pool), |
87 | 88 | ): |
| 89 | + if current_user is not None: |
| 90 | + user_role = current_user.get("role", "") |
| 91 | + if not check_create_permission(user_role): |
| 92 | + return JSONResponse( |
| 93 | + status_code=status.HTTP_403_FORBIDDEN, |
| 94 | + content={ |
| 95 | + "code": 403, |
| 96 | + "type": "error", |
| 97 | + "message": "Insufficient privileges.", |
| 98 | + }, |
| 99 | + ) |
| 100 | + |
88 | 101 | try: |
89 | 102 | require_json_content_type(request) |
90 | 103 |
|
@@ -163,6 +176,18 @@ async def create_datastream_for_thing( |
163 | 176 | current_user=user, |
164 | 177 | pool=Depends(get_pool_w) if POSTGRES_PORT_WRITE else Depends(get_pool), |
165 | 178 | ): |
| 179 | + if current_user is not None: |
| 180 | + user_role = current_user.get("role", "") |
| 181 | + if not check_create_permission(user_role): |
| 182 | + return JSONResponse( |
| 183 | + status_code=status.HTTP_403_FORBIDDEN, |
| 184 | + content={ |
| 185 | + "code": 403, |
| 186 | + "type": "error", |
| 187 | + "message": "Insufficient privileges.", |
| 188 | + }, |
| 189 | + ) |
| 190 | + |
166 | 191 | try: |
167 | 192 | require_json_content_type(request) |
168 | 193 |
|
@@ -246,6 +271,18 @@ async def create_datastream_for_sensor( |
246 | 271 | current_user=user, |
247 | 272 | pool=Depends(get_pool_w) if POSTGRES_PORT_WRITE else Depends(get_pool), |
248 | 273 | ): |
| 274 | + if current_user is not None: |
| 275 | + user_role = current_user.get("role", "") |
| 276 | + if not check_create_permission(user_role): |
| 277 | + return JSONResponse( |
| 278 | + status_code=status.HTTP_403_FORBIDDEN, |
| 279 | + content={ |
| 280 | + "code": 403, |
| 281 | + "type": "error", |
| 282 | + "message": "Insufficient privileges.", |
| 283 | + }, |
| 284 | + ) |
| 285 | + |
249 | 286 | try: |
250 | 287 | require_json_content_type(request) |
251 | 288 |
|
@@ -332,12 +369,24 @@ async def create_datastream_for_observed_property( |
332 | 369 | current_user=user, |
333 | 370 | pool=Depends(get_pool_w) if POSTGRES_PORT_WRITE else Depends(get_pool), |
334 | 371 | ): |
| 372 | + if current_user is not None: |
| 373 | + user_role = current_user.get("role", "") |
| 374 | + if not check_create_permission(user_role): |
| 375 | + return JSONResponse( |
| 376 | + status_code=status.HTTP_403_FORBIDDEN, |
| 377 | + content={ |
| 378 | + "code": 403, |
| 379 | + "type": "error", |
| 380 | + "message": "Insufficient privileges.", |
| 381 | + }, |
| 382 | + ) |
| 383 | + |
335 | 384 | try: |
336 | 385 | require_json_content_type(request) |
337 | 386 |
|
338 | 387 | if not observed_property_id: |
339 | 388 | raise Exception("Observed Property ID is required.") |
340 | | - |
| 389 | + |
341 | 390 | payload["ObservedProperty"] = {"@iot.id": observed_property_id} |
342 | 391 |
|
343 | 392 | validate_payload_keys(payload, ALLOWED_KEYS) |
|
0 commit comments