Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 25 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,11 @@ The SDK's organizational architecture strictly mirrors the Rust version:
* **Facet Aggregation & Grouping**: Out-of-the-box support for multi-dimensional facet aggregations, group-bys, and hierarchical data processing.
* **Provider Support**: Highly extensible asynchronous database connectivity (integrating third-party async drivers like `aiosqlite` through a unified Transport layer).
* **Context & Logging Management**: Built-in support for lifecycle context passing, end-to-end tracing, and SQL execution log interception and dispatch.
* **Governed Mutation Policy**: An application-owned policy can review an
immutable whole-graph plan after Checker/Fix and before the first provider
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.
* **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
Expand Down Expand Up @@ -104,5 +109,25 @@ development-only acknowledgement are maintained in the canonical
Opaque tokens never replace the backend's authorization, tenant, ownership,
role, or optimistic-version checks.

### Mutation Policy installation

Policy implementations are installed from trusted application startup through
`UserContext`; request JSON cannot select or replace them. Built-in SQL and TFP
providers enter the same governed boundary.

```python
context = (
UserContext.new()
.with_mutation_policy_registry(policy_registry)
.with_mutation_policy_approval_provider(approval_provider)
.with_mutation_governance_sink(warning_sink)
)
```

Generated graph saves call `preflight_mutation(...)` for every operation before
the first provider write. See the repeatable
[`examples/mutation-policy`](examples/mutation-policy) example for allow,
approval, audit propagation, and zero-write denial evidence.

---
To run test validations and business logic simulations locally, simply run `pytest` in the project root.
16 changes: 16 additions & 0 deletions examples/mutation-policy/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
# Mutation Policy example

This focused example installs an application-owned policy and exact approval
through `UserContext`, preflights an entire two-entity graph after Checker/Fix,
and proves that a denied graph reaches no persistent provider mutation. The
successful audit events retain the same governance snapshot.

Run it against the local runtime under development:

```bash
PYTHONPATH=src python examples/mutation-policy/main.py
```

The deterministic in-memory transaction keeps the example independent of a
database. Built-in SQL and TFP providers exercise the same policy boundary in
their focused runtime tests.
174 changes: 174 additions & 0 deletions examples/mutation-policy/main.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,174 @@
"""Focused Mutation Policy example with an atomic in-memory transaction."""

import asyncio
from datetime import datetime, timezone

from teaql.core.mutation import InsertCommand, MutationRequest, TraceNode
from teaql.data_service import DataServiceOperation, ExecutionMetadata, MutationResult
from teaql.runtime import (
DelegatingMutationPolicyApprovalProvider,
DelegatingMutationPolicyRegistry,
MutationDecision,
MutationPolicyApproval,
MutationPolicyApprovalStatus,
MutationPolicyError,
MutationPolicyIdentity,
UserContext,
)
from teaql.runtime.audit import MutationAuditKind, RawAuditEvent


class OrderPolicy:
identity = MutationPolicyIdentity(
"order-submission", "1", "sha256:order-submission-v1"
)

def review(self, context, plan):
if any(
operation.changed_values.get("name").try_text() == "DENIED"
for operation in plan.operations
if operation.changed_values.get("name") is not None
):
return MutationDecision.denied(
"ORDER_DENIED", "the order policy rejected this graph", "Order.name"
)
return MutationDecision.allowed()


class AuditRecorder:
def __init__(self):
self.events = []

async def on_safe_event(self, context, event):
self.events.append(event)


class MemoryTransaction:
def __init__(self, provider):
self.provider = provider
self.pending = []

async def mutate(self, context, request):
with context.mutation_policy_execution(request):
command = request._data
self.pending.append(command.entity)
await context.send_audit_event(
RawAuditEvent(
MutationAuditKind.CREATED,
command.entity,
command.values.get("id"),
(),
tuple(request.trace_chain()),
context.user_identifier(),
"mutation-policy-example",
context.current_mutation_governance(),
)
)
now = datetime.now(timezone.utc)
return MutationResult(
affected_rows=1,
generated_values={},
persisted_record={
key: value.val for key, value in command.values.items()
},
metadata=ExecutionMetadata(
backend="memory-example",
operation=DataServiceOperation.Insert,
started_at=now,
ended_at=now,
affected_rows=1,
),
)

