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:
ExceptionA referenced class does not pre-exist as the required kind.
- exception kapps_semantic_middleware.registration.OperationResolutionError[source]#
Bases:
ExceptionAn 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/andhttp://host:8000would 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 othersvc:addressand heartbeat).linedrather than a hash. Production IRIs remain back-resolvable.IRI.from_linedreads the address off the node IRI.addressis required. A default would silently reinstate the shared node.- Raises:
ValueError – If
addressis 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_iriorclass_iriis a subclass ofbase_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_classis notsvc:Serviceor 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) andsvc:endpoint. The Capability ownssvc: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) andsvc:endpoint. The Capability ownssvc: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:lastHeartbeatvia 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:addressbut either has nosvc:lastHeartbeator its latest heartbeat is older thanmax_age_seconds. Services are identified bysvc:addresspresence (real services are subclass-typed, sordf:type svc:Servicewould 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_queueneeds 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 itsrdf:type, the single-valuedcfc:implementsCapabilitylink, andsvc:operationStatus(plus any svc:-domaindata). 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:implementsCapabilitywithrdfs:domain cfc:Operation). The scenario/demo ontologies do. The Resource->Capabilitycfc:hasCapabilitylink written at registration is a separate, multi-valued append. It stays on the low-level path untilkapps_ogmgrows 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_classrealized by a Workflow whose Service has a livesvc:address:Capability --realizedByWorkflow--> Workflow --isWorkflowOf--> Service --address. Iftarget_resourceis given, the Service is pinned to that resource (svc:isServiceOf). The bound Capability instance is what the dispatched Operation willcfc: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.deleteis not implemented inkapps_ogmyet. This removes exactly whatcreate_operationwrote. Literal properties receive clear through the OGM commit path first (commit needs therdf:typetriple present to resolve the class). Thecfc:implementsCapabilityandrdf:typetriples then receive removal via the low-levelogm.db.triple_delete(a sanctioned kapps_triplestore_interface write).
- exception kapps_semantic_middleware.registration.OperationQueueEmpty[source]#
Bases:
ExceptionRaised 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 intosvc:failureStatein 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:operationStatusis one ofstatuses.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:isServiceOfhas 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
failedwith a reason via the OGM atomic commit path.Use to reclaim an orphaned
runningOperation (a resource crashed mid-execution) or to sweep a dead resource stranded Operations. Recovery is never a silent physical replay. The Operation goes tofailed. A planner decides whether to re-dispatch.
- exception kapps_semantic_middleware.registration.HandoverPreconditionError[source]#
Bases:
ExceptionA 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:PossessionStatethe 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_iriholdsworkpiece_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:complementsis 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:PossessionStatereceives mint for the new possessor. The new possessorcfc:hasPossessorreceives 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 workpiececfc:hasPossessedWorkpiecethen receives re-point via a singleOGM.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.