Initial production analytics foundation

This commit is contained in:
2026-09-04 05:34:52 +02:00
commit f017d187eb
30 changed files with 1219 additions and 0 deletions
+1
View File
@@ -0,0 +1 @@
"""Derived production analytics between ENLYZE and Grafana."""
@@ -0,0 +1,5 @@
"""Versioned, pure calculation operators and their contracts."""
from .base import Calculation
__all__ = ["Calculation"]
@@ -0,0 +1,15 @@
"""Calculation contracts; concrete operators follow verified source semantics."""
from collections.abc import Iterable
from typing import Protocol
from production_analytics.domain import CalculatedMetric, ProcessEvent
class Calculation(Protocol):
"""A testable operator that emits derived results from supplied input."""
calculation_type: str
version: str
def calculate(self, samples: Iterable[object]) -> Iterable[CalculatedMetric | ProcessEvent]: ...
+1
View File
@@ -0,0 +1 @@
"""Command-line entry points; API exploration is the first planned command."""
+147
View File
@@ -0,0 +1,147 @@
"""Command-line entry point for safe, read-only API exploration."""
from __future__ import annotations
import argparse
import json
import os
import sys
from collections.abc import Sequence
from pathlib import Path
from production_analytics.enlyze.exploration import (
ExplorationClient,
ExplorationError,
ExplorationSettings,
load_secret_file,
)
from production_analytics.enlyze.sanitize import sanitize
def _query_item(value: str) -> tuple[str, str]:
if "=" not in value:
raise argparse.ArgumentTypeError("query parameters must use KEY=VALUE format")
key, item_value = value.split("=", 1)
if not key:
raise argparse.ArgumentTypeError("query parameter key must not be empty")
return key, item_value
def _parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(prog="production-analytics")
namespaces = parser.add_subparsers(dest="namespace", required=True)
enlyze = namespaces.add_parser("enlyze", help="Read-only ENLYZE API exploration")
commands = enlyze.add_subparsers(dest="command", required=True)
raw = commands.add_parser("raw", help="GET an operator-verified relative API path")
raw.add_argument("path", help="Relative API path beginning with '/'")
raw.add_argument("--query", action="append", type=_query_item, default=[], metavar="KEY=VALUE")
raw.add_argument("--pretty", action="store_true", help="Pretty-print sanitized JSON")
raw.add_argument("--verbose", action="store_true", help="Show safe request progress on stderr")
raw.add_argument("--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/")
raw.add_argument("--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/")
raw.add_argument("--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS)
raw.add_argument("--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS)
raw.add_argument(
"--secrets-file",
type=Path,
default=Path("secrets/enlyze.env"),
help="Local dotenv-style secret file (default: secrets/enlyze.env)",
)
timeseries = commands.add_parser("timeseries", help="Read-only POST /v2/timeseries exploration")
timeseries.add_argument("--machine", required=True, help="Machine UUID")
timeseries.add_argument("--start", required=True, help="ISO 8601 datetime with timezone")
timeseries.add_argument("--end", required=True, help="ISO 8601 datetime with timezone")
timeseries.add_argument("--variable", required=True, help="Variable UUID")
timeseries.add_argument("--resampling-interval", type=int, help="Seconds; schema range is 10..604800")
timeseries.add_argument(
"--resampling-method",
choices=["first", "last", "max", "min", "count", "sum", "avg", "median", "std", "q5", "q25", "q75", "q95"],
help="Optional schema-defined method for this variable",
)
timeseries.add_argument("--pretty", action="store_true", help="Pretty-print sanitized JSON")
timeseries.add_argument("--verbose", action="store_true", help="Show safe request progress on stderr")
timeseries.add_argument("--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/")
timeseries.add_argument("--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/")
timeseries.add_argument("--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS)
timeseries.add_argument("--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS)
timeseries.add_argument(
"--secrets-file",
type=Path,
default=Path("secrets/enlyze.env"),
help="Local dotenv-style secret file (default: secrets/enlyze.env)",
)
return parser
def _write_fixture(directory: Path, name: str, payload: object) -> Path:
candidate = Path(name)
if candidate.name != name or candidate.suffix.lower() != ".json":
raise ExplorationError("Fixture name must be a plain filename ending in .json.")
directory.mkdir(parents=True, exist_ok=True)
destination = directory / candidate
destination.write_text(json.dumps(payload, indent=2, sort_keys=True) + "\n", encoding="utf-8")
return destination
def main(argv: Sequence[str] | None = None) -> int:
args = _parser().parse_args(argv)
if args.namespace != "enlyze" or args.command not in {"raw", "timeseries"}:
return 2
try:
environment = dict(os.environ)
environment.update(load_secret_file(args.secrets_file))
settings = ExplorationSettings.from_environment(environment)
if args.verbose:
authentication = "configured" if settings.api_key else "not configured"
operation = "GET" if args.command == "raw" else "POST"
target = args.path if args.command == "raw" else "/v2/timeseries"
print(
f"Requesting {operation} {target} (timeout={settings.timeout_seconds:g}s; authentication {authentication}).",
file=sys.stderr,
)
client = ExplorationClient(settings)
if args.command == "raw":
response = client.get(args.path, dict(args.query))
request = {"method": "GET", "path": response.path}
else:
if args.resampling_interval is not None and not 10 <= args.resampling_interval <= 604800:
raise ExplorationError("--resampling-interval must be between 10 and 604800 seconds.")
variable: dict[str, str] = {"uuid": args.variable}
if args.resampling_method:
variable["resampling_method"] = args.resampling_method
request_body: dict[str, object] = {
"machine": args.machine,
"start": args.start,
"end": args.end,
"variables": [variable],
}
if args.resampling_interval is not None:
request_body["resampling_interval"] = args.resampling_interval
response = client.post_json("/v2/timeseries", request_body)
request = {"method": "POST", "path": response.path, "body": request_body}
if args.verbose:
print(f"Received HTTP {response.status_code} for {response.path}.", file=sys.stderr)
raw_payload = {
"request": request,
"response": {
"status_code": response.status_code,
"headers": response.headers,
"body": response.body,
},
}
payload = sanitize(raw_payload)
indent = 2 if args.pretty or args.save_fixture else None
print(json.dumps(payload, indent=indent, sort_keys=True))
if args.save_fixture:
print(f"Sanitized fixture written to {_write_fixture(args.fixture_dir, args.save_fixture, payload)}")
if args.save_raw:
print(f"Raw response written to {_write_fixture(args.raw_dir, args.save_raw, raw_payload)}")
except ExplorationError as error:
print(f"ENLYZE exploration error: {error}")
return 1
return 0
if __name__ == "__main__":
raise SystemExit(main())
@@ -0,0 +1,19 @@
"""Stable domain models independent of API and database implementations."""
from .models import (
CalculatedMetric,
CalculationRun,
MachineState,
ProcessEvent,
ProductionOrder,
SourceRange,
)
__all__ = [
"CalculatedMetric",
"CalculationRun",
"MachineState",
"ProcessEvent",
"ProductionOrder",
"SourceRange",
]
+72
View File
@@ -0,0 +1,72 @@
"""Small shared models for derived-result attribution and provenance."""
from dataclasses import dataclass
from datetime import datetime
from enum import StrEnum
from typing import Mapping
@dataclass(frozen=True, slots=True)
class SourceRange:
"""Inclusive/exclusive raw-data interval used by a calculation."""
start: datetime
end: datetime
def __post_init__(self) -> None:
if self.end < self.start:
raise ValueError("source range end must not precede start")
@dataclass(frozen=True, slots=True)
class ProductionOrder:
source_id: str
machine_id: str
article_number: str | None = None
start: datetime | None = None
end: datetime | None = None
class MachineState(StrEnum):
OPERATING = "operating"
DOWNTIME = "downtime"
UNKNOWN = "unknown"
@dataclass(frozen=True, slots=True)
class CalculatedMetric:
name: str
value: float
unit: str
machine_id: str
observed_at: datetime
source_range: SourceRange
calculation_type: str
calculation_version: str
production_order_id: str | None = None
article_number: str | None = None
@dataclass(frozen=True, slots=True)
class ProcessEvent:
event_type: str
occurred_at: datetime
machine_id: str
source_range: SourceRange
calculation_type: str
calculation_version: str
value: float | None = None
unit: str | None = None
production_order_id: str | None = None
article_number: str | None = None
attributes: Mapping[str, str | float | int | bool | None] | None = None
@dataclass(frozen=True, slots=True)
class CalculationRun:
calculation_id: str
calculation_type: str
calculation_version: str
source_range: SourceRange
started_at: datetime
completed_at: datetime | None = None
@@ -0,0 +1,6 @@
"""ENLYZE adapter boundary; concrete API behavior awaits verification."""
from .exploration import ExplorationClient, ExplorationSettings, load_secret_file
from .ports import EnlyzeGateway
__all__ = ["EnlyzeGateway", "ExplorationClient", "ExplorationSettings", "load_secret_file"]
@@ -0,0 +1,165 @@
"""Small, read-only HTTP boundary for observing an unverified ENLYZE API."""
from __future__ import annotations
import json
import os
import shlex
from dataclasses import dataclass
from typing import Any, Mapping
from urllib.error import HTTPError, URLError
from urllib.parse import urlencode, urljoin, urlsplit
from urllib.request import Request, urlopen
class ExplorationError(RuntimeError):
"""Base error whose message is safe to show to an operator."""
class ConfigurationError(ExplorationError):
"""Raised when local exploration configuration is incomplete or unsafe."""
class AuthenticationError(ExplorationError):
"""Raised for HTTP 401/403 without exposing credentials or response bodies."""
class HttpResponseError(ExplorationError):
"""Raised for a non-successful HTTP response."""
class NonJsonResponseError(ExplorationError):
"""Raised when a successful response cannot be decoded as JSON."""
@dataclass(frozen=True, slots=True)
class ExplorationSettings:
"""Local HTTP settings for documented ENLYZE Bearer-token access."""
base_url: str
timeout_seconds: float = 20.0
api_key: str | None = None
@classmethod
def from_environment(cls, environment: Mapping[str, str] | None = None) -> ExplorationSettings:
environment = environment or os.environ
base_url = environment.get("ENLYZE_BASE_URL", "").strip()
api_key = environment.get("ENLYZE_API_KEY", "").strip() or None
timeout_value = environment.get("ENLYZE_HTTP_TIMEOUT_SECONDS", "20")
if not base_url:
raise ConfigurationError("ENLYZE_BASE_URL must be set for API exploration.")
try:
timeout_seconds = float(timeout_value)
except ValueError as error:
raise ConfigurationError("ENLYZE_HTTP_TIMEOUT_SECONDS must be a number.") from error
if timeout_seconds <= 0:
raise ConfigurationError("ENLYZE_HTTP_TIMEOUT_SECONDS must be greater than zero.")
return cls(base_url.rstrip("/") + "/", timeout_seconds, api_key)
@dataclass(frozen=True, slots=True)
class ExplorationResponse:
"""JSON response plus safe metadata useful for API investigation."""
status_code: int
path: str
headers: Mapping[str, str]
body: Any
class ExplorationClient:
"""Client for documented, read-only exploration operations."""
def __init__(self, settings: ExplorationSettings) -> None:
self._settings = settings
def get(self, path: str, query: Mapping[str, str] | None = None) -> ExplorationResponse:
"""Request one relative path and decode its JSON response."""
return self._request_json("GET", path, query=query)
def post_json(self, path: str, payload: Mapping[str, Any]) -> ExplorationResponse:
"""Issue a documented read-only POST operation with a JSON body."""
return self._request_json("POST", path, payload=payload)
def _request_json(
self,
method: str,
path: str,
*,
query: Mapping[str, str] | None = None,
payload: Mapping[str, Any] | None = None,
) -> ExplorationResponse:
parsed_path = urlsplit(path)
if (
not path.startswith("/")
or path.startswith("//")
or parsed_path.scheme
or parsed_path.netloc
or ".." in parsed_path.path.split("/")
):
raise ConfigurationError("PATH must be a relative path beginning with one slash.")
encoded_query = urlencode(query or {})
request_path = f"{path}?{encoded_query}" if encoded_query else path
request_url = urljoin(self._settings.base_url, request_path.lstrip("/"))
request_headers = {"Accept": "application/json"}
request_body = None
if payload is not None:
request_headers["Content-Type"] = "application/json"
request_body = json.dumps(payload).encode("utf-8")
if self._settings.api_key:
request_headers["Authorization"] = f"Bearer {self._settings.api_key}"
request = Request(request_url, data=request_body, headers=request_headers, method=method)
try:
with urlopen(request, timeout=self._settings.timeout_seconds) as response: # noqa: S310
raw_body = response.read()
status_code = response.status
response_headers = dict(response.headers.items())
except HTTPError as error:
if error.code in {401, 403}:
raise AuthenticationError(f"Authentication or authorization failed (HTTP {error.code}).") from error
raise HttpResponseError(f"HTTP request failed with status {error.code}.") from error
except URLError as error:
raise HttpResponseError("HTTP request could not be completed.") from error
try:
body = json.loads(raw_body)
except (UnicodeDecodeError, json.JSONDecodeError) as error:
content_type = response_headers.get("Content-Type", "unknown")
raise NonJsonResponseError(
f"Expected a JSON response; received Content-Type {content_type!r}."
) from error
return ExplorationResponse(status_code, request_path, response_headers, body)
def load_secret_file(path: str | os.PathLike[str]) -> dict[str, str]:
"""Read simple dotenv-style assignments without executing local secret content."""
secrets: dict[str, str] = {}
try:
lines = open(path, encoding="utf-8")
except FileNotFoundError:
return secrets
with lines:
for line_number, line in enumerate(lines, start=1):
stripped = line.strip()
if not stripped or stripped.startswith("#"):
continue
if stripped.startswith("export "):
stripped = stripped.removeprefix("export ").lstrip()
if "=" not in stripped:
raise ConfigurationError(f"Invalid secret-file assignment on line {line_number}.")
key, raw_value = stripped.split("=", 1)
key = key.strip()
if not key.isidentifier():
raise ConfigurationError(f"Invalid secret-file variable name on line {line_number}.")
try:
parsed = shlex.split(raw_value, comments=True, posix=True)
except ValueError as error:
raise ConfigurationError(f"Invalid secret-file value on line {line_number}.") from error
if len(parsed) > 1:
raise ConfigurationError(f"Invalid secret-file value on line {line_number}.")
secrets[key] = parsed[0] if parsed else ""
return secrets
+7
View File
@@ -0,0 +1,7 @@
"""Ports intentionally avoid unverified ENLYZE endpoint or payload assumptions."""
from typing import Protocol
class EnlyzeGateway(Protocol):
"""Marker for a verified ENLYZE adapter, defined after API exploration."""
@@ -0,0 +1,62 @@
"""Conservative sanitization for reviewable API fixtures and CLI output."""
from __future__ import annotations
import re
from collections.abc import Mapping
from dataclasses import dataclass, field
from typing import Any
from urllib.parse import urlsplit, urlunsplit
_SECRET_KEY = re.compile(
r"(?:authorization|token|api[_-]?key|password|secret|cookie|credential|session)", re.IGNORECASE
)
_IDENTIFIER_KEY = re.compile(r"(?:^|[_-])(id|uuid|guid)(?:$|[_-])", re.IGNORECASE)
_SENSITIVE_TEXT_KEY = re.compile(
r"(?:email|phone|first[_-]?name|last[_-]?name|full[_-]?name|user[_-]?name|customer|site|tenant|name)",
re.IGNORECASE,
)
_HOST_KEY = re.compile(r"(?:host|hostname|base[_-]?url|url|uri|endpoint)", re.IGNORECASE)
@dataclass
class _SanitizationContext:
identifiers: dict[tuple[str, str], str | int] = field(default_factory=dict)
def replacement_identifier(self, key: str, value: object) -> str | int:
original = str(value)
identity = (key.lower(), original)
if identity not in self.identifiers:
number = len(self.identifiers) + 1
self.identifiers[identity] = 1000 + number if isinstance(value, int) else f"redacted-{key}-{number}"
return self.identifiers[identity]
def sanitize(value: Any) -> Any:
"""Return a structurally equivalent value with common sensitive fields replaced."""
return _sanitize(value, _SanitizationContext())
def _sanitize(value: Any, context: _SanitizationContext(), key: str | None = None) -> Any:
if isinstance(value, Mapping):
return {str(item_key): _sanitize(item_value, context, str(item_key)) for item_key, item_value in value.items()}
if isinstance(value, list):
return [_sanitize(item, context, key) for item in value]
if isinstance(value, tuple):
return [_sanitize(item, context, key) for item in value]
if key and _SECRET_KEY.search(key):
return "<redacted-secret>"
if key and _IDENTIFIER_KEY.search(key) and isinstance(value, (str, int)) and not isinstance(value, bool):
return context.replacement_identifier(key, value)
if key and _HOST_KEY.search(key) and isinstance(value, str):
return _redact_host(value)
if key and _SENSITIVE_TEXT_KEY.search(key) and isinstance(value, str):
return f"<redacted-{key}>"
return value
def _redact_host(value: str) -> str:
parsed = urlsplit(value)
if parsed.scheme and parsed.netloc:
return urlunsplit((parsed.scheme, "redacted-host.invalid", parsed.path, parsed.query, parsed.fragment))
return "<redacted-host>"
@@ -0,0 +1,5 @@
"""Persistence ports for derived data and incremental calculation state."""
from .ports import DerivedResultRepository
__all__ = ["DerivedResultRepository"]
@@ -0,0 +1,13 @@
"""Database-independent interfaces; no raw ENLYZE sample persistence."""
from typing import Protocol
from production_analytics.domain import CalculatedMetric, ProcessEvent
class DerivedResultRepository(Protocol):
"""Stores only derived results; implementation follows schema design."""
def save_metric(self, metric: CalculatedMetric) -> None: ...
def save_event(self, event: ProcessEvent) -> None: ...
@@ -0,0 +1 @@
"""Future FastAPI composition layer; intentionally endpoint-free for now."""