diff --git a/CHANGELOG.md b/CHANGELOG.md index 124c57d..6609b26 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,13 @@ This project uses Semantic Versioning as **interpreted for MemNet**: package `a. ## [Unreleased] +### Fixed +- **Serve request isolation (MN-REQ-06.13)** — concurrent `memnet serve` commands no longer share a process-global stdout/stderr swap. Each request captures its own streams on contextvars. A caller cannot receive another session's records, and a failed mutate cannot come back as another call's success. No serve-wide lock. TCP listen backlog is 128 so a burst is queued rather than refused. +- **Mutate line split** — leftover pipe stdin and GQL statement breaks split on LF only (optional trailing CR stripped), the same rule as snapshot load. VT, FF, FS, GS, RS, NEL, LS, and PS stay inside the value. +- **Snapshot extras and a SCHEMA without `id`** — undeclared properties widen the snapshot SCHEMA even when the CREATE label's case differs from the tag. A map that omits `id` still saves: the snapshot SCHEMA appends `id` and a minted nickname. The live map is unchanged until load. An unloadable SCHEMA key fails closed. +- **Labelled MATCH** — `MATCH (n:Label {…}) SET` case-folds the label to the SCHEMA tag. A different known label that matches nothing is `not_found`. +- **Housekeep endpoints (MN-REQ-04.12)** — orphan, dangling, and stale counts resolve edge ends by hidden element id and by nickname. A GQL graph is not reported as all orphans. `prune orphans` and `prune stale` with `--apply` recompute and refuse `prune_referenced` when a node is still an endpoint. No automatic `prune --apply`. + ### Changed - **Invent only — ClusterRoute vs SliceHandCarry (#191 / #47 cousin)** — `MemNetTwoMoves` outside `MemNetSystem` (`MN-REQ-06.9` + `MN-REQ-06.10` / `MN-VER-06-S08`). ClusterRoute = where the session lives (`MemNetLanMcpFront`; one owner; `pin_map` / `find` SHALL NOT span backends). SliceHandCarry = explicit copy into another session (`export_pin_map` or `session_save` → LAN file copy → dest import/`session_load`; `import_slice` same-serve only). Not a live hop. `import_slice(from_url)` not shipped. tip≠face. `inventOnly=true`; `implemented=false`; no engine code; no SemVer bump. Wire: [`docs/operations/cluster-route-vs-slice-hand-carry.md`](docs/operations/cluster-route-vs-slice-hand-carry.md). - **Invent only — LAN MCP front over several serves (#191)** — `MemNetLanMcpFront` outside `MemNetSystem` (`MN-REQ-06.9` / `MN-VER-06-S07`). One MCP catalogue, N LAN `memnet serve` backends; `SessionOwnerRegistry` is owner (explicit pin allowed; silent hash is not sole routing). One owner per session; `pin_map` / `find` SHALL NOT span backends. Cousin of #47 (peer sid handoff), not the same invent. tip≠face. `inventOnly=true`; `implemented=false`; no engine code; no SemVer bump. Wire: [`docs/operations/memnet-lan-mcp-front.md`](docs/operations/memnet-lan-mcp-front.md). diff --git a/docs/cap-contract.md b/docs/cap-contract.md index f082d59..1a29a55 100644 --- a/docs/cap-contract.md +++ b/docs/cap-contract.md @@ -159,7 +159,15 @@ Raised at map load (`memnet/tag_map.py`). Session open fails; nothing is stored. Snapshot emit escapes every Python `str.splitlines()` separator (LF, CR, VT, FF, FS/GS/RS, NEL, LS, PS) plus `|` and `\`. Load splits records on LF only (not `str.splitlines()`). -**Undeclared properties.** GQL mutate may store keys that are absent from the live tag SCHEMA (GraphElement extras; Path-B locators such as `qname`). Those keys stay in RAM. Snapshot save persists them by widening the **snapshot** SCHEMA (live session SCHEMA is unchanged). After load, the restored map includes the extra columns. Extras on fixed tags `EDG` / `LAW`, or a widened SCHEMA over `max_fields`, refuse `snapshot_unsaveable` and write no file. +Mutate stdin, leftover pipe batches, and GQL statement breaks use that same LF-only split (one optional trailing CR is stripped). A raw VT, FF, FS, GS, RS, NEL, LS, or PS inside a value stays inside the value. It does not start another statement. An embedded CR is not a break. A raw LF still separates statements. This matches snapshot load (#202): the separator is kept literally, not refused. + +**Undeclared properties.** GQL mutate may store keys that are absent from the live tag SCHEMA (GraphElement extras; Path-B locators such as `qname`). Those keys stay in RAM. Snapshot save persists them by widening the **snapshot** SCHEMA (live session SCHEMA is unchanged). The widen keys extras by the canonical SCHEMA tag, so a CREATE label whose spelling differs in case from the SCHEMA tag still persists. After load, the restored map includes the extra columns. Extras on fixed tags `EDG` / `LAW`, or a widened SCHEMA over `max_fields`, refuse `snapshot_unsaveable` and write no file. + +A live SCHEMA that omits `id` is accepted at open (`id` need not be present or first). Save still succeeds: the snapshot SCHEMA appends an `id` column and the row gets a minted nickname. The live map is unchanged until that snapshot is loaded. A property key that does not survive the SCHEMA whitespace split (for example `foo bar`) cannot reload, so save refuses `snapshot_unsaveable` and writes no file. A 32-field map with no `id` becomes 33 columns and refuses `fields|33/32`. + +**Labelled MATCH and edge create.** `MATCH (n:Label {…}) SET` is intended to match. The label is case-folded to the SCHEMA tag (`cst`, `Cst`, and `CST` are the same tag). A label that is a different known tag and matches nothing is `@ERR: not_found`, not a silent success. `MATCH (a) MATCH (b) CREATE (a)-[:rel]->(b)` is `@ERR: parse_error` (`unsupported MATCH continuation`). The comma form `MATCH (a), (b) CREATE (a)-[:rel]->(b)` with `--allow-new-relation` creates the edge. The MATCH-MATCH form is not implemented. + +**Housekeep orphans (MN-REQ-04.12).** GQL edges store endpoints as hidden element ids (`_elN`). Leftover pipe may store the nickname. Orphan, dangling, and stale counts resolve both. A node is an orphan only when no edge names it. `housekeep stats` on a fully linked GQL graph reports `orphans` 0. `prune orphans` and `prune stale` with `--apply` recompute and refuse `@ERR: prune_referenced` (nothing deleted) when a candidate node is still an endpoint. The `stale_orphans` warning suggests `prune orphans --apply` only for that true-orphan count. No sweep, expire job, or MCP tool runs `prune --apply`. The only caller is the CLI when a person passes `--apply`. `housekeep_stats` on MCP is read-only. Snapshot save verifies every emitted row can parse back to the same values. If any row cannot, save refuses `@ERR: snapshot_unsaveable|{tag} nick={nick} …` and writes no file. Expire-save in that case (or any other save failure, such as an unwritable disk or directory) emits `@WRN: expire_snapshot_failed|{code}` on every sweep or access that retries expiry, writes no file, and **keeps the session in RAM**. It still counts against `MEMNET_MAX_SESSIONS`. Access after TTL is `@ERR: session_expired|overdue`. An explicit `session save` that succeeds, or an explicit `session close`, ends that hold. When save-on-expire is off, TTL still drops RAM. Snapshots written by 0.19.18 (pipe and backslash escapes only) still load. @@ -292,6 +300,8 @@ Source: `memnet/catalog_snap.py` `snap_model` / `_precheck_plan`; `memnet/serve. | Default TTL | `60` minutes (`MEMNET_SESSION_TTL_MINUTES`). Open `--ttl` / MCP `ttl` | | Legal range | `1..1440` else `@ERR: bad_ttl\|ttl must be 1..1440` | | Live access | Sliding TTL: each `get_session` extends `expires_at` by the original minutes | +| Unit | Minutes, not seconds. `--ttl 60` is 60 minutes. 130 s idle is inside that window | +| Sweep | No background timer. `purge_expired` runs on access (`get_session`, open, list, count, save, close). A successful access slides `expires_at` forward by another `ttl_minutes`. After real expiry, access refuses `session_expired` (`snap_missing`, or `overdue` when expire-save kept RAM) | | Expire, save off (default) | Session dropped from memory. First access of the still-registered expired id: `@ERR: session_expired\|snap_missing` (exit 2). After purge already ran: `@ERR: session_not_found\|unknown session` | | Expire, `MEMNET_SAVE_ON_EXPIRE` truthy + `MEMNET_EXPIRE_SNAPSHOT_DIR` set | Snapshot `{dir}/{sid}.snap` (filename only; do not log it). Next use: `@ERR: session_expired\|snap_available`. Restore: `session_load` with that id | | Save-on-expire on, dir unset | `@WRN: save_on_expire_no_dir\|dir unset`, then drop; `snap_missing` | diff --git a/docs/operations/one-session-per-document.md b/docs/operations/one-session-per-document.md index 908ba51..4c2cfd2 100644 --- a/docs/operations/one-session-per-document.md +++ b/docs/operations/one-session-per-document.md @@ -16,7 +16,7 @@ Settings this gate uses: | Knob | Value | |------|--------| -| TTL | 60 minutes (`MEMNET_SESSION_TTL_MINUTES` or `session open --ttl 60`) | +| TTL | 60 **minutes** (`MEMNET_SESSION_TTL_MINUTES` or `session open --ttl 60`). Not 60 seconds. 130 s idle is still inside the window | | Save on expire | on (`MEMNET_SAVE_ON_EXPIRE=1`) plus `MEMNET_EXPIRE_SNAPSHOT_DIR` | | Concurrent sessions | 1024 (`MEMNET_MAX_SESSIONS`) | | Document size | about 1 800 part nodes plus a few opaque `USR` text nodes | @@ -39,6 +39,10 @@ Length-prefixed UTF-8 JSON on TCP `127.0.0.1` (default port 18765): There is no MCP `errors` array and no `session_id` field on this envelope. Session id appears only as `@SESSION:` on stdout. Hard refuse is stderr `@ERR: {code}|{message}` (exit 1 or 2). Mutate also prints `ok=N fail=M` on stderr. +Concurrent commands on one serve do not share that envelope (MN-REQ-06.13). Each request captures its own stdout, stderr, and exit code. A caller does not receive another session's records, and a failed mutate does not come back as `ok=1 fail=0`. Capture is per request. There is no serve-wide lock. + +`housekeep stats` counts orphans by resolving each edge end as a hidden element id or a nickname. A linked GQL graph is not a set of orphans. `prune orphans --apply` / `prune stale --apply` refuse `prune_referenced` rather than delete a node an edge still names. Nothing in the serve sweep or MCP calls `prune --apply` on its own. + Client helper: `memnet.serve.send_command(args, stdin=…, host=…, port=…)`. Wait default is 30 s (`SERVE_CLIENT_TIMEOUT_S`); a 1 800-node mutate batch may need a longer `timeout=` from the gate. ## Snapshot diff --git a/docs/operations/product-gateway-contract.md b/docs/operations/product-gateway-contract.md index c8f9f98..9ba4792 100644 --- a/docs/operations/product-gateway-contract.md +++ b/docs/operations/product-gateway-contract.md @@ -1,74 +1,22 @@ # Product gateway contract -**Status:** implemented for the behaviours in [This cut](#this-cut). **MN-REQ-06.12** / **MN-VER-06-S10**. Model: `MemNetProductGateway` nested on `MemNetLanMcpFront` (`sysml-models/models/deploy.sysml`). Code: `parts/memnet-mcp/software/memnet_mcp/product_gateway.py`. No SemVer bump. +**Status:** implemented and deployed. **MN-REQ-06.12** / **MN-VER-06-S10**. Model: `MemNetProductGateway` nested on `MemNetLanMcpFront` (`sysml-models/models/deploy.sysml`). Code: `parts/memnet-mcp/software/memnet_mcp/product_gateway.py`. The gateway is `memnet-mcp --transport gateway`. It is not a separate droplet shim. The #191 catalogue (MCP tool union, SnapshotHandCarry) stays invent-only on the parent part. -Online products reach a `memnet serve` only through this process. Endleaf now. Atelier and SysMLEdge later, each with their own credential and backend list. +Online products reach a `memnet serve` only through this process. Sample pin below is **0.19.20**. Backend ids are arbitrary strings. The live Endleaf backend id is `pi-endleaf`. -## Today +stdio and streamable-http do not read the registry. With no `MEMNET_GATEWAY_CONFIG`, one memnet-mcp process stays a single in-process or streamable-http server. -Checked on tags **v0.19.6** (the droplet `memnet-mcp-http`) and **v0.19.18** (the Pi serve). This checkout matches tag v0.19.18 for the MCP and serve files cited below. +## Credentials -### (1) One MCP, one serve +The gateway process reads `MEMNET_GATEWAY_CONFIG` (a JSON file). It stores each product credential as a SHA-256 hex digest and compares `Authorization: Bearer ` to that digest. -Neither tag can point one MCP process at several serves. +The product side names its bearer `MEMNET_GATEWAY_TOKEN`, for example in `/etc//memnet-gateway.env`. The gateway process does not read `MEMNET_GATEWAY_TOKEN`. The caller puts that value in the `Authorization` header. -| Tag | How the backend is chosen | -|-----|---------------------------| -| v0.19.6 | `MEMNET_MCP_TRANSPORT` is `inprocess` (default) or `tcp` / `serve`. TCP calls `probe()` and `send_command()` with no host argument, so both use `MEMNET_SERVE_HOST` (default `127.0.0.1`) and `MEMNET_SERVE_PORT` (default `18765`). | -| v0.19.18 | Same. The only client change is that `session expire-status` is left off the session flag. | +`MEMNET_MCP_HTTP_TOKEN` is the streamable-http shared secret only. It is unrelated to this gateway. The gateway does not accept an `mn_tip_` key. -Citations: - -- v0.19.6 `parts/memnet-mcp/software/memnet_mcp/client.py` lines 105–139 (`_transport`, `run_memnet`) -- v0.19.6 `parts/common/memnet/memnet/config.py` lines 90–95 (`serve_host`, `serve_port`) -- v0.19.18 `parts/memnet-mcp/software/memnet_mcp/client.py` lines 106–140 -- v0.19.18 `parts/common/memnet/memnet/config.py` lines 96–101 - -There is no backend id, no product table, and no second host. Setting `MEMNET_SERVE_HOST=100.118.79.40` and `MEMNET_SERVE_PORT=18795` would aim the single TCP client at the Endleaf Pi. It would not route a second product elsewhere. - -### (2) What auth exists in front - -| Check | Where | What is checked | -|-------|--------|-----------------| -| Optional shared bearer | streamable-http only | `MEMNET_MCP_HTTP_TOKEN`. Exact `Authorization: Bearer `. Empty or unset means auth off. | -| stdio | default `memnet-mcp` | No bearer, no tip key, no product credential. | -| `caller` | MCP tool argument, both tags | Optional `--caller` on pin_map / mutate / find when that session's CapsPolicy ACL is enabled. Not a gateway. Off unless the session turned ACL on. | -| Tip `mn_tip_…` | v0.19.18 portal only. Absent at v0.19.6. | SHA-256 of the key, invite not revoked. nginx `auth_request` in front of one public `/mcp`. Not scoped to a product's backends. Not a version pin. Not inside `memnet-mcp`. | - -Citations: - -- v0.19.6 and v0.19.18 `parts/memnet-mcp/software/memnet_mcp/http_transport.py` lines 70–76 and 230–264 (`mcp_http_token`, `SharedBearerASGI`). The file is identical across the two tags. -- v0.19.6 `parts/memnet-mcp/software/memnet_mcp/server.py` lines 233–234 (`--caller` only when the tool argument is set). The same argument exists on v0.19.18. -- v0.19.6 has no `ops/tip_access_portal/`. -- v0.19.18 `ops/tip_access_portal/tip_access_portal/store.py` lines 22–23 (`hash_secret`), 158–161 (`mn_tip_` + SHA-256 at insert), 206–210 (`key_is_active`) -- v0.19.18 `ops/tip_access_portal/tip_access_portal/gate.py` lines 8–12 - -**Reuse decision.** The portal already hashes and revokes a bearer. It does not map a key to a backend id, a house, or a pinned serve version. This gateway copies that hash-and-revoke pattern into its own config. It does not call the portal store and it does not accept an `mn_tip_` key as a product credential. `MEMNET_MCP_HTTP_TOKEN` stays the streamable-http shared secret and is not a product credential either. - -### (3) Can a 0.19.6 MCP talk to a 0.19.18 serve? - -Yes, for the argv that v0.19.6 MCP actually sends. The length-prefixed JSON envelope is the same shape. - -| | v0.19.6 | v0.19.18 | -|--|---------|----------| -| Request | `{"args": [...]}` plus optional `"stdin"` (`serve.py` `send_command` lines 206–217; `_handle_request` lines 76–108) | Same, plus optional `admin_token` and `admin_usage` only when the client sets them (`serve.py` lines 81–126 and 263–282). v0.19.6 never sets them. | -| Frame | 4-byte big-endian length + JSON | Same | -| Reply | `{exit_code, stdout, stderr}` | Same. `parse.py` is identical across the tags, so session-id and `@ERR` extraction do not change. | -| Version line | `@VER: memnet\|` (`cli.py` line 215) | Same shape (`cli.py` line 225). The text is `0.19.6` or `0.19.18`. | - -v0.19.18 CLI adds optional flags `--max-edges`, `--product`, and `--token`. It removes none. New commands the old MCP does not send: `admin usage-report`, `session expire-status`. - -A v0.19.6 MCP `session open`, `pin_map`, `mutate`, `session load --file`, and `session save` argv is still accepted. The old MCP cannot select a backend other than the one `MEMNET_SERVE_HOST` / `MEMNET_SERVE_PORT`. - -The other direction is not the same. A v0.19.18 MCP `snap_model` / `ingest_*` call always sends `--max-edges`. A v0.19.6 serve CLI has no such flag and will refuse that argv. That does not block a 0.19.6 MCP talking to a 0.19.18 serve. - -## This cut - -`memnet-mcp --transport gateway` reads `MEMNET_GATEWAY_CONFIG` (a JSON file). stdio and streamable-http do not read it. With no registry, one memnet-mcp process behaves as it does today. - -### Endpoint +## Endpoint `POST ` (default `/gateway`) with `Content-Type: application/json` and: @@ -96,7 +44,7 @@ With neither `house` nor `backend`, a new session walks `products..backends` A later call for a known session goes to the owning backend only. A `house` or `backend` that names a different backend is `gateway_forbidden_session`. The session is not moved. -### Reply +## Reply ```json {"exit_code": 0, "stdout": "…", "stderr": "…"} @@ -104,61 +52,62 @@ A later call for a known session goes to the owning backend only. A `house` or ` Engine `exit_code`, stdout, and stderr are the backend's bytes, including `@ERR`, `@WRN`, and `@STAT`, with one exception: `session list` drops `@SESSION` rows this product does not own and rewrites `@STAT: sessions|n/max` so `n` is the filtered count. `max` stays the backend's figure (the first successful backend's max). Other control lines are copied. If any allowed backend is down or the pin does not match, stderr also carries a `gateway_` line and `exit_code` is 2, so the list does not look complete. -HTTP status is a hint. The body is always the envelope. +HTTP status is a hint. The body is always the envelope. Backend engine errors stay HTTP 200 with the engine `exit_code`. -| HTTP | When | -|------|------| -| 200 | Backend reply, including engine `exit_code` other than 0 | -| 401 | `gateway_auth` | -| 403 | `gateway_forbidden_session` or `gateway_forbidden` | -| 409 | `gateway_backend_version_mismatch` | -| 413 | `gateway_body_too_large` | -| 502 | `gateway_backend_unreachable` | +| HTTP | Code | When | +|------|------|------| +| 200 | engine lines | Backend reply, including engine `exit_code` other than 0 | +| 200 | list, some backends down | `session list` when at least one backend answered. `exit_code` is 2 if any backend was a gateway error | +| 400 | `gateway_bad_request` | Body is not a JSON object; `args` is not a list of strings; `stdin` is not a string; `session` field and `--session` disagree; `house` or `backend` is not a string; `session` is required and missing | +| 401 | `gateway_auth` | Missing, unknown, or revoked bearer; namespace does not match the credential; admin bearer missing or wrong | +| 403 | `gateway_forbidden_session` | Unknown session, another product's session, or a pin that would move the owner. The detail does not say which | +| 403 | `gateway_forbidden` | `house` or `backend` is outside the product, or argv is `admin` or `serve` | +| 404 | `gateway_bad_request` | Method or path is not `POST `. The detail is `POST ` plus the configured path | +| 404 | `gateway_unconfigured` | `GET /admin/counts` when `admin.sha256` is unset | +| 409 | `gateway_backend_version_mismatch` | `memnet version` is not the product pin | +| 413 | `gateway_body_too_large` | Body above `body_max_bytes` | +| 502 | `gateway_backend_unreachable` | TCP probe failed, the version command did not answer, the client wait expired, or `session list` and no backend answered | -### Gateway errors +No HTTP response (the process exits 2 instead): -```text -@ERR: gateway_| -``` +| Exit | Code | When | +|------|------|------| +| 2 | `gateway_unconfigured` | `MEMNET_GATEWAY_CONFIG` is missing | +| 2 | `gateway_bad_request` | The config file does not load, or the bind address is refused | + +`exit_code` on HTTP refusals is 2. They are not engine lines. An engine `@ERR: not_found|…` stays `not_found`. -| Code | Meaning | -|------|---------| -| `gateway_auth` | Missing, unknown, or revoked bearer, or namespace does not match the credential | -| `gateway_forbidden_session` | Unknown session, another product's session, or a pin that would move the owner. The detail does not say which. | -| `gateway_forbidden` | `house` or `backend` is outside the product, or argv is `admin` or `serve` | -| `gateway_backend_unreachable` | TCP probe failed, the version command did not answer, or the client wait expired | -| `gateway_backend_version_mismatch` | `memnet version` is not the product pin | -| `gateway_body_too_large` | Body above `body_max_bytes` | -| `gateway_bad_request` | JSON, `args`, or disagreeing session fields | -| `gateway_unconfigured` | Admin counts requested and no admin hash is set; also process exit when `MEMNET_GATEWAY_CONFIG` is missing | +## Concurrency -`exit_code` on these refusals is 2. They are not engine lines. An engine `@ERR: not_found|…` stays `not_found`. +The gateway is a threaded HTTP server. Each POST is one `send_command` to the owning `memnet serve`. The serve is a threaded TCP server. MN-REQ-06.13: each command's stdout, stderr, and `exit_code` are captured on that request's context, not by swapping process-global `sys.stdout`. One session's records do not appear in another session's envelope. A failed mutate does not come back as another call's `ok=1 fail=0`. -### Limits +There is no serve-wide lock. Requests on different sessions still overlap. The cost is a context-local lookup on each write, not a queue. On this VM, 12 successful TCP mutates at once took 0.494 s against 0.525 s one after another (ratio 0.94). The work is CPU-bound under the GIL, so the overlap is small. A lock would queue every command and could not beat that serial time. The TCP listen backlog is 128 so a burst is queued instead of refused at `accept`. -The gateway does not enforce `MEMNET_MAX_ROWS`, `MEMNET_MAX_BATCH_LINES`, or `MEMNET_MAX_VALUE_BYTES`. The Endleaf Pi runs 0.19.18 with `MEMNET_MAX_ROWS=10000`. Engine defaults elsewhere are `MEMNET_MAX_ROWS` 5000, batch 1000 lines, value 4096 decoded bytes (`parts/common/memnet/memnet/config.py`). A product may send `--max-rows 10000` and a mutate stdin up to the body ceiling. The backend still applies its own caps. Those refusals come back verbatim. +## Limits + +The gateway does not enforce `MEMNET_MAX_ROWS`, `MEMNET_MAX_BATCH_LINES`, or `MEMNET_MAX_VALUE_BYTES`. Engine defaults are `MEMNET_MAX_ROWS` 5000, batch 1000 lines, value 4096 decoded bytes (`parts/common/memnet/memnet/config.py`). A product may send `--max-rows` and a mutate stdin up to the body ceiling. The backend still applies its own caps. Those refusals come back verbatim. The gateway body ceiling is `body_max_bytes`, default **4194304** (4 MiB), the same figure as `MEMNET_SERVE_MAX_FRAME_BYTES`. State it in the config. Leave the gateway process's `MEMNET_SERVE_MAX_FRAME_BYTES` at least that large so the TCP client does not refuse a body the contract already accepted. -### Version pin +## Version pin -Each product has `pinned_version` (for Endleaf, `0.19.18`). Before a forward, the gateway calls `version` on that backend and requires `@VER: memnet|`. A mismatch is `gateway_backend_version_mismatch` and the argv is not sent. Upgrading the Pi serve without changing the pin refuses the product. Changing the pin is how the owner is notified: edit the config and restart the gateway. +Each product has `pinned_version` (sample below, `0.19.20`). Before a forward, the gateway calls `version` on that backend and requires `@VER: memnet|`. A mismatch is `gateway_backend_version_mismatch` and the argv is not sent. Upgrading a serve without changing the pin refuses the product. Changing the pin is how the owner is notified: edit the config and restart the gateway. `version_cache_s` defaults to 15. Set `0` to check every call. A serve upgraded inside a non-zero window can still be reached until the cache expires. -### Admin counts +## Admin counts `GET /admin/counts` with `Authorization: Bearer `. -The admin secret is stored as `admin.sha256` only. The JSON lists each product's `requests`, `gateway_refusals`, and `backend_errors`. It has no graph text and no session id. Counts are in memory and reset when the process stops. +The admin secret is stored as `admin.sha256` only. The JSON lists each product's `requests`, `gateway_refusals`, and `backend_errors`. It has no graph text and no session id. Counts are in memory and reset when the process stops. If `admin.sha256` is unset, the route is `gateway_unconfigured` (HTTP 404). -### Listener +## Listener -Default bind `127.0.0.1`. Tailnet addresses in `100.64.0.0/10` (for example `100.118.79.40`) and Tailscale IPv6 `fd7a:115c:a1e0::/48` are allowed. `0.0.0.0` and other public addresses are refused unless `MEMNET_GATEWAY_ALLOW_PUBLIC=1`. Endleaf on the droplet should not set that. Bind the droplet loopback, or the droplet's own tailnet address. Do not publish the port. +Default bind `127.0.0.1`. Tailnet addresses in `100.64.0.0/10` and Tailscale IPv6 `fd7a:115c:a1e0::/48` are allowed. `0.0.0.0` and other public addresses are refused unless `MEMNET_GATEWAY_ALLOW_PUBLIC=1`. Bind the loopback, or the host's own tailnet address. Do not publish the port. ## Examples -Credential and session ids below are placeholders. +Credential and session ids below are placeholders. The backend id `pi-endleaf` is the live Endleaf id; any other string is valid if the config uses it consistently. ### session open @@ -170,7 +119,7 @@ Content-Type: application/json { "args": ["session", "open", "--map-file", "/var/memnet/schema.sysml.example.txt"], "namespace": "endleaf", - "house": "syson" + "backend": "pi-endleaf" } ``` @@ -182,7 +131,7 @@ Content-Type: application/json } ``` -The gateway records `mn_example` → backend `endleaf-rpi5-syson`. The map path is on the Pi, because `args` are not rewritten. +The gateway records `mn_example` → backend `pi-endleaf`. The map path is on that serve, because `args` are not rewritten. ### MATCH … SET @@ -195,33 +144,31 @@ The gateway records `mn_example` → backend `endleaf-rpi5-syson`. The map path } ``` -Reply `exit_code`, stdout, and stderr are the Pi's. A missing element comes back as the engine line `@ERR: not_found|…`, not a `gateway_` line. +Reply `exit_code`, stdout, and stderr are the serve's. A missing element comes back as the engine line `@ERR: not_found|…`, not a `gateway_` line. ### Edge delete -On 0.19.18 the form that reaches an edge delete is: - ```json { "args": ["mutate", "--stdin", "--session", "mn_example"], - "stdin": "MATCH (n WHERE true)-[r {id: 'E_gw'}]->() DELETE r\n", + "stdin": "MATCH ()-[r {id: 'E_gw'}]->() DELETE r\n", "session": "mn_example" } ``` -`MATCH ()-[r {id: 'E_gw'}]-() DELETE r` is still forwarded byte for byte. On this engine it lowers as a node delete and the Pi replies `@ERR: not_found|DELETE matched no element`. The gateway does not rewrite it. See [Open decisions](#open-decisions). +On 0.19.20 this empty-paren form deletes the edge (see closed decision 11). The gateway still forwards the bytes unchanged. It does not rewrite the statement. `MATCH (n WHERE true)-[r {id: 'E_gw'}]->() DELETE r` remains a valid spelling as well. ## Endleaf configuration -Pi serve (already running, not changed by this repo): `rpi5-syson`, tailnet `100.118.79.40:18795`, MemNet 0.19.18, `MEMNET_MAX_ROWS=10000`, no token and no ACL, accepting loopback and the droplet. - -On the droplet, hash the product credential and an admin credential. Do not commit the plaintext. +Hash the product credential and an admin credential. Do not commit the plaintext. ```bash python -c "import hashlib,sys; print(hashlib.sha256(sys.argv[1].encode()).hexdigest())" 'REPLACE_WITH_ENDLEAF_CREDENTIAL' ``` -`/etc/memnet/gateway.json` (mode 0600, not in git): +The product host exports the plaintext bearer as `MEMNET_GATEWAY_TOKEN` (for example from `/etc/endleaf/memnet-gateway.env`). That file is not read by `memnet-mcp`. `MEMNET_MCP_HTTP_TOKEN` is a different variable and does not authenticate this route. + +`/etc/memnet/gateway.json` (mode 0600, not in git). Backend ids are chosen by the operator. This sample uses the live id: ```json { @@ -234,16 +181,16 @@ python -c "import hashlib,sys; print(hashlib.sha256(sys.argv[1].encode()).hexdig "backend_timeout_s": 30, "admin": {"sha256": ""}, "backends": { - "endleaf-rpi5-syson": { + "pi-endleaf": { "host": "100.118.79.40", "port": 18795 } }, "products": { "endleaf": { - "backends": ["endleaf-rpi5-syson"], - "houses": {"syson": "endleaf-rpi5-syson"}, - "pinned_version": "0.19.18", + "backends": ["pi-endleaf"], + "houses": {"syson": "pi-endleaf"}, + "pinned_version": "0.19.20", "credentials": [ {"id": "endleaf-1", "sha256": "", "revoked": false} ] @@ -259,9 +206,11 @@ memnet-mcp --transport gateway Rotate by appending a new `{id, sha256, revoked: false}` and setting the old object's `revoked` to true. Restart is required for a config edit. The owner file remembers session → backend across restarts. It holds session ids. The admin counts endpoint does not. -To add a standby later, append another backend id to `backends` and to `endleaf.backends`, in failover order. Do not put Atelier's backend on Endleaf's list. +To add a standby later, append another backend id to `backends` and to `endleaf.backends`, in failover order. Do not put another product's backend on Endleaf's list. + +## Closed decisions -The droplet's existing 0.19.6 `memnet-mcp-http` does not speak this contract. This cut is not deployed. +11. **Empty-paren edge delete.** On 0.19.18, `MATCH ()-[r {id}]->() DELETE r` did not delete an edge. On 0.19.20 it does (honoured; `tests/test_doc_gate_readiness.py`). The gateway still does not rewrite the statement. Closed for this pin. ## Open decisions @@ -271,9 +220,8 @@ The droplet's existing 0.19.6 `memnet-mcp-http` does not speak this contract. Th 4. **List `max`.** The filtered stat keeps the backend's `max`. Whether to hide `max` is open. 5. **Unknown and foreign sessions** share `gateway_forbidden_session` so the caller cannot tell them apart. 6. **Failover** is only for a new session with no `house` and no `backend`. There is no automatic return to a recovered preferred backend, and no SnapshotHandCarry (same sid, new owner). That relocate stays the #191 invent. -7. **Public bind** exists only behind `MEMNET_GATEWAY_ALLOW_PUBLIC=1`. Endleaf should leave it unset. +7. **Public bind** exists only behind `MEMNET_GATEWAY_ALLOW_PUBLIC=1`. Leave it unset on a public host. 8. **Bind host** must be an IP in loopback or tailnet, or the name `localhost`. MagicDNS names are not resolved. 9. **Tip keys.** Federating `mn_tip_` into this credential store is open. This cut does not. -10. **Version cache.** Default 15 seconds. Endleaf's sample config sets `0`. -11. **Empty-paren edge delete** (`MATCH ()-[r {id}]->() DELETE r`) does not delete an edge on 0.19.18. Fixing that engine lower is open. The gateway will not rewrite it. -12. **`session list` on a shared backend** shows only sids this gateway recorded for that product. A session opened on the Pi by a path other than this gateway is invisible here and is refused as unknown. Whether the gateway should adopt a sid the product already holds on its pinned backend is open. +10. **Version cache.** Default 15 seconds. The sample config sets `0`. +12. **`session list` on a shared backend** shows only sids this gateway recorded for that product. A session opened on the serve by a path other than this gateway is invisible here and is refused as unknown. Whether the gateway should adopt a sid the product already holds on its pinned backend is open. diff --git a/parts/common/memnet/memnet/cli.py b/parts/common/memnet/memnet/cli.py index dd1db35..995955b 100644 --- a/parts/common/memnet/memnet/cli.py +++ b/parts/common/memnet/memnet/cli.py @@ -1005,22 +1005,33 @@ def relations_list( emit_stdout(f"@REL: {rel}") +def _decode_ingest_text(raw: str | bytes) -> str: + if isinstance(raw, bytes): + try: + return raw.decode("utf-8") + except UnicodeDecodeError as exc: + raise MemNetError("encoding", "input must be UTF-8") from exc + return raw + + def _read_ingest_input( line: str | None, file: Path | None, stdin: bool, caps: Caps, ) -> list[str]: + from memnet.wire import split_lf_lines + raw_lines: list[str] = [] if line: - raw_lines = [line] + raw_lines = split_lf_lines(line) if line else [] elif file: - raw_lines = file.read_bytes().splitlines() + raw_lines = split_lf_lines(_decode_ingest_text(file.read_bytes())) elif stdin: if hasattr(sys.stdin, "buffer"): - raw_lines = sys.stdin.buffer.read().splitlines() + raw_lines = split_lf_lines(_decode_ingest_text(sys.stdin.buffer.read())) else: - raw_lines = sys.stdin.read().splitlines() + raw_lines = split_lf_lines(sys.stdin.read()) else: raise MemNetError("no_input", "provide line, --file, or --stdin") if len(raw_lines) > caps.max_batch_lines: @@ -1826,7 +1837,10 @@ def prune_orphans( ss, lock = _load_session(session, exclusive=apply) with lock: rows = orphan_rows(ss, tag=tag) - _prune_rows(ss, rows, apply, "orphans") + try: + _prune_rows(ss, rows, apply, "orphans") + except MemNetError as exc: + _handle_error(exc) @prune_app.command("dangling") @@ -1848,22 +1862,25 @@ def prune_recyclable( def _prune_kind(session: str | None, kind: str, apply: bool) -> None: ss, lock = _load_session(session, exclusive=apply) with lock: - if kind == "stale": - rows = stale_rows(ss) - if apply: - deleted = prune_stale(ss) + try: + if kind == "stale": + rows = stale_rows(ss) + if apply: + deleted = prune_stale(ss) + else: + deleted = [] + elif kind == "recyclable": + rows = recyclable_rows(ss) + deleted = prune_rows(ss, rows) if apply else [] + elif kind == "dangling": + rows = dangling_rows(ss) + deleted = prune_rows(ss, rows) if apply else [] else: + rows = [] deleted = [] - elif kind == "recyclable": - rows = recyclable_rows(ss) - deleted = prune_rows(ss, rows) if apply else [] - elif kind == "dangling": - rows = dangling_rows(ss) - deleted = prune_rows(ss, rows) if apply else [] - else: - rows = [] - deleted = [] - _prune_rows(ss, rows, apply, kind, deleted=deleted) + _prune_rows(ss, rows, apply, kind, deleted=deleted) + except MemNetError as exc: + _handle_error(exc) def _prune_rows( diff --git a/parts/common/memnet/memnet/gql.py b/parts/common/memnet/memnet/gql.py index d769f07..05dc5b6 100644 --- a/parts/common/memnet/memnet/gql.py +++ b/parts/common/memnet/memnet/gql.py @@ -21,6 +21,7 @@ ) from memnet.models import SHAPE_DROP_KEYS from memnet.tier_a import Document, EdgeRec, Field, NodeRec, Op, Section +from memnet.wire import split_lf_lines _IDENT = r"[A-Za-z_][A-Za-z0-9_]*" _LABEL = r"[A-Za-z_][A-Za-z0-9_]*" @@ -544,7 +545,8 @@ def _split_statements(text: str) -> list[tuple[int, str]]: statements: list[tuple[int, str]] = [] buf: list[str] = [] start_line = 1 - for line_no, raw in enumerate(text.splitlines(), start=1): + + for line_no, raw in enumerate(split_lf_lines(text), start=1): s = raw.strip() if not s or s.startswith("#"): continue diff --git a/parts/common/memnet/memnet/housekeep.py b/parts/common/memnet/memnet/housekeep.py index f46d591..94b7dc5 100644 --- a/parts/common/memnet/memnet/housekeep.py +++ b/parts/common/memnet/memnet/housekeep.py @@ -12,6 +12,7 @@ from dataclasses import dataclass, field from memnet.config import ORPHAN_EXEMPT_TAGS +from memnet.exceptions import MemNetError from memnet.models import Record from memnet.session import SessionStore @@ -23,6 +24,51 @@ class _Buckets: orphans: list[Record] = field(default_factory=list) rows_non_law: int = 0 edges: int = 0 + referenced: set[str] = field(default_factory=set) + + +def _node_tokens(rec: Record, dict_key: str) -> set[str]: + """Identities an edge endpoint may use: hid, dict key, optional nickname.""" + tokens = {rec.hid, dict_key} + nick = rec.fields.get("id", "") + if nick: + tokens.add(nick) + return tokens + + +def _resolve_endpoint(store: SessionStore, token: str, token_to_hid: dict[str, str]) -> str | None: + """Map an edge src/dist token to a live node hid. + + GQL stores ``_elN``. Leftover pipe and older snapshots may store the + nickname. Both count. An unresolved token is a dangling end. + """ + if not token: + return None + hid = token_to_hid.get(token) + if hid is not None: + return hid + found = store.store.resolve_one(token) + if found is None or found.tag == "EDG" or found.tag == "LAW": + return None + return found.hid + + +def _refuse_referenced(rows: list[Record], referenced: set[str]) -> None: + """Do not delete a node that any edge still names.""" + blocked = [ + rec + for rec in rows + if rec.tag not in ("EDG", "LAW") and rec.kind == "node" and rec.hid in referenced + ] + if not blocked: + return + sample = blocked[0] + nick = sample.fields.get("id") or sample.hid + raise MemNetError( + "prune_referenced", + (f"{len(blocked)} node(s) still referenced by an edge|refusing delete {nick}"), + example="housekeep dangling; do not prune nodes an edge still names", + ) def _categorise( @@ -31,10 +77,14 @@ def _categorise( orphan_tag: str | None = None, orphan_include_tags: set[str] | None = None, ) -> _Buckets: - """One pass over the store; emits all housekeep buckets + counts.""" + """Classify rows. Edges are resolved before orphan decisions. + + Endpoint tokens match a node by hidden element id (``_elN``) or by + the optional ``id`` nickname. A node is an orphan only when no edge + names it. An edge is dangling when either end does not resolve. + """ by_id = store.store._by_hid - node_ids: set[str] = set() - incident: set[str] = set() + token_to_hid: dict[str, str] = {} records: list[Record] = [] buckets = _Buckets() for rid, rec in by_id.items(): @@ -42,31 +92,33 @@ def _categorise( continue buckets.rows_non_law += 1 records.append(rec) - if rec.tag == "EDG": - buckets.edges += 1 - incident.add(rec.fields.get("src", "")) - incident.add(rec.fields.get("dist", "")) - elif rec.kind == "node": - node_ids.add(rid) - incident.discard("") + if rec.tag != "EDG" and rec.kind == "node": + for token in _node_tokens(rec, rid): + token_to_hid.setdefault(token, rec.hid) exempt = ORPHAN_EXEMPT_TAGS - (orphan_include_tags or set()) tag_filter = orphan_tag.upper() if orphan_tag else None for rec in records: if rec.is_recyclable(): buckets.recyclable.append(rec) - if rec.tag == "EDG": - src = rec.fields.get("src", "") - dist = rec.fields.get("dist", "") - if src not in node_ids or dist not in node_ids: - buckets.dangling.append(rec) + if rec.tag != "EDG": continue - if rec.kind != "node": + buckets.edges += 1 + src = _resolve_endpoint(store, rec.fields.get("src", ""), token_to_hid) + dist = _resolve_endpoint(store, rec.fields.get("dist", ""), token_to_hid) + if src is None or dist is None: + buckets.dangling.append(rec) + if src is not None: + buckets.referenced.add(src) + if dist is not None: + buckets.referenced.add(dist) + for rec in records: + if rec.tag == "EDG" or rec.kind != "node": continue if rec.tag in exempt: continue if tag_filter and rec.tag != tag_filter: continue - if rec.id not in incident: + if rec.hid not in buckets.referenced: buckets.orphans.append(rec) buckets.dangling.sort(key=lambda r: r.id) buckets.orphans.sort(key=lambda r: r.id) @@ -96,8 +148,8 @@ def stale_rows(store: SessionStore) -> list[Record]: combined: list[Record] = [] for group in (buckets.recyclable, buckets.dangling, buckets.orphans): for rec in group: - if rec.id not in seen: - seen.add(rec.id) + if rec.hid not in seen: + seen.add(rec.hid) combined.append(rec) return combined @@ -115,22 +167,37 @@ def stats(store: SessionStore) -> dict[str, int]: def prune_rows(store: SessionStore, rows: list[Record]) -> list[Record]: + """Delete ``rows`` after a fresh categorisation. + + Refuses, and deletes nothing, when a node is still an edge endpoint. + """ + buckets = _categorise(store) + _refuse_referenced(rows, buckets.referenced) deleted: list[Record] = [] for rec in rows: - if store.store.delete(rec.id): + if store.store.delete(rec.hid): deleted.append(rec) return deleted def prune_stale(store: SessionStore) -> list[Record]: + """Recompute stale rows, then delete them. + + Same refusal as ``prune_rows``: a node any edge still references is + not deleted, and neither is the rest of the batch. + """ buckets = _categorise(store) seen: set[str] = set() - out: list[Record] = [] + rows: list[Record] = [] for group in (buckets.recyclable, buckets.dangling, buckets.orphans): for rec in group: - if rec.id in seen: + if rec.hid in seen: continue - if store.store.delete(rec.id): - seen.add(rec.id) - out.append(rec) + seen.add(rec.hid) + rows.append(rec) + _refuse_referenced(rows, buckets.referenced) + out: list[Record] = [] + for rec in rows: + if store.store.delete(rec.hid): + out.append(rec) return out diff --git a/parts/common/memnet/memnet/in_process_engine.py b/parts/common/memnet/memnet/in_process_engine.py index 164b7b9..05ffe49 100644 --- a/parts/common/memnet/memnet/in_process_engine.py +++ b/parts/common/memnet/memnet/in_process_engine.py @@ -2,11 +2,11 @@ from __future__ import annotations -import io import os -import sys from typing import Any +from memnet.output import capture_request_stdio + def run_argv(argv: list[str], *, stdin: str | None = None) -> dict[str, Any]: """Execute memnet CLI argv in-process; return JSON-shaped envelope.""" @@ -14,25 +14,20 @@ def run_argv(argv: list[str], *, stdin: str | None = None) -> dict[str, Any]: os.environ["MEMNET_SERVE_INTERNAL"] = "1" from memnet.cli import app - out = io.StringIO() - err = io.StringIO() - old_out, old_err, old_in = sys.stdout, sys.stderr, sys.stdin - sys.stdout, sys.stderr = out, err - if stdin: - sys.stdin = io.StringIO(stdin) code = 0 - try: - result = app(argv, prog_name="memnet", standalone_mode=False) - if isinstance(result, int) and result != 0: - code = result - except SystemExit as exc: - code = int(exc.code) if isinstance(exc.code, int) else 1 - except Exception as exc: - code = 1 - err.write(f"@ERR: internal|{type(exc).__name__}: {exc}\n") - finally: - sys.stdout, sys.stderr, sys.stdin = old_out, old_err, old_in - return {"exit_code": code, "stdout": out.getvalue(), "stderr": err.getvalue()} + with capture_request_stdio(stdin) as (out, err): + try: + result = app(argv, prog_name="memnet", standalone_mode=False) + if isinstance(result, int) and result != 0: + code = result + except SystemExit as exc: + code = int(exc.code) if isinstance(exc.code, int) else 1 + except Exception as exc: + code = 1 + err.write(f"@ERR: internal|{type(exc).__name__}: {exc}\n") + stdout = out.getvalue() + stderr = err.getvalue() + return {"exit_code": code, "stdout": stdout, "stderr": stderr} class InProcessEngine: diff --git a/parts/common/memnet/memnet/mutate_gate.py b/parts/common/memnet/memnet/mutate_gate.py index b0e099b..318dc08 100644 --- a/parts/common/memnet/memnet/mutate_gate.py +++ b/parts/common/memnet/memnet/mutate_gate.py @@ -852,6 +852,8 @@ def _node_to_record(self, node: NodeRec) -> Record: if kind and not tag_def: known = ",".join(self.ss.tag_map.tag_names()) raise MemNetError("unknown_tag", f"{kind} not in schema known: {known}") + if tag_def is not None: + kind = tag_def.tag fields: dict[str, str] = {} base = dict(bound.fields) if bound else {} diff --git a/parts/common/memnet/memnet/output.py b/parts/common/memnet/memnet/output.py index 60c1cb8..84c344c 100644 --- a/parts/common/memnet/memnet/output.py +++ b/parts/common/memnet/memnet/output.py @@ -2,19 +2,24 @@ from __future__ import annotations +import io import sys +import threading +from collections.abc import Iterator +from contextlib import contextmanager +from contextvars import ContextVar +from typing import Any, TextIO from memnet.exceptions import MemNetError from memnet.models import Record, TagDef, TagMap from memnet.wire import join_payload _MAX_WRN = 12 -_WARN_EMITTED = 0 +_warn_emitted: ContextVar[int] = ContextVar("memnet_warn_emitted", default=0) def reset_warn_budget() -> None: - global _WARN_EMITTED - _WARN_EMITTED = 0 + _warn_emitted.set(0) def emit_stdout(line: str) -> None: @@ -39,11 +44,11 @@ def format_wrn( *, force: bool = False, ) -> str | None: - global _WARN_EMITTED if not force: - if _WARN_EMITTED >= _MAX_WRN: + emitted = _warn_emitted.get() + if emitted >= _MAX_WRN: return None - _WARN_EMITTED += 1 + _warn_emitted.set(emitted + 1) msg = message.replace("|", " ") if example: return f"@WRN: {code}|{msg}|{example}" @@ -109,3 +114,141 @@ def parse_err_line(line: str) -> tuple[str, str, str | None]: def values_from_record(record: Record, tag_def: TagDef) -> list[str]: return [record.fields.get(f, "") for f in tag_def.fields] + + +# Per-request stdio. A process-global swap of sys.stdout races on a threaded +# serve (MN-REQ-06.13). The proxy stays installed only while at least one +# capture is active; each request's buffers live on contextvars. +_stdout_buf: ContextVar[io.StringIO | None] = ContextVar("memnet_stdout_buf", default=None) +_stderr_buf: ContextVar[io.StringIO | None] = ContextVar("memnet_stderr_buf", default=None) +_stdin_buf: ContextVar[io.StringIO | None] = ContextVar("memnet_stdin_buf", default=None) +_capture_guard = threading.Lock() +_capture_depth = 0 +_saved_stdio: tuple[TextIO, TextIO, TextIO] | None = None +_fallback_out: list[TextIO] = [sys.stdout] +_fallback_err: list[TextIO] = [sys.stderr] +_fallback_in: list[TextIO] = [sys.stdin] + + +class _TextAsBuffer: + """Bytes view of a text stdin so ``sys.stdin.buffer.read`` stays in-request.""" + + def __init__(self, textio: TextIO) -> None: + self._textio = textio + + def read(self, n: int = -1) -> bytes: + text = self._textio.read() if n < 0 else self._textio.read(n) + return text.encode("utf-8") + + def readline(self, n: int = -1) -> bytes: + text = self._textio.readline() if n < 0 else self._textio.readline(n) + return text.encode("utf-8") + + +class _RoutingStream: + def __init__(self, var: ContextVar[io.StringIO | None], fallback: list[TextIO]) -> None: + self._var = var + self._fallback = fallback + + def _target(self) -> TextIO: + buf = self._var.get() + if buf is not None: + return buf + return self._fallback[0] + + def write(self, s: str) -> int: + return self._target().write(s) + + def flush(self) -> None: + flush = getattr(self._target(), "flush", None) + if flush: + flush() + + def isatty(self) -> bool: + fn = getattr(self._target(), "isatty", None) + return bool(fn()) if fn else False + + def read(self, n: int = -1) -> str: + return self._target().read(n) + + def readline(self, n: int = -1) -> str: + return self._target().readline(n) + + def readlines(self, hint: int = -1) -> list[str]: + return self._target().readlines(hint) + + def __iter__(self) -> Iterator[str]: + return iter(self._target()) + + @property + def encoding(self) -> str: + return getattr(self._target(), "encoding", None) or "utf-8" + + @property + def errors(self) -> str | None: + return getattr(self._target(), "errors", None) + + @property + def buffer(self) -> Any: + target = self._target() + buf = getattr(target, "buffer", None) + if buf is not None: + return buf + return _TextAsBuffer(target) + + def __getattr__(self, name: str) -> Any: + return getattr(self._target(), name) + + +_OUT_PROXY = _RoutingStream(_stdout_buf, _fallback_out) +_ERR_PROXY = _RoutingStream(_stderr_buf, _fallback_err) +_IN_PROXY = _RoutingStream(_stdin_buf, _fallback_in) + + +def _enter_capture() -> None: + global _capture_depth, _saved_stdio + with _capture_guard: + if _capture_depth == 0: + _saved_stdio = (sys.stdout, sys.stderr, sys.stdin) + _fallback_out[0] = sys.stdout + _fallback_err[0] = sys.stderr + _fallback_in[0] = sys.stdin + sys.stdout = _OUT_PROXY + sys.stderr = _ERR_PROXY + sys.stdin = _IN_PROXY + _capture_depth += 1 + + +def _leave_capture() -> None: + global _capture_depth, _saved_stdio + with _capture_guard: + _capture_depth -= 1 + if _capture_depth == 0 and _saved_stdio is not None: + sys.stdout, sys.stderr, sys.stdin = _saved_stdio + _saved_stdio = None + + +@contextmanager +def capture_request_stdio( + stdin_text: str | None = None, +) -> Iterator[tuple[io.StringIO, io.StringIO]]: + """Bind this context's stdout, stderr, and stdin for one serve command. + + Overlapping requests do not share buffers. The proxy is process-global + only as a router; the bytes live on contextvars, so a threaded serve + does not need a serve-wide lock. + """ + out = io.StringIO() + err = io.StringIO() + inp = io.StringIO(stdin_text or "") + _enter_capture() + tok_out = _stdout_buf.set(out) + tok_err = _stderr_buf.set(err) + tok_in = _stdin_buf.set(inp) + try: + yield out, err + finally: + _stdout_buf.reset(tok_out) + _stderr_buf.reset(tok_err) + _stdin_buf.reset(tok_in) + _leave_capture() diff --git a/parts/common/memnet/memnet/serve.py b/parts/common/memnet/memnet/serve.py index 161216b..ea9672d 100644 --- a/parts/common/memnet/memnet/serve.py +++ b/parts/common/memnet/memnet/serve.py @@ -11,7 +11,6 @@ from __future__ import annotations -import io import ipaddress import json import logging @@ -102,28 +101,25 @@ def _handle_request(payload: dict[str, Any]) -> dict[str, Any]: } os.environ["MEMNET_SERVE_INTERNAL"] = "1" from memnet.cli import app + from memnet.output import capture_request_stdio - out = io.StringIO() - err = io.StringIO() - old_out, old_err, old_in = sys.stdout, sys.stderr, sys.stdin - sys.stdout, sys.stderr = out, err - if stdin_text: - sys.stdin = io.StringIO(stdin_text) code = 0 token_ctx = set_caller_token(token_s) - try: - result = app(argv, prog_name="memnet", standalone_mode=False) - if isinstance(result, int) and result != 0: - code = result - except SystemExit as exc: - code = int(exc.code) if isinstance(exc.code, int) else 1 - except Exception as exc: - code = 1 - err.write(f"@ERR: internal|{type(exc).__name__}: {exc}\n") - finally: - reset_caller_token(token_ctx) - sys.stdout, sys.stderr, sys.stdin = old_out, old_err, old_in - return {"exit_code": code, "stdout": out.getvalue(), "stderr": err.getvalue()} + with capture_request_stdio(stdin_text if isinstance(stdin_text, str) else None) as (out, err): + try: + result = app(argv, prog_name="memnet", standalone_mode=False) + if isinstance(result, int) and result != 0: + code = result + except SystemExit as exc: + code = int(exc.code) if isinstance(exc.code, int) else 1 + except Exception as exc: + code = 1 + err.write(f"@ERR: internal|{type(exc).__name__}: {exc}\n") + finally: + reset_caller_token(token_ctx) + stdout = out.getvalue() + stderr = err.getvalue() + return {"exit_code": code, "stdout": stdout, "stderr": stderr} class _Handler(socketserver.BaseRequestHandler): @@ -226,6 +222,9 @@ def _read_exact(self, n: int) -> bytes | None: class _Server(socketserver.ThreadingTCPServer): allow_reuse_address = True daemon_threads = True + # Default listen backlog is 5. A burst of clients (and the gateway's + # probe plus command) would be refused before a thread accepts them. + request_queue_size = 128 def run_serve(host: str | None = None, port: int | None = None) -> None: diff --git a/parts/common/memnet/memnet/snapshot.py b/parts/common/memnet/memnet/snapshot.py index 5869c58..8eedc96 100644 --- a/parts/common/memnet/memnet/snapshot.py +++ b/parts/common/memnet/memnet/snapshot.py @@ -46,6 +46,25 @@ def _nick_of(rec: Record) -> str: return rec.fields.get("id") or rec.tag +def _canonical_tag(ss: SessionStore, tag: str) -> str: + """SCHEMA key for a stored label. TagMap lookup is case-folding.""" + td = ss.tag_map.get(tag) + if td is not None: + return td.tag + return tag.upper() + + +def _append_id_column(tag: str, fields: list[str]) -> list[str]: + """Snapshot SCHEMA gains ``id`` when the live map omitted it (MN-REQ-01.9). + + Live SCHEMA is not rewritten. Fixed tags keep their declared columns. + ``id`` is appended, not forced first. + """ + if tag in FIXED_TAGS or "id" in fields: + return fields + return [*fields, "id"] + + def _snapshot_emit_tag_map(ss: SessionStore) -> TagMap: """Widen SCHEMA with undeclared RAM keys so extras round-trip (MN-REQ-01.9).""" extras: dict[str, list[str]] = {} @@ -53,14 +72,15 @@ def _snapshot_emit_tag_map(ss: SessionStore) -> TagMap: rec = ss.store._by_hid.get(rid) if not rec: continue - td = ss.tag_map.get(rec.tag) + canon = _canonical_tag(ss, rec.tag) + td = ss.tag_map.get(canon) base = list(td.fields) if td else ["id"] known = set(base) - extra = extras.setdefault(rec.tag, []) + extra = extras.setdefault(canon, []) for key in rec.fields: if key in known: continue - if rec.tag in FIXED_TAGS: + if canon in FIXED_TAGS or rec.tag in FIXED_TAGS: raise MemNetError( "snapshot_unsaveable", f"{rec.tag} nick={_nick_of(rec)} field={key} fixed_tag extra", @@ -70,7 +90,7 @@ def _snapshot_emit_tag_map(ss: SessionStore) -> TagMap: tags: dict[str, TagDef] = {} max_fields = ss.caps.max_fields for tag, td in ss.tag_map.tags.items(): - fields = list(td.fields) + extras.get(tag, []) + fields = _append_id_column(tag, list(td.fields) + extras.get(tag, [])) if len(fields) > max_fields: raise MemNetError( "snapshot_unsaveable", @@ -80,7 +100,7 @@ def _snapshot_emit_tag_map(ss: SessionStore) -> TagMap: for tag, extra in extras.items(): if tag in tags or not extra: continue - fields = ["id"] + extra + fields = _append_id_column(tag, ["id"] + extra) if len(fields) > max_fields: raise MemNetError( "snapshot_unsaveable", @@ -112,7 +132,7 @@ def _emit_record_lines( token = fields.get(key, "") if token in hid_to_nick: fields[key] = hid_to_nick[token] - clone = rec.model_copy(update={"fields": fields}) + clone = rec.model_copy(update={"fields": fields, "tag": _canonical_tag(ss, rec.tag)}) rec_lines.append(emit_record(clone, emit_map)) return rec_lines, hid_to_nick, emit_map @@ -150,7 +170,23 @@ def _verify_emitted_snapshot( ) -> None: """Refuse save if the formatted blob would not load as the same values.""" used_nicks: set[str] = set() - _, _, _, rec_lines = _parse_sections(split_snapshot_lines(text)) + _, map_lines, _, rec_lines = _parse_sections(split_snapshot_lines(text)) + try: + loaded_map = load_persisted_map_from_lines(map_lines, ss.caps) + except MemNetError as exc: + raise MemNetError( + "snapshot_unsaveable", + f"map {exc.code} {exc.message}", + ) from exc + for tag, td in emit_map.tags.items(): + if tag in FIXED_TAGS: + continue + loaded = loaded_map.get(tag) + if loaded is None or list(loaded.fields) != list(td.fields): + raise MemNetError( + "snapshot_unsaveable", + f"{tag} nick=- schema not reloadable", + ) rec_iter = iter(rec_lines) for rid in ss.store.write_order: rec = ss.store._by_hid.get(rid) @@ -164,7 +200,7 @@ def _verify_emitted_snapshot( f"{rec.tag} nick={_nick_of(rec)} missing emit row", ) from exc try: - parsed = parse_line(line, emit_map, ss.caps, used_nicks=used_nicks) + parsed = parse_line(line, loaded_map, ss.caps, used_nicks=used_nicks) except MemNetError as exc: raise MemNetError( "snapshot_unsaveable", diff --git a/parts/common/memnet/memnet/warnings.py b/parts/common/memnet/memnet/warnings.py index 98e99b2..9d45da3 100644 --- a/parts/common/memnet/memnet/warnings.py +++ b/parts/common/memnet/memnet/warnings.py @@ -47,6 +47,8 @@ def emit_stale_warnings(store: SessionStore) -> None: if d: emit_wrn("stale_dangling", f"{d}|housekeep dangling or prune dangling --apply") if o: + # Only true orphans (no edge endpoint, hid or nickname) reach here. + # Do not suggest prune while edges still name the nodes. emit_wrn("stale_orphans", f"{o}|housekeep orphans or prune orphans --apply") if r or d or o: emit_wrn("stale_graph", f"{r}|{d}|{o}|housekeep stale") diff --git a/parts/common/memnet/memnet/wire.py b/parts/common/memnet/memnet/wire.py index 7edbc6c..57b7db2 100644 --- a/parts/common/memnet/memnet/wire.py +++ b/parts/common/memnet/memnet/wire.py @@ -22,12 +22,15 @@ _HEX = frozenset("0123456789abcdefABCDEF") -def split_snapshot_lines(text: str) -> list[str]: - """Record split for leftover snapshots: LF only, optional CRLF trim. +def split_lf_lines(text: str) -> list[str]: + """Split on LF only and strip one trailing CR from each line. - MUST NOT use str.splitlines() — that splits on CR / VT / FF / NEL / LS / PS - before unescape and yields FIELD_COUNT. + Empty text is no lines. MUST NOT use str.splitlines() — that also + breaks on CR, VT, FF, FS, GS, RS, NEL, LS, and PS. Snapshot records + and mutate statement/stdin lines share this rule. """ + if text == "": + return [] if text.endswith("\n"): text = text[:-1] lines: list[str] = [] @@ -38,6 +41,17 @@ def split_snapshot_lines(text: str) -> list[str]: return lines +def split_snapshot_lines(text: str) -> list[str]: + """Record split for leftover snapshots: LF only, optional CRLF trim. + + MUST NOT use str.splitlines() — that splits on CR / VT / FF / NEL / LS / PS + before unescape and yields FIELD_COUNT. + """ + if text == "": + return [""] + return split_lf_lines(text) + + def split_payload(payload: str) -> list[str]: if "\\" not in payload: return payload.split("|") diff --git a/sysml-models/models/deploy.sysml b/sysml-models/models/deploy.sysml index 2f282d4..0d6d97b 100644 --- a/sysml-models/models/deploy.sysml +++ b/sysml-models/models/deploy.sysml @@ -433,6 +433,8 @@ package MemNet { port errOut : ErrorSignalOutPort; port libraryLocatorIn : LibraryLocatorInPort; attribute libraryIngestThisCut : Boolean = true; + attribute statementSplitLfOnly : Boolean = true; + attribute labelMatchCaseFold : Boolean = true; part allocator : IdAllocator; part caps : CapsPolicy; attribute honourWhereOrRefuse : Boolean = true; @@ -1053,12 +1055,19 @@ package MemNet { not TTL expire-save and not a second store. usageLookOut is a read-only stats snapshot for MemNetUsageDashboard — not a manage/prune command from the page. + Orphan / dangling / stale resolve edge ends by hid and nickname + (MN-REQ-04.12). prune --apply recomputes and refuses a node an + edge still names. */ port statsOut : GqlWireOutPort; port usageLookOut : ServeUsageLookOutPort; port errOut : ErrorSignalOutPort; + attribute endpointIdentityResolvesHid : Boolean = true; + attribute pruneApplyRefusesReferenced : Boolean = true; + attribute autoPruneApply : Boolean = false; satisfy MN_REQ_00_MissionBridge::MN_REQ_04_SliceEconomy::MN_REQ_04_2_RecycleHiddenFromPinMap; + satisfy MN_REQ_00_MissionBridge::MN_REQ_04_SliceEconomy::MN_REQ_04_12_HousekeepEndpointIdentity; } part def SnapshotStore { @@ -1105,6 +1114,7 @@ package MemNet { attribute splitlinesSeparatorsEscaped : Boolean = true; attribute recordSplitLfOnly : Boolean = true; attribute persistUndeclaredProperties : Boolean = true; + attribute idColumnWidenedWhenAbsent : Boolean = true; attribute lineBytesOnEmittedLine : Boolean = true; satisfy MN_REQ_00_MissionBridge::MN_REQ_01_SessionLifecycle::MN_REQ_01_4_SaveSessionSnapshot; @@ -1317,9 +1327,12 @@ package MemNet { part adminUsage : AdminUsageReport; attribute host : String = "127.0.0.1"; attribute listenPort : Integer = 18765; + attribute requestOutputIsolated : Boolean = true; + attribute serialisesRequests : Boolean = false; satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_3_TcpServeMigration; satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_11_AdminServeUsageReport; + satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_13_ServeRequestIsolation; satisfy MN_REQ_00_MissionBridge::MN_REQ_12_MultitaskMode::MN_REQ_12_2_SharedStoreTransport; } diff --git a/sysml-models/models/requirements.sysml b/sysml-models/models/requirements.sysml index 40d8038..251101a 100644 --- a/sysml-models/models/requirements.sysml +++ b/sysml-models/models/requirements.sysml @@ -20,6 +20,8 @@ never default goldfish MN-REQ-04.11 pin_map order is a function of Shape observables, not hid / nickname id / CREATE order (gauge) + MN-REQ-04.12 Housekeep orphan/dangling/stale resolve edge ends by + hid and nickname; prune --apply refuses a referenced node MN-REQ-06.5 Human usage look (read-only dashboard; not agent wire) MN-REQ-06.6 Device MemNet services; one MemNet MCP (tip/ops) at the droplet; product face is sysmledge (tip≠face) MN-REQ-06.7 Allocate SSOT parts to live code in this checkout (one Hatch wheel, many hosts; sysmledge not in wheel) @@ -33,6 +35,7 @@ MN-REQ-05.3 One decoded value_bytes cap on pipe mutate, GQL mutate, snapshot load MN-REQ-06.11 Admin-only serve usage report (counts/caps; opaque session alias; not agent MCP) MN-REQ-06.12 Product gateway on memnet-mcp: per-product route to owning serve (not mn_tip_) + MN-REQ-06.13 Serve request isolation: each command's stdout, stderr, and exit_code MN-REQ-11.17 Catalog Snap / model Snap — session strata (0.15) MN-REQ-11.17.1 Cross-cut satisfy locators on the catalog; SysMLEdge pin_map is a different tool (not this catalog) @@ -225,10 +228,19 @@ package MemNetRequirements { emit SHALL escape those losslessly; load SHALL decode them. Record split SHALL be LF-only (MUST NOT use str.splitlines(), which splits on CR/VT/FF/NEL/LS/PS - before unescape). Properties mutate stores that are absent - from the live tag SCHEMA SHALL persist by widening the - snapshot SCHEMA (GraphElement extras and Path-B locators - such as qname). Live session SCHEMA is unchanged. Fixed + before unescape). Mutate stdin, leftover pipe batches, + and GQL statement breaks SHALL use that same LF-only + split (optional trailing CR stripped). A raw VT, FF, + FS, GS, RS, NEL, LS, or PS inside a value SHALL stay + inside that value and SHALL NOT start another statement. + Properties mutate stores that are absent from the live + tag SCHEMA SHALL persist by widening the snapshot SCHEMA + (GraphElement extras and Path-B locators such as qname), + including when the CREATE label spelling differs in case + from the SCHEMA tag. A live SCHEMA that omits the id + column SHALL still save: the snapshot SCHEMA gains an + id column and a minted nickname. Live session SCHEMA is + unchanged until load. Fixed tags EDG/LAW SHALL NOT widen; extras there SHALL fail closed. Save SHALL fail closed (snapshot_unsaveable naming the row) rather than write an unloadable file, @@ -652,6 +664,25 @@ package MemNetRequirements { */ attribute requirementId : String = "MN-REQ-04.11"; } + + requirement def MN_REQ_04_12_HousekeepEndpointIdentity { + doc /* + Housekeep orphan, dangling, and stale counts SHALL resolve + each edge endpoint by hidden element id (_elN) and by the + optional id nickname. A GQL edge stores _elN. Comparing that + token only to the nickname SHALL NOT mark every node an + orphan. A node is an orphan only when no edge names it. + An edge is dangling only when an end does not resolve. + prune orphans and prune stale with --apply SHALL recompute + with that rule and SHALL refuse prune_referenced, deleting + nothing, when a candidate node is still an edge endpoint. + The stale_orphans warning SHALL NOT suggest prune while + edges still name those nodes. No sweep, housekeep job, or + MCP tool SHALL run prune --apply; only an explicit CLI + --apply does. + */ + attribute requirementId : String = "MN-REQ-04.12"; + } } // ----- Hard caps ----- @@ -1087,6 +1118,23 @@ package MemNetRequirements { */ attribute requirementId : String = "MN-REQ-06.12"; } + + requirement def MN_REQ_06_13_ServeRequestIsolation { + doc /* + Concurrent memnet serve commands (TCP and the product + gateway, which forwards to that serve) SHALL isolate + each request's stdout, stderr, and exit_code. One + caller's records SHALL NOT appear in another caller's + envelope. A failed write SHALL NOT come back with + another call's success text or exit_code 0. + Capture SHALL be per request (context-local streams). + A process-wide swap of sys.stdout / sys.stderr is not + sufficient on a threaded serve. Requests on different + sessions MAY still run together; isolation SHALL NOT + require a serve-wide lock. + */ + attribute requirementId : String = "MN-REQ-06.13"; + } } // ----- MCP agent boundary ----- @@ -2079,6 +2127,8 @@ package MemNetRequirements { : MN_REQ_00_MissionBridge::MN_REQ_04_SliceEconomy::MN_REQ_04_10_PeakLLastResort; requirement pinMapOrderFromObservablesReq : MN_REQ_00_MissionBridge::MN_REQ_04_SliceEconomy::MN_REQ_04_11_PinMapOrderFromObservables; + requirement housekeepEndpointIdentityReq + : MN_REQ_00_MissionBridge::MN_REQ_04_SliceEconomy::MN_REQ_04_12_HousekeepEndpointIdentity; requirement hardCapsReq : MN_REQ_00_MissionBridge::MN_REQ_05_HardCaps; @@ -2115,6 +2165,8 @@ package MemNetRequirements { : MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_11_AdminServeUsageReport; requirement productGatewayReq : MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_12_ProductGateway; + requirement serveRequestIsolationReq + : MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_13_ServeRequestIsolation; requirement mcpAgentBoundaryReq : MN_REQ_00_MissionBridge::MN_REQ_07_McpAgentBoundary; diff --git a/sysml-models/models/verify.sysml b/sysml-models/models/verify.sysml index 62a1b72..317d2ae 100644 --- a/sysml-models/models/verify.sysml +++ b/sysml-models/models/verify.sysml @@ -1632,6 +1632,47 @@ package MemNetVerification { } } + verification def MN_VER_04_S06_HousekeepEndpointIdentity { + doc /* + MN-REQ-04.12 — housekeep resolves edge endpoints by hid and + nickname. prune orphans / prune stale --apply recomputes and + refuses a referenced node. Nothing runs prune --apply by itself. + */ + attribute verificationId : String = "MN-VER-04-S06"; + + subject housekeep : HousekeepSettle; + + objective housekeepEndpointIdentity { + verify housekeepEndpointIdentityReq; + require constraint { + housekeep.endpointIdentityResolvesHid == true + and housekeep.pruneApplyRefusesReferenced == true + and housekeep.autoPruneApply == false + } + } + } + + verification def MN_VER_06_S11_ServeRequestIsolation { + doc /* + MN-REQ-06.13 — each memnet serve command keeps its own + stdout, stderr, and exit_code. TCP and the product gateway + (which forwards to that serve) share this rule. Capture is + per request. Different sessions MAY run together + (serialisesRequests=false). + */ + attribute verificationId : String = "MN-VER-06-S11"; + + subject tcp : TcpServeBridge; + + objective serveRequestIsolation { + verify serveRequestIsolationReq; + require constraint { + tcp.requestOutputIsolated == true + and tcp.serialisesRequests == false + } + } + } + verification def MN_VER_13_S01_RecallCommitTwoOperators { doc /* MN-REQ-13.1 — two-operator cut (math skeleton; leftover goldfish). @@ -1960,8 +2001,14 @@ package MemNetVerification { .recordSplitLfOnly == true and system.core.transport.inProcess.memory.sessions.snap .persistUndeclaredProperties == true + and system.core.transport.inProcess.memory.sessions.snap + .idColumnWidenedWhenAbsent == true and system.core.transport.inProcess.memory.sessions.snap .lineBytesOnEmittedLine == true + and system.core.transport.inProcess.memory.sessions.recallCommit + .commit.mutate.statementSplitLfOnly == true + and system.core.transport.inProcess.memory.sessions.recallCommit + .commit.mutate.labelMatchCaseFold == true and system.core.cli.sessionSave.failClosedUnsaveable == true and system.core.cli.sessionLoad.legacy01918Load == true and system.core.transport.inProcess.memory.sessions.caps @@ -2148,4 +2195,8 @@ package MemNetVerification { : MN_VER_06_S07_LanMcpFrontSeveralServes; verification clusterRouteVsSliceHandCarryVerify : MN_VER_06_S08_ClusterRouteVsSliceHandCarry; + verification serveRequestIsolationVerify + : MN_VER_06_S11_ServeRequestIsolation; + verification housekeepEndpointIdentityVerify + : MN_VER_04_S06_HousekeepEndpointIdentity; } diff --git a/tests/test_endleaf_keeper.py b/tests/test_endleaf_keeper.py new file mode 100644 index 0000000..5ef42c6 --- /dev/null +++ b/tests/test_endleaf_keeper.py @@ -0,0 +1,230 @@ +"""Endleaf Keeper items 2–4: LF-only mutate, snapshot extras, MATCH and TTL.""" + +from __future__ import annotations + +from datetime import UTC, datetime, timedelta +from pathlib import Path + +import pytest + +from memnet.cli import _read_ingest_input +from memnet.config import Caps +from memnet.exceptions import MemNetError +from memnet.mutate_gate import MutateGate +from memnet.session import get_session, open_session, set_now_override +from memnet.snapshot import load_snapshot, write_snapshot +from memnet.wire import SPLITLINES_SEPARATORS + +_MAP = ["SCHEMA CST ; fields=id name role"] +_INJECT_SEPS = [sep for sep in SPLITLINES_SEPARATORS if sep != "\n"] + + +def _nicks(ss) -> set[str]: + return {rec.fields.get("id", "") for rec in ss.store._by_hid.values()} + + +@pytest.mark.parametrize("sep", _INJECT_SEPS) +def test_pipe_and_gql_keep_splitlines_separators_inside_values(memnet_temp, sep: str, monkeypatch): + literal = f"@CST: GOOD|ok{sep}x|r" + monkeypatch.setattr("sys.stdin", __import__("io").StringIO(literal)) + lines = _read_ingest_input(None, None, True, Caps()) + assert lines == [literal] + ss = open_session(map_lines=_MAP) + MutateGate(ss).apply(lines, mode="add") + got = ss.store.get("GOOD") + assert got is not None + assert got.fields["name"] == f"ok{sep}x" + + injected = f"@CST: GOOD2|ok|r{sep}@CST: EVIL|injected|r" + monkeypatch.setattr("sys.stdin", __import__("io").StringIO(injected)) + split = _read_ingest_input(None, None, True, Caps()) + assert len(split) == 1 + assert sep in split[0] + try: + MutateGate(ss).apply(split, mode="add") + except MemNetError: + pass + assert "EVIL" not in _nicks(ss) + + gql_keep = f"CREATE (:CST {{id: 'GKEEP', name: 'a{sep}b', role: 'r'}})" + MutateGate(ss).apply([gql_keep], mode="mutate") + kept = ss.store.get("GKEEP") + assert kept is not None + assert kept.fields["name"] == f"a{sep}b" + + gql_inject = ( + "CREATE (:CST {id: 'G2', name: 'ok', role: 'r'})" + + sep + + "CREATE (:CST {id: 'EVIL', name: 'x', role: 'r'})" + ) + try: + MutateGate(ss).apply([gql_inject], mode="mutate") + except MemNetError as exc: + assert exc.code == "parse_error" + assert "EVIL" not in _nicks(ss) + + +def test_lf_splits_statements_and_embedded_cr_does_not(memnet_temp): + ss = open_session(map_lines=_MAP) + MutateGate(ss).apply( + [ + "CREATE (:CST {id: 'GOOD', name: 'ok', role: 'r'})\n" + "CREATE (:CST {id: 'NEXT', name: 'n', role: 'r'})" + ], + mode="mutate", + ) + assert ss.store.get("GOOD") is not None + assert ss.store.get("NEXT") is not None + cr = ( + "CREATE (:CST {id: 'CR1', name: 'ok', role: 'r'})\r" + "CREATE (:CST {id: 'EVIL', name: 'x', role: 'r'})" + ) + try: + MutateGate(ss).apply([cr], mode="mutate") + except MemNetError as exc: + assert exc.code == "parse_error" + assert ss.store.get("EVIL") is None + + +def test_snapshot_mixed_case_label_persists_extra(memnet_temp, tmp_path: Path): + ss = open_session(map_lines=_MAP) + MutateGate(ss).apply( + ["CREATE (:Cst {id: 'N_X', name: 'keep', role: 'r', extra_k: 'extra_v'})"], + mode="add", + ) + rec = ss.store.get("N_X") + assert rec is not None + assert rec.tag == "CST" + assert rec.fields["extra_k"] == "extra_v" + rec.tag = "Cst" + path = tmp_path / "case.snap" + write_snapshot(ss, path) + text = path.read_text(encoding="utf-8") + assert "SCHEMA CST ; fields=id name role extra_k" in text + loaded = load_snapshot(path) + got = loaded.store.get("N_X") + assert got is not None + assert got.fields["extra_k"] == "extra_v" + + +def test_schema_without_id_saves_by_widening_snapshot_only(memnet_temp, tmp_path: Path): + ss = open_session(map_lines=["SCHEMA SEC ; fields=title body status"]) + assert ss.tag_map.get("SEC").fields == ["title", "body", "status"] + MutateGate(ss).apply( + ["CREATE (:SEC {title: 'T', body: 'B', status: 'open'})"], + mode="add", + ) + live = next(rec for rec in ss.store._by_hid.values() if rec.tag == "SEC") + assert "id" not in live.fields + path = tmp_path / "noid.snap" + write_snapshot(ss, path) + text = path.read_text(encoding="utf-8") + assert "SCHEMA SEC ; fields=title body status id" in text + assert ss.tag_map.get("SEC").fields == ["title", "body", "status"] + loaded = load_snapshot(path) + got = next(rec for rec in loaded.store._by_hid.values() if rec.tag == "SEC") + assert got.fields["title"] == "T" + assert got.fields["body"] == "B" + assert got.fields["status"] == "open" + assert got.fields["id"] + assert "id" in loaded.tag_map.get("SEC").fields + + +def test_unloadable_schema_key_fails_closed(memnet_temp, tmp_path: Path): + ss = open_session(map_lines=_MAP) + MutateGate(ss).apply( + ["CREATE (:CST {id: 'N_S', name: 'keep', role: 'r'})"], + mode="add", + ) + rec = ss.store.get("N_S") + assert rec is not None + rec.fields["foo bar"] = "x" + path = tmp_path / "space.snap" + with pytest.raises(MemNetError) as ei: + write_snapshot(ss, path) + assert ei.value.code == "snapshot_unsaveable" + assert not path.exists() + + +def test_labelled_match_set_folds_case_and_miss_is_not_found(memnet_temp): + ss = open_session( + map_lines=[ + "SCHEMA CST ; fields=id name role", + "SCHEMA SEC ; fields=id title", + ] + ) + MutateGate(ss).apply( + ["CREATE (:CST {id: 'N1', name: 'a', role: 'r'})"], + mode="add", + ) + for label in ("cst", "Cst", "CST"): + MutateGate(ss).apply( + [f"MATCH (n:{label} {{id: 'N1'}}) SET n.name = '{label}'"], + mode="mutate", + ) + assert ss.store.get("N1").fields["name"] == label + MutateGate(ss).apply( + ["CREATE (:Cst {id: 'N2', name: 'a', role: 'r'})"], + mode="add", + ) + assert ss.store.get("N2").tag == "CST" + MutateGate(ss).apply( + ["MATCH (n:Cst {id: 'N2'}) SET n.name = 'mixed'"], + mode="mutate", + ) + MutateGate(ss).apply( + ["MATCH (n:CST {id: 'N2'}) SET n.name = 'upper'"], + mode="mutate", + ) + assert ss.store.get("N2").fields["name"] == "upper" + with pytest.raises(MemNetError) as ei: + MutateGate(ss).apply( + ["MATCH (n:SEC {id: 'N1'}) SET n.name = 'no'"], + mode="mutate", + ) + assert ei.value.code == "not_found" + assert ss.store.get("N1").fields["name"] == "CST" + + +def test_match_match_edge_is_parse_error_comma_form_works(memnet_temp): + ss = open_session(map_lines=_MAP) + MutateGate(ss).apply( + [ + "CREATE (:CST {id: 'A', name: 'a', role: 'r'})", + "CREATE (:CST {id: 'B', name: 'b', role: 'r'})", + ], + mode="add", + ) + with pytest.raises(MemNetError) as ei: + MutateGate(ss).apply( + ["MATCH (a {id: 'A'}) MATCH (b {id: 'B'}) CREATE (a)-[:knows]->(b)"], + mode="mutate", + allow_new_relation=True, + ) + assert ei.value.code == "parse_error" + assert "MATCH" in ei.value.message + assert not any(rec.tag == "EDG" for rec in ss.store._by_hid.values()) + MutateGate(ss).apply( + ["MATCH (a {id: 'A'}), (b {id: 'B'}) CREATE (a)-[:knows]->(b)"], + mode="mutate", + allow_new_relation=True, + ) + assert any(rec.tag == "EDG" for rec in ss.store._by_hid.values()) + + +def test_ttl_is_minutes_and_access_after_expiry_refuses(memnet_temp): + t0 = datetime(2026, 10, 8, 12, 0, tzinfo=UTC) + set_now_override(t0) + ss = open_session(map_lines=_MAP, ttl_minutes=60) + sid = ss.session_id + set_now_override(t0 + timedelta(seconds=130)) + got = get_session(sid) + assert got.session_id == sid + set_now_override(t0) + fresh = open_session(map_lines=_MAP, ttl_minutes=60) + fresh_id = fresh.session_id + set_now_override(t0 + timedelta(minutes=61)) + with pytest.raises(MemNetError) as ei: + get_session(fresh_id) + assert ei.value.code == "session_expired" + set_now_override(None) diff --git a/tests/test_housekeep_endpoints.py b/tests/test_housekeep_endpoints.py new file mode 100644 index 0000000..0fcb636 --- /dev/null +++ b/tests/test_housekeep_endpoints.py @@ -0,0 +1,163 @@ +"""MN-REQ-04.12 — edge endpoints are hid or nickname, not nickname alone.""" + +from __future__ import annotations + +from pathlib import Path + +from typer.testing import CliRunner + +from memnet.cli import app +from memnet.exceptions import MemNetError +from memnet.housekeep import orphan_rows, prune_rows, prune_stale, stats +from memnet.mutate_gate import MutateGate +from memnet.session import open_session +from memnet.snapshot import load_snapshot, write_snapshot + +runner = CliRunner() +_MAP = ["SCHEMA PRT ; fields=id name role"] + + +def _gql_chain(ss, n: int = 4) -> None: + creates = [f"CREATE (:PRT {{id: 'P{i}', name: 'part{i}', role: 'r'}})" for i in range(n)] + edges = [] + for i in range(n - 1): + edges.append( + "MATCH (a {id: 'P" + + str(i) + + "'}), (b {id: 'P" + + str(i + 1) + + "'}) CREATE (a)-[:links {id: 'E" + + str(i) + + "'}]->(b)" + ) + MutateGate(ss).apply(creates + edges, mode="mutate", allow_new_relation=True) + + +def _ids(ss) -> set[str]: + return {rec.fields.get("id", "") for rec in ss.store._by_hid.values() if rec.tag == "PRT"} + + +def test_gql_graph_has_no_orphans_and_survives_reload(memnet_temp, tmp_path: Path): + ss = open_session(map_lines=_MAP) + _gql_chain(ss, 4) + got = stats(ss) + assert got["edges"] == 3 + assert got["dangling"] == 0 + assert got["orphans"] == 0 + assert orphan_rows(ss) == [] + edge = next(rec for rec in ss.store._by_hid.values() if rec.tag == "EDG") + assert edge.fields["src"].startswith("_el") + assert edge.fields["dist"].startswith("_el") + path = tmp_path / "gql.snap" + write_snapshot(ss, path) + loaded = load_snapshot(path) + again = stats(loaded) + assert again["edges"] == 3 + assert again["dangling"] == 0 + assert again["orphans"] == 0 + + +def test_pipe_graph_keeps_nickname_endpoints_and_true_orphans(memnet_temp): + ss = open_session(map_lines=_MAP) + MutateGate(ss).apply( + [ + "@PRT: P1|a|r", + "@PRT: P2|b|r", + "@PRT: P3|alone|r", + "@EDG: E1|P1|links|P2|||", + ], + mode="add", + allow_new_relation=True, + ) + got = stats(ss) + assert got["edges"] == 1 + assert got["dangling"] == 0 + assert got["orphans"] == 1 + assert [rec.fields["id"] for rec in orphan_rows(ss)] == ["P3"] + edge = next(rec for rec in ss.store._by_hid.values() if rec.tag == "EDG") + edge.fields["src"] = "P1" + edge.fields["dist"] = "P2" + again = stats(ss) + assert again["dangling"] == 0 + assert again["orphans"] == 1 + assert [rec.fields["id"] for rec in orphan_rows(ss)] == ["P3"] + + +def test_half_missing_end_is_dangling_and_the_live_end_is_not_an_orphan(memnet_temp): + ss = open_session(map_lines=_MAP) + MutateGate(ss).apply( + ["CREATE (:PRT {id: 'P1', name: 'a', role: 'r'})"], + mode="add", + ) + p1 = ss.store.get("P1") + assert p1 is not None + MutateGate(ss).apply( + [f"@EDG: E_d|{p1.hid}|links|_el_missing|||"], + mode="add", + allow_new_relation=True, + ) + got = stats(ss) + assert got["dangling"] == 1 + assert got["orphans"] == 0 + + +def test_prune_apply_on_gql_graph_deletes_nothing_referenced(memnet_temp): + ss = open_session(map_lines=_MAP) + _gql_chain(ss, 4) + sid = ss.session_id + dry = runner.invoke(app, ["housekeep", "prune", "orphans", "--session", sid]) + assert dry.exit_code == 0, dry.stderr + assert "would-delete 0" in dry.stderr + assert "prune orphans --apply" not in dry.stderr + applied = runner.invoke( + app, + ["housekeep", "prune", "orphans", "--apply", "--session", sid], + ) + assert applied.exit_code == 0, applied.stderr + assert "deleted 0" in applied.stderr + assert _ids(ss) == {"P0", "P1", "P2", "P3"} + nodes = [rec for rec in ss.store._by_hid.values() if rec.tag == "PRT"] + try: + prune_rows(ss, nodes) + raise AssertionError("expected prune_referenced") + except MemNetError as exc: + assert exc.code == "prune_referenced" + assert _ids(ss) == {"P0", "P1", "P2", "P3"} + p0 = ss.store.get("P0") + assert p0 is not None + p0.fields["recycle"] = "delete_on_settle" + refused = runner.invoke( + app, + ["housekeep", "prune", "stale", "--apply", "--session", sid], + ) + assert refused.exit_code != 0 + assert "prune_referenced" in refused.stderr + assert ss.store.get("P0") is not None + try: + prune_stale(ss) + raise AssertionError("expected prune_stale to refuse") + except MemNetError as exc: + assert exc.code == "prune_referenced" + assert _ids(ss) == {"P0", "P1", "P2", "P3"} + + +def test_true_orphan_still_prunes_and_unreferenced_only(memnet_temp): + ss = open_session(map_lines=_MAP) + MutateGate(ss).apply( + [ + "CREATE (:PRT {id: 'P1', name: 'a', role: 'r'})", + "CREATE (:PRT {id: 'P2', name: 'b', role: 'r'})", + "CREATE (:PRT {id: 'LONE', name: 'z', role: 'r'})", + ], + mode="add", + ) + MutateGate(ss).apply( + ["MATCH (a {id: 'P1'}), (b {id: 'P2'}) CREATE (a)-[:links {id: 'E1'}]->(b)"], + mode="mutate", + allow_new_relation=True, + ) + assert [rec.fields["id"] for rec in orphan_rows(ss)] == ["LONE"] + deleted = prune_rows(ss, orphan_rows(ss)) + assert [rec.fields["id"] for rec in deleted] == ["LONE"] + assert ss.store.get("P1") is not None + assert ss.store.get("LONE") is None diff --git a/tests/test_serve_request_isolation.py b/tests/test_serve_request_isolation.py new file mode 100644 index 0000000..72cf1a2 --- /dev/null +++ b/tests/test_serve_request_isolation.py @@ -0,0 +1,352 @@ +"""MN-REQ-06.13 — concurrent serve commands do not share stdout or exit_code.""" + +from __future__ import annotations + +import hashlib +import json +import os +import socket +import subprocess +import sys +import threading +import time +import urllib.error +import urllib.request +from pathlib import Path + +from memnet import __version__ +from memnet.serve import probe, send_command + +ROOT = Path(__file__).resolve().parents[1] +N = 12 +MAP_LINE = "SCHEMA CST ; fields=id name role" + + +def _free_port() -> int: + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock: + sock.bind(("127.0.0.1", 0)) + return int(sock.getsockname()[1]) + + +def _sha(text: str) -> str: + return hashlib.sha256(text.encode()).hexdigest() + + +def _sid(blob: str) -> str: + for line in blob.splitlines(): + if line.startswith("@SESSION:"): + return line.split("|", 1)[0].replace("@SESSION:", "").strip() + raise AssertionError(blob) + + +def _start_serve(port: int, log_path: Path) -> subprocess.Popen[bytes]: + env = os.environ.copy() + for key in ( + "MEMNET_TEST_INLINE", + "MEMNET_SERVE_INTERNAL", + "MEMNET_SESSION", + "MEMNET_GATEWAY_CONFIG", + ): + env.pop(key, None) + env["MEMNET_SERVE_HOST"] = "127.0.0.1" + env["MEMNET_SERVE_PORT"] = str(port) + log = log_path.open("w", encoding="utf-8") + proc = subprocess.Popen( + [sys.executable, "-m", "memnet", "serve", "--host", "127.0.0.1", "--port", str(port)], + env=env, + cwd=str(ROOT), + stdout=log, + stderr=subprocess.STDOUT, + ) + log.close() + deadline = time.time() + 20 + while time.time() < deadline: + if proc.poll() is not None: + raise AssertionError(log_path.read_text(encoding="utf-8")) + if probe(host="127.0.0.1", port=port): + return proc + time.sleep(0.05) + proc.terminate() + raise AssertionError(log_path.read_text(encoding="utf-8")) + + +def _stop(proc: subprocess.Popen[bytes] | None) -> None: + if proc is None or proc.poll() is not None: + return + proc.terminate() + try: + proc.wait(timeout=5) + except subprocess.TimeoutExpired: + proc.kill() + proc.wait(timeout=5) + + +def _open(host: str, port: int) -> str: + resp = send_command( + ["session", "open", "--map", MAP_LINE], + host=host, + port=port, + ) + assert resp["exit_code"] == 0, resp + return _sid(resp.get("stdout") or "") + + +def _check(kind: str, marker: str, resp: dict, others: list[str]) -> None: + blob = (resp.get("stdout") or "") + "\n" + (resp.get("stderr") or "") + if kind == "ok": + assert resp["exit_code"] == 0, resp + assert marker in blob + assert "ok=1 fail=0" in (resp.get("stderr") or "") + assert "@ERR:" not in (resp.get("stderr") or "") + else: + assert resp["exit_code"] != 0, resp + assert "@ERR:" in (resp.get("stderr") or "") + assert "ok=1 fail=0" not in (resp.get("stderr") or "") + assert marker not in blob + for other in others: + if other != marker: + assert other not in blob, resp + + +def _burst(fns: list) -> tuple[list, float]: + barrier = threading.Barrier(len(fns)) + out: list = [None] * len(fns) + errors: list = [None] * len(fns) + + def run(i: int, fn) -> None: + barrier.wait(timeout=30) + try: + out[i] = fn() + except Exception as exc: # noqa: BLE001 — collected for the assertion + errors[i] = exc + + threads = [threading.Thread(target=run, args=(i, fn)) for i, fn in enumerate(fns)] + started = time.perf_counter() + for thread in threads: + thread.start() + for thread in threads: + thread.join(timeout=60) + elapsed = time.perf_counter() - started + for i, thread in enumerate(threads): + assert not thread.is_alive(), f"thread {i} still running" + assert errors[i] is None, errors[i] + return out, elapsed + + +def test_tcp_serve_isolates_concurrent_mutate(tmp_path: Path): + port = _free_port() + proc = _start_serve(port, tmp_path / "serve.log") + try: + sids = [_open("127.0.0.1", port) for _ in range(N)] + markers = [f"ISO_MARK_{i:02d}_ZXQ" for i in range(N)] + kinds = ["ok" if i % 2 == 0 else "fail" for i in range(N)] + + def call(i: int): + if kinds[i] == "ok": + stdin = ( + "CREATE (:CST {id: 'N" + + f"{i:02d}" + + "', name: '" + + markers[i] + + "', role: 'r'})\n" + ) + else: + stdin = "CREATE (:NOPE {id: 'BAD" + f"{i:02d}" + "', name: 'x'})\n" + return send_command( + ["mutate", "--stdin", "--session", sids[i]], + stdin=stdin, + host="127.0.0.1", + port=port, + ) + + results, _mixed_s = _burst([lambda i=i: call(i) for i in range(N)]) + for i, resp in enumerate(results): + own = markers[i] if kinds[i] == "ok" else f"ISO_FAIL_{i:02d}_ZXQ" + _check(kinds[i], own if kinds[i] == "ok" else "ISO_MARK_", resp, markers) + + time_sids = [_open("127.0.0.1", port) for _ in range(N)] + + def ok_call(i: int): + return send_command( + ["mutate", "--stdin", "--session", time_sids[i]], + stdin=( + "CREATE (:CST {id: 'T" + + f"{i:02d}" + + "', name: 'TIME" + + f"{i:02d}" + + "', role: 'r'})\n" + ), + host="127.0.0.1", + port=port, + ) + + _timed, parallel_s = _burst([lambda i=i: ok_call(i) for i in range(N)]) + for resp in _timed: + assert resp["exit_code"] == 0, resp + serial_sids = [_open("127.0.0.1", port) for _ in range(N)] + serial_started = time.perf_counter() + for i, sid in enumerate(serial_sids): + resp = send_command( + ["mutate", "--stdin", "--session", sid], + stdin=( + "CREATE (:CST {id: 'S" + + f"{i:02d}" + + "', name: 'SER" + + f"{i:02d}" + + "', role: 'r'})\n" + ), + host="127.0.0.1", + port=port, + ) + assert resp["exit_code"] == 0, resp + serial_s = time.perf_counter() - serial_started + note = ( + f"tcp parallel_s={parallel_s:.4f} serial_s={serial_s:.4f} " + f"ratio={parallel_s / serial_s if serial_s else 0:.3f}\n" + ) + path = Path("/tmp/memnet-isolation-throughput.txt") + path.write_text(note, encoding="utf-8") + art = Path("/opt/cursor/artifacts") + if art.is_dir(): + (art / "isolation-throughput.txt").write_text(note, encoding="utf-8") + finally: + _stop(proc) + + +def _post(url: str, token: str, payload: dict) -> dict: + data = json.dumps(payload).encode("utf-8") + req = urllib.request.Request( + url, + data=data, + method="POST", + headers={ + "Content-Type": "application/json", + "Authorization": f"Bearer {token}", + }, + ) + try: + with urllib.request.urlopen(req, timeout=30) as resp: + return json.loads(resp.read().decode("utf-8")) + except urllib.error.HTTPError as exc: + return json.loads(exc.read().decode("utf-8")) + + +def test_product_gateway_isolates_concurrent_mutate(tmp_path: Path): + serve_port = _free_port() + gw_port = _free_port() + serve = _start_serve(serve_port, tmp_path / "serve.log") + token = "endleaf-secret" + cfg = tmp_path / "gateway.json" + state = tmp_path / "owners.json" + cfg.write_text( + json.dumps( + { + "bind": "127.0.0.1", + "port": gw_port, + "path": "/gateway", + "body_max_bytes": 4194304, + "state_path": str(state), + "version_cache_s": 0, + "backend_timeout_s": 20, + "admin": {"sha256": _sha("admin-secret")}, + "backends": {"pi-endleaf": {"host": "127.0.0.1", "port": serve_port}}, + "products": { + "endleaf": { + "backends": ["pi-endleaf"], + "houses": {"syson": "pi-endleaf"}, + "pinned_version": __version__, + "credentials": [ + {"id": "endleaf-1", "sha256": _sha(token), "revoked": False} + ], + } + }, + } + ), + encoding="utf-8", + ) + env = os.environ.copy() + for key in ("MEMNET_TEST_INLINE", "MEMNET_SERVE_INTERNAL", "MEMNET_SESSION"): + env.pop(key, None) + env["MEMNET_GATEWAY_CONFIG"] = str(cfg) + log_path = tmp_path / "gateway.log" + log = log_path.open("w", encoding="utf-8") + gw = subprocess.Popen( + [ + sys.executable, + "-m", + "memnet_mcp.server", + "--transport", + "gateway", + "--host", + "127.0.0.1", + "--port", + str(gw_port), + ], + env=env, + cwd=str(ROOT), + stdout=log, + stderr=subprocess.STDOUT, + ) + log.close() + url = f"http://127.0.0.1:{gw_port}/gateway" + try: + deadline = time.time() + 20 + while time.time() < deadline: + if gw.poll() is not None: + raise AssertionError(log_path.read_text(encoding="utf-8")) + try: + ready = _post(url, token, {"args": ["version"], "namespace": "endleaf"}) + except (urllib.error.URLError, ConnectionError, TimeoutError): + time.sleep(0.05) + continue + if ready.get("exit_code") == 0 and "@VER:" in (ready.get("stdout") or ""): + break + time.sleep(0.05) + else: + raise AssertionError(log_path.read_text(encoding="utf-8")) + + def open_one() -> str: + body = _post( + url, + token, + { + "args": ["session", "open", "--map", MAP_LINE], + "namespace": "endleaf", + "backend": "pi-endleaf", + }, + ) + assert body["exit_code"] == 0, body + return _sid(body.get("stdout") or "") + + sids = [open_one() for _ in range(N)] + markers = [f"GW_MARK_{i:02d}_ZXQ" for i in range(N)] + kinds = ["ok" if i % 2 == 0 else "fail" for i in range(N)] + + def call(i: int) -> dict: + if kinds[i] == "ok": + stdin = ( + "CREATE (:CST {id: 'G" + + f"{i:02d}" + + "', name: '" + + markers[i] + + "', role: 'r'})\n" + ) + else: + stdin = "CREATE (:NOPE {id: 'BAD" + f"{i:02d}" + "', name: 'x'})\n" + return _post( + url, + token, + { + "args": ["mutate", "--stdin", "--session", sids[i]], + "stdin": stdin, + "session": sids[i], + "namespace": "endleaf", + }, + ) + + results, _elapsed = _burst([lambda i=i: call(i) for i in range(N)]) + for i, resp in enumerate(results): + _check(kinds[i], markers[i] if kinds[i] == "ok" else "GW_MARK_", resp, markers) + finally: + _stop(gw) + _stop(serve) diff --git a/tests/test_sysml_product_gateway.py b/tests/test_sysml_product_gateway.py index 7b78777..fd4245c 100644 --- a/tests/test_sysml_product_gateway.py +++ b/tests/test_sysml_product_gateway.py @@ -53,7 +53,7 @@ def test_requirement_verify_and_allocate(): def test_contract_is_not_draft_and_no_semver_bump(): contract = CONTRACT.read_text(encoding="utf-8") assert "MN-REQ-06.12" in contract - assert "implemented for the behaviours" in contract + assert "implemented and deployed" in contract lowered = contract.lower() assert "draft" not in lowered assert "100.118.79.40" in contract