Middleware#

The entry point. Extends transitional_sync_middleware.Middleware with graph registration, discovery and execution.

KAPPS Semantic Middleware.

Extends transitional_sync_middleware.Middleware with knowledge-graph registration, discovery, and execution capabilities. Supports three modes – see modes.Mode for what each one means and which are implemented.

Operation execution: an Operation resolves via its implemented Capability to a Workflow endpoint, which the system then invokes over HTTP.

class kapps_semantic_middleware.middleware.SemanticMiddleware(*, mode: Mode = Mode.RESOURCE, resource_iri: str | None = None, service_class: str | None = None, ogm: Any = None, host: str = '127.0.0.1', port: int = 8000, address: str | None = None, named_graph: str | None = None, heartbeat_interval: float | None = 30.0, staleness_threshold: float = 90.0, sweep_interval: float = 30.0, class_scope: Any = None, autoregister_connectors: bool = True, connector_sync_direction: SyncDirection = SyncDirection.BIDIRECTIONAL, connector_registry: SemanticConnectorRegistry | None = None, ensure_transport: Callable[[str, int], None] | None = None, activity_feed: bool = False, activity_capacity: int = 200)[source]#

Bases: Middleware

KAPPS Semantic Middleware extending transitional_sync_middleware.Middleware.

Adds knowledge-graph registration, discovery, and execution. It supports the three modes of Mode. Resource mode registers the Service/Workflow/Capability instances on startup, and deregisters them on shutdown. It is the only mode with runtime consequence today.

mode accepts a Mode constant or the equivalent bare string – Mode is a str subclass, so existing mode="resource" callers are unaffected.

Operation execution: an Operation resolves via its implemented Capability to a Workflow endpoint, which is then invoked over HTTP.

APP_TITLE = 'semantic-middleware'#

What this product calls itself. The base class names the framework it is built on, which is accurate for the framework and wrong for anything hitting our address – a visitor landing on a middleware instance should be told what it reached. The dependency policy allows only bugfixes in sibling repos, so the name is corrected here rather than there: the base class stays as it is, and every override lives on this side.

property app#

The FastAPI app, with its root route renamed and a quiet favicon route added.

The base class registers GET / inside its own app property, so the route exists the moment the app does. Starlette matches routes in order and a second registration would never be reached, so the original is removed rather than shadowed. Guarded by a flag: the base property caches, but this one is consulted on every access, and neither patch below may run twice on the same app object.

The favicon route answers GET /favicon.ico with a bare 204: this product ships no icon asset, and every browser that requests one otherwise logs a 404 – the only console error a control station or a unit middleware’s page would otherwise show on a clean load.

The nested root function handles GET / requests. It returns the welcome message defined in APP_WELCOME.

The nested favicon function handles GET /favicon.ico requests. It returns an empty 204 response to suppress browser console errors.

async emit_heartbeat() → None[source]#

Refresh this service’s heartbeat once (also useful for tests/manual pings).

async sweep() → List[str][source]#

Deregister every stale service once. Returns the swept service IRIs (as str).

A watchdog-mode instance calls this on its sweep_interval. It is also a plain method, so a test or operator triggers a sweep directly.

workflow(*args, capability_class: Any = None, workflow_class: Any = None, **kwargs)[source]#

Register a function as a REST-invokable workflow with KG registration.

In resource mode, requires capability_class and workflow_class as keyword-only IRIs (both classes must pre-exist in the ontology). Registers the REST endpoint via the base class, then schedules KG registration of the Workflow and Capability instances on startup (after the service is registered).

Raises:
  • RuntimeError – If called outside resource mode.

  • ValueError – If capability_class or workflow_class is missing.

state(*, capability_class: Any = None, state_property_class: Any = None, name: str | None = None)[source]#

Expose a getter as a GET-readable state property with KG registration.

Parallel to workflow() but GET-only, for a readable, potentially high-frequency-changing value (e.g. a door’s status). In resource mode, it requires capability_class and state_property_class as keyword-only IRIs (both classes must pre-exist in the ontology). Registers a GET endpoint at /state/{name} that calls the decorated getter on demand. The live value is NEVER written to the graph. Only the stable endpoint triple is written, at registration.

Raises:
  • RuntimeError – If called outside resource mode.

  • ValueError – If capability_class or state_property_class is missing.

