feat: add per-node metadata with auto-collection, admin API, and cons… (#318)

* 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)
This commit is contained in:
Patrick Buckley
2026-04-06 03:08:19 -07:00
committed by GitHub
parent 5cbc4bc87c
commit 24f59a6c53
19 changed files with 1571 additions and 9 deletions
+3 -1
View File
@@ -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(
+137
View File
@@ -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
+185
View File
@@ -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"}
+84
View File
@@ -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)
+34
View File
@@ -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)
+41
View File
@@ -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",
+8 -1
View File
@@ -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,
+258 -6
View File
@@ -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(
+234
View File
@@ -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 = '<div class="dashboard-empty">Loading\u2026</div>';
// 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 =
'<div class="dashboard-empty">Failed to load node metadata</div>';
});
}
function _renderNodeMetadata() {
var container = document.getElementById("admin-node-metadata-content");
if (!container) return;
var nodeIds = Object.keys(_nodeMetaCache).sort();
if (!nodeIds.length) {
container.innerHTML =
'<div class="dashboard-empty">No nodes registered</div>';
return;
}
var html = "";
nodeIds.forEach(function (nid) {
var meta = _nodeMetaCache[nid] || [];
html +=
'<div class="settings-section" data-section="nm-' +
escapeHtml(nid) +
'" data-collapsed>';
html +=
'<div class="settings-section-header" onclick="_toggleSettingsSection(this)" ';
html += 'onkeydown="_onSettingsHeaderKey(event,this)" ';
html += 'role="button" tabindex="0" aria-expanded="false" ';
html += 'aria-controls="nm-body-' + escapeHtml(nid) + '">';
html +=
"<span>" +
escapeHtml(nid) +
" <small>(" +
meta.length +
" keys)</small></span>";
html += "</div>";
html +=
'<div class="settings-section-body" id="nm-body-' +
escapeHtml(nid) +
'">';
// Table of metadata — all values passed through escapeHtml()
if (meta.length) {
html += '<table class="nm-table">';
html +=
'<caption class="sr-only">Metadata for node ' +
escapeHtml(nid) +
"</caption>";
html += '<thead><tr><th scope="col">Key</th>';
html += '<th scope="col">Value</th>';
html += '<th scope="col">Source</th>';
html +=
'<th scope="col"><span class="sr-only">Actions</span></th></tr></thead><tbody>';
meta.forEach(function (m) {
var valStr =
typeof m.value === "object"
? JSON.stringify(m.value)
: String(m.value);
var isAuto = m.source === "auto";
html += "<tr>";
html += '<td class="nm-key">' + escapeHtml(m.key) + "</td>";
html +=
'<td class="nm-val" title="' +
escapeHtml(valStr) +
'">' +
escapeHtml(valStr) +
"</td>";
html +=
'<td><span class="nm-source-badge nm-source-' +
escapeHtml(m.source) +
'">' +
escapeHtml(m.source) +
"</span></td>";
html += "<td>";
if (!isAuto) {
html +=
'<button class="admin-btn-danger nm-del-btn" aria-label="Delete ' +
escapeHtml(m.key) +
'" data-node="' +
escapeHtml(nid) +
'" data-key="' +
escapeHtml(m.key) +
'">Del</button>';
}
html += "</td></tr>";
});
html += "</tbody></table>";
} else {
html +=
'<div class="dashboard-empty" style="padding:8px">No metadata</div>';
}
// Add metadata form
html += '<div class="nm-add-row">';
html +=
'<input id="nm-key-' +
escapeHtml(nid) +
'" type="text" placeholder="key" aria-label="Metadata key">';
html +=
'<input id="nm-val-' +
escapeHtml(nid) +
'" type="text" placeholder="value (JSON or string)" aria-label="Metadata value">';
html +=
'<button class="admin-btn-action nm-add-btn" data-node="' +
escapeHtml(nid) +
'" style="white-space:nowrap">Add</button>';
html += "</div>";
html += "</div></div>";
});
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");
});
},
);
}
+58
View File
@@ -967,6 +967,7 @@ function drillDownToNode(nodeId, serverUrl) {
'<div class="dashboard-empty">Loading workstreams...</div>';
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 */
});
}
+17
View File
@@ -53,6 +53,12 @@
<span class="dash-col dash-col-ctx">CTX</span>
</div>
<div id="node-ws-table" class="dash-table" role="group" aria-label="Workstreams" aria-live="polite"></div>
<div id="node-metadata-section" style="margin-top:16px;display:none">
<div class="dash-header">
<span class="dash-header-title">METADATA</span>
</div>
<div id="node-metadata-table" style="font-size:.85rem"></div>
</div>
<a id="node-link" class="node-link">Open node UI</a>
</div>
@@ -112,6 +118,7 @@
<div class="admin-sidebar-group" data-group="system" role="group" aria-label="System">
<div class="admin-sidebar-group-label" aria-hidden="true">System</div>
<button id="tab-models" class="admin-nav" data-tab="models" role="tab" aria-selected="false" aria-controls="admin-models" tabindex="-1" onclick="switchAdminTab('models')">Models</button>
<button id="tab-node-metadata" class="admin-nav" data-tab="node-metadata" role="tab" aria-selected="false" aria-controls="admin-node-metadata" tabindex="-1" onclick="switchAdminTab('node-metadata')">Nodes</button>
<button id="tab-settings" class="admin-nav" data-tab="settings" role="tab" aria-selected="false" aria-controls="admin-settings" tabindex="-1" onclick="switchAdminTab('settings')">Settings</button>
<button id="tab-tls" class="admin-nav" data-tab="tls" role="tab" aria-selected="false" aria-controls="admin-tls" tabindex="-1" onclick="switchAdminTab('tls')">TLS</button>
</div>
@@ -655,6 +662,16 @@
</div>
</div>
<!-- Node Metadata Tab -->
<div id="admin-node-metadata" class="admin-panel" role="tabpanel" aria-labelledby="tab-node-metadata" style="display:none">
<div class="admin-toolbar">
<span class="section-header" style="margin:0">NODE METADATA</span>
</div>
<div id="admin-node-metadata-content" role="list" aria-label="Node metadata" aria-live="polite">
<div class="dashboard-empty">Loading&hellip;</div>
</div>
</div>
<!-- Settings Tab -->
<div id="admin-settings" class="admin-panel" role="tabpanel" aria-labelledby="tab-settings" style="display:none">
<div class="admin-toolbar">
+85
View File
@@ -2499,3 +2499,88 @@ h3.skill-spec-heading { font-size: inherit; margin-block: 0; }
.mcp-reg-card-repo { transition: none; }
.mcp-sync-pending, .model-sync-pending { animation: none; }
}
/* Node metadata */
.nm-source-badge {
display: inline-block;
padding: 1px 6px;
border-radius: var(--radius-sm);
font-size: .75rem;
font-weight: 600;
text-transform: uppercase;
letter-spacing: 0.04em;
}
.nm-source-auto { background: var(--green-glow); color: var(--green); }
.nm-source-user { background: var(--cyan-glow); color: var(--cyan); }
.nm-source-config { background: var(--yellow-glow); color: var(--yellow); }
.nm-table { width: 100%; border-collapse: collapse; }
.nm-table th {
font-family: var(--font-display);
font-size: 10px;
font-weight: 600;
text-transform: uppercase;
letter-spacing: 0.08em;
color: var(--fg-dim);
padding: 4px 8px;
text-align: left;
border-bottom: 1px solid var(--border);
}
.nm-table td {
padding: 4px 8px;
font-size: 12px;
color: var(--fg);
border-bottom: 1px solid var(--border);
}
.nm-key {
font-family: var(--font-mono);
color: var(--fg-bright);
}
.nm-val {
max-width: 300px;
overflow: hidden;
text-overflow: ellipsis;
white-space: nowrap;
}
.nm-add-row {
display: flex;
gap: 8px;
align-items: center;
padding: 8px 0;
}
.nm-add-row input[type="text"] {
padding: 5px 8px;
background: var(--bg);
border: 1px solid var(--border-strong);
border-radius: var(--radius-sm);
color: var(--fg);
font: inherit;
font-size: 12px;
transition: border-color 0.15s, box-shadow 0.15s;
}
.nm-add-row input[type="text"]:first-of-type { width: 120px; }
.nm-add-row input[type="text"]:nth-of-type(2) { flex: 1; }
.nm-add-row input[type="text"]:focus {
border-color: var(--accent);
outline: none;
box-shadow: 0 0 0 3px var(--accent-dim);
}
.nm-add-row input[type="text"]::placeholder {
color: var(--fg-dim);
opacity: 0.6;
}
.nm-add-row input[type="text"]:disabled {
opacity: 0.55;
cursor: not-allowed;
}
@media (max-width: 700px) {
.nm-add-row { flex-wrap: wrap; }
.nm-add-row input[type="text"] { width: 100% !important; flex: none; }
.nm-val { max-width: 150px; }
}
@media (prefers-reduced-motion: reduce) {
.nm-add-row input[type="text"] { transition: none; }
}
+69
View File
@@ -0,0 +1,69 @@
"""Collect auto-populated node metadata using stdlib only."""
from __future__ import annotations
import logging
import os
import platform
import socket
from typing import Any
log = logging.getLogger(__name__)
def _is_loopback_or_link_local(addr: str) -> bool:
"""Return True for loopback and link-local addresses."""
return addr.startswith("127.") or addr == "::1" or addr.startswith("fe80:")
def _collect_interfaces() -> dict[str, list[str]]:
"""Best-effort host IP collection using stdlib.
Returns a mapping from hostname to non-loopback IP addresses.
Without psutil/netifaces, per-interface resolution is not available
from stdlib alone, so we report resolved host addresses honestly.
"""
result: dict[str, list[str]] = {}
try:
hostname = socket.gethostname()
addrs = socket.getaddrinfo(hostname, None, proto=socket.IPPROTO_TCP)
ips = sorted({str(a[4][0]) for a in addrs if not _is_loopback_or_link_local(str(a[4][0]))})
if ips:
result[hostname] = ips
except OSError:
log.debug("node_info: interface collection failed", exc_info=True)
return result
def collect_node_info() -> dict[str, Any]:
"""Collect auto-populated node metadata.
Returns a dict of ``{key: value}`` where values are JSON-serializable.
Each field is collected independently one failure does not block others.
"""
info: dict[str, Any] = {}
for key, fn in (
("hostname", socket.gethostname),
("fqdn", socket.getfqdn),
("os", platform.system),
("os_release", platform.release),
("arch", platform.machine),
("python", platform.python_version),
("cpu_count", os.cpu_count),
):
try:
val = fn()
if val is not None:
info[key] = val
except Exception:
log.debug("node_info: failed to collect %s", key, exc_info=True)
try:
ifaces = _collect_interfaces()
if ifaces:
info["interfaces"] = ifaces
except Exception:
log.debug("node_info: failed to collect interfaces", exc_info=True)
return info
+114
View File
@@ -1337,6 +1337,120 @@ class PostgreSQLBackend:
conn.commit()
return result.rowcount > 0
# -- Node metadata ---------------------------------------------------------
def get_node_metadata(self, node_id: str) -> list[dict[str, Any]]:
from turnstone.core.storage._schema import node_metadata
with self._conn() as conn:
rows = conn.execute(
sa.select(node_metadata)
.where(node_metadata.c.node_id == node_id)
.order_by(node_metadata.c.key)
).fetchall()
return [dict(r._mapping) for r in rows]
def get_all_node_metadata(self) -> dict[str, list[dict[str, Any]]]:
from turnstone.core.storage._schema import node_metadata
with self._conn() as conn:
rows = conn.execute(
sa.select(node_metadata).order_by(node_metadata.c.node_id, node_metadata.c.key)
).fetchall()
result: dict[str, list[dict[str, Any]]] = {}
for r in rows:
d = dict(r._mapping)
result.setdefault(d["node_id"], []).append(d)
return result
def set_node_metadata(self, node_id: str, key: str, value: str, source: str = "user") -> None:
from sqlalchemy.dialects.postgresql import insert as pg_insert
from turnstone.core.storage._schema import node_metadata
now = datetime.now(UTC).strftime("%Y-%m-%dT%H:%M:%S")
stmt = pg_insert(node_metadata).values(
node_id=node_id,
key=key,
value=value,
source=source,
created=now,
updated=now,
)
stmt = stmt.on_conflict_do_update(
index_elements=[node_metadata.c.node_id, node_metadata.c.key],
set_={"value": value, "source": source, "updated": now},
)
with self._conn() as conn:
conn.execute(stmt)
conn.commit()
def set_node_metadata_bulk(self, node_id: str, entries: list[tuple[str, str, str]]) -> None:
from sqlalchemy.dialects.postgresql import insert as pg_insert
from turnstone.core.storage._schema import node_metadata
now = datetime.now(UTC).strftime("%Y-%m-%dT%H:%M:%S")
with self._conn() as conn:
for key, value, source in entries:
stmt = pg_insert(node_metadata).values(
node_id=node_id,
key=key,
value=value,
source=source,
created=now,
updated=now,
)
stmt = stmt.on_conflict_do_update(
index_elements=[node_metadata.c.node_id, node_metadata.c.key],
set_={"value": value, "source": source, "updated": now},
)
conn.execute(stmt)
conn.commit()
def delete_node_metadata(self, node_id: str, key: str) -> bool:
from turnstone.core.storage._schema import node_metadata
with self._conn() as conn:
result = conn.execute(
sa.delete(node_metadata).where(
(node_metadata.c.node_id == node_id) & (node_metadata.c.key == key)
)
)
conn.commit()
return result.rowcount > 0
def delete_node_metadata_by_source(self, node_id: str, source: str) -> int:
from turnstone.core.storage._schema import node_metadata
with self._conn() as conn:
result = conn.execute(
sa.delete(node_metadata).where(
(node_metadata.c.node_id == node_id) & (node_metadata.c.source == source)
)
)
conn.commit()
return result.rowcount
def filter_nodes_by_metadata(self, filters: dict[str, str]) -> set[str]:
from turnstone.core.storage._schema import node_metadata
if not filters:
return set()
conditions = [
sa.and_(node_metadata.c.key == k, node_metadata.c.value == v)
for k, v in filters.items()
]
stmt = (
sa.select(node_metadata.c.node_id)
.where(sa.or_(*conditions))
.group_by(node_metadata.c.node_id)
.having(sa.func.count() == len(filters))
)
with self._conn() as conn:
rows = conn.execute(stmt).fetchall()
return {r[0] for r in rows}
# -- Hash ring routing -----------------------------------------------------
def list_ring_buckets(self) -> list[dict[str, Any]]:
+30
View File
@@ -480,6 +480,36 @@ class StorageBackend(Protocol):
"""Remove a service registration. Returns True if existed."""
...
# -- Node metadata ---------------------------------------------------------
def get_node_metadata(self, node_id: str) -> list[dict[str, Any]]:
"""Return all metadata rows for a node."""
...
def get_all_node_metadata(self) -> dict[str, list[dict[str, Any]]]:
"""Return metadata grouped by node_id for all nodes."""
...
def set_node_metadata(self, node_id: str, key: str, value: str, source: str = "user") -> None:
"""Upsert a single metadata key for a node."""
...
def set_node_metadata_bulk(self, node_id: str, entries: list[tuple[str, str, str]]) -> None:
"""Upsert multiple (key, value, source) entries for a node. Atomic."""
...
def delete_node_metadata(self, node_id: str, key: str) -> bool:
"""Delete a single metadata key. Returns True if existed."""
...
def delete_node_metadata_by_source(self, node_id: str, source: str) -> int:
"""Delete all metadata for a node with the given source. Returns count."""
...
def filter_nodes_by_metadata(self, filters: dict[str, str]) -> set[str]:
"""Return node_ids where ALL key=value filters match (exact match)."""
...
# -- Hash ring routing ---
def list_ring_buckets(self) -> list[dict[str, Any]]:
+18
View File
@@ -233,6 +233,24 @@ services = sa.Table(
sa.Index("idx_services_type_heartbeat", services.c.service_type, services.c.last_heartbeat)
# ---------------------------------------------------------------------------
# Node metadata (per-node key/value with source tracking)
# ---------------------------------------------------------------------------
node_metadata = sa.Table(
"node_metadata",
metadata,
sa.Column("node_id", sa.Text, nullable=False),
sa.Column("key", sa.Text, nullable=False),
sa.Column("value", sa.Text, nullable=False),
sa.Column("source", sa.Text, nullable=False, server_default="user"),
sa.Column("created", sa.Text, nullable=False),
sa.Column("updated", sa.Text, nullable=False),
sa.PrimaryKeyConstraint("node_id", "key"),
)
sa.Index("idx_node_metadata_key", node_metadata.c.key)
# ---------------------------------------------------------------------------
# Hash ring routing tables
# ---------------------------------------------------------------------------
+114
View File
@@ -1414,6 +1414,120 @@ class SQLiteBackend:
conn.commit()
return result.rowcount > 0
# -- Node metadata ---------------------------------------------------------
def get_node_metadata(self, node_id: str) -> list[dict[str, Any]]:
from turnstone.core.storage._schema import node_metadata
with self._conn() as conn:
rows = conn.execute(
sa.select(node_metadata)
.where(node_metadata.c.node_id == node_id)
.order_by(node_metadata.c.key)
).fetchall()
return [dict(r._mapping) for r in rows]
def get_all_node_metadata(self) -> dict[str, list[dict[str, Any]]]:
from turnstone.core.storage._schema import node_metadata
with self._conn() as conn:
rows = conn.execute(
sa.select(node_metadata).order_by(node_metadata.c.node_id, node_metadata.c.key)
).fetchall()
result: dict[str, list[dict[str, Any]]] = {}
for r in rows:
d = dict(r._mapping)
result.setdefault(d["node_id"], []).append(d)
return result
def set_node_metadata(self, node_id: str, key: str, value: str, source: str = "user") -> None:
from sqlalchemy.dialects.sqlite import insert as sqlite_insert
from turnstone.core.storage._schema import node_metadata
now = datetime.now(UTC).strftime("%Y-%m-%dT%H:%M:%S")
stmt = sqlite_insert(node_metadata).values(
node_id=node_id,
key=key,
value=value,
source=source,
created=now,
updated=now,
)
stmt = stmt.on_conflict_do_update(
index_elements=["node_id", "key"],
set_={"value": value, "source": source, "updated": now},
)
with self._conn() as conn:
conn.execute(stmt)
conn.commit()
def set_node_metadata_bulk(self, node_id: str, entries: list[tuple[str, str, str]]) -> None:
from sqlalchemy.dialects.sqlite import insert as sqlite_insert
from turnstone.core.storage._schema import node_metadata
now = datetime.now(UTC).strftime("%Y-%m-%dT%H:%M:%S")
with self._conn() as conn:
for key, value, source in entries:
stmt = sqlite_insert(node_metadata).values(
node_id=node_id,
key=key,
value=value,
source=source,
created=now,
updated=now,
)
stmt = stmt.on_conflict_do_update(
index_elements=["node_id", "key"],
set_={"value": value, "source": source, "updated": now},
)
conn.execute(stmt)
conn.commit()
def delete_node_metadata(self, node_id: str, key: str) -> bool:
from turnstone.core.storage._schema import node_metadata
with self._conn() as conn:
result = conn.execute(
sa.delete(node_metadata).where(
(node_metadata.c.node_id == node_id) & (node_metadata.c.key == key)
)
)
conn.commit()
return result.rowcount > 0
def delete_node_metadata_by_source(self, node_id: str, source: str) -> int:
from turnstone.core.storage._schema import node_metadata
with self._conn() as conn:
result = conn.execute(
sa.delete(node_metadata).where(
(node_metadata.c.node_id == node_id) & (node_metadata.c.source == source)
)
)
conn.commit()
return result.rowcount
def filter_nodes_by_metadata(self, filters: dict[str, str]) -> set[str]:
from turnstone.core.storage._schema import node_metadata
if not filters:
return set()
conditions = [
sa.and_(node_metadata.c.key == k, node_metadata.c.value == v)
for k, v in filters.items()
]
stmt = (
sa.select(node_metadata.c.node_id)
.where(sa.or_(*conditions))
.group_by(node_metadata.c.node_id)
.having(sa.func.count() == len(filters))
)
with self._conn() as conn:
rows = conn.execute(stmt).fetchall()
return {r[0] for r in rows}
# -- Hash ring routing -----------------------------------------------------
def list_ring_buckets(self) -> list[dict[str, Any]]:
@@ -0,0 +1,50 @@
"""Add node_metadata table for per-node key/value metadata.
Revision ID: 035
Revises: 034
Create Date: 2026-04-05
"""
import sqlalchemy as sa
from alembic import op
revision = "035"
down_revision = "034"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.create_table(
"node_metadata",
sa.Column("node_id", sa.Text, nullable=False),
sa.Column("key", sa.Text, nullable=False),
sa.Column("value", sa.Text, nullable=False),
sa.Column("source", sa.Text, nullable=False, server_default="user"),
sa.Column("created", sa.Text, nullable=False),
sa.Column("updated", sa.Text, nullable=False),
sa.PrimaryKeyConstraint("node_id", "key"),
)
op.create_index("idx_node_metadata_key", "node_metadata", ["key"])
# Grant admin.nodes permission to the built-in admin role
conn = op.get_bind()
conn.execute(
sa.text(
"UPDATE roles SET permissions = permissions || ',admin.nodes' "
"WHERE role_id = 'builtin-admin' "
"AND permissions NOT LIKE '%admin.nodes%'"
)
)
def downgrade() -> None:
conn = op.get_bind()
conn.execute(
sa.text(
"UPDATE roles SET permissions = REPLACE(permissions, ',admin.nodes', '') "
"WHERE role_id = 'builtin-admin'"
)
)
op.drop_index("idx_node_metadata_key", table_name="node_metadata")
op.drop_table("node_metadata")
+32 -1
View File
@@ -2978,6 +2978,30 @@ async def _lifespan(app: Starlette) -> AsyncGenerator[None, None]:
_svc_storage.register_service("server", _svc_node_id, _svc_url)
log.info("server.service_registered", node_id=_svc_node_id, url=_svc_url)
# Collect and store node metadata (auto + config)
try:
from turnstone.core.config import load_config as _load_meta_config
from turnstone.core.node_info import collect_node_info
_auto_info = collect_node_info()
_meta_entries: list[tuple[str, str, str]] = [
(k, json.dumps(v), "auto") for k, v in _auto_info.items()
]
_cfg_meta = _load_meta_config("metadata")
_meta_entries.extend((k, json.dumps(v), "config") for k, v in _cfg_meta.items())
if _meta_entries:
# Clear stale auto/config rows from a prior run before upserting
_svc_storage.delete_node_metadata_by_source(_svc_node_id, "auto")
_svc_storage.delete_node_metadata_by_source(_svc_node_id, "config")
_svc_storage.set_node_metadata_bulk(_svc_node_id, _meta_entries)
log.info(
"server.node_metadata_stored",
node_id=_svc_node_id,
count=len(_meta_entries),
)
except Exception:
log.warning("server.node_metadata_failed", node_id=_svc_node_id, exc_info=True)
async def _heartbeat_loop() -> None:
"""Periodically update service heartbeat."""
from turnstone.core.storage._registry import StorageUnavailableError
@@ -3001,7 +3025,14 @@ async def _lifespan(app: Starlette) -> AsyncGenerator[None, None]:
from turnstone.core.storage import get_storage as _get_svc_dereg
try:
await asyncio.to_thread(_get_svc_dereg().deregister_service, "server", _svc_node_id)
_dereg_storage = _get_svc_dereg()
await asyncio.to_thread(_dereg_storage.deregister_service, "server", _svc_node_id)
await asyncio.to_thread(
_dereg_storage.delete_node_metadata_by_source, _svc_node_id, "auto"
)
await asyncio.to_thread(
_dereg_storage.delete_node_metadata_by_source, _svc_node_id, "config"
)
log.info("server.service_deregistered", node_id=_svc_node_id)
except Exception:
log.exception("server.deregister_failed")