Comprehensive user management integration system for the Open Source Site Tracking platform, including authentication integration, user synchronization, and API integration.
The user management integration system provides comprehensive tools for integrating with external authentication systems, synchronizing user data, and managing user lifecycle through APIs and webhooks.
# OAuth 2.0 integration
class OAuth2Integration:
def __init__(self, provider_config):
self.provider_config = provider_config
self.token_store = TokenStore()
async def get_authorization_url(self, redirect_uri, scope=None, state=None):
"""Generate OAuth authorization URL"""
auth_params = {
'client_id': self.provider_config['client_id'],
'redirect_uri': redirect_uri,
'response_type': 'code',
'scope': scope or self.provider_config['default_scope'],
'state': state or self.generate_state()
}
return f"{self.provider_config['auth_url']}?{urlencode(auth_params)}"
async def exchange_code_for_token(self, code, redirect_uri):
"""Exchange authorization code for access token"""
token_data = {
'grant_type': 'authorization_code',
'client_id': self.provider_config['client_id'],
'client_secret': self.provider_config['client_secret'],
'code': code,
'redirect_uri': redirect_uri
}
response = await self.http_client.post(
self.provider_config['token_url'],
data=token_data
)
token_info = response.json()
# Store token
await self.token_store.store_token(token_info)
return token_info
async def get_user_info(self, access_token):
"""Get user information from OAuth provider"""
headers = {'Authorization': f'Bearer {access_token}'}
response = await self.http_client.get(
self.provider_config['user_info_url'],
headers=headers
)
return response.json()
# SAML integration
class SAMLIntegration:
def __init__(self, saml_config):
self.saml_config = saml_config
self.idp_metadata = self.load_idp_metadata()
async def create_auth_request(self, relay_state=None):
"""Create SAML authentication request"""
auth_request = self.create_saml_auth_request()
if relay_state:
auth_request.relay_state = relay_state
# Sign request if configured
if self.saml_config.get('sign_requests'):
auth_request.sign()
# Redirect to IdP
redirect_url = self.build_redirect_url(auth_request)
return redirect_url
async def process_response(self, saml_response):
"""Process SAML response from IdP"""
# Validate response
if not self.validate_response(saml_response):
raise ValueError("Invalid SAML response")
# Extract attributes
attributes = self.extract_attributes(saml_response)
# Create or update user
user = await self.create_or_update_user(attributes)
return user
def validate_response(self, saml_response):
"""Validate SAML response"""
# Check signature
if not self.verify_signature(saml_response):
return False
# Check conditions
if not self.check_conditions(saml_response):
return False
# Check audience
if not self.check_audience(saml_response):
return False
return True
# User synchronization service
class UserSynchronizationService:
def __init__(self, sync_config):
self.sync_config = sync_config
self.user_store = UserStore()
self.logger = logging.getLogger(__name__)
async def sync_users(self, source_system, sync_type='full'):
"""Synchronize users from source system"""
try:
# Get users from source system
source_users = await self.get_source_users(source_system)
# Get existing users
existing_users = await self.get_existing_users()
# Determine sync actions
sync_actions = self.determine_sync_actions(
source_users, existing_users, sync_type
)
# Execute sync actions
results = await self.execute_sync_actions(sync_actions)
# Log results
await self.log_sync_results(results)
return results
except Exception as e:
self.logger.error(f"User synchronization failed: {e}")
raise
def determine_sync_actions(self, source_users, existing_users, sync_type):
"""Determine synchronization actions"""
actions = []
for source_user in source_users:
existing_user = existing_users.get(source_user['id'])
if not existing_user:
# New user - create
actions.append({
'action': 'create',
'user_data': source_user
})
elif sync_type == 'full' or self.has_user_changed(source_user, existing_user):
# Existing user - update
actions.append({
'action': 'update',
'user_id': source_user['id'],
'user_data': source_user,
'existing_data': existing_user
})
if sync_type == 'full':
# Handle deactivated users
for existing_user_id, existing_user in existing_users.items():
if existing_user_id not in [u['id'] for u in source_users]:
actions.append({
'action': 'deactivate',
'user_id': existing_user_id
})
return actions
async def execute_sync_actions(self, sync_actions):
"""Execute synchronization actions"""
results = {
'created': 0,
'updated': 0,
'deactivated': 0,
'errors': []
}
for action in sync_actions:
try:
if action['action'] == 'create':
await self.create_user(action['user_data'])
results['created'] += 1
elif action['action'] == 'update':
await self.update_user(
action['user_id'],
action['user_data'],
action['existing_data']
)
results['updated'] += 1
elif action['action'] == 'deactivate':
await self.deactivate_user(action['user_id'])
results['deactivated'] += 1
except Exception as e:
results['errors'].append({
'action': action['action'],
'user_id': action.get('user_id'),
'error': str(e)
})
return results
# REST API integration client
class UserManagementAPIClient:
def __init__(self, api_config):
self.api_config = api_config
self.http_client = httpx.AsyncClient(
base_url=api_config['base_url'],
headers={
'Authorization': f"Bearer {api_config['api_token']}",
'Content-Type': 'application/json'
}
)
async def create_user(self, user_data):
"""Create user via API"""
try:
response = await self.http_client.post(
'/api/v1/users',
json=user_data
)
response.raise_for_status()
return response.json()
except httpx.HTTPError as e:
raise APIError(f"Failed to create user: {e}")
async def update_user(self, user_id, user_data):
"""Update user via API"""
try:
response = await self.http_client.put(
f'/api/v1/users/{user_id}',
json=user_data
)
response.raise_for_status()
return response.json()
except httpx.HTTPError as e:
raise APIError(f"Failed to update user: {e}")
async def get_user(self, user_id):
"""Get user via API"""
try:
response = await self.http_client.get(f'/api/v1/users/{user_id}')
response.raise_for_status()
return response.json()
except httpx.HTTPError as e:
raise APIError(f"Failed to get user: {e}")
async def delete_user(self, user_id):
"""Delete user via API"""
try:
response = await self.http_client.delete(f'/api/v1/users/{user_id}')
response.raise_for_status()
return True
except httpx.HTTPError as e:
raise APIError(f"Failed to delete user: {e}")
async def list_users(self, filters=None, pagination=None):
"""List users via API"""
params = {}
if filters:
params.update(filters)
if pagination:
params.update(pagination)
try:
response = await self.http_client.get('/api/v1/users', params=params)
response.raise_for_status()
return response.json()
except httpx.HTTPError as e:
raise APIError(f"Failed to list users: {e}")
# Webhook system for user management
class UserManagementWebhookSystem:
def __init__(self, webhook_config):
self.webhook_config = webhook_config
self.event_handlers = {}
self.retry_policy = webhook_config.get('retry_policy', {})
def register_event_handler(self, event_type, handler):
"""Register event handler"""
if event_type not in self.event_handlers:
self.event_handlers[event_type] = []
self.event_handlers[event_type].append(handler)
async def trigger_webhook(self, event_type, payload):
"""Trigger webhook for event"""
# Get webhook URLs for event type
webhook_urls = self.get_webhook_urls(event_type)
if not webhook_urls:
return
# Prepare webhook payload
webhook_payload = {
'event_type': event_type,
'timestamp': datetime.utcnow().isoformat(),
'data': payload
}
# Send webhooks
tasks = [
self.send_webhook(url, webhook_payload)
for url in webhook_urls
]
results = await asyncio.gather(*tasks, return_exceptions=True)
return results
async def send_webhook(self, url, payload, attempt=1):
"""Send webhook with retry logic"""
try:
async with httpx.AsyncClient() as client:
response = await client.post(
url,
json=payload,
headers={
'Content-Type': 'application/json',
'X-Webhook-Signature': self.generate_signature(payload)
},
timeout=self.webhook_config.get('timeout', 30)
)
response.raise_for_status()
return {'url': url, 'status': 'success', 'attempt': attempt}
except Exception as e:
if attempt < self.retry_policy.get('max_attempts', 3):
await asyncio.sleep(self.get_retry_delay(attempt))
return await self.send_webhook(url, payload, attempt + 1)
return {
'url': url,
'status': 'error',
'error': str(e),
'attempt': attempt
}
def generate_signature(self, payload):
"""Generate webhook signature"""
secret = self.webhook_config.get('secret')
if not secret:
return None
payload_string = json.dumps(payload, sort_keys=True)
signature = hmac.new(
secret.encode(),
payload_string.encode(),
hashlib.sha256
).hexdigest()
return f"sha256={signature}"