async def commit(self, context):
self.provider.persisted.extend(self.pending)

async def rollback(self, context):
self.pending.clear()


class MemoryProvider:
def __init__(self):
self.persisted = []

async def begin(self, context):
return MemoryTransaction(self)


def insert(entity, entity_id, name):
command = (
InsertCommand.new(entity)
.value("id", entity_id)
.value("version", 1)
.value("name", name)
)
command.trace_chain.append(TraceNode(entity, entity_id, f"create {entity}"))
return command


def context_for(provider, audit):
policy = OrderPolicy()
return (
UserContext.new()
.insert_resource("dataService", provider)
.with_trace_id("python-mutation-policy-example")
.with_app_audit_event_sink(audit)
.with_mutation_policy_registry(
DelegatingMutationPolicyRegistry(lambda request_key: policy)
)
.with_mutation_policy_approval_provider(
DelegatingMutationPolicyApprovalProvider(
lambda identity: MutationPolicyApproval(
identity, "security-owner", datetime.now(timezone.utc)
)
)
)
)


async def save(context, *commands):
async def graph():
transaction = context.require_resource("dataService")
for command in commands:
context.preflight_mutation(command)
for command in commands:
await transaction.mutate(context, MutationRequest(command))

await context.execute_graph_save(graph)


async def main():
provider = MemoryProvider()
audit = AuditRecorder()
allowed = context_for(provider, audit)
await save(
allowed,
insert("Order", 42, "SUBMITTED"),
insert("OrderLine", 99, "LINE-1"),
)

denied_provider = MemoryProvider()
denied = context_for(denied_provider, AuditRecorder())
try:
await save(denied, insert("Order", 43, "DENIED"))
except MutationPolicyError as error:
assert "ORDER_DENIED" in str(error)
else:
raise AssertionError("denied graph unexpectedly persisted")

assert provider.persisted == ["Order", "OrderLine"]
assert denied_provider.persisted == []
assert len(audit.events) == 2
assert all(
event.mutation_governance.approval_status
== MutationPolicyApprovalStatus.APPROVED
for event in audit.events
)
print(
"PYTHON_MUTATION_POLICY_PASS "
"allowed_operations=2 denied_provider_mutations=0 audit_events=2"
)


if __name__ == "__main__":
asyncio.run(main())
3 changes: 2 additions & 1 deletion scripts/verify-examples.sh
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
set -euo pipefail

repo="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
expected=(conformance order-management school-management task_board)
expected=(conformance mutation-policy order-management 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
Expand All @@ -24,6 +24,7 @@ PYTHONPATH="$repo/examples/conformance:$repo/src" python -m app.main
PYTHONPATH="$repo/src" python -m unittest discover -s "$repo/examples/conformance" -p 'test_sql_log_intent.py' -v
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"
order_management_tmp="$(mktemp -d)"
task_board_tmp="$(mktemp -d)"
trap 'rm -rf "$order_management_tmp" "$task_board_tmp"' EXIT
Expand Down
47 changes: 26 additions & 21 deletions src/teaql/provider/tfp_client/__init__.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
from __future__ import annotations

from contextlib import nullcontext
from dataclasses import dataclass, field
from datetime import datetime
from typing import Any, Awaitable, Callable, Dict, Mapping, Optional
Expand Down Expand Up @@ -164,27 +165,31 @@ async def query(self, context: Any, request: QueryRequest) -> QueryResult:
)

async def mutate(self, context: Any, request: MutationRequest) -> MutationResult:
started_at = datetime.now()
federal = _mutation_request(request)
data = await self.federal_client.execute_mutation(federal)
records = data.get("data") or []
generated = records[0] if records else {}
affected = int(data.get("affectedRows", 0))
operation = {
"Create": DataServiceOperation.Insert,
"Update": DataServiceOperation.Update,
"Delete": DataServiceOperation.Delete,
"Recover": DataServiceOperation.Recover,
}[federal.action]
return MutationResult(
affected_rows=affected, generated_values=generated,
persisted_record=generated or None,
metadata=ExecutionMetadata(
backend="teaql-federal", operation=operation,
started_at=started_at, ended_at=datetime.now(), affected_rows=affected,
comment=federal.comment,
),
)
if context is not None and not context.consume_mutation_checked(request._data):
context.check_and_fix_mutation(request._data)
scope = context.mutation_policy_execution(request) if context is not None else nullcontext()
with scope:
started_at = datetime.now()
federal = _mutation_request(request)
data = await self.federal_client.execute_mutation(federal)
records = data.get("data") or []
generated = records[0] if records else {}
affected = int(data.get("affectedRows", 0))
operation = {
"Create": DataServiceOperation.Insert,
"Update": DataServiceOperation.Update,
"Delete": DataServiceOperation.Delete,
"Recover": DataServiceOperation.Recover,
}[federal.action]
return MutationResult(
affected_rows=affected, generated_values=generated,
persisted_record=generated or None,
metadata=ExecutionMetadata(
backend="teaql-federal", operation=operation,
started_at=started_at, ended_at=datetime.now(), affected_rows=affected,
comment=federal.comment,
),
)


