This commit is contained in:
0
backend/infrastructure/__init__.py
Normal file
0
backend/infrastructure/__init__.py
Normal file
43
backend/infrastructure/paperless_client.py
Normal file
43
backend/infrastructure/paperless_client.py
Normal file
@@ -0,0 +1,43 @@
|
||||
# backend/infrastructure/paperless_client.py
|
||||
|
||||
import requests
|
||||
from django.conf import settings
|
||||
|
||||
class PaperlessGateway:
|
||||
"""
|
||||
Gateway za komunikaciju s Paperless-ngx API-jem.
|
||||
"""
|
||||
def __init__(self):
|
||||
self.base_url = settings.PAPERLESS_API_URL
|
||||
self.headers = {
|
||||
"Authorization": f"Token {settings.PAPERLESS_API_TOKEN}",
|
||||
"Content-Type": "application/json"
|
||||
}
|
||||
|
||||
def upload_document(self, file_content, filename):
|
||||
"""Šalje dokument u Paperless."""
|
||||
url = f"{self.base_url}/documents/post_document/"
|
||||
files = {'document': (filename, file_content)}
|
||||
response = requests.post(url, headers={"Authorization": self.headers["Authorization"]}, files=files)
|
||||
response.raise_for_status()
|
||||
return response.json()
|
||||
|
||||
def get_document_status(self, task_id):
|
||||
"""Provjerava status procesiranja dokumenta."""
|
||||
url = f"{self.base_url}/tasks/{task_id}/"
|
||||
response = requests.get(url, headers=self.headers)
|
||||
return response.json()
|
||||
|
||||
def get_document_metadata(self, document_id: int):
|
||||
"""Dohvaća metapodatke dokumenta iz Paperless-a."""
|
||||
url = f"{self.base_url}/api/documents/{document_id}/"
|
||||
response = requests.get(url, headers=self.headers)
|
||||
response.raise_for_status()
|
||||
return response.json()
|
||||
|
||||
def fetch_unprocessed_documents(self):
|
||||
"""Dohvaća listu dokumenata koji čekaju obradu."""
|
||||
# Filtriramo dokumente koji nemaju naš 'processed' tag
|
||||
url = f"{self.base_url}/api/documents/?query=is_tagged:false"
|
||||
response = requests.get(url, headers=self.headers)
|
||||
return response.json().get('results', [])
|
||||
221
backend/infrastructure/pusher_service.py
Normal file
221
backend/infrastructure/pusher_service.py
Normal file
@@ -0,0 +1,221 @@
|
||||
# backend/modules/fleet/services/pusher_service.py
|
||||
"""
|
||||
Pusher Integration Service for Real-time Notifications
|
||||
Handles broadcasting notifications to users via WebSocket channels
|
||||
"""
|
||||
|
||||
import logging
|
||||
import json
|
||||
from typing import Optional, Dict, Any
|
||||
from django.conf import settings
|
||||
from django.core.serializers.json import DjangoJSONEncoder
|
||||
|
||||
try:
|
||||
import pusher
|
||||
PUSHER_AVAILABLE = True
|
||||
except ImportError:
|
||||
PUSHER_AVAILABLE = False
|
||||
pusher = None
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class PusherService:
|
||||
"""
|
||||
Wrapper service for Pusher integration.
|
||||
Handles event broadcasting, channel authentication, and error handling.
|
||||
"""
|
||||
|
||||
_instance = None
|
||||
_client = None
|
||||
|
||||
def __new__(cls):
|
||||
"""Singleton pattern to reuse Pusher client"""
|
||||
if cls._instance is None:
|
||||
cls._instance = super(PusherService, cls).__new__(cls)
|
||||
return cls._instance
|
||||
|
||||
@classmethod
|
||||
def get_client(cls):
|
||||
"""Get or initialize Pusher client"""
|
||||
if cls._client is None:
|
||||
if not PUSHER_AVAILABLE:
|
||||
logger.warning("Pusher library not installed. Install with: pip install pusher")
|
||||
return None
|
||||
|
||||
try:
|
||||
pusher_config = {
|
||||
'app_id': settings.PUSHER_APP_ID,
|
||||
'key': settings.PUSHER_KEY,
|
||||
'secret': settings.PUSHER_SECRET,
|
||||
'cluster': settings.PUSHER_CLUSTER,
|
||||
'ssl': True,
|
||||
'json_encoder': DjangoJSONEncoder,
|
||||
}
|
||||
cls._client = pusher.Pusher(**pusher_config)
|
||||
except AttributeError as e:
|
||||
logger.error(f"Pusher configuration missing: {e}")
|
||||
return None
|
||||
|
||||
return cls._client
|
||||
|
||||
@staticmethod
|
||||
def is_available() -> bool:
|
||||
"""Check if Pusher is available and configured"""
|
||||
if not PUSHER_AVAILABLE:
|
||||
return False
|
||||
|
||||
required_values = [
|
||||
settings.PUSHER_APP_ID,
|
||||
settings.PUSHER_KEY,
|
||||
settings.PUSHER_SECRET,
|
||||
settings.PUSHER_CLUSTER,
|
||||
]
|
||||
return all(required_values)
|
||||
|
||||
@staticmethod
|
||||
def get_private_channel_name(user_id: int) -> str:
|
||||
"""Generate private channel name for user"""
|
||||
return f"private-user-{user_id}"
|
||||
|
||||
@staticmethod
|
||||
def get_presence_channel_name() -> str:
|
||||
"""Get the presence channel name for fleet activity"""
|
||||
return "presence-fleet"
|
||||
|
||||
@classmethod
|
||||
def broadcast_notification(
|
||||
cls,
|
||||
user_id: int,
|
||||
title: str,
|
||||
message: str,
|
||||
level: str = 'info',
|
||||
notification_id: Optional[Any] = None,
|
||||
metadata: Optional[Dict[str, Any]] = None
|
||||
) -> bool:
|
||||
"""
|
||||
Broadcast a notification to a specific user via private channel.
|
||||
|
||||
Args:
|
||||
user_id: ID of the user to notify
|
||||
title: Notification title
|
||||
message: Notification message
|
||||
level: Notification level (info, warning, critical)
|
||||
notification_id: ID of the VehicleNotification record
|
||||
metadata: Additional metadata to send with notification
|
||||
|
||||
Returns:
|
||||
True if broadcast successful, False otherwise
|
||||
"""
|
||||
if not cls.is_available():
|
||||
logger.debug("Pusher not available, skipping broadcast")
|
||||
return False
|
||||
|
||||
client = cls.get_client()
|
||||
if not client:
|
||||
return False
|
||||
|
||||
channel_name = cls.get_private_channel_name(user_id)
|
||||
event_name = "vehicle.notification.created"
|
||||
|
||||
from django.utils.timezone import now
|
||||
event_data = {
|
||||
'id': cls._to_json_safe(notification_id),
|
||||
'title': title,
|
||||
'message': message,
|
||||
'level': level,
|
||||
'is_read': False,
|
||||
'created_at': now().isoformat(),
|
||||
}
|
||||
|
||||
if metadata:
|
||||
event_data['metadata'] = cls._to_json_safe(metadata)
|
||||
|
||||
try:
|
||||
response = client.trigger(
|
||||
channels=channel_name,
|
||||
event_name=event_name,
|
||||
data=event_data
|
||||
)
|
||||
logger.info(f"Notification broadcasted to {channel_name}: {event_name}")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.exception(f"Error broadcasting notification to {channel_name}: {e}")
|
||||
return False
|
||||
|
||||
@staticmethod
|
||||
def _to_json_safe(value: Any) -> Any:
|
||||
"""
|
||||
Convert complex Python/Django values (UUID, date, Decimal, etc.) into
|
||||
JSON-serializable primitives using DjangoJSONEncoder.
|
||||
"""
|
||||
return json.loads(json.dumps(value, cls=DjangoJSONEncoder))
|
||||
|
||||
@classmethod
|
||||
def broadcast_to_multiple_users(
|
||||
cls,
|
||||
user_ids: list,
|
||||
title: str,
|
||||
message: str,
|
||||
level: str = 'info',
|
||||
notification_ids: Optional[Dict[int, int]] = None
|
||||
) -> Dict[int, bool]:
|
||||
"""
|
||||
Broadcast the same notification to multiple users.
|
||||
|
||||
Args:
|
||||
user_ids: List of user IDs to notify
|
||||
title: Notification title
|
||||
message: Notification message
|
||||
level: Notification level
|
||||
notification_ids: Dict mapping user_id to notification_id
|
||||
|
||||
Returns:
|
||||
Dict of user_id -> success boolean
|
||||
"""
|
||||
results = {}
|
||||
for user_id in user_ids:
|
||||
notification_id = notification_ids.get(user_id) if notification_ids else None
|
||||
results[user_id] = cls.broadcast_notification(
|
||||
user_id=user_id,
|
||||
title=title,
|
||||
message=message,
|
||||
level=level,
|
||||
notification_id=notification_id
|
||||
)
|
||||
return results
|
||||
|
||||
@classmethod
|
||||
def authenticate_channel(cls, socket_id: str, channel_name: str, user_id: int) -> Optional[Dict[str, Any]]:
|
||||
"""
|
||||
Authenticate a user for a private or presence channel.
|
||||
|
||||
Args:
|
||||
socket_id: Pusher socket ID from client
|
||||
channel_name: Channel name to authenticate
|
||||
user_id: Current user ID
|
||||
|
||||
Returns:
|
||||
Signed authentication data or None if authentication fails
|
||||
"""
|
||||
if not cls.is_available():
|
||||
logger.debug("Pusher not available for authentication")
|
||||
return None
|
||||
|
||||
client = cls.get_client()
|
||||
if not client:
|
||||
return None
|
||||
|
||||
# Only allow subscription to user's own private channel
|
||||
expected_channel = cls.get_private_channel_name(user_id)
|
||||
if channel_name != expected_channel and channel_name != cls.get_presence_channel_name():
|
||||
logger.warning(f"Unauthorized channel subscription attempt: {channel_name} by user {user_id}")
|
||||
return None
|
||||
|
||||
try:
|
||||
auth_data = client.authenticate(channel=channel_name, socket_id=socket_id)
|
||||
logger.info(f"Channel {channel_name} authenticated for user {user_id}")
|
||||
return auth_data
|
||||
except Exception as e:
|
||||
logger.exception(f"Error authenticating channel {channel_name}: {e}")
|
||||
return None
|
||||
0
backend/infrastructure/webhooks.py
Normal file
0
backend/infrastructure/webhooks.py
Normal file
Reference in New Issue
Block a user