File size: 84,326 Bytes
2aa8b3a 38cbc41 2aa8b3a 38cbc41 2aa8b3a 1d4a48e 95200a7 3c241d8 95200a7 3c241d8 1d4a48e 3c241d8 1d4a48e 3c241d8 1d4a48e 3c241d8 1d4a48e 3c241d8 1d4a48e 95200a7 2aa8b3a 3c241d8 2aa8b3a 3c241d8 95200a7 3c241d8 2aa8b3a 3c241d8 95200a7 3c241d8 95200a7 3c241d8 95200a7 3c241d8 95200a7 3c241d8 2aa8b3a 467fa9e 3c241d8 2aa8b3a 1d4a48e 2aa8b3a 1d4a48e 2aa8b3a 467fa9e 2aa8b3a 3c241d8 2aa8b3a 3c241d8 2aa8b3a 3c241d8 2aa8b3a 3c241d8 2aa8b3a 3c241d8 2aa8b3a 3c241d8 2aa8b3a 3c241d8 2aa8b3a 1d4a48e 2aa8b3a 467fa9e 2aa8b3a 1d4a48e 2aa8b3a 95200a7 2aa8b3a 7a2798b 2aa8b3a 1d4a48e 2aa8b3a 3c241d8 2aa8b3a 3c241d8 2aa8b3a 1d4a48e 2aa8b3a | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 915 916 917 918 919 920 921 922 923 924 925 926 927 928 929 930 931 932 933 934 935 936 937 938 939 940 941 942 943 944 945 946 947 948 949 950 951 952 953 954 955 956 957 958 959 960 961 962 963 964 965 966 967 968 969 970 971 972 973 974 975 976 977 978 979 980 981 982 983 984 985 986 987 988 989 990 991 992 993 994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 1026 1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075 1076 1077 1078 1079 1080 1081 1082 1083 1084 1085 1086 1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146 1147 1148 1149 1150 1151 1152 1153 1154 1155 1156 1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188 1189 1190 1191 1192 1193 1194 1195 1196 1197 1198 1199 1200 1201 1202 1203 1204 1205 1206 1207 1208 1209 1210 1211 1212 1213 1214 1215 1216 1217 1218 1219 1220 1221 1222 1223 1224 1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237 1238 1239 1240 1241 1242 1243 1244 1245 1246 1247 1248 1249 1250 1251 1252 1253 1254 1255 1256 1257 1258 1259 1260 1261 1262 1263 1264 1265 1266 1267 1268 1269 1270 1271 1272 1273 1274 1275 1276 1277 1278 1279 1280 1281 1282 1283 1284 1285 1286 1287 1288 1289 1290 1291 1292 1293 1294 1295 1296 1297 1298 1299 1300 1301 1302 1303 1304 1305 1306 1307 1308 1309 1310 1311 1312 1313 1314 1315 1316 1317 1318 1319 1320 1321 1322 1323 1324 1325 1326 1327 1328 1329 1330 1331 1332 1333 1334 1335 1336 1337 1338 1339 1340 1341 1342 1343 1344 1345 1346 1347 1348 1349 1350 1351 1352 1353 1354 1355 1356 1357 1358 1359 1360 1361 1362 1363 1364 1365 1366 1367 1368 1369 1370 1371 1372 1373 1374 1375 1376 1377 1378 1379 1380 1381 1382 1383 1384 1385 1386 1387 1388 1389 1390 1391 1392 1393 1394 1395 1396 1397 1398 1399 1400 1401 1402 1403 1404 1405 1406 1407 1408 1409 1410 1411 1412 1413 1414 1415 1416 1417 1418 1419 1420 1421 1422 1423 1424 1425 1426 1427 1428 1429 1430 1431 1432 1433 1434 1435 1436 1437 1438 1439 1440 1441 1442 1443 1444 1445 1446 1447 1448 1449 1450 1451 1452 1453 1454 1455 1456 1457 1458 1459 1460 1461 1462 1463 1464 1465 1466 1467 1468 1469 1470 1471 1472 1473 1474 1475 1476 1477 1478 1479 1480 1481 1482 1483 1484 1485 1486 1487 1488 1489 1490 1491 1492 1493 1494 1495 1496 1497 1498 1499 1500 1501 1502 1503 1504 1505 1506 1507 1508 1509 1510 1511 1512 1513 1514 1515 1516 1517 1518 1519 1520 1521 1522 1523 1524 1525 1526 1527 1528 1529 1530 1531 1532 1533 1534 1535 1536 1537 1538 1539 1540 1541 1542 1543 1544 1545 1546 1547 1548 1549 1550 1551 1552 1553 1554 1555 1556 1557 1558 1559 1560 1561 1562 1563 1564 1565 1566 1567 1568 1569 1570 1571 1572 1573 1574 1575 1576 1577 1578 1579 1580 1581 1582 1583 1584 1585 1586 1587 1588 1589 1590 1591 1592 1593 1594 1595 1596 1597 1598 1599 1600 1601 1602 1603 1604 1605 1606 1607 1608 1609 1610 1611 1612 1613 1614 1615 1616 1617 1618 1619 1620 1621 1622 1623 1624 1625 1626 1627 1628 1629 1630 1631 1632 1633 1634 1635 1636 1637 1638 1639 1640 1641 1642 1643 1644 1645 1646 1647 1648 1649 1650 1651 1652 1653 1654 1655 1656 1657 1658 1659 1660 1661 1662 1663 1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678 1679 1680 1681 1682 1683 1684 1685 1686 1687 1688 1689 1690 1691 1692 1693 1694 1695 1696 1697 1698 1699 1700 1701 1702 1703 1704 1705 1706 1707 1708 1709 1710 1711 1712 1713 1714 1715 1716 1717 1718 1719 1720 1721 1722 1723 1724 1725 1726 1727 1728 1729 1730 1731 1732 1733 1734 1735 1736 1737 1738 1739 1740 1741 1742 1743 1744 1745 1746 1747 1748 1749 1750 1751 1752 1753 1754 1755 1756 1757 1758 1759 1760 1761 1762 1763 1764 1765 1766 1767 1768 1769 1770 1771 1772 1773 1774 1775 1776 1777 1778 1779 1780 1781 1782 1783 1784 1785 1786 1787 1788 1789 1790 1791 1792 1793 1794 1795 1796 1797 1798 1799 1800 1801 1802 1803 1804 1805 1806 1807 1808 1809 1810 1811 1812 1813 1814 1815 1816 1817 1818 1819 1820 1821 1822 1823 1824 1825 1826 1827 1828 1829 1830 1831 1832 1833 1834 1835 1836 1837 1838 1839 1840 1841 1842 1843 1844 1845 1846 1847 1848 1849 1850 1851 1852 1853 1854 1855 1856 1857 1858 1859 1860 1861 1862 1863 1864 1865 1866 1867 1868 1869 1870 1871 1872 1873 1874 1875 1876 1877 1878 1879 1880 1881 1882 1883 1884 1885 1886 1887 1888 1889 1890 1891 1892 1893 1894 1895 1896 1897 1898 1899 1900 1901 1902 1903 1904 1905 1906 1907 1908 1909 1910 1911 1912 1913 1914 1915 1916 1917 1918 1919 1920 1921 1922 1923 1924 1925 1926 1927 1928 1929 1930 1931 1932 1933 1934 1935 1936 1937 1938 1939 1940 1941 1942 1943 1944 1945 1946 1947 1948 1949 1950 1951 1952 1953 1954 1955 1956 1957 1958 1959 1960 1961 1962 1963 1964 1965 1966 1967 1968 1969 1970 1971 1972 1973 1974 1975 1976 1977 1978 1979 1980 1981 1982 1983 1984 1985 1986 1987 1988 1989 1990 1991 1992 1993 1994 1995 1996 1997 1998 1999 2000 2001 2002 2003 2004 2005 2006 2007 2008 2009 2010 2011 2012 2013 2014 2015 2016 2017 2018 2019 2020 2021 2022 2023 2024 2025 2026 2027 2028 2029 2030 2031 2032 2033 2034 2035 2036 | """Minimal CLI and injectable transport entry point for the local worker."""
from __future__ import annotations
import argparse
import ctypes
import json
import math
import os
import platform
import re
import sys
import threading
import time
import uuid
from collections.abc import Callable, Iterable
from dataclasses import replace
from pathlib import Path
from urllib.parse import urlsplit
# Before any project import: the project uses 3.10 syntax at module level
# (PEP 604 unions in isinstance), so on an older interpreter the failure is
# an unrelated-looking TypeError deep in an import. A person on Windows whose
# ``python`` is a stray 3.9 deserves the actual sentence instead.
if sys.version_info < (3, 10): # pragma: no cover - cannot run under test
sys.exit(
"distinct needs Python 3.10 or newer; this is "
f"{sys.version_info.major}.{sys.version_info.minor} ({sys.executable}).\n"
"On Windows, run: py -3.12 -m distinct_agent ...\n"
"Elsewhere, run: python3.12 -m distinct_agent ..."
)
from distinct_protocol import (
MODE_AGENT,
MODE_SIMPLE,
AgentCapabilities,
JobResult,
JobSpec,
JobStatus,
)
from distinct_protocol.netpolicy import (
ALLOW_PLAINTEXT_SERVER_ENV_VAR,
ALLOW_REMOTE_SERVER_ENV_VAR,
DEFAULT_SERVER_URL,
announce_remote_server,
is_loopback_host,
)
from distinct_tools import (
INSTALLABLE_SPECS,
REPOSITORY_SPECS,
REPOSITORY_STATUS,
operator_policy,
registry_notes,
)
from distinct_tools.approval import APPROVED_TOOLS_ENV_VAR
from .access import announce_access_code, new_access_code
from .energy import detect_energy_meter
from .fanout import MultiServerTransport, ServerLink
from .isolation import describe_request_isolation, describe_support
from .models import DiscoveredModel, ModelVerificationError
from .queue import DEFAULT_POLL_INTERVAL_SECONDS, InMemoryJobQueue, OfferDecision
from .rlm import RlmUnavailable, build_inference_harness, prepare_rlm_runtime
from .runners import DeterministicDemoRunner, LlamaCppRunner
from .sandbox import GeneratedExecutionPolicy, check_deno_sandbox
from .server_runner import LlamaServerRunner
from .servers import (
REASON_DENIED_BY_AGENT,
AllowanceError,
AllowanceGate,
Refusal,
ServerRegistry,
)
from .tools import advertised_tool_refs
from .transport import (
AgentTransportError,
GradioAgentTransport,
PredictClient,
)
from .weights import (
DEFAULT_CACHE_GB,
WeightsCache,
WeightsUnavailable,
cache_bytes,
describe_catalogue,
fetchable_manifests,
shared_cache_directory,
)
from .worker import AgentTransport, LocalInferenceExecutor, WorkerLoop, WorkerLoopSettings
class ConsoleTransport:
"""One-job transport used only for local smoke tests."""
def __init__(self, offer: JobSpec | None = None) -> None:
self.offer = offer
self.done = threading.Event()
def poll_offer(self, snapshot):
print(json.dumps({"snapshot": snapshot.to_dict()}, sort_keys=True))
offer, self.offer = self.offer, None
return offer
def poll_cancellations(self) -> Iterable[str]:
return ()
def acknowledge_offer(self, decision: OfferDecision) -> None:
print(json.dumps({"offer": decision.to_dict()}, sort_keys=True))
if not decision.accepted:
self.done.set()
def publish_status(self, job_id: str, status: JobStatus, reason=None) -> None:
print(
json.dumps(
{"status": {"job_id": job_id, "status": status.value, "reason": reason}},
sort_keys=True,
)
)
if status in {JobStatus.CANCELLED, JobStatus.FAILED}:
self.done.set()
def publish_result(self, result: JobResult) -> None:
print(json.dumps({"result": result.to_dict()}, sort_keys=True))
self.done.set()
def _runner_problem(args) -> str:
"""Why this machine cannot run a model yet, in words, or empty if it can.
Written as an instruction rather than a diagnosis. "llama-server was not
found. Supply --llama-server" told somebody who had just waited out a nine
gigabyte download the name of a flag, not where to get the program that
flag points at, and read as though the flag took no value — which it does,
so the next thing they saw was an argparse usage dump.
"""
if args.llama_cli and not args.llama_server:
from shutil import which
if which(args.llama_cli) or Path(args.llama_cli).is_file():
return ""
return (
f"llama-cli was not found at {args.llama_cli!r}. Give the path to the "
"binary, or run with --demo-runner to start without a model."
)
if args.llama_server and args.llama_server is not True:
if Path(str(args.llama_server)).is_file():
return ""
return (
f"llama-server was not found at {str(args.llama_server)!r}. Give the "
"path to the binary, or run with --demo-runner to start without a model."
)
if _resolve_llama_server(args):
return ""
from . import runtime as _runtime
where = _runtime.platform_key()
if all(_runtime.pin_for(key) is None for key in _runtime.preferred_keys()):
# Nothing recorded for this platform, so nothing was attempted. Say
# which platform, because the answer is "add a pin for it" and that is
# a thing somebody can do.
return (
"This worker has no llama.cpp to run a model with, and none is recorded\n"
f"for {where}, so nothing was downloaded.\n"
"\n"
" An unverified build of the program that runs the models is the one\n"
" thing this project will not install for you. Either point at a build\n"
" you already have:\n"
"\n"
" --llama-server /path/to/llama-server\n"
"\n"
" or start without a model. The worker still joins, and still runs every\n"
" tool and skill in the library; only the model's own replies are absent:\n"
"\n"
" --demo-runner\n"
"\n"
f" To record a build for {where}, run:\n"
f" python scripts/record_llama_pin.py {where}\n"
)
return (
"This worker has no llama.cpp to run a model with, and the pinned build\n"
"could not be installed. The reason is above.\n"
"\n"
" Point at a build you already have:\n"
"\n"
" --llama-server /path/to/llama-server\n"
"\n"
" or start without a model. The worker still joins, and still runs every\n"
" tool and skill in the library; only the model's own replies are absent:\n"
"\n"
" --demo-runner\n"
)
def _resolve_llama_server(args) -> str:
"""The build this machine should run, installing it if that is what it takes.
ORDER MATTERS, AND THE OLD ORDER PINNED A MACHINE TO ITS FIRST INSTALL.
It was "anything already installed, then fetch". A volunteer whose first
run predated accelerated builds -- or who joined before their card was
recorded -- had a CPU build in ``runtime/``, so the answer was always yes,
and they stayed on the CPU for ever. That is the state the benchmark was
measured in: a CUDA laptop at three tokens a second.
So the first question is now "is the *best* build for this machine already
here", and only a no sends it to fetch. Anything installed, and then PATH,
remain as the last resorts they always were.
"""
from . import runtime as _runtime
for key in _runtime.preferred_keys():
found = _runtime.installed_server(key=key)
if found:
return found
# Nothing installed for the best key this machine can use. Fetching it
# is the point; if it cannot be fetched, the loop falls through to the
# next-best key, and _obtain_llama_server reports why in passing.
break
return _obtain_llama_server(args) or _find_llama_server()
def _find_llama_server() -> str:
"""Look in runtime/ first, then PATH.
``runtime/`` is where :func:`distinct_agent.runtime.ensure` installs the
pinned build, so a worker that fetched one on a previous run finds it here
and does not fetch again.
"""
from . import runtime as _runtime
installed = _runtime.best_installed()
if installed:
return installed
from shutil import which
for name in ("llama-server", "llama-server.exe"):
found = which(name)
if found:
return found
return ""
def _obtain_llama_server(args) -> str:
"""Fetch and verify the pinned llama.cpp for this machine, or say why not.
Returns the path, or an empty string when nothing could be installed. The
reasons are printed as they happen rather than collected, because the
download is the slow part and a volunteer watching a progress line wants to
know what it is for.
"""
from . import runtime as _runtime
if getattr(args, "no_fetch_runtime", False):
return ""
wanted = _runtime.preferred_keys()
pin = next((found for found in map(_runtime.pin_for, wanted) if found is not None), None)
if pin is None:
return ""
how = "CPU" if pin.accelerator == "cpu" else f"{pin.accelerator.upper()}-accelerated"
print(
f"No llama.cpp on this machine. Fetching the pinned {how} build "
f"({pin.tag}, {pin.bytes / (1024 * 1024):.0f} MB) and verifying it against "
"its recorded SHA-256. This happens once.",
file=sys.stderr,
flush=True,
)
last = [-1]
def report(fraction: float, message: str) -> None:
percent = int(fraction * 100)
if percent // 10 != last[0] // 10 or fraction >= 1.0:
last[0] = percent
print(f" llama.cpp: {message} ({percent}%)", file=sys.stderr, flush=True)
try:
return _runtime.ensure_best(progress=report)
except _runtime.RuntimeVerificationError as exc:
# A digest mismatch is not a transient failure and must not be retried
# around: it is what a substituted build looks like.
print(f"\n{exc}\n", file=sys.stderr, flush=True)
return ""
except _runtime.RuntimeUnavailable as exc:
print(f" llama.cpp could not be installed: {exc}", file=sys.stderr, flush=True)
return ""
def build_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(description="Run a Distinct local inference worker")
parser.add_argument(
"--server",
action="append",
nargs="?",
const=DEFAULT_SERVER_URL,
default=[],
metavar="URL",
help=(
"URL of a Distinct server this worker will accept work from. "
f"Defaults to {DEFAULT_SERVER_URL} when given without a value. "
"Repeatable: one machine can serve several servers at once, "
"sharing one queue, so the capacity it advertises is the capacity "
"it has. Give a --pair code for each, in the same order. A server "
"on another machine is supported but must be opted into with "
"--allow-remote-server, and must use HTTPS."
),
)
parser.add_argument(
"--list-servers",
action="store_true",
help="print the servers this worker has paired with, then exit",
)
parser.add_argument(
"--forget-server",
action="append",
default=[],
metavar="URL",
help=(
"remove a paired server by URL and exit. Its credentials and the "
"operator approval bound to them are deleted; joining it again "
"means pairing and approving again. Repeatable."
),
)
parser.add_argument(
"--servers-file",
default=None,
metavar="PATH",
help=(
"where paired servers are remembered between runs. Defaults to a "
"per-user file beside the weights cache. Contains credentials, so "
"it is written readable only by you."
),
)
parser.add_argument(
"--allow-plaintext-server",
action="store_true",
help=(
"permit a remote --server that is not HTTPS. Separate from "
"--allow-remote-server because it is a separate risk: that one is "
"about which machine, this one is about whether anything between "
"you and it can read your prompts. Use it when the link is already "
"encrypted underneath, such as a Tailscale or WireGuard address, "
"and not on open wifi. Equivalent to "
f"{ALLOW_PLAINTEXT_SERVER_ENV_VAR}=1."
),
)
parser.add_argument(
"--server-api",
choices=("gradio", "http"),
default="gradio",
help=(
"how to reach the server. 'gradio' is the desktop server, which "
"serves its agent API through Gradio's client protocol. 'http' is "
"for servers that are not Gradio, such as the phone app, which "
"serve the same five signed endpoints as plain JSON over HTTP. "
"The signing and the response verification are identical either "
"way; only the framing differs."
),
)
parser.add_argument(
"--allow-remote-server",
action="store_true",
help=(
"permit --server to name a machine other than this one. Its users "
"can send prompts to this worker and will see its hardware and "
f"energy telemetry. Equivalent to {ALLOW_REMOTE_SERVER_ENV_VAR}=1."
),
)
parser.add_argument(
"--pair",
action="append",
default=[],
metavar="CODE",
help=(
"short-lived one-use pairing code from the server. Repeatable, "
"matched to --server in the order both are given."
),
)
parser.add_argument("--name", default=platform.node() or "Distinct worker")
parser.add_argument("--agent-id", help="optional fixed agent id; the server assigns one by default")
# Models are named, never located. See distinct_agent.weights for why the
# old ``--model ID=PATH`` form was removed rather than kept alongside this.
parser.add_argument(
"--models",
action="append",
default=[],
metavar="ID[,ID...]",
help=(
"model ids this worker is willing to run, from the catalogue "
"(repeatable, or comma-separated). Defaults to every model that can "
"be fetched and verified. Weights are downloaded on demand and cached."
),
)
parser.add_argument(
"--tools",
action="append",
default=[],
metavar="REF[,REF...]",
help=(
"library members this worker will offer: exact refs such as "
"create_pdf@1, or the bundle name 'local' for every installed "
"member that opens no socket. Nothing is offered unless named here, "
"and naming one is the operator approval the broker checks."
),
)
parser.add_argument(
"--list-tools",
action="store_true",
help="print the library members this worker could offer, then exit",
)
parser.add_argument(
"--fetch-on-demand",
action="store_true",
help=(
"do not download the chosen models at start-up; fetch each one the "
"first time a run needs it instead. The first such run then waits "
"for the download."
),
)
parser.add_argument(
"--list-models",
action="store_true",
help="print the models this worker could offer, and why any are withheld, then exit",
)
parser.add_argument(
"--cache-dir",
default=None,
help=(
"where fetched weights are cached. Defaults to a shared per-user "
"location so every agent and tool on this machine reuses the same "
"download rather than each keeping its own copy: "
"DISTINCT_MODELS_DIR if set, else the Hugging Face cache named by "
"HF_HOME, else the platform user cache."
),
)
parser.add_argument(
"--models-dir",
dest="cache_dir_legacy",
default=None,
help="deprecated alias for --cache-dir",
)
parser.add_argument(
"--model-cache-gb",
type=float,
default=DEFAULT_CACHE_GB,
metavar="GB",
help=(
"ceiling on the weights cache. When a fetch would exceed it, the "
"least recently used model is evicted first. "
f"Default {DEFAULT_CACHE_GB:g} GB."
),
)
parser.add_argument("--model-id", help="model used by --demo-prompt")
# Present only so that an operator running the old command gets a sentence
# rather than argparse's "ambiguous option: --model".
parser.add_argument(
"--model",
action="append",
default=[],
metavar="ID=PATH",
help=argparse.SUPPRESS,
)
parser.add_argument(
"--without-energy-measurement",
action="store_true",
help=(
"start even though this machine has no readable energy counter. "
"Every run will report no energy figure. Off by default, because a "
"worker that cannot measure cannot do the one thing this network is "
"for."
),
)
parser.add_argument(
"--no-fetch-runtime",
action="store_true",
help=(
"never download llama.cpp. The worker uses one already on this "
"machine or refuses to run models, which is what an operator on a "
"metered or air-gapped connection wants."
),
)
parser.add_argument(
"--llama-server",
nargs="?",
const=True,
help=(
"path to llama-server. This is the default runner: clean output, "
"real token counts, streaming, and a resident model across queued "
"requests. Auto-detected in runtime/ or on PATH."
),
)
parser.add_argument(
"--llama-cli",
nargs="?",
const=True,
help="legacy per-request runner. Only used when --llama-server is absent.",
)
parser.add_argument("--queue-capacity", type=int, default=4)
parser.add_argument(
"--poll-interval",
type=float,
default=None,
help=(
"seconds between polls; defaults to the interval the server advertises at "
"pairing, clamped to [1, 30]. The poll carries the snapshot, so this is "
"also the heartbeat rate."
),
)
parser.add_argument("--demo-prompt", help="run one deterministic local smoke-test job")
parser.add_argument(
"--demo-runner",
action="store_true",
help="use deterministic inference (explicit test mode; no real model inference)",
)
parser.add_argument(
"--mode",
dest="harness",
choices=("agent", "simple", "dspy-rlm", "structured"),
help=(
"how this worker serves a request. 'agent' (the default) runs a "
"DSPy RLM agent in the Deno sandbox: it reasons, calls tools and "
"loops. 'simple' does one planning pass, runs whatever tools the "
"model named, then one answering pass, which holds up on models too "
"small to sustain an agent loop. 'dspy-rlm' and 'structured' are the "
"internal names and remain accepted."
),
)
parser.add_argument(
"--harness",
dest="harness",
choices=("agent", "simple", "dspy-rlm", "structured"),
help=argparse.SUPPRESS,
)
parser.add_argument(
"--prepare-rlm",
action="store_true",
help="prepare the lock-pinned RLM cache, verify cached-only execution, and exit",
)
parser.add_argument(
"--no-request-sandbox",
action="store_true",
help=(
"run every request in the worker's own directories instead of "
"giving each one a fresh, disposable sandbox. Only for diagnosing "
"a sandbox problem: it removes the boundary between one user's "
"request and the next."
),
)
parser.add_argument(
"--allow-generated-code",
action="store_true",
help="enable only after the strict Deno permission probe passes",
)
parser.add_argument(
"--guard-model",
metavar="MODEL_ID",
help=(
"model used for the input/output safety screen, chosen by the "
"agent operator at setup. Must be one of this worker's installed "
"models. Defaults to the model already loaded for the run."
),
)
parser.add_argument(
"--no-guard",
action="store_true",
help="disable the input/output safety screen entirely (operator's choice)",
)
parser.add_argument(
"--dashboard",
action="store_true",
help=(
"draw the live status screen instead of scrolling a log. This is "
"what a double-clicked executable uses; it is a flag rather than a "
"separate program so that a screen and a log are the same worker "
"with the same checks, differing only in how they report."
),
)
parser.add_argument(
"--approve",
action="store_true",
help=(
"record operator approval for the full choice set (the server's "
"catalogue intersected with what this worker offers). Without it, "
"an interactive console prompts; a non-interactive worker starts "
"with nothing approved and refuses every job, which is the honest "
"default rather than a failure."
),
)
return parser
def _open_registry(args) -> "ServerRegistry":
"""The persisted server registry, with a permissive declared allowance.
``--list-servers`` and ``--forget-server`` read and edit what is already
on file; they do not pair with anything, so the declared allowance they
are opened with bounds nothing and is deliberately empty rather than a
guess at what this boot would have offered.
"""
from distinct_protocol.handshake import Allowance
return ServerRegistry(Allowance(models=(), tools=()), path=_servers_file(args.servers_file))
def _list_servers(args) -> int:
registry = _open_registry(args)
ids = registry.server_ids()
if not ids:
print("This worker has not paired with any server.")
print(f"Nothing is stored at {_servers_file(args.servers_file)}.")
return 0
print(f"{len(ids)} server(s), remembered in {_servers_file(args.servers_file)}:")
for server_id in ids:
record = registry.require(server_id)
approval = getattr(record, "approval", None)
if approval is None:
terms = "no approval on file, so it would be refused every job"
else:
models = ", ".join(approval.allowance.models) or "no models"
tools = ", ".join(approval.allowance.tools) or "no tools"
terms = f"models {models}; tools {tools}"
print(f" {server_id}")
print(f" url: {record.url}")
print(f" approved: {terms}")
print("\nUse --forget-server URL to leave one.")
return 0
def _forget_servers(args) -> int:
"""Leave one or more servers, by URL, and say what was dropped.
By URL rather than by server id, because the URL is what the operator
typed and the id is something the server chose. A URL that matches
nothing is an error rather than a silent success: an operator who
mistypes it and sees nothing would reasonably believe they had left.
"""
registry = _open_registry(args)
by_url = {registry.require(server_id).url: server_id for server_id in registry.server_ids()}
unknown = [url for url in args.forget_server if url not in by_url]
if unknown:
print(f"Not paired with: {', '.join(unknown)}", file=sys.stderr)
if by_url:
print("Known servers:", file=sys.stderr)
for url in sorted(by_url):
print(f" {url}", file=sys.stderr)
else:
print("This worker has not paired with any server.", file=sys.stderr)
return 2
for url in args.forget_server:
server_id = by_url[url]
registry.forget(server_id)
print(
f"Left {url} ({server_id}). Its credentials and the approval bound to "
"them are gone; joining again means pairing and approving again."
)
return 0
def _no_meter_message() -> str:
"""Why this machine cannot be a worker, and what would change that."""
from .energy import energy_backend_report
lines = [
"This worker will not start, because it cannot measure the energy it uses.",
"",
" distinct exists to put a metered figure beside every answer. A worker",
" with no readable counter would serve runs that report nothing, which is",
" worse than not serving them.",
"",
" Backends tried on this machine:",
]
for label, ok, reason in energy_backend_report():
mark = "works" if ok else "no"
lines.append(f" [{mark:>5}] {label}")
if not ok and reason:
lines.append(f" {reason}")
lines += [
"",
" The usual fixes:",
" * NVIDIA GPU: install the library with pip install nvidia-ml-py",
" * AMD GPU: install ROCm, which ships the amdsmi package",
" * Linux CPU: RAPL lives in /sys/class/powercap and is often root-only.",
" Containers usually hide it; run on the host, or grant the path.",
" * Windows CPU: the Energy Meter counter needs an Intel or AMD chip",
" with a working EMI driver. A VM will not have one.",
" * A laptop: unplug it. The battery is itself a meter, and this can",
" read it while the machine is discharging.",
" * A server: ipmitool, if the BMC implements DCMI.",
"",
" To run anyway, knowing every run will report no energy:",
" --without-energy-measurement",
]
return "\n".join(lines)
def _server_pairs(urls, codes):
"""Match each --server to its --pair, or explain the mismatch.
With no ``--pair`` values at all, every server is paired the other way
round: the worker invents a claim code, prints it, and waits for somebody
signed in to enter it. That is the path the website documents, because the
old one -- fetch a code from the page, carry it to a terminal -- reliably
produced a person typing a value of their own where the code went and being
told one program later that their code was invalid.
``--pair`` is kept for scripted setups, where a code is minted and consumed
by the same program and nobody reads anything. Positional matching, because
it is the only scheme that reads correctly aloud: "this server, this code;
that server, that code". A mismatch is an error rather than a best guess,
since guessing would mean sending one server's code to another, which is a
credential leak dressed up as convenience.
"""
if not codes:
# NO --pair AT ALL IS THE ORDINARY CASE NOW.
#
# Every server gets its own freshly invented claim code, and the
# operator types it into that server's page. `None` for a code is the
# signal further down to register-and-wait rather than pair.
return [(url, None) for url in urls]
if len(codes) != len(urls):
print(
f"{len(urls)} --server value(s) and {len(codes)} --pair value(s). "
"Either give no --pair at all, and each server will print a code for "
"you to enter on its page, or give one code per server in the same "
"order. Codes are never interchangeable between servers.",
file=sys.stderr,
)
return None
seen: set[str] = set()
pairs = []
for url, code in zip(urls, codes, strict=True):
if url in seen:
print(f"--server {url} was given twice.", file=sys.stderr)
return None
seen.add(url)
pairs.append((url, code))
return pairs
def _servers_file(explicit: str | None) -> Path:
"""Where paired servers are remembered between runs.
Beside the weights cache, because both are per-user worker state and an
operator who knows where one lives can find the other. It holds
credentials, so :mod:`distinct_agent.servers` writes it private.
"""
if explicit:
return Path(explicit)
return shared_cache_directory().parent / "servers.json"
def _register_and_wait(
url,
capabilities,
*,
access_code,
client_factory,
allow_remote,
allow_plaintext,
signing_public_key,
poll_seconds: float = 3.0,
announce=None,
sleep=None,
):
"""Print a claim code for this server and wait until somebody enters it.
Returns the transport once the worker has been claimed, or ``None`` if it
was not claimed in time or the server refused the registration.
The wait is deliberately visible. A worker that registered silently and sat
there would look identical to one that had hung, and the whole point of
this direction is that the person is looking at the terminal when the code
appears.
"""
import time as _time
from .access import new_access_code as _new_code
say = announce or (lambda text: print(text, file=sys.stderr, flush=True))
nap = sleep or _time.sleep
claim_code = _new_code()
try:
transport = GradioAgentTransport.register(
url,
claim_code,
capabilities,
access_code=access_code,
client_factory=client_factory,
allow_remote=allow_remote or None,
allow_plaintext=allow_plaintext or None,
signing_public_key=signing_public_key,
)
except (AgentTransportError, TypeError, ValueError) as exc:
say(f"Registering with {url} failed: {exc}")
return None
say("")
say(" ┌─────────────────────────────────────────────────────────┐")
say(" │ This worker is waiting to be claimed. │")
say(" └─────────────────────────────────────────────────────────┘")
say("")
say(f" Sign in at {url}")
say(' open "Run a community agent", and enter this code:')
say("")
say(f" {claim_code}")
say("")
say(" Nobody can use this machine until you do, and the code works")
say(" only once. Waiting…")
say("")
deadline = _time.monotonic() + 15 * 60
while _time.monotonic() < deadline:
try:
state = transport.claim_state()
except (AgentTransportError, TypeError, ValueError):
# A poll that fails is not a claim that failed. Keep waiting; a
# server that is really gone will still be gone at the deadline.
state = "waiting"
if state == "claimed":
say(f"Claimed. {capabilities.name!r} is now serving {url} as {transport.agent_id}.")
return transport
if state == "gone":
say(
"This worker's registration expired before anybody claimed it. "
"Nothing was left running on the server; start the worker again "
"for a fresh code."
)
return None
nap(poll_seconds)
say(
"Nobody claimed this worker within fifteen minutes, so it has stopped "
"waiting. Start it again for a fresh code."
)
return None
def _join_server(
url,
code,
*,
capabilities,
manifests,
client_factory,
allow_remote,
allow_plaintext,
approve_flag,
registry,
access_code="",
):
"""Pair with one server, consent to it, and return a link and its gate.
Returns ``(None, None)`` on any failure, having said why. A server that
cannot be joined is not fatal to a worker that has others: the caller
carries on, and refuses to start only if nothing was joined at all.
"""
from distinct_protocol.handshake import (
Allowance,
OperatorApproval,
ed25519_public_key,
generate_ed25519_key,
)
from .transport import decode_agent_secret as _decode
# Say what a remote server is being granted before pairing gives it
# anything, not after.
host = urlsplit(url).hostname or ""
if host and not is_loopback_host(host):
announce_remote_server(url)
declared = Allowance(
models=tuple(manifest.id for manifest in manifests),
tools=capabilities.tools,
)
identity_key = generate_ed25519_key()
if code is None:
# THE CODE GOES FROM HERE TO THE BROWSER, NOT THE OTHER WAY.
transport = _register_and_wait(
url,
capabilities,
access_code=access_code,
client_factory=client_factory,
allow_remote=allow_remote,
allow_plaintext=allow_plaintext,
signing_public_key=ed25519_public_key(identity_key),
)
if transport is None:
return None, None
else:
try:
transport = GradioAgentTransport.pair(
url,
code,
capabilities,
access_code=access_code,
client_factory=client_factory,
allow_remote=allow_remote or None,
allow_plaintext=allow_plaintext or None,
signing_public_key=ed25519_public_key(identity_key),
)
except (AgentTransportError, TypeError, ValueError) as exc:
print(f"Pairing with {url} failed: {exc}", file=sys.stderr)
return None, None
print(
f"Paired {capabilities.name!r} with {url} as {transport.agent_id}.",
file=sys.stderr,
)
# -------- Consent handshake: catalogue, approval, gate ----------
#
# Per server, and deliberately not shared. The server's signed catalogue
# is pulled fresh, the operator approves a subset bound to that
# catalogue's digest, the approved set is advertised in every signed
# snapshot, and the AllowanceGate refuses any job outside it whatever the
# server says.
try:
catalogue = transport.fetch_catalogue()
registry.register(
catalogue.server_id,
url=url,
agent_id=transport.agent_id,
# Trust-on-first-use identity for this pairing: the exact
# permission set the operator is about to be shown.
pin=catalogue.digest(),
secret=_decode(transport.credential.secret),
offer=declared,
identity_key=identity_key,
# Rejoining a server already on file. Pairing codes are one-use,
# so this path always carries fresh credentials, and the approval
# bound to the old ones is dropped and re-taken below.
replace_existing=True,
)
choice = registry.record_catalogue(catalogue.server_id, catalogue)
approved = _operator_approval(choice, approve_flag=approve_flag)
if not approved:
# Approving nothing remains a legitimate answer. Staying joined
# afterwards is not: a worker that will refuse every job from this
# server has nothing to contribute to it, and appearing online
# while refusing everything is the failure mode that is invisible
# from both ends.
print(
f"Nothing approved for {url}; not joining it.",
file=sys.stderr,
flush=True,
)
close = getattr(transport, "close", None)
if callable(close):
close()
return None, None
registry.approve(
catalogue.server_id,
OperatorApproval(
server_id=catalogue.server_id,
catalogue_digest=catalogue.digest(),
allowance=choice,
approved_at=time.time(),
method="--approve flag" if approve_flag else "operator-console",
),
)
print(
f"Operator approved for {catalogue.server_id!r}: "
f"models {list(choice.models)}, tools {list(choice.tools)}.",
file=sys.stderr,
flush=True,
)
gate = AllowanceGate(registry, catalogue.server_id)
transport.snapshot_signer = (
lambda snapshot: registry.sign_snapshot(catalogue.server_id, snapshot)
)
except (AgentTransportError, AllowanceError, TypeError, ValueError) as exc:
print(f"Consent handshake with {url} failed: {exc}", file=sys.stderr)
return None, None
return ServerLink(name=catalogue.server_id, url=url, transport=transport), gate
class _CombinedGate:
"""The right server's allowance gate for whichever server sent the job.
Each server has its own approved set, and a job must be checked against
the set approved for *its* server. A single merged gate would let a model
approved for one server be run for another, which is exactly the
disclosure the per-server consent handshake exists to prevent.
"""
def __init__(self, gates, transport) -> None:
self._gates = dict(gates)
self._transport = transport
def evaluate(self, job):
link = self._transport._link_for(job.id) # noqa: SLF001 - one module's pair
gate = self._gates.get(link.name) if link is not None else None
if gate is None:
# A job whose server cannot be identified is refused. Falling back
# to any gate would mean checking it against permissions somebody
# granted to a different server.
return Refusal(REASON_DENIED_BY_AGENT, "this job has no identifiable server")
return gate.evaluate(job)
def _advertised_mode(harness_name: str) -> str:
"""Which of the two agent modes a harness actually is.
Four harness names exist internally and two of them are the same thing to
a person choosing a worker: whether the model gets to see a tool's result
before deciding the next call. That is the distinction worth advertising,
so it is the one that crosses the wire.
"""
return MODE_SIMPLE if harness_name in {"simple", "structured"} else MODE_AGENT
def main(
argv: list[str] | None = None,
*,
transport: AgentTransport | None = None,
client_factory: Callable[[str], PredictClient] | None = None,
) -> int:
# No arguments at all, on a console somebody can answer: open the screen
# rather than print an argument grammar. See distinct_agent/desktop.py for
# why that is the right answer for this program and the wrong one for most.
#
# Guarded on the injection points as well as on argv: a test or an embedder
# that hands in its own transport has been specific about what it wants,
# and reading its terminal to decide would make that unpredictable.
if argv is None and transport is None and client_factory is None:
from .desktop import run as run_desktop
from .desktop import wants_desktop
if wants_desktop():
return run_desktop()
args = build_parser().parse_args(argv)
if client_factory is None and getattr(args, "server_api", "gradio") == "http":
# Chosen, never inferred. A silent fall back from one transport to the
# other would turn "this server speaks something else" into a mystery
# timeout, and the two are not interchangeable from the server's side.
from .httpapi import HttpPredictClient
client_factory = HttpPredictClient
if args.prepare_rlm:
try:
prepare_rlm_runtime()
except (ImportError, OSError, RlmUnavailable, RuntimeError, ValueError) as exc:
print(f"RLM runtime preparation failed: {exc}", file=sys.stderr)
return 2
print("RLM runtime prepared; cached-only verification passed.")
return 0
if args.server and args.demo_prompt is not None:
print("--demo-prompt is a local smoke test and cannot be used with --server.", file=sys.stderr)
return 2
if args.demo_prompt is not None and not args.demo_prompt.strip():
# Checked here rather than left to the protocol, which raises. An
# operator who typed an empty string got a traceback ending in
# "prompt cannot be empty", which is the right fact delivered as a
# crash. Argparse-level mistakes deserve argparse-level answers.
print("--demo-prompt needs something to ask; it was empty.", file=sys.stderr)
return 2
if args.list_servers:
return _list_servers(args)
if args.forget_server:
return _forget_servers(args)
if args.pair and not args.server and transport is None:
print("--pair requires --server.", file=sys.stderr)
return 2
# Checked here, before anything is downloaded. A worker told to join two
# servers with one code has been mis-invoked, and finding that out after a
# four-hundred-megabyte download is a waste of somebody's evening.
if args.server and transport is None and _server_pairs(args.server, args.pair) is None:
return 2
if not 1 <= args.queue_capacity <= 64:
print("--queue-capacity must be between 1 and 64.", file=sys.stderr)
return 2
if args.poll_interval is not None and (
not math.isfinite(args.poll_interval) or not 0 < args.poll_interval <= 60
):
print("--poll-interval must be finite and between 0 and 60 seconds.", file=sys.stderr)
return 2
if args.model:
print(
"--model ID=PATH has been removed. Weights are no longer configured by "
"path: this worker resolves each model to its pinned repository and "
"revision, downloads it on first use, and verifies it against the "
"published digest before loading it. A path could do none of that.\n"
"Use --models ID[,ID...] to choose which models to offer, and "
"--list-models to see what is available.",
file=sys.stderr,
)
return 2
# The operator's library approval, applied before anything reads the
# registry. It is written into the environment the approval layer already
# reads rather than into a second mechanism beside it: one place decides
# what is approved, and it is still a human naming it.
if args.tools:
os.environ[APPROVED_TOOLS_ENV_VAR] = ",".join(_requested_model_ids(args.tools))
if args.list_tools:
offered = set(advertised_tool_refs())
named = {
f"{ref.tool_id}@{ref.version}" for ref in operator_policy().approved
}
def show(ref: str, kind: str) -> None:
# Three states, not two. A ref the operator approved and did not
# get is not "not approved", and telling them it is sends them to
# fix the one thing that is already right.
if ref in offered:
mark = "offered"
elif ref in named:
mark = "approved, not installed"
else:
mark = "not approved"
print(f" {ref:<28} {kind:<6} {mark}")
print("Library members, approvable one by one or all at once with 'local':")
for spec in sorted(INSTALLABLE_SPECS, key=lambda item: (item.kind, item.tool_id)):
show(f"{spec.tool_id}@{spec.version}", spec.kind)
# The next two groups are deliberately outside the 'local' bundle:
# each one has to be named in full. distinct_tools' module docstring
# says why, and the short version is that one word should not approve
# a skill read off the disk or a tool that writes files.
print("\nSkills from the on-disk repository, each named in full to approve:")
if REPOSITORY_SPECS:
for spec in sorted(REPOSITORY_SPECS, key=lambda item: item.tool_id):
show(f"{spec.tool_id}@{spec.version}", spec.kind)
else:
print(f" none available: {REPOSITORY_STATUS}")
print("\nMCP servers, each tool named in full to approve:")
try:
from distinct_mcp import FILE_ROOT_ENV_VAR, catalogue_refs
except ImportError as exc: # pragma: no cover - part of this repository
print(f" none available: distinct_mcp could not be imported ({exc})")
else:
for ref in catalogue_refs():
show(ref, "mcp")
print(
f" the files.* tools also need {FILE_ROOT_ENV_VAR} set to the one "
"directory they may use; there is no default"
)
notes = registry_notes()
if notes:
print("\nApproved but not installed, so every call to it fails closed:")
for note in notes:
print(f" {note}")
if not offered:
print(
"\nNothing is approved, so every tool call fails closed. Pass "
"--tools local to offer every member that opens no socket, or "
"name exact refs."
)
return 0
offerable, withheld = describe_catalogue()
if args.list_models:
print("Models this worker can offer (weights fetched on demand):")
for line in offerable:
print(f" {line}")
if withheld:
print("\nWithheld, and why:")
for line in withheld:
print(f" {line}")
return 0
# The catalogue, not the filesystem, decides what may be offered. A model
# is advertisable when its repository, revision and digest are all pinned;
# the file itself arrives the first time a run needs it.
available = {manifest.id: manifest for manifest in fetchable_manifests()}
requested = _requested_model_ids(args.models)
if requested:
unknown = [model_id for model_id in requested if model_id not in available]
if unknown:
print(
f"Not offerable: {', '.join(unknown)}. Run --list-models to see what "
"this worker can offer and why anything is withheld.",
file=sys.stderr,
)
return 2
manifests = tuple(available[model_id] for model_id in requested)
else:
manifests = tuple(available.values())
if not manifests:
print(
"No model in the catalogue has a pinned repository, revision and digest, "
"so this worker has nothing it can honestly offer.",
file=sys.stderr,
)
return 2
if withheld:
# Printed, not filtered away in silence: an operator who expected a
# model to be here should read why it is not on the same screen.
print("Models withheld from this worker:", file=sys.stderr)
for line in withheld:
print(f" {line}", file=sys.stderr)
cache_directory = args.cache_dir or args.cache_dir_legacy or shared_cache_directory()
try:
weights = WeightsCache(
cache_directory,
limit_bytes=cache_bytes(args.model_cache_gb),
manifests=manifests,
)
models = weights.installed()
except ModelVerificationError as exc:
print(f"Model verification failed: {exc}", file=sys.stderr)
return 2
except (OSError, TypeError, ValueError) as exc:
print(f"Invalid weights cache configuration: {exc}", file=sys.stderr)
return 2
if args.demo_runner or args.demo_prompt is not None:
# A smoke test performs no inference, so it must not download hundreds
# of megabytes to prove that. The deterministic runner never opens the
# path, and the cache is withheld from the executor so nothing can.
models = tuple(
DiscoveredModel(manifest, Path(manifest.filename), 0, False)
for manifest in manifests
)
weights = None
usage_source = weights
pending = (
[manifest.id for manifest in manifests if not weights.is_cached(manifest)]
if usage_source is not None
else []
)
if usage_source is not None:
usage = usage_source.usage()
print(
f"Weights cache: {cache_directory} "
f"({usage.used_bytes / (1024**3):.2f} of {args.model_cache_gb:g} GB used). "
f"Verified and ready: {', '.join(usage.model_ids) or 'none'}. "
f"Fetched on first use: {', '.join(pending) or 'none'}.",
flush=True,
)
# DOWNLOAD AT START-UP, NOT MID-RUN.
#
# The operator's gesture is "I will host these models"; the download is
# what that gesture means, so it happens when they make it. Deferring it
# to the first job puts a several-hundred-megabyte wait inside somebody
# else's request, on a worker that advertised itself as ready, and if it
# fails it fails as a failed run rather than as a worker that could not
# start.
# CHECK FOR THE THING THAT RUNS THE MODEL BEFORE FETCHING THE MODEL.
#
# The runner was resolved after this block, so a machine without
# llama.cpp downloaded nine gigabytes of weights, verified both digests,
# printed its sandbox report, and only then said it had nothing to run
# them with. Every second of that was avoidable: whether `llama-server`
# exists is knowable before the first byte.
if not (args.demo_runner or args.demo_prompt is not None):
problem = _runner_problem(args)
if problem:
print(problem, file=sys.stderr)
return 2
if pending and not args.fetch_on_demand and not (args.demo_runner or args.demo_prompt is not None):
print(
f"Fetching {len(pending)} model(s) now: {', '.join(pending)}. "
"This happens once per model.",
flush=True,
)
for manifest in manifests:
if manifest.id not in pending:
continue
report = _download_reporter(manifest.id)
try:
weights.ensure(manifest, progress=report)
except ModelVerificationError as exc:
print(f"Weights verification failed: {exc}", file=sys.stderr)
return 2
except WeightsUnavailable as exc:
print(
f"Could not obtain {manifest.id}: {exc}\n"
"Nothing was substituted. Fix the network or drop this model "
"from --models, then start the worker again.",
file=sys.stderr,
)
return 2
print(f" {manifest.id}: verified against its published digest.", flush=True)
models = weights.installed()
model_id = args.model_id or manifests[0].id
if model_id not in {manifest.id for manifest in manifests}:
print(f"Requested model is not offered by this worker: {model_id}", file=sys.stderr)
return 2
# Say what the OS will and will not enforce, every start-up. Silence here
# would let an operator assume the inference child is contained when on
# most hosts it is not.
request_report = describe_request_isolation()
if args.no_request_sandbox:
print(
"Per-request sandbox: DISABLED by --no-request-sandbox. Every request "
"runs in this worker's own directories, and what one request leaves "
"behind is visible to the next.",
file=sys.stderr,
flush=True,
)
elif request_report.mechanism == "none":
print(
f"Per-request sandbox: NONE. {request_report.reason} Each request still "
"gets a fresh working directory and a cleared environment, but nothing "
"at the OS layer stops it reading what your account can read.",
file=sys.stderr,
flush=True,
)
else:
# This used to say each request got "a new process, no network and a
# cleared environment". It did not. `RequestSandbox.spawn`, the method
# that applies those, has no caller outside the tests: the worker
# builds a sandbox, takes a scratch directory and a report from it, and
# runs inference through the ordinary runner. The directory and its
# erasure were real; the process, the network denial and the cleared
# environment were a sentence.
#
# A start-up banner is the one place an operator looks to find out what
# is protecting their machine, so it says what is applied and names
# what is not. Applying the rest is a change to the run path, not to
# this line, and until somebody makes it this must not read as though
# somebody has.
print(
f"Per-request sandbox: {request_report.mechanism} is available, and is "
"used for the fresh working directory each request gets, erased when "
"the request ends. It does NOT yet wrap the model process: inference "
"runs in the ordinary way, so nothing at the OS layer stops it "
"reading what your account can read or reaching the network.",
file=sys.stderr,
flush=True,
)
isolation_report = describe_support()
if isolation_report.applied:
# flush=True: stdout is block-buffered when redirected to a file or a
# log, and a security-relevant status line that appears minutes late,
# or not at all if the process is killed, is worse than useless.
print(
f"Process restriction: {isolation_report.mechanism} active "
f"(max {isolation_report.limits.get('max_processes')} process, "
f"{isolation_report.limits.get('memory_bytes', 0) // (1024 * 1024)} MiB, "
"killed with the agent). Containment and resource control, not a "
"security sandbox: the model process still runs under your user account.",
flush=True,
)
else:
print(
f"Process restriction: NONE. {isolation_report.reason} "
"The local model process runs with your full user privileges.",
file=sys.stderr,
flush=True,
)
# Detected before the runner is built: the server runner uses it to price
# model loading, which is charged to residency rather than to a run.
energy_meter = detect_energy_meter()
if not energy_meter.available and not args.without_energy_measurement:
# A WORKER THAT CANNOT MEASURE DOES NOT START.
#
# This project's claim is that the energy figure beside an answer was
# metered on the machine that produced it. A worker with no readable
# counter cannot honour that: every run it serves would carry an
# absence, and a network of those is a network that measures nothing
# while looking like one that does.
#
# So it stops here, before pairing, before downloading nine gigabytes
# of weights, and before anybody has been told this machine is
# available. And it stops with the whole list rather than a verdict,
# because "no energy counter" is not something an operator can act on
# and "the EMI driver is absent but the Energy Meter counter would work
# if you were not in a container" is.
print(_no_meter_message(), file=sys.stderr, flush=True)
return 2
if energy_meter.available:
print(f"Energy: {energy_meter.scope} via {energy_meter.provider}.",
file=sys.stderr, flush=True)
else:
print(
"Energy: NOT MEASURED on this machine, and you asked to continue anyway.\n"
" Every run this worker serves will report no energy figure. It is\n"
" counted as a run and never as zero, and the server shows the absence.",
file=sys.stderr, flush=True,
)
deno_status = check_deno_sandbox() if args.allow_generated_code else None
generated_policy = GeneratedExecutionPolicy(
enabled=args.allow_generated_code,
status=deno_status,
)
if args.allow_generated_code and not generated_policy.allowed:
print(
f"Generated execution remains disabled: {deno_status.reason if deno_status else 'not checked'}",
file=sys.stderr,
)
return 2
demo_runner = bool(args.demo_runner or args.demo_prompt is not None)
if demo_runner:
runner = DeterministicDemoRunner()
elif args.llama_server or not args.llama_cli:
# `--llama-server` with no value means "look for it"; the finder is
# what a bare flag asks for, and passing True into a path would have
# been a confusing failure two screens later.
if args.llama_server is True:
args.llama_server = ""
# llama-server is the default. It returns clean output with no chat
# furniture, reports real token counts, streams, keeps the model
# resident across a queue, and speaks the OpenAI-compatible protocol
# DSPy expects. --llama-cli remains for operators who want the old
# per-request subprocess.
runner = LlamaServerRunner(
args.llama_server or _resolve_llama_server(args),
energy_meter=energy_meter,
)
# SAY WHICH DEVICE WILL DO THE WORK, AND SAY THAT IT WAS MEASURED.
#
# A volunteer with a graphics card who sees it sitting idle assumes
# something is broken. On the laptop this was developed on the card
# really is the wrong device -- 0.50 tokens per second against 2.93 on
# the CPU, because a partly offloaded model pays a round trip per
# token -- and the first time a model is loaded the worker spends a
# few minutes finding that out. Both are worth saying out loud.
runner._notify = lambda message: print(f" {message}", file=sys.stderr, flush=True)
print(f"Runtime: {runner.executable}", file=sys.stderr, flush=True)
if not runner.available:
print(
"llama-server was not found. Supply --llama-server, or --llama-cli "
"for the legacy per-request runner, or --demo-runner.",
file=sys.stderr,
)
return 2
else:
runner = LlamaCppRunner("" if args.llama_cli is True else args.llama_cli)
if not runner.available:
print(
"llama-cli was not found; supply --llama-cli or use --demo-runner.",
file=sys.stderr,
)
return 2
harness_name = args.harness or ("structured" if demo_runner else "dspy-rlm")
try:
inference_harness = build_inference_harness(harness_name)
readiness_check = getattr(inference_harness, "check_ready", None)
if callable(readiness_check):
readiness_check()
except (RlmUnavailable, RuntimeError, ValueError) as exc:
if harness_name == "dspy-rlm" and args.harness is None:
# The default harness is unavailable on this machine. Degrade to
# the structured harness with the reason stated up front, rather
# than refusing to start. An operator who explicitly asked for
# --harness dspy-rlm still gets the hard failure below, because an
# explicit request must never be silently substituted.
print(
f"DSPy RLM harness unavailable: {exc}\n"
"Falling back to the structured tool harness for this worker. "
"Run --prepare-rlm (with Deno available) to enable the RLM harness.",
file=sys.stderr,
flush=True,
)
harness_name = "structured"
inference_harness = build_inference_harness(harness_name)
else:
print(f"Inference harness unavailable: {exc}", file=sys.stderr)
return 2
requested_agent_id = args.agent_id or (
"" if args.server and transport is None else _default_agent_id()
)
try:
capabilities = AgentCapabilities(
agent_id=requested_agent_id,
name=args.name,
os=platform.system() or "unknown",
arch=platform.machine() or "unknown",
cpu=platform.processor() or "unknown",
ram_gb=_ram_gb(),
models=tuple(manifest.id for manifest in manifests),
tools=advertised_tool_refs(),
energy_provider=energy_meter.provider,
# Advertised, so the person choosing a worker is choosing how
# their request will be answered rather than finding out
# afterwards. One mode, because a worker runs one harness: the
# honest advertisement is what it will actually do, not what it
# could be reconfigured to do.
modes=(_advertised_mode(harness_name),),
queue_capacity=args.queue_capacity,
max_concurrency=1,
)
except (TypeError, ValueError) as exc:
print(f"Invalid agent capabilities: {exc}", file=sys.stderr)
return 2
# The guard and the deny list are wired only on the real server path: a
# local smoke test has no user identity to protect anyone from. The deny
# list is agent-local, in memory, and never leaves this machine.
deny_list = None
if args.server:
from .guard import DenyList
deny_list = DenyList()
guard_model_id = None
if args.guard_model:
if args.guard_model not in {manifest.id for manifest in manifests}:
print(
f"--guard-model {args.guard_model!r} is not among the models this "
"worker offers.",
file=sys.stderr,
)
return 2
guard_model_id = args.guard_model
# The read-only view a sandboxed request is given: the bundled llama.cpp
# runtime and the interpreter's own tree, and nothing else. It is shared
# and immutable, so no per-request copy is ever made.
runtime_paths = [
path
for path in (
Path(__file__).resolve().parent.parent / "runtime",
Path(sys.prefix),
)
if path.exists()
]
executor = LocalInferenceExecutor(
models=models,
manifests=manifests,
weights=weights,
runner=runner,
energy_meter=energy_meter,
generated_execution=generated_policy,
tool_harness=inference_harness,
guard_enabled=bool(args.server) and not args.no_guard,
guard_model_id=guard_model_id,
deny_list=deny_list,
request_isolation=not args.no_request_sandbox,
sandbox_runtime_paths=runtime_paths,
)
gate = None
from distinct_protocol.handshake import Allowance
# DIRECT MODE, ADDED TO WHATEVER ELSE THIS WORKER IS DOING.
#
# DIRECT MODE HAS BEEN REMOVED, AND THE REASON IS THE WHOLE ACCESS MODEL.
#
# `--direct` made the worker hold its own listening socket and answer a
# browser with no server in between. It had no authentication of any kind:
# a red-team pass submitted a job and read the answer back with no Hugging
# Face account, no access code and no server, using nothing but curl. On
# the loopback default that meant any local user on a shared machine; with
# `--direct 0.0.0.0` it meant anyone on the network; behind a tunnel it
# meant the internet.
#
# That was not a defect in the code, which said plainly what it did. It was
# a second front door, built for the case where the operator and the user
# are the same person, and it is incompatible with the requirement that
# nobody reaches a volunteer's machine without redeeming an access code and
# signing in. A worker now has exactly one way in: the server it paired
# with, pulling work it agreed to accept.
#
# `tests/test_loopback_policy.py` no longer exempts any listener, which is
# the check that stops this coming back by accident.
if transport is None:
if args.server:
# ONE MACHINE, SEVERAL SERVERS, ONE QUEUE.
#
# Each server is paired and consented to separately: its own
# catalogue, its own operator approval, its own credentials and its
# own allowance gate. That separation is the point. Joining a
# second server must not disclose to it what the first was
# offered, and approving a model for one is not approving it for
# the other.
#
# What is shared is the machine: one queue, one capacity figure,
# one set of weights. See distinct_agent/fanout.py for why running
# a loop per server instead would have every server believing it
# had the whole machine to itself.
pairs = _server_pairs(args.server, args.pair)
if pairs is None:
return 2
# One registry for the machine, holding a separate record, secret
# and approval per server. Persisted, so a worker restarted
# tomorrow still knows which servers it agreed to serve and on
# what terms.
server_registry = ServerRegistry(
Allowance(
models=tuple(manifest.id for manifest in manifests),
tools=capabilities.tools,
),
path=_servers_file(args.servers_file),
)
# One code for this worker, generated here and given to every
# server it joins. Printed once, loudly, because it is the only
# time anybody sees it: the servers keep a digest and this process
# does not write it anywhere.
access_code = new_access_code()
# PRINTED AFTER THE JOIN, NOT BEFORE IT.
#
# This used to be announced here, which put two high-entropy codes
# on screen at once: the access code to share, and then the claim
# code to enter. Two codes visible together is the same confusion
# this whole change exists to remove -- the person scrolls back,
# finds *a* code, and uses it for the wrong thing.
#
# One at a time, in the order they are needed: the claim code while
# the worker waits, and the access code once it is actually running
# and there is something worth sharing.
# Kept so the dashboard can show it. It is printed once and stored
# nowhere on disk, so a screen that could not show it would leave
# somebody who scrolled past it with no way to get it back short of
# restarting the worker, which issues a different code and retires
# this one.
args._access_code = access_code
links: list[ServerLink] = []
gates: dict[str, AllowanceGate] = {}
for url, code in pairs:
link, gate_for = _join_server(
url,
code,
capabilities=capabilities,
manifests=manifests,
client_factory=client_factory,
allow_remote=args.allow_remote_server,
allow_plaintext=args.allow_plaintext_server,
access_code=access_code,
approve_flag=args.approve,
registry=server_registry,
)
if link is None:
# Refused, unreachable or unapproved. Named by
# _join_server; the worker carries on with the rest,
# because one server declining is not a reason to abandon
# the others.
continue
links.append(link)
gates[link.name] = gate_for
capabilities = replace(capabilities, agent_id=link.transport.agent_id)
if not links:
print(
"Not starting: no server was joined. Nothing is wrong with your "
"machine and nothing was left running.",
file=sys.stderr,
flush=True,
)
return 2
# `--server` is repeatable, so this is a list. One worker can serve
# several servers with the same code, and the invitation names them
# all rather than picking one and quietly misdirecting the people
# sent to the others.
print(
announce_access_code(access_code, servers=tuple(args.server or ())),
file=sys.stderr,
flush=True,
)
if len(links) == 1:
# One server: the transport it paired with, unwrapped. The
# fan-out is a real object with real behaviour and there is no
# reason to put it in the path of the common case.
transport = links[0].transport
gate = gates[links[0].name]
else:
transport = MultiServerTransport(
links, report=lambda message: print(message, file=sys.stderr, flush=True)
)
gate = _CombinedGate(gates, transport)
print(
f"Serving {len(links)} servers from one queue:\n"
+ transport.describe(),
file=sys.stderr,
flush=True,
)
else:
offer = None
if args.demo_prompt is not None:
offer = JobSpec(
id=f"demo-{uuid.uuid4().hex}",
session_id="demo-session",
conversation_id="demo-conversation",
parent_job_id=None,
prompt=args.demo_prompt,
model_id=model_id,
allowed_tools=(),
target_agent_id=capabilities.agent_id,
created_at=time.time(),
)
transport = ConsoleTransport(offer)
else:
transported_agent_id = getattr(transport, "agent_id", None)
if args.agent_id and transported_agent_id and args.agent_id != transported_agent_id:
print("--agent-id does not match the injected transport credential.", file=sys.stderr)
return 2
if isinstance(transported_agent_id, str) and transported_agent_id:
capabilities = replace(capabilities, agent_id=transported_agent_id)
elif not capabilities.agent_id:
print("Injected transport must expose an agent id.", file=sys.stderr)
return 2
# One cadence, and the agent honours the one the server advertised, so the
# rate the UI describes is the rate that actually happens. An explicit
# --poll-interval still wins, for operators who need to throttle.
if args.poll_interval is not None:
poll_interval = args.poll_interval
elif isinstance(transport, GradioAgentTransport):
poll_interval = transport.poll_interval_s
else:
poll_interval = DEFAULT_POLL_INTERVAL_SECONDS
worker = WorkerLoop(
capabilities=capabilities,
transport=transport,
executor=executor,
queue=InMemoryJobQueue(
capacity=args.queue_capacity,
max_active=1,
agent_id=capabilities.agent_id,
),
settings=WorkerLoopSettings(poll_interval_seconds=poll_interval),
# The allowlist boundary. Never None on the real server path: even a
# worker whose operator approved nothing gets a gate, and that gate
# refuses everything, which is the designed meaning of "no approval".
gate=gate,
# Agent-local refusals, keyed on the server-derived pseudonymous user
# key. Checked before the gate: whether this machine serves this
# person is prior to what they asked for.
deny_list=deny_list,
)
if isinstance(transport, ConsoleTransport):
worker.run_cycle()
if args.demo_prompt is None:
return 0
deadline = time.monotonic() + 30.0
while not transport.done.is_set() and time.monotonic() < deadline:
worker.run_cycle()
transport.done.wait(poll_interval)
worker.shutdown()
return 0 if transport.done.is_set() else 1
if args.dashboard:
return _run_with_dashboard(
worker,
transport=transport,
capabilities=capabilities,
args=args,
containment=isolation_report,
request_sandbox=request_report,
)
exit_code = 0
try:
worker.run_forever()
except KeyboardInterrupt:
pass
except AgentTransportError as exc:
print(f"Agent transport stopped: {exc}", file=sys.stderr)
exit_code = 1
finally:
worker.shutdown()
close = getattr(transport, "close", None)
if callable(close):
close()
return exit_code
def _run_with_dashboard(
worker,
*,
transport,
capabilities,
args,
containment,
request_sandbox,
) -> int:
"""Run the same worker, with the screen drawing instead of a log scrolling.
The worker runs on a thread and the screen draws on this one, rather than
the other way round, for one reason: `Ctrl+C` is delivered to the main
thread. Putting the worker there and the screen on a background thread
would mean an interrupt landing inside the poll loop, which is the one
place a half-finished job could be left in an unclear state.
Every exit path stops the worker and closes the transport, including the
path where the screen itself fails: a dashboard that crashed and left a
worker serving invisibly would be worse than no dashboard.
"""
from .dashboard import ActivityLog, WorkerView
try:
from .tui import DashboardApp
except ImportError as exc: # pragma: no cover - only without the dependency
# A traceback about a module nobody asked for is not an answer. The
# worker itself does not need Textual, so this is the one place its
# absence matters and the one place worth explaining it.
raise SystemExit(
f"The live status screen needs the 'textual' package ({exc}). "
"Install this project's dependencies with `pip install -e .` in the "
"folder you cloned, or drop --dashboard to run with a scrolling log."
) from exc
log = getattr(sys.stderr, "_distinct_activity_log", None)
if log is None:
log = ActivityLog()
base = WorkerView(
name=capabilities.name,
agent_id=capabilities.agent_id,
servers=tuple(args.server or ()),
connected=True,
platform_name=f"{platform.system()} {platform.release()} ({capabilities.arch})".strip(),
cpu=capabilities.cpu,
ram_gb=capabilities.ram_gb,
models=tuple(capabilities.models),
tools=tuple(capabilities.tools),
queue_capacity=capabilities.queue_capacity,
energy_provider=capabilities.energy_provider,
containment=_report_sentence(containment),
request_sandbox=_report_sentence(request_sandbox),
access_code=getattr(args, "_access_code", "") or "",
)
stop = threading.Event()
failure: list[str] = []
def serve() -> None:
try:
worker.run_forever(stop_event=stop)
except AgentTransportError as exc:
failure.append(f"The connection stopped: {exc}")
except Exception as exc: # noqa: BLE001 - reported on the screen
failure.append(f"The worker stopped: {type(exc).__name__}: {exc}")
finally:
stop.set()
thread = threading.Thread(target=serve, name="distinct-worker", daemon=True)
thread.start()
# From here the terminal belongs to the application, so the log stops
# echoing to it and keeps the lines for the activity pane instead.
log.detach()
try:
DashboardApp(worker=worker, base=base, activity=log, stop=stop).run()
except KeyboardInterrupt:
pass
finally:
stop.set()
thread.join(timeout=10)
worker.shutdown()
close = getattr(transport, "close", None)
if callable(close):
close()
# Restore the console before anything else prints, or the summary below
# is written into a screen that is no longer being redrawn.
if getattr(sys.stderr, "_distinct_activity_log", None) is not None:
sys.stderr = sys.__stderr__
for message in failure:
print(message, file=sys.stderr)
return 1 if failure else 0
def _report_sentence(report) -> str:
"""One line from an isolation report, in the report's own words.
Reads `mechanism` and `reason` rather than reformatting them, because those
strings are already written to be honest about what is and is not applied,
and a dashboard that paraphrased them would be a second, laxer description
of the same thing.
"""
if report is None:
return "unknown"
mechanism = str(getattr(report, "mechanism", "") or "")
applied = bool(getattr(report, "applied", False))
reason = str(getattr(report, "reason", "") or "")
if applied and mechanism:
return mechanism
if mechanism and mechanism != "none":
return f"{mechanism} available, not applied"
return reason or "none"
def _operator_approval(choice, *, approve_flag: bool) -> bool:
"""Decide whether the operator approves the pending choice set.
``--approve`` is an explicit gesture given on the command line by the
person starting the worker; it approves the full choice set. Otherwise, a
console that can ask, asks. A console that cannot ask approves nothing,
because auto-accepting on a headless box is exactly the path this
handshake exists to close.
**Every path out of here says what happened and what to do about it.** The
previous version printed the question and, when the answer was empty or the
console could not be read, printed "Nothing approved" on the same line. An
operator saw their prompt collide with a refusal and had no way to tell
whether they had pressed the wrong key, whether the worker had failed, or
what to type next. A refusal is a legitimate outcome; an unexplained one is
not.
"""
if not choice:
print(
"\nNothing to approve. The server's catalogue and this worker's offer "
"have no models or library members in common, so there is no set to "
"agree on. Check --list-models against the server's catalogue.",
file=sys.stderr,
flush=True,
)
return False
if approve_flag:
return True
if not sys.stdin.isatty():
print(
"\nApproval was NOT requested, because this console cannot be read "
"(the worker is running without an interactive terminal, for example "
"from a service, a scheduled task, or a batch file whose input is "
"redirected).\n"
"Approving by default on a machine nobody is watching is exactly what "
"this handshake exists to prevent, so the worker approved nothing.\n"
"To run unattended, restart it with --approve, which approves the set "
"below for this server:\n"
f" models: {', '.join(choice.models) or 'none'}\n"
f" library (tools and skills): {', '.join(choice.tools) or 'none'}",
file=sys.stderr,
flush=True,
)
return False
print(
"\nThis server may ask this worker to run:\n"
f" models: {', '.join(choice.models) or 'none'}\n"
f" library (tools and skills): {', '.join(choice.tools) or 'none'}\n"
"Approving binds this exact set to this server's current catalogue. The "
"worker refuses anything outside it, whatever the server later sends.\n"
"Answer 'n', or press Enter, to approve nothing and stop.",
file=sys.stderr,
flush=True,
)
try:
answer = input("Approve this exact set? [y/N] ")
except (EOFError, KeyboardInterrupt):
print(
"\nNo answer was given, so nothing was approved. Start the worker "
"again and answer 'y', or pass --approve.",
file=sys.stderr,
flush=True,
)
return False
if answer.strip().lower() in {"y", "yes"}:
return True
print(
"\nDeclined. Nothing was approved, so this worker will advertise nothing "
"and refuse every job. That is a working, safe state, not an error: it is "
"simply a worker that has agreed to do nothing yet.\n"
"To contribute, start it again and answer 'y', or pass --approve. To offer "
"a narrower set, restart with --models naming only the models you want.",
file=sys.stderr,
flush=True,
)
return False
def _download_reporter(model_id: str) -> Callable[[float, str], None]:
"""Print download progress in tenths, so a long fetch is visibly alive.
Bound outside the loop it is used in: a closure over the loop variable
would report every model under the last model's name.
"""
last = [-1]
def report(fraction: float, message: str) -> None:
percent = int(fraction * 100)
if percent >= last[0] + 10:
last[0] = percent
print(f" {model_id}: {percent}%", flush=True)
return report
def _requested_model_ids(values: Iterable[str]) -> tuple[str, ...]:
"""Parse ``--models`` into an ordered, de-duplicated tuple of ids.
Accepts the flag repeated and comma-separated within one value, because
both read naturally and refusing one of them is an argument with the
operator rather than a safety property.
"""
seen: list[str] = []
for value in values or ():
for part in str(value).split(","):
model_id = part.strip()
if model_id and model_id not in seen:
seen.append(model_id)
return tuple(seen)
def _default_agent_id() -> str:
host = re.sub(r"[^A-Za-z0-9_.-]+", "-", platform.node()).strip("-.")
return (host or "local-agent")[:128]
def _ram_gb() -> float:
try:
if os.name == "nt":
class MemoryStatus(ctypes.Structure):
_fields_ = [
("length", ctypes.c_ulong),
("memory_load", ctypes.c_ulong),
("total_phys", ctypes.c_ulonglong),
("avail_phys", ctypes.c_ulonglong),
("total_page_file", ctypes.c_ulonglong),
("avail_page_file", ctypes.c_ulonglong),
("total_virtual", ctypes.c_ulonglong),
("avail_virtual", ctypes.c_ulonglong),
("avail_extended_virtual", ctypes.c_ulonglong),
]
status = MemoryStatus()
status.length = ctypes.sizeof(MemoryStatus)
if ctypes.windll.kernel32.GlobalMemoryStatusEx(ctypes.byref(status)):
return round(status.total_phys / (1024**3), 2)
pages = os.sysconf("SC_PHYS_PAGES")
page_size = os.sysconf("SC_PAGE_SIZE")
return round((pages * page_size) / (1024**3), 2)
except (AttributeError, OSError, ValueError):
return 0.0
if __name__ == "__main__": # pragma: no cover
raise SystemExit(main())
|