def _federal_query_payload(query: FederalQuery) -> Dict[str, Any]:
Expand Down
32 changes: 32 additions & 0 deletions src/teaql/runtime/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,12 +10,44 @@
from .module import RuntimeModule, DefaultEntityDataServiceBehavior
from .store import DataStore
from .audit import RawAuditEvent, SafeAuditEvent, MutationAuditKind
from .mutation_policy import (
MISSING_APPROVAL,
MISSING_POLICY,
DelegatingMutationGovernanceSink,
DelegatingMutationPolicyApprovalProvider,
DelegatingMutationPolicyRegistry,
MutationDecision,
MutationGovernanceEvent,
MutationGovernanceSnapshot,
MutationOperation,
MutationOperationKind,
MutationOperationSummary,
MutationPlan,
MutationPolicyApproval,
MutationPolicyApprovalStatus,
MutationPolicyError,
MutationPolicyIdentity,
MutationPolicySource,
MutationVerdict,
)
from .i18n import CheckException, CheckResult, I18nCatalog, JsonFieldNamingProfile, Locale, ObjectLocation, UnsupportedLocaleError
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", "ContextTools", "ExecutableHttpTool", "HTTP_TOOL", "HttpIntentPhase", "HttpTool", "HttpToolProvider", "ToolDeniedError", "ToolError", "ToolPolicy", "ToolRisk", "Tools", "ToolToken", "ToolUnavailableError"]

__all__ += [
"MISSING_APPROVAL", "MISSING_POLICY",
"DelegatingMutationGovernanceSink",
"DelegatingMutationPolicyApprovalProvider",
"DelegatingMutationPolicyRegistry", "MutationDecision",
"MutationGovernanceEvent", "MutationGovernanceSnapshot",
"MutationOperation", "MutationOperationKind", "MutationOperationSummary",
"MutationPlan", "MutationPolicyApproval", "MutationPolicyApprovalStatus",
"MutationPolicyError", "MutationPolicyIdentity", "MutationPolicySource",
"MutationVerdict",
]


def __getattr__(name):
# Keep provider construction lazy: importing a SQL provider loads runtime
Expand Down
4 changes: 3 additions & 1 deletion src/teaql/runtime/audit.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ class RawAuditEvent:
trace_chain: tuple[Any, ...] = field(default_factory=tuple)
actor: Optional[str] = None
category: Optional[str] = None
mutation_governance: Any = None

def safe(self, mask_fields: List[str], max_length: Optional[int]) -> "SafeAuditEvent":
from .log_privacy import REDACTED, credential_name, payload_has_credentials, plaintext_enabled, scrub, value_strings
Expand All @@ -51,7 +52,7 @@ def safe(self, mask_fields: List[str], max_length: Optional[int]) -> "SafeAuditE
intent_values = secrets + value_strings(self.entity_id)
return SafeAuditEvent(
self.kind, self.entity, self.entity_id, scrub(tuple(fields), secrets), scrub(self.trace_chain, intent_values),
scrub(self.actor, intent_values), self.category,
scrub(self.actor, intent_values), self.category, self.mutation_governance,
)


Expand All @@ -72,6 +73,7 @@ class SafeAuditEvent:
trace_chain: tuple[Any, ...] = field(default_factory=tuple)
actor: Optional[str] = None
category: Optional[str] = None
mutation_governance: Any = None


def _mask(value: str) -> str:
Expand Down
Loading
Loading