Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion ai-docs/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -322,7 +322,7 @@ store. Epoch IDs are compared as big-endian integers.

### Centralized

A single `RealCentralizedKms` instance holds all key material. No MPC; keys
A single `CentralizedKms` instance holds all key material. No MPC; keys
live in the configured vault backend. Preprocessing / reshare RPCs are not
applicable.

Expand Down
2 changes: 1 addition & 1 deletion charts/kms-core/templates/kms-core-configmap.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -210,7 +210,7 @@ data:

[rate_limiter_conf]
bucket_size = {{ .Values.kmsCore.rateLimiter.bucketSize }}
pub_decrypt = 1
pub_decrypt = {{ int .Values.kmsCore.rateLimiter.pubDecrypt }}
user_decrypt = 1
crsgen = 100
preproc = 50000
Expand Down
3 changes: 2 additions & 1 deletion charts/kms-core/values.schema.json
Original file line number Diff line number Diff line change
Expand Up @@ -304,7 +304,8 @@
"rateLimiter": {
"type": "object",
"properties": {
"bucketSize": { "type": "integer", "minimum": 1 }
"bucketSize": { "type": "integer", "minimum": 1 },
"pubDecrypt": { "type": "integer", "minimum": 1 }
}
},
"migration": {
Expand Down
3 changes: 3 additions & 0 deletions charts/kms-core/values.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,9 @@ kmsCore:
backupVaultKeychainAWSKMSRootKeyID: KMS_CORE__BACKUP_VAULT__KEYCHAIN__AWS_KMS__ROOT_KEY_ID
rateLimiter:
bucketSize: 50000
# Tokens each public decryption holds while in flight, so at most bucketSize / pubDecrypt
# run at once. Each one keeps its ciphertext and MPC session in enclave memory until done.
pubDecrypt: 1
# One-time migration input associating existing epochs with the context they belong to.
# Consumed once by the epoch-data migration on startup; idempotent, so it can be left in place.
# Each epoch belongs to exactly one context, so a given epoch ID must appear under a single entry.
Expand Down
5 changes: 3 additions & 2 deletions ci/perf-testing/PERF_TEST_README.md
Original file line number Diff line number Diff line change
Expand Up @@ -75,9 +75,10 @@ The example shows one sync rung; the configuration file contains both complete s
The ladders run sequentially to avoid competing workloads. Each ladder starts independently,
even when a preceding ladder exceeds its limits. A failed rung skips only higher rates in its own ladder.

Sync rates are PDEC 1,200/1,400/1,700 and UDEC 1,500/1,900/2,300 requests/s.
Sync rates are PDEC 1,200/1,400/1,700/2,000/2,400/2,800 and UDEC 1,500/1,900/2,300/2,700/3,100/3,500 requests/s.
The first two rungs of each sync ladder use the default limits and fail the run on a breach.
The third rung has relaxed limits and `allowfail = true`, so a breach produces a warning.
The remaining four rungs have relaxed limits and `allowfail = true`, so a breach produces a warning
and skips the higher rungs of that ladder.

To preview the fully-expanded workflow locally:

Expand Down
54 changes: 33 additions & 21 deletions ci/perf-testing/perf-scenarios.toml
Original file line number Diff line number Diff line change
Expand Up @@ -21,38 +21,50 @@ maxshed = 0 # max shed (rate-limited) requests, % of offered
pct = 98 # min achieved/target rate, %
allowfail = false # false → a limit breach fails the run; true → warns only

[scenarios.pdec-async]
key = "udec-key-gen"
after = ["crs-gen"]
rates = [
{ rate = 1100 },
{ rate = 1300, maxfail = 1, maxshed = 1, pct = 90 },
{ rate = 1500, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
]
# EXPERIMENT: the async ladders are commented out; only pdec-sync followed by udec-sync runs,
# to test whether pdec-sync leaves state behind that slows udec-sync. Restore before merging.
# [scenarios.pdec-async]
# key = "udec-key-gen"
# after = ["crs-gen"]
# rates = [
# { rate = 1100 },
# { rate = 1300, maxfail = 1, maxshed = 1, pct = 90 },
# { rate = 1500, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
# ]

[scenarios.udec-async]
key = "udec-key-gen" # task whose request-id becomes each rate's key_id
after = ["pdec-async"] # extra deps for the first rate (run after pdec)
rates = [
{ rate = 2400 },
{ rate = 2800, maxfail = 1, maxshed = 1, pct = 95 },
{ rate = 3200, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
]
# [scenarios.udec-async]
# key = "udec-key-gen" # task whose request-id becomes each rate's key_id
# after = ["pdec-async"] # extra deps for the first rate (run after pdec)
# rates = [
# { rate = 2400 },
# { rate = 2800, maxfail = 1, maxshed = 1, pct = 95 },
# { rate = 3200, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
# ]

[scenarios.pdec-sync]
key = "udec-key-gen"
after = ["udec-async"]
after = ["crs-gen"] # EXPERIMENT: normally after = ["udec-async"]
rates = [
{ rate = 1200 },
{ rate = 1400 },
{ rate = 1200, allowfail = true },
{ rate = 1400, allowfail = true },
{ rate = 1700, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
{ rate = 2000, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
{ rate = 2400, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
{ rate = 2800, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
]

[scenarios.udec-sync]
key = "udec-key-gen"
after = ["pdec-sync"]
rates = [
{ rate = 1500 },
{ rate = 1900 },
{ rate = 1500, allowfail = true },
{ rate = 1900, allowfail = true },
{ rate = 2300, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
{ rate = 2700, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
{ rate = 3100, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
{ rate = 3500, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
{ rate = 4000, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
{ rate = 4500, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
{ rate = 5000, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
{ rate = 5500, maxfail = 10, maxshed = 25, pct = 70, allowfail = true },
]
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,13 @@ kmsPeers:
count: 1

kmsCore:
# EXPERIMENT: bound public decryptions in flight per core to 100000 / 12 = 8333 (about 16-33 GiB
# at the 2-4 MiB each measured under backlog), so an overloaded core sheds load instead of running
# its 96 GiB enclave out of memory. With a bucket twice the preprocessing cost (50000),
# preprocessing no longer needs the whole bucket to itself.
rateLimiter:
bucketSize: 100000
pubDecrypt: 12
thresholdMode:
enabled: true
decCapacity: 30000
Expand Down
103 changes: 103 additions & 0 deletions ci/scripts/sample_core_logs.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
"""Streams kms-core container logs for the whole run.

Collecting logs once at the end loses everything that container log rotation has already
discarded, and everything a restarted container printed before it died. Following each
container from the start keeps both.
"""

import argparse
import gzip
import os
import subprocess
import threading
from pathlib import Path

from perf_common import cli, kube_json, timestamp

REATTACH_DELAY_SECS = 5
POD_DISCOVERY_INTERVAL_SECS = 30


def containers():
# In enclave mode kms-server runs inside the enclave and its output reaches the pod through
# the logger container; the kms-core container runs the enclave and the liveness probe.
if os.environ.get("DEPLOYMENT_TYPE") == "thresholdWithEnclave":
return ["kms-core-enclave-logger", "kms-core"]
return ["kms-core"]


def follow(namespace, pod, container, path, stop, processes):
"""Follows one container, reattaching after restarts from the last timestamp seen."""
since = None
with gzip.open(path, "at") as out:
while not stop.is_set():
origin = f"--since-time={since}" if since else "--tail=-1"
out.write(f"### {timestamp()} following {pod}/{container} ({origin})\n")
process = subprocess.Popen(
["kubectl", "-n", namespace, "logs", "-f", "--timestamps", origin, pod, "-c", container],
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
)
processes.append(process)
for count, line in enumerate(process.stdout, 1):
out.write(line)
stamp = line.split(" ", 1)[0]
if stamp[:4].isdigit():
since = stamp
if count % 1000 == 0:
out.flush()
process.wait()
out.write(f"### {timestamp()} stream ended with status {process.returncode}\n")
out.flush()
stop.wait(REATTACH_DELAY_SECS)


def stream(namespace, output):
output.mkdir(parents=True, exist_ok=True)
stop = threading.Event()
processes, threads, followed = [], [], set()
try:
while True:
pods = kube_json(namespace, "get", "pods", "-l", "app=kms-core").get("items", [])
for pod in pods:
name = pod.get("metadata", {}).get("name")
for container in containers():
if not name or (name, container) in followed:
continue
followed.add((name, container))
thread = threading.Thread(
target=follow,
args=(
namespace,
name,
container,
output / f"{name}-{container}.log.gz",
stop,
processes,
),
daemon=True,
)
thread.start()
threads.append(thread)
if stop.wait(POD_DISCOVERY_INTERVAL_SECS):
break
finally:
stop.set()
for process in processes:
if process.poll() is None:
process.terminate()
for thread in threads:
thread.join(timeout=30)


def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("namespace", nargs="?", default="kms-ci")
parser.add_argument("output", nargs="?", type=Path, default=Path("perf-diagnostics/core-logs"))
args = parser.parse_args()
stream(args.namespace, args.output)


if __name__ == "__main__":
cli(main)
92 changes: 90 additions & 2 deletions ci/scripts/sample_perf_diagnostics.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
"""Coordinates CPU, application metrics, ENA lifecycle, and pod placement samples."""
"""Coordinates CPU, application metrics, ENA lifecycle, pod placement, and kms-core log and
restart samples."""

import argparse
import time
Expand Down Expand Up @@ -59,6 +60,73 @@ def lifecycle_rows(daemonset, pods, stamp):
return rows


def core_lifecycle_rows(pods, stamp):
"""One row per kms-core container: readiness, restarts, and why it last terminated."""
rows = []
for pod in pods.get("items", []):
name = pod.get("metadata", {}).get("name")
for container in pod.get("status", {}).get("containerStatuses") or []:
state = container.get("state") or {}
last = (container.get("lastState") or {}).get("terminated") or {}
rows.append(
[
stamp,
name,
container.get("name"),
container.get("ready", False),
container.get("restartCount", 0),
state.get("running", {}).get("startedAt", ""),
state.get("waiting", {}).get("reason", ""),
last.get("reason", ""),
last.get("exitCode", ""),
last.get("startedAt", ""),
last.get("finishedAt", ""),
]
)
return rows


def finish_core(namespace, output):
"""Records why kms-core containers restarted, while the pods still exist."""
output.mkdir(parents=True, exist_ok=True)

def capture(args, path, timeout=120):
result = best_effort(
["kubectl", f"--request-timeout={timeout}s", "-n", namespace, *args],
timeout=timeout + 5,
)
with path.open("w") as stream:
stream.write(result.stdout + result.stderr)

capture(["describe", "pods", "-l", "app=kms-core"], output / "describe-pods.txt")
capture(["get", "events", "--sort-by=.lastTimestamp", "-o", "wide"], output / "events.txt")
for pod in kube_json(namespace, "get", "pods", "-l", "app=kms-core").get("items", []):
name = pod.get("metadata", {}).get("name")
for container in pod.get("status", {}).get("containerStatuses") or []:
if container.get("restartCount", 0):
capture(
["logs", "--previous", "--timestamps", name, "-c", container.get("name")],
output / f"{name}-{container.get('name')}-previous.log",
)
if any(
c.get("name") == "kms-core-enclave-logger"
for c in pod.get("status", {}).get("containerStatuses") or []
):
capture(
[
"exec",
name,
"-c",
"kms-core",
"--",
"sh",
"-c",
"nitro-cli describe-enclaves; cat /var/log/nitro_enclaves/*.log",
],
output / f"{name}-nitro.log",
)


def finish(namespace, output):
def capture(args, stream):
result = best_effort(
Expand Down Expand Up @@ -137,14 +205,33 @@ def sample(namespace, output):
("sample_pod_placement.py", "pod-placement.tsv"),
]:
children.append(start_sampler(script, [namespace], output / name))
with (output / "ena-lifecycle.log").open("w") as stream:
children.append(
start_sampler(
"sample_core_logs.py",
[namespace, str(output / "core-logs")],
output / "core-logs-sampler.log",
)
)
with (
(output / "ena-lifecycle.log").open("w") as stream,
(output / "core-lifecycle.tsv").open("w") as core_lifecycle,
):
stream.write("# daemonset: timestamp kind created desired current ready unavailable\n")
stream.write(
"# pod: timestamp kind pod created started node phase ready restarts "
"container_started waiting_reason terminated_reason\n"
)
core_lifecycle.write(
"# timestamp pod container ready restarts started waiting_reason "
"last_terminated_reason last_exit_code last_started last_finished\n"
)
while True:
stamp = timestamp()
for row in core_lifecycle_rows(
kube_json(namespace, "get", "pods", "-l", "app=kms-core"), stamp
):
core_lifecycle.write(tsv(row) + "\n")
core_lifecycle.flush()
rows = lifecycle_rows(
kube_json(namespace, "get", "daemonset/ena-probe"),
kube_json(namespace, "get", "pods", "-l", "app=ena-probe"),
Expand All @@ -156,6 +243,7 @@ def sample(namespace, output):
time.sleep(10)
finally:
stop_samplers(children)
finish_core(namespace, output / "core-restarts")
finish(namespace, output)


Expand Down
Loading
Loading