mirror of
https://github.com/xcat2/confluent.git
synced 2026-09-29 00:31:09 +00:00
Users can request and revoke API keys
Proivde support for API keys. Users can register by doing a POST like:
/confluent-api/sessions/current/apikey/create with body of '{"expiration":null}
API keys can be reset by doing a POST to /confluent-api/sessions/current/apikey/revokeall
A secret is provisioned per user to allow them to clear all keys, while remaining mostly stateless to avoid having to manage every key server side.
This commit is contained in:
@@ -0,0 +1,147 @@
|
||||
# vim: tabstop=4 shiftwidth=4 softtabstop=4
|
||||
|
||||
# Copyright 2026 Lenovo
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
import base64
|
||||
import hashlib
|
||||
import hmac
|
||||
import json
|
||||
import math
|
||||
import secrets
|
||||
import time
|
||||
|
||||
|
||||
_SECRET_ATTRIBUTE = 'secret.apikey'
|
||||
|
||||
|
||||
def _base64url(value):
|
||||
return base64.urlsafe_b64encode(value).rstrip(b'=').decode('ascii')
|
||||
|
||||
|
||||
def _decode_base64url(value):
|
||||
if not isinstance(value, str):
|
||||
raise ValueError('Invalid bearer token')
|
||||
return base64.urlsafe_b64decode(value + '=' * (-len(value) % 4))
|
||||
|
||||
|
||||
def _make_jws(username, secret, expiration=None):
|
||||
header = {'alg': 'HS256', 'typ': 'JWT'}
|
||||
payload = {'sub': username, 'iat': int(time.time())}
|
||||
if expiration is not None:
|
||||
payload['exp'] = payload['iat'] + int(expiration * 86400)
|
||||
encoded_header = _base64url(json.dumps(
|
||||
header, separators=(',', ':')).encode('utf8'))
|
||||
encoded_payload = _base64url(json.dumps(
|
||||
payload, separators=(',', ':')).encode('utf8'))
|
||||
signing_input = '{0}.{1}'.format(encoded_header, encoded_payload).encode('ascii')
|
||||
signature = hmac.new(secret, signing_input, hashlib.sha256).digest()
|
||||
return '{0}.{1}'.format(signing_input.decode('ascii'), _base64url(signature))
|
||||
|
||||
|
||||
def _username_string(username):
|
||||
if isinstance(username, bytes):
|
||||
return username.decode('utf8')
|
||||
return str(username)
|
||||
|
||||
|
||||
def _get_expiration(reqbody):
|
||||
if not reqbody:
|
||||
return None
|
||||
if isinstance(reqbody, bytes):
|
||||
reqbody = reqbody.decode('utf8')
|
||||
try:
|
||||
body = json.loads(reqbody)
|
||||
except (TypeError, UnicodeDecodeError, json.JSONDecodeError):
|
||||
raise ValueError('Request body must be JSON')
|
||||
if not isinstance(body, dict):
|
||||
raise ValueError('Request body must be a JSON object')
|
||||
if 'expiration' not in body:
|
||||
raise ValueError('expiration is required parameter')
|
||||
expiration = body['expiration']
|
||||
if not expiration: # a false-y expiration means opt out of expiration
|
||||
return None
|
||||
if not isinstance(expiration, (int, float)) or not math.isfinite(expiration) or expiration < 0:
|
||||
raise ValueError('Invalid number specified, must be either number of days or false/null')
|
||||
return expiration
|
||||
|
||||
|
||||
def validate_bearer_token(token, cfgmgr):
|
||||
"""Return the token subject when a bearer token is valid."""
|
||||
try:
|
||||
encoded_header, encoded_payload, encoded_signature = token.split('.')
|
||||
header = json.loads(_decode_base64url(encoded_header))
|
||||
payload = json.loads(_decode_base64url(encoded_payload))
|
||||
signature = _decode_base64url(encoded_signature)
|
||||
except (AttributeError, ValueError, UnicodeDecodeError, json.JSONDecodeError,
|
||||
TypeError, base64.binascii.Error):
|
||||
return None
|
||||
if header != {'alg': 'HS256', 'typ': 'JWT'}:
|
||||
return None
|
||||
username = payload.get('sub')
|
||||
if not isinstance(username, str) or not username:
|
||||
return None
|
||||
expiration = payload.get('exp')
|
||||
if expiration is not None and (
|
||||
isinstance(expiration, bool) or
|
||||
not isinstance(expiration, (int, float)) or
|
||||
not math.isfinite(expiration) or expiration <= time.time()):
|
||||
return None
|
||||
user = cfgmgr.get_user(username, decrypt=True)
|
||||
if not user or not user.get(_SECRET_ATTRIBUTE):
|
||||
return None
|
||||
secret = user[_SECRET_ATTRIBUTE]['value']
|
||||
if not isinstance(secret, bytes):
|
||||
secret = secret.encode('utf8')
|
||||
signing_input = '{0}.{1}'.format(
|
||||
encoded_header, encoded_payload).encode('ascii')
|
||||
expected = hmac.new(secret, signing_input, hashlib.sha256).digest()
|
||||
if not hmac.compare_digest(signature, expected):
|
||||
return None
|
||||
return username
|
||||
|
||||
async def handle_api_request(url, username, cfgmgr, reqbody):
|
||||
"""Handle an authenticated API-key request.
|
||||
|
||||
The HTTP layer supplies the authenticated request context. The return
|
||||
value is a ``(status, payload)`` pair for ``httpapi`` to serialize.
|
||||
"""
|
||||
username = _username_string(username)
|
||||
operation = url.removeprefix('/sessions/current/apikey/')
|
||||
if operation == 'create':
|
||||
try:
|
||||
expiration = _get_expiration(reqbody)
|
||||
except ValueError as error:
|
||||
return 400, {'error': str(error)}
|
||||
user = cfgmgr.get_user(username, decrypt=True)
|
||||
if user is None:
|
||||
await cfgmgr.create_user(username, role='Stub')
|
||||
user = cfgmgr.get_user(username, decrypt=True)
|
||||
secret = user.get(_SECRET_ATTRIBUTE, {}).get('value', None)
|
||||
if not secret:
|
||||
secret = secrets.token_bytes(32)
|
||||
secret = _base64url(secret)
|
||||
await cfgmgr.set_user(username, {_SECRET_ATTRIBUTE: secret})
|
||||
if isinstance(secret, str):
|
||||
secret = secret.encode('utf8')
|
||||
return 200, {'jws': _make_jws(username, secret, expiration)}
|
||||
if operation == 'revokeall':
|
||||
user = cfgmgr.get_user(username, decrypt=True)
|
||||
if user is None:
|
||||
return 200, {'revoked': True, 'msg': "User doesn't exist"}
|
||||
if not user.get(_SECRET_ATTRIBUTE):
|
||||
return 200, {'revoked': True, 'msg': "No API keys to revoke"}
|
||||
await cfgmgr.set_user(username, {_SECRET_ATTRIBUTE: None})
|
||||
return 200, {'revoked': True}
|
||||
return 404, {'error': 'Unknown API-key operation'}
|
||||
@@ -1645,7 +1645,7 @@ class ConfigManager(object):
|
||||
except KeyError:
|
||||
return []
|
||||
|
||||
def get_user(self, name):
|
||||
def get_user(self, name, decrypt=False):
|
||||
"""Get user information from DB
|
||||
|
||||
:param name: Name of the user
|
||||
@@ -1657,7 +1657,13 @@ class ConfigManager(object):
|
||||
|
||||
"""
|
||||
try:
|
||||
return copy.deepcopy(self._cfgstore['users'][name])
|
||||
ret = copy.deepcopy(self._cfgstore['users'][name])
|
||||
if decrypt:
|
||||
for key in ret:
|
||||
if isinstance(ret[key], dict) and 'cryptvalue' in ret[key]:
|
||||
ret[key]['value'] = decrypt_value(ret[key]['cryptvalue'])
|
||||
return ret
|
||||
|
||||
except KeyError:
|
||||
return None
|
||||
|
||||
@@ -1787,6 +1793,12 @@ class ConfigManager(object):
|
||||
pw = pw.encode('utf-8')
|
||||
crypted = hashlib.pbkdf2_hmac('sha256', pw, salt, 10000, dklen=32)
|
||||
user['cryptpass'] = (salt, crypted)
|
||||
elif attribute.startswith('secret.'):
|
||||
if attributemap[attribute] is None:
|
||||
if attribute in user:
|
||||
del user[attribute]
|
||||
continue
|
||||
user[attribute] = {'cryptvalue': crypt_value(attributemap[attribute])}
|
||||
else:
|
||||
user[attribute] = attributemap[attribute]
|
||||
_mark_dirtykey('users', name, self.tenant)
|
||||
@@ -2814,7 +2826,7 @@ class ConfigManager(object):
|
||||
displayname = ucfg.get('displayname', None)
|
||||
role = ucfg.get('role', None)
|
||||
await self.create_user(user, uid=uid, displayname=displayname, role=role)
|
||||
for attrname in ('webauthid', 'authenticators', 'cryptpass'):
|
||||
for attrname in ('secret.apikey', 'webauthid', 'authenticators', 'cryptpass'):
|
||||
if attrname in tmpconfig[confarea][user]:
|
||||
self._cfgstore['users'][user][attrname] = tmpconfig[confarea][user][attrname]
|
||||
_mark_dirtykey('users', user, self.tenant)
|
||||
@@ -2854,9 +2866,8 @@ class ConfigManager(object):
|
||||
for attribute in self._cfgstore[confarea][element]:
|
||||
if 'inheritedfrom' in dumpdata[confarea][element][attribute]:
|
||||
del dumpdata[confarea][element][attribute]
|
||||
elif (attribute == 'cryptpass' or
|
||||
'cryptvalue' in
|
||||
dumpdata[confarea][element][attribute]):
|
||||
elif (attribute == 'cryptpass' or (isinstance(dumpdata[confarea][element][attribute], dict) and
|
||||
'cryptvalue' in dumpdata[confarea][element][attribute])):
|
||||
if redact is not None:
|
||||
dumpdata[confarea][element][attribute] = '*REDACTED*'
|
||||
else:
|
||||
|
||||
@@ -31,6 +31,7 @@ except ImportError:
|
||||
webauthn = None
|
||||
import asyncio
|
||||
from aiohttp import web, WSMsgType
|
||||
import confluent.apikey as apikey
|
||||
import confluent.auth as auth
|
||||
import confluent.config.attributes as attribs
|
||||
import confluent.config.configmanager as configmanager
|
||||
@@ -338,24 +339,39 @@ async def _authorize_request(req, operation, reqbody):
|
||||
# of a CSRF
|
||||
return {'code': 401}
|
||||
return ('logout',)
|
||||
if req.headers['Authorization'].startswith('MultiBasic '):
|
||||
authorization = req.headers['Authorization']
|
||||
if authorization.startswith('Bearer '):
|
||||
name = apikey.validate_bearer_token(
|
||||
authorization[7:], configmanager.ConfigManager(None))
|
||||
if name:
|
||||
authdata = auth.authorize(
|
||||
name, element=element, operation=operation)
|
||||
if authdata is False:
|
||||
return {'code': 403}
|
||||
elif not authdata:
|
||||
return {'code': 401}
|
||||
else:
|
||||
return {'code': 401}
|
||||
elif authorization.startswith('MultiBasic '):
|
||||
name, passphrase = base64.b64decode(
|
||||
req.headers['Authorization'].replace('MultiBasic ', '')).split(b':', 1)
|
||||
authorization.replace('MultiBasic ', '')).split(b':', 1)
|
||||
passphrase = json.loads(passphrase)
|
||||
else:
|
||||
name, passphrase = base64.b64decode(
|
||||
req.headers['Authorization'].replace('Basic ', '')).split(b':', 1)
|
||||
try:
|
||||
authdata = await auth.check_user_passphrase(name, passphrase, operation=operation, element=element)
|
||||
except Exception as e:
|
||||
if hasattr(e, 'prompts'):
|
||||
return {'code': 403, 'prompts': e.prompts}
|
||||
raise
|
||||
if authdata is False:
|
||||
return {'code': 403}
|
||||
elif not authdata:
|
||||
return {'code': 401}
|
||||
sessid = _establish_http_session(req, authdata, name, cookie)
|
||||
authorization.replace('Basic ', '')).split(b':', 1)
|
||||
if not authdata and not authorization.startswith('Bearer '):
|
||||
try:
|
||||
authdata = await auth.check_user_passphrase(name, passphrase, operation=operation, element=element)
|
||||
except Exception as e:
|
||||
if hasattr(e, 'prompts'):
|
||||
return {'code': 403, 'prompts': e.prompts}
|
||||
raise
|
||||
if authdata is False:
|
||||
return {'code': 403}
|
||||
elif not authdata:
|
||||
return {'code': 401}
|
||||
if authdata:
|
||||
sessid = _establish_http_session(req, authdata, name, cookie)
|
||||
if authdata and element and element.startswith('/sessions/current/webauthn/validate/'):
|
||||
if not webauthn:
|
||||
raise exc.NotFoundException('WebAuthn support not available')
|
||||
@@ -1060,6 +1076,16 @@ async def resourcehandler_backend(req, make_response):
|
||||
tlvdata.unicode_dictvalues(sessinfo)
|
||||
await rsp.write(json.dumps(sessinfo).encode('utf8'))
|
||||
return rsp
|
||||
elif url.startswith('/sessions/current/apikey/'):
|
||||
if operation == 'retrieve':
|
||||
rsp = await make_response('text/plain', 405, 'Method Not Allowed')
|
||||
return rsp
|
||||
status, apidata = await apikey.handle_api_request(
|
||||
url, authorized['username'], cfgmgr, reqbody)
|
||||
rsp = await make_response('application/json', status,
|
||||
cookies=cookies)
|
||||
await rsp.write(json.dumps(apidata).encode('utf8'))
|
||||
return rsp
|
||||
elif url.startswith('/sessions/current/webauthn/'):
|
||||
if not webauthn:
|
||||
rsp = await make_response('text/plain', 501, 'Not Implemented')
|
||||
|
||||
Reference in New Issue
Block a user