request(*, capability_class: Any, operation_class: Any, operation_iri: str | None = None, target_resource: str | None = None)[source]#

Caller-side dispatch as a transaction context manager.

Usage:

with mw.request(capability_class=cap, operation_class=op_cls) as op:
    ...  # populate op.data with the operation's domain fields (if any)

The body populates the yielded draft. On clean exit the middleware atomically creates the Operation (status queued, addressed to a reachable Service via its Capability through discovery), and triggers that Service’s event trigger over REST. If the trigger delivery fails, the system reverts the created Operation atomically. A body exception aborts before anything is written.

This is the in-process caller face — not REST-exposed.

claim_next(scope: Any = None)[source]#

Pull-and-run the next queued Operation as a transaction context manager.

Usage:

with mw.claim_next(scope) as claimed:
    claimed.result = do_the_work(claimed.operation)

__enter__ pops the next queued Operation from this instance’s in-memory queue (FIFO order, raises OperationQueueEmpty if none). It marks the Operation running, and re-fetches it under the domain-supplied ClassScope, so the body gets exactly the object shape it needs. The body runs the work, and may set claimed.result. On clean exit the CM atomically records done plus provenance (executedByWorkflow / executionTimestamp / executionResult). On a body exception it records failed plus the exception message, and re-raises. Provenance is folded into the terminal transition. On failure the CM also dumps the resource datamodel and stores it as the Operation’s svc:failureState, so the state that produced the failure survives the exception.

This is the in-process receiver face — not REST-exposed.

handover(*, mode: Any, workpiece: Any, counterpart: Any)[source]#

Change-of-possession primitive as a transaction context manager.

Usage:

with mw.handover(mode=MES.Pass, workpiece=wp, counterpart=other):
    ...  # domain-owned physical transport / counterpart coordination

__enter__ runs exactly two precondition checks OUTSIDE the transaction. First, the caller currently possesses the workpiece (a cfc:PossessionState the workpiece points to, with this resource as possessor). Second, the counterpart carries the mes:hasHandoverAbility COMPLEMENTARY to mode. There is deliberately no destination-free, universal maxCount-1 check. Possession is not always single. The body is domain-owned: physical transport, and any counterpart coordination, for example through a request(…) dispatch. The Core never references the event trigger. On clean exit, __exit__ switches possession to the counterpart. The middleware appends the counterpart’s cfc:hasPossessor to a fresh PossessionState (an atomic insertion, because a resource can possess several workpieces). The middleware then re-points the workpiece’s cardinality-1 cfc:hasPossessedWorkpiece, in a single OGM DELETE/INSERT along the update path. The workpiece thus points to exactly one PossessionState throughout the process. On an exception, the middleware aborts with no switch. Possession is Core’s reified model (cfc:PossessionState). The “possessed by exactly one” cardinality is Core’s own Workpiece restriction, the commit-time SHACL backstop.

register_callback(callback: Any, scope: Any = None) → None[source]#

Register a domain work callback fired on enqueue.

callback is the work function callback(operation) -> result. It runs the Operation’s work, and returns the value recorded as svc:executionResult. On each enqueue, The middleware drives a background pull-and-run around it: with claim_next(scope) as c: c.result = callback(c.operation). The middleware handles status (running->`done`/failed) and provenance. The domain writes only the work. Opt-in: with no callback registered, an enqueued Operation stays queued for a manual claim_next. scope is the domain ClassScope used to re-fetch the Operation.

exception kapps_semantic_middleware.middleware.OperationResolutionError[source]#

Bases: Exception

An Operation cannot resolve to a reachable Workflow endpoint.

class kapps_semantic_middleware.middleware.Mode(*values)[source]#

Bases: str, Enum

A middleware instance mode. Each member below says what it means.

RESOURCE = 'resource'#

Wrap one resource_iri. Register a Service. Serve its workflows and parameters. Hold a heartbeat. This is the only mode with runtime consequence today.

SERVER = 'server'#

Reserved. Serve data with no physical resource. Not implemented. Construct one raises. Out of scope for the scenario-3 controller: a controller consumes a graph, and it does not serve one.

WATCHDOG = 'watchdog'#

Reserved. Sweep liveness from a central point. Sweep stale Services. Register nothing of its own.