From 2b7cfec1e2f40f52bac525d09d326b2e71fc60c9 Mon Sep 17 00:00:00 2001 From: Pearson White Date: Thu, 11 Jun 2026 09:32:06 -0400 Subject: [PATCH] Add logoscore router with endpoints * /logoscore/init Save the issued logoscore token by getting it from the K8s secret * /logoscore/call Make a logoscore module function call using the given params --- Dockerfile | 16 +++- routers/logoscore.py | 203 +++++++++++++++++++++++++++++++++++++++++++ routers/registry.py | 2 + 3 files changed, 219 insertions(+), 2 deletions(-) create mode 100644 routers/logoscore.py diff --git a/Dockerfile b/Dockerfile index 53daff2..81d5778 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,8 +1,20 @@ -FROM python:3.11.9-alpine AS base -WORKDIR /app +FROM pearsonwhite/logos-core-dst:wip1-amd AS base + +USER root +RUN apt-get update && \ + apt-get install -y --no-install-recommends \ + python-is-python3 && \ + rm -rf /var/lib/apt/lists/* + +USER nixuser + +# This is typically created when running logoscore. +# Since we're adding the tokens before running logoscore, we must manually create it. +RUN mkdir -p /home/nixuser/.logoscore/client/ COPY requirements.txt . RUN pip install --no-cache-dir --break-system-packages -r requirements.txt +RUN pip install --no-cache-dir --break-system-packages git+https://github.com/logos-co/logos-logoscore-py.git@5842bd32b9b5529abf017a1e9032dee988d8180b COPY api_requester.py \ utils.py \ configs.py \ diff --git a/routers/logoscore.py b/routers/logoscore.py new file mode 100644 index 0000000..bfca014 --- /dev/null +++ b/routers/logoscore.py @@ -0,0 +1,203 @@ +"""Integration for logoscore interaction.""" + +import base64 +import json +from pathlib import Path +from typing import Annotated, Awaitable, Callable, Optional, Union + +from fastapi import APIRouter, Depends, Request +from kubernetes import client +from logoscore import LogoscoreClient +from pydantic import BaseModel, Field + +from common import get_pod_infos +from kube_client import core_v1 +from routers.deps import TargetConfig, TargetName, endpoint_error_handler, unwrap_arg +from schemas import NotFoundError, TargetPodInfo +from utils import setup_logger + +logger = setup_logger(__file__) + +CONFIG_DIR = Path("~/.logoscore/").expanduser() + + +def get_token_from_secret( + secret_name: str, + namespace: str, + file_name: Optional[str] = None, +) -> dict: + """ + :param secret_name: eg. "alicesecret" + :param namespace: eg. "zerotesting" + :param file_name: Key inside the secret json containing the token. eg. "alice.json". Leave None to use the first key. + """ + HOST = "https://kubernetes.default.svc" + CA_CERT = "/var/run/secrets/kubernetes.io/serviceaccount/ca.crt" + TOKEN_FILE = "/var/run/secrets/kubernetes.io/serviceaccount/token" + + with open(TOKEN_FILE, "r") as f: + token = f.read().strip() + + configuration = client.Configuration() + configuration.host = HOST + configuration.ssl_ca_cert = CA_CERT + configuration.api_key["authorization"] = token + configuration.api_key_prefix["authorization"] = "Bearer" + + api_client = client.ApiClient(configuration) + api_client.default_headers["Authorization"] = f"Bearer {token}" + auth_header = api_client.default_headers.get("Authorization") + logger.debug(f"Client Header: {auth_header}") + + secret = core_v1.read_namespaced_secret(name=secret_name, namespace=namespace) + + name = file_name or next(iter(secret.data)) + + full_token_str = base64.b64decode(secret.data[name]).strip().decode("utf-8") + return json.loads(full_token_str) + + +def save_token(token: dict, config_dir: str, file_name: str): + config_file = Path(config_dir) / "client" / file_name + with open(config_file, "w", encoding="utf-8") as f: + json.dump(token, f, indent=4) + + +def make_config(token_file_name: str, ip: str): + config = { + "version": 2, + "token_file": token_file_name, + "daemon": { + "core_service": {"transport": "tcp", "host": ip, "port": 8645}, + "capability_module": {"transport": "tcp", "host": ip, "port": 8646}, + }, + } + with open(CONFIG_DIR / "client" / "config.json", "w") as config_file: + json.dump(config, config_file, indent=4) + + +def secret_for(target: str) -> str: + return f"{target}-logoscore-secret" + + +def token_file_for(target) -> str: + return f"{target}_token.json" + + +def init_token(target, namespace): + result_data = {} + try: + secret = secret_for(target) + token_file = token_file_for(target) + token = get_token_from_secret(secret, namespace) + save_token(token, CONFIG_DIR, token_file) + + result_data["response"] = { + "status_code": 200, + "text": json.dumps( + {"config_dir": CONFIG_DIR.as_posix(), "token_file": token_file}, + ), + } + except Exception as e: + result_data["exception"] = str(e) + + logger.debug(f"init_token result: {result_data}") + return result_data + + +def make_call(target: TargetPodInfo, params: dict): + namespace = target.pod.metadata.namespace + pod_name = target.pod_name + token_file = token_file_for(pod_name) + target_host = f"{pod_name}.{namespace}.svc.cluster.local" + debug_dict = {"target": target_host, "params": params} + logger.debug(f"make_call {json.dumps(debug_dict, indent=4)}") + target_ip = target.pod.status.pod_ip + make_config(token_file, target_ip) + + result_data = {"request": debug_dict} + + lclient = LogoscoreClient(config_dir=CONFIG_DIR.as_posix()) + + module = params["module"] + function = params["function"] + func_params = params["params"] + + if isinstance(func_params, dict): + func_params = json.dumps(func_params) + + try: + result = lclient.call(module, function, func_params) + + result_data["response"] = { + "status_code": 200 if result["success"] else 500, + "text": json.dumps(result), + } + except Exception as e: + result_data["exception"] = e + + logger.debug(f"make_call result: {result_data}") + return result_data + + +class LogosCoreRequestData(BaseModel): + target: Annotated[Union[TargetName, TargetConfig], Field(discriminator="kind")] + params: Optional[dict] = None + + +def create_router(get_config: Callable[[], Awaitable[dict]]) -> APIRouter: + router = APIRouter() + + @router.post("/logoscore/init") + @endpoint_error_handler + async def init( + request: Request, + data: LogosCoreRequestData, + config=Depends(get_config), + ): + target = unwrap_arg(data.target, "targets", config) + + try: + pods = get_pod_infos( + targets=[target], + namespace=request.app.state.namespace, + cache=request.app.state.cache, + ) + pod_info = next(iter(pods)) + except StopIteration as e: + raise NotFoundError(f"Target not found. Target: {target}") from e + + namespace = pod_info.pod.metadata.namespace + return init_token(pod_info.pod_name, namespace) + + @router.post("/logoscore/call") + @endpoint_error_handler + async def call( + request: Request, + data: LogosCoreRequestData, + config=Depends(get_config), + ): + target = unwrap_arg(data.target, "targets", config) + params = data.params + + try: + pods = get_pod_infos( + targets=[target], + namespace=request.app.state.namespace, + cache=request.app.state.cache, + ) + pod_info = next(iter(pods)) + except StopIteration as e: + raise NotFoundError(f"Target not found. Target: {target}") from e + + result = make_call(pod_info, params) + + # TODO: do this. we want to redact on message call results, not other call results (but we can EAFP) + # Remove values for "payload", since it is a large amount of generated bytes. + # result = redact_keys(result, keys_to_redact=("payload",)) + # configEndpoint = result["request"]["configEndpoint"] + # configEndpoint.params = redact_keys(configEndpoint.params, keys_to_redact=("payload",)) + + return result + + return router diff --git a/routers/registry.py b/routers/registry.py index 70174fa..a5c59a5 100644 --- a/routers/registry.py +++ b/routers/registry.py @@ -3,9 +3,11 @@ from typing import Iterable from fastapi import APIRouter from routers.generic import create_router as create_generic_router +from routers.logoscore import create_router as create_logoscore_router from routers.waku import create_router as create_waku_router def build_routers(get_config) -> Iterable[APIRouter]: yield create_generic_router(get_config) yield create_waku_router(get_config) + yield create_logoscore_router(get_config)