From 0f54aa4e18794487f683286fe5bd5edb4415d804 Mon Sep 17 00:00:00 2001 From: Philip Z Date: Thu, 1 Oct 2026 04:03:42 +0800 Subject: [PATCH] feat: add governed Business ID lifecycle Closes #41 --- README.md | 15 +- examples/business-id/README.md | 9 + examples/business-id/main.py | 63 +++++++ scripts/verify-examples.sh | 3 +- src/teaql/core/__init__.py | 16 +- src/teaql/core/business_id.py | 163 ++++++++++++++++++ src/teaql/runtime/__init__.py | 7 +- src/teaql/runtime/business_id.py | 171 +++++++++++++++++++ src/teaql/runtime/context.py | 23 +++ src/teaql/sql/__init__.py | 3 +- src/teaql/sql/business_id.py | 114 +++++++++++++ src/teaql/sql/executor.py | 6 + tests/provider/test_business_id_allocator.py | 86 ++++++++++ tests/runtime/test_business_id_lifecycle.py | 87 ++++++++++ tests/runtime/test_sql_mask_lifecycle.py | 4 +- 15 files changed, 759 insertions(+), 11 deletions(-) create mode 100644 examples/business-id/README.md create mode 100644 examples/business-id/main.py create mode 100644 src/teaql/sql/business_id.py create mode 100644 tests/provider/test_business_id_allocator.py create mode 100644 tests/runtime/test_business_id_lifecycle.py diff --git a/README.md b/README.md index 79dbb2b..a92bfc4 100644 --- a/README.md +++ b/README.md @@ -70,9 +70,10 @@ The SDK's organizational architecture strictly mirrors the Rust version: mutation. Exact policy identity and approval state are retained with audit evidence. Missing customer policy or approval emits stable warnings without changing persistence semantics; an explicit denial fails closed. -* **Portable Business ID Encoding**: Core scope/key types and the runtime - `daily-permuted-v1` encoder execute the canonical cross-language golden - vectors without exposing the internal sequence. +* **Governed Business ID Lifecycle**: Core model contracts, a context-owned + profile/key/service boundary, retry-safe assignment, an in-memory allocator, + and explicit-schema durable SQLite allocation extend the portable + `daily-permuted-v1` encoder without exposing its internal sequence. * **TeaQL Federal Protocol Client**: `TeaQLFederalClient` and `TfpHttpProvider` execute governed canonical TFP v1 queries and audited mutations against a remote TeaQL endpoint such as Rust. Direct query execution returns @@ -150,9 +151,11 @@ key = BusinessIdEncodingKey(1, key_from_secret_manager) code = encode_business_id_permutation_v1(0, scope, key) ``` -Durable concurrent allocation, aggregate retry reuse, and typed lookup are -separate lifecycle capabilities. Secret key material is application-owned and -must not be placed in KSML or generated source. +The retained [`examples/business-id`](examples/business-id) flow proves durable +concurrent allocation infrastructure and aggregate retry reuse. Generated +strongly typed fields and external lookup remain a separate generator +capability. Secret key material is application-owned and must not be placed in +KSML or generated source. ### Mutation Policy installation diff --git a/examples/business-id/README.md b/examples/business-id/README.md new file mode 100644 index 0000000..7e9cea4 --- /dev/null +++ b/examples/business-id/README.md @@ -0,0 +1,9 @@ +# Governed Business ID example + +This focused example proves explicit SQLite schema installation, a +context-owned business date and key provider, durable allocation, and +idempotent reuse of an already assigned Business ID. + +```bash +PYTHONPATH=src python examples/business-id/main.py +``` diff --git a/examples/business-id/main.py b/examples/business-id/main.py new file mode 100644 index 0000000..14b77d2 --- /dev/null +++ b/examples/business-id/main.py @@ -0,0 +1,63 @@ +import asyncio +from datetime import datetime, timezone +from pathlib import Path +from tempfile import TemporaryDirectory + +from teaql.core import BusinessIdDefinition, BusinessIdEncodingKey +from teaql.provider.sqlite import create_sqlite_service +from teaql.runtime import ( + DefaultBusinessIdProfileFactory, + DefaultBusinessIdService, + FixedBusinessClock, + StaticBusinessIdKeyProvider, + UserContext, +) +from teaql.sql import SqlBusinessIdAllocator + + +class OrderNumberSlot: + def __init__(self): + self.value = None + + def current_value(self): + return self.value + + def new_aggregate(self): + return True + + def assign_canonical_value(self, value): + self.value = value + + +async def main(): + with TemporaryDirectory() as directory: + service = create_sqlite_service(str(Path(directory) / "business-id.db")) + allocator = SqlBusinessIdAllocator(service.dialect, service.transport) + context = ( + UserContext.new() + .insert_resource("dataService", service) + .with_business_clock(FixedBusinessClock( + datetime(2026, 10, 1, 8, 30, tzinfo=timezone.utc) + )) + .with_business_id_key_provider(StaticBusinessIdKeyProvider( + BusinessIdEncodingKey(1, bytes(range(32))) + )) + .with_business_id_profile_factory(DefaultBusinessIdProfileFactory()) + .with_business_id_service(DefaultBusinessIdService(allocator)) + ) + await context.ensure_schema() + slot = OrderNumberSlot() + definition = BusinessIdDefinition.daily_permuted( + "order_number", "ORD", "order_number" + ) + first = await context.ensure_business_id( + definition, "commerce", "commerce_order", slot + ) + retry = await context.ensure_business_id( + definition, "commerce", "commerce_order", slot + ) + assert first == retry and slot.value.startswith("ORD-20261001-") + print("PASS Python governed Business ID lifecycle example") + + +asyncio.run(main()) diff --git a/scripts/verify-examples.sh b/scripts/verify-examples.sh index 0f0e578..9c5ae48 100755 --- a/scripts/verify-examples.sh +++ b/scripts/verify-examples.sh @@ -2,7 +2,7 @@ set -euo pipefail repo="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" -expected=(business-clock conformance mutation-policy opaque-entity-reference order-management query-policy school-management task_board) +expected=(business-clock business-id conformance mutation-policy opaque-entity-reference order-management query-policy school-management task_board) mapfile -t actual < <(find "$repo/examples" -mindepth 1 -maxdepth 1 -type d -printf '%f\n' | sort) if [[ "${actual[*]}" != "${expected[*]}" ]]; then echo "example inventory changed; update scripts/verify-examples.sh: ${actual[*]}" >&2 @@ -26,6 +26,7 @@ PYTHONPATH="$repo/examples/school-management:$repo/src" python -m app.main PYTHONPATH="$repo/src" python -m unittest discover -s "$repo/examples/school-management" -p 'test_sql_log_intent.py' -v PYTHONPATH="$repo/src" python "$repo/examples/mutation-policy/main.py" PYTHONPATH="$repo/src" python "$repo/examples/business-clock/main.py" +PYTHONPATH="$repo/src" python "$repo/examples/business-id/main.py" PYTHONPATH="$repo/src" python "$repo/examples/query-policy/main.py" PYTHONPATH="$repo/src" python "$repo/examples/opaque-entity-reference/main.py" order_management_tmp="$(mktemp -d)" diff --git a/src/teaql/core/__init__.py b/src/teaql/core/__init__.py index ff5691c..c82407c 100644 --- a/src/teaql/core/__init__.py +++ b/src/teaql/core/__init__.py @@ -25,10 +25,21 @@ from .safe_expression import SafeExpression from .xls import XlsWorkbook, XlsPage, XlsBlock, XlsBlockBuildContext from .business_id import ( + BusinessIdAllocation, + BusinessIdAllocator, + BusinessIdDefinition, BusinessIdEncodingKey, BusinessIdError, BusinessIdErrorCode, + BusinessIdGenerationRequest, + BusinessIdKeyProvider, + BusinessIdPlan, + BusinessIdProfile, + BusinessIdProfileFactory, + BusinessIdService, + BusinessIdSlot, BusinessIdScope, + BusinessIdValue, ) __all__ = [ @@ -45,8 +56,11 @@ "GraphNode", "EntityDescriptor", "PropertyDescriptor", "SmartList", "TeaQLPage", "LoadState", "EvalResult", "SafeExpression", + "BusinessIdAllocation", "BusinessIdAllocator", "BusinessIdDefinition", "BusinessIdEncodingKey", "BusinessIdError", "BusinessIdErrorCode", - "BusinessIdScope", + "BusinessIdGenerationRequest", "BusinessIdKeyProvider", "BusinessIdPlan", + "BusinessIdProfile", "BusinessIdProfileFactory", "BusinessIdService", + "BusinessIdSlot", "BusinessIdScope", "BusinessIdValue", "XlsWorkbook", "XlsPage", "XlsBlock", "XlsBlockBuildContext" ] import builtins diff --git a/src/teaql/core/business_id.py b/src/teaql/core/business_id.py index 55431ea..6f2682f 100644 --- a/src/teaql/core/business_id.py +++ b/src/teaql/core/business_id.py @@ -3,12 +3,23 @@ from __future__ import annotations from dataclasses import dataclass +from datetime import date from enum import Enum +from typing import Protocol, TYPE_CHECKING + +if TYPE_CHECKING: + from teaql.runtime.context import UserContext class BusinessIdErrorCode(str, Enum): + PROFILE_NOT_FOUND = "BUSINESS_ID_PROFILE_NOT_FOUND" DEFINITION_INVALID = "BUSINESS_ID_DEFINITION_INVALID" RANGE_EXHAUSTED = "BUSINESS_ID_RANGE_EXHAUSTED" + ALLOCATION_RETRY_EXHAUSTED = "BUSINESS_ID_ALLOCATION_RETRY_EXHAUSTED" + FORMAT_INVALID = "BUSINESS_ID_FORMAT_INVALID" + DUPLICATE = "BUSINESS_ID_DUPLICATE" + IMMUTABLE = "BUSINESS_ID_IMMUTABLE" + KEY_NOT_FOUND = "BUSINESS_ID_KEY_NOT_FOUND" ENCODING_FAILED = "BUSINESS_ID_ENCODING_FAILED" @@ -38,6 +49,17 @@ def __post_init__(self) -> None: f"{name} must not be blank", ) + def canonical_key(self) -> str: + def escape(value: str) -> str: + return value.replace("%", "%25").replace("|", "%7C") + + return "|".join(escape(value) for value in ( + self.domain_root_key, + self.aggregate_type, + self.namespace, + self.period_key, + )) + @dataclass(frozen=True) class BusinessIdEncodingKey: @@ -63,3 +85,144 @@ def __post_init__(self) -> None: "Business ID V1 key must contain exactly 32 bytes", ) object.__setattr__(self, "key", bytes(self.key)) + + +@dataclass(frozen=True) +class BusinessIdDefinition: + field_name: str + profile: str + prefix: str + date_format: str + reset: str + digits: int + separator: str + namespace: str + policy_version: int + + DEFAULT_PROFILE = "daily-permuted-v1" + DEFAULT_DATE_FORMAT = "yyyyMMdd" + DEFAULT_DIGITS = 6 + + def __post_init__(self) -> None: + for name in ( + "field_name", "profile", "prefix", "date_format", "reset", + "separator", "namespace", + ): + value = getattr(self, name) + if not isinstance(value, str) or not value.strip(): + raise BusinessIdError( + BusinessIdErrorCode.DEFINITION_INVALID, + f"{name} must not be blank", + ) + if self.profile == self.DEFAULT_PROFILE and self.digits != self.DEFAULT_DIGITS: + raise BusinessIdError( + BusinessIdErrorCode.DEFINITION_INVALID, + "daily-permuted-v1 requires exactly 6 digits", + ) + if not isinstance(self.digits, int) or isinstance(self.digits, bool) or not 1 <= self.digits <= 18: + raise BusinessIdError( + BusinessIdErrorCode.DEFINITION_INVALID, + "digits must be between 1 and 18", + ) + if not isinstance(self.policy_version, int) or isinstance(self.policy_version, bool) or self.policy_version < 1: + raise BusinessIdError( + BusinessIdErrorCode.DEFINITION_INVALID, + "policy_version must be positive", + ) + + @classmethod + def daily_permuted(cls, field_name: str, prefix: str, namespace: str) -> "BusinessIdDefinition": + return cls( + field_name, cls.DEFAULT_PROFILE, prefix, cls.DEFAULT_DATE_FORMAT, + "daily", cls.DEFAULT_DIGITS, "-", namespace, 1, + ) + + def maximum_sequence(self) -> int: + if self.profile == self.DEFAULT_PROFILE: + return 2_176_782_335 + return 10 ** self.digits - 1 + + +@dataclass(frozen=True) +class BusinessIdGenerationRequest: + definition: BusinessIdDefinition + domain_root_key: str + aggregate_type: str + business_date: date + + +@dataclass(frozen=True) +class BusinessIdPlan: + definition: BusinessIdDefinition + scope: BusinessIdScope + business_date: date + date_text: str + initial_sequence: int + maximum_sequence: int + + def __post_init__(self) -> None: + if self.initial_sequence < 0 or self.maximum_sequence < self.initial_sequence: + raise BusinessIdError( + BusinessIdErrorCode.DEFINITION_INVALID, + "Business ID allocation range must satisfy 0 <= initial_sequence <= maximum_sequence", + ) + + +@dataclass(frozen=True) +class BusinessIdAllocation: + scope: BusinessIdScope + sequence: int + + +@dataclass(frozen=True) +class BusinessIdValue: + value: str + profile: str + policy_version: int + + def __post_init__(self) -> None: + if not isinstance(self.value, str) or not self.value.strip(): + raise BusinessIdError( + BusinessIdErrorCode.FORMAT_INVALID, + "Business ID value must not be blank", + ) + + +class BusinessIdSlot(Protocol): + def current_value(self) -> str | None: ... + def new_aggregate(self) -> bool: ... + def assign_canonical_value(self, value: str) -> None: ... + + +class BusinessIdAllocator(Protocol): + async def allocate(self, plan: BusinessIdPlan) -> BusinessIdAllocation: ... + + +class BusinessIdProfile(Protocol): + def plan(self, request: BusinessIdGenerationRequest) -> BusinessIdPlan: ... + def format(self, plan: BusinessIdPlan, allocation: BusinessIdAllocation) -> BusinessIdValue: ... + def validate(self, definition: BusinessIdDefinition, value: str) -> BusinessIdValue: ... + + +class BusinessIdProfileFactory(Protocol): + def create(self, context: "UserContext", definition: BusinessIdDefinition) -> BusinessIdProfile: ... + + +class BusinessIdKeyProvider(Protocol): + def current_key( + self, + context: "UserContext", + definition: BusinessIdDefinition, + scope: BusinessIdScope, + ) -> BusinessIdEncodingKey: ... + + +class BusinessIdService(Protocol): + async def ensure( + self, + context: "UserContext", + definition: BusinessIdDefinition, + domain_root_key: str, + aggregate_type: str, + slot: BusinessIdSlot, + ) -> BusinessIdValue: ... diff --git a/src/teaql/runtime/__init__.py b/src/teaql/runtime/__init__.py index 63d9573..44f2f21 100644 --- a/src/teaql/runtime/__init__.py +++ b/src/teaql/runtime/__init__.py @@ -25,6 +25,11 @@ BUSINESS_ID_PERMUTATION_V1_DOMAIN_SIZE, BUSINESS_ID_PERMUTATION_V1_MAX_SEQUENCE, BUSINESS_ID_PERMUTATION_V1_WIDTH, + DefaultBusinessIdProfileFactory, + DefaultBusinessIdService, + InMemoryBusinessIdAllocator, + PermutedDailyBusinessIdProfile, + StaticBusinessIdKeyProvider, encode_business_id_permutation_v1, ) from .mutation_policy import ( @@ -51,7 +56,7 @@ from .wire_fields import NormalizedWireInput, WireEntityMetadata, WireFieldMetadata, WireInputError, create_wire_entity_metadata, encode_wire_output, normalize_wire_input, retain_submitted_paths from teaql.core.entity import EntityKey, EntityChangeSet, EntityRoot -__all__ = ["WireFieldMetadata", "WireEntityMetadata", "NormalizedWireInput", "WireInputError", "create_wire_entity_metadata", "normalize_wire_input", "encode_wire_output", "retain_submitted_paths", "EntityKey", "EntityChangeSet", "EntityRoot", "ContextEntityRef", "ContextRootError", "CheckException", "CheckResult", "I18nCatalog", "JsonFieldNamingProfile", "Locale", "ObjectLocation", "UnsupportedLocaleError", "UserContext", "TeaqlRuntime", "SqlLogEntry", "SqlLogOperation", "DiagnosticSqlLogSink", "TextDiagnosticSqlLogSink", "ServiceRuntimeFromEnv", "RuntimeModule", "DataStore", "RawAuditEvent", "SafeAuditEvent", "MutationAuditKind", "BusinessClock", "FixedBusinessClock", "SystemBusinessClock", "AeadEntityReferenceCodec", "EntityReferenceClaims", "EntityReferenceCodec", "EntityReferenceTokenError", "ENTITY_REFERENCE_AAD", "UNSAFE_RAW_ENTITY_REFERENCES_ACKNOWLEDGEMENT", "UNSAFE_RAW_ENTITY_REFERENCES_ENVIRONMENT", "BUSINESS_ID_PERMUTATION_V1_ALPHABET", "BUSINESS_ID_PERMUTATION_V1_DOMAIN_SIZE", "BUSINESS_ID_PERMUTATION_V1_MAX_SEQUENCE", "BUSINESS_ID_PERMUTATION_V1_WIDTH", "encode_business_id_permutation_v1", "ContextTools", "ExecutableHttpTool", "HTTP_TOOL", "HttpIntentPhase", "HttpTool", "HttpToolProvider", "ToolDeniedError", "ToolError", "ToolPolicy", "ToolRisk", "Tools", "ToolToken", "ToolUnavailableError"] +__all__ = ["WireFieldMetadata", "WireEntityMetadata", "NormalizedWireInput", "WireInputError", "create_wire_entity_metadata", "normalize_wire_input", "encode_wire_output", "retain_submitted_paths", "EntityKey", "EntityChangeSet", "EntityRoot", "ContextEntityRef", "ContextRootError", "CheckException", "CheckResult", "I18nCatalog", "JsonFieldNamingProfile", "Locale", "ObjectLocation", "UnsupportedLocaleError", "UserContext", "TeaqlRuntime", "SqlLogEntry", "SqlLogOperation", "DiagnosticSqlLogSink", "TextDiagnosticSqlLogSink", "ServiceRuntimeFromEnv", "RuntimeModule", "DataStore", "RawAuditEvent", "SafeAuditEvent", "MutationAuditKind", "BusinessClock", "FixedBusinessClock", "SystemBusinessClock", "AeadEntityReferenceCodec", "EntityReferenceClaims", "EntityReferenceCodec", "EntityReferenceTokenError", "ENTITY_REFERENCE_AAD", "UNSAFE_RAW_ENTITY_REFERENCES_ACKNOWLEDGEMENT", "UNSAFE_RAW_ENTITY_REFERENCES_ENVIRONMENT", "BUSINESS_ID_PERMUTATION_V1_ALPHABET", "BUSINESS_ID_PERMUTATION_V1_DOMAIN_SIZE", "BUSINESS_ID_PERMUTATION_V1_MAX_SEQUENCE", "BUSINESS_ID_PERMUTATION_V1_WIDTH", "DefaultBusinessIdProfileFactory", "DefaultBusinessIdService", "InMemoryBusinessIdAllocator", "PermutedDailyBusinessIdProfile", "StaticBusinessIdKeyProvider", "encode_business_id_permutation_v1", "ContextTools", "ExecutableHttpTool", "HTTP_TOOL", "HttpIntentPhase", "HttpTool", "HttpToolProvider", "ToolDeniedError", "ToolError", "ToolPolicy", "ToolRisk", "Tools", "ToolToken", "ToolUnavailableError"] __all__ += [ "MISSING_APPROVAL", "MISSING_POLICY", diff --git a/src/teaql/runtime/business_id.py b/src/teaql/runtime/business_id.py index d4535ec..4b98b24 100644 --- a/src/teaql/runtime/business_id.py +++ b/src/teaql/runtime/business_id.py @@ -5,12 +5,24 @@ import hashlib import hmac import struct +import asyncio +from datetime import datetime +import re from teaql.core.business_id import ( + BusinessIdAllocation, + BusinessIdAllocator, + BusinessIdDefinition, BusinessIdEncodingKey, BusinessIdError, BusinessIdErrorCode, + BusinessIdGenerationRequest, + BusinessIdKeyProvider, + BusinessIdPlan, + BusinessIdProfile, + BusinessIdSlot, BusinessIdScope, + BusinessIdValue, ) @@ -21,6 +33,165 @@ _MAGIC = b"teaql-business-id-fp-v1\0" +class InMemoryBusinessIdAllocator: + """Single-process allocator for tests and development.""" + + def __init__(self) -> None: + self._counters: dict[BusinessIdScope, int] = {} + self._lock = asyncio.Lock() + + async def allocate(self, plan: BusinessIdPlan) -> BusinessIdAllocation: + async with self._lock: + current = self._counters.get(plan.scope, plan.initial_sequence) + if current > plan.maximum_sequence: + raise BusinessIdError( + BusinessIdErrorCode.RANGE_EXHAUSTED, + f"Business ID range exhausted for {plan.scope.canonical_key()}", + ) + self._counters[plan.scope] = current + 1 + return BusinessIdAllocation(plan.scope, current) + + +class StaticBusinessIdKeyProvider: + def __init__(self, key: BusinessIdEncodingKey) -> None: + self._key = key + + def current_key(self, context, definition, scope) -> BusinessIdEncodingKey: + return self._key + + +class PermutedDailyBusinessIdProfile: + def __init__(self, context, key_provider: BusinessIdKeyProvider) -> None: + self._context = context + self._key_provider = key_provider + + def plan(self, request: BusinessIdGenerationRequest) -> BusinessIdPlan: + definition = request.definition + if ( + definition.profile != BusinessIdDefinition.DEFAULT_PROFILE + or definition.reset != "daily" + or definition.date_format != BusinessIdDefinition.DEFAULT_DATE_FORMAT + or definition.digits != BUSINESS_ID_PERMUTATION_V1_WIDTH + ): + raise BusinessIdError( + BusinessIdErrorCode.DEFINITION_INVALID, + "daily-permuted-v1 requires reset=daily, date_format=yyyyMMdd and digits=6", + ) + date_text = request.business_date.strftime("%Y%m%d") + scope = BusinessIdScope( + request.domain_root_key, + request.aggregate_type, + definition.namespace, + date_text, + ) + return BusinessIdPlan( + definition, scope, request.business_date, date_text, 0, + BUSINESS_ID_PERMUTATION_V1_MAX_SEQUENCE, + ) + + def format( + self, plan: BusinessIdPlan, allocation: BusinessIdAllocation + ) -> BusinessIdValue: + if allocation.scope != plan.scope: + raise BusinessIdError( + BusinessIdErrorCode.DEFINITION_INVALID, + "Allocation scope does not match Business ID plan", + ) + key = self._key_provider.current_key( + self._context, plan.definition, plan.scope + ) + if key is None: + raise BusinessIdError( + BusinessIdErrorCode.KEY_NOT_FOUND, + "Business ID key provider returned no current key", + ) + code = encode_business_id_permutation_v1( + allocation.sequence, plan.scope, key + ) + definition = plan.definition + return BusinessIdValue( + definition.separator.join((definition.prefix, plan.date_text, code)), + definition.profile, + definition.policy_version, + ) + + def validate( + self, definition: BusinessIdDefinition, value: str + ) -> BusinessIdValue: + escaped = re.escape(definition.separator) + pattern = re.compile( + rf"^{re.escape(definition.prefix)}{escaped}(\d{{8}}){escaped}([0-9A-Z]{{6}})$" + ) + match = pattern.fullmatch(value or "") + if match is None: + raise BusinessIdError( + BusinessIdErrorCode.FORMAT_INVALID, + f"Invalid daily-permuted-v1 Business ID: {value}", + ) + try: + datetime.strptime(match.group(1), "%Y%m%d") + except ValueError as error: + raise BusinessIdError( + BusinessIdErrorCode.FORMAT_INVALID, + f"Invalid daily-permuted-v1 Business ID: {value}", + ) from error + return BusinessIdValue(value, definition.profile, definition.policy_version) + + +class DefaultBusinessIdProfileFactory: + def create(self, context, definition: BusinessIdDefinition) -> BusinessIdProfile: + if definition.profile != BusinessIdDefinition.DEFAULT_PROFILE: + raise BusinessIdError( + BusinessIdErrorCode.PROFILE_NOT_FOUND, + f"Business ID profile is not registered: {definition.profile}", + ) + key_provider = context.get_resource("business_id_key_provider") + if key_provider is None: + raise BusinessIdError( + BusinessIdErrorCode.KEY_NOT_FOUND, + "Business ID key provider is not registered", + ) + return PermutedDailyBusinessIdProfile(context, key_provider) + + +class DefaultBusinessIdService: + def __init__(self, allocator: BusinessIdAllocator) -> None: + self._allocator = allocator + + async def ensure( + self, + context, + definition: BusinessIdDefinition, + domain_root_key: str, + aggregate_type: str, + slot: BusinessIdSlot, + ) -> BusinessIdValue: + factory = context.get_resource("business_id_profile_factory") + if factory is None: + raise BusinessIdError( + BusinessIdErrorCode.PROFILE_NOT_FOUND, + "Business ID profile factory is not registered", + ) + profile = factory.create(context, definition) + current = slot.current_value() + if current is not None and current.strip(): + return profile.validate(definition, current) + if not slot.new_aggregate(): + raise BusinessIdError( + BusinessIdErrorCode.IMMUTABLE, + "An established Aggregate cannot be assigned a new Business ID", + ) + plan = profile.plan(BusinessIdGenerationRequest( + definition, + domain_root_key, + aggregate_type, + context.business_date(), + )) + value = profile.format(plan, await self._allocator.allocate(plan)) + slot.assign_canonical_value(value.value) + return value + + def encode_business_id_permutation_v1( sequence: int, scope: BusinessIdScope, key: BusinessIdEncodingKey ) -> str: diff --git a/src/teaql/runtime/context.py b/src/teaql/runtime/context.py index b5573f1..a758daf 100644 --- a/src/teaql/runtime/context.py +++ b/src/teaql/runtime/context.py @@ -536,6 +536,29 @@ def with_internal_id_generator(self, gen: Any) -> 'UserContext': def set_internal_id_generator(self, gen: Any): self.insert_resource("internal_id_generator", gen) + def with_business_id_profile_factory(self, factory: Any) -> 'UserContext': + if factory is None or not callable(getattr(factory, "create", None)): + raise TypeError("business ID profile factory must provide create()") + return self.insert_resource("business_id_profile_factory", factory) + + def with_business_id_key_provider(self, provider: Any) -> 'UserContext': + if provider is None or not callable(getattr(provider, "current_key", None)): + raise TypeError("business ID key provider must provide current_key()") + return self.insert_resource("business_id_key_provider", provider) + + def with_business_id_service(self, service: Any) -> 'UserContext': + if service is None or not callable(getattr(service, "ensure", None)): + raise TypeError("business ID service must provide ensure()") + return self.insert_resource("business_id_service", service) + + async def ensure_business_id( + self, definition: Any, domain_root_key: str, aggregate_type: str, slot: Any + ) -> Any: + service = self.require_resource("business_id_service") + return await service.ensure( + self, definition, domain_root_key, aggregate_type, slot + ) + def with_schema_provider(self, provider: Any) -> 'UserContext': self.insert_resource("schema_provider", provider) return self diff --git a/src/teaql/sql/__init__.py b/src/teaql/sql/__init__.py index 5a594c1..1bc7995 100644 --- a/src/teaql/sql/__init__.py +++ b/src/teaql/sql/__init__.py @@ -11,6 +11,7 @@ SqlTransport, SqlExecutorError, CompileError, TransportError, SchemaProvider, SqlDataServiceExecutor ) +from .business_id import SqlBusinessIdAllocator __all__ = [ "DatabaseKind", "CompiledQuery", "SqlCompileError", @@ -21,5 +22,5 @@ "InvalidSubQueryOperatorError", "SqlDialect", "quote_identifier_if_needed", "SqlTransport", "SqlExecutorError", "CompileError", "TransportError", - "SchemaProvider", "SqlDataServiceExecutor" + "SchemaProvider", "SqlDataServiceExecutor", "SqlBusinessIdAllocator" ] diff --git a/src/teaql/sql/business_id.py b/src/teaql/sql/business_id.py new file mode 100644 index 0000000..a456fc5 --- /dev/null +++ b/src/teaql/sql/business_id.py @@ -0,0 +1,114 @@ +"""Durable SQL allocation for externally visible Business IDs.""" + +from __future__ import annotations + +import asyncio +from datetime import datetime, time, timezone + +from teaql.core.business_id import ( + BusinessIdAllocation, + BusinessIdError, + BusinessIdErrorCode, + BusinessIdPlan, +) +from teaql.core.value import Value +from teaql.sql.types import CompiledQuery + + +class SqlBusinessIdAllocator: + """Portable optimistic allocator. Schema installation remains explicit.""" + + MAX_ATTEMPTS = 100 + + def __init__(self, dialect, transport) -> None: + self._dialect = dialect + self._transport = transport + + async def allocate(self, plan: BusinessIdPlan) -> BusinessIdAllocation: + scope_key = plan.scope.canonical_key() + updated_at = int(datetime.combine( + plan.business_date, time.min, tzinfo=timezone.utc + ).timestamp() * 1000) + placeholders = [self._dialect.placeholder(index) for index in range(1, 6)] + last_conflict: Exception | None = None + + for _attempt in range(1, self.MAX_ATTEMPTS + 1): + try: + rows = await self._transport.fetch_all_sql(CompiledQuery( + "SELECT current_value, version FROM teaql_business_id_space " + f"WHERE scope_key = {placeholders[0]}", + [Value.from_any(scope_key)], + )) + except Exception as error: + raise RuntimeError( + "Business ID allocation requires explicit ensure_schema() " + "and an accessible teaql_business_id_space table" + ) from error + if not rows: + try: + changed, _ = await self._transport.execute_sql(CompiledQuery( + "INSERT INTO teaql_business_id_space " + "(scope_key, current_value, version, updated_at) VALUES " + f"({placeholders[0]}, {placeholders[1]}, 1, {placeholders[2]})", + [ + Value.from_any(scope_key), + Value.from_any(plan.initial_sequence), + Value.from_any(updated_at), + ], + )) + if changed == 1: + return BusinessIdAllocation( + plan.scope, plan.initial_sequence + ) + raise RuntimeError( + f"Business ID insert for {scope_key} changed {changed} rows" + ) + except Exception as error: + last_conflict = error + observed = await self._transport.fetch_all_sql(CompiledQuery( + "SELECT current_value, version FROM teaql_business_id_space " + f"WHERE scope_key = {placeholders[0]}", + [Value.from_any(scope_key)], + )) + if not observed: + raise RuntimeError( + "Business ID allocation requires explicit ensure_schema()" + ) from error + else: + current = int(rows[0]["current_value"]) + version = int(rows[0]["version"]) + if current < 0 or version < 1: + raise RuntimeError( + f"Invalid Business ID sequence row for {scope_key}" + ) + if current >= plan.maximum_sequence: + raise BusinessIdError( + BusinessIdErrorCode.RANGE_EXHAUSTED, + f"Business ID range exhausted for {scope_key}", + ) + next_value = current + 1 + changed, _ = await self._transport.execute_sql(CompiledQuery( + "UPDATE teaql_business_id_space SET current_value = " + f"{placeholders[0]}, version = version + 1, updated_at = {placeholders[1]} " + f"WHERE scope_key = {placeholders[2]} AND version = {placeholders[3]} " + f"AND current_value = {placeholders[4]}", + [ + Value.from_any(next_value), + Value.from_any(updated_at), + Value.from_any(scope_key), + Value.from_any(version), + Value.from_any(current), + ], + )) + if changed == 1: + return BusinessIdAllocation(plan.scope, next_value) + if changed != 0: + raise RuntimeError( + f"Business ID update for {scope_key} changed {changed} rows" + ) + await asyncio.sleep(0.001) + + raise BusinessIdError( + BusinessIdErrorCode.ALLOCATION_RETRY_EXHAUSTED, + f"Business ID allocation did not converge for {scope_key}: {last_conflict}", + ) diff --git a/src/teaql/sql/executor.py b/src/teaql/sql/executor.py index 59b00b6..2f4207e 100644 --- a/src/teaql/sql/executor.py +++ b/src/teaql/sql/executor.py @@ -1032,6 +1032,12 @@ async def _ensure_schema(self, context: 'UserContext', capability: object) -> No "CREATE TABLE IF NOT EXISTS teaql_id_space (" "type_name VARCHAR(255) NOT NULL PRIMARY KEY, " "current_level BIGINT NOT NULL)", [])) + await self.transport.execute_sql(CompiledQuery( + "CREATE TABLE IF NOT EXISTS teaql_business_id_space (" + "scope_key VARCHAR(512) NOT NULL PRIMARY KEY, " + "current_value BIGINT NOT NULL, " + "version BIGINT NOT NULL, " + "updated_at BIGINT NOT NULL)", [])) async def begin(self, context: 'UserContext') -> 'teaql.data_service.Transaction': if not isinstance(self.transport, SqlTransactionTransport): diff --git a/tests/provider/test_business_id_allocator.py b/tests/provider/test_business_id_allocator.py new file mode 100644 index 0000000..1bb9cfe --- /dev/null +++ b/tests/provider/test_business_id_allocator.py @@ -0,0 +1,86 @@ +import asyncio +from datetime import date + +import pytest + +from teaql.core import ( + BusinessIdDefinition, + BusinessIdError, + BusinessIdErrorCode, + BusinessIdPlan, + BusinessIdScope, +) +from teaql.provider.sqlite import create_sqlite_service +from teaql.runtime import UserContext +from teaql.sql import SqlBusinessIdAllocator + + +def plan(scope, maximum=100): + definition = BusinessIdDefinition.daily_permuted( + "order_number", "ORD", "order_number" + ) + return BusinessIdPlan( + definition, + scope, + date(2026, 10, 1), + "20261001", + 0, + maximum, + ) + + +@pytest.mark.asyncio +async def test_sqlite_business_id_allocator_requires_explicit_schema(tmp_path): + service = create_sqlite_service(str(tmp_path / "explicit.db")) + allocator = SqlBusinessIdAllocator(service.dialect, service.transport) + + with pytest.raises(RuntimeError, match="explicit ensure_schema"): + await allocator.allocate(plan(BusinessIdScope( + "root", "commerce_order", "order_number", "20261001" + ))) + + +@pytest.mark.asyncio +async def test_sqlite_business_id_allocator_is_shared_concurrent_and_restart_safe(tmp_path): + path = str(tmp_path / "business-ids.db") + first_service = create_sqlite_service(path) + second_service = create_sqlite_service(path) + await UserContext.new().insert_resource("dataService", first_service).ensure_schema() + first = SqlBusinessIdAllocator(first_service.dialect, first_service.transport) + second = SqlBusinessIdAllocator(second_service.dialect, second_service.transport) + shared = plan(BusinessIdScope( + "root", "commerce_order", "order_number", "20261001" + )) + + allocated = await asyncio.gather(*[ + (first if index % 2 == 0 else second).allocate(shared) + for index in range(40) + ]) + assert sorted(value.sequence for value in allocated) == list(range(40)) + + restarted_service = create_sqlite_service(path) + restarted = SqlBusinessIdAllocator( + restarted_service.dialect, restarted_service.transport + ) + assert (await restarted.allocate(shared)).sequence == 40 + + other_scope = plan(BusinessIdScope( + "other-root", "commerce_order", "order_number", "20261001" + )) + assert (await restarted.allocate(other_scope)).sequence == 0 + + +@pytest.mark.asyncio +async def test_sqlite_business_id_allocator_reports_capacity_exhaustion(tmp_path): + service = create_sqlite_service(str(tmp_path / "capacity.db")) + await UserContext.new().insert_resource("dataService", service).ensure_schema() + allocator = SqlBusinessIdAllocator(service.dialect, service.transport) + bounded = plan(BusinessIdScope( + "root", "commerce_order", "tiny", "20261001" + ), maximum=1) + + assert (await allocator.allocate(bounded)).sequence == 0 + assert (await allocator.allocate(bounded)).sequence == 1 + with pytest.raises(BusinessIdError) as failure: + await allocator.allocate(bounded) + assert failure.value.code == BusinessIdErrorCode.RANGE_EXHAUSTED diff --git a/tests/runtime/test_business_id_lifecycle.py b/tests/runtime/test_business_id_lifecycle.py new file mode 100644 index 0000000..fe28d53 --- /dev/null +++ b/tests/runtime/test_business_id_lifecycle.py @@ -0,0 +1,87 @@ +from datetime import datetime, timezone + +import pytest + +from teaql.core import ( + BusinessIdDefinition, + BusinessIdEncodingKey, + BusinessIdError, + BusinessIdErrorCode, +) +from teaql.runtime import ( + DefaultBusinessIdProfileFactory, + DefaultBusinessIdService, + FixedBusinessClock, + InMemoryBusinessIdAllocator, + StaticBusinessIdKeyProvider, + UserContext, +) + + +class Slot: + def __init__(self, value=None, new=True): + self.value = value + self.new = new + + def current_value(self): + return self.value + + def new_aggregate(self): + return self.new + + def assign_canonical_value(self, value): + self.value = value + + +def context_with(allocator): + return ( + UserContext.new() + .with_business_clock( + FixedBusinessClock(datetime(2026, 10, 1, 8, 30, tzinfo=timezone.utc)) + ) + .with_business_id_key_provider( + StaticBusinessIdKeyProvider(BusinessIdEncodingKey(1, bytes(range(32)))) + ) + .with_business_id_profile_factory(DefaultBusinessIdProfileFactory()) + .with_business_id_service(DefaultBusinessIdService(allocator)) + ) + + +@pytest.mark.asyncio +async def test_context_owned_business_id_lifecycle_is_idempotent_and_uses_business_date(): + allocator = InMemoryBusinessIdAllocator() + context = context_with(allocator) + definition = BusinessIdDefinition.daily_permuted( + "order_number", "ORD", "order_number" + ) + slot = Slot() + + first = await context.ensure_business_id( + definition, "commerce", "commerce_order", slot + ) + retry = await context.ensure_business_id( + definition, "commerce", "commerce_order", slot + ) + following = await context.ensure_business_id( + definition, "commerce", "commerce_order", Slot() + ) + + assert first == retry + assert slot.value == first.value + assert first.value.startswith("ORD-20261001-") + assert following.value != first.value + + +@pytest.mark.asyncio +async def test_established_aggregate_cannot_acquire_missing_business_id(): + context = context_with(InMemoryBusinessIdAllocator()) + definition = BusinessIdDefinition.daily_permuted( + "order_number", "ORD", "order_number" + ) + + with pytest.raises(BusinessIdError) as failure: + await context.ensure_business_id( + definition, "commerce", "commerce_order", Slot(new=False) + ) + + assert failure.value.code == BusinessIdErrorCode.IMMUTABLE diff --git a/tests/runtime/test_sql_mask_lifecycle.py b/tests/runtime/test_sql_mask_lifecycle.py index 300d11b..a90b60e 100644 --- a/tests/runtime/test_sql_mask_lifecycle.py +++ b/tests/runtime/test_sql_mask_lifecycle.py @@ -158,7 +158,9 @@ async def test_schema_failure_does_not_print_driver_error(fixture, monkeypatch, class SchemaFaultTransport(FaultTransport): async def execute_sql(self, query): - if 'CREATE TABLE' in query.sql and 'teaql_id_space' not in query.sql: + support_table = ('teaql_id_space' in query.sql + or 'teaql_business_id_space' in query.sql) + if 'CREATE TABLE' in query.sql and not support_table: raise RuntimeError('DRIVER-CANARY Riverside PASSWORD-CANARY') return 0, None