"""Authentication command business logic."""
from __future__ import annotations
import base64
import json
from collections.abc import Callable
from dataclasses import asdict, dataclass
from typing import Optional, cast
from urllib.parse import urlparse
import structlog
from packaging.version import InvalidVersion, Version
from kelvin.sdk.exceptions import AuthenticationError
from kelvin.sdk.services.auth import AuthService, DeviceAuthResponse, OAuthFlowError, TokenResponse
from kelvin.sdk.services.credential_store import CredentialStore, TokenCredentials
from kelvin.sdk.services.docker import DockerService
from kelvin.sdk.services.session import SessionService
from kelvin.sdk.services.session_models import PlatformMetadata, SessionInfo
from kelvin.sdk.version import version as cli_version
logger = cast(structlog.stdlib.BoundLogger, structlog.get_logger(__name__))
# ============ Models ============
[docs]
@dataclass
class DockerLoginInfo:
"""Docker registry login status from auth flow."""
login_url: str
is_custom_registry: bool = False
login_failed: bool = False
[docs]
@dataclass
class VersionInfo:
"""CLI version compatibility info."""
current: str
recommended: str
[docs]
@dataclass
class LoginResult:
"""Result of login operation."""
url: str
docker: Optional[DockerLoginInfo] = None
version: Optional[VersionInfo] = None
#: Logged-in user's identity (the email, falling back to ``preferred_username``
#: / ``sub``). Used as the analytics identity and sent as a user property.
username: Optional[str] = None
[docs]
@dataclass
class TokenResult:
"""Result of token operation."""
access_token: str
expires_at: int
full_credentials: Optional[dict[str, object]] = None
# ============ Exceptions ============
[docs]
class AuthCommandsError(AuthenticationError):
"""Base exception for auth command operations."""
[docs]
class MissingSessionError(AuthCommandsError):
"""Raised when no active session is found."""
[docs]
class MissingCredentialsError(AuthCommandsError):
"""Raised when credentials are missing for a session."""
[docs]
class TokenRefreshError(AuthCommandsError):
"""Raised when token refresh is not possible."""
# ============ Commands ============
[docs]
class AuthCommands:
"""Authentication command business logic.
Receives complete, validated inputs from CLI layer.
Orchestrates multiple services to implement workflows.
"""
def __init__(
self,
auth: AuthService,
credentials: CredentialStore,
session: SessionService,
docker: DockerService,
) -> None:
self.auth: AuthService = auth
self.credentials: CredentialStore = credentials
self.session: SessionService = session
self.docker: DockerService = docker
[docs]
def login_browser(
self,
url: str,
open_browser: bool = True,
on_auth_url: Optional[Callable[[str], None]] = None,
) -> LoginResult:
"""Browser SSO login workflow."""
logger.info("starting browser login", url=url)
tokens = self.auth.login_browser(url, open_browser=open_browser, on_auth_url=on_auth_url)
_ = self._store_tokens(url, tokens, oauth_client_id=AuthService.BROWSER_CLIENT_ID)
metadata = self._fetch_metadata(url)
docker_username = self._get_token_username(tokens.access_token)
identity = self._get_token_identity(tokens.access_token)
session = self.session.set_current_session(url, metadata=metadata, username=identity)
return LoginResult(
url=session.url,
docker=self._login_docker(docker_username, tokens.access_token),
version=self._get_version_info(metadata),
username=identity,
)
[docs]
def login_device_code(
self,
url: str,
open_browser: bool = True,
on_device_code: Optional[Callable[[DeviceAuthResponse], None]] = None,
) -> LoginResult:
"""Device code login workflow (RFC 8628)."""
logger.info("starting device code login", url=url)
tokens = self.auth.login_device_code(url, open_browser=open_browser, on_device_code=on_device_code)
_ = self._store_tokens(url, tokens, oauth_client_id=AuthService.BROWSER_CLIENT_ID)
metadata = self._fetch_metadata(url)
docker_username = self._get_token_username(tokens.access_token)
identity = self._get_token_identity(tokens.access_token)
session = self.session.set_current_session(url, metadata=metadata, username=identity)
return LoginResult(
url=session.url,
docker=self._login_docker(docker_username, tokens.access_token),
version=self._get_version_info(metadata),
username=identity,
)
[docs]
def login_password(
self,
url: str,
username: str,
password: str,
totp: Optional[str] = None,
) -> LoginResult:
"""Password login workflow."""
logger.info("starting password login", url=url, username=username)
tokens = self.auth.login_password(url, username, password, totp)
logger.debug("password login successful", url=url, username=username)
_ = self._store_tokens(url, tokens, oauth_client_id=AuthService.USER_CLIENT_ID)
metadata = self._fetch_metadata(url)
identity = self._get_token_identity(tokens.access_token)
session = self.session.set_current_session(url, metadata=metadata, username=identity)
return LoginResult(
url=session.url,
docker=self._login_docker(username, password),
version=self._get_version_info(metadata),
username=identity,
)
[docs]
def login_with_stored_credentials(self, url: str) -> Optional[LoginResult]:
"""Attempt login using stored keyring credentials.
Checks for stored credentials, validates or refreshes tokens,
and verifies against Keycloak. Returns None if stored credentials
are unavailable or invalid, signaling the caller to prompt.
"""
logger.info("attempting login with stored credentials", url=url)
creds = self.credentials.retrieve(url)
if creds is None:
logger.debug("no stored credentials found", url=url)
return None
if self.credentials.is_token_expired(creds):
logger.debug("access token expired, attempting refresh", url=url)
if not creds.refresh_token or self.credentials.is_refresh_token_expired(creds):
logger.debug("refresh token unavailable or expired", url=url)
return None
try:
tokens = self.auth.refresh_tokens(url, creds.refresh_token, oauth_client_id=creds.oauth_client_id)
creds = self._store_tokens(url, tokens, oauth_client_id=creds.oauth_client_id)
except OAuthFlowError:
logger.debug("token refresh failed", url=url)
return None
docker_username: Optional[str] = None
verified_username: Optional[str] = None
verified_identity: Optional[str] = None
if not creds.is_service_account:
try:
token_claims = self.auth.verify_token(url, creds.access_token)
verified_username = str(token_claims.get("preferred_username") or "") or None
# Identity is the email, falling back to preferred_username / sub.
verified_identity = (
str(token_claims.get("email") or "") or verified_username or str(token_claims.get("sub") or "")
) or None
docker_username = verified_username or str(token_claims.get("sub") or "") or ""
except OAuthFlowError:
logger.debug("stored token is invalid on server", url=url)
return None
else:
logger.debug("service account token, skipping remote verification", url=url)
metadata = self._fetch_metadata(url)
# Prefer the server-verified claims; fall back to local JWT decode.
identity = verified_identity or self._get_token_identity(creds.access_token)
session = self.session.set_current_session(url, metadata=metadata, username=identity)
logger.info("login with stored credentials successful", url=url)
if docker_username is None:
docker_username = self._get_token_username(creds.access_token)
docker = self._login_docker(docker_username, creds.access_token)
return LoginResult(
url=session.url,
docker=docker,
version=self._get_version_info(metadata),
username=identity,
)
[docs]
def reset_url(self, url: str) -> None:
"""Clear stored credentials for a specific URL."""
logger.info("resetting stored credentials for URL", url=url)
_ = self.credentials.clear(url)
[docs]
def login_client_credentials(
self,
url: str,
client_id: str,
client_secret: str,
) -> LoginResult:
"""Client credentials login workflow."""
logger.info("starting client credentials login", url=url, client_id=client_id)
tokens = self.auth.login_client_credentials(url, client_id, client_secret)
_ = self._store_tokens(url, tokens)
metadata = self._fetch_metadata(url)
docker_username = self._get_token_username(tokens.access_token)
identity = self._get_token_identity(tokens.access_token)
session = self.session.set_current_session(url, metadata=metadata, username=identity)
return LoginResult(
url=session.url,
docker=self._login_docker(docker_username, tokens.access_token),
version=self._get_version_info(metadata),
username=identity,
)
[docs]
def logout(self) -> None:
"""Logout workflow."""
session = self._require_session()
creds = self.credentials.retrieve(session.url)
if creds and creds.refresh_token:
try:
self.auth.logout(session.url, creds.refresh_token, oauth_client_id=creds.oauth_client_id)
except Exception as exc:
logger.warning("failed to logout from server", url=session.url, error=str(exc))
_ = self.credentials.clear(session.url)
self.session.clear_session()
logger.info("logout completed", url=session.url)
[docs]
def get_token(self, full: bool = False) -> TokenResult:
"""Get current token, refreshing if needed."""
session = self._require_session()
creds = self._require_credentials(session)
if self.credentials.is_token_expired(creds):
logger.info("access token expired, refreshing", url=session.url)
creds = self._refresh_tokens(session, creds)
full_credentials = asdict(creds) if full else None
return TokenResult(
access_token=creds.access_token,
expires_at=creds.access_expires_at,
full_credentials=full_credentials,
)
[docs]
def reset(self) -> None:
"""Full reset of auth state."""
session = self.session.get_current_session()
if session:
_ = self.credentials.clear(session.url)
self.session.clear_session()
logger.info("auth state reset")
def _store_tokens(self, url: str, tokens: TokenResponse, oauth_client_id: Optional[str] = None) -> TokenCredentials:
"""Store tokens in the credential store."""
logger.debug("storing credentials", url=url, has_refresh=bool(tokens.refresh_token))
return self.credentials.store_token_response(url, tokens, oauth_client_id=oauth_client_id)
def _fetch_metadata(self, url: str) -> Optional[PlatformMetadata]:
"""Fetch platform metadata from the API."""
logger.debug("fetching platform metadata", url=url)
try:
return self.auth.fetch_platform_metadata(url)
except Exception as exc:
logger.warning("failed to fetch platform metadata", url=url, error=str(exc))
return None
def _is_custom_docker_registry(self) -> bool:
"""Check if the Docker registry is a custom (non-instance) registry.
Compares the docker_url domain with the session URL domain.
If they differ, the Kelvin instance credentials will not work
for the docker registry.
"""
session = self.session.get_current_session()
if session is None or not session.docker_url:
return False
instance_host = urlparse(session.url).hostname or ""
docker_host = session.docker_url
return instance_host != docker_host
def _login_docker(self, username: str, password: str) -> Optional[DockerLoginInfo]:
"""Attempt login to Docker registry using user credentials.
Returns DockerLoginInfo when there is something to warn about,
None on silent success or when no registry is configured.
Does NOT raise on failure — the auth login should succeed regardless.
"""
docker_login_url = self.session.get_docker_login_url()
if not docker_login_url:
logger.info("docker registry not configured, skipping login")
return None
if self._is_custom_docker_registry():
logger.info(
"custom docker registry detected, skipping automatic login",
docker_url=docker_login_url,
)
return DockerLoginInfo(login_url=docker_login_url, is_custom_registry=True)
try:
self.docker.login(docker_login_url, username, password)
return None
except Exception as exc:
logger.warning(
"docker registry login failed",
registry=docker_login_url,
error=str(exc),
)
return DockerLoginInfo(login_url=docker_login_url, login_failed=True)
@staticmethod
def _get_token_username(access_token: str) -> str:
"""Extract the username from a Keycloak JWT access token.
Decodes the JWT payload (no signature verification — Keycloak already validated it)
and returns `preferred_username`, falling back to `sub`.
"""
try:
payload = access_token.split(".")[1]
payload += "=" * (-len(payload) % 4)
raw: object = json.loads(base64.urlsafe_b64decode(payload))
if not isinstance(raw, dict):
return ""
claims = cast(dict[str, object], raw)
preferred = claims.get("preferred_username")
sub = claims.get("sub")
value = preferred or sub or ""
return str(value) if isinstance(value, str) else ""
except Exception as exc:
logger.warning("failed to decode token username from JWT", error=str(exc))
return ""
@staticmethod
def _get_token_subject(access_token: str) -> Optional[str]:
"""Extract the opaque subject id (``sub``) from a JWT access token.
Kept as a fallback identity when the ``email`` claim is unavailable.
"""
try:
payload = access_token.split(".")[1]
payload += "=" * (-len(payload) % 4)
raw: object = json.loads(base64.urlsafe_b64decode(payload))
if not isinstance(raw, dict):
return None
sub = cast(dict[str, object], raw).get("sub")
return str(sub) if isinstance(sub, str) and sub else None
except Exception: # noqa: BLE001 - identity is best-effort, never fatal
return None
@classmethod
def _get_token_identity(cls, access_token: str) -> Optional[str]:
"""Extract the analytics identity from a JWT access token.
Uses the user's ``email`` claim as the identity so usage is attributed
to a human-readable address in Heap, falling back to ``preferred_username``
and then the opaque ``sub`` when the email claim is absent.
"""
try:
payload = access_token.split(".")[1]
payload += "=" * (-len(payload) % 4)
raw: object = json.loads(base64.urlsafe_b64decode(payload))
if isinstance(raw, dict):
claims = cast(dict[str, object], raw)
email = claims.get("email")
if isinstance(email, str) and email:
return email
preferred = claims.get("preferred_username")
if isinstance(preferred, str) and preferred:
return preferred
except Exception: # noqa: BLE001 - identity is best-effort, never fatal
pass
# Fall back to the opaque subject id.
return cls._get_token_subject(access_token)
def _get_version_info(self, metadata: Optional[PlatformMetadata]) -> Optional[VersionInfo]:
"""Check CLI version against platform requirements.
Returns VersionInfo when the CLI version is outside the recommended range,
None otherwise.
"""
if metadata is None:
return None
if not metadata.ksdk_minimum_version or not metadata.ksdk_latest_version:
return None
try:
current = Version(cli_version)
minimum = Version(metadata.ksdk_minimum_version)
latest = Version(metadata.ksdk_latest_version)
except InvalidVersion:
return None
# Skip version checks for dev/pre-release builds
if current.is_devrelease or current.is_prerelease:
return None
if current < minimum or current > latest:
return VersionInfo(
current=str(current),
recommended=metadata.ksdk_latest_version,
)
return None
def _refresh_tokens(self, session: SessionInfo, creds: TokenCredentials) -> TokenCredentials:
"""Refresh access token and store updated credentials."""
if self.credentials.is_refresh_token_expired(creds):
raise TokenRefreshError("Refresh token expired. Please log in again.")
if not creds.refresh_token:
raise TokenRefreshError("No refresh token available. Please log in again.")
tokens = self.auth.refresh_tokens(session.url, creds.refresh_token, oauth_client_id=creds.oauth_client_id)
return self._store_tokens(session.url, tokens, oauth_client_id=creds.oauth_client_id)
def _require_session(self) -> SessionInfo:
"""Get current session or raise."""
session = self.session.get_current_session()
if session is None:
raise MissingSessionError("No active session. Please log in.")
return session
def _require_credentials(self, session: SessionInfo) -> TokenCredentials:
"""Get stored credentials for session or raise."""
creds = self.credentials.retrieve(session.url)
if creds is None:
raise MissingCredentialsError("No stored credentials. Please log in.")
return creds