Files
Django/gvsdsdk/msync/client.py

276 lines
8.8 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""MSYNC 服务客户端模块
提供 ServiceClient 类,用于微服务间发起认证请求。
支持自动令牌获取与刷新、OGM 解析、服务发现。
本模块自包含实现(使用 urllib不依赖 Django。
"""
import json
import time
import logging
import urllib.request
import urllib.error
logger = logging.getLogger(__name__)
MSYNC_AUTH_PATH = '/api/msync/auth/token/'
MSYNC_AUTH_REFRESH_PATH = '/api/msync/auth/token/refresh/'
MSYNC_OGM_RESOLVE_PATH = '/api/msync/ogm/{global_id}/'
MSYNC_REGISTRY_PATH = '/api/msync/registry/{name}/'
from gvsdsdk.msync.auth import SignPayload, BuildAuthPayload
def _GenerateChallenge():
return f'{int(time.time() * 1000)}-{id(object())}'
class ServiceClientError(Exception):
def __init__(self, Message, StatusCode=None, ResponseBody=None):
super().__init__(Message)
self.StatusCode = StatusCode
self.ResponseBody = ResponseBody
class _TokenStore:
"""内存中的令牌存储"""
def __init__(self):
self._AccessToken = None
self._RefreshToken = None
self._ExpiresAt = 0.0
self._Scope = []
@property
def AccessToken(self):
return self._AccessToken
@property
def RefreshToken(self):
return self._RefreshToken
@property
def Scope(self):
return self._Scope
def IsValid(self):
return self._AccessToken is not None and time.time() < self._ExpiresAt
def Store(self, AccessToken, RefreshToken, ExpiresAt, Scope=None):
self._AccessToken = AccessToken
self._RefreshToken = RefreshToken
self._ExpiresAt = ExpiresAt
self._Scope = Scope or []
def Clear(self):
self._AccessToken = None
self._RefreshToken = None
self._ExpiresAt = 0.0
self._Scope = []
class ServiceClient:
"""微服务通信客户端——用于向其他服务发起认证请求
用法:
client = ServiceClient(
ServiceName='gerp',
ServiceSecret='your-secret',
)
client.Authenticate()
# 获取其他服务信息
info = client.DiscoverService('crb')
# 解析 OGM
ogm = client.ResolveOGM('user:abc123')
# 向其他服务发请求
response = client.Request('GET', 'https://crb.gvsds.com/api/wallet', Scope=['ledger:read'])
"""
def __init__(self, ServiceName, ServiceSecret, MasterURL=None):
self.ServiceName = ServiceName
self.ServiceSecret = ServiceSecret
if MasterURL is None:
MasterURL = 'https://api.gvsds.com'
self.MasterURL = MasterURL.rstrip('/')
self._tokens = _TokenStore()
@property
def IsAuthenticated(self):
return self._tokens.IsValid()
@property
def AccessToken(self):
return self._tokens.AccessToken
@property
def RefreshToken(self):
return self._tokens.RefreshToken
def Authenticate(self):
if not self._Authenticate():
raise ServiceClientError('认证失败')
return self
def Refresh(self):
if not self._RefreshToken():
raise ServiceClientError('刷新失败')
return self
def _Authenticate(self):
challenge = _GenerateChallenge()
payload = BuildAuthPayload(self.ServiceName, challenge)
signature = SignPayload(payload, self.ServiceSecret)
req_body = json.dumps({
'ServiceName': self.ServiceName,
'Payload': payload,
'Signature': signature,
}).encode('utf-8')
req = urllib.request.Request(
f'{self.MasterURL}{MSYNC_AUTH_PATH}',
data=req_body,
headers={'Content-Type': 'application/json'},
method='POST',
)
try:
with urllib.request.urlopen(req, timeout=30) as resp:
data = json.loads(resp.read().decode('utf-8'))
except urllib.error.HTTPError as e:
body = e.read().decode('utf-8', errors='replace')
logger.error(f'认证失败 HTTP {e.code}: {body}')
return False
except Exception as e:
logger.error(f'认证异常: {e}')
return False
if data.get('code') not in (200, 0):
logger.error(f'认证业务错误: {data}')
return False
token_data = data.get('data') or {}
self._tokens.Store(
AccessToken=token_data.get('AccessToken'),
RefreshToken=token_data.get('RefreshToken'),
ExpiresAt=token_data.get('ExpiresAt', 0),
Scope=token_data.get('Scope', []),
)
return self._tokens.AccessToken is not None
def _RefreshToken(self):
if not self._tokens.RefreshToken:
return False
req_body = json.dumps({
'ServiceName': self.ServiceName,
'RefreshToken': self._tokens.RefreshToken,
}).encode('utf-8')
req = urllib.request.Request(
f'{self.MasterURL}{MSYNC_AUTH_REFRESH_PATH}',
data=req_body,
headers={'Content-Type': 'application/json'},
method='POST',
)
try:
with urllib.request.urlopen(req, timeout=30) as resp:
data = json.loads(resp.read().decode('utf-8'))
except Exception as e:
logger.error(f'刷新令牌异常: {e}')
return False
if data.get('code') not in (200, 0):
return False
token_data = data.get('data') or {}
self._tokens.Store(
AccessToken=token_data.get('AccessToken'),
RefreshToken=token_data.get('RefreshToken'),
ExpiresAt=token_data.get('ExpiresAt', 0),
Scope=token_data.get('Scope', []),
)
return self._tokens.AccessToken is not None
def _EnsureAuthenticated(self):
if self.IsAuthenticated:
return
if self._tokens.RefreshToken:
self._RefreshToken()
return
self._Authenticate()
def Request(self, Method, URL, Data=None, Headers=None, Scope=None, Timeout=30):
self._EnsureAuthenticated()
req_headers = {
'Authorization': f'Service {self.AccessToken}',
'Content-Type': 'application/json',
'X-Source-Service': self.ServiceName,
}
if Scope:
req_headers['X-Required-Scope'] = ','.join(Scope)
if Headers:
req_headers.update(Headers)
body = None
if Data is not None:
if isinstance(Data, dict):
body = json.dumps(Data).encode('utf-8')
elif isinstance(Data, bytes):
body = Data
elif isinstance(Data, str):
body = Data.encode('utf-8')
req = urllib.request.Request(URL, data=body, headers=req_headers, method=Method)
try:
with urllib.request.urlopen(req, timeout=Timeout) as resp:
response_data = resp.read().decode('utf-8')
content_type = resp.headers.get('Content-Type', '')
if 'application/json' in content_type:
return json.loads(response_data)
return response_data
except urllib.error.HTTPError as e:
body_text = e.read().decode('utf-8', errors='replace')
if e.code == 401:
logger.info('令牌可能已失效,尝试刷新')
self._tokens.Clear()
self._EnsureAuthenticated()
req_headers['Authorization'] = f'Service {self.AccessToken}'
retry_req = urllib.request.Request(URL, data=body, headers=req_headers, method=Method)
try:
with urllib.request.urlopen(retry_req, timeout=Timeout) as resp2:
response_data = resp2.read().decode('utf-8')
content_type = resp2.headers.get('Content-Type', '')
if 'application/json' in content_type:
return json.loads(response_data)
return response_data
except urllib.error.HTTPError as e2:
body2 = e2.read().decode('utf-8', errors='replace')
raise ServiceClientError(
f'重试后仍失败: HTTP {e2.code}',
StatusCode=e2.code,
ResponseBody=body2,
)
raise ServiceClientError(
f'请求失败: HTTP {e.code}',
StatusCode=e.code,
ResponseBody=body_text,
)
def ResolveOGM(self, GlobalID):
url = f'{self.MasterURL}{MSYNC_OGM_RESOLVE_PATH.format(global_id=GlobalID)}'
return self.Request('GET', url, Scope=['ogm:read'])
def DiscoverService(self, ServiceName):
url = f'{self.MasterURL}{MSYNC_REGISTRY_PATH.format(name=ServiceName)}'
return self.Request('GET', url, Scope=['registry:read'])