diff --git a/emulator/data/827PD2708_demo.csv b/emulator/data/827PD2708_demo.csv new file mode 100644 index 0000000..5eb7581 --- /dev/null +++ b/emulator/data/827PD2708_demo.csv @@ -0,0 +1,151 @@ +timestamp,value +2026-05-23T12:00:00+01:00,20.098790 +2026-05-23T12:01:00+01:00,19.893183 +2026-05-23T12:02:00+01:00,19.128400 +2026-05-23T12:03:00+01:00,20.899973 +2026-05-23T12:04:00+01:00,20.424315 +2026-05-23T12:05:00+01:00,21.179129 +2026-05-23T12:06:00+01:00,21.512134 +2026-05-23T12:07:00+01:00,19.301748 +2026-05-23T12:08:00+01:00,19.509385 +2026-05-23T12:09:00+01:00,19.161866 +2026-05-23T12:10:00+01:00,18.846243 +2026-05-23T12:11:00+01:00,19.438751 +2026-05-23T12:12:00+01:00,20.980735 +2026-05-23T12:13:00+01:00,20.570343 +2026-05-23T12:14:00+01:00,18.920447 +2026-05-23T12:15:00+01:00,20.168629 +2026-05-23T12:16:00+01:00,20.990209 +2026-05-23T12:17:00+01:00,19.524368 +2026-05-23T12:18:00+01:00,19.984025 +2026-05-23T12:19:00+01:00,20.509962 +2026-05-23T12:20:00+01:00,20.598713 +2026-05-23T12:21:00+01:00,21.497972 +2026-05-23T12:22:00+01:00,20.777731 +2026-05-23T12:23:00+01:00,20.665232 +2026-05-23T12:24:00+01:00,20.669859 +2026-05-23T12:25:00+01:00,21.216786 +2026-05-23T12:26:00+01:00,20.895265 +2026-05-23T12:27:00+01:00,21.679000 +2026-05-23T12:28:00+01:00,20.941076 +2026-05-23T12:29:00+01:00,20.639256 +2026-05-23T12:30:00+01:00,20.805094 +2026-05-23T12:31:00+01:00,20.563094 +2026-05-23T12:32:00+01:00,19.323915 +2026-05-23T12:33:00+01:00,21.916487 +2026-05-23T12:34:00+01:00,20.212019 +2026-05-23T12:35:00+01:00,19.381460 +2026-05-23T12:36:00+01:00,19.110708 +2026-05-23T12:37:00+01:00,20.289606 +2026-05-23T12:38:00+01:00,20.047132 +2026-05-23T12:39:00+01:00,19.891442 +2026-05-23T12:40:00+01:00,21.089875 +2026-05-23T12:41:00+01:00,19.441301 +2026-05-23T12:42:00+01:00,22.249905 +2026-05-23T12:43:00+01:00,21.201351 +2026-05-23T12:44:00+01:00,18.932340 +2026-05-23T12:45:00+01:00,19.404163 +2026-05-23T12:46:00+01:00,21.012815 +2026-05-23T12:47:00+01:00,20.813325 +2026-05-23T12:48:00+01:00,19.727619 +2026-05-23T12:49:00+01:00,18.990339 +2026-05-23T12:50:00+01:00,19.417946 +2026-05-23T12:51:00+01:00,20.618113 +2026-05-23T12:52:00+01:00,18.831698 +2026-05-23T12:53:00+01:00,19.265574 +2026-05-23T12:54:00+01:00,19.230057 +2026-05-23T12:55:00+01:00,19.870521 +2026-05-23T12:56:00+01:00,18.971600 +2026-05-23T12:57:00+01:00,18.619893 +2026-05-23T12:58:00+01:00,18.879316 +2026-05-23T12:59:00+01:00,19.495739 +2026-05-23T13:00:00+01:00,-0.557326 +2026-05-23T13:01:00+01:00,-0.596192 +2026-05-23T13:02:00+01:00,-0.591885 +2026-05-23T13:03:00+01:00,-0.456344 +2026-05-23T13:04:00+01:00,-0.569352 +2026-05-23T13:05:00+01:00,-0.527288 +2026-05-23T13:06:00+01:00,-0.584083 +2026-05-23T13:07:00+01:00,-0.595779 +2026-05-23T13:08:00+01:00,-0.600629 +2026-05-23T13:09:00+01:00,-0.597288 +2026-05-23T13:10:00+01:00,-0.593904 +2026-05-23T13:11:00+01:00,-0.593569 +2026-05-23T13:12:00+01:00,-0.597801 +2026-05-23T13:13:00+01:00,-0.596797 +2026-05-23T13:14:00+01:00,-0.595079 +2026-05-23T13:15:00+01:00,-0.593314 +2026-05-23T13:16:00+01:00,-0.595152 +2026-05-23T13:17:00+01:00,-0.590297 +2026-05-23T13:18:00+01:00,-0.593574 +2026-05-23T13:19:00+01:00,-0.598075 +2026-05-23T13:20:00+01:00,-0.595300 +2026-05-23T13:21:00+01:00,-0.595021 +2026-05-23T13:22:00+01:00,-0.598563 +2026-05-23T13:23:00+01:00,-0.594525 +2026-05-23T13:24:00+01:00,-0.595928 +2026-05-23T13:25:00+01:00,-0.592931 +2026-05-23T13:26:00+01:00,-0.598961 +2026-05-23T13:27:00+01:00,-0.598864 +2026-05-23T13:28:00+01:00,-0.599080 +2026-05-23T13:29:00+01:00,-0.592581 +2026-05-23T13:30:00+01:00,20.746060 +2026-05-23T13:31:00+01:00,20.574841 +2026-05-23T13:32:00+01:00,21.186586 +2026-05-23T13:33:00+01:00,18.677824 +2026-05-23T13:34:00+01:00,19.498506 +2026-05-23T13:35:00+01:00,21.057821 +2026-05-23T13:36:00+01:00,19.428906 +2026-05-23T13:37:00+01:00,19.718633 +2026-05-23T13:38:00+01:00,20.679431 +2026-05-23T13:39:00+01:00,18.775636 +2026-05-23T13:40:00+01:00,21.303574 +2026-05-23T13:41:00+01:00,19.944833 +2026-05-23T13:42:00+01:00,19.987845 +2026-05-23T13:43:00+01:00,19.433113 +2026-05-23T13:44:00+01:00,20.264441 +2026-05-23T13:45:00+01:00,22.184703 +2026-05-23T13:46:00+01:00,19.246954 +2026-05-23T13:47:00+01:00,18.714419 +2026-05-23T13:48:00+01:00,22.007063 +2026-05-23T13:49:00+01:00,19.590064 +2026-05-23T13:50:00+01:00,20.329044 +2026-05-23T13:51:00+01:00,18.683035 +2026-05-23T13:52:00+01:00,20.272753 +2026-05-23T13:53:00+01:00,19.364303 +2026-05-23T13:54:00+01:00,19.976835 +2026-05-23T13:55:00+01:00,20.413678 +2026-05-23T13:56:00+01:00,21.437960 +2026-05-23T13:57:00+01:00,19.686194 +2026-05-23T13:58:00+01:00,18.483160 +2026-05-23T13:59:00+01:00,21.120720 +2026-05-23T14:00:00+01:00,20.759554 +2026-05-23T14:01:00+01:00,19.664443 +2026-05-23T14:02:00+01:00,19.837872 +2026-05-23T14:03:00+01:00,19.018561 +2026-05-23T14:04:00+01:00,19.692430 +2026-05-23T14:05:00+01:00,19.814016 +2026-05-23T14:06:00+01:00,19.441306 +2026-05-23T14:07:00+01:00,19.762021 +2026-05-23T14:08:00+01:00,18.786839 +2026-05-23T14:09:00+01:00,20.697803 +2026-05-23T14:10:00+01:00,20.860491 +2026-05-23T14:11:00+01:00,18.999857 +2026-05-23T14:12:00+01:00,21.310931 +2026-05-23T14:13:00+01:00,20.317433 +2026-05-23T14:14:00+01:00,21.124394 +2026-05-23T14:15:00+01:00,20.165350 +2026-05-23T14:16:00+01:00,20.181338 +2026-05-23T14:17:00+01:00,19.682744 +2026-05-23T14:18:00+01:00,19.523780 +2026-05-23T14:19:00+01:00,19.845955 +2026-05-23T14:20:00+01:00,19.576923 +2026-05-23T14:21:00+01:00,19.062933 +2026-05-23T14:22:00+01:00,20.332759 +2026-05-23T14:23:00+01:00,19.544526 +2026-05-23T14:24:00+01:00,19.267434 +2026-05-23T14:25:00+01:00,21.256525 +2026-05-23T14:26:00+01:00,19.767486 +2026-05-23T14:27:00+01:00,19.110618 +2026-05-23T14:28:00+01:00,20.604229 +2026-05-23T14:29:00+01:00,19.982214 diff --git a/nurse/README.md b/nurse/README.md index 27b854e..abeca9b 100644 --- a/nurse/README.md +++ b/nurse/README.md @@ -1,35 +1,277 @@ # mbaigo System: Nurse -The Nurse is a system that reacts to signals that exceed a threshold, which triggers a maintenance request. +## Purpose -## Status -As with the other systems, this is a prototype that shows that the mbaigo library can be used with ease. +The *Nurse* is the condition-monitoring consumer of an Arrowhead local cloud. +Like a duty nurse on a ward, it watches a set of physical signals against a +range, raises a maintenance request when one of them goes — and stays — out +of bounds, and stops bothering anyone about that signal until the requested +work is done. -## Compiling -After cloning the code, initialize the *go.mod* file with ``` go mod init github.com/sdoque/systems/nurse``` before running *go mod tidy*. +The Nurse's value-add over a thermostat-style controller is its **diagnostic +context**: it knows, for every sensor it polls, which downstream **actuator** +that sensor diagnoses (resolved from a knowledge graph), and it raises the +work order against the *actuator*, not the sensor that observed the symptom. +A pressure-differential sensor reads low → the work order targets the valve +the sensor is observing, with the sensor identity preserved in the request +for traceability. -The reason the *go.mod* file is not included in the repository is that when developing the mbaigo module, a replace statement needs to be included to point to the development code. +## Architecture -To run the code, one just needs to type in ```go run .``` within a terminal or at a command prompt. +Four moving pieces: -It is **important** to start the program from within its own directory (and each system should have their own directory) because it looks for its configuration file there. If it does not find it there, it will generate one and shutdown to allow the configuration file to be updated. +1. **Per-signal samplers.** One goroutine per signal in `Signals[]`, each + polling every `samplingPeriod` seconds and counting consecutive + out-of-range readings per source. Five in a row trips the alarm. +2. **Range-based threshold.** Each signal carries a `lowerThreshold` and + `upperThreshold`. Alert fires on `value < lower || value > upper`. Resets + when the value comes back in range. Direction-agnostic — caught equally + by a collapse or a spike. +3. **Sensor → actuator resolver.** At first discovery of each new provider + node, the Nurse queries a GraphDB triple store with the sensor's name and + caches the actuator's functional-location tag. A node whose actuator + cannot be resolved is marked **unresolvable** and skipped permanently — + raising a misdirected work order is worse than raising none. +4. **Two endpoints back from the Sapper.** Beyond the polling loop, the + Nurse exposes `/monitor` (the standard `SignalMonitoring` callback for + TECO completion) and `/enrichment` (the `EnrichmentNotification` + endpoint where the human planner's submitted operations payload + arrives at REL). Both are dispatched at the same `HealthTracker` + asset. -The configuration and operation of the system can be verified using the system's web server using a standard web browser, whose address is provided by the system at startup. +## How it fits in the cloud -To build the software for one's own machine, -```go build -o nurse```. +``` +Emulator/sensor ──poll──► Nurse ─┬─► Sapper (raise CRTD order) + │ + └─► GraphDB (publish sensor-side facts) + ◄───────── REL enrichment ─────── Sapper + ◄───────── TECO completion ────── Sapper +``` +The Nurse polls a signal provider (e.g. the `emulator` replaying a CSV), and +when it fires it (a) discovers the Sapper via Arrowhead's `MaintenanceOrder` +service definition and POSTs the order, and (b) publishes the sensor-side +context to GraphDB so the order's full lineage is queryable. Later, the +Sapper's REL and TECO events arrive at `/enrichment` and `/monitor` +respectively, and the Nurse logs the enrichment payload and restores the +signal to operational status. -## Cross compiling/building -The following commands enable one to build for different platforms: -- Intel Mac: ```GOOS=darwin GOARCH=amd64 go build -o nurse_imac``` -- ARM Mac: ```GOOS=darwin GOARCH=arm64 go build -o nurse_amac ``` -- Windows 64: ```GOOS=windows GOARCH=amd64 go build -o nurse.exe``` -- Raspberry Pi 64: ```GOOS=linux GOARCH=arm64 go build -o nurse_rpi64``` -- Linux: ```GOOS=linux GOARCH=amd64 go build -o nurse_linux``` +## Sequence diagram -One can find a complete list of platform by typing *‌go tool dist list* at the command prompt +```mermaid +sequenceDiagram + autonumber + participant E as Emulator + participant N as Nurse + participant OR as Orchestrator + participant G as GraphDB + participant S as Sapper -If one wants to secure copy it to a Raspberry pi, -`scp nurse_rpi64 jan@192.168.1.10:rpiExec/nurse/` where user is the *username* @ the *IP address* of the Raspberry Pi with a relative (to the user's home directory) target *rpiExec/nurse/* directory.nurse + Note over N: first poll + N->>OR: discover signal providers + OR-->>N: Emulator URL + N->>G: resolve sensor → valve (one query per new node) + G-->>N: valve FL tag (e.g. "827-PV2708-200") + loop every samplingPeriod + N->>E: GET /emulator/{sensor}/access + E-->>N: SignalA_v1a {value, unit, timestamp} + end + + Note over N: value out of range for 5 consecutive samples + N->>OR: discover MaintenanceOrder + OR-->>N: Sapper URL + N->>S: POST /sapper/SAPSimulator/maintenanceorders + S-->>N: 201 Created {orderId, status:"CRTD"} + N->>G: INSERT (bySensor, targetFLTag, reason) + + Note over N: signal in WorkRequested state — no more polls for this node + S-->>N: POST /nurse/HealthTracker/enrichment (planner submitted) + N->>N: log enrichment payload + S-->>N: POST /nurse/HealthTracker/monitor (TECO) + N->>N: signal restored to operational, counter reset +``` + +## Services + +### Provided + +| Service definition | Subpath | Methods | Description | +|--------------------|---------|---------|-------------| +| `SignalMonitoring` | `monitor` | `GET` | Plain-text status of every monitored signal: range, consecutive out-of-range counts per node, operational flag | +| `SignalMonitoring` | `monitor` | `POST` | TECO completion callback (called by the Sapper). Marks the signal back to operational | +| `EnrichmentNotification` | `enrichment` | `POST` | REL enrichment notification (called by the Sapper). Logs the planner's operations payload | + +### Consumed (via Arrowhead orchestration) + +| Service definition | Used for | +|--------------------|----------| +| `` (one per entry in `signals[]`, e.g. `pressure`) | The numeric signal the Nurse polls | +| `MaintenanceOrder` | Discovered fresh at the moment a work order needs to be raised; replaces the previous hardcoded `sap_url` | + +The Nurse also queries a GraphDB triple store directly. That URL **is** +configured (it identifies the policy source, not a peer system), so it +isn't part of the Arrowhead orchestration. + +## Configuration + +```json +{ + "systemname": "nurse", + "unit_assets": [ + { + "name": "HealthTracker", + "services": [ + { "definition": "SignalMonitoring", "subpath": "monitor", "registrationPeriod": 22 }, + { "definition": "EnrichmentNotification", "subpath": "enrichment", "registrationPeriod": 22, + "details": { "Forms": ["application/json"] } } + ], + "traits": [ + { + "graphdb_url": "http://:7200/repositories/", + "signals": [ + { + "serviceDefinition": "pressure", + "details": { "Unit": ["kPa"] }, + "samplingPeriod": 4, + "lowerThreshold": 10.0, + "upperThreshold": 25.0, + "spareParts": [] + } + ] + } + ] + } + ], + "protocolsNports": { "coap": 0, "http": 20181, "https": 0 }, + "coreSystems": [ /* serviceregistrar, orchestrator, ca, maitreD */ ] +} +``` + +### Trait reference + +| Field | Type | Description | +|-------|------|-------------| +| `graphdb_url` | string | Base URL of the GraphDB repository (no `/statements`). Probed at startup with an ASK; the Nurse refuses to run if it cannot reach it | +| `signals[]` | array | One entry per measurement the Nurse monitors | + +### Per-signal fields + +| Field | Type | Description | +|-------|------|-------------| +| `serviceDefinition` | string | The Arrowhead service definition the Nurse discovers and polls | +| `samplingPeriod` | integer | Seconds between successive polls | +| `lowerThreshold` | float | Lower edge of the allowed range | +| `upperThreshold` | float | Upper edge of the allowed range | +| `spareParts` | array | Optional SAP-side parts list passed through to the maintenance order | + +A reading outside `[lowerThreshold, upperThreshold]` for **five consecutive +polls on the same source** fires the alarm. The counter resets the moment a +value comes back in range. + +## GraphDB prerequisites + +The Nurse refuses to start without a reachable GraphDB and refuses to monitor +any sensor whose actuator it cannot resolve. The graph must contain, for +every sensor the Nurse will poll: + +```turtle + afo:hasName "" . + afo:diagnosesActuator . + arrowhead:functionalLocation "" . +``` + +The third triple typically already exists in upstream data (DEXPI / CFIHOS +imports). The first two come from the cloud's RDF — the +[kgrapher](../kgrapher/) emits `afo:hasName` automatically; the +`afo:diagnosesActuator` link is added once per sensor with a SPARQL +INSERT…WHERE that joins by literal name: + +```sparql +PREFIX afo: +PREFIX arrowhead: +INSERT { + GRAPH { + ?sensor afo:diagnosesActuator ?valveFL . + } +} +WHERE { + ?sensor afo:hasName "" . + ?valveFL arrowhead:functionalLocation "" . +} +``` + +The Nurse will issue the equivalent SELECT at runtime; if it returns no rows +the corresponding node is marked unresolvable in the log and skipped. + +## Order request shape + +When a signal fires, the Nurse builds and POSTs a `MaintenanceOrderEvent` to +the discovered Sapper URL. Pre-populated from the cached resolution: + +```json +{ + "equipmentId": "827PD2708", + "functionalLocation": "827-PV2708-200", + "plant": "1000", + "description": "Signal pressure out of range [10.00, 25.00]", + "priority": "3", + "maintenanceOrderType": "PM01", + "plannedStartTime": "", + "plannedEndTime": "", + "operations": [ { "text": "...", "workCenter": "MAINT-WC01", "duration": 4, "durationUnit": "H" } ] +} +``` + +After a successful response, the Nurse pushes a small set of sensor-side +context triples to GraphDB — `ex:bySensor`, `ex:targetFLTag`, `ex:reason` — +keyed by the same order IRI the Sapper writes the lifecycle to. + +## Building and running + +```bash +# Run from source (development) +go run . + +# Build a binary for the current machine +go build -o nurse_amac . + +# Cross-compile for a 64-bit Raspberry Pi +GOOS=linux GOARCH=arm64 go build -o nurse_rpi64 . + +# Deploy +scp nurse_rpi64 jan@:oslo/nurse/ +``` + +Run the binary from **inside its own directory** so it can find (or +auto-generate) `systemconfig.json`. + +## Startup order + +``` +Arrowhead core + GraphDB → Sapper → Nurse → signal providers (Emulator, ...) +``` + +The Nurse's GraphDB probe is fail-loud — if the triple store doesn't answer, +the Nurse won't start. The Sapper discovery is fail-tolerant — if the Sapper +isn't yet registered when an order needs to be raised, the call logs the +failure and the signal stays in `WorkRequested` until the next process +restart. Signal providers can join after the Nurse is already running; the +poll loop discovers them on the next tick. + +## Development with a local mbaigo clone + +Add both modules to the workspace `go.work` at the repository root: + +``` +use ./mbaigo +use ./systems/nurse +``` + +Or add a `replace` directive to `go.mod`: + +``` +require github.com/sdoque/mbaigo v0.x.x +replace github.com/sdoque/mbaigo => ../../mbaigo +``` diff --git a/nurse/nurse.go b/nurse/nurse.go index a589075..3004c6a 100644 --- a/nurse/nurse.go +++ b/nurse/nurse.go @@ -105,6 +105,13 @@ func serving(t *Traits, w http.ResponseWriter, r *http.Request, servicePath stri default: http.Error(w, "Method is not supported.", http.StatusMethodNotAllowed) } + case "enrichment": + switch r.Method { + case http.MethodPost: + t.enrichment(w, r) + default: + http.Error(w, "Method is not supported.", http.StatusMethodNotAllowed) + } default: http.Error(w, "Invalid service request [Do not modify the services subpath in the configuration file]", http.StatusBadRequest) } diff --git a/nurse/thing.go b/nurse/thing.go index eedc8c8..8f7b7fa 100644 --- a/nurse/thing.go +++ b/nurse/thing.go @@ -19,12 +19,10 @@ package main import ( "bytes" "context" - "crypto/tls" "encoding/json" "fmt" "io" "log" - "net" "net/http" "net/url" "path" @@ -55,10 +53,10 @@ type SignalT struct { // Traits holds the configurable parameters for the nurse unit asset. type Traits struct { - SAP_URL string `json:"sap_url"` - GraphDB_URL string `json:"graphdb_url"` - Signals []SignalT `json:"signals"` - pendingOrders map[string]string // orderID → signalName; not serialized + GraphDB_URL string `json:"graphdb_url"` + Signals []SignalT `json:"signals"` + pendingOrders map[string]string // orderID → signalName; not serialized + sapper *components.Cervice // discovers the Sapper's MaintenanceOrder service owner *components.System ua *components.UnitAsset } @@ -74,15 +72,22 @@ func initTemplate() *components.UnitAsset { RegPeriod: 22, Description: "monitors the value of the consumed service signal (GET)", } + enrichmentService := components.Service{ + Definition: "EnrichmentNotification", + SubPath: "enrichment", + Details: map[string][]string{"Forms": {"application/json"}}, + RegPeriod: 22, + Description: "receives planner enrichment + release notifications (POST)", + } return &components.UnitAsset{ Name: "HealthTracker", Details: map[string][]string{}, ServicesMap: components.Services{ - monitorService.SubPath: &monitorService, + monitorService.SubPath: &monitorService, + enrichmentService.SubPath: &enrichmentService, }, Traits: &Traits{ - SAP_URL: "http://192.168.1.108:20191/sapper/SAPSimulator/orders", GraphDB_URL: "http://13.79.36.131:7200/repositories/arrowhead-skoghall-v2", Signals: []SignalT{ { @@ -122,18 +127,10 @@ func newResource(configuredAsset usecases.ConfigurableAsset, sys *components.Sys t.Signals[i].UnresolvableNodes = make(map[string]bool) } - // Derive the health endpoint from the SAP URL (replace the API path with /health) - sapHealthURL := t.SAP_URL - if u, err := url.Parse(t.SAP_URL); err == nil { - u.Path = "/health" - sapHealthURL = u.String() - } - r := CheckServerUp(sapHealthURL, 2*time.Second) - if r.Up { - fmt.Printf("SAP server is up (status=%d, in %s)\n", r.StatusCode, r.Duration) - } else { - fmt.Printf("SAP server is down (%v, in %s)\n", r.Err, r.Duration) - } + // The Sapper is discovered via Arrowhead at order-creation time — no + // startup healthcheck. If the Sapper isn't yet registered when an order + // needs to be raised, requestMaintenanceOrder logs the failure and the + // signal stays in WorkRequested until the next process restart. // GraphDB is load-bearing: without it the nurse cannot resolve a sensor // to its associated asset and any maintenance order it raised would be @@ -158,6 +155,15 @@ func newResource(configuredAsset usecases.ConfigurableAsset, sys *components.Sys } cervices[cSignal.Definition] = &cSignal } + // Cervice used to discover the Sapper's MaintenanceOrder service at + // runtime, replacing the previous hardcoded sap_url. + t.sapper = &components.Cervice{ + Definition: "MaintenanceOrder", + Protos: sProtocols, + Nodes: make(map[string][]components.NodeInfo), + Mode: "set", + } + cervices[t.sapper.Definition] = t.sapper ua := &components.UnitAsset{ Name: configuredAsset.Name, @@ -296,6 +302,9 @@ func (t *Traits) sigMon(name string, period time.Duration) error { if orderID != "" { t.pendingOrders[orderID] = name log.Printf("Maintenance order %s created for signal %s/%s\n", orderID, name, node) + reason := fmt.Sprintf("Signal %s out of range [%.2f, %.2f]", + sig.Name, sig.LowerThreshold, sig.UpperThreshold) + go t.pushOrderContext(orderID, equipmentID, location, reason) } else { log.Printf("SAP order failed for signal %s/%s; monitoring paused until system restart\n", name, node) } @@ -335,6 +344,29 @@ func (t *Traits) state(w http.ResponseWriter) { w.Write([]byte(text)) } +// enrichment receives the planner's release-and-enrichment notification from +// the Sapper. For now the Nurse simply logs the payload to the terminal so +// the demo audience can see the planner's decision arriving downstream of the +// originating signal. Future work could attach it to the signal's state for +// the /monitor GET view. +func (t *Traits) enrichment(w http.ResponseWriter, r *http.Request) { + body, err := io.ReadAll(r.Body) + if err != nil { + http.Error(w, "failed to read request body", http.StatusBadRequest) + return + } + defer r.Body.Close() + + var pretty bytes.Buffer + if json.Indent(&pretty, body, "", " ") == nil { + log.Printf("← enrichment from %s\n%s\n", r.RemoteAddr, pretty.String()) + } else { + log.Printf("← enrichment from %s\n%s\n", r.RemoteAddr, string(body)) + } + + w.WriteHeader(http.StatusNoContent) +} + // update handles an incoming completion notification and restores the signal to operational. // It accepts both the nurse's MaintenanceDoneEvent format (orderId) and the SAP adaptor's // notifyDigitalTwin format (maintenanceOrder), so either party can call this endpoint. @@ -399,58 +431,38 @@ func (t *Traits) update(w http.ResponseWriter, r *http.Request) { w.Write([]byte(ttl)) } -// //-------------------------------------SAP server interaction functions - -type CheckResult struct { - Up bool - StatusCode int - Err error - Duration time.Duration -} - -func CheckServerUp(rawURL string, timeout time.Duration) CheckResult { - start := time.Now() - - parsed, err := url.Parse(rawURL) - if err != nil { - return CheckResult{Up: false, Err: fmt.Errorf("invalid url: %w", err), Duration: time.Since(start)} - } - - dialer := &net.Dialer{ - Timeout: timeout, - KeepAlive: 30 * time.Second, - } - - tr := &http.Transport{ - Proxy: http.ProxyFromEnvironment, - DialContext: func(ctx context.Context, network, addr string) (net.Conn, error) { - return dialer.DialContext(ctx, network, addr) - }, - TLSHandshakeTimeout: timeout, - ResponseHeaderTimeout: timeout, - IdleConnTimeout: 30 * time.Second, - TLSClientConfig: &tls.Config{ - MinVersion: tls.VersionTLS12, - }, - } - - client := &http.Client{ - Timeout: timeout, - Transport: tr, +// pushOrderContext POSTs a SPARQL INSERT linking the Sapper-issued Order IRI +// to the sensor-side facts only the Nurse knows: which sensor raised the +// order, which actuator FL tag it diagnoses, and the threshold-breach reason. +// The Order IRI shape matches the Sapper's, so the three sets of triples — +// Sapper CRTD, Nurse context, Sapper TECO — merge on the same subject in the +// graph. Graph publication is a side effect; failures are logged but do not +// block the maintenance loop. +func (t *Traits) pushOrderContext(orderID, sensorName, valveTag, reason string) { + if t.GraphDB_URL == "" { + return } + orderURI := "https://sinetiq.se/sap/MaintenanceOrder/" + orderID + sparql := fmt.Sprintf(`PREFIX ex: +INSERT DATA { + GRAPH { + <%s> ex:bySensor %q ; + ex:targetFLTag %q ; + ex:reason %q . + } +}`, orderURI, sensorName, valveTag, reason) - req, err := http.NewRequest(http.MethodGet, parsed.String(), nil) - if err != nil { - return CheckResult{Up: false, Err: fmt.Errorf("build request: %w", err), Duration: time.Since(start)} - } + log.Printf("→ GraphDB INSERT (context) order=%s\n%s\n", orderID, sparql) - resp, err := client.Do(req) + endpoint := strings.TrimRight(t.GraphDB_URL, "/") + "/statements" + resp, err := http.Post(endpoint, "application/sparql-update", strings.NewReader(sparql)) if err != nil { - return CheckResult{Up: false, Err: err, Duration: time.Since(start)} + log.Printf("pushOrderContext: POST failed: %v\n", err) + return } defer resp.Body.Close() - - return CheckResult{Up: true, StatusCode: resp.StatusCode, Duration: time.Since(start)} + msg, _ := io.ReadAll(resp.Body) + log.Printf("← GraphDB %s body=%s\n", resp.Status, string(msg)) } // resolveActuators looks up the actuator (e.g. valve) functional-location tag @@ -506,7 +518,13 @@ func (t *Traits) resolveActuators(sig *SignalT, nodes map[string][]components.No // diagnosesActuator triple, has more than one match (ambiguous), or the // endpoint is unreachable. func resolveActuatorTag(endpoint, sensorName string, timeout time.Duration) (string, error) { - query := fmt.Sprintf(`PREFIX afo: + // afo: is the producer's namespace — the same URI the kgrapher emits and + // the same URI under which the sensor's triples live in the remote graph. + // arrowhead: is the upstream (Skoghall) namespace under which the valve's + // functional-location object is published. The mixed prefixes are + // intentional: each side keeps its own vocabulary; alignment ontologies + // bridge them at reasoning time. + query := fmt.Sprintf(`PREFIX afo: PREFIX arrowhead: SELECT ?valveTag WHERE { ?sensor afo:hasName %q . @@ -570,7 +588,10 @@ func CheckGraphDBUp(endpoint string, timeout time.Duration) error { ctx, cancel := context.WithTimeout(context.Background(), timeout) defer cancel() - const probe = `SELECT (1 AS ?ok) WHERE {}` + // ASK with a trivial triple pattern is the most universally-accepted probe: + // no projection expressions (some GraphDB versions misparse `(N AS ?v)`), + // no aggregation, just a yes/no. + const probe = `ASK { ?s ?p ?o }` req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, strings.NewReader(probe)) if err != nil { return fmt.Errorf("build request: %w", err) @@ -637,12 +658,34 @@ func (t *Traits) requestMaintenanceOrder(sig *SignalT, equipmentID string, locat return "" } - log.Printf("→ SAP POST %s\n%s\n", t.SAP_URL, string(bodyBytes)) + // Discover the Sapper's MaintenanceOrder service via Arrowhead. The + // previous design hardcoded a sap_url; orchestration removes that + // topology dependency from the config. + if len(t.sapper.Nodes) == 0 { + if err := usecases.Search4Services(t.sapper, t.owner); err != nil { + log.Printf("requestMaintenanceOrder: discovery failed: %v\n", err) + return "" + } + } + var sapURL string + for _, nodes := range t.sapper.Nodes { + if len(nodes) > 0 { + sapURL = nodes[0].URL + break + } + } + if sapURL == "" { + log.Printf("requestMaintenanceOrder: no MaintenanceOrder provider found\n") + t.sapper.Nodes = make(map[string][]components.NodeInfo) // force re-discovery next time + return "" + } + + log.Printf("→ SAP POST %s\n%s\n", sapURL, string(bodyBytes)) ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) defer cancel() - req, err := http.NewRequestWithContext(ctx, http.MethodPost, t.SAP_URL, bytes.NewReader(bodyBytes)) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, sapURL, bytes.NewReader(bodyBytes)) if err != nil { log.Printf("requestMaintenanceOrder: new request: %v\n", err) return "" diff --git a/nurse/thing_test.go b/nurse/thing_test.go index 8f624f7..22202a8 100644 --- a/nurse/thing_test.go +++ b/nurse/thing_test.go @@ -43,8 +43,8 @@ func TestInitTemplate(t *testing.T) { if !ok { t.Fatal("Traits should be *Traits") } - if tr.SAP_URL == "" { - t.Error("SAP_URL default should not be empty") + if tr.GraphDB_URL == "" { + t.Error("GraphDB_URL default should not be empty") } if len(tr.Signals) == 0 { t.Error("template should include at least one signal") @@ -99,40 +99,6 @@ func TestAssetNameFromURL(t *testing.T) { } } -// ------------------------------------- CheckServerUp - -func TestCheckServerUp_Up(t *testing.T) { - srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - w.WriteHeader(http.StatusOK) - })) - defer srv.Close() - - r := CheckServerUp(srv.URL, 2*time.Second) - if !r.Up { - t.Errorf("expected Up=true, got false (%v)", r.Err) - } - if r.StatusCode != http.StatusOK { - t.Errorf("StatusCode = %d, want 200", r.StatusCode) - } -} - -func TestCheckServerUp_Down(t *testing.T) { - r := CheckServerUp("http://127.0.0.1:1/health", 500*time.Millisecond) - if r.Up { - t.Error("expected Up=false for unreachable server") - } - if r.Err == nil { - t.Error("expected a non-nil error") - } -} - -func TestCheckServerUp_InvalidURL(t *testing.T) { - r := CheckServerUp("://not a url", 1*time.Second) - if r.Up { - t.Error("expected Up=false for invalid URL") - } -} - // ------------------------------------- ConvertAckJSONToTurtle func TestConvertAckJSONToTurtle_Valid(t *testing.T) { @@ -224,7 +190,6 @@ func TestConvertMaintenanceDoneEventToTurtle_DefaultBaseIRI(t *testing.T) { func newTraitsWithPendingOrder(orderID, signalName string) *Traits { return &Traits{ - SAP_URL: "http://localhost", Signals: []SignalT{ { Name: signalName, @@ -322,7 +287,6 @@ func TestUpdate_InvalidJSON(t *testing.T) { func TestState(t *testing.T) { tr := &Traits{ - SAP_URL: "http://localhost", Signals: []SignalT{ {Name: "temperature", LowerThreshold: 0.0, UpperThreshold: 75.0, Operational: true, TOverCount: map[string]int{}, WorkRequested: map[string]bool{}}, }, diff --git a/sapper/README.md b/sapper/README.md index 3b90bd1..8b39ef1 100644 --- a/sapper/README.md +++ b/sapper/README.md @@ -2,60 +2,95 @@ ## Purpose -The *Sapper* system bridges the Arrowhead local cloud with **SAP Plant Maintenance** (SAP Hana PM module). The name comes from the French *Sapeur* — a combat engineer who keeps equipment operational under field conditions. It is an apt metaphor: the sapper detects that a component has failed, dispatches a repair order, and reports back when the equipment is back in service. - -In the Arrowhead architecture the sapper: - -1. **Accepts maintenance order requests** from any authorized consumer (e.g. the Nurse) over a registered `MaintenanceOrder` service. -2. **Manages the order lifecycle** — progressing the status from `CRTD` (created) → `REL` (released) → `TECO` (technically completed). -3. **Notifies the consumer** when the order reaches `TECO` by discovering the consumer's monitoring endpoint at runtime via Arrowhead orchestration — no hard-coded callback URL required. - -In simulation mode the sapper embeds the full SAP lifecycle internally; for production it can be adapted to forward orders to a real SAP system. +The *Sapper* bridges the Arrowhead local cloud with **SAP Plant Maintenance** +(SAP Hana PM module). The name is from the French *sapeur* — a combat engineer +who keeps equipment operational under field conditions. The Sapper's role is +the same: accept a maintenance request raised by some condition-monitoring +consumer, route it through a SAP-style lifecycle that includes a **human +planner step**, and report back when the work is technically completed so the +consumer can resume monitoring. + +In simulation mode the Sapper embeds the full lifecycle internally; for +production it can be adapted to forward orders to a real SAP system through +Alex Chiquito's +[SAP Maintenance Order Adaptor](https://github.com/AlexChiquito/SAP-Maintenance-order-adaptor). + +## Architecture + +Three responsibilities: + +1. **Accept** `MaintenanceOrder` POSTs from any authorised consumer (e.g. the + [Nurse](../nurse/)). New orders are created in **CRTD** (created) status — + they do **not** auto-progress. +2. **Show the work to a planner** via the `firefighting` web service. The + planner picks an order, adjusts the operations JSON, and submits. That + submission transitions the order to **REL** (released) and starts a + countdown to **TECO** (technically completed). +3. **Notify the consumer** at TECO — the Sapper discovers the consumer's + `SignalMonitoring` endpoint via Arrowhead orchestration and POSTs a + completion event. No callback URL is configured anywhere. + +If a GraphDB triple store is configured, each lifecycle transition pushes a +SPARQL update so the full audit trail +(`CRTD → enrichment → REL → TECO`) is queryable. ## How it fits with the Nurse -The **Nurse** system continuously monitors physical signals (e.g. pressure, temperature). When a signal repeatedly exceeds its threshold, the Nurse requests a maintenance order from the Sapper. Once the Sapper completes the lifecycle it discovers the Nurse's `SignalMonitoring` service via orchestration and POSTs a completion event, which causes the Nurse to restore that signal to operational status. +The [Nurse](../nurse/) monitors physical signals against a range. When a +signal stays out of range for five consecutive samples, the Nurse raises a +maintenance request on the Sapper. The order sits in CRTD until a human +planner opens the Sapper's firefighting page, enriches the operations payload, +and clicks Submit. The Sapper then notifies the Nurse with the enrichment +payload (so the Nurse log shows the planner's decision), runs its TECO +countdown, and finally posts a completion event back to the Nurse. ``` -Emulator ──► Nurse ──► Sapper ──► (SAP / simulator) - ◄────────────── - completion callback - (discovered via Arrowhead) +Emulator ──► Nurse ──► Sapper ─┬──► CRTD (waits) + ▲ │ + │ ▼ Planner opens /firefighting + │ REL → TECO countdown + │ │ + │ ▼ + └────── completion callback (TECO) + + enrichment notification (REL) ``` ## Sequence diagram ```mermaid sequenceDiagram - participant E as Emulator + autonumber participant N as Nurse participant OR as Orchestrator participant S as Sapper - - Note over E,N: Normal monitoring loop (every samplingPeriod seconds) - loop Signal polling - N->>E: GET /emulator/{signal}/access - E-->>N: SignalA_v1a {value, unit, timestamp} - end - - Note over N: value > threshold for 5 consecutive readings - N->>N: mark signal non-operational - N->>S: POST /sapper/SAPSimulator/orders
{equipmentId, plant, description, operations, ...} - S-->>N: 201 Created
{maintenanceOrder, maintenanceNotification, status:"CRTD", createdAt} - N->>N: store pendingOrders[orderID] = signalName - - Note over S: Order lifecycle (completionDelay seconds) - S->>S: status → REL (at ½ completionDelay) - S->>S: status → TECO (at completionDelay) - - Note over S: Discover callback endpoint at runtime - S->>OR: POST /orchestrator/orchestration/squest
{serviceDefinition:"SignalMonitoring"} - OR-->>S: ServicePoint_v1 {url: http://nurse-host/nurse/HealthTracker/monitor} - - S->>N: POST /nurse/HealthTracker/monitor
{orderId, status:"TECO", completedAt, actualWorkHours} - N->>N: look up pendingOrders[orderId] → signalName - N->>N: restore signal: Operational=true, WorkRequested=false, TOverCount=0 - N-->>S: 200 OK + participant P as Planner (browser) + participant G as GraphDB + + Note over N: out-of-range for 5 consecutive samples + N->>OR: discover MaintenanceOrder + OR-->>N: Sapper URL + N->>S: POST /sapper/SAPSimulator/maintenanceorders + S->>G: INSERT (CRTD status, equipment, description) + S-->>N: 201 Created {orderId, status:"CRTD"} + N->>G: INSERT (bySensor, targetFLTag, reason) + + Note over P,S: Order stays CRTD until a planner intervenes + P->>S: GET /sapper/SAPSimulator/firefighting + S-->>P: HTML page (CRTD list + textarea) + Note over P: pick order, edit operations JSON + P->>S: POST /sapper/SAPSimulator/firefighting
{orderId, enrichment JSON} + S->>G: INSERT (REL status, releasedAt) + S->>OR: discover EnrichmentNotification + OR-->>S: Nurse URL + S->>N: POST /nurse/HealthTracker/enrichment
{orderId, status:"REL", releasedAt, enrichment} + S-->>P: 303 See Other → GET firefighting (flash banner) + + Note over S: TECO countdown (completionDelay seconds) + S->>G: INSERT (TECO status, completedAt, actualWorkHours) + S->>OR: discover SignalMonitoring + OR-->>S: Nurse URL + S->>N: POST /nurse/HealthTracker/monitor
{orderId, status:"TECO", completedAt} + N->>N: restore signal to operational ``` ## Services @@ -64,16 +99,22 @@ sequenceDiagram | Service definition | Subpath | Methods | Description | |--------------------|---------|---------|-------------| -| `MaintenanceOrder` | `orders` | `POST` | Create a maintenance order; returns order ID and notification ID | -| `MaintenanceOrder` | `orders` | `GET ?id=` | Query the current status of a known order | +| `MaintenanceOrder` | `maintenanceorders` | `POST` | Create a new CRTD maintenance order | +| `MaintenanceOrder` | `maintenanceorders` | `GET ?id=` | Query the current state of a known order | +| `firefighting` | `firefighting` | `GET` | Planner UI: live-refreshed list of CRTD orders + editable enrichment textarea | +| `firefighting` | `firefighting` | `POST` | Submit handler: attaches enrichment, transitions to REL, redirects (303) to the GET view | + +### Consumed (via Arrowhead orchestration) -### Consumed +| Service definition | Used for | +|--------------------|----------| +| `SignalMonitoring` | TECO completion event — POSTed to the consumer that originally raised the order | +| `EnrichmentNotification` | REL release event — POSTed to the same consumer when the planner clicks Submit | -| Service definition | Purpose | -|--------------------|---------| -| `SignalMonitoring` | Callback endpoint discovered at runtime to report order completion | +Both consumed endpoints are discovered fresh at the moment of need; the Sapper +holds no configured URLs for them. -## Configuration (`systemconfig.json`) +## Configuration ```json { @@ -81,33 +122,34 @@ sequenceDiagram "unit_assets": [ { "name": "SAPSimulator", - "details": { - "Plant": ["1000"] - }, + "details": { "Plant": ["1000"] }, "services": [ { "definition": "MaintenanceOrder", - "subpath": "orders", + "subpath": "maintenanceorders", "details": { "Forms": ["application/json"] }, "registrationPeriod": 30 + }, + { + "definition": "firefighting", + "subpath": "firefighting", + "details": { "Forms": ["text/html"] }, + "registrationPeriod": 30 } ], "traits": [ { - "completionDelay": 30 + "completionDelay": 30, + "graphDbUrl": "http://:7200/repositories/" } ] } ], - "protocolsNports": { - "coap": 0, - "http": 20191, - "https": 0 - }, + "protocolsNports": { "coap": 0, "http": 20191, "https": 0 }, "coreSystems": [ - { "coreSystem": "serviceregistrar", "url": "http://192.168.1.109:20102/serviceregistrar/registry" }, - { "coreSystem": "orchestrator", "url": "http://192.168.1.109:20103/orchestrator/orchestration" }, - { "coreSystem": "ca", "url": "http://192.168.1.109:20100/ca/certification" }, + { "coreSystem": "serviceregistrar", "url": "http://:20102/serviceregistrar/registry" }, + { "coreSystem": "orchestrator", "url": "http://:20103/orchestrator/orchestration" }, + { "coreSystem": "ca", "url": "http://:20100/ca/certification" }, { "coreSystem": "maitreD", "url": "http://localhost:20101/maitreD/maitreD" } ] } @@ -117,32 +159,65 @@ sequenceDiagram | Field | Type | Default | Description | |-------|------|---------|-------------| -| `completionDelay` | integer (seconds) | `30` | Time from order creation to `TECO`. Set to `5` for fast demos, `3600` for realistic 1-hour jobs. | - -### Wiring the Nurse - -Point the Nurse's `sap_url` at the sapper's `orders` endpoint: - -```json -"sap_url": "http://:20191/sapper/SAPSimulator/orders" +| `completionDelay` | integer (seconds) | `30` | Time from REL (planner submit) to TECO. Set `5` for fast demos, `3600` for realistic 1-hour jobs | +| `graphDbUrl` | string | `""` | Base URL of a GraphDB repository (no `/statements` suffix). Empty disables all SPARQL pushes; lifecycle still works in-memory | + +## The firefighting UI + +`GET /sapper/SAPSimulator/firefighting` renders an HTML page with two halves: + +- **Top half** — a table of work orders currently in CRTD status, refreshed + every 3 seconds via background `fetch` so new orders appear (and released + ones disappear) without the planner reloading the page. +- **Bottom half** — a textarea for the operations JSON, empty on initial load. + When the planner selects an order, the textarea is populated with that + order's *suggested enrichment* (a per-order field on `Order`, currently + defaulted to a seals/gaskets repair template). The planner edits in place + and clicks **Submit**. + +Submission uses **Post/Redirect/Get** (HTTP 303) so browser reload after a +submit re-issues a GET, not a re-POST. A flash banner above the table +confirms success or reports a validation failure. + +## Order ID continuity (graph-primed counter) + +Order IDs follow the SAP convention `4XXXXXXXX` — zero-padded with a leading +`4`. The Sapper's in-memory counter normally starts at `0` and increments per +order, so a fresh process would produce `400000001`, `400000002`, etc. — and +would **collide** with historical IDs if the workorders graph already has +those subjects. + +To prevent this, the Sapper queries GraphDB on the **first** order creation +(once per process, behind `sync.Once`): + +```sparql +SELECT (MAX(?id) AS ?max) WHERE { + ?wo a step:WorkOrder ; + workorder:Id ?idNode . + FILTER(STRSTARTS(STR(?wo), "https://sinetiq.se/sap/")) + ?idNode identifier:Id ?id . +} ``` -The sapper discovers the Nurse's callback endpoint automatically via Arrowhead orchestration — no further configuration is needed on the Nurse side. +The returned literal (e.g. `"400000017"`) primes the counter so the next +order becomes `400000018`. The peek is failure-tolerant: if GraphDB is +unconfigured, unreachable, empty, or returns garbage, the counter stays at +`0` and the Sapper produces `400000001` (the previous behaviour). ## Order payload -### Request (POST `/orders`) +### Request (POST `/maintenanceorders`) ```json { - "equipmentId": "10000045", - "functionalLocation": "FL100-200-300", + "equipmentId": "827PD2708", + "functionalLocation": "827-PV2708-200", "plant": "1000", - "description": "Signal pressure exceeded threshold 75.00", + "description": "Signal pressure out of range [10.00, 25.00]", "priority": "3", "maintenanceOrderType": "PM01", - "plannedStartTime": "2026-03-27T14:00:00+01:00", - "plannedEndTime": "2026-03-27T22:00:00+01:00", + "plannedStartTime": "2026-05-30T06:00:00Z", + "plannedEndTime": "2026-05-30T14:00:00Z", "operations": [ { "text": "Inspect and service equipment for signal pressure", @@ -158,43 +233,74 @@ The sapper discovers the Nurse's callback endpoint automatically via Arrowhead o ```json { - "maintenanceOrder": "400000001", - "maintenanceNotification": "200000001", + "maintenanceOrder": "400000018", + "maintenanceNotification": "200000018", "status": "CRTD", "message": "Maintenance order created successfully", - "createdAt": "2026-03-27T13:58:21Z" + "createdAt": "2026-05-30T13:58:21Z" +} +``` + +### Enrichment notification (POST to consumer's `/enrichment`) + +Sent right after the planner clicks Submit, before the TECO countdown: + +```json +{ + "orderId": "400000018", + "status": "REL", + "releasedAt": "2026-05-30T14:00:01Z", + "enrichment": { "operations": [ /* the planner's submitted JSON */ ] } } ``` -### Completion callback (POST to consumer's monitor) +### Completion callback (POST to consumer's `/monitor`) + +Sent when TECO fires: ```json { - "orderId": "400000001", + "orderId": "400000018", "status": "TECO", - "completedAt": "2026-03-27T14:28:21Z", + "completedAt": "2026-05-30T14:00:31Z", "actualWorkHours": 0.0083, "notes": "Completed by SAP simulator" } ``` +## GraphDB triples (when `graphDbUrl` is set) + +All three lifecycle transitions land in the named graph +``, keyed by the same Order +IRI ``: + +- **CRTD** — full STEP/WorkOrder block: `step:WorkOrder`, `step:WorkRequest`, + `step:WorkRequestAssignment`, `dcterms:created`, `ex:status "CRTD"`. +- **REL** — `ex:status "REL"`, `ex:releasedAt`. +- **TECO** — `ex:status "TECO"`, `ex:completedAt`, `ex:actualWorkHours`. + +A consumer (the Nurse) adds sensor-side context to the same order subject — +`ex:bySensor`, `ex:targetFLTag`, `ex:reason` — so querying any order URI +returns the full lifecycle in a single result set. + ## Building and running ```bash -# Run in place (for development) +# Run from source (development) go run . -# Build for the current machine -go build -o sapper_local +# Build a binary for the current machine (with a host-specific suffix so it's gitignored) +go build -o sapper_amac . -# Cross-compile for Raspberry Pi 64-bit -GOOS=linux GOARCH=arm64 go build -o sapper_rpi64 +# Cross-compile for a 64-bit Raspberry Pi +GOOS=linux GOARCH=arm64 go build -o sapper_rpi64 . -# Copy to a Raspberry Pi -scp sapper_rpi64 user@192.168.1.10:mbaigo/sapper/ +# Deploy +scp sapper_rpi64 jan@:oslo/sapper/ ``` -Run the binary from **inside its own directory** so it can find (or create) `systemconfig.json`. +Run the binary from **inside its own directory** so it can find (or +auto-generate) `systemconfig.json`. ## Startup order @@ -202,7 +308,11 @@ Run the binary from **inside its own directory** so it can find (or create) `sys Arrowhead core systems → Sapper → Nurse (or any other consumer) ``` -The Sapper must be registered before the Nurse starts monitoring, so the Orchestrator can resolve the `MaintenanceOrder` service. The completion callback works in the opposite direction — the Nurse must be registered before any order reaches `TECO`. +The Sapper must be registered before any consumer attempts to raise a +maintenance order, so the Orchestrator can resolve `MaintenanceOrder`. The +two callback paths work in the opposite direction — the consumer must be +registered before the Sapper transitions to REL (for the enrichment +notification) or TECO (for the completion callback). ## Development with a local mbaigo clone diff --git a/sapper/models.go b/sapper/models.go index b0019ad..442223a 100644 --- a/sapper/models.go +++ b/sapper/models.go @@ -15,7 +15,10 @@ package main -import "time" +import ( + "encoding/json" + "time" +) // OrderRequest is the body sent by the nurse (or any consumer) to create a maintenance order. // Field names match the nurse's MaintenanceOrderEvent JSON tags so no translation is needed. @@ -62,12 +65,18 @@ type OrderResponse struct { } // Order is the in-memory representation of a maintenance order managed by the sapper. +// An order is created in CRTD status by the nurse and stays there until a planner +// enriches it via the firefighting UI; that submission flips it to REL and starts +// the TECO countdown. type Order struct { - ID string - Notification string - Status string // CRTD → REL → TECO - CreatedAt time.Time - Request OrderRequest + ID string + Notification string + Status string // CRTD → REL → TECO + CreatedAt time.Time + ReleasedAt time.Time // zero until enrichAndRelease is called + Request OrderRequest + SuggestedEnrichment string // template the firefighting UI prefills when this order is selected; defaults to the seals/gaskets example + Enrichment json.RawMessage // planner's actual submission; nil until released } // CompletionEvent is POSTed to the consumer's monitor endpoint when an order reaches TECO. diff --git a/sapper/sapper.go b/sapper/sapper.go index 155fb51..2482591 100644 --- a/sapper/sapper.go +++ b/sapper/sapper.go @@ -95,7 +95,7 @@ func main() { // serving dispatches incoming HTTP requests to the appropriate handler. func serving(t *Traits, w http.ResponseWriter, r *http.Request, servicePath string) { switch servicePath { - case "orders": + case "maintenanceorders": switch r.Method { case http.MethodPost: t.createOrderHandler(w, r) @@ -104,6 +104,8 @@ func serving(t *Traits, w http.ResponseWriter, r *http.Request, servicePath stri default: http.Error(w, "Method not supported", http.StatusMethodNotAllowed) } + case "firefighting": + t.firefightingHandler(w, r) default: http.Error(w, "Invalid service path [do not modify subpath in configuration file]", http.StatusBadRequest) } diff --git a/sapper/thing.go b/sapper/thing.go index 4307840..e3a9d4d 100644 --- a/sapper/thing.go +++ b/sapper/thing.go @@ -17,12 +17,17 @@ package main import ( "bytes" + "context" "crypto/md5" "encoding/json" "fmt" + "html/template" "io" "log" "net/http" + "net/url" + "sort" + "strconv" "strings" "sync" "sync/atomic" @@ -32,6 +37,184 @@ import ( "github.com/sdoque/mbaigo/usecases" ) +// firefightingHTML is the html/template body served at +// /sapper/SAPSimulator/firefighting. Inlined here (rather than as a separate +// .html file) to keep the system self-contained — same convention as the +// clerk system's orderPage. +const firefightingHTML = ` + + + + Sapper — Firefighting + + + +

Sapper · Firefighting

+

Work orders awaiting planner enrichment and release. Pick an order, edit + the operations JSON to match the field decision, then submit. The order + transitions CRTD → REL and the 30 s TECO countdown + begins.

+ + {{if .Flash}}
{{.Flash}}
{{end}} + +
+

Work orders (CRTD)

+ + + + + + + + + {{range .Orders}} + + + + + + + + + {{else}} + + {{end}} + +
OrderSensorFunctional locationDescriptionCreated
{{.ID}}{{.Request.EquipmentID}}{{.Request.FunctionalLocation}}{{.Request.Description}}{{.CreatedAt.Format "15:04:05"}}
No CRTD orders awaiting release.
+ +

Enrichment payload

+ + +

+
+ + + + +` + +// defaultEnrichmentJSON pre-fills the textarea with Alex's seals/gaskets +// example payload. Edit-in-place is the planner's main job. +const defaultEnrichmentJSON = `{ + "operations": [ + { + "operation": "0010", + "description": "Isolate valve and verify closure", + "plannedWorkQuantity": "1.0", + "workQuantityUnit": "H", + "workCenter": "VALVE-WC01", + "plant": "1000" + }, + { + "operation": "0020", + "description": "Disassemble valve body", + "plannedWorkQuantity": "1.5", + "workQuantityUnit": "H", + "workCenter": "VALVE-WC01", + "plant": "1000" + }, + { + "operation": "0030", + "description": "Replace gaskets and stem", + "plannedWorkQuantity": "1.0", + "workQuantityUnit": "H", + "workCenter": "VALVE-WC01", + "plant": "1000", + "components": [ + { + "material": "VALVE-GASKET-V50", + "description": "Valve gasket set", + "requiredQuantity": "2.0", + "unit": "EA", + "plant": "1000", + "storageLocation": "0001" + }, + { + "material": "VALVE-STEM-KIT", + "description": "Valve stem repair kit", + "requiredQuantity": "1.0", + "unit": "EA", + "plant": "1000", + "storageLocation": "0001" + } + ] + }, + { + "operation": "0040", + "description": "Reassemble and pressure test valve", + "plannedWorkQuantity": "1.5", + "workQuantityUnit": "H", + "workCenter": "VALVE-WC01", + "plant": "1000", + "components": [ + { + "material": "THREAD-SEALANT", + "description": "Thread sealant compound", + "requiredQuantity": "1.0", + "unit": "EA", + "plant": "1000", + "storageLocation": "0001" + } + ] + } + ] +} +` + +var firefightingTemplate = template.Must(template.New("firefighting").Parse(firefightingHTML)) + //-------------------------------------Define the unit asset // Traits holds the configurable parameters for the sapper unit asset. @@ -41,7 +224,9 @@ type Traits struct { orders map[string]*Order mu sync.Mutex seq atomic.Int64 // monotonic counter for order IDs + primeOnce sync.Once // guards a single graph-peek before the first order is allocated monitor *components.Cervice + enrichment *components.Cervice // discovers the nurse's enrichment endpoint owner *components.System ua *components.UnitAsset } @@ -52,17 +237,25 @@ type Traits struct { func initTemplate() *components.UnitAsset { ordersService := components.Service{ Definition: "MaintenanceOrder", - SubPath: "orders", + SubPath: "maintenanceorders", Details: map[string][]string{"Forms": {"application/json"}}, RegPeriod: 30, Description: "creates (POST) and queries (GET ?id=) maintenance orders", } + firefightingService := components.Service{ + Definition: "firefighting", + SubPath: "firefighting", + Details: map[string][]string{"Forms": {"text/html"}}, + RegPeriod: 30, + Description: "planner UI (GET) to enrich and release CRTD work orders (POST)", + } return &components.UnitAsset{ Name: "SAPSimulator", Details: map[string][]string{"Plant": {"1000"}}, ServicesMap: components.Services{ - ordersService.SubPath: &ordersService, + ordersService.SubPath: &ordersService, + firefightingService.SubPath: &firefightingService, }, Traits: &Traits{ CompletionDelay: 30, // 30 × time.Second = 30 s @@ -85,8 +278,9 @@ func newResource(configuredAsset usecases.ConfigurableAsset, sys *components.Sys } } - // Build a cervice so the sapper can discover who consumes MaintenanceOrders - // (i.e., the nurse's "SignalMonitoring" monitor endpoint) at runtime. + // Build cervices so the sapper can discover, at runtime, the nurse's + // callback endpoints: SignalMonitoring (TECO completion) and + // EnrichmentNotification (planner's release with enrichment payload). sProtocols := components.SProtocols(sys.Husk.ProtoPort) monitorCervice := &components.Cervice{ Definition: "SignalMonitoring", @@ -94,10 +288,18 @@ func newResource(configuredAsset usecases.ConfigurableAsset, sys *components.Sys Nodes: make(map[string][]components.NodeInfo), Mode: "get", } + enrichmentCervice := &components.Cervice{ + Definition: "EnrichmentNotification", + Protos: sProtocols, + Nodes: make(map[string][]components.NodeInfo), + Mode: "set", + } t.monitor = monitorCervice + t.enrichment = enrichmentCervice cervices := components.Cervices{ - "monitor": monitorCervice, + "monitor": monitorCervice, + "enrichment": enrichmentCervice, } ua := &components.UnitAsset{ @@ -120,55 +322,177 @@ func newResource(configuredAsset usecases.ConfigurableAsset, sys *components.Sys //-------------------------------------Order lifecycle -// nextID generates a zero-padded order number using a monotonic counter. +// nextOrderID generates a zero-padded order number using a monotonic counter. +// On the very first call, it queries GraphDB for the highest order number +// already in the workorders graph and jumps the counter past it, so a +// restarted Sapper does not produce IDs that collide with historical orders. func (t *Traits) nextOrderID() string { + t.primeOnce.Do(t.primeCounterFromGraph) n := t.seq.Add(1) return fmt.Sprintf("4%08d", n) } +// primeCounterFromGraph fast-forwards the order counter to one past the +// highest existing ID in the workorders graph. Silently no-ops if the graph +// is unconfigured, unreachable, empty, or returns garbage — in any of those +// cases the counter stays at 0 and the next order will be 400000001 (the +// previous, restart-resetting behaviour). +func (t *Traits) primeCounterFromGraph() { + if t.GraphDBURL == "" { + return + } + maxN, ok := t.peekMaxOrderID() + if !ok { + return + } + // CAS loop so a concurrent caller can't drag the counter backwards if + // they raced ahead between our load and store. + for { + cur := t.seq.Load() + if cur >= maxN { + return + } + if t.seq.CompareAndSwap(cur, maxN) { + log.Printf("sapper: counter primed from graph at %d (next order will be 4%08d)", maxN, maxN+1) + return + } + } +} + +// peekMaxOrderID queries GraphDB for the highest order number already +// recorded as a step:WorkOrder under our IRI namespace. Returns (n, true) on +// success, (0, true) when the graph is empty, (0, false) on any error. +func (t *Traits) peekMaxOrderID() (int64, bool) { + const query = `PREFIX step: +PREFIX identifier: +PREFIX workorder: +SELECT (MAX(?id) AS ?max) WHERE { + ?wo a step:WorkOrder ; + workorder:Id ?idNode . + FILTER(STRSTARTS(STR(?wo), "https://sinetiq.se/sap/")) + ?idNode identifier:Id ?id . +}` + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, t.GraphDBURL, strings.NewReader(query)) + if err != nil { + log.Printf("sapper: peekMaxOrderID: build request: %v", err) + return 0, false + } + req.Header.Set("Content-Type", "application/sparql-query") + req.Header.Set("Accept", "application/sparql-results+json") + + resp, err := http.DefaultClient.Do(req) + if err != nil { + log.Printf("sapper: peekMaxOrderID: POST failed: %v", err) + return 0, false + } + defer resp.Body.Close() + + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + body, _ := io.ReadAll(resp.Body) + log.Printf("sapper: peekMaxOrderID: HTTP %s: %s", resp.Status, strings.TrimSpace(string(body))) + return 0, false + } + + var result struct { + Results struct { + Bindings []map[string]struct { + Value string `json:"value"` + } `json:"bindings"` + } `json:"results"` + } + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + log.Printf("sapper: peekMaxOrderID: decode: %v", err) + return 0, false + } + if len(result.Results.Bindings) == 0 { + return 0, true // aggregation returned no rows: empty graph + } + binding, ok := result.Results.Bindings[0]["max"] + if !ok || binding.Value == "" { + return 0, true // no MAX binding: graph has no WorkOrder triples under our namespace + } + if len(binding.Value) < 1 || binding.Value[0] != '4' { + log.Printf("sapper: peekMaxOrderID: unexpected ID format %q", binding.Value) + return 0, false + } + n, err := strconv.ParseInt(binding.Value[1:], 10, 64) + if err != nil { + log.Printf("sapper: peekMaxOrderID: parse %q: %v", binding.Value, err) + return 0, false + } + return n, true +} + func (t *Traits) nextNotifID() string { n := t.seq.Load() return fmt.Sprintf("2%08d", n) } -// createOrder stores a new order and starts its lifecycle goroutine. +// createOrder stores a new order in CRTD status. It does NOT start the +// lifecycle countdown — the order waits in CRTD until a planner enriches and +// releases it through the firefighting UI. func (t *Traits) createOrder(req OrderRequest) *Order { o := &Order{ - ID: t.nextOrderID(), - Notification: t.nextNotifID(), - Status: "CRTD", - CreatedAt: time.Now(), - Request: req, + ID: t.nextOrderID(), + Notification: t.nextNotifID(), + Status: "CRTD", + CreatedAt: time.Now(), + Request: req, + SuggestedEnrichment: defaultEnrichmentJSON, } t.mu.Lock() t.orders[o.ID] = o t.mu.Unlock() - - go t.runLifecycle(o) return o } -// runLifecycle advances order status CRTD → REL → TECO and then notifies the consumer. -func (t *Traits) runLifecycle(o *Order) { +// enrichAndRelease attaches the planner's enrichment to the order, transitions +// it to REL, records the release in GraphDB, notifies the nurse with the +// enrichment payload, and starts the TECO countdown. Returns false if the +// order does not exist or is not in CRTD status. +func (t *Traits) enrichAndRelease(orderID string, enrichment json.RawMessage) (*Order, error) { + t.mu.Lock() + o, ok := t.orders[orderID] + if !ok { + t.mu.Unlock() + return nil, fmt.Errorf("order %s not found", orderID) + } + if o.Status != "CRTD" { + t.mu.Unlock() + return nil, fmt.Errorf("order %s is %s, only CRTD orders can be released", orderID, o.Status) + } + o.Status = "REL" + o.ReleasedAt = time.Now() + o.Enrichment = enrichment + t.mu.Unlock() + log.Printf("order %s → REL (planner authorised)\n", o.ID) + + go t.insertReleaseToGraphDB(o) + go t.notifyEnrichment(o) + go t.runTECOCountdown(o) + return o, nil +} + +// runTECOCountdown sleeps the configured completion delay and then transitions +// the order to TECO. Called from enrichAndRelease — the countdown is measured +// from the planner's release, not from order creation. +func (t *Traits) runTECOCountdown(o *Order) { delay := t.CompletionDelay * time.Second if delay <= 0 { delay = 30 * time.Second // safe default } + time.Sleep(delay) - // Advance to REL at half-time. - time.Sleep(delay / 2) - t.mu.Lock() - o.Status = "REL" - t.mu.Unlock() - log.Printf("order %s → REL\n", o.ID) - - // Advance to TECO at full delay. - time.Sleep(delay / 2) t.mu.Lock() o.Status = "TECO" t.mu.Unlock() log.Printf("order %s → TECO\n", o.ID) + go t.insertCompletionToGraphDB(o) t.notifyConsumer(o) } @@ -179,6 +503,12 @@ var discoverMonitor = func(c *components.Cervice, sys *components.System) error return usecases.Search4Services(c, sys) } +// discoverEnrichment mirrors discoverMonitor for the nurse's enrichment endpoint. +var discoverEnrichment = func(c *components.Cervice, sys *components.System) error { + c.Nodes = make(map[string][]components.NodeInfo) + return usecases.Search4Services(c, sys) +} + // notifyConsumer discovers the SignalMonitoring endpoint via Arrowhead and POSTs // the completion event. func (t *Traits) notifyConsumer(o *Order) { @@ -302,19 +632,133 @@ where ) } +// buildRELInsertSPARQL records the planner's release of an order: the new REL +// status triple and the release timestamp. The enrichment JSON itself is not +// stored as a literal in the graph — it ships to the nurse over the +// EnrichmentNotification path and is captured in the nurse's log; embedding +// it as a SPARQL string literal couples graph state to textarea encoding and +// produces malformed queries on CRLF input. +func (t *Traits) buildRELInsertSPARQL(o *Order) string { + orderURI := "https://sinetiq.se/sap/MaintenanceOrder/" + o.ID + releasedAt := o.ReleasedAt.UTC().Format(time.RFC3339Nano) + return fmt.Sprintf(`PREFIX ex: +PREFIX xsd: +INSERT DATA { + GRAPH { + <%s> ex:status "REL" ; + ex:releasedAt "%s"^^xsd:dateTime . + } +}`, orderURI, releasedAt) +} + +// insertReleaseToGraphDB pushes the REL transition to GraphDB. +func (t *Traits) insertReleaseToGraphDB(o *Order) { + t.postSPARQL("REL", o.ID, t.buildRELInsertSPARQL(o)) +} + +// notifyEnrichment posts the planner's enrichment payload to the nurse's +// enrichment endpoint discovered via Arrowhead. Independent from +// notifyConsumer (which fires on TECO). +func (t *Traits) notifyEnrichment(o *Order) { + if t.enrichment == nil { + log.Printf("notifyEnrichment: no enrichment cervice for order %s\n", o.ID) + return + } + if err := discoverEnrichment(t.enrichment, t.owner); err != nil { + log.Printf("notifyEnrichment: discovery failed for order %s: %v\n", o.ID, err) + return + } + var callbackURL string + for _, nodes := range t.enrichment.Nodes { + if len(nodes) > 0 { + callbackURL = nodes[0].URL + break + } + } + if callbackURL == "" { + log.Printf("notifyEnrichment: no EnrichmentNotification endpoint found for order %s\n", o.ID) + return + } + + // The body bundles the order ID with the planner's payload so the nurse + // can correlate without parsing the URL. + envelope, err := json.Marshal(struct { + OrderID string `json:"orderId"` + Status string `json:"status"` + ReleasedAt time.Time `json:"releasedAt"` + Enrichment json.RawMessage `json:"enrichment"` + }{o.ID, o.Status, o.ReleasedAt, o.Enrichment}) + if err != nil { + log.Printf("notifyEnrichment: marshal error: %v\n", err) + return + } + + log.Printf("→ enrichment %s order=%s\n", callbackURL, o.ID) + resp, err := http.Post(callbackURL, "application/json", bytes.NewReader(envelope)) + if err != nil { + log.Printf("notifyEnrichment: POST failed: %v\n", err) + return + } + defer resp.Body.Close() + msg, _ := io.ReadAll(resp.Body) + log.Printf("← enrichment %s body=%s\n", resp.Status, string(msg)) +} + +// buildTECOInsertSPARQL constructs the SPARQL UPDATE statement that records an +// order's transition to TECO. It appends new triples (status, completion time, +// actual work hours) to the existing Order IRI created by buildSPARQL — the +// graph retains both the CRTD and TECO status triples as an audit trail. +func (t *Traits) buildTECOInsertSPARQL(o *Order) string { + orderURI := "https://sinetiq.se/sap/MaintenanceOrder/" + o.ID + completedAt := time.Now().UTC().Format(time.RFC3339Nano) + actualHours := float64(t.CompletionDelay) / 3600.0 // CompletionDelay is in seconds + + return fmt.Sprintf(`PREFIX ex: +PREFIX xsd: +INSERT DATA { + GRAPH { + <%s> ex:status "TECO" ; + ex:completedAt "%s"^^xsd:dateTime ; + ex:actualWorkHours "%g"^^xsd:decimal . + } +}`, orderURI, completedAt, actualHours) +} + +// updateEndpoint returns the SPARQL update URL by appending /statements to the +// repository base URL the operator configured. Empty in, empty out. +func (t *Traits) updateEndpoint() string { + if t.GraphDBURL == "" { + return "" + } + return strings.TrimRight(t.GraphDBURL, "/") + "/statements" +} + // insertToGraphDB prints the SPARQL UPDATE to the terminal and, when GraphDBURL // is configured, POSTs it to GraphDB. func (t *Traits) insertToGraphDB(o *Order) { - sparql := t.buildSPARQL(o) - log.Printf("→ GraphDB INSERT order=%s\n%s\n", o.ID, sparql) + t.postSPARQL("CRTD", o.ID, t.buildSPARQL(o)) +} - if t.GraphDBURL == "" { +// insertCompletionToGraphDB records the order's TECO transition in the graph. +// Called from runLifecycle once the in-memory status flips to TECO. +func (t *Traits) insertCompletionToGraphDB(o *Order) { + t.postSPARQL("TECO", o.ID, t.buildTECOInsertSPARQL(o)) +} + +// postSPARQL is the common path for both lifecycle pushes. Graph publication is +// a side effect of the SAP lifecycle, not load-bearing for it, so failures here +// are logged and dropped — the maintenance loop continues either way. +func (t *Traits) postSPARQL(stage, orderID, sparql string) { + log.Printf("→ GraphDB INSERT (%s) order=%s\n%s\n", stage, orderID, sparql) + + endpoint := t.updateEndpoint() + if endpoint == "" { return } - resp, err := http.Post(t.GraphDBURL, "application/sparql-update", strings.NewReader(sparql)) + resp, err := http.Post(endpoint, "application/sparql-update", strings.NewReader(sparql)) if err != nil { - log.Printf("insertToGraphDB: POST failed: %v\n", err) + log.Printf("postSPARQL (%s): POST failed: %v\n", stage, err) return } defer resp.Body.Close() @@ -324,6 +768,74 @@ func (t *Traits) insertToGraphDB(o *Order) { //-------------------------------------HTTP handlers +// firefightingHandler routes the firefighting UI: GET renders the planner +// page (work-order list + enrichment textarea), POST processes a Submit then +// redirects (303 See Other) back to GET so browser reload doesn't re-POST. +func (t *Traits) firefightingHandler(w http.ResponseWriter, r *http.Request) { + switch r.Method { + case http.MethodGet: + q := r.URL.Query() + t.renderFirefighting(w, q.Get("msg"), q.Get("kind")) + case http.MethodPost: + t.submitFirefighting(w, r) + default: + http.Error(w, "Method not supported", http.StatusMethodNotAllowed) + } +} + +// renderFirefighting writes the firefighting HTML, optionally with a flash +// message ("ok" or "err") after a Submit. +func (t *Traits) renderFirefighting(w http.ResponseWriter, flash, kind string) { + t.mu.Lock() + crtd := make([]*Order, 0, len(t.orders)) + for _, o := range t.orders { + if o.Status == "CRTD" { + crtd = append(crtd, o) + } + } + t.mu.Unlock() + sort.Slice(crtd, func(i, j int) bool { return crtd[i].ID < crtd[j].ID }) + + data := struct { + Orders []*Order + Flash string + FlashKind string + }{crtd, flash, kind} + + w.Header().Set("Content-Type", "text/html; charset=utf-8") + if err := firefightingTemplate.Execute(w, data); err != nil { + log.Printf("firefighting: template error: %v\n", err) + } +} + +// submitFirefighting handles the planner's Submit: validate, release the +// order, then 303-redirect to GET with a flash message in the query string. +// Post/Redirect/Get prevents browser reload from re-submitting the form. +func (t *Traits) submitFirefighting(w http.ResponseWriter, r *http.Request) { + msg, kind := t.processSubmit(r) + q := url.Values{"msg": {msg}, "kind": {kind}} + http.Redirect(w, r, "firefighting?"+q.Encode(), http.StatusSeeOther) +} + +func (t *Traits) processSubmit(r *http.Request) (string, string) { + if err := r.ParseForm(); err != nil { + return "Could not parse form: " + err.Error(), "err" + } + orderID := strings.TrimSpace(r.FormValue("orderId")) + enrichment := strings.TrimSpace(r.FormValue("enrichment")) + if orderID == "" { + return "Pick a work order before submitting.", "err" + } + if !json.Valid([]byte(enrichment)) { + return "Enrichment is not valid JSON.", "err" + } + o, err := t.enrichAndRelease(orderID, json.RawMessage(enrichment)) + if err != nil { + return err.Error(), "err" + } + return fmt.Sprintf("Order %s released. TECO in %d s.", o.ID, t.CompletionDelay), "ok" +} + // createOrderHandler handles POST /orders — creates a new maintenance order. func (t *Traits) createOrderHandler(w http.ResponseWriter, r *http.Request) { var req OrderRequest diff --git a/sapper/thing_test.go b/sapper/thing_test.go index e4ef379..f0559e2 100644 --- a/sapper/thing_test.go +++ b/sapper/thing_test.go @@ -31,7 +31,7 @@ import ( // ── helpers ─────────────────────────────────────────────────────────────────── // newTestTraits builds a minimal Traits suitable for unit tests. -// CompletionDelay is set to 0 so runLifecycle is never used in handler tests. +// CompletionDelay is set to 0 so the TECO countdown is never used in handler tests. func newTestTraits() *Traits { return &Traits{ CompletionDelay: 0, @@ -40,6 +40,10 @@ func newTestTraits() *Traits { Definition: "SignalMonitoring", Nodes: make(map[string][]components.NodeInfo), }, + enrichment: &components.Cervice{ + Definition: "EnrichmentNotification", + Nodes: make(map[string][]components.NodeInfo), + }, } } @@ -52,6 +56,14 @@ func withDiscoverMonitor(t *testing.T, fn func(*components.Cervice, *components. t.Cleanup(func() { discoverMonitor = orig }) } +// withDiscoverEnrichment is the equivalent shim for the nurse's enrichment endpoint. +func withDiscoverEnrichment(t *testing.T, fn func(*components.Cervice, *components.System) error) { + t.Helper() + orig := discoverEnrichment + discoverEnrichment = fn + t.Cleanup(func() { discoverEnrichment = orig }) +} + // validOrderBody returns a JSON-encoded minimal valid OrderRequest. func validOrderBody(t *testing.T) *bytes.Buffer { t.Helper() @@ -72,8 +84,8 @@ func TestInitTemplate(t *testing.T) { if ua.GetName() != "SAPSimulator" { t.Errorf("name = %q, want %q", ua.GetName(), "SAPSimulator") } - if _, ok := ua.GetServices()["orders"]; !ok { - t.Error("expected 'orders' service in ServicesMap") + if _, ok := ua.GetServices()["maintenanceorders"]; !ok { + t.Error("expected 'maintenanceorders' service in ServicesMap") } tr, ok := ua.GetTraits().(*Traits) if !ok { @@ -125,7 +137,7 @@ func TestNewResource(t *testing.T) { cfgAsset := usecases.ConfigurableAsset{ Name: "SAPSimulator", Traits: []json.RawMessage{traitJSON}, - Services: []components.Service{{Definition: "MaintenanceOrder", SubPath: "orders"}}, + Services: []components.Service{{Definition: "MaintenanceOrder", SubPath: "maintenanceorders"}}, } ua, cleanup := newResource(cfgAsset, &sys) @@ -137,8 +149,8 @@ func TestNewResource(t *testing.T) { if ua.ServingFunc == nil { t.Error("ServingFunc must be set") } - if _, ok := ua.GetServices()["orders"]; !ok { - t.Error("expected 'orders' service") + if _, ok := ua.GetServices()["maintenanceorders"]; !ok { + t.Error("expected 'maintenanceorders' service") } if ua.CervicesMap["monitor"] == nil { t.Error("expected 'monitor' cervice") @@ -206,10 +218,13 @@ func TestCreateOrder(t *testing.T) { } } -// ── runLifecycle ────────────────────────────────────────────────────────────── +// ── enrichAndRelease + runTECOCountdown ────────────────────────────────────── +// +// The lifecycle is now split: the order sits in CRTD until a planner calls +// enrichAndRelease, which flips status to REL and starts the TECO countdown. +// This test drives that flow against fake monitor + enrichment endpoints. -func TestRunLifecycle(t *testing.T) { - // Set up a fake monitor endpoint that records whether a callback arrived. +func TestEnrichAndRelease(t *testing.T) { callbackReceived := make(chan struct{}, 1) monitor := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { callbackReceived <- struct{}{} @@ -217,14 +232,24 @@ func TestRunLifecycle(t *testing.T) { })) defer monitor.Close() - // Swap discovery so it returns the fake monitor URL instead of contacting Arrowhead. + enrichmentReceived := make(chan struct{}, 1) + enrichmentSrv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + enrichmentReceived <- struct{}{} + w.WriteHeader(http.StatusNoContent) + })) + defer enrichmentSrv.Close() + withDiscoverMonitor(t, func(c *components.Cervice, _ *components.System) error { c.Nodes = map[string][]components.NodeInfo{"testNode": {{URL: monitor.URL}}} return nil }) + withDiscoverEnrichment(t, func(c *components.Cervice, _ *components.System) error { + c.Nodes = map[string][]components.NodeInfo{"testNode": {{URL: enrichmentSrv.URL}}} + return nil + }) tr := newTestTraits() - tr.CompletionDelay = 1 // 1 × time.Second = 1 s total lifecycle + tr.CompletionDelay = 1 // 1 s TECO countdown after release o := &Order{ ID: "400000001", @@ -237,30 +262,50 @@ func TestRunLifecycle(t *testing.T) { tr.orders[o.ID] = o tr.mu.Unlock() - go tr.runLifecycle(o) + released, err := tr.enrichAndRelease(o.ID, json.RawMessage(`{"operations":[]}`)) + if err != nil { + t.Fatalf("enrichAndRelease: unexpected error: %v", err) + } + if released.Status != "REL" { + t.Errorf("status after release = %q, want REL", released.Status) + } + if released.ReleasedAt.IsZero() { + t.Error("ReleasedAt must be set after enrichAndRelease") + } - // After ~0.6 s the order should be REL. - time.Sleep(600 * time.Millisecond) - tr.mu.Lock() - statusMid := o.Status - tr.mu.Unlock() - if statusMid != "REL" { - t.Errorf("status after half delay = %q, want REL", statusMid) + // Enrichment should arrive at the nurse essentially immediately. + select { + case <-enrichmentReceived: + case <-time.After(500 * time.Millisecond): + t.Error("expected enrichment to be received by nurse server") } - // After ~1.1 s the order should be TECO and the callback sent. - time.Sleep(600 * time.Millisecond) + // After ~1.2 s the TECO countdown should have fired and the monitor callback delivered. + time.Sleep(1200 * time.Millisecond) tr.mu.Lock() statusFinal := o.Status tr.mu.Unlock() if statusFinal != "TECO" { - t.Errorf("final status = %q, want TECO", statusFinal) + t.Errorf("status after TECO countdown = %q, want TECO", statusFinal) } select { case <-callbackReceived: - // good default: - t.Error("expected callback to be received by monitor server") + t.Error("expected TECO callback to be received by monitor server") + } +} + +// TestEnrichAndRelease_RejectsNonCRTD verifies the release path refuses to act +// on orders that aren't waiting. +func TestEnrichAndRelease_RejectsNonCRTD(t *testing.T) { + tr := newTestTraits() + o := &Order{ID: "400000002", Status: "TECO", CreatedAt: time.Now()} + tr.mu.Lock() + tr.orders[o.ID] = o + tr.mu.Unlock() + + if _, err := tr.enrichAndRelease(o.ID, json.RawMessage(`{}`)); err == nil { + t.Error("expected error releasing a non-CRTD order, got nil") } } @@ -405,28 +450,28 @@ func TestQueryOrderHandler(t *testing.T) { func TestServing(t *testing.T) { tr := newTestTraits() - t.Run("POST to orders creates order", func(t *testing.T) { - req := httptest.NewRequest(http.MethodPost, "/sapper/SAPSimulator/orders", validOrderBody(t)) + t.Run("POST to maintenanceorders creates order", func(t *testing.T) { + req := httptest.NewRequest(http.MethodPost, "/sapper/SAPSimulator/maintenanceorders", validOrderBody(t)) w := httptest.NewRecorder() - serving(tr, w, req, "orders") + serving(tr, w, req, "maintenanceorders") if w.Code != http.StatusCreated { t.Errorf("status = %d, want 201", w.Code) } }) - t.Run("GET to orders without id returns 400", func(t *testing.T) { - req := httptest.NewRequest(http.MethodGet, "/sapper/SAPSimulator/orders", nil) + t.Run("GET to maintenanceorders without id returns 400", func(t *testing.T) { + req := httptest.NewRequest(http.MethodGet, "/sapper/SAPSimulator/maintenanceorders", nil) w := httptest.NewRecorder() - serving(tr, w, req, "orders") + serving(tr, w, req, "maintenanceorders") if w.Code != http.StatusBadRequest { t.Errorf("status = %d, want 400", w.Code) } }) t.Run("unsupported method returns 405", func(t *testing.T) { - req := httptest.NewRequest(http.MethodDelete, "/sapper/SAPSimulator/orders", nil) + req := httptest.NewRequest(http.MethodDelete, "/sapper/SAPSimulator/maintenanceorders", nil) w := httptest.NewRecorder() - serving(tr, w, req, "orders") + serving(tr, w, req, "maintenanceorders") if w.Code != http.StatusMethodNotAllowed { t.Errorf("status = %d, want 405", w.Code) }