Registration#

Putting an asset and its capabilities into the graph.

Graph-write core of the KAPPS Semantic Middleware.

Every knowledge-graph write goes through the OGM (OGM.create / OGM.commit). This is the architecture single validated write path (paper §4.4.2). No write issues raw SPARQL UPDATE. Reads may use the triple-store access module (ogm.db) directly. The architecture permits this.

Only one direction of each inverse relation receives materialization. The direction belongs to the newly-created instance (e.g. a workflow svc:isWorkflowOf, a capability svc:realizedByWorkflow). The container-side inverses (svc:hasWorkflow etc.) are OWL-inferable. They receive no write. Each registration is a clean single-instance OGM.create rather than a read-modify-write append onto a growing multi-valued property. Queries in this module therefore use the materialized (instance-owned) direction.

exception kapps_semantic_middleware.registration.OntologyGroundTruthError[source]#

Bases: Exception

A referenced class does not pre-exist as the required kind.

exception kapps_semantic_middleware.registration.OperationResolutionError[source]#

Bases: Exception

An Operation cannot resolve to a reachable Workflow endpoint.

kapps_semantic_middleware.registration.normalize_address(address: str) → str[source]#

Normalize an instance address. Trivially different spellings become one deployment.

Lowercase the scheme and the netloc (host names are case-insensitive, paths are not). Strip a trailing /. Without this, http://Host:8000/ and http://host:8000 would mint two Service nodes for one process. This discriminator prevents exactly that orphaning.

kapps_semantic_middleware.registration.mint_service_iri(resource_iri: IRI, address: str) → IRI[source]#

Mint the Service IRI for one middleware instance: {resource_iri}_service_{address}.

The discriminator is the instance normalized address mangled with IRI.lined. This satisfies two requirements at once. It is stable across restarts of the same deployment (a returning instance re-adopts its node). It is distinct between concurrent instances (a controller and a monitor on one resource do not overwrite each other svc:address and heartbeat). lined rather than a hash. Production IRIs remain back-resolvable. IRI.from_lined reads the address off the node IRI.

address is required. A default would silently reinstate the shared node.

Raises:

ValueError – If address is not an absolute http(s) URL.

kapps_semantic_middleware.registration.mint_workflow_iri(service_iri: IRI, name: str) → IRI[source]#

Mint the deterministic Workflow IRI: {service_iri}_workflow_{name}.

kapps_semantic_middleware.registration.mint_capability_iri(resource_iri: IRI, name: str) → IRI[source]#

Mint the deterministic Capability IRI: {resource_iri}_capability_{name}.

kapps_semantic_middleware.registration.mint_state_property_iri(service_iri: IRI, name: str) → IRI[source]#

Mint the deterministic StateProperty IRI: {service_iri}_state_{name}.

kapps_semantic_middleware.registration.build_workflow_endpoint(address: str, workflow_name: str) → str[source]#

Build the callable endpoint URL for a workflow (POST /workflows/{name}/execute).

kapps_semantic_middleware.registration.build_state_endpoint(address: str, state_name: str) → str[source]#

Build the GET endpoint URL for a state property (GET /state/{name}).

kapps_semantic_middleware.registration.assert_class_registered(ogm: OGM, class_iri: IRI, base_iri: IRI, named_graph: IRI | None = None) → None[source]#

Assert a class pre-exists as the required kind.

Valid iff class_iri == base_iri or class_iri is a subclass of base_iri. A read against the access module. A non-existent class makes the subclass check False. This also catches a missing class.

Raises:

OntologyGroundTruthError – If the class is missing or not the required subclass.

kapps_semantic_middleware.registration.register_service(ogm: OGM, *, resource_iri: IRI, service_iri: IRI, service_class: IRI, address: str, named_graph: IRI | None = None) → None[source]#

Register a Service instance for a resource. Set its base address.

Raises OntologyGroundTruthError if service_class is not svc:Service or a subclass.

kapps_semantic_middleware.registration.register_workflow(ogm: OGM, *, resource_iri: IRI, service_iri: IRI, workflow_iri: IRI, workflow_class: IRI, capability_iri: IRI, capability_class: IRI, endpoint: str, named_graph: IRI | None = None) → None[source]#

