From 24f59a6c532c84903a25b0b0096a20029d99d4c6 Mon Sep 17 00:00:00 2001 From: Patrick Buckley Date: Mon, 6 Apr 2026 03:08:19 -0700 Subject: [PATCH] =?UTF-8?q?feat:=20add=20per-node=20metadata=20with=20auto?= =?UTF-8?q?-collection,=20admin=20API,=20and=20cons=E2=80=A6=20(#318)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat: add per-node metadata with auto-collection, admin API, and console UI Adds a normalized node_metadata table for structured per-node key/value metadata with source tracking (auto/user/config). Auto-populated fields (hostname, OS, arch, interfaces, cpu_count) are collected at server startup via stdlib; user-defined fields are managed through the admin API, CLI, or config.toml [metadata] section. Storage: migration 035, 7 new protocol methods (get, get_all, set, set_bulk, delete, delete_by_source, filter), both SQLite and PostgreSQL backends. Filtering uses single-query GROUP BY/HAVING for efficiency. Console API: GET/PUT/DELETE endpoints under /admin/nodes/{node_id}/metadata with auto-source protection. cluster_nodes gains meta.* query param filtering; cluster_node_detail attaches metadata to responses. Frontend: new Nodes admin tab with collapsible per-node sections, inline add form, delete with confirmation. Read-only metadata panel in node detail drill-down. Proper design token usage, accessibility (ARIA, keyboard nav, screen reader labels), and mobile responsiveness. CLI: turnstone-admin list-node-metadata, set-node-metadata, and delete-node-metadata subcommands. 64 tests (25 storage, 19 node_info, 20 existing unaffected). * fix: resolve CI typecheck and test failures - Fix mypy error: use %-style format string instead of structlog kwargs for standard Logger.warning() in console server - Fix test_get_nodes assertion to include new node_ids=None parameter - Add debug logging to _collect_interfaces empty except block * fix: address Copilot review feedback on node metadata - Clear stale auto/config metadata before upserting on startup - Wrap metadata filter in try/except with graceful fallback - Add metadata field to NodeDetailResponse schema - Use _VALID_NODE_ID regex for consistent node_id validation - Defensive JSON decode in admin_get_node_metadata - Switch to read_json_or_400 and require_storage_or_503 helpers - Add SetNodeMetadataValueRequest for single-key PUT endpoint - Add bulk GET /admin/node-metadata endpoint (replaces N+1 fetches) - Update frontend to use single bulk metadata fetch * feat: add admin.nodes permission scope for node metadata - Add admin.nodes to builtin-admin role via migration 035 - Switch all node metadata handlers from admin.settings to admin.nodes - Register admin.nodes in the admin panel permission set - Node detail metadata panel fetches from cluster endpoint (no admin permission needed) instead of admin endpoint * fix: address second round of Copilot feedback - Replace inline onclick handlers with data-* attributes and event delegation to prevent JS string context XSS - Move NodeMetadataEntry before NodeDetailResponse and use it as the typed metadata field (was list[dict[str, Any]]) - Clean up config metadata on shutdown (was only cleaning auto) --- tests/test_console.py | 4 +- tests/test_node_info.py | 137 +++++++++ tests/test_node_metadata_storage.py | 185 ++++++++++++ turnstone/admin.py | 84 ++++++ turnstone/api/console_schemas.py | 34 +++ turnstone/api/console_spec.py | 41 +++ turnstone/console/collector.py | 9 +- turnstone/console/server.py | 264 +++++++++++++++++- turnstone/console/static/admin.js | 234 ++++++++++++++++ turnstone/console/static/app.js | 58 ++++ turnstone/console/static/index.html | 17 ++ turnstone/console/static/style.css | 85 ++++++ turnstone/core/node_info.py | 69 +++++ turnstone/core/storage/_postgresql.py | 114 ++++++++ turnstone/core/storage/_protocol.py | 30 ++ turnstone/core/storage/_schema.py | 18 ++ turnstone/core/storage/_sqlite.py | 114 ++++++++ .../migrations/versions/035_node_metadata.py | 50 ++++ turnstone/server.py | 33 ++- 19 files changed, 1571 insertions(+), 9 deletions(-) create mode 100644 tests/test_node_info.py create mode 100644 tests/test_node_metadata_storage.py create mode 100644 turnstone/core/node_info.py create mode 100644 turnstone/core/storage/migrations/versions/035_node_metadata.py diff --git a/tests/test_console.py b/tests/test_console.py index e9a0bea7..51bb5c4e 100644 --- a/tests/test_console.py +++ b/tests/test_console.py @@ -757,7 +757,9 @@ class TestConsoleHTTPEndpoints: assert status == 200 assert len(data["nodes"]) == 1 assert data["total"] == 1 - mock_collector.get_nodes.assert_called_once_with(sort_by="activity", limit=10, offset=0) + mock_collector.get_nodes.assert_called_once_with( + sort_by="activity", limit=10, offset=0, node_ids=None + ) def test_get_workstreams(self, client, mock_collector): status, data = self._get( diff --git a/tests/test_node_info.py b/tests/test_node_info.py new file mode 100644 index 00000000..c2b7f845 --- /dev/null +++ b/tests/test_node_info.py @@ -0,0 +1,137 @@ +"""Tests for auto-populated node metadata collection.""" + +from __future__ import annotations + +import json +from unittest.mock import patch + +from turnstone.core.node_info import ( + _collect_interfaces, + _is_loopback_or_link_local, + collect_node_info, +) + + +class TestCollectNodeInfo: + def test_returns_dict(self): + info = collect_node_info() + assert isinstance(info, dict) + + def test_expected_keys_present(self): + info = collect_node_info() + # These should always be available on any platform + assert "hostname" in info + assert "os" in info + assert "arch" in info + assert "python" in info + + def test_values_json_serializable(self): + info = collect_node_info() + for _key, value in info.items(): + serialized = json.dumps(value) + assert isinstance(serialized, str) + + def test_hostname_is_string(self): + info = collect_node_info() + assert isinstance(info["hostname"], str) + assert len(info["hostname"]) > 0 + + def test_cpu_count_is_int(self): + info = collect_node_info() + if "cpu_count" in info: + assert isinstance(info["cpu_count"], int) + assert info["cpu_count"] > 0 + + def test_interfaces_is_dict(self): + info = collect_node_info() + if "interfaces" in info: + assert isinstance(info["interfaces"], dict) + for iface, ips in info["interfaces"].items(): + assert isinstance(iface, str) + assert isinstance(ips, list) + + def test_one_field_failure_does_not_block_others(self): + """Individual field failures must not prevent other fields from collecting.""" + with patch("turnstone.core.node_info.socket.gethostname", side_effect=OSError("boom")): + info = collect_node_info() + assert "hostname" not in info + # Other fields should still be present + assert "os" in info + assert "arch" in info + assert "python" in info + + def test_none_value_excluded(self): + with patch("turnstone.core.node_info.os.cpu_count", return_value=None): + info = collect_node_info() + assert "cpu_count" not in info + assert "hostname" in info + + def test_interface_failure_does_not_block_fields(self): + """Interface collection failure must not prevent scalar fields.""" + with patch( + "turnstone.core.node_info._collect_interfaces", + side_effect=RuntimeError("boom"), + ): + info = collect_node_info() + assert "interfaces" not in info + assert "hostname" in info + assert "os" in info + + +class TestCollectInterfaces: + def test_returns_dict(self): + result = _collect_interfaces() + assert isinstance(result, dict) + + def test_values_are_string_lists(self): + result = _collect_interfaces() + for label, ips in result.items(): + assert isinstance(label, str) + assert isinstance(ips, list) + for ip in ips: + assert isinstance(ip, str) + + def test_no_loopback_in_results(self): + result = _collect_interfaces() + for _label, ips in result.items(): + for ip in ips: + assert not ip.startswith("127.") + assert ip != "::1" + assert not ip.startswith("fe80:") + + def test_getaddrinfo_oserror_returns_empty(self): + with patch( + "turnstone.core.node_info.socket.getaddrinfo", + side_effect=OSError("no network"), + ): + result = _collect_interfaces() + assert result == {} + + def test_all_loopback_returns_empty(self): + import socket + + mock_addrs = [ + (socket.AF_INET, socket.SOCK_STREAM, 6, "", ("127.0.0.1", 0)), + (socket.AF_INET6, socket.SOCK_STREAM, 6, "", ("::1", 0, 0, 0)), + ] + with patch("turnstone.core.node_info.socket.getaddrinfo", return_value=mock_addrs): + result = _collect_interfaces() + assert result == {} + + +class TestIsLoopbackOrLinkLocal: + def test_ipv4_loopback(self): + assert _is_loopback_or_link_local("127.0.0.1") is True + assert _is_loopback_or_link_local("127.0.1.1") is True + + def test_ipv6_loopback(self): + assert _is_loopback_or_link_local("::1") is True + + def test_link_local(self): + assert _is_loopback_or_link_local("fe80::1") is True + assert _is_loopback_or_link_local("fe80:abc::def") is True + + def test_normal_addresses(self): + assert _is_loopback_or_link_local("10.0.0.5") is False + assert _is_loopback_or_link_local("192.168.1.1") is False + assert _is_loopback_or_link_local("2001:db8::1") is False diff --git a/tests/test_node_metadata_storage.py b/tests/test_node_metadata_storage.py new file mode 100644 index 00000000..7f1dce64 --- /dev/null +++ b/tests/test_node_metadata_storage.py @@ -0,0 +1,185 @@ +"""Tests for node_metadata storage methods.""" + +from __future__ import annotations + +import json + + +class TestNodeMetadata: + def test_set_and_get(self, storage): + storage.set_node_metadata("node-1", "rack", json.dumps("us-east-1a")) + rows = storage.get_node_metadata("node-1") + assert len(rows) == 1 + assert rows[0]["key"] == "rack" + assert json.loads(rows[0]["value"]) == "us-east-1a" + assert rows[0]["source"] == "user" + + def test_set_with_source(self, storage): + storage.set_node_metadata("node-1", "hostname", json.dumps("web-01"), source="auto") + rows = storage.get_node_metadata("node-1") + assert rows[0]["source"] == "auto" + + def test_upsert_overwrites(self, storage): + storage.set_node_metadata("node-1", "rack", json.dumps("old")) + storage.set_node_metadata("node-1", "rack", json.dumps("new")) + rows = storage.get_node_metadata("node-1") + assert len(rows) == 1 + assert json.loads(rows[0]["value"]) == "new" + + def test_complex_value(self, storage): + val = {"model": "A100", "count": 4} + storage.set_node_metadata("node-1", "gpu", json.dumps(val)) + rows = storage.get_node_metadata("node-1") + assert json.loads(rows[0]["value"]) == val + + def test_list_value(self, storage): + val = ["inference", "eval"] + storage.set_node_metadata("node-1", "roles", json.dumps(val)) + rows = storage.get_node_metadata("node-1") + assert json.loads(rows[0]["value"]) == val + + def test_get_empty(self, storage): + rows = storage.get_node_metadata("nonexistent") + assert rows == [] + + def test_get_all_node_metadata(self, storage): + storage.set_node_metadata("node-1", "rack", json.dumps("a")) + storage.set_node_metadata("node-2", "rack", json.dumps("b")) + storage.set_node_metadata("node-2", "os", json.dumps("Linux")) + result = storage.get_all_node_metadata() + assert "node-1" in result + assert "node-2" in result + assert len(result["node-1"]) == 1 + assert len(result["node-2"]) == 2 + node2_keys = {r["key"] for r in result["node-2"]} + assert node2_keys == {"rack", "os"} + + def test_get_all_empty(self, storage): + result = storage.get_all_node_metadata() + assert result == {} + + def test_bulk_set(self, storage): + entries = [ + ("hostname", json.dumps("web-01"), "auto"), + ("os", json.dumps("Linux"), "auto"), + ("rack", json.dumps("us-east-1a"), "config"), + ] + storage.set_node_metadata_bulk("node-1", entries) + rows = storage.get_node_metadata("node-1") + assert len(rows) == 3 + keys = {r["key"] for r in rows} + assert keys == {"hostname", "os", "rack"} + + def test_bulk_set_upsert(self, storage): + storage.set_node_metadata("node-1", "rack", json.dumps("old"), source="config") + entries = [("rack", json.dumps("new"), "config")] + storage.set_node_metadata_bulk("node-1", entries) + rows = storage.get_node_metadata("node-1") + assert len(rows) == 1 + assert json.loads(rows[0]["value"]) == "new" + + def test_delete(self, storage): + storage.set_node_metadata("node-1", "rack", json.dumps("a")) + deleted = storage.delete_node_metadata("node-1", "rack") + assert deleted is True + assert storage.get_node_metadata("node-1") == [] + + def test_delete_nonexistent(self, storage): + deleted = storage.delete_node_metadata("node-1", "nope") + assert deleted is False + + def test_delete_by_source(self, storage): + storage.set_node_metadata("node-1", "hostname", json.dumps("h"), source="auto") + storage.set_node_metadata("node-1", "os", json.dumps("Linux"), source="auto") + storage.set_node_metadata("node-1", "rack", json.dumps("a"), source="user") + count = storage.delete_node_metadata_by_source("node-1", "auto") + assert count == 2 + rows = storage.get_node_metadata("node-1") + assert len(rows) == 1 + assert rows[0]["key"] == "rack" + + def test_delete_by_source_empty(self, storage): + count = storage.delete_node_metadata_by_source("node-1", "auto") + assert count == 0 + + def test_filter_single_key(self, storage): + storage.set_node_metadata("node-1", "rack", json.dumps("us-east-1a")) + storage.set_node_metadata("node-2", "rack", json.dumps("us-west-2a")) + result = storage.filter_nodes_by_metadata({"rack": json.dumps("us-east-1a")}) + assert result == {"node-1"} + + def test_filter_multiple_keys(self, storage): + storage.set_node_metadata("node-1", "rack", json.dumps("a")) + storage.set_node_metadata("node-1", "os", json.dumps("Linux")) + storage.set_node_metadata("node-2", "rack", json.dumps("a")) + storage.set_node_metadata("node-2", "os", json.dumps("Windows")) + result = storage.filter_nodes_by_metadata( + { + "rack": json.dumps("a"), + "os": json.dumps("Linux"), + } + ) + assert result == {"node-1"} + + def test_filter_no_match(self, storage): + storage.set_node_metadata("node-1", "rack", json.dumps("a")) + result = storage.filter_nodes_by_metadata({"rack": json.dumps("z")}) + assert result == set() + + def test_filter_empty_filters(self, storage): + result = storage.filter_nodes_by_metadata({}) + assert result == set() + + def test_filter_partial_intersection_eliminates_all(self, storage): + """First filter matches 2 nodes, second filter matches neither.""" + storage.set_node_metadata("node-1", "rack", json.dumps("a")) + storage.set_node_metadata("node-2", "rack", json.dumps("a")) + storage.set_node_metadata("node-1", "os", json.dumps("Linux")) + storage.set_node_metadata("node-2", "os", json.dumps("Linux")) + result = storage.filter_nodes_by_metadata( + {"rack": json.dumps("a"), "region": json.dumps("eu")} + ) + assert result == set() + + def test_upsert_preserves_created(self, storage): + storage.set_node_metadata("node-1", "rack", json.dumps("old")) + rows = storage.get_node_metadata("node-1") + first_created = rows[0]["created"] + + storage.set_node_metadata("node-1", "rack", json.dumps("new")) + rows = storage.get_node_metadata("node-1") + assert rows[0]["created"] == first_created + assert json.loads(rows[0]["value"]) == "new" + + def test_bulk_set_empty_list(self, storage): + storage.set_node_metadata_bulk("node-1", []) + rows = storage.get_node_metadata("node-1") + assert rows == [] + + def test_ordered_by_key(self, storage): + storage.set_node_metadata("node-1", "zz", json.dumps("last")) + storage.set_node_metadata("node-1", "aa", json.dumps("first")) + rows = storage.get_node_metadata("node-1") + assert rows[0]["key"] == "aa" + assert rows[1]["key"] == "zz" + + def test_upsert_changes_source(self, storage): + storage.set_node_metadata("node-1", "rack", json.dumps("a"), source="auto") + storage.set_node_metadata("node-1", "rack", json.dumps("a"), source="user") + rows = storage.get_node_metadata("node-1") + assert rows[0]["source"] == "user" + + def test_delete_by_source_does_not_affect_other_nodes(self, storage): + storage.set_node_metadata("node-1", "hostname", json.dumps("h1"), source="auto") + storage.set_node_metadata("node-2", "hostname", json.dumps("h2"), source="auto") + storage.delete_node_metadata_by_source("node-1", "auto") + rows = storage.get_node_metadata("node-2") + assert len(rows) == 1 + assert rows[0]["key"] == "hostname" + + def test_filter_returns_multiple_matches(self, storage): + storage.set_node_metadata("node-1", "rack", json.dumps("a")) + storage.set_node_metadata("node-2", "rack", json.dumps("a")) + storage.set_node_metadata("node-3", "rack", json.dumps("b")) + result = storage.filter_nodes_by_metadata({"rack": json.dumps("a")}) + assert result == {"node-1", "node-2"} diff --git a/turnstone/admin.py b/turnstone/admin.py index 57a575a7..7bceb18e 100644 --- a/turnstone/admin.py +++ b/turnstone/admin.py @@ -293,6 +293,74 @@ def _cmd_tls_list(args: argparse.Namespace) -> None: print(f"{c['domain']:<30s} {c['issued_at']:<22s} {c['expires_at']:<22s}") +def _cmd_list_node_metadata(args: argparse.Namespace) -> None: + """List metadata for a node.""" + import json + + storage = _get_storage() + rows = storage.get_node_metadata(args.node_id) + if not rows: + print(f"No metadata for node: {args.node_id}") + return + + print(f"{'KEY':<20s} {'VALUE':<40s} {'SOURCE':<8s} {'UPDATED':<20s}") + print("-" * 88) + for r in rows: + val = r["value"] + try: + parsed = json.loads(val) + val_str = json.dumps(parsed) if isinstance(parsed, (dict, list)) else str(parsed) + except (json.JSONDecodeError, TypeError): + val_str = val + if len(val_str) > 38: + val_str = val_str[:35] + "..." + key_str = r["key"] + if len(key_str) > 18: + key_str = key_str[:15] + "..." + print(f"{key_str:<20s} {val_str:<40s} {r['source']:<8s} {r['updated']:<20s}") + + +def _cmd_set_node_metadata(args: argparse.Namespace) -> None: + """Set a metadata key on a node.""" + import json + + storage = _get_storage() + + # Check for auto-source conflict + existing = storage.get_node_metadata(args.node_id) + for r in existing: + if r["key"] == args.key and r["source"] == "auto": + print(f"Error: cannot overwrite auto-populated key: {args.key}", file=sys.stderr) + sys.exit(1) + + # Try JSON parse, fall back to string + try: + value = json.loads(args.value) + except (json.JSONDecodeError, TypeError): + value = args.value + + storage.set_node_metadata(args.node_id, args.key, json.dumps(value), source="user") + print(f"Set {args.key}={json.dumps(value)} on {args.node_id}") + + +def _cmd_delete_node_metadata(args: argparse.Namespace) -> None: + """Delete a metadata key from a node.""" + storage = _get_storage() + + existing = storage.get_node_metadata(args.node_id) + for r in existing: + if r["key"] == args.key and r["source"] == "auto": + print(f"Error: cannot delete auto-populated key: {args.key}", file=sys.stderr) + sys.exit(1) + + deleted = storage.delete_node_metadata(args.node_id, args.key) + if deleted: + print(f"Deleted {args.key} from {args.node_id}") + else: + print(f"Key not found: {args.key} on {args.node_id}", file=sys.stderr) + sys.exit(1) + + def _discover_console_url() -> str: """Discover console URL from the services table.""" from turnstone.core.storage import get_storage @@ -378,6 +446,19 @@ def main() -> None: p_tlslist = sub.add_parser("tls-list", help="List issued certificates") p_tlslist.add_argument("--console-url", default="", help="Console URL") + # Node metadata commands + p_lnm = sub.add_parser("list-node-metadata", help="List metadata for a node") + p_lnm.add_argument("node_id", help="Node ID") + + p_snm = sub.add_parser("set-node-metadata", help="Set a metadata key on a node") + p_snm.add_argument("node_id", help="Node ID") + p_snm.add_argument("key", help="Metadata key") + p_snm.add_argument("value", help="Value (JSON or plain string)") + + p_dnm = sub.add_parser("delete-node-metadata", help="Delete a metadata key from a node") + p_dnm.add_argument("node_id", help="Node ID") + p_dnm.add_argument("key", help="Metadata key") + args = parser.parse_args() if not args.command: parser.print_help() @@ -393,5 +474,8 @@ def main() -> None: "tls-issue": _cmd_tls_issue, "tls-ca-cert": _cmd_tls_ca_cert, "tls-list": _cmd_tls_list, + "list-node-metadata": _cmd_list_node_metadata, + "set-node-metadata": _cmd_set_node_metadata, + "delete-node-metadata": _cmd_delete_node_metadata, } dispatch[args.command](args) diff --git a/turnstone/api/console_schemas.py b/turnstone/api/console_schemas.py index 8f6bc803..621db941 100644 --- a/turnstone/api/console_schemas.py +++ b/turnstone/api/console_schemas.py @@ -90,6 +90,12 @@ class ClusterWorkstreamsResponse(BaseModel): # --------------------------------------------------------------------------- +class NodeMetadataEntry(BaseModel): + key: str + value: Any + source: str = "user" + + class NodeDetailResponse(BaseModel): node_id: str server_url: str = "" @@ -97,6 +103,7 @@ class NodeDetailResponse(BaseModel): workstreams: list[ClusterWorkstreamInfo] = [] aggregate: dict[str, int] = Field(default_factory=dict) reachable: bool = True + metadata: list[NodeMetadataEntry] = Field(default_factory=list) # --------------------------------------------------------------------------- @@ -898,3 +905,30 @@ class RouteCreateResponse(BaseModel): ws_id: str = "" node_url: str = "" node_id: str = "" + + +# --------------------------------------------------------------------------- +# Node metadata +# --------------------------------------------------------------------------- + + +class NodeMetadataResponse(BaseModel): + node_id: str + metadata: list[NodeMetadataEntry] = Field(default_factory=list) + + +class SetNodeMetadataValueRequest(BaseModel): + """Request body for PUT /admin/nodes/{node_id}/metadata/{key}.""" + + value: Any + + +class SetNodeMetadataRequest(BaseModel): + """Single entry in a bulk metadata set.""" + + key: str + value: Any + + +class BulkSetNodeMetadataRequest(BaseModel): + entries: list[SetNodeMetadataRequest] = Field(default_factory=list) diff --git a/turnstone/api/console_spec.py b/turnstone/api/console_spec.py index b24627f1..3e92e3ff 100644 --- a/turnstone/api/console_spec.py +++ b/turnstone/api/console_spec.py @@ -12,6 +12,7 @@ from turnstone.api.console_schemas import ( AssignRoleRequest, AuditEventInfo, AvailableModelInfo, + BulkSetNodeMetadataRequest, ChannelUserInfo, ClusterNodesResponse, ClusterOverviewResponse, @@ -55,6 +56,7 @@ from turnstone.api.console_schemas import ( ModelDefinitionInfo, ModelReloadResponse, NodeDetailResponse, + NodeMetadataResponse, OrgInfo, OutputAssessmentInfo, RegistryInstallRequest, @@ -62,6 +64,7 @@ from turnstone.api.console_schemas import ( RoleInfo, RouteCreateResponse, RouteResponse, + SetNodeMetadataValueRequest, SettingInfo, SettingSchemaInfo, SkillDiscoverResponse, @@ -977,6 +980,44 @@ CONSOLE_ENDPOINTS: list[EndpointSpec] = [ error_codes=[404], tags=["Admin"], ), + # --- Admin: Node metadata --- + EndpointSpec( + "/v1/api/admin/node-metadata", + "GET", + "Get metadata for all nodes (bulk)", + tags=["Admin"], + ), + EndpointSpec( + "/v1/api/admin/nodes/{node_id}/metadata", + "GET", + "Get all metadata for a node", + response_model=NodeMetadataResponse, + error_codes=[400], + tags=["Admin"], + ), + EndpointSpec( + "/v1/api/admin/nodes/{node_id}/metadata", + "PUT", + "Bulk set user metadata for a node", + request_model=BulkSetNodeMetadataRequest, + error_codes=[400], + tags=["Admin"], + ), + EndpointSpec( + "/v1/api/admin/nodes/{node_id}/metadata/{key}", + "PUT", + "Set a single metadata key for a node", + request_model=SetNodeMetadataValueRequest, + error_codes=[400], + tags=["Admin"], + ), + EndpointSpec( + "/v1/api/admin/nodes/{node_id}/metadata/{key}", + "DELETE", + "Delete a single metadata key for a node", + error_codes=[400, 404], + tags=["Admin"], + ), # --- Admin: TLS / ACME --- EndpointSpec( "/v1/api/admin/tls/ca", diff --git a/turnstone/console/collector.py b/turnstone/console/collector.py index e4729b30..fa110b25 100644 --- a/turnstone/console/collector.py +++ b/turnstone/console/collector.py @@ -609,15 +609,22 @@ class ClusterCollector: } def get_nodes( - self, sort_by: str = "activity", limit: int | None = 100, offset: int = 0 + self, + sort_by: str = "activity", + limit: int | None = 100, + offset: int = 0, + node_ids: set[str] | None = None, ) -> tuple[list[dict[str, Any]], int]: """Return sorted, paginated node list with per-node counts. Pass ``limit=None`` to return all nodes (no pagination). + Pass ``node_ids`` to restrict results to the given set. """ with self._lock: items = [] for node in self._nodes.values(): + if node_ids is not None and node.node_id not in node_ids: + continue ws_states = { "running": 0, "thinking": 0, diff --git a/turnstone/console/server.py b/turnstone/console/server.py index e40c6195..c1bc84f1 100644 --- a/turnstone/console/server.py +++ b/turnstone/console/server.py @@ -250,7 +250,35 @@ async def cluster_nodes(request: Request) -> JSONResponse: sort_by = params.get("sort", "activity") limit = _parse_int(params, "limit", 100, minimum=1, maximum=1000) offset = _parse_int(params, "offset", 0) - nodes, total = collector.get_nodes(sort_by=sort_by, limit=limit, offset=offset) + + # Extract meta.* filters for node metadata filtering + meta_filters = {k[5:]: v for k, v in params.items() if k.startswith("meta.") and k[5:]} + node_ids: set[str] | None = None + if meta_filters: + import json as _mf_json + + storage = getattr(request.app.state, "auth_storage", None) + if storage is not None: + # Values in the DB are JSON-encoded. Try to use raw value if it is + # already valid JSON (e.g. meta.cpu_count=4), otherwise wrap as string. + encoded = {} + for mk, mv in meta_filters.items(): + try: + _mf_json.loads(mv) + encoded[mk] = mv + except (ValueError, TypeError): + encoded[mk] = _mf_json.dumps(mv) + try: + node_ids = storage.filter_nodes_by_metadata(encoded) + except Exception: + log.warning("cluster.metadata_filter_failed", exc_info=True) + node_ids = None # fall back to unfiltered + if node_ids is not None and not node_ids: + return JSONResponse({"nodes": [], "total": 0}) + + nodes, total = collector.get_nodes( + sort_by=sort_by, limit=limit, offset=offset, node_ids=node_ids + ) return JSONResponse({"nodes": nodes, "total": total}) @@ -286,12 +314,34 @@ async def cluster_workstreams(request: Request) -> JSONResponse: async def cluster_node_detail(request: Request) -> JSONResponse: collector: ClusterCollector = request.app.state.collector node_id = request.path_params["node_id"] - if not node_id or "/" in node_id or len(node_id) > 256: - return JSONResponse({"error": "Invalid node ID"}, status_code=400) + nv = _validate_node_id(node_id) + if nv: + return nv detail = collector.get_node_detail(node_id) - if detail: - return JSONResponse(detail) - return JSONResponse({"error": "Node not found"}, status_code=404) + if not detail: + return JSONResponse({"error": "Node not found"}, status_code=404) + + # Attach metadata if available + import json as _nd_json + + storage = getattr(request.app.state, "auth_storage", None) + if storage is not None: + try: + raw = storage.get_node_metadata(node_id) + entries = [] + for r in raw: + try: + val = _nd_json.loads(r["value"]) + except (ValueError, TypeError): + val = r["value"] + entries.append({"key": r["key"], "value": val, "source": r["source"]}) + detail["metadata"] = entries + except Exception: + log.warning("cluster.node_metadata_load_failed node_id=%s", node_id, exc_info=True) + detail["metadata"] = [] + else: + detail["metadata"] = [] + return JSONResponse(detail) async def cluster_snapshot(request: Request) -> JSONResponse: @@ -2244,6 +2294,7 @@ _VALID_PERMISSIONS = frozenset( "admin.watches", "admin.judge", "admin.memories", + "admin.nodes", "admin.settings", "admin.mcp", "admin.models", @@ -6917,6 +6968,189 @@ async def admin_validate_regex(request: Request) -> JSONResponse: return JSONResponse({"valid": True}) +def _validate_node_id(node_id: str) -> JSONResponse | None: + """Return an error response if node_id is invalid, else None.""" + if not node_id or len(node_id) > 256 or not _VALID_NODE_ID.match(node_id): + return JSONResponse({"error": "Invalid node ID"}, status_code=400) + return None + + +async def admin_get_all_node_metadata(request: Request) -> JSONResponse: + """GET /v1/api/admin/node-metadata — metadata for all nodes.""" + import json as _anm_json + + from turnstone.core.auth import require_permission + from turnstone.core.web_helpers import require_storage_or_503 + + err = require_permission(request, "admin.nodes") + if err: + return err + storage, serr = require_storage_or_503(request) + if serr: + return serr + all_meta = storage.get_all_node_metadata() + result: dict[str, list[dict[str, Any]]] = {} + for nid, rows in all_meta.items(): + entries = [] + for r in rows: + try: + val = _anm_json.loads(r["value"]) + except (ValueError, TypeError): + val = r["value"] + entries.append({"key": r["key"], "value": val, "source": r["source"]}) + result[nid] = entries + return JSONResponse({"nodes": result}) + + +async def admin_get_node_metadata(request: Request) -> JSONResponse: + """GET /v1/api/admin/nodes/{node_id}/metadata — all metadata for a node.""" + import json as _nm_json + + from turnstone.core.auth import require_permission + from turnstone.core.web_helpers import require_storage_or_503 + + err = require_permission(request, "admin.nodes") + if err: + return err + node_id = request.path_params["node_id"] + nv = _validate_node_id(node_id) + if nv: + return nv + storage, serr = require_storage_or_503(request) + if serr: + return serr + rows = storage.get_node_metadata(node_id) + metadata = [] + for r in rows: + try: + val = _nm_json.loads(r["value"]) + except (ValueError, TypeError): + val = r["value"] + metadata.append({"key": r["key"], "value": val, "source": r["source"]}) + return JSONResponse({"node_id": node_id, "metadata": metadata}) + + +async def admin_set_node_metadata(request: Request) -> JSONResponse: + """PUT /v1/api/admin/nodes/{node_id}/metadata — bulk set user metadata.""" + import json as _nm_json + + from turnstone.core.auth import require_permission + from turnstone.core.web_helpers import read_json_or_400, require_storage_or_503 + + err = require_permission(request, "admin.nodes") + if err: + return err + node_id = request.path_params["node_id"] + nv = _validate_node_id(node_id) + if nv: + return nv + body = await read_json_or_400(request) + if isinstance(body, JSONResponse): + return body + entries = body.get("entries", []) + if not entries: + return JSONResponse({"error": "No entries provided"}, status_code=400) + + storage, serr = require_storage_or_503(request) + if serr: + return serr + + # Validate entries + existing = {r["key"]: r["source"] for r in storage.get_node_metadata(node_id)} + for e in entries: + key = e.get("key", "") + if not key: + return JSONResponse({"error": "Empty key"}, status_code=400) + if len(key) > 128: + return JSONResponse( + {"error": f"Key too long (max 128): {key[:32]}..."}, status_code=400 + ) + if "value" not in e: + return JSONResponse({"error": f"Missing value for key: {key}"}, status_code=400) + if existing.get(key) == "auto": + return JSONResponse( + {"error": f"Cannot overwrite auto-populated key: {key}"}, + status_code=400, + ) + + bulk = [(e["key"], _nm_json.dumps(e["value"]), "user") for e in entries] + storage.set_node_metadata_bulk(node_id, bulk) + return JSONResponse({"ok": True, "count": len(bulk)}) + + +async def admin_set_node_metadata_key(request: Request) -> JSONResponse: + """PUT /v1/api/admin/nodes/{node_id}/metadata/{key} — set single key.""" + import json as _nm_json + + from turnstone.core.auth import require_permission + from turnstone.core.web_helpers import read_json_or_400, require_storage_or_503 + + err = require_permission(request, "admin.nodes") + if err: + return err + node_id = request.path_params["node_id"] + nv = _validate_node_id(node_id) + if nv: + return nv + key = request.path_params["key"] + if not key: + return JSONResponse({"error": "Empty key"}, status_code=400) + if len(key) > 128: + return JSONResponse({"error": "Key too long (max 128)"}, status_code=400) + + storage, serr = require_storage_or_503(request) + if serr: + return serr + existing = storage.get_node_metadata(node_id) + for r in existing: + if r["key"] == key and r["source"] == "auto": + return JSONResponse( + {"error": f"Cannot overwrite auto-populated key: {key}"}, + status_code=400, + ) + + body = await read_json_or_400(request) + if isinstance(body, JSONResponse): + return body + if "value" not in body: + return JSONResponse({"error": "Missing value"}, status_code=400) + storage.set_node_metadata(node_id, key, _nm_json.dumps(body["value"]), source="user") + return JSONResponse({"ok": True}) + + +async def admin_delete_node_metadata_key(request: Request) -> JSONResponse: + """DELETE /v1/api/admin/nodes/{node_id}/metadata/{key} — delete single key.""" + from turnstone.core.auth import require_permission + from turnstone.core.web_helpers import require_storage_or_503 + + err = require_permission(request, "admin.nodes") + if err: + return err + node_id = request.path_params["node_id"] + nv = _validate_node_id(node_id) + if nv: + return nv + key = request.path_params["key"] + if not key: + return JSONResponse({"error": "Empty key"}, status_code=400) + + storage, serr = require_storage_or_503(request) + if serr: + return serr + existing = storage.get_node_metadata(node_id) + for r in existing: + if r["key"] == key and r["source"] == "auto": + return JSONResponse( + {"error": f"Cannot delete auto-populated key: {key}"}, + status_code=400, + ) + + deleted = storage.delete_node_metadata(node_id, key) + if not deleted: + return JSONResponse({"error": "Key not found"}, status_code=404) + return JSONResponse({"ok": True}) + + async def admin_ring_status(request: Request) -> JSONResponse: """GET /v1/api/admin/ring/status — hash ring rebalancer status.""" from turnstone.core.auth import require_permission @@ -7449,6 +7683,24 @@ def create_app( admin_rescan_skill, methods=["POST"], ), + # Node metadata + Route("/api/admin/node-metadata", admin_get_all_node_metadata), + Route( + "/api/admin/nodes/{node_id}/metadata/{key}", + admin_set_node_metadata_key, + methods=["PUT"], + ), + Route( + "/api/admin/nodes/{node_id}/metadata/{key}", + admin_delete_node_metadata_key, + methods=["DELETE"], + ), + Route("/api/admin/nodes/{node_id}/metadata", admin_get_node_metadata), + Route( + "/api/admin/nodes/{node_id}/metadata", + admin_set_node_metadata, + methods=["PUT"], + ), # Hash ring Route("/api/admin/ring/status", admin_ring_status), Route( diff --git a/turnstone/console/static/admin.js b/turnstone/console/static/admin.js index d38c5a2c..f7bb2752 100644 --- a/turnstone/console/static/admin.js +++ b/turnstone/console/static/admin.js @@ -239,6 +239,7 @@ function switchAdminTab(tab) { "audit", "memories", "models", + "node-metadata", "settings", "tls", "mcp", @@ -265,6 +266,7 @@ function switchAdminTab(tab) { } if (tab === "memories") loadAdminMemories(); if (tab === "models") loadAdminModels(); + if (tab === "node-metadata") loadAdminNodeMetadata(); if (tab === "settings") loadSettings(); if (tab === "tls") loadTlsCerts(); if (tab === "mcp") loadAdminMcp(); @@ -5006,3 +5008,235 @@ function reloadModelNodes() { btn.textContent = "Sync to Nodes"; }); } + +// --------------------------------------------------------------------------- +// Node Metadata tab +// --------------------------------------------------------------------------- + +var _nodeMetaCache = {}; + +function loadAdminNodeMetadata() { + var container = document.getElementById("admin-node-metadata-content"); + if (!container) return; + container.innerHTML = '
Loading\u2026
'; + + // Single bulk fetch for all node metadata + authFetch("/v1/api/admin/node-metadata") + .then(function (r) { + if (!r.ok) throw new Error("Failed"); + return r.json(); + }) + .then(function (data) { + _nodeMetaCache = data.nodes || {}; + _renderNodeMetadata(); + }) + .catch(function () { + container.innerHTML = + '
Failed to load node metadata
'; + }); +} + +function _renderNodeMetadata() { + var container = document.getElementById("admin-node-metadata-content"); + if (!container) return; + var nodeIds = Object.keys(_nodeMetaCache).sort(); + if (!nodeIds.length) { + container.innerHTML = + '
No nodes registered
'; + return; + } + + var html = ""; + nodeIds.forEach(function (nid) { + var meta = _nodeMetaCache[nid] || []; + html += + '
'; + html += + '
'; + html += + "" + + escapeHtml(nid) + + " (" + + meta.length + + " keys)"; + html += "
"; + html += + '
'; + + // Table of metadata — all values passed through escapeHtml() + if (meta.length) { + html += ''; + html += + '"; + html += ''; + html += ''; + html += ''; + html += + ''; + meta.forEach(function (m) { + var valStr = + typeof m.value === "object" + ? JSON.stringify(m.value) + : String(m.value); + var isAuto = m.source === "auto"; + html += ""; + html += '"; + html += + '"; + html += + '"; + html += ""; + }); + html += "
Metadata for node ' + + escapeHtml(nid) + + "
KeyValueSourceActions
' + escapeHtml(m.key) + "' + + escapeHtml(valStr) + + "' + + escapeHtml(m.source) + + ""; + if (!isAuto) { + html += + ''; + } + html += "
"; + } else { + html += + '
No metadata
'; + } + + // Add metadata form + html += '
'; + html += + ''; + html += + ''; + html += + ''; + html += "
"; + + html += "
"; + }); + container.innerHTML = html; + + // Bind button handlers (data-* attrs carry node/key context) + var delBtns = container.querySelectorAll(".nm-del-btn"); + for (var d = 0; d < delBtns.length; d++) { + delBtns[d].addEventListener("click", function () { + _deleteNodeMeta( + this.getAttribute("data-node"), + this.getAttribute("data-key"), + ); + }); + } + var addBtns = container.querySelectorAll(".nm-add-btn"); + for (var a = 0; a < addBtns.length; a++) { + addBtns[a].addEventListener("click", function () { + _addNodeMeta(this.getAttribute("data-node")); + }); + } +} + +function _addNodeMeta(nodeId) { + var keyEl = document.getElementById("nm-key-" + nodeId); + var valEl = document.getElementById("nm-val-" + nodeId); + if (!keyEl || !valEl) return; + var key = keyEl.value.trim(); + var rawVal = valEl.value.trim(); + if (!key) { + showToast("Key is required", "error"); + return; + } + + var value; + try { + value = JSON.parse(rawVal); + } catch (e) { + value = rawVal; + } + + authFetch( + "/v1/api/admin/nodes/" + + encodeURIComponent(nodeId) + + "/metadata/" + + encodeURIComponent(key), + { + method: "PUT", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ value: value }), + }, + ) + .then(function (r) { + if (!r.ok) + return r + .json() + .catch(function () { + return {}; + }) + .then(function (d) { + throw new Error(d.error || "Failed"); + }); + showToast("Metadata set"); + loadAdminNodeMetadata(); + }) + .catch(function (e) { + showToast(e.message, "error"); + }); +} + +function _deleteNodeMeta(nodeId, key) { + showConfirmModal( + "Delete Metadata", + 'Delete key "' + key + '" from node ' + nodeId + "?", + "Delete", + function () { + authFetch( + "/v1/api/admin/nodes/" + + encodeURIComponent(nodeId) + + "/metadata/" + + encodeURIComponent(key), + { + method: "DELETE", + }, + ) + .then(function (r) { + if (!r.ok) + return r + .json() + .catch(function () { + return {}; + }) + .then(function (d) { + throw new Error(d.error || "Failed"); + }); + showToast("Metadata deleted"); + loadAdminNodeMetadata(); + }) + .catch(function (e) { + showToast(e.message, "error"); + }); + }, + ); +} diff --git a/turnstone/console/static/app.js b/turnstone/console/static/app.js index f0be85fc..c3f8af1b 100644 --- a/turnstone/console/static/app.js +++ b/turnstone/console/static/app.js @@ -967,6 +967,7 @@ function drillDownToNode(nodeId, serverUrl) { '
Loading workstreams...
'; loadNodeDetail(nodeId); } + _loadNodeMetadataPanel(nodeId); document.getElementById("breadcrumb-home").focus(); if (!_navigatingFromPopstate) history.pushState( @@ -1484,3 +1485,60 @@ function _ensureSSE() { history.replaceState({ view: "overview" }, ""); initLogin(); loadOverview(); + +// --- Node Metadata Panel (read-only in node detail view) --- +function _loadNodeMetadataPanel(nodeId) { + var section = document.getElementById("node-metadata-section"); + var table = document.getElementById("node-metadata-table"); + if (!section || !table) return; + section.style.display = "none"; + table.textContent = ""; + authFetch("/v1/api/cluster/node/" + encodeURIComponent(nodeId)) + .then(function (r) { + return r.ok ? r.json() : null; + }) + .then(function (data) { + if (!data || !data.metadata || !data.metadata.length) return; + section.style.display = ""; + var tbl = document.createElement("table"); + tbl.className = "nm-table"; + var thead = document.createElement("thead"); + var hr = document.createElement("tr"); + ["Key", "Value", "Source"].forEach(function (h) { + var th = document.createElement("th"); + th.setAttribute("scope", "col"); + th.textContent = h; + hr.appendChild(th); + }); + thead.appendChild(hr); + tbl.appendChild(thead); + var tbody = document.createElement("tbody"); + data.metadata.forEach(function (m) { + var tr = document.createElement("tr"); + var tdKey = document.createElement("td"); + tdKey.className = "nm-key"; + tdKey.textContent = m.key; + tr.appendChild(tdKey); + var tdVal = document.createElement("td"); + tdVal.className = "nm-val"; + tdVal.textContent = + typeof m.value === "object" + ? JSON.stringify(m.value) + : String(m.value); + tdVal.title = tdVal.textContent; + tr.appendChild(tdVal); + var tdSrc = document.createElement("td"); + var badge = document.createElement("span"); + badge.className = "nm-source-badge nm-source-" + m.source; + badge.textContent = m.source; + tdSrc.appendChild(badge); + tr.appendChild(tdSrc); + tbody.appendChild(tr); + }); + tbl.appendChild(tbody); + table.appendChild(tbl); + }) + .catch(function () { + /* silent — metadata is supplementary */ + }); +} diff --git a/turnstone/console/static/index.html b/turnstone/console/static/index.html index c7bcfde7..3019c337 100644 --- a/turnstone/console/static/index.html +++ b/turnstone/console/static/index.html @@ -53,6 +53,12 @@ CTX
+ Open node UI @@ -112,6 +118,7 @@
+
@@ -655,6 +662,16 @@ + + +