diff --git a/MODULE.bazel.lock b/MODULE.bazel.lock index 40ea6b08f7..4a464e5aa9 100644 --- a/MODULE.bazel.lock +++ b/MODULE.bazel.lock @@ -1311,7 +1311,7 @@ "usagesDigest": "xr+U7navw2+SHogogmHRboKLbXAnkZfsctPQuyjmArw=", "recordedFileInputs": { "@@//doc/requirements.txt": "69647040f4c4bdd93fd8369b245316b08cfabd17a23693d833081b5785c0f131", - "@@//tools/env/pip3/requirements.txt": "d2f5e829afd46ab2db03085eebc6f2271a42ab5cc46a4d0973ba2c9a537862b9", + "@@//tools/env/pip3/requirements.txt": "4fe27a162735033dae27bd4edee0747fc0f4e7b15f45287ad45b9fe6589b29a4", "@@//tools/lint/python/requirements.txt": "c6eb43e931b4200ae71e62fc65bc78d72224495a1648dbc0a565f09571724bd8", "@@rules_fuzzing+//fuzzing/requirements.txt": "ab04664be026b632a0d2a2446c4f65982b7654f5b6851d2f9d399a19b7242a5b", "@@rules_python+//tools/publish/requirements_darwin.txt": "095d4a4f3d639dce831cd493367631cd51b53665292ab20194bac2c0c6458fa8", @@ -3722,6 +3722,15 @@ ] } }, + "scion_python_deps_312_bcrypt": { + "repoRuleId": "@@rules_python+//python/private/pypi:whl_library.bzl%whl_library", + "attributes": { + "dep_template": "@scion_python_deps//{name}:{target}", + "python_interpreter_target": "@@rules_python++python+python_3_12_host//:python", + "repo": "scion_python_deps_312", + "requirement": "bcrypt==5.0.0 --hash=sha256:046ad6db88edb3c5ece4369af997938fb1c19d6a699b9c1b27b0db432faae4c4 --hash=sha256:0c418ca99fd47e9c59a301744d63328f17798b5947b0f791e9af3c1c499c2d0a --hash=sha256:0c8e093ea2532601a6f686edbc2c6b2ec24131ff5c52f7610dd64fa4553b5464 --hash=sha256:0cae4cb350934dfd74c020525eeae0a5f79257e8a201c0c176f4b84fdbf2a4b4 --hash=sha256:137c5156524328a24b9fac1cb5db0ba618bc97d11970b39184c1d87dc4bf1746 --hash=sha256:200af71bc25f22006f4069060c88ed36f8aa4ff7f53e67ff04d2ab3f1e79a5b2 --hash=sha256:212139484ab3207b1f0c00633d3be92fef3c5f0af17cad155679d03ff2ee1e41 --hash=sha256:2b732e7d388fa22d48920baa267ba5d97cca38070b69c0e2d37087b381c681fd --hash=sha256:35a77ec55b541e5e583eb3436ffbbf53b0ffa1fa16ca6782279daf95d146dcd9 --hash=sha256:38cac74101777a6a7d3b3e3cfefa57089b5ada650dce2baf0cbdd9d65db22a9e --hash=sha256:3abeb543874b2c0524ff40c57a4e14e5d3a66ff33fb423529c88f180fd756538 --hash=sha256:3ca8a166b1140436e058298a34d88032ab62f15aae1c598580333dc21d27ef10 --hash=sha256:3cf67a804fc66fc217e6914a5635000259fbbbb12e78a99488e4d5ba445a71eb --hash=sha256:4870a52610537037adb382444fefd3706d96d663ac44cbb2f37e3919dca3d7ef --hash=sha256:48f753100931605686f74e27a7b49238122aa761a9aefe9373265b8b7aa43ea4 --hash=sha256:4bfd2a34de661f34d0bda43c3e4e79df586e4716ef401fe31ea39d69d581ef23 --hash=sha256:560ddb6ec730386e7b3b26b8b4c88197aaed924430e7b74666a586ac997249ef --hash=sha256:5b1589f4839a0899c146e8892efe320c0fa096568abd9b95593efac50a87cb75 --hash=sha256:5feebf85a9cefda32966d8171f5db7e3ba964b77fdfe31919622256f80f9cf42 --hash=sha256:611f0a17aa4a25a69362dcc299fda5c8a3d4f160e2abb3831041feb77393a14a --hash=sha256:61afc381250c3182d9078551e3ac3a41da14154fbff647ddf52a769f588c4172 --hash=sha256:64d7ce196203e468c457c37ec22390f1a61c85c6f0b8160fd752940ccfb3a683 --hash=sha256:64ee8434b0da054d830fa8e89e1c8bf30061d539044a39524ff7dec90481e5c2 --hash=sha256:6b8f520b61e8781efee73cba14e3e8c9556ccfb375623f4f97429544734545b4 --hash=sha256:741449132f64b3524e95cd30e5cd3343006ce146088f074f31ab26b94e6c75ba --hash=sha256:744d3c6b164caa658adcb72cb8cc9ad9b4b75c7db507ab4bc2480474a51989da --hash=sha256:79cfa161eda8d2ddf29acad370356b47f02387153b11d46042e93a0a95127493 --hash=sha256:7aeef54b60ceddb6f30ee3db090351ecf0d40ec6e2abf41430997407a46d2254 --hash=sha256:7edda91d5ab52b15636d9c30da87d2cc84f426c72b9dba7a9b4fe142ba11f534 --hash=sha256:7f277a4b3390ab4bebe597800a90da0edae882c6196d3038a73adf446c4f969f --hash=sha256:7f4c94dec1b5ab5d522750cb059bb9409ea8872d4494fd152b53cca99f1ddd8c --hash=sha256:801cad5ccb6b87d1b430f183269b94c24f248dddbbc5c1f78b6ed231743e001c --hash=sha256:83e787d7a84dbbfba6f250dd7a5efd689e935f03dd83b0f919d39349e1f23f83 --hash=sha256:89042e61b5e808b67daf24a434d89bab164d4de1746b37a8d173b6b14f3db9ff --hash=sha256:92864f54fb48b4c718fc92a32825d0e42265a627f956bc0361fe869f1adc3e7d --hash=sha256:9d52ed507c2488eddd6a95bccee4e808d3234fa78dd370e24bac65a21212b861 --hash=sha256:9fffdb387abe6aa775af36ef16f55e318dcda4194ddbf82007a6f21da29de8f5 --hash=sha256:a28bc05039bdf3289d757f49d616ab3efe8cf40d8e8001ccdd621cd4f98f4fc9 --hash=sha256:a5393eae5722bcef046a990b84dff02b954904c36a194f6cfc817d7dca6c6f0b --hash=sha256:a71f70ee269671460b37a449f5ff26982a6f2ba493b3eabdd687b4bf35f875ac --hash=sha256:b17366316c654e1ad0306a6858e189fc835eca39f7eb2cafd6aaca8ce0c40a2e --hash=sha256:baade0a5657654c2984468efb7d6c110db87ea63ef5a4b54732e7e337253e44f --hash=sha256:c2388ca94ffee269b6038d48747f4ce8df0ffbea43f31abfa18ac72f0218effb --hash=sha256:c58b56cdfb03202b3bcc9fd8daee8e8e9b6d7e3163aa97c631dfcfcc24d36c86 --hash=sha256:cde08734f12c6a4e28dc6755cd11d3bdfea608d93d958fffbe95a7026ebe4980 --hash=sha256:d79e5c65dcc9af213594d6f7f1fa2c98ad3fc10431e7aa53c176b441943efbdd --hash=sha256:d8d65b564ec849643d9f7ea05c6d9f0cd7ca23bdd4ac0c2dbef1104ab504543d --hash=sha256:db99dca3b1fdc3db87d7c57eac0c82281242d1eabf19dcb8a6b10eb29a2e72d1 --hash=sha256:dcd58e2b3a908b5ecc9b9df2f0085592506ac2d5110786018ee5e160f28e0911 --hash=sha256:dd19cf5184a90c873009244586396a6a884d591a5323f0e8a5922560718d4993 --hash=sha256:ddb4e1500f6efdd402218ffe34d040a1196c072e07929b9820f363a1fd1f4191 --hash=sha256:e3cf5b2560c7b5a142286f69bde914494b6d8f901aaa71e453078388a50881c4 --hash=sha256:ed2e1365e31fc73f1825fa830f1c8f8917ca1b3ca6185773b349c20fd606cec2 --hash=sha256:edfcdcedd0d0f05850c52ba3127b1fce70b9f89e0fe5ff16517df7e81fa3cbb8 --hash=sha256:f0ce778135f60799d89c9693b9b398819d15f1921ba15fe719acb3178215a7db --hash=sha256:f2347d3534e76bf50bca5500989d6c1d05ed64b440408057a37673282c654927 --hash=sha256:f3c08197f3039bec79cee59a606d62b96b16669cff3949f21e74796b6e3cd2be --hash=sha256:f632fd56fc4e61564f78b46a2269153122db34988e78b6be8b32d28507b7eaeb --hash=sha256:f6984a24db30548fd39a44360532898c33528b74aedf81c26cf29c51ee47057e --hash=sha256:f70aadb7a809305226daedf75d90379c397b094755a710d7014b8b117df1ebbf --hash=sha256:f748f7c2d6fd375cc93d3fba7ef4a9e3a092421b8dbf34d8d4dc06be9492dfdd --hash=sha256:f8429e1c410b4073944f03bd778a9e066e7fad723564a52ff91841d278dfc822 --hash=sha256:fc746432b951e92b58317af8e0ca746efe93e66555f1b40888865ef5bf56446b" + } + }, "scion_python_deps_312_certifi": { "repoRuleId": "@@rules_python+//python/private/pypi:whl_library.bzl%whl_library", "attributes": { @@ -4458,6 +4467,7 @@ "repo_name": "scion_python_deps", "extra_hub_aliases": {}, "whl_map": { + "bcrypt": "{\"scion_python_deps_312_bcrypt\":[{\"version\":\"3.12\"}]}", "certifi": "{\"scion_python_deps_312_certifi\":[{\"version\":\"3.12\"}]}", "charset_normalizer": "{\"scion_python_deps_312_charset_normalizer\":[{\"version\":\"3.12\"}]}", "idna": "{\"scion_python_deps_312_idna\":[{\"version\":\"3.12\"}]}", @@ -4473,6 +4483,7 @@ "urllib3": "{\"scion_python_deps_312_urllib3\":[{\"version\":\"3.12\"}]}" }, "packages": [ + "bcrypt", "certifi", "charset_normalizer", "idna", diff --git a/marketplace/db/BUILD.bazel b/marketplace/db/BUILD.bazel index bcad1f4137..69cfb32984 100644 --- a/marketplace/db/BUILD.bazel +++ b/marketplace/db/BUILD.bazel @@ -1,6 +1,9 @@ load("@rules_go//go:def.bzl", "go_library") load("//tools:go.bzl", "go_test") +# The topology generator needs the schema when it prepopulates a marketplace. +exports_files(["schema.sql"]) + go_library( name = "go_default_library", srcs = [ diff --git a/marketplace/tools/populate_marketplace.py b/marketplace/tools/populate_marketplace.py index d84f5e74ef..731900ecb2 100755 --- a/marketplace/tools/populate_marketplace.py +++ b/marketplace/tools/populate_marketplace.py @@ -1,17 +1,11 @@ #!/usr/bin/env python3 -import base64 -import json -import math -import sqlite3 -import hashlib -import bcrypt -import argparse -import re -from datetime import datetime, timedelta, timezone -from pathlib import Path +"""Rebuilds the database of a marketplace. + +The topology generator already prepopulates the database of the marketplace it +adds with -m/--marketplace, so this script is only needed to load a different set +of entries, or to reset the database of a running topology. -""" An account without a scope, or with an empty one, is the main account of its user; there can be only one per user. Assets and reservations given an "owner" are owned by the main account of that user. @@ -172,417 +166,49 @@ } """ -DEFAULT_PASSWORD = "1234" -DEFAULT_USERS = ["alice", "bob"] -DEFAULT_USER_BALANCE = 1000000 -DEFAULT_AS_BALANCE = 0 -DEFAULT_ASSET_BANDWIDTH = 1000000 -DEFAULT_ASSET_BANDWIDTH_MIN = 10 -DEFAULT_ASSET_BANDWIDTH_MAX = 1000000 -DEFAULT_ASSET_PRICE = 1 -DEFAULT_ASSET_TIME_GRANULARITY = 10 -DEFAULT_ASSET_TIME_MIN_DURATION = 10 -DEFAULT_ASSET_DURATION = timedelta(days=100) -DEFAULT_DELEGATION_RES_ID_LIMIT = 1000 -# The flyover carries the bandwidth as a 10 bit codepoint into the points -# published by the AS, so there is one point per codepoint. The values must be -# the ones the router decodes, see ConvertBW in router/tokenbucket/tokenbucket.go. -BW_CODEPOINTS = 1 << 10 -MIN_BW_KBPS = 10 -MAX_BW_KBPS = 10_000_000 -BW_LOG_ENCODING_START = 60 - -# Derivation of the Hummingbird secret value of an AS from its master key, see -# DeriveSecretValue in pkg/slayers/path/hummingbird/mac.go. -SECRET_VALUE_SALT = b"Derive hbird sv" -SECRET_VALUE_ITERATIONS = 1000 -SECRET_VALUE_LENGTH = 16 - -def loadJson(path): - with open(path, "r", encoding="utf-8") as file: - data = json.load(file) - return data - -def loadTopology(genDir): - """Reads the IA, interface IDs and directory of every AS of a topology.""" - topologies = sorted(Path(genDir).glob("AS*/topology.json")) - if not topologies: - raise FileNotFoundError(f"no AS*/topology.json found in {genDir}") - ases = [] - for path in topologies: - topo = loadJson(path) - ia = topo.get("isd_as") - if ia is None: - raise ValueError(f"{path} has no isd_as") - ifids = sorted( - int(ifid) - for router in topo.get("border_routers", {}).values() - for ifid in router.get("interfaces", {}) - ) - ases.append((ia, ifids, path.parent)) - return ases - -def encodingPoints(): - """The bandwidth points of an AS, in kbps, one per codepoint. - - The first BW_LOG_ENCODING_START points are one kbps apart, starting at - MIN_BW_KBPS, and the rest grow geometrically up to MAX_BW_KBPS. The router - decodes a flyover with the very same points, so publishing anything else - would sell a bandwidth that the router does not enforce. - """ - step = (MAX_BW_KBPS / (MIN_BW_KBPS + BW_LOG_ENCODING_START)) ** ( - 1.0 / (BW_CODEPOINTS - BW_LOG_ENCODING_START - 1) - ) - points = [] - for codepoint in range(BW_CODEPOINTS): - if codepoint < BW_LOG_ENCODING_START: - points.append(MIN_BW_KBPS + codepoint) - else: - kbps = (MIN_BW_KBPS + BW_LOG_ENCODING_START) * step ** ( - codepoint - BW_LOG_ENCODING_START - ) - points.append(math.ceil(kbps)) - return points - -def secretValue(asDir): - """Derives the Hummingbird secret value of an AS from its master key. - - Handing it to the marketplace is what lets the marketplace redeem the assets - of that AS, i.e. derive the authenticator of a flyover on its behalf. - """ - path = asDir / "keys" / "master0.key" - try: - master = base64.b64decode(path.read_text(encoding="utf-8").strip(), validate=True) - except OSError as e: - raise FileNotFoundError(f"cannot read the master key {path}: {e}") - except ValueError as e: - raise ValueError(f"{path} is not a base64 encoded key: {e}") - return hashlib.pbkdf2_hmac( - "sha256", master, SECRET_VALUE_SALT, SECRET_VALUE_ITERATIONS, SECRET_VALUE_LENGTH - ) - -def interfacePairs(ifids): - """The (ingress, egress) pairs of an AS. - - Interface 0 stands for no interface, i.e. a flyover that starts or ends in - this AS, so that ASes with a single interface get assets as well. - """ - pairs = [(i, e) for i in ifids for e in ifids if i != e] - pairs += [(0, e) for e in ifids] - pairs += [(i, 0) for i in ifids] - return pairs - -def defaultEntries(genDir, now=None): - """Builds the default users, ASes and assets for the topology in genDir.""" - if now is None: - now = datetime.now(timezone.utc) - startsAt = now.replace(microsecond=0).strftime("%Y-%m-%dT%H:%M:%SZ") - stopsAt = (now + DEFAULT_ASSET_DURATION).replace(microsecond=0).strftime("%Y-%m-%dT%H:%M:%SZ") - - ases = loadTopology(genDir) - assets = [ - { - "ia": ia, - "bandwidth": DEFAULT_ASSET_BANDWIDTH, - "bandwidth_min": DEFAULT_ASSET_BANDWIDTH_MIN, - "bandwidth_max": DEFAULT_ASSET_BANDWIDTH_MAX, - "price": DEFAULT_ASSET_PRICE, - "time_granularity": DEFAULT_ASSET_TIME_GRANULARITY, - "time_min_duration": DEFAULT_ASSET_TIME_MIN_DURATION, - "starts_at": startsAt, - "stops_at": stopsAt, - "ingress": ingress, - "egress": egress, - } - for ia, ifids, _ in ases - for ingress, egress in interfacePairs(ifids) - ] - return { - "users": [ - {"name": name, "password": DEFAULT_PASSWORD} - for name in DEFAULT_USERS - ], - "accounts": [ - {"user": name, "balance": DEFAULT_USER_BALANCE} - for name in DEFAULT_USERS - ], - "ases": [ - {"ia": ia, "password": DEFAULT_PASSWORD, "balance": DEFAULT_AS_BALANCE} - for ia, _, _ in ases - ], - "assets": assets, - # Every AS delegates the redemption of its assets to the marketplace. - # Without this the marketplace cannot derive the flyover authenticators, - # and the assets it sells are worthless. - "delegations": [ - { - "ia": ia, - "res_id_limit": DEFAULT_DELEGATION_RES_ID_LIMIT, - "expiration": stopsAt, - "paid_until": stopsAt, - "key": secretValue(asDir).hex(), - "encodings": encodingPoints(), - } - for ia, _, asDir in ases - ], - } - -def dropTables(db): - """Drops every table, so that the schema and the entries are recreated. - - Indexes are dropped along with their table, and so are the AUTOINCREMENT - counters that SQLite keeps in its internal sqlite_sequence table. - """ - tables = [ - name for (name,) in db.execute( - "SELECT name FROM sqlite_master WHERE type = 'table' AND name NOT LIKE 'sqlite_%'" - ) - ] - for table in tables: - db.execute(f'DROP TABLE "{table}"') - return tables - -def insertAll(db, data) -> bool: - if not insertUsers(db, data.get("users")): - return False - if not insertAccounts(db, data.get("accounts")): - return False - if not insertASes(db, data.get("ases")): - return False - if not insertAssets(db, data.get("assets")): - return False - if not insertReservations(db, data.get("reservations")): - return False - return insertDelegations(db, data.get("delegations")) - -def hash_file(filename): - with open(filename, "r", encoding="utf-8") as f: - sql = f.read() - normalized = re.sub(r"\s+", " ", sql).strip() - return hashlib.sha256(normalized.encode("utf-8")).hexdigest() - -def verifyScheme(schemePath, version): - actualVersion = hash_file(schemePath) - if version != actualVersion: - print("scheme version mismatch. Expected: ", actualVersion) - return False - return True - -def hash_password(password: str) -> str: - if password is None: - return None - hashed = bcrypt.hashpw( - password.encode("utf-8"), - bcrypt.gensalt(rounds=12) - ) - return hashed.decode("utf-8") - -def insertAccounts(db, accounts) -> bool: - if accounts is None: - return True - # Updated instead of replaced: a replace would delete the conflicting row and - # insert it under a new id, orphaning the assets and reservations that refer - # to it. - db.executemany( - """ - INSERT INTO Accounts (scope, balance, user_id) - VALUES (?, ?, ( - SELECT id - FROM Users - WHERE name = ? - LIMIT 1 - )) - ON CONFLICT(user_id, scope) DO UPDATE SET balance = excluded.balance - """, - [ - (u.get("scope") or "", u.get("balance"), u.get("user")) - for u in accounts - ], - ) - return True - -def insertUsers(db,users) -> bool: - if users is None: - return True - for u in users: - u["pw_hash"] = hash_password(u.get("password")) - db.executemany( - """ - INSERT INTO Users (name, pw_hash) - VALUES (?, ?) - ON CONFLICT(name) DO UPDATE SET pw_hash = excluded.pw_hash - """, - [ - (u.get("name"), u.get("pw_hash")) - for u in users - ], - ) - return True - -def parseIA(ia_string): - parts = ia_string.split("-") - isd_id = int(parts[0]) - as_string = parts[1] - as_parts = as_string.split(":") - if len(as_parts) == 1: - return isd_id, int(as_string) - if len(as_parts) != 3: - print("invalid IA", ia_string) - return 0,0 - as_id = 0 - for i in range(3): - as_id <<= 16 - as_id |= int(as_parts[i], 16) - return isd_id, as_id - -def insertAssets(db, assets) -> bool: - if assets is None: - return True - for a in assets: - isd_id, as_id = parseIA(a.get("ia")) - if isd_id == 0 and as_id == 0: - return False - a["isd_id"] = isd_id - a["as_id"] = as_id - db.executemany( - """ - INSERT INTO Assets (isd_id, as_id, bandwidth, bandwidth_min, bandwidth_max, price, time_granularity, time_min_duration, starts_at, stops_at, ingress, egress, account_id) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ( - SELECT a.id - FROM Accounts a - JOIN Users u ON u.ID = a.user_id - WHERE u.name = ? AND a.scope = '' - LIMIT 1 - )) - """, - [ - (u.get("isd_id"), u.get("as_id"), u.get("bandwidth"), u.get("bandwidth_min"), u.get("bandwidth_max"), u.get("price"), u.get("time_granularity"), u.get("time_min_duration"), u.get("starts_at"), u.get("stops_at"), u.get("ingress"), u.get("egress"), u.get("owner"),) - for u in assets - ], - ) - return True +import argparse +import sys +from pathlib import Path -def insertASes(db, ases) -> bool: - if ases is None: - return True - for a in ases: - isd_id, as_id = parseIA(a.get("ia")) - if isd_id == 0 and as_id == 0: - return False - a["isd_id"] = isd_id - a["as_id"] = as_id - a["pw_hash"] = hash_password(a.get("password")) - if a.get("pw_hash") == None: - a["pw_hash"] = "" - if a.get("jwt_version") == None: - a["jwt_version"] = 0 - if a.get("balance") == None: - a["balance"] = 0 - db.executemany( - """ - INSERT OR REPLACE INTO Ases (isd_id, as_id, pw_hash, jwt_version, balance) - VALUES (?, ?, ?, ?, ?) - """, - [ - (u.get("isd_id"), u.get("as_id"), u.get("pw_hash"), u.get("jwt_version"), u.get("balance")) - for u in ases - ], - ) - return True +# We need to import from the topology generator, which is not in sys.path. +# This is unavoidable unless we move either the topology generation, +# or the marketplace tools. +sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "tools")) -def insertReservations(db, reservations): - if reservations is None: - return True - for a in reservations: - isd_id, as_id = parseIA(a.get("ia")) - if isd_id == 0 and as_id == 0: - return False - a["isd_id"] = isd_id - a["as_id"] = as_id - a["key_bytes"] = bytes.fromhex(a.get("key")) - for u in reservations: - db.execute( - """ - INSERT OR REPLACE INTO Reservations (reservation_id, isd_id, as_id, ingress, egress, bandwidth, bw_encoded, starts_at, stops_at, key, account_id) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ( - SELECT id - FROM Users - WHERE name = ? - LIMIT 1 - )) - """, (u.get("id") , u.get("isd_id"), u.get("as_id"), u.get("ingress"), u.get("egress"), u.get("bandwidth"), u.get("bw_encoded"), u.get("starts_at"), u.get("stops_at"), u.get("key_bytes"),u.get("owner")) - ) - return True +# Linters like flake8 would complain about the import not being top-level. Avoid with noqa: e402 +from topology.marketplace import ( # noqa: E402 + MARKETPLACE_DB_NAME, + MARKETPLACE_SCHEMA, + LOCAL_CACHE_DIR, + MarketplaceError, + defaultEntries, + loadJson, + populateDB, + verifySchema, +) -def insertDelegations(db, delegations): - if delegations is None: - return True - for a in delegations: - isd_id, as_id = parseIA(a.get("ia")) - if isd_id == 0 and as_id == 0: - return False - a["isd_id"] = isd_id - a["as_id"] = as_id - a["key_bytes"] = bytes.fromhex(a.get("key")) - a["encoding_bytes"] = b"".join( - i.to_bytes(4, byteorder="little") - for i in a.get("encodings") - ) - db.executemany( - """ - INSERT OR REPLACE INTO Redemption_Delegations (isd_id, as_id, res_id_limit, expiration, paid_until, key, encodings) - VALUES (?, ?, ?, ?, ?, ?, ?) - """, - [ - (u.get("isd_id"), u.get("as_id"), u.get("res_id_limit"), u.get("expiration"), u.get("paid_until"), u.get("key_bytes"),u.get("encoding_bytes")) - for u in delegations - ], - ) - return True +DEFAULT_DB = str(Path(LOCAL_CACHE_DIR) / MARKETPLACE_DB_NAME) -def applyScheme(db, schemePath): - with open(schemePath, "r", encoding="utf-8") as f: - db.executescript(f.read()) -def main(args): - data = {} - if args.file is not None: - data = loadJson(args.file) - if not verifyScheme(args.schema, data.get("version")): - return +def main(args) -> None: if args.default_entries: + data = defaultEntries(args.gen_dir) + else: try: - data = defaultEntries(args.gen_dir) + data = loadJson(args.file) except (OSError, ValueError) as e: - print("cannot read the topology:", e) - return - try: - conn = sqlite3.connect(args.db) - except sqlite3.Error as e: - print("cannot open the database:", e) - return - cursor = conn.cursor() - try: - dropped = dropTables(cursor) - if dropped: - print("dropped", ", ".join(dropped)) - applyScheme(cursor, args.schema) - if insertAll(cursor, data): - conn.commit() - print("stored in database") - # The database is the one bind mounted into the containers, but a - # running marketplace keeps its handles on the tables just dropped. - if (Path(args.gen_dir) / "scion-dc.yml").is_file(): - print("this topology runs on docker: restart the marketplace service " - "if it is already running") - else: - conn.rollback() - print("changes rolled back") - except Exception as e: - print("error", e) - conn.rollback() - finally: - conn.close() + raise MarketplaceError("cannot read %s: %s" % (args.file, e)) + verifySchema(args.schema, data.get("version")) + + dropped = populateDB(args.db, args.schema, data) + if dropped: + print("dropped", ", ".join(dropped)) + print("stored in database") + # The database is the one bind-mounted into the containers, + # but a running marketplace keeps its handles on the tables just dropped. + if (Path(args.gen_dir) / "scion-dc.yml").is_file(): + print("this topology runs on docker: restart the marketplace service if it is running") + if __name__ == "__main__": parser = argparse.ArgumentParser( @@ -591,13 +217,14 @@ def main(args): ) parser.add_argument( "--db", - default="gen-cache/marketplace.db", - help="Path to the SQLite database. Its tables are dropped and recreated on every run." + default=DEFAULT_DB, + help="Path to the SQLite database. Its tables are dropped and recreated on every run. " + "(default: %s)" % DEFAULT_DB ) parser.add_argument( "--schema", - default="marketplace/db/schema.sql", - help="Path to the schema.sql file." + default=MARKETPLACE_SCHEMA, + help="Path to the schema.sql file (default: %s)." % MARKETPLACE_SCHEMA ) source = parser.add_mutually_exclusive_group(required=True) source.add_argument( @@ -608,7 +235,8 @@ def main(args): "--default-entries", action="store_true", help="Fill the database with the default users, ASes, assets and " - "redemption delegations for the topology in --gen-dir." + "redemption delegations for the topology in --gen-dir. This is what " + "topogen -m does." ) parser.add_argument( "--gen-dir", @@ -616,5 +244,12 @@ def main(args): help="Path to the generated topology, read by --default-entries (default: gen)." ) - args = parser.parse_args() - main(args) \ No newline at end of file + try: + main(parser.parse_args()) + except MarketplaceError as e: + print("error: %s" % e, file=sys.stderr) + sys.exit(1) + except Exception as e: + # Most likely a malformed entry in --file. Don't traceback. + print("error: %s: %s" % (type(e).__name__, e), file=sys.stderr) + sys.exit(1) diff --git a/marketplace/tools/setup_marketplace.py b/marketplace/tools/setup_marketplace.py index 1fffbecc71..5a19e65e4d 100755 --- a/marketplace/tools/setup_marketplace.py +++ b/marketplace/tools/setup_marketplace.py @@ -2,88 +2,65 @@ """Adds a marketplace instance to an already generated SCION topology. -Both backends are supported, and the one to use is detected the same way -scion.sh does it: gen/scion-dc.yml means the topology runs on docker, and the -marketplace becomes a compose service; otherwise it is a supervisord program. +`topogen -m ISD-AS` (i.e. `./scion.sh topology -m ISD-AS`) does this as part of +generating the topology, and prepopulates the database on top. This script is for +the topologies that are already there: it adds the marketplace to one of their +ASes without regenerating anything. The database is left alone, so +populate_marketplace.py is still needed to fill it. + +Both backends are supported, and the one to use is detected the same way scion.sh +does it: gen/scion-dc.yml means the topology runs on docker, and the marketplace +becomes a compose service; otherwise it is a supervisord program. The script is idempotent: running it again for the same IA reports what is already in place instead of adding a second marketplace. """ import argparse -import json +import configparser +import os import re import sys -from collections import namedtuple +from io import StringIO from pathlib import Path import yaml -# -, with the AS either as a hex triple (underscores or colons, the -# form used by gen/AS) or as a plain decimal (BGP-style) AS number. -IA_RE = re.compile( - r"^(?P\d{1,5})-" - r"(?P[0-9a-fA-F]{1,4}[_:][0-9a-fA-F]{1,4}[_:][0-9a-fA-F]{1,4}|\d{1,10})$" +# The marketplace is described once by the topology generator, +# so that a marketplace added here is the same one that a generated topology would have. +sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "tools")) + +# Linters like flake8 would complain about the import not being top-level. Avoid with noqa: e402 +from topology.common import TopoID # noqa: E402 +from topology.marketplace import ( # noqa: E402 + CONTAINER_CONFIG_DIR, + CONTAINER_DB, + CONTAINER_TOML, + LOCAL_CACHE_DIR, + LOCAL_HOST, + MARKETPLACE_CONFIG_NAME, + MARKETPLACE_DB_NAME, + MARKETPLACE_PORT, + PROGRAM_NAME, + STATIC_INFO_CONFIG_NAME, + Endpoints, + MarketplaceError, + advertiseMarketplace, + checkDispatchedPorts, + dockerService, + dockerServiceName, + hostPort, + marketplaceToml, + supervisordProgram, ) -MAX_ISD = (1 << 16) - 1 -MAX_BGP_AS = (1 << 32) - 1 - -# The port of both APIs, the TCP one serving the ConnectRPC API and the web app -# from a single mux, and the SCION/QUIC one, which is UDP. They do not collide. -# It has to be inside the dispatched_ports range of the AS, otherwise only the -# shim dispatcher can deliver SCION packets to it. -MARKETPLACE_PORT = 31888 - -# Where the marketplace of a docker topology finds its files inside the container. -CONTAINER_CONFIG_DIR = "/etc/scion" -CONTAINER_DB = "/share/cache/marketplace.db" -CONTAINER_TOML = f"{CONTAINER_CONFIG_DIR}/marketplace.toml" - DOCKER_CONF = "scion-dc.yml" +SUPERVISOR_CONF = "supervisord.conf" DEFAULT_IMAGE = "scion/marketplace:latest" -# What the marketplace binds, and what staticInfoConfig.json advertises for it. -Endpoints = namedtuple("Endpoints", "host api_port scion_port") - -MARKETPLACE_TOML = """[general] -id = "marketplace" -config_dir = "{config_dir}" - -[log.console] -level = "debug" - -[marketplace] -api_addr = "{api_addr}" -scion_api_addr = "{scion_api_addr}" -currency = "CHF" -currency_exponent = 2 -supports_redemption_delegation = true -delegation_hourly_fee = 10000 -transaction_fee_relative = 0.01 -transaction_fee_absolute = 10 -split_combine_fee_absolute = 10 - -[marketplace_db] -connection = "{db_path}" -""" - -MARKETPLACE_PROGRAM = """[program:marketplace] -autostart = false -autorestart = false -environment = TZ=UTC,GODEBUG="cgocheck=0" -stdout_logfile = logs/marketplace.log -redirect_stderr = True -startretries = 0 -startsecs = 5 -priority = 100 -command = bin/marketplace --config {config} - -""" - # The [program:marketplace] section, i.e. up to the next section or EOF. PROGRAM_SECTION_RE = re.compile( - r"^\[program:marketplace\]\s*$.*?(?=^\[|\Z)", + r"^\[program:%s\]\s*$.*?(?=^\[|\Z)" % PROGRAM_NAME, re.MULTILINE | re.DOTALL, ) PROGRAM_CONFIG_RE = re.compile(r"^command\s*=.*?--config\s+(?P\S+)", re.MULTILINE) @@ -92,136 +69,24 @@ SERVICE_RE = re.compile(r"marketplace(?P\d+)-(?P[0-9a-fA-F_]+)") -class SetupError(Exception): - """An error that aborts the setup with a message, but without a traceback.""" - - -def parseIA(ia_string: str): - """Splits an IA into (isd, as_id), with the AS id in the gen/ underscore form.""" - match = IA_RE.match(ia_string.strip()) - if match is None: - raise SetupError( - f"invalid ISD-AS {ia_string!r}; expected e.g. 1-ff00_0_111, " - "1-ff00:0:111 or 1-64512" - ) - isd = int(match.group("isd")) - if isd < 1 or isd > MAX_ISD: - raise SetupError(f"invalid ISD in {ia_string!r}: must be in [1, {MAX_ISD}]") - as_id = match.group("as") - if as_id.isdigit() and int(as_id) > MAX_BGP_AS: - raise SetupError(f"invalid AS number in {ia_string!r}: must be <= {MAX_BGP_AS}") - return isd, as_id.replace(":", "_") - - -def hostPort(host: str, port: int) -> str: - """host:port, with an IPv6 host in brackets.""" - if ":" in host: - return f"[{host}]:{port}" - return f"{host}:{port}" - - -def marketplaceEntries(ia_colons: str, endpoints: Endpoints): - """The staticInfoConfig entries advertising this marketplace. - - They describe what the marketplace really binds: the web app is served by the - same mux as the TCP API, so both live on the API port. - """ - website = f"https://{hostPort(endpoints.host, endpoints.api_port)}" - return [ - { - "name": "Test Market", - "protocol": "connectrpc/TLS/QUIC/SCION", - "api": f"[{ia_colons},{endpoints.host}]:{endpoints.scion_port}", - "website": website, - }, - { - "name": "Test Market", - "protocol": "connectrpc/TLS/TCP", - "api": website, - "website": website, - }, - ] - - -def marketplaceToml(config_dir: str, endpoints: Endpoints, db_path: str) -> str: - return MARKETPLACE_TOML.format( - config_dir=config_dir, - api_addr=hostPort(endpoints.host, endpoints.api_port), - scion_api_addr=hostPort(endpoints.host, endpoints.scion_port), - db_path=db_path, - ) - - -def checkDispatchedPorts(marketplace_dir: Path, port: int) -> None: - """Checks that the SCION API port is one the border router delivers directly.""" - topology = marketplace_dir / "topology.json" - if not topology.is_file(): - return +def parseIA(ia_string: str) -> TopoID: try: - ports = json.loads(topology.read_text(encoding="utf-8")).get("dispatched_ports") - except (OSError, json.JSONDecodeError) as e: - raise SetupError(f"cannot read {topology}: {e}") - if not ports or ports == "all": - return - match = re.fullmatch(r"(\d+)-(\d+)", str(ports).strip()) - if match is None: - return - start, end = int(match.group(1)), int(match.group(2)) - if not start <= port <= end: - raise SetupError( - f"the SCION API port {port} is outside the dispatched_ports range " - f"{ports} of {topology}; the marketplace would only receive SCION " - "packets through a shim dispatcher" - ) - - -def loadStaticInfo(path: Path): - """Reads staticInfoConfig.json and its embedded note object. - - Returns (static_info, note), both dicts. Missing files yield empty ones. - """ - if not path.exists(): - return {}, {} - - try: - static_info = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError) as e: - raise SetupError(f"cannot read {path}: {e}") - if not isinstance(static_info, dict): - raise SetupError(f"{path}: expected a JSON object, got {type(static_info).__name__}") - - note = static_info.get("note", "") - if note == "": - return static_info, {} - try: - note = json.loads(note) - except (TypeError, json.JSONDecodeError): - raise SetupError( - f"{path}: the 'note' field is not a JSON object; refusing to overwrite it" - ) - if not isinstance(note, dict): - raise SetupError( - f"{path}: the 'note' field is not a JSON object; refusing to overwrite it" - ) - return static_info, note - - -def mergeEntries(existing, entries): - """Merges the marketplace entries into the existing ones, keyed by name+protocol.""" - merged = [e for e in existing if isinstance(e, dict)] - changed = len(merged) != len(existing) - for entry in entries: - key = (entry["name"], entry["protocol"]) - for i, old in enumerate(merged): - if (old.get("name"), old.get("protocol")) == key: - if old != entry: - merged[i] = entry - changed = True - break - else: - merged.append(entry) - changed = True - return merged, changed + return TopoID(ia_string.strip()) + except ValueError: + raise MarketplaceError( + "invalid ISD-AS %r; expected e.g. 1-ff00_0_111, 1-ff00:0:111 or 1-64512" + % ia_string) + + +def programSection(config_path: str) -> str: + """The [program:marketplace] section of a supervisord config, as text.""" + config = configparser.ConfigParser(interpolation=None) + config["program:%s" % PROGRAM_NAME] = { + key: str(value) for key, value in supervisordProgram(config_path).items() + } + text = StringIO() + config.write(text) + return text.getvalue() def addToGroup(text: str, group_name: str): @@ -235,46 +100,46 @@ def addToGroup(text: str, group_name: str): if match is None: groups = re.findall(r"^\[group:(.+)\]$", text, re.MULTILINE) known = ", ".join(groups) if groups else "none" - raise SetupError( + raise MarketplaceError( f"no [group:{group_name}] section in the supervisord config " f"(known groups: {known})" ) programs = [p.strip() for p in match.group(2).split(",") if p.strip()] - if "marketplace" in programs: + if PROGRAM_NAME in programs: return text, False - programs.append("marketplace") + programs.append(PROGRAM_NAME) return text[:match.start()] + match.group(1) + ",".join(programs) + text[match.end():], True -def planSupervisord(gen_dir: Path, isd: int, as_id: str, as_dir: str, - api_port: int, scion_port: int): +def planSupervisord(gen_dir: Path, topo_id: TopoID, api_port: int, scion_port: int): """Plans the marketplace of a supervisord topology. Returns the endpoints it will bind, the contents of its marketplace.toml, and a function applying the changes to gen/supervisord.conf. """ - supervisord_conf = gen_dir / "supervisord.conf" - group_name = f"as{isd}-{as_id}" # e.g. as1-ff00_0_111 - config_path = f"{gen_dir / as_dir}/marketplace.toml" + supervisord_conf = gen_dir / SUPERVISOR_CONF + group_name = "as%s" % topo_id.file_fmt() # e.g. as1-ff00_0_111 + as_base = str(gen_dir / topo_id.AS_file()) + config_path = os.path.join(as_base, MARKETPLACE_CONFIG_NAME) try: config = supervisord_conf.read_text(encoding="utf-8") except OSError as e: - raise SetupError(f"cannot read {supervisord_conf}: {e}") + raise MarketplaceError(f"cannot read {supervisord_conf}: {e}") # There is a single [program:marketplace]: if it exists, it must be ours. section = PROGRAM_SECTION_RE.search(config) if section is not None: command = PROGRAM_CONFIG_RE.search(section.group(0)) if command is None: - raise SetupError( - f"{supervisord_conf} already has a [program:marketplace] section with an " + raise MarketplaceError( + f"{supervisord_conf} already has a [program:{PROGRAM_NAME}] section with an " "unrecognized command; remove it before running this script" ) if Path(command.group("config")) != Path(config_path): - raise SetupError( + raise MarketplaceError( f"{supervisord_conf} already runs a marketplace with " f"{command.group('config')}; only one marketplace is supported" ) @@ -283,21 +148,22 @@ def planSupervisord(gen_dir: Path, isd: int, as_id: str, as_dir: str, # applied, so that a missing group is not a partial setup. config, group_added = addToGroup(config, group_name) - endpoints = Endpoints("127.0.0.1", api_port, scion_port) - toml_content = marketplaceToml(str(gen_dir / as_dir), endpoints, "gen-cache/marketplace.db") + endpoints = Endpoints(LOCAL_HOST, api_port, scion_port) + db_path = os.path.join(LOCAL_CACHE_DIR, MARKETPLACE_DB_NAME) + toml_content = marketplaceToml(as_base, endpoints, db_path) def apply(changes): text = config if section is None: - text = MARKETPLACE_PROGRAM.format(config=config_path) + text - changes.append(f"added [program:marketplace] to {supervisord_conf}") + text = programSection(config_path) + text + changes.append(f"added [program:{PROGRAM_NAME}] to {supervisord_conf}") if group_added: - changes.append(f"added marketplace to [group:{group_name}]") + changes.append(f"added {PROGRAM_NAME} to [group:{group_name}]") if section is None or group_added: try: supervisord_conf.write_text(text, encoding="utf-8") except OSError as e: - raise SetupError(f"cannot write {supervisord_conf}: {e}") + raise MarketplaceError(f"cannot write {supervisord_conf}: {e}") return endpoints, toml_content, apply @@ -307,13 +173,13 @@ def loadCompose(path: Path): try: compose = yaml.safe_load(path.read_text(encoding="utf-8")) except (OSError, yaml.YAMLError) as e: - raise SetupError(f"cannot read {path}: {e}") + raise MarketplaceError(f"cannot read {path}: {e}") if not isinstance(compose, dict) or not isinstance(compose.get("services"), dict): - raise SetupError(f"{path} has no services, it is not a compose file") + raise MarketplaceError(f"{path} has no services, it is not a compose file") return compose -def controlService(compose, path: Path, isd: int, as_id: str): +def controlService(compose, path: Path, topo_id: TopoID): """The control service of an AS, and the address its containers share. The marketplace joins the network namespace of the control service, so it ends @@ -321,26 +187,24 @@ def controlService(compose, path: Path, isd: int, as_id: str): pinned on the dispatcher, which is the container that owns the namespace. """ services = compose["services"] - pattern = re.compile(rf"cs{isd}-{re.escape(as_id)}-\d+") + pattern = re.compile(r"cs%s-\d+" % re.escape(topo_id.file_fmt())) names = sorted(n for n in services if pattern.fullmatch(n)) if not names: - raise SetupError( - f"{path} has no control service of {isd}-{as_id.replace('_', ':')}; " - "is that AS part of it?" - ) + raise MarketplaceError( + f"{path} has no control service of {topo_id}; is that AS part of it?") control = names[0] dispatcher = f"disp_{control}" if dispatcher not in services: - raise SetupError(f"{path} has no {dispatcher}; regenerate the topology") + raise MarketplaceError(f"{path} has no {dispatcher}; regenerate the topology") networks = services[dispatcher].get("networks") if not isinstance(networks, dict) or len(networks) != 1: - raise SetupError(f"{path}: expected exactly one network on {dispatcher}") + raise MarketplaceError(f"{path}: expected exactly one network on {dispatcher}") (_, addresses), = networks.items() for key in ("ipv4_address", "ipv6_address"): if key in addresses: return control, dispatcher, addresses[key] - raise SetupError(f"{path}: {dispatcher} has no address of its own") + raise MarketplaceError(f"{path}: {dispatcher} has no address of its own") def marketplaceImage(control): @@ -353,41 +217,26 @@ def marketplaceImage(control): return f"{registry}/marketplace{colon}{tag}" -def composeService(path: Path, as_dir: str, control_name: str, control, dispatcher: str): - """The compose service running the marketplace of an AS. +def marketplaceVolumes(path: Path, as_dir: str, control_name: str, control): + """The volumes of the marketplace, i.e. the ones of its control service. - It shares the network namespace of the control service, the way the control - service and the hummingbird service already share the one of their dispatcher. - It therefore carries no address of its own, and needs none. + Unlike the other services, the marketplace writes into its config dir: + its TLS certificate and its JWT signature keys are created on first run. """ volumes = [] for volume in control.get("volumes") or []: if not isinstance(volume, str): - raise SetupError(f"{path}: expected string volumes on {control_name}") - # Unlike the other services, the marketplace writes into its config dir: - # its TLS certificate and its JWT signature keys are created on first run. + raise MarketplaceError(f"{path}: expected string volumes on {control_name}") volumes.append(volume[:-2] + "rw" if volume.endswith(":ro") else volume) if not any(v.endswith(f"/{as_dir}:{CONTAINER_CONFIG_DIR}:rw") for v in volumes): - raise SetupError( + raise MarketplaceError( f"{path}: {control_name} does not mount {as_dir} at {CONTAINER_CONFIG_DIR}; " "regenerate the topology" ) - - entry = { - "command": ["--config", CONTAINER_TOML], - # The marketplace asks the control service for trust material on startup. - "depends_on": {control_name: {"condition": "service_healthy"}}, - "image": marketplaceImage(control), - "network_mode": f"service:{dispatcher}", - "volumes": volumes, - } - if "user" in control: - entry["user"] = control["user"] - return entry + return volumes -def planDocker(gen_dir: Path, isd: int, as_id: str, as_dir: str, - api_port: int, scion_port: int): +def planDocker(gen_dir: Path, topo_id: TopoID, api_port: int, scion_port: int): """Plans the marketplace of a docker topology. Returns the endpoints it will bind, the contents of its marketplace.toml, and @@ -396,7 +245,7 @@ def planDocker(gen_dir: Path, isd: int, as_id: str, as_dir: str, dc_file = gen_dir / DOCKER_CONF compose = loadCompose(dc_file) services = compose["services"] - name = f"marketplace{isd}-{as_id}" + name = dockerServiceName(topo_id) # There is a single marketplace: if one exists, it must be ours. for other, service in services.items(): @@ -404,21 +253,28 @@ def planDocker(gen_dir: Path, isd: int, as_id: str, as_dir: str, continue match = SERVICE_RE.fullmatch(other) if match is not None: - raise SetupError( + raise MarketplaceError( f"{dc_file} already runs a marketplace for " f"{match.group('isd')}-{match.group('as').replace('_', ':')}; " "only one marketplace is supported" ) if CONTAINER_TOML in (service.get("command") or []): - raise SetupError( + raise MarketplaceError( f"{dc_file} already has a marketplace service named {other}; " "remove it before running this script" ) - control_name, dispatcher, address = controlService(compose, dc_file, isd, as_id) + control_name, dispatcher, address = controlService(compose, dc_file, topo_id) + control = services[control_name] endpoints = Endpoints(address, api_port, scion_port) toml_content = marketplaceToml(CONTAINER_CONFIG_DIR, endpoints, CONTAINER_DB) - entry = composeService(dc_file, as_dir, control_name, services[control_name], dispatcher) + entry = dockerService( + marketplaceImage(control), + control_name, + dispatcher, + marketplaceVolumes(dc_file, topo_id.AS_file(), control_name, control), + control.get("user"), + ) def apply(changes): if services.get(name) == entry: @@ -428,53 +284,45 @@ def apply(changes): try: dc_file.write_text(yaml.dump(compose, default_flow_style=False), encoding="utf-8") except OSError as e: - raise SetupError(f"cannot write {dc_file}: {e}") + raise MarketplaceError(f"cannot write {dc_file}: {e}") changes.append(f"{verb} the {name} service of {dc_file}") return endpoints, toml_content, apply def main(args) -> None: - isd, as_id = parseIA(args.ia) - ia_colons = f"{isd}-{as_id.replace('_', ':')}" - as_dir = f"AS{as_id}" # e.g. ASff00_0_111 + topo_id = parseIA(args.ia) # Paths gen_dir = Path(args.gen_dir) - marketplace_dir = gen_dir / as_dir + marketplace_dir = gen_dir / topo_id.AS_file() certs_dir = marketplace_dir / "certs" - marketplace_toml = marketplace_dir / "marketplace.toml" - static_info_json = marketplace_dir / "staticInfoConfig.json" + marketplace_toml = marketplace_dir / MARKETPLACE_CONFIG_NAME + static_info_json = marketplace_dir / STATIC_INFO_CONFIG_NAME # 0. Validate the topology before touching anything, so that a typo in the # IA does not leave a half configured AS behind. if not gen_dir.is_dir(): - raise SetupError(f"{gen_dir} does not exist; generate the topology first") + raise MarketplaceError(f"{gen_dir} does not exist; generate the topology first") if not marketplace_dir.is_dir(): - raise SetupError(f"{marketplace_dir} does not exist; {ia_colons} is not part of {gen_dir}") + raise MarketplaceError( + f"{marketplace_dir} does not exist; {topo_id} is not part of {gen_dir}") # The backend is the one the topology was generated for, detected the way # scion.sh does it. if (gen_dir / DOCKER_CONF).is_file(): plan = planDocker - elif (gen_dir / "supervisord.conf").is_file(): + elif (gen_dir / SUPERVISOR_CONF).is_file(): plan = planSupervisord else: - raise SetupError( - f"neither {gen_dir / DOCKER_CONF} nor {gen_dir / 'supervisord.conf'} exists; " + raise MarketplaceError( + f"neither {gen_dir / DOCKER_CONF} nor {gen_dir / SUPERVISOR_CONF} exists; " "generate the topology first" ) checkDispatchedPorts(marketplace_dir, args.scion_api_port) endpoints, toml_content, apply = plan( - gen_dir, isd, as_id, as_dir, args.api_port, args.scion_api_port) - - static_info, note = loadStaticInfo(static_info_json) - hummingbird = note.get("hummingbird", []) - if not isinstance(hummingbird, list): - raise SetupError( - f"{static_info_json}: 'note'.hummingbird is not a list; refusing to overwrite it" - ) + gen_dir, topo_id, args.api_port, args.scion_api_port) changes = [] @@ -500,25 +348,20 @@ def main(args) -> None: # 6. Advertise the marketplace in staticInfoConfig.json, without dropping # whatever else is already configured there. - hummingbird, changed = mergeEntries(hummingbird, marketplaceEntries(ia_colons, endpoints)) - if changed: - note["hummingbird"] = hummingbird - static_info["note"] = json.dumps(note) - try: - static_info_json.write_text(json.dumps(static_info, indent=2), encoding="utf-8") - except OSError as e: - raise SetupError(f"cannot write {static_info_json}: {e}") + advertised = advertiseMarketplace(static_info_json, topo_id, endpoints) + if advertised: changes.append(f"advertised the marketplace in {static_info_json}") if not changes: - print(f"Marketplace already configured for {ia_colons}, nothing to do") + print(f"Marketplace already configured for {topo_id}, nothing to do") return for change in changes: print(change) - print(f"Marketplace config added for {ia_colons}") - if changed: + print(f"Marketplace config added for {topo_id}") + if advertised: print("The control service reads staticInfoConfig.json on startup: restart the " "topology for other ASes to discover the marketplace") + print("Its database is not touched here: run populate_marketplace.py to fill it") if __name__ == "__main__": @@ -551,6 +394,6 @@ def main(args) -> None: try: main(parser.parse_args()) - except SetupError as e: + except MarketplaceError as e: print(f"error: {e}", file=sys.stderr) sys.exit(1) diff --git a/scion.sh b/scion.sh index bd9768954c..3875fed8c5 100755 --- a/scion.sh +++ b/scion.sh @@ -198,8 +198,10 @@ cmd_help() { topology subcommand. Usage: - $PROGRAM topology [-d] [-c TOPOFILE] + $PROGRAM topology [-d] [-c TOPOFILE] [-m ISD-AS] Create topology, configuration, and execution files. + With -m, the given AS also runs a marketplace, and its database is + pre-populated with default users, assets and redemption delegations. All arguments or options are passed to tools/topogen.py $PROGRAM run Run network. diff --git a/tools/BUILD.bazel b/tools/BUILD.bazel index e174aec941..25f53b5168 100644 --- a/tools/BUILD.bazel +++ b/tools/BUILD.bazel @@ -50,6 +50,7 @@ py_binary( name = "topogen", srcs = ["topogen.py"], data = [ + "//marketplace/db:schema.sql", "//scion-pki/cmd/scion-pki", "//tools:docker_ip", ], @@ -61,6 +62,7 @@ py_binary( deps = [ "//tools/topology:py_default_library", "@bazel_tools//tools/python/runfiles", + requirement("bcrypt"), requirement("toml"), requirement("plumbum"), requirement("pyyaml"), diff --git a/tools/end2end/BUILD.bazel b/tools/end2end/BUILD.bazel index 16f7dd9cc1..e805066bf5 100644 --- a/tools/end2end/BUILD.bazel +++ b/tools/end2end/BUILD.bazel @@ -11,7 +11,7 @@ go_library( "//pkg/addr:go_default_library", "//pkg/daemon:go_default_library", "//pkg/daemon/types:go_default_library", - "//pkg/hummingbird:go_default_library", + "//pkg/hummingbird/marketplace:go_default_library", "//pkg/hummingbird/redemption:go_default_library", "//pkg/log:go_default_library", "//pkg/private/common:go_default_library", diff --git a/tools/end2end/main.go b/tools/end2end/main.go index bb3278fc01..1afe18a584 100644 --- a/tools/end2end/main.go +++ b/tools/end2end/main.go @@ -31,6 +31,7 @@ import ( "errors" "flag" "fmt" + "math" "net" "os" "path/filepath" @@ -44,7 +45,7 @@ import ( "github.com/scionproto/scion/pkg/addr" "github.com/scionproto/scion/pkg/daemon" daemontypes "github.com/scionproto/scion/pkg/daemon/types" - hummpkg "github.com/scionproto/scion/pkg/hummingbird" + marketclient "github.com/scionproto/scion/pkg/hummingbird/marketplace" "github.com/scionproto/scion/pkg/hummingbird/redemption" "github.com/scionproto/scion/pkg/log" "github.com/scionproto/scion/pkg/private/common" @@ -89,6 +90,7 @@ var ( hummingbird string // e.g. "1,5s" or "1,5s,2" hummKeysDir string // for testing purposes only hummParams hummingbirdParameters // derived from the string in hummingbird + marketParams marketplaceParameters // derived from the environment ) func main() { @@ -126,9 +128,14 @@ func addFlags() { flag.Var(timeout, "timeout", "The timeout for each attempt") flag.BoolVar(&epic, "epic", false, "Enable EPIC") flag.StringVar(&hummingbird, "hummingbird", "", - "Enable Hummingbird with BW,dur[,reverseBW] (e.g. '3,5s' or '3,5s,2')") + "Enable Hummingbird with BW,dur[,reverseBW]. With -hummKeysDir the bandwidths are "+ + "a bandwidth class, without a unit (e.g. '3,5s' or '3,5s,2'); without it they "+ + "are bandwidths and need a unit, kbps|mbps|gbps "+ + "(e.g. '100kbps,20s' or '100kbps,20s,1mbps')") flag.StringVar(&hummKeysDir, "hummKeysDir", "", - "Root directory containing AS*/keys/master0.key files for Hummingbird") + "Root directory containing AS*/keys/master0.key files for Hummingbird. "+ + "Without it, the reservations are bought from the marketplace configured "+ + "in "+envMarketplaceURL+" and "+envMarketplaceJWT) } func validateFlags() { @@ -148,13 +155,29 @@ func validateFlags() { } if hummingbird != "" { var err error - hummParams, err = parseHummingbirdFlag(hummingbird) + // With the secret values of the ASes the bandwidths are a bandwidth + // class; bought from a marketplace they are bandwidths, with a unit. + hummParams, err = parseHummingbirdFlag(hummingbird, hummKeysDir == "") if err != nil { integration.LogFatal("bad hummingbird flag", "value", hummingbird, "err", err) } + if hummKeysDir != "" { + // Flyovers derived from the AS secret values carry a 16 bit bandwidth class. + if hummParams.Bw > math.MaxUint16 || hummParams.ReverseBw > math.MaxUint16 { + integration.LogFatal("hummingbird bandwidth must fit in 16 bits", + "bw", hummParams.Bw, "reverse_bw", hummParams.ReverseBw) + } + } else if integration.Mode == integration.ModeClient { + // Without the secret values the reservations are bought from a + // marketplace, and the bandwidths are in kbps. + marketParams, err = marketplaceParametersFromEnv() + if err != nil { + integration.LogFatal("bad marketplace configuration", "err", err) + } + } } log.Info("Flags", "timeout", timeout, "epic", epic, "hummingbird", hummingbird, - "humm_keys_dir", hummKeysDir, "remote", remote) + "humm_keys_dir", hummKeysDir, "marketplace_url", marketParams.URL, "remote", remote) } type server struct{} @@ -272,6 +295,10 @@ type client struct { hummKeysDir string hummParams hummingbirdParameters hummSVByIA map[addr.IA][]byte + // Specific to Hummingbird without secret values, i.e. buying the flyovers: + topo snet.Topology + marketParams marketplaceParameters + marketClient *marketclient.MarketplaceClient } func (c *client) run() int { @@ -301,6 +328,8 @@ func (c *client) run() int { c.hummKeysDir = hummKeysDir c.hummParams = hummParams c.hummSVByIA = make(map[addr.IA][]byte) + c.topo = topo + c.marketParams = marketParams log.Info("Send", "local", fmt.Sprintf("%v,[%v] -> %v,[%v]", integration.Local.IA, integration.Local.Host, @@ -456,7 +485,7 @@ func (c *client) configureRemotePath(ctx context.Context, path snet.Path) error if c.hummKeysDir != "" { reservation, err = c.buildReservationWithSecretValues(ctx, path, time.Now()) } else { - reservation, err = c.buildReservationWithRedemptions(ctx, path, time.Now()) + reservation, err = c.buildReservationWithMarketplace(ctx, path, time.Now()) } if err != nil { return err @@ -469,19 +498,106 @@ func (c *client) configureRemotePath(ctx context.Context, path snet.Path) error return nil } -func (c *client) buildReservationWithRedemptions( +// buildReservationWithMarketplace buys the flyovers for the path from a +// marketplace and binds them to the forward path. If a reverse bandwidth was +// requested, the reverse flyovers are bought as well and travel as an extension +// on the forward reservation. +func (c *client) buildReservationWithMarketplace( ctx context.Context, path snet.Path, now time.Time, ) (*snetpath.Reservation, error) { - return redemption.OneShotReservation(ctx, c.sdConn, integration.Local.Host.IP, path, - hummpkg.RedemptionRequestNoHop{ - StartTime: uint32(now.Unix()), - Bw: hummParams.Bw, - Duration: hummParams.Duration, - }, - c.hummParams.ReverseBw, + scionPath, ok := path.Dataplane().(snetpath.SCION) + if !ok { + return nil, serrors.New("provided path must be of type scion") + } + market, err := c.marketplaceClient(ctx) + if err != nil { + return nil, err + } + // Whole seconds: the flyover carries a start time in seconds and a duration + // in seconds, and the marketplace matches the assets it sold on exactly the + // timestamps it was asked for. + startsAt := now.Add(hummStartOffset).Truncate(time.Second) + stopsAt := startsAt.Add(time.Duration(c.hummParams.Duration) * time.Second) + log.Debug("Buying Hummingbird reservation from marketplace", + "url", c.marketParams.URL, + "bandwidth_kbps", c.hummParams.Bw, + "starts_at", startsAt, + "stops_at", stopsAt) + + flyovers, err := market.ObtainReservationsFullPath(ctx, path, c.hummParams.Bw, + startsAt, stopsAt, marketplaceMaxPrice, marketplaceBuyMode, + marketplaceFetchReservations, marketplaceCombineAssets, marketplaceRetries) + if err != nil { + return nil, serrors.Wrap("obtaining reservations from marketplace", err) + } + // One hop per AS on the path, which is what the marketplace client buys for. + expected := len(snetpath.InterfacesToBaseHops(path.Metadata().Interfaces)) + if err := checkFlyovers(flyovers, expected); err != nil { + return nil, serrors.Wrap("checking bought reservations", err) + } + reservation, err := snetpath.NewReservation( + snetpath.WithDataplanePath(scionPath, path.Destination(), flyovers), ) + if err != nil || c.hummParams.ReverseBw == 0 { + return reservation, err + } + + reversePairs := reverseInterfacePairs(path.Metadata().Interfaces) + log.Debug("Buying reverse Hummingbird reservation from marketplace", + "bandwidth_kbps", c.hummParams.ReverseBw) + reverseFlyovers, err := market.ObtainReservationsForInterfacePairs(ctx, reversePairs, + c.hummParams.ReverseBw, startsAt, stopsAt, marketplaceMaxPrice, marketplaceBuyMode, + marketplaceFetchReservations, marketplaceCombineAssets, marketplaceRetries) + if err != nil { + return nil, serrors.Wrap("obtaining reverse reservations from marketplace", err) + } + if err := checkFlyovers(reverseFlyovers, len(reversePairs)); err != nil { + return nil, serrors.Wrap("checking bought reverse reservations", err) + } + extn, err := redemption.BuildReverseReservationExtn(scionPath, path.Source(), reverseFlyovers) + if err != nil { + return nil, err + } + reservation.SetReverseReservationExtn(extn) + return reservation, nil +} + +// checkFlyovers rejects hops that carry no flyover. The marketplace client +// returns a bare hop when redeeming its asset failed, and a path built from +// those would silently travel as a plain SCION path. +func checkFlyovers(hops []*snetpath.Hop, expected int) error { + if len(hops) != expected { + return serrors.New("unexpected number of hops", "expected", expected, "actual", len(hops)) + } + for _, hop := range hops { + if hop == nil { + return serrors.New("missing hop") + } + if hop.Flyover == nil { + return serrors.New("hop without flyover, the asset could not be redeemed", + "ia", hop.IA, "ingress", hop.Ingress, "egress", hop.Egress) + } + } + return nil +} + +// marketplaceClient returns the marketplace client, creating it on first use. +func (c *client) marketplaceClient(ctx context.Context) (*marketclient.MarketplaceClient, error) { + if c.marketClient != nil { + return c.marketClient, nil + } + // The querier and the topology are only used for marketplaces reached over + // SCION, but they are always available here. + querier := daemon.Querier{Connector: c.sdConn, IA: c.topo.LocalIA} + market, err := marketclient.NewMarketplaceClient(ctx, c.marketParams.URL, c.marketParams.JWT, + querier, c.topo, marketplaceInsecure) + if err != nil { + return nil, serrors.Wrap("creating marketplace client", err, "url", c.marketParams.URL) + } + c.marketClient = market + return market, nil } func (c *client) buildReservationWithSecretValues( @@ -498,7 +614,7 @@ func (c *client) buildReservationWithSecretValues( scionPath, path.Destination(), baseHops, - c.hummParams.Bw, + uint16(c.hummParams.Bw), now, ) if err != nil || c.hummParams.ReverseBw == 0 { @@ -506,7 +622,7 @@ func (c *client) buildReservationWithSecretValues( } reverseFlyovers, err := c.deriveFlyoversFromSecretValues( reverseBaseHops(baseHops), - c.hummParams.ReverseBw, + uint16(c.hummParams.ReverseBw), now, ) if err != nil { @@ -584,18 +700,70 @@ func readFrom(conn *snet.Conn, pld []byte) (int, net.Addr, error) { } +// hummingbirdParameters holds the parsed -hummingbird flag. The bandwidths are +// a 16 bit bandwidth class when the flyovers are derived from the AS secret +// values, and kbps when they are bought from a marketplace. type hummingbirdParameters struct { - Bw uint16 + Bw uint32 Duration uint16 - ReverseBw uint16 + ReverseBw uint32 +} + +// bandwidthUnits are the units a bandwidth of the -hummingbird flag can carry, +// and what one of them is in kbps. +var bandwidthUnits = []struct { + suffix string + kbps uint64 +}{ + {suffix: "kbps", kbps: 1}, + {suffix: "mbps", kbps: 1000}, + {suffix: "gbps", kbps: 1000 * 1000}, } -func parseHummingbirdFlag(raw string) (hummingbirdParameters, error) { +// parseBandwidth parses one bandwidth of the -hummingbird flag into kbps. A +// bandwidth bought from a marketplace is a bandwidth and carries a unit; one +// derived from the secret values of the ASes is a bandwidth class, and has none. +func parseBandwidth(raw string, withUnit bool) (uint32, error) { + value := strings.ToLower(strings.TrimSpace(raw)) + for _, unit := range bandwidthUnits { + number, hasUnit := strings.CutSuffix(value, unit.suffix) + if !hasUnit { + continue + } + if !withUnit { + return 0, serrors.New("bandwidth class must not carry a unit", "value", raw) + } + parsed, err := strconv.ParseUint(strings.TrimSpace(number), 10, 32) + if err != nil { + return 0, serrors.Wrap("parsing bandwidth", err, "value", raw) + } + kbps := parsed * unit.kbps + if kbps > math.MaxUint32 { + return 0, serrors.New("bandwidth too large", "value", raw, + "max_kbps", uint64(math.MaxUint32)) + } + return uint32(kbps), nil + } + if withUnit { + return 0, serrors.New("bandwidth must carry a unit", "value", raw, + "units", "kbps|mbps|gbps") + } + parsed, err := strconv.ParseUint(value, 10, 32) + if err != nil { + return 0, serrors.Wrap("parsing bandwidth class", err, "value", raw) + } + return uint32(parsed), nil +} + +// parseHummingbirdFlag parses the -hummingbird flag. The bandwidths carry a unit +// exactly when the reservations are bought from a marketplace, i.e. when no +// -hummKeysDir is given. +func parseHummingbirdFlag(raw string, withUnits bool) (hummingbirdParameters, error) { parts := strings.Split(raw, ",") if len(parts) != 2 && len(parts) != 3 { return hummingbirdParameters{}, serrors.New("expected BW,dur[,reverseBW]") } - bw, err := strconv.ParseUint(parts[0], 10, 16) + bw, err := parseBandwidth(parts[0], withUnits) if err != nil { return hummingbirdParameters{}, serrors.Wrap("parsing hummingbird bandwidth", err, "value", parts[0]) @@ -612,16 +780,68 @@ func parseHummingbirdFlag(raw string) (hummingbirdParameters, error) { ) } params := hummingbirdParameters{ - Bw: uint16(bw), + Bw: bw, Duration: uint16(dur.Seconds()), } if len(parts) == 3 { - reverseBw, err := strconv.ParseUint(parts[2], 10, 16) + reverseBw, err := parseBandwidth(parts[2], withUnits) if err != nil { return hummingbirdParameters{}, serrors.Wrap("parsing reverse hummingbird bandwidth", err, "value", parts[2]) } - params.ReverseBw = uint16(reverseBw) + params.ReverseBw = reverseBw + } + return params, nil +} + +// Environment variables configuring the marketplace used by the -hummingbird +// mode when no -hummKeysDir is given. Both are mandatory. +const ( + envMarketplaceURL = "SCION_MARKETPLACE_URL" + envMarketplaceJWT = "SCION_MARKETPLACE_JWT" +) + +// How the reservations are bought. These are the only values this test needs, so +// they are not configurable. +const ( + // The marketplace serves a self-signed certificate in the local topologies, + // so verifying it cannot succeed. + marketplaceInsecure = true + + // The money is play money, and a purchase that is too expensive is more + // confusing than useful in a test. + marketplaceMaxPrice = uint64(math.MaxUint64) + + // A path is only useful whole, so give up as soon as one AS of it cannot be + // reserved. + marketplaceBuyMode = marketclient.FailOnError + + // Reuse the reservations of an earlier attempt instead of buying again. + marketplaceFetchReservations = true + + // The assets of a local topology span far more than one test run, so there + // is nothing to combine along the time axis. + marketplaceCombineAssets = false + + // Somebody else may buy an asset between searching for it and paying it. + marketplaceRetries = 3 +) + +type marketplaceParameters struct { + URL string + JWT string +} + +func marketplaceParametersFromEnv() (marketplaceParameters, error) { + params := marketplaceParameters{ + URL: os.Getenv(envMarketplaceURL), + JWT: os.Getenv(envMarketplaceJWT), + } + if params.URL == "" { + return params, serrors.New("missing marketplace url", "env", envMarketplaceURL) + } + if params.JWT == "" { + return params, serrors.New("missing marketplace token", "env", envMarketplaceJWT) } return params, nil } @@ -705,6 +925,21 @@ func (c *client) deriveFlyoversFromSecretValues( return flyovers, nil } +// reverseInterfacePairs returns the interface pairs of the reverse direction of +// a path, which is what the reverse reservation must be bought for. +func reverseInterfacePairs(ifaces []snet.PathInterface) []marketclient.InterfacePair { + hops := reverseBaseHops(snetpath.InterfacesToBaseHops(ifaces)) + pairs := make([]marketclient.InterfacePair, 0, len(hops)) + for _, hop := range hops { + pairs = append(pairs, marketclient.InterfacePair{ + IA: uint64(hop.IA), + Ingress: uint32(hop.Ingress), + Egress: uint32(hop.Egress), + }) + } + return pairs +} + func reverseBaseHops(hops []snetpath.BaseHop) []snetpath.BaseHop { reversed := make([]snetpath.BaseHop, len(hops)) for i, hop := range hops { diff --git a/tools/end2end/main_test.go b/tools/end2end/main_test.go index 73b891b6b8..496f14fbeb 100644 --- a/tools/end2end/main_test.go +++ b/tools/end2end/main_test.go @@ -106,3 +106,91 @@ func requireManualDerivationMatchesClient(t *testing.T, dirs [2]string, clientSV manualSV := hummlib.DeriveSecretValue(master0) require.Equal(t, manualSV, clientSV) } + +// TestParseHummingbirdFlag covers the bandwidths of the -hummingbird flag: with +// a marketplace they are bandwidths and carry a unit, with the secret values of +// the ASes they are bandwidth classes and carry none. +func TestParseHummingbirdFlag(t *testing.T) { + t.Parallel() + + testCases := map[string]struct { + raw string + withUnits bool + expected hummingbirdParameters + assertErr require.ErrorAssertionFunc + }{ + "class, no reverse": { + raw: "3,5s", + expected: hummingbirdParameters{Bw: 3, Duration: 5}, + assertErr: require.NoError, + }, + "class, with reverse": { + raw: "3,5s,2", + expected: hummingbirdParameters{Bw: 3, Duration: 5, ReverseBw: 2}, + assertErr: require.NoError, + }, + "class must not carry a unit": { + raw: "3kbps,5s", + assertErr: require.Error, + }, + "kbps": { + raw: "100kbps,20s", + withUnits: true, + expected: hummingbirdParameters{Bw: 100, Duration: 20}, + assertErr: require.NoError, + }, + "mbps and gbps": { + raw: "1mbps,20s,2gbps", + withUnits: true, + expected: hummingbirdParameters{Bw: 1000, Duration: 20, ReverseBw: 2000000}, + assertErr: require.NoError, + }, + "units are case insensitive": { + raw: "100KBPS,20s,1MBps", + withUnits: true, + expected: hummingbirdParameters{Bw: 100, Duration: 20, ReverseBw: 1000}, + assertErr: require.NoError, + }, + "bandwidth needs a unit": { + raw: "100,20s", + withUnits: true, + assertErr: require.Error, + }, + "reverse bandwidth needs a unit too": { + raw: "100kbps,20s,100", + withUnits: true, + assertErr: require.Error, + }, + "unknown unit": { + raw: "100tbps,20s", + withUnits: true, + assertErr: require.Error, + }, + "bandwidth beyond 32 bits of kbps": { + raw: "5000gbps,20s", + withUnits: true, + assertErr: require.Error, + }, + "missing duration": { + raw: "100kbps", + withUnits: true, + assertErr: require.Error, + }, + "duration beyond 16 bits of seconds": { + raw: "100kbps,100000s", + withUnits: true, + assertErr: require.Error, + }, + } + + for name, tc := range testCases { + t.Run(name, func(t *testing.T) { + t.Parallel() + params, err := parseHummingbirdFlag(tc.raw, tc.withUnits) + tc.assertErr(t, err) + if err == nil { + require.Equal(t, tc.expected, params) + } + }) + } +} diff --git a/tools/env/pip3/requirements.in b/tools/env/pip3/requirements.in index 0a073e7ac2..b3a8f83613 100644 --- a/tools/env/pip3/requirements.in +++ b/tools/env/pip3/requirements.in @@ -7,3 +7,4 @@ prometheus-client==0.24.1 requests==2.32.5 # use latest SIX six==1.15.0 +bcrypt==5.0.0 diff --git a/tools/env/pip3/requirements.txt b/tools/env/pip3/requirements.txt index 1283484305..2dd0903940 100644 --- a/tools/env/pip3/requirements.txt +++ b/tools/env/pip3/requirements.txt @@ -4,6 +4,71 @@ # # bazel run //tools/env/pip3:requirements.update # +bcrypt==5.0.0 \ + --hash=sha256:046ad6db88edb3c5ece4369af997938fb1c19d6a699b9c1b27b0db432faae4c4 \ + --hash=sha256:0c418ca99fd47e9c59a301744d63328f17798b5947b0f791e9af3c1c499c2d0a \ + --hash=sha256:0c8e093ea2532601a6f686edbc2c6b2ec24131ff5c52f7610dd64fa4553b5464 \ + --hash=sha256:0cae4cb350934dfd74c020525eeae0a5f79257e8a201c0c176f4b84fdbf2a4b4 \ + --hash=sha256:137c5156524328a24b9fac1cb5db0ba618bc97d11970b39184c1d87dc4bf1746 \ + --hash=sha256:200af71bc25f22006f4069060c88ed36f8aa4ff7f53e67ff04d2ab3f1e79a5b2 \ + --hash=sha256:212139484ab3207b1f0c00633d3be92fef3c5f0af17cad155679d03ff2ee1e41 \ + --hash=sha256:2b732e7d388fa22d48920baa267ba5d97cca38070b69c0e2d37087b381c681fd \ + --hash=sha256:35a77ec55b541e5e583eb3436ffbbf53b0ffa1fa16ca6782279daf95d146dcd9 \ + --hash=sha256:38cac74101777a6a7d3b3e3cfefa57089b5ada650dce2baf0cbdd9d65db22a9e \ + --hash=sha256:3abeb543874b2c0524ff40c57a4e14e5d3a66ff33fb423529c88f180fd756538 \ + --hash=sha256:3ca8a166b1140436e058298a34d88032ab62f15aae1c598580333dc21d27ef10 \ + --hash=sha256:3cf67a804fc66fc217e6914a5635000259fbbbb12e78a99488e4d5ba445a71eb \ + --hash=sha256:4870a52610537037adb382444fefd3706d96d663ac44cbb2f37e3919dca3d7ef \ + --hash=sha256:48f753100931605686f74e27a7b49238122aa761a9aefe9373265b8b7aa43ea4 \ + --hash=sha256:4bfd2a34de661f34d0bda43c3e4e79df586e4716ef401fe31ea39d69d581ef23 \ + --hash=sha256:560ddb6ec730386e7b3b26b8b4c88197aaed924430e7b74666a586ac997249ef \ + --hash=sha256:5b1589f4839a0899c146e8892efe320c0fa096568abd9b95593efac50a87cb75 \ + --hash=sha256:5feebf85a9cefda32966d8171f5db7e3ba964b77fdfe31919622256f80f9cf42 \ + --hash=sha256:611f0a17aa4a25a69362dcc299fda5c8a3d4f160e2abb3831041feb77393a14a \ + --hash=sha256:61afc381250c3182d9078551e3ac3a41da14154fbff647ddf52a769f588c4172 \ + --hash=sha256:64d7ce196203e468c457c37ec22390f1a61c85c6f0b8160fd752940ccfb3a683 \ + --hash=sha256:64ee8434b0da054d830fa8e89e1c8bf30061d539044a39524ff7dec90481e5c2 \ + --hash=sha256:6b8f520b61e8781efee73cba14e3e8c9556ccfb375623f4f97429544734545b4 \ + --hash=sha256:741449132f64b3524e95cd30e5cd3343006ce146088f074f31ab26b94e6c75ba \ + --hash=sha256:744d3c6b164caa658adcb72cb8cc9ad9b4b75c7db507ab4bc2480474a51989da \ + --hash=sha256:79cfa161eda8d2ddf29acad370356b47f02387153b11d46042e93a0a95127493 \ + --hash=sha256:7aeef54b60ceddb6f30ee3db090351ecf0d40ec6e2abf41430997407a46d2254 \ + --hash=sha256:7edda91d5ab52b15636d9c30da87d2cc84f426c72b9dba7a9b4fe142ba11f534 \ + --hash=sha256:7f277a4b3390ab4bebe597800a90da0edae882c6196d3038a73adf446c4f969f \ + --hash=sha256:7f4c94dec1b5ab5d522750cb059bb9409ea8872d4494fd152b53cca99f1ddd8c \ + --hash=sha256:801cad5ccb6b87d1b430f183269b94c24f248dddbbc5c1f78b6ed231743e001c \ + --hash=sha256:83e787d7a84dbbfba6f250dd7a5efd689e935f03dd83b0f919d39349e1f23f83 \ + --hash=sha256:89042e61b5e808b67daf24a434d89bab164d4de1746b37a8d173b6b14f3db9ff \ + --hash=sha256:92864f54fb48b4c718fc92a32825d0e42265a627f956bc0361fe869f1adc3e7d \ + --hash=sha256:9d52ed507c2488eddd6a95bccee4e808d3234fa78dd370e24bac65a21212b861 \ + --hash=sha256:9fffdb387abe6aa775af36ef16f55e318dcda4194ddbf82007a6f21da29de8f5 \ + --hash=sha256:a28bc05039bdf3289d757f49d616ab3efe8cf40d8e8001ccdd621cd4f98f4fc9 \ + --hash=sha256:a5393eae5722bcef046a990b84dff02b954904c36a194f6cfc817d7dca6c6f0b \ + --hash=sha256:a71f70ee269671460b37a449f5ff26982a6f2ba493b3eabdd687b4bf35f875ac \ + --hash=sha256:b17366316c654e1ad0306a6858e189fc835eca39f7eb2cafd6aaca8ce0c40a2e \ + --hash=sha256:baade0a5657654c2984468efb7d6c110db87ea63ef5a4b54732e7e337253e44f \ + --hash=sha256:c2388ca94ffee269b6038d48747f4ce8df0ffbea43f31abfa18ac72f0218effb \ + --hash=sha256:c58b56cdfb03202b3bcc9fd8daee8e8e9b6d7e3163aa97c631dfcfcc24d36c86 \ + --hash=sha256:cde08734f12c6a4e28dc6755cd11d3bdfea608d93d958fffbe95a7026ebe4980 \ + --hash=sha256:d79e5c65dcc9af213594d6f7f1fa2c98ad3fc10431e7aa53c176b441943efbdd \ + --hash=sha256:d8d65b564ec849643d9f7ea05c6d9f0cd7ca23bdd4ac0c2dbef1104ab504543d \ + --hash=sha256:db99dca3b1fdc3db87d7c57eac0c82281242d1eabf19dcb8a6b10eb29a2e72d1 \ + --hash=sha256:dcd58e2b3a908b5ecc9b9df2f0085592506ac2d5110786018ee5e160f28e0911 \ + --hash=sha256:dd19cf5184a90c873009244586396a6a884d591a5323f0e8a5922560718d4993 \ + --hash=sha256:ddb4e1500f6efdd402218ffe34d040a1196c072e07929b9820f363a1fd1f4191 \ + --hash=sha256:e3cf5b2560c7b5a142286f69bde914494b6d8f901aaa71e453078388a50881c4 \ + --hash=sha256:ed2e1365e31fc73f1825fa830f1c8f8917ca1b3ca6185773b349c20fd606cec2 \ + --hash=sha256:edfcdcedd0d0f05850c52ba3127b1fce70b9f89e0fe5ff16517df7e81fa3cbb8 \ + --hash=sha256:f0ce778135f60799d89c9693b9b398819d15f1921ba15fe719acb3178215a7db \ + --hash=sha256:f2347d3534e76bf50bca5500989d6c1d05ed64b440408057a37673282c654927 \ + --hash=sha256:f3c08197f3039bec79cee59a606d62b96b16669cff3949f21e74796b6e3cd2be \ + --hash=sha256:f632fd56fc4e61564f78b46a2269153122db34988e78b6be8b32d28507b7eaeb \ + --hash=sha256:f6984a24db30548fd39a44360532898c33528b74aedf81c26cf29c51ee47057e \ + --hash=sha256:f70aadb7a809305226daedf75d90379c397b094755a710d7014b8b117df1ebbf \ + --hash=sha256:f748f7c2d6fd375cc93d3fba7ef4a9e3a092421b8dbf34d8d4dc06be9492dfdd \ + --hash=sha256:f8429e1c410b4073944f03bd778a9e066e7fad723564a52ff91841d278dfc822 \ + --hash=sha256:fc746432b951e92b58317af8e0ca746efe93e66555f1b40888865ef5bf56446b + # via -r tools/env/pip3/requirements.in certifi==2026.1.4 \ --hash=sha256:9943707519e4add1115f44c2bc244f782c0249876bf51b6599fee1ffbedd685c \ --hash=sha256:ac726dd470482006e014ad384921ed6438c457018f4b3d204aea4281258b2120 diff --git a/tools/topogen.py b/tools/topogen.py index db365f43fb..31d0be4eae 100755 --- a/tools/topogen.py +++ b/tools/topogen.py @@ -42,6 +42,9 @@ def add_arguments(parser): help='IPv6 network to create subnets in (E.g. "fd00:f00d:cafe::7f00:0000/104"') parser.add_argument('-o', '--output-dir', default=GEN_PATH, help='Output directory') + parser.add_argument('-m', '--marketplace', metavar='ISD-AS', + help='Add a marketplace to the given AS of the topology, and fill its \ + database with the default users, assets and redemption delegations') parser.add_argument('-t', '--topology-jsons-only', action='store_true', help='Create only topology.json files') parser.add_argument('--random-ifids', action='store_true', diff --git a/tools/topology/BUILD.bazel b/tools/topology/BUILD.bazel index 1a394a4b95..6df4f6fa7d 100644 --- a/tools/topology/BUILD.bazel +++ b/tools/topology/BUILD.bazel @@ -7,6 +7,7 @@ py_library( name = "py_default_library", srcs = glob(["**/*.py"]), deps = [ + requirement("bcrypt"), requirement("pyyaml"), ], ) diff --git a/tools/topology/config.py b/tools/topology/config.py index 8460b56411..7a743f3fbd 100644 --- a/tools/topology/config.py +++ b/tools/topology/config.py @@ -45,6 +45,11 @@ SubnetGenerator, DEFAULT_NETWORK, ) +from topology.marketplace import ( + MarketplaceError, + MarketplaceGenArgs, + MarketplaceGenerator, +) from topology.monitoring import MonitoringGenArgs, MonitoringGenerator from topology.supervisor import SupervisorGenArgs, SupervisorGenerator from topology.topo import TopoGenArgs, TopoGenerator @@ -96,6 +101,7 @@ def generate_all(self): """ self._ensure_uniq_ases() self._canonicalize_isd_asns() + self._ensure_marketplace_as() topo_dicts, self.all_networks = self._generate_topology() if not self.args.topology_jsons_only: self.networks = remove_v4_nets(self.all_networks) @@ -112,6 +118,22 @@ def _ensure_uniq_ases(self): sys.exit(1) seen.add(ia.as_str()) + def _ensure_marketplace_as(self): + if not self.args.marketplace: + return + try: + ia = str(ISD_AS(self.args.marketplace)) + except ValueError: + logging.critical("Invalid ISD-AS '%s' for the marketplace", self.args.marketplace) + sys.exit(1) + if ia not in self.topo_config["ASes"]: + logging.critical("The marketplace AS '%s' is not part of %s", ia, + self.args.topo_config) + sys.exit(1) + # The marketplace is canonicalized like the ASes it is checked against, so + # that the generators can compare it to a TopoID. + self.args.marketplace = ia + def _canonicalize_isd_asns(self): canonicalized = {} for asStr, value in self.topo_config["ASes"].items(): @@ -126,6 +148,21 @@ def _generate_with_topo(self, topo_dicts): self._generate_supervisor(topo_dicts) self._generate_monitoring_conf(topo_dicts) self._generate_certs_trcs(topo_dicts) + if self.args.marketplace: + self._generate_marketplace(topo_dicts) + + def _generate_marketplace(self, topo_dicts): + # Last, because the redemption delegations of the marketplace are derived + # from the master keys that _generate_certs_trcs writes. + gen = MarketplaceGenerator(self._marketplace_args(topo_dicts)) + try: + gen.generate() + except MarketplaceError as e: + logging.critical("Cannot add the marketplace to %s: %s", self.args.marketplace, e) + sys.exit(1) + + def _marketplace_args(self, topo_dicts): + return MarketplaceGenArgs(self.args, topo_dicts, self.networks) def _generate_certs_trcs(self, topo_dicts): certgen = CertGenerator(self._cert_args()) diff --git a/tools/topology/docker.py b/tools/topology/docker.py index a96abcc56f..811a616afd 100644 --- a/tools/topology/docker.py +++ b/tools/topology/docker.py @@ -30,6 +30,11 @@ sciond_name, ) from topology.docker_utils import DockerUtilsGenArgs, DockerUtilsGenerator +from topology.marketplace import ( + dockerService, + dockerServiceName, + hosts_marketplace, +) from topology.net import NetworkDescription, IPNetwork from topology.sig import SIGGenArgs, SIGGenerator @@ -93,6 +98,7 @@ def _gen_topo(self, topo_id, topo, base): self._br_conf(topo_id, topo, base) self._control_service_conf(topo_id, topo, base) self._hummingbird_conf(topo_id, topo, base) + self._marketplace_conf(topo_id, topo, base) self._sciond_conf(topo_id, base) def _gen_sig(self): @@ -222,6 +228,26 @@ def _hummingbird_conf(self, topo_id, topo, base): } self.dc_conf['services'][name] = entry + def _marketplace_conf(self, topo_id, topo, base): + if not hosts_marketplace(self.args, topo_id): + return + for k in topo.get("control_service", {}): + if not k.endswith("-1"): + continue + self.dc_conf['services'][dockerServiceName(topo_id)] = dockerService( + docker_image(self.args, 'marketplace'), + k, + 'disp_%s' % k, + [ + self._cache_vol(), + # The marketplace writes into its config dir: + # its TLS certificate and its JWT signature keys are created on first run. + '%s:/etc/scion:rw' % base, + ], + self.user, + ) + break + def _dispatcher_conf(self, topo_id, topo, base): image = 'dispatcher' base_entry = { diff --git a/tools/topology/marketplace.py b/tools/topology/marketplace.py new file mode 100644 index 0000000000..03db1fd233 --- /dev/null +++ b/tools/topology/marketplace.py @@ -0,0 +1,786 @@ +# Copyright 2026 ETH Zurich +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +""" +:mod:`marketplace` --- SCION topology marketplace generator +=========================================================== + +Everything a marketplace needs to run on a generated topology: +its config file, the entries advertising it to the other ASes, +the service or program running it, +and the entries its database starts out with. + +Besides the topology generator, this module is imported by marketplace/tools/setup_marketplace.py, +which adds a marketplace to an already generated topology, +and by marketplace/tools/populate_marketplace.py, +which rebuilds its database. +Both run on the system interpreter, so the module keeps to the standard library and bcrypt. +""" +# Stdlib +import base64 +import hashlib +import json +import math +import os +import re +import shlex +import sqlite3 +from collections import namedtuple +from datetime import datetime, timedelta, timezone +from pathlib import Path + +# External packages +import bcrypt + +# SCION +from topology.common import ArgsTopoDicts, TopoID +from topology.scion_addr import ISD_AS +from topology.util import write_file + +MARKETPLACE_CONFIG_NAME = 'marketplace.toml' +STATIC_INFO_CONFIG_NAME = 'staticInfoConfig.json' +MARKETPLACE_DB_NAME = 'marketplace.db' + +# The schema the marketplace applies itself on startup. Prepopulating the +# database means applying it here, before the marketplace has ever run. +MARKETPLACE_SCHEMA = os.path.join('marketplace', 'db', 'schema.sql') + +# The port of both APIs, the TCP one serving the ConnectRPC API and the web app +# from a single mux, and the SCION/QUIC one, which is UDP. They do not collide. +# It has to be inside the dispatched_ports range of the AS, otherwise only the +# shim dispatcher can deliver SCION packets to it. See DefaultAddr in +# marketplace/config.go. +MARKETPLACE_PORT = 31888 + +# Where the marketplace of a docker topology finds its files inside the container. +CONTAINER_CONFIG_DIR = '/etc/scion' +CONTAINER_CACHE_DIR = '/share/cache' +CONTAINER_TOML = '%s/%s' % (CONTAINER_CONFIG_DIR, MARKETPLACE_CONFIG_NAME) +CONTAINER_DB = '%s/%s' % (CONTAINER_CACHE_DIR, MARKETPLACE_DB_NAME) + +# The marketplace of a supervisord topology binds the loopback address, and keeps +# its database next to the ones of the other services. +LOCAL_HOST = '127.0.0.1' +LOCAL_CACHE_DIR = 'gen-cache' + +# The marketplace program of a supervisord config. There is a single marketplace, +# so unlike the per-AS services it needs no IA in its name. +PROGRAM_NAME = 'marketplace' + +# What the marketplace binds, and what staticInfoConfig.json advertises for it. +Endpoints = namedtuple('Endpoints', 'host api_port scion_port') + + +class MarketplaceError(Exception): + """An error that aborts the marketplace setup with a message.""" + + +# The config of the marketplace, as a template rather than a dict, so that +# setup_marketplace.py can tell an unchanged file from a locally edited one. +MARKETPLACE_TOML = """[general] +id = "marketplace" +config_dir = "{config_dir}" + +[log.console] +level = "debug" + +[marketplace] +api_addr = "{api_addr}" +scion_api_addr = "{scion_api_addr}" +currency = "CHF" +currency_exponent = 2 +supports_redemption_delegation = true +delegation_hourly_fee = 10000 +transaction_fee_relative = 0.01 +transaction_fee_absolute = 10 +split_combine_fee_absolute = 10 + +[marketplace_db] +connection = "{db_path}" +""" + + +def hostPort(host: str, port: int) -> str: + """host:port, with an IPv6 host in brackets.""" + if ":" in host: + return "[%s]:%d" % (host, port) + return "%s:%d" % (host, port) + + +def iaNumbers(ia): + """The ISD and AS numbers of an IA, the way the marketplace database stores them.""" + try: + isd_as = ISD_AS(str(ia)) + except ValueError: + raise MarketplaceError("invalid ISD-AS %r" % str(ia)) + as_str = isd_as.as_str() + if ":" not in as_str: + return int(isd_as.isd_str()), int(as_str) + asn = 0 + for part in as_str.split(":"): + asn = (asn << 16) | int(part, 16) + return int(isd_as.isd_str()), asn + + +def marketplaceToml(config_dir: str, endpoints: Endpoints, db_path: str) -> str: + return MARKETPLACE_TOML.format( + config_dir=config_dir, + api_addr=hostPort(endpoints.host, endpoints.api_port), + scion_api_addr=hostPort(endpoints.host, endpoints.scion_port), + db_path=db_path, + ) + + +def marketplaceEntries(ia, endpoints: Endpoints): + """The staticInfoConfig entries advertising this marketplace. + + They describe what the marketplace really binds: the web app is served by the + same mux as the TCP API, so both live on the API port. + """ + website = "https://%s" % hostPort(endpoints.host, endpoints.api_port) + return [ + { + "name": "Test Market", + "protocol": "connectrpc/TLS/QUIC/SCION", + "api": "[%s,%s]:%d" % (ISD_AS(str(ia)), endpoints.host, endpoints.scion_port), + "website": website, + }, + { + "name": "Test Market", + "protocol": "connectrpc/TLS/TCP", + "api": website, + "website": website, + }, + ] + + +def loadStaticInfo(path): + """Reads staticInfoConfig.json and its embedded note object. + + Returns (static_info, note), both dicts. A missing file yields empty ones. + The note is a JSON document of its own, encoded as a string, because the + control service copies it verbatim into the beacons it propagates. + """ + path = Path(path) + if not path.exists(): + return {}, {} + + try: + static_info = json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError) as e: + raise MarketplaceError("cannot read %s: %s" % (path, e)) + if not isinstance(static_info, dict): + raise MarketplaceError( + "%s: expected a JSON object, got %s" % (path, type(static_info).__name__)) + + note = static_info.get("note", "") + if note == "": + return static_info, {} + try: + note = json.loads(note) + except (TypeError, json.JSONDecodeError): + raise MarketplaceError( + "%s: the 'note' field is not a JSON object; refusing to overwrite it" % path) + if not isinstance(note, dict): + raise MarketplaceError( + "%s: the 'note' field is not a JSON object; refusing to overwrite it" % path) + return static_info, note + + +def mergeEntries(existing, entries): + """Merges the marketplace entries into the existing ones, keyed by name+protocol.""" + merged = [e for e in existing if isinstance(e, dict)] + changed = len(merged) != len(existing) + for entry in entries: + key = (entry["name"], entry["protocol"]) + for i, old in enumerate(merged): + if (old.get("name"), old.get("protocol")) == key: + if old != entry: + merged[i] = entry + changed = True + break + else: + merged.append(entry) + changed = True + return merged, changed + + +def advertiseMarketplace(path, ia, endpoints: Endpoints) -> bool: + """Advertises the marketplace in the staticInfoConfig.json at path. + + Whatever else the file already configures is kept, and so are the entries of + other marketplaces. Returns whether the file was written. + """ + static_info, note = loadStaticInfo(path) + hummingbird = note.get("hummingbird", []) + if not isinstance(hummingbird, list): + raise MarketplaceError( + "%s: 'note'.hummingbird is not a list; refusing to overwrite it" % path) + + hummingbird, changed = mergeEntries(hummingbird, marketplaceEntries(ia, endpoints)) + if not changed: + return False + note["hummingbird"] = hummingbird + static_info["note"] = json.dumps(note) + try: + write_file(str(path), json.dumps(static_info, indent=2)) + except OSError as e: + raise MarketplaceError("cannot write %s: %s" % (path, e)) + return True + + +def checkDispatchedPorts(as_dir, port: int) -> None: + """Checks that the SCION API port is one the border router delivers directly.""" + topology = Path(as_dir) / "topology.json" + if not topology.is_file(): + return + try: + ports = json.loads(topology.read_text(encoding="utf-8")).get("dispatched_ports") + except (OSError, json.JSONDecodeError) as e: + raise MarketplaceError("cannot read %s: %s" % (topology, e)) + if not ports or ports == "all": + return + match = re.fullmatch(r"(\d+)-(\d+)", str(ports).strip()) + if match is None: + return + start, end = int(match.group(1)), int(match.group(2)) + if not start <= port <= end: + raise MarketplaceError( + "the SCION API port %d is outside the dispatched_ports range %s of %s; the " + "marketplace would only receive SCION packets through a shim dispatcher" + % (port, ports, topology)) + + +def dockerServiceName(topo_id) -> str: + return 'marketplace%s' % TopoID(str(topo_id)).file_fmt() + + +def dockerService(image, control_name, dispatcher_name, volumes, user=None): + """The compose service running the marketplace of an AS. + + It shares the network namespace of the control service, the way the control + service and the hummingbird service already share the one of their dispatcher. + It therefore carries no address of its own, and needs none. + """ + entry = { + 'command': ['--config', CONTAINER_TOML], + # The marketplace asks the control service for trust material on startup. + 'depends_on': {control_name: {'condition': 'service_healthy'}}, + 'image': image, + 'network_mode': 'service:%s' % dispatcher_name, + 'volumes': volumes, + } + if user is not None: + entry['user'] = user + return entry + + +def supervisordProgram(config_path: str): + """The supervisord settings of the marketplace program.""" + return { + 'autostart': 'false', + 'autorestart': 'false', + # Like the router, the marketplace is a cgo binary. + 'environment': 'TZ=UTC,GODEBUG="cgocheck=0"', + 'stdout_logfile': 'logs/%s.log' % PROGRAM_NAME, + 'redirect_stderr': True, + 'startretries': 0, + 'startsecs': 5, + 'priority': 100, + 'command': ' '.join( + shlex.quote(a) for a in ['bin/marketplace', '--config', config_path]), + } + + +# +# The default entries of the database. +# + +DEFAULT_PASSWORD = "1234" +DEFAULT_USERS = ["alice", "bob"] +DEFAULT_USER_BALANCE = 1000000 +DEFAULT_AS_BALANCE = 0 +DEFAULT_ASSET_BANDWIDTH = 1000000 +DEFAULT_ASSET_BANDWIDTH_MIN = 10 +DEFAULT_ASSET_BANDWIDTH_MAX = 1000000 +DEFAULT_ASSET_PRICE = 1 +DEFAULT_ASSET_TIME_GRANULARITY = 10 +DEFAULT_ASSET_TIME_MIN_DURATION = 10 +DEFAULT_ASSET_DURATION = timedelta(days=100) +DEFAULT_DELEGATION_RES_ID_LIMIT = 1000 + +# The flyover carries the bandwidth as a 10 bit codepoint into the points +# published by the AS, so there is one point per codepoint. The values must be +# the ones the router decodes, see ConvertBW in router/tokenbucket/tokenbucket.go. +BW_CODEPOINTS = 1 << 10 +MIN_BW_KBPS = 10 +MAX_BW_KBPS = 10_000_000 +BW_LOG_ENCODING_START = 60 + +# Derivation of the Hummingbird secret value of an AS from its master key, see +# DeriveSecretValue in pkg/slayers/path/hummingbird/mac.go. +SECRET_VALUE_SALT = b"Derive hbird sv" +SECRET_VALUE_ITERATIONS = 1000 +SECRET_VALUE_LENGTH = 16 + + +def loadJson(path): + with open(path, "r", encoding="utf-8") as file: + return json.load(file) + + +def loadTopology(gen_dir): + """Reads the IA, interface IDs and directory of every AS of a topology.""" + topologies = sorted(Path(gen_dir).glob("AS*/topology.json")) + if not topologies: + raise MarketplaceError("no AS*/topology.json found in %s" % gen_dir) + ases = [] + for path in topologies: + try: + topo = loadJson(path) + except (OSError, json.JSONDecodeError) as e: + raise MarketplaceError("cannot read %s: %s" % (path, e)) + ia = topo.get("isd_as") + if ia is None: + raise MarketplaceError("%s has no isd_as" % path) + ifids = sorted( + int(ifid) + for router in topo.get("border_routers", {}).values() + for ifid in router.get("interfaces", {}) + ) + ases.append((ia, ifids, path.parent)) + return ases + + +def encodingPoints(): + """The bandwidth points of an AS, in kbps, one per codepoint. + + The first BW_LOG_ENCODING_START points are one kbps apart, starting at + MIN_BW_KBPS, and the rest grow geometrically up to MAX_BW_KBPS. The router + decodes a flyover with the very same points, so publishing anything else + would sell a bandwidth that the router does not enforce. + """ + step = (MAX_BW_KBPS / (MIN_BW_KBPS + BW_LOG_ENCODING_START)) ** ( + 1.0 / (BW_CODEPOINTS - BW_LOG_ENCODING_START - 1) + ) + points = [] + for codepoint in range(BW_CODEPOINTS): + if codepoint < BW_LOG_ENCODING_START: + points.append(MIN_BW_KBPS + codepoint) + else: + kbps = (MIN_BW_KBPS + BW_LOG_ENCODING_START) * step ** ( + codepoint - BW_LOG_ENCODING_START + ) + points.append(math.ceil(kbps)) + return points + + +def secretValue(as_dir): + """Derives the Hummingbird secret value of an AS from its master key. + + Handing it to the marketplace is what lets the marketplace redeem the assets + of that AS, i.e. derive the authenticator of a flyover on its behalf. + """ + path = Path(as_dir) / "keys" / "master0.key" + try: + master = base64.b64decode(path.read_text(encoding="utf-8").strip(), validate=True) + except OSError as e: + raise MarketplaceError("cannot read the master key %s: %s" % (path, e)) + except ValueError as e: + raise MarketplaceError("%s is not a base64 encoded key: %s" % (path, e)) + return hashlib.pbkdf2_hmac( + "sha256", master, SECRET_VALUE_SALT, SECRET_VALUE_ITERATIONS, SECRET_VALUE_LENGTH + ) + + +def interfacePairs(ifids): + """The (ingress, egress) pairs of an AS. + + Interface 0 stands for no interface, i.e. a flyover that starts or ends in + this AS, so that ASes with a single interface get assets as well. + """ + pairs = [(i, e) for i in ifids for e in ifids if i != e] + pairs += [(0, e) for e in ifids] + pairs += [(i, 0) for i in ifids] + return pairs + + +def defaultEntries(gen_dir, now=None): + """Builds the default users, ASes, assets and delegations for a topology.""" + if now is None: + now = datetime.now(timezone.utc) + startsAt = now.replace(microsecond=0).strftime("%Y-%m-%dT%H:%M:%SZ") + stopsAt = (now + DEFAULT_ASSET_DURATION).replace(microsecond=0).strftime("%Y-%m-%dT%H:%M:%SZ") + + ases = loadTopology(gen_dir) + assets = [ + { + "ia": ia, + "bandwidth": DEFAULT_ASSET_BANDWIDTH, + "bandwidth_min": DEFAULT_ASSET_BANDWIDTH_MIN, + "bandwidth_max": DEFAULT_ASSET_BANDWIDTH_MAX, + "price": DEFAULT_ASSET_PRICE, + "time_granularity": DEFAULT_ASSET_TIME_GRANULARITY, + "time_min_duration": DEFAULT_ASSET_TIME_MIN_DURATION, + "starts_at": startsAt, + "stops_at": stopsAt, + "ingress": ingress, + "egress": egress, + } + for ia, ifids, _ in ases + for ingress, egress in interfacePairs(ifids) + ] + return { + "users": [ + {"name": name, "password": DEFAULT_PASSWORD} + for name in DEFAULT_USERS + ], + "accounts": [ + {"user": name, "balance": DEFAULT_USER_BALANCE} + for name in DEFAULT_USERS + ], + "ases": [ + {"ia": ia, "password": DEFAULT_PASSWORD, "balance": DEFAULT_AS_BALANCE} + for ia, _, _ in ases + ], + "assets": assets, + # Every AS delegates the redemption of its assets to the marketplace. + # Without this the marketplace cannot derive the flyover authenticators, + # and the assets it sells are worthless. + "delegations": [ + { + "ia": ia, + "res_id_limit": DEFAULT_DELEGATION_RES_ID_LIMIT, + "expiration": stopsAt, + "paid_until": stopsAt, + "key": secretValue(as_dir).hex(), + "encodings": encodingPoints(), + } + for ia, _, as_dir in ases + ], + } + + +# +# Filling the database. +# + +def schemaHash(schema_path): + """The hash identifying a schema, i.e. the version an import has to match.""" + try: + with open(schema_path, "r", encoding="utf-8") as f: + sql = f.read() + except OSError as e: + raise MarketplaceError("cannot read the schema %s: %s" % (schema_path, e)) + normalized = re.sub(r"\s+", " ", sql).strip() + return hashlib.sha256(normalized.encode("utf-8")).hexdigest() + + +def verifySchema(schema_path, version): + actual = schemaHash(schema_path) + if version != actual: + raise MarketplaceError("schema version mismatch; expected %s" % actual) + + +def applySchema(db, schema_path): + try: + with open(schema_path, "r", encoding="utf-8") as f: + db.executescript(f.read()) + except OSError as e: + raise MarketplaceError("cannot read the schema %s: %s" % (schema_path, e)) + + +def hashPassword(password): + if password is None: + return None + return bcrypt.hashpw(password.encode("utf-8"), bcrypt.gensalt(rounds=12)).decode("utf-8") + + +def dropTables(db): + """Drops every table, so that the schema and the entries are recreated. + + Indexes are dropped along with their table, and so are the AUTOINCREMENT + counters that SQLite keeps in its internal sqlite_sequence table. + """ + tables = [ + name for (name,) in db.execute( + "SELECT name FROM sqlite_master WHERE type = 'table' AND name NOT LIKE 'sqlite_%'" + ) + ] + for table in tables: + db.execute('DROP TABLE "%s"' % table) + return tables + + +def insertUsers(db, users): + if users is None: + return + db.executemany( + """ + INSERT INTO Users (name, pw_hash) + VALUES (?, ?) + ON CONFLICT(name) DO UPDATE SET pw_hash = excluded.pw_hash + """, + [(u.get("name"), hashPassword(u.get("password"))) for u in users], + ) + + +def insertAccounts(db, accounts): + if accounts is None: + return + # Updated instead of replaced: a replace would delete the conflicting row and + # insert it under a new id, orphaning the assets and reservations that refer + # to it. + db.executemany( + """ + INSERT INTO Accounts (scope, balance, user_id) + VALUES (?, ?, ( + SELECT id + FROM Users + WHERE name = ? + LIMIT 1 + )) + ON CONFLICT(user_id, scope) DO UPDATE SET balance = excluded.balance + """, + [(a.get("scope") or "", a.get("balance"), a.get("user")) for a in accounts], + ) + + +def insertASes(db, ases): + if ases is None: + return + rows = [] + for a in ases: + isd_id, as_id = iaNumbers(a.get("ia")) + rows.append(( + isd_id, + as_id, + hashPassword(a.get("password")) or "", + a.get("jwt_version") or 0, + a.get("balance") or 0, + )) + db.executemany( + """ + INSERT OR REPLACE INTO Ases (isd_id, as_id, pw_hash, jwt_version, balance) + VALUES (?, ?, ?, ?, ?) + """, + rows, + ) + + +def insertAssets(db, assets): + if assets is None: + return + rows = [] + for a in assets: + isd_id, as_id = iaNumbers(a.get("ia")) + rows.append(( + isd_id, as_id, a.get("bandwidth"), a.get("bandwidth_min"), a.get("bandwidth_max"), + a.get("price"), a.get("time_granularity"), a.get("time_min_duration"), + a.get("starts_at"), a.get("stops_at"), a.get("ingress"), a.get("egress"), + a.get("owner"), + )) + db.executemany( + """ + INSERT INTO Assets (isd_id, as_id, bandwidth, bandwidth_min, bandwidth_max, price, + time_granularity, time_min_duration, starts_at, stops_at, + ingress, egress, account_id) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ( + SELECT a.id + FROM Accounts a + JOIN Users u ON u.ID = a.user_id + WHERE u.name = ? AND a.scope = '' + LIMIT 1 + )) + """, + rows, + ) + + +def insertReservations(db, reservations): + if reservations is None: + return + rows = [] + for r in reservations: + isd_id, as_id = iaNumbers(r.get("ia")) + rows.append(( + r.get("id"), isd_id, as_id, r.get("ingress"), r.get("egress"), r.get("bandwidth"), + r.get("bw_encoded"), r.get("starts_at"), r.get("stops_at"), + bytes.fromhex(r.get("key")), r.get("owner"), + )) + db.executemany( + """ + INSERT OR REPLACE INTO Reservations (reservation_id, isd_id, as_id, ingress, egress, + bandwidth, bw_encoded, starts_at, stops_at, + key, account_id) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ( + SELECT id + FROM Users + WHERE name = ? + LIMIT 1 + )) + """, + rows, + ) + + +def insertDelegations(db, delegations): + if delegations is None: + return + rows = [] + for d in delegations: + isd_id, as_id = iaNumbers(d.get("ia")) + encodings = b"".join( + i.to_bytes(4, byteorder="little") for i in d.get("encodings") + ) + rows.append(( + isd_id, as_id, d.get("res_id_limit"), d.get("expiration"), d.get("paid_until"), + bytes.fromhex(d.get("key")), encodings, + )) + db.executemany( + """ + INSERT OR REPLACE INTO Redemption_Delegations (isd_id, as_id, res_id_limit, expiration, + paid_until, key, encodings) + VALUES (?, ?, ?, ?, ?, ?, ?) + """, + rows, + ) + + +def insertAll(db, data): + insertUsers(db, data.get("users")) + insertAccounts(db, data.get("accounts")) + insertASes(db, data.get("ases")) + insertAssets(db, data.get("assets")) + insertReservations(db, data.get("reservations")) + insertDelegations(db, data.get("delegations")) + + +def populateDB(db_path, schema_path, data): + """Rebuilds the marketplace database from the given entries. + + Every table is dropped first, so that both the schema and the entries are the + ones of this run, and not the leftovers of a previous one. Returns the names + of the tables that were dropped. + """ + Path(db_path).parent.mkdir(parents=True, exist_ok=True) + try: + conn = sqlite3.connect(db_path) + except sqlite3.Error as e: + raise MarketplaceError("cannot open the database %s: %s" % (db_path, e)) + try: + cursor = conn.cursor() + dropped = dropTables(cursor) + applySchema(cursor, schema_path) + insertAll(cursor, data) + conn.commit() + return dropped + except sqlite3.Error as e: + conn.rollback() + raise MarketplaceError("cannot fill the database %s: %s" % (db_path, e)) + except Exception: + conn.rollback() + raise + finally: + conn.close() + + +# +# The generator, run as part of the topology generation. +# + +class MarketplaceGenArgs(ArgsTopoDicts): + def __init__(self, args, topo_dicts, networks): + """ + :param object args: Contains the passed command line arguments as named attributes. + :param dict topo_dicts: The generated topo dicts from TopoGenerator. + :param dict networks: The generated networks from SubnetGenerator. + """ + super().__init__(args, topo_dicts) + self.networks = networks + + +class MarketplaceGenerator(object): + def __init__(self, args): + """ + :param MarketplaceGenArgs args: Contains the passed command line arguments, + the topo dicts and the networks. + """ + self.args = args + self.output_base = os.environ.get('SCION_OUTPUT_BASE', os.getcwd()) + + def generate(self): + """Writes the config of the marketplace and fills its database. + + The service or program running the marketplace is added by the backend + generator, so that it lands in the compose file or the supervisord config + along with the services of its AS. + """ + topo_id = marketplace_topo_id(self.args) + base = topo_id.base_dir(self.args.output_dir) + endpoints = self._endpoints(topo_id) + checkDispatchedPorts(base, endpoints.scion_port) + + if self.args.docker: + config_dir, db_path = CONTAINER_CONFIG_DIR, CONTAINER_DB + else: + config_dir = base + db_path = os.path.join(LOCAL_CACHE_DIR, MARKETPLACE_DB_NAME) + write_file(os.path.join(base, MARKETPLACE_CONFIG_NAME), + marketplaceToml(config_dir, endpoints, db_path)) + + # The control service copies this into the beacons it propagates, which is + # how the other ASes learn where to buy the assets of this one. + advertiseMarketplace( + os.path.join(base, STATIC_INFO_CONFIG_NAME), topo_id, endpoints) + + # The marketplace applies the schema itself, but it cannot invent the + # entries: they are the ones a local topology is expected to start with. + populateDB( + os.path.join(self.output_base, LOCAL_CACHE_DIR, MARKETPLACE_DB_NAME), + MARKETPLACE_SCHEMA, + defaultEntries(self.args.output_dir), + ) + + def _endpoints(self, topo_id) -> Endpoints: + if not self.args.docker: + return Endpoints(LOCAL_HOST, MARKETPLACE_PORT, MARKETPLACE_PORT) + # The marketplace joins the network namespace of the control service, so + # it binds the very same address, and only its ports set it apart. + name = control_service_name(topo_id) + for net_desc in self.args.networks.values(): + if name in net_desc.ip_net: + return Endpoints(str(net_desc.ip_net[name].ip), + MARKETPLACE_PORT, MARKETPLACE_PORT) + raise MarketplaceError("no address generated for %s" % name) + + +def control_service_name(topo_id) -> str: + """The control service the marketplace of an AS shares its address with. + + Only a single control service instance per AS is currently supported, and the + dispatcher owning the network namespace is named after it. + """ + return 'cs%s-1' % TopoID(str(topo_id)).file_fmt() + + +def marketplace_topo_id(args): + """The TopoID of the AS hosting the marketplace, or None if there is none.""" + if not getattr(args, 'marketplace', None): + return None + return TopoID(args.marketplace) + + +def hosts_marketplace(args, topo_id) -> bool: + """Whether topo_id is the AS hosting the marketplace.""" + return marketplace_topo_id(args) == topo_id diff --git a/tools/topology/supervisor.py b/tools/topology/supervisor.py index 9542c745fe..318aaf18a6 100644 --- a/tools/topology/supervisor.py +++ b/tools/topology/supervisor.py @@ -23,6 +23,12 @@ from io import StringIO # SCION +from topology.marketplace import ( + MARKETPLACE_CONFIG_NAME, + PROGRAM_NAME as MARKETPLACE_PROGRAM_NAME, + hosts_marketplace, + supervisordProgram, +) from topology.util import write_file from topology.common import ( ArgsTopoDicts, @@ -71,6 +77,7 @@ def _as_entries(self, topo_id, topo): entries.extend(self._br_entries(topo, "bin/router", base)) entries.extend(self._control_service_entries(topo, base)) entries.extend(self._hummingbird_entries(topo_id, topo, base)) + entries.extend(self._marketplace_entries(topo_id, base)) entries.append(self._sciond_entry(topo_id, base)) return entries @@ -130,6 +137,12 @@ def _hummingbird_entries(self, topo_id, topo, base): entries.append((name, self._common_entry(name, cmd_args))) return entries + def _marketplace_entries(self, topo_id, base): + if not hosts_marketplace(self.args, topo_id): + return [] + conf = os.path.join(base, MARKETPLACE_CONFIG_NAME) + return [(MARKETPLACE_PROGRAM_NAME, supervisordProgram(conf))] + def _add_dispatcher(self, config): name, entry = self._dispatcher_entry() self._add_prog(config, name, entry)