Register a Workflow instance, its Capability instance, and their links.

The Workflow owns svc:isWorkflowOf (-> service) and svc:endpoint. The Capability owns svc:realizedByWorkflow (-> workflow). The resource links to the Capability it provides (cfc:hasCapability, Resource -> Capability). The resource provides the Capability. Raises OntologyGroundTruthError if the classes are not the required subclasses.

kapps_semantic_middleware.registration.register_state_property(ogm: OGM, *, resource_iri: IRI, service_iri: IRI, state_property_iri: IRI, state_property_class: IRI, capability_iri: IRI, capability_class: IRI, endpoint: str, named_graph: IRI | None = None) → None[source]#

Register a StateProperty instance, its Capability instance, and their links.

The StateProperty owns svc:isStatePropertyOf (-> service) and svc:endpoint. The Capability owns svc:providedByStateProperty (-> state property). The resource links to the Capability it provides (cfc:hasCapability). Only the stable endpoint receives write. The live value receives no persist.

kapps_semantic_middleware.registration.deregister_service(ogm: OGM, service_iri: IRI, named_graph: IRI | None = None) → None[source]#

Remove a Service reachability (address + all workflow/state endpoints).

Endpoints receive removal via OGM.commit (the validated write path). The workflows/state- properties to clear are found via a read on the instance-owned inverse (svc:isWorkflowOf / svc:isStatePropertyOf). Structural triples and rdf:type are preserved.

kapps_semantic_middleware.registration.update_heartbeat(ogm: OGM, service_iri: IRI, timestamp: datetime | None = None, named_graph: IRI | None = None) → None[source]#

Refresh a Service svc:lastHeartbeat via OGM.commit (atomic replace).

kapps_semantic_middleware.registration.find_stale_services(ogm: OGM, max_age_seconds: float, *, now: datetime | None = None, named_graph: IRI | None = None) → list[IRI][source]#

Find reachable-but-silent Services whose heartbeat has gone stale (a read).

A service is stale if it currently has an svc:address but either has no svc:lastHeartbeat or its latest heartbeat is older than max_age_seconds. Services are identified by svc:address presence (real services are subclass-typed, so rdf:type svc:Service would not match).

kapps_semantic_middleware.registration.sweep_stale_services(ogm: OGM, max_age_seconds: float, *, now: datetime | None = None, named_graph: IRI | None = None) → list[IRI][source]#

Deregister every stale Service. Fail its resource stranded Operations.

Alongside removal of a dead resource reachability, the watchdog marks that resource stranded Operations (queued/running) failed. Work addressed to a resource that will never return does not hang forever. Returns the swept Service IRIs.

A resource may carry several Services. This is no longer quite right. A stale monitor drags its live sibling stranded Operations down with it. Fail them only when no sibling survives is also wrong. A surviving monitor realizes no Workflow. It would shield a dead controller queue forever. The correct predicate is per-Operation (“does any surviving Service realize this Operation Capability”). _reconstruct_queue needs the same treatment. Neither is done yet.

kapps_semantic_middleware.registration.EVENT_TRIGGER_WORKFLOW_NAME = 'event_trigger'#

Reserved built-in Workflow name for the receiver event-trigger REST endpoint.

kapps_semantic_middleware.registration.build_event_trigger_url(address: str) → str[source]#

Build the receiver event-trigger endpoint (POST /workflows/event_trigger/execute).

kapps_semantic_middleware.registration.mint_operation_iri(operation_class: IRI) → IRI[source]#

Mint a unique Operation instance IRI under the operation-class namespace.

Operations are per-dispatch. Each dispatch mints a fresh IRI (unlike the deterministic registration IRIs).

kapps_semantic_middleware.registration.create_operation(ogm: OGM, *, operation_iri: IRI, operation_class: IRI, capability_iri: IRI, status: str = 'queued', data: dict | None = None, named_graph: IRI | None = None) → None[source]#

Create an Operation individual for dispatch. Address it to a target Capability.

The whole Operation receives creation in ONE OGM.create. This includes its rdf:type, the single-valued cfc:implementsCapability link, and svc:operationStatus (plus any svc:-domain data). This is the single validated write path (paper §4.3/4.4.2, “every write originates as an OGM commit”). This requires the loaded ontology to declare the Core Operation-property domains (cfc:implementsCapability with rdfs:domain cfc:Operation). The scenario/demo ontologies do. The Resource->Capability cfc:hasCapability link written at registration is a separate, multi-valued append. It stays on the low-level path until kapps_ogm grows a validated single-triple append.

kapps_semantic_middleware.registration.resolve_dispatch_target(ogm: OGM, capability_class: IRI, target_resource: IRI | None = None, named_graph: IRI | None = None) → Tuple[IRI, IRI, str][source]#

Resolve a reachable receiver for a Capability class.

Bind a Capability instance of capability_class realized by a Workflow whose Service has a live svc:address: Capability --realizedByWorkflow--> Workflow --isWorkflowOf--> Service --address. If target_resource is given, the Service is pinned to that resource (svc:isServiceOf). The bound Capability instance is what the dispatched Operation will cfc:implementsCapability.

Returns:

Tuple of (capability_instance_iri, service_iri, service_address).

Raises:

OperationResolutionError – If no reachable Service realizes the capability class.

kapps_semantic_middleware.registration.revert_operation(ogm: OGM, operation_iri: IRI, *, operation_class: IRI, capability_iri: IRI, data: dict | None = None, named_graph: IRI | None = None) → None[source]#

Remove a just-created Operation whose event trigger failed to deliver.

Atomic create-and-notify: a failed notify reverts the created Operation. OGM.delete is not implemented in kapps_ogm yet. This removes exactly what create_operation wrote. Literal properties receive clear through the OGM commit path first (commit needs the rdf:type triple present to resolve the class). The cfc:implementsCapability and rdf:type triples then receive removal via the low-level ogm.db.triple_delete (a sanctioned kapps_triplestore_interface write).

exception kapps_semantic_middleware.registration.OperationQueueEmpty[source]#

Bases: Exception

Raised when claim_next finds no queued Operation to pull.

kapps_semantic_middleware.registration.resolve_operation_workflow(ogm: OGM, operation_iri: IRI, named_graph: IRI | None = None) → IRI[source]#

Resolve an Operation to the Workflow that realizes its Capability, WITHOUT require of a live endpoint. This is the provenance resolver used by pull-and-run so executedByWorkflow receives record even if the workflow svc:endpoint was deregistered mid-run.

Raises:

OperationResolutionError – If no Workflow realizes the Operation Capability.

kapps_semantic_middleware.registration.set_operation_status(ogm: OGM, *, operation_iri: IRI, status: str, named_graph: IRI | None = None) → None[source]#

Set svc:operationStatus on an Operation via the OGM atomic commit path.

This covers the queued->running transition and any other single-status transition.

kapps_semantic_middleware.registration.record_terminal_status(ogm: OGM, *, operation_iri: IRI, workflow_iri: IRI, status: str, result: str | None = None, failure_state: str | None = None, timestamp: datetime | None = None, named_graph: IRI | None = None) → None[source]#

Record terminal status (done/failed) with execution provenance in ONE commit.

Write status + provenance atomically (single commit). An Operation is never terminal-without-provenance. On failure, failure_state (a JSON snapshot of the resource datamodel) receives write into svc:failureState in the same commit. The failed status and the state that produced it are never separable.

kapps_semantic_middleware.registration.find_resource_operations(ogm: OGM, resource_iri: IRI, statuses: list[str], named_graph: IRI | None = None) → list[IRI][source]#

Find a resource own Operations whose svc:operationStatus is one of statuses.

Ontology-backed traversal: an Operation implements a Capability (cfc:implementsCapability). The resource provides that Capability (cfc:hasCapability). Both links persist across a restart. This works at early startup before re-registration. Explicitly does NOT match on IRI names.

Parameters:
  • ogm – The OGM instance.

  • resource_iri – The resource whose queue to inspect.

  • statuses – Iterable of status strings to filter (e.g. ["queued", "running"]).

  • named_graph – Optional named graph for the read.

Returns:

List of Operation IRIs matching the criteria.

kapps_semantic_middleware.registration.services_of_resource(ogm: OGM, resource_iri: IRI, *, reachable_only: bool = False, named_graph: IRI | None = None) → list[IRI][source]#

Every Service bound to a resource, via the instance-owned svc:isServiceOf.

A Service IRI carries an instance discriminator. It can no longer receive reconstruction from a resource IRI alone. This read replaces that reconstruction. svc:isServiceOf has always been many-to-one. Consumers must no longer assume the answer has exactly one element. A controller and a read-only monitor on one resource are two Services.

Parameters:
  • ogm – The OGM instance.

  • resource_iri – The resource whose services to find.

  • reachable_only – Keep only Services that currently carry an svc:address. A deregistered instance keeps its individual but loses its address. This is the reachable-now set rather than the ever-registered set.

  • named_graph – Optional named graph for the read.

Returns:

List of Service IRIs.

kapps_semantic_middleware.registration.mark_operation_failed(ogm: OGM, operation_iri: IRI, reason: str, named_graph: IRI | None = None) → None[source]#

Transition an Operation to failed with a reason via the OGM atomic commit path.

Use to reclaim an orphaned running Operation (a resource crashed mid-execution) or to sweep a dead resource stranded Operations. Recovery is never a silent physical replay. The Operation goes to failed. A planner decides whether to re-dispatch.

exception kapps_semantic_middleware.registration.HandoverPreconditionError[source]#

Bases: Exception

A handover precondition failed (caller lacks possession, or counterpart lacks the complementary handover ability). Rejection occurs before any physical work begins.

kapps_semantic_middleware.registration.mint_possession_state_iri(workpiece_iri: IRI) → IRI[source]#

Mint a fresh PossessionState IRI. Each change of possession makes a new state. The previous one is kept as implicit history.

kapps_semantic_middleware.registration.create_possession(ogm: OGM, *, workpiece_iri: IRI, possessor_iri: IRI, named_graph: IRI | None = None) → IRI[source]#

Establish an initial possession over Core reified model (verified vs Core 0.9.0).

A cfc:PossessionState the possessor holds (cfc:hasPossessor, appended as an atomic insertion since a resource may possess several workpieces) and the workpiece is possessed by (cfc:hasPossessedWorkpiece, set through the OGM commit/update path). The scenario ontology must declare these Core terms domains so the OGM commit validates.

Returns:

The minted PossessionState IRI.

kapps_semantic_middleware.registration.find_possession_state(ogm: OGM, workpiece_iri: IRI, possessor_iri: IRI, named_graph: IRI | None = None) → IRI | None[source]#

The current PossessionState in which possessor_iri holds workpiece_iri, or None.

An ontology-backed read (workpiece cfc:hasPossessedWorkpiece ?ps . possessor cfc:hasPossessor ?ps). Never match on IRI names.

kapps_semantic_middleware.registration.counterpart_has_complementary_ability(ogm: OGM, counterpart_iri: IRI, mode_ability_iri: IRI, named_graph: IRI | None = None) → bool[source]#

True iff the counterpart carries the handover ability COMPLEMENTARY to mode_ability_iri.

mes:complements is symmetric. This is a read (mode_ability mes:complements ?a . counterpart mes:hasHandoverAbility ?a).

kapps_semantic_middleware.registration.switch_possession(ogm: OGM, *, workpiece_iri: IRI, new_possessor_iri: IRI, named_graph: IRI | None = None) → IRI[source]#

Atomically change possession of a workpiece to a new possessor.

A fresh cfc:PossessionState receives mint for the new possessor. The new possessor cfc:hasPossessor receives APPEND as an atomic insertion (kapps_triplestore_interface). A resource may possess several workpieces. This must not disturb its existing possessions, because possession is not universally maxCount 1. The workpiece cfc:hasPossessedWorkpiece then receives re-point via a single OGM.commit. This is the update path, an atomic DELETE/INSERT. It keeps the workpiece pointed to exactly one PossessionState throughout. Core Workpiece cardinality-1 (the commit-time SHACL backstop) is never transiently violated. Append the possessor link BEFORE the re-point. A failed re-point leaves the prior possession fully intact. The previous PossessionState is kept as implicit history.

Returns:

The new PossessionState IRI.