Distributed Systems Security and Trust Questions
Security problems that exist because a system is distributed: keeping trust state correct while it propagates across many services, clusters and regions. Covers credential, token and certificate revocation under eventual consistency and network partitions; fleet-wide rotation of signing keys, secrets and trust anchors without outages (key rollover and grace windows, canary rotation, recovery from a compromised root); distributed authorization (replicated policy decision points, cached decisions, fail-open versus fail-closed when an auth dependency degrades); propagating caller identity and permissions through service call chains; tamper-evident audit trails across services and regions (hash chains, Merkle proofs, ordering events with imperfect clocks); Byzantine and partially trusted participants; cross-cluster and cross-organization trust federation; securing shared distributed components such as caches and message brokers against injection, replay and cross-tenant access; protecting data in transit across region boundaries; and tenant isolation as a security blast-radius boundary. Steady-state mTLS, service-mesh identity and network segmentation mechanics are covered by zero-trust service-to-service security; single-system cryptography and KMS basics by applied cryptography.
Design a secure service-to-service authentication and authorization system across multiple clusters and regions. Cover token issuance, short-lived credentials, certificate/key rotation, trust boundaries, least-privilege authorization, fallback when an identity provider is down, and operational practices for key management.
Sample Answer
Direct answer
Give every workload a cryptographic identity issued by its own region's certificate authority after the platform attests what the workload is, and make every credential short-lived so rotation is continuous rather than an event. Authentication happens with mutual TLS (both sides present certificates) using those identities; authorization is a default-deny policy keyed on the caller's identity, evaluated locally. The design principle that answers "what if the identity provider is down" is: verification never needs the issuer online, and issuance is regional with enough credential lifetime left to ride out an outage.
Terms
- Workload identity: a name for a running service, for example the SPIFFE ID
spiffe://prod.eu.example/ns/payments/sa/api(SPIFFE, the Secure Production Identity Framework for Everyone, is an open standard for naming workloads and issuing them certificates, called SVIDs, short for SPIFFE Verifiable Identity Document: the actual certificate or token a workload presents). - Attestation: proving what a workload is from facts the platform controls (its Kubernetes service account, node, image) rather than from something the workload says.
- Certificate chain: a certificate is signed by the authority above it (its issuer); that issuer's own certificate is signed by the authority above it, and so on up to a root. A workload's certificate, at the bottom of that chain, is called a leaf certificate. To trust a leaf, a verifier walks the chain up to a root it already has in its trust bundle; if it has never been told about one of the intermediates in between, it cannot validate anything that intermediate signed, no matter how correctly the leaf itself was issued. That is what "chained to a CA" means, and it is why the rotation order below matters.
- Trust domain / trust boundary: the set of workloads that share a root of trust; crossing a boundary requires explicit federation: the two sides must each deliberately configure trust for specific identities in the other domain, rather than automatically trusting everything the other domain issues.
- Trust bundle: the set of CA (certificate authority) certificates a verifier accepts.
Architecture
flowchart TB
ROOT[Offline root CA] --> ICA1[Intermediate CA: region EU]
ROOT --> ICA2[Intermediate CA: region US]
ICA1 --> AG1[Node agents EU clusters]
ICA2 --> AG2[Node agents US clusters]
AG1 --> W1[Workload certs, 24h]
AG2 --> W2[Workload certs, 24h]
PAP[Policy repo] --> PE1[Local policy check EU]
PAP --> PE2[Local policy check US]
Token and credential issuance
- A node agent (a small process the platform runs on every host, distinct from the application itself) attests the workload using facts it can verify locally: its Kubernetes service account, its namespace (a Kubernetes grouping boundary that scopes and isolates related workloads), and its image digest (a cryptographic hash identifying the exact container image build, so a tampered or swapped image gets a different identity). It then requests a certificate from the regional intermediate CA.
- The workload receives an X.509 certificate (X.509 is the standard format essentially every TLS certificate uses) valid for 24 hours, with its SPIFFE ID in the subject alternative name, a field in the certificate that states which identity or identities it is valid for.
- For calls that also need application-level claims (user on whose behalf, audience), a regional token service issues JWTs (JSON Web Tokens: compact, signed tokens that carry a set of claims, such as who the caller is and who the token is for) valid for 5 minutes with an audience claim, signed with a key whose public half is distributed alongside the trust bundle.
Short-lived credentials and outage tolerance, with numbers
Agents renew a certificate when 50% of its lifetime remains. The worst case for an issuer outage is that it starts at the moment renewal was due, so the tolerance is the remaining lifetime at the renewal point:
| Certificate lifetime | Renew at | Issuer outage the fleet survives | Useful life of a stolen cert |
|---|---|---|---|
| 1 hour | 30 min left | 30 minutes | up to 1 hour |
| 24 hours | 12 h left | 12 hours | up to 24 hours |
| 7 days | 3.5 days left | 3.5 days | up to 7 days |
Commit to 24 hours for certificates. Twelve hours of outage tolerance is long enough to fix a broken CA across a night, and a day is a tolerable exposure window when combined with the other controls. Load is trivial: 50,000 workloads renewing every 12 hours is 50,000 / 43,200 s, about 1.2 issuances per second. The JWT layer stays at 5 minutes because it carries user context, where faster expiry matters more.
Certificate and key rotation
- Leaf certificates rotate continuously by construction.
- Intermediate CAs rotate every few months by overlap: add the new intermediate to trust bundles first, wait until every verifier has the updated bundle, then start issuing from it, and remove the old one only after the last certificate it signed has expired (24 hours plus margin). Reversing that order causes an outage because verifiers reject certificates chained to a CA they have not been told about.
- The root is kept offline in an HSM (hardware security module) and used only to sign intermediates. Root rotation is the same add-then-switch-then-remove sequence stretched over weeks.
- JWT signing keys rotate the same way, publishing the new public key before signing with it.
Trust boundaries
- One trust domain per region (or per environment), not one global domain. A compromised EU intermediate lets an attacker mint EU identities only.
- Cross-region calls use federation: US verifiers are explicitly configured to accept the EU bundle, for specific identities. The US payments service accepts
spiffe://prod.eu.example/ns/checkout/sa/apiand nothing else from EU. - Staging and production never share a root. This is the most common real-world boundary failure.
Least-privilege authorization
Policies are written against identities, not network addresses: "payments accepts calls from checkout on POST /charges; everything else is denied". Policies live in version control, are compiled and pushed to a local enforcement point in every workload (a sidecar, a helper process deployed alongside the application container, or a library linked into the application itself), and are evaluated without a network call. New services start with zero inbound permissions and receive them through review. Every allow and deny is logged with caller identity and policy version.
Fallback when the identity provider is down
Walk through what depends on the issuer:
| Function | Depends on issuer? | Behaviour during outage |
|---|---|---|
| Verifying an existing certificate or token | No: needs only the cached trust bundle and public keys | Unaffected |
| Renewing a workload certificate | Yes | Existing certs keep working for up to 12 hours; alert immediately |
| New pods starting (a pod is the smallest deployable unit in Kubernetes, one or more containers scheduled together on a host) | Yes | Scale-ups and deploys stall; freeze deploys in that region |
| Minting 5-minute JWTs | Yes, regional token service | Fail over to a second token-service instance in the same region; cross-region failover only if the other region's issuer is in the verifiers' federated bundle |
What we deliberately do not do: fall back to a long-lived static shared secret, or disable mTLS. Both turn an availability incident into a security incident. If the outage outlasts the 12-hour tolerance, the controlled response is an emergency extension: issue longer-lived certificates from a standby intermediate held for this purpose, logged and time-boxed.
Operational practices for key management
- Private keys for leaf certs are generated on the host and never leave it; only CSRs (certificate signing requests) travel.
- CA keys live in HSMs or a cloud KMS (key management service); root operations require two people (dual control) and are recorded.
- Monitor: certificate time-to-expiry across the fleet (alert when any workload is under 25% of lifetime), issuance error rates, and bundle version per verifier.
- Rehearse a compromised intermediate: revoke it by removing it from bundles, re-issue the region's workloads from a fresh intermediate, and measure how long that takes. That measured time is your real recovery objective.
Trade-offs and pitfalls
- Short lifetimes replace revocation lists. Certificate revocation lists and OCSP (Online Certificate Status Protocol) checks are rarely reliable at mesh scale (the scale of a service mesh, the network of many services all calling each other over these short-lived certificates); expiry is. The cost is issuer availability, which the table above quantifies.
- Regional isolation costs federation configuration. It is worth it: it bounds a CA compromise to one region.
- Clock skew breaks short-lived credentials. Allow a small skew tolerance (a minute or two) and monitor time sync.
- Common wrong turn: authorizing by IP or namespace label. Those change with every reschedule; identities do not.
Your platform needs a secrets and API key management system for services sitting behind an API gateway. It must support zero-downtime key rotation, immediate revocation, audit logging, and minimal blast radius if a key leaks. Explain the storage model, how keys get distributed to gateways and services, how you would orchestrate rotation, and the emergency revocation flow.
Sample Answer
Direct answer
Split the problem into two kinds of secret. Client API keys (presented by callers to the gateway) are stored only as hashes in a central key registry and pushed to gateways as a versioned, in-memory snapshot plus a change stream (a continuous feed of create, scope-change and revoke events, so a gateway only has to apply what changed rather than re-fetch everything), so a gateway can validate a key locally and a revocation reaches every gateway in seconds. Service secrets (database passwords, third-party credentials) live in a secrets manager, encrypted under keys held in a KMS (key management service) or HSM (hardware security module), and are fetched at runtime by workloads that authenticate with their platform identity, never with a secret baked into an image. Rotation always runs through a window where old and new are both valid; revocation skips that window and is pushed, acknowledged and audited.
Terms
- Blast radius: how much one leaked key can do. Shrink it with per-client, per-environment, narrowly scoped, short-lived credentials.
- Envelope encryption: each secret is encrypted with its own data key, and that data key is encrypted with a master key that never leaves the KMS/HSM.
- Secret zero: the first credential a brand-new workload needs in order to fetch all its other secrets. How you bootstrap it decides whether the whole design is sound.
Storage model
flowchart LR
ADM[Key admin API] --> REG[(Key registry: hashes + metadata)]
REG --> STR[Change stream]
STR --> G1[Gateway region A]
STR --> G2[Gateway region B]
SM[(Secrets manager)] --> KMS[KMS or HSM master key]
W[Service workload] -->|platform identity| SM
REG --> AUD[(Append-only audit log)]
SM --> AUD
Client API keys. Generate 256 random bits, format as pk_live_<keyid>_<secret>. Store keyid, SHA-256(secret), owner, scopes, environment, status, created and expiry times. Never store the key itself: it is shown once at creation. A fast hash is correct here (unlike passwords): passwords are hashed with deliberately slow algorithms like bcrypt because people choose low-entropy, guessable values, so the hash must be expensive enough to make guessing infeasible. A 256-bit random secret has no such weakness: it cannot be brute-forced regardless of hash speed, so bcrypt's deliberate slowness buys nothing and would add latency on every request. The keyid prefix lets gateways look up the record in one step and lets secret-scanning tools (automated scanners that watch code repositories for accidentally committed credentials) recognise a leaked key in a public repository.
Service secrets. Stored in the secrets manager under envelope encryption, organised by path per service and environment (prod/orders/db), with a policy granting each workload identity exactly its own paths. Prefer dynamic secrets where the backend supports them: the secrets manager creates a unique database user per workload with a lease (a grant that expires automatically unless the workload renews it before the deadline) of, say, 1 hour, so there is nothing long-lived to leak and each credential is attributable to one workload.
Sizing the gateway snapshot. With 2,000,000 active client keys at about 32 bytes of hash plus 100 bytes of metadata each: 2,000,000 x 132 = 264,000,000 bytes, about 252 MiB. That fits in gateway memory, so every validation is a local lookup with no network call. If it did not fit, gateways would cache on demand with a short TTL (time-to-live) and fall back to the registry on a miss.
Distribution to gateways and services
- Gateways load a full snapshot at start-up (tagged with a version number), then apply a change stream of
created / scope_changed / revokedevents in order. Each gateway reports the version it has applied, so the control plane (the system that manages and coordinates gateway configuration, as distinct from the data plane that carries live request traffic) always knows the laggards, the gateways that have fallen behind the latest version. If the stream breaks, the gateway re-fetches a snapshot; if it cannot, it keeps serving from its last snapshot and alerts, because refusing all traffic is worse than a snapshot a few minutes old, except for revocations (see below). - Services fetch their secrets at start-up through a local agent (a sidecar: a helper process deployed alongside the application container that does this fetching and renewal on the app's behalf), which renews leases before expiry and writes secrets to an in-memory file the app re-reads, so rotation does not require a redeploy.
Secret zero: bootstrapping a brand-new workload
The failure pattern is putting a long-lived token in CI (continuous integration) variables or the container image so the new container can log in to the secrets manager. Instead, let the platform vouch for the workload:
- CI/CD pipeline: the CI provider issues a short-lived OIDC (OpenID Connect) token describing the exact repository, branch and job; the secrets manager or the cloud provider's security token service (STS) trusts that issuer and exchanges the token for credentials scoped to that pipeline. No static secret exists in CI at all.
- Ephemeral containers: Kubernetes projects a service-account token that is audience-bound (it can only be redeemed by the exact service it was issued to, not replayed elsewhere) and expires; the agent presents it, the secrets manager validates it against the cluster and returns a short-lived token for that service's paths. On a cloud VM the equivalent is the instance identity document: a document signed by the cloud provider itself, attesting which specific VM instance is making the request.
- When no platform identity exists, use a single-use, short-lived wrapped token delivered by the deploy system (Vault calls this response wrapping: the secrets manager seals the real credential in an envelope that can be opened, or "unwrapped," exactly once): if an attacker used it first, the legitimate workload's unwrap fails and that failure is itself the alarm.
Zero-downtime rotation orchestration
For client API keys, rotation is the client's action, so the platform must allow two active keys per client:
- Client (or automation) creates key B; A and B are both valid.
- Client deploys B; the gateway emits
last_usedper key ID, so both sides can see traffic migrate from A to B. - When A has had zero traffic for a set period (for example 7 days), A is revoked. Keys also carry an expiry (for example 1 year), with warnings at 30 and 7 days.
For service secrets:
- Create the new credential in the backend (new DB password or user) while the old one stays valid.
- Publish it to the secrets manager as a new version; agents pick it up on their next refresh.
- Wait until every consumer reports the new version (or until one full refresh interval plus margin has passed).
- Disable the old credential, then delete it after a cooldown. If error rates rise after step 4, re-enabling the old credential is the rollback, which is why you disable before deleting.
Emergency revocation flow
Leak detected (secret scanner hit, anomaly alert or customer report):
- Revoke in the registry: status
revoked, with reason and actor recorded in the audit log. - Push: the revocation event goes on the change stream with priority; gateways apply it and acknowledge by version. Target: 99% of gateways within 5 seconds.
- Bound the tail: a gateway that has not acknowledged within the deadline is marked unhealthy and removed from the load balancer until it catches up. This turns "eventually revoked" into "revoked, or not serving traffic".
- Scope the damage: query the audit log for every request made with that key ID since the suspected leak time.
- Replace: issue a new key to the owner; for a service secret, rotate immediately without the normal overlap window, accepting a brief error spike.
Audit logging
Every key creation, scope change, revocation, secret read and failed authentication is written to an append-only store in a separate account that application operators cannot delete from. Log key IDs, never key material. Alert on reads from an unexpected identity, a spike in failed validations for one key ID (someone guessing), and use of a key from a new network or country.
Trade-offs and pitfalls
- Local validation vs instant revocation: local snapshots make validation fast and resilient but make revocation a distributed-propagation problem. The acknowledgement-plus-eviction step is what closes that gap; without it you only have "probably revoked".
- Rotation that requires a restart is rotation nobody runs. Hot reload of secrets is a prerequisite.
- Scoping is the cheapest blast-radius control: a key that can only read
ordersinstagingis a far smaller incident than a master key. - Common wrong turn: encrypting API keys reversibly "so support can look them up". If the platform can decrypt them, so can whoever compromises the platform.
Why are tamper-evident and append-only audit logs important for a distributed system? As an SRE, describe how you would design audit logging across services and regions so logs can't be silently altered or deleted, how you'd keep events correctly ordered per entity, and what operational controls you would put in place to preserve logs during an incident.
Sample Answer
Direct answer
Audit logs are the record you use to answer "who did what, when" after something goes wrong. That makes them the first thing an attacker (or a panicking insider) wants to edit. Append-only means records can be added but never changed; tamper-evident means that if someone does change or delete a record, verification detects it. Across services and regions, the design is: each service writes events with a per-entity sequence number to a replicated log (an ordered sequence of events copied to multiple machines, so no single machine failure loses it and no single operator holds the only copy), events are chained by hashes, the chain head is regularly anchored somewhere the writers cannot modify, and the storage itself enforces retention that no operator can override.
Why it matters in a distributed system
- Forensics: after a breach you must reconstruct what happened across many services. If the log can be silently edited, you cannot trust any conclusion.
- Compliance: standards such as PCI DSS (Payment Card Industry Data Security Standard) and SOC 2 (a widely used security-controls audit report) require audit trails protected from modification.
- Distributed makes it harder: many writers, many regions, clocks that disagree, and many operators with access to some part of the pipeline.
Design
1. Making changes detectable: hash chains
Each record stores the hash of the previous record, so changing any record changes every hash after it. This short program shows what a chain catches and what it does not:
import hashlib, json
GENESIS = "0" * 64 # placeholder "previous hash" used only for the first entry, which has no real predecessor
def entry_hash(prev_hash, event):
body = json.dumps(event, sort_keys=True, separators=(",", ":"))
return hashlib.sha256((prev_hash + body).encode()).hexdigest()
def append(log, event):
prev = log[-1]["hash"] if log else GENESIS
log.append({"event": event, "prev": prev, "hash": entry_hash(prev, event)})
def verify(log):
prev = GENESIS
for i, rec in enumerate(log):
if rec["prev"] != prev or rec["hash"] != entry_hash(prev, rec["event"]):
return f"broken at entry {i}"
prev = rec["hash"]
return "ok"
log = []
append(log, {"seq": 1, "actor": "alice", "action": "role.grant", "target": "bob:admin"})
append(log, {"seq": 2, "actor": "bob", "action": "export", "target": "customers.csv"})
append(log, {"seq": 3, "actor": "alice", "action": "role.revoke", "target": "bob:admin"})
head = log[-1]["hash"]
print("clean log:", verify(log), "| head", head[:12])
# 1) edit an event in place
log[1]["event"]["target"] = "nothing.csv"
print("after edit:", verify(log))
log[1]["event"]["target"] = "customers.csv"
# 2) delete the incriminating entry and re-link the next one
forged = [log[0], dict(log[2], prev=log[0]["hash"])]
print("after delete + relink:", verify(forged))
# 3) delete it and recompute every later hash: the chain verifies again...
rebuilt = []
for rec in (log[0], log[2]):
append(rebuilt, rec["event"])
print("after full rewrite:", verify(rebuilt), "| head", rebuilt[-1]["hash"][:12])
# ...which is why the head hash must be anchored somewhere the writer cannot change
print("head matches anchored copy:", rebuilt[-1]["hash"] == head)
clean log: ok | head 1537d687119a
after edit: broken at entry 1
after delete + relink: broken at entry 1
after full rewrite: ok | head 8cf5f38f565b
head matches anchored copy: False
The third case is the important lesson: a hash chain alone does not stop someone with write access from rebuilding it. It only becomes tamper-evident when the head hash is anchored: copied every few minutes to a place the log writers cannot change (a separate account's locked storage, another organization, or a public transparency log, a shared append-only log, such as the one behind Certificate Transparency, that many independent parties can audit). This is a scale optimisation, not a change to the core idea above: instead of anchoring every single event's hash (expensive at high volume), events are batched into a Merkle tree (a tree of hashes whose single root commits to every event in the batch) and only that one root is anchored. A Merkle proof then shows a specific event was in a batch using about log2(n) hashes, which is 20 hashes for a batch of one million events, instead of handing over all one million. If your system's volume is low enough to anchor every event's hash directly, you do not need this step at all.
2. Making changes impossible at the storage layer
- Write to storage that enforces immutability, for example Amazon S3 Object Lock in compliance mode, where no user, including the account root, can delete or overwrite an object before its retention date.
- Put the audit store in a separate account owned by the security team, so compromising a production account does not grant delete rights.
- Replicate to a second region so losing one region does not lose the trail.
3. Keeping events ordered per entity
Wall clocks across servers drift by milliseconds or more, so timestamps alone cannot order two events that happened close together. Instead:
- Route all events for one entity (a user, an account, a resource) through one partition of the log. A partition is one independently ordered slice of the log; the log is split into many partitions (think separate lanes) so it can scale, and every event for a given entity ID is routed, by that ID, to the same lane every time. Within a partition, order is total: every event in that one lane has an unambiguous before-and-after position. Order is not directly comparable across two different partitions.
- The writer for that entity assigns a monotonic sequence number (1, 2, 3...). A gap reveals a missing event; a duplicate reveals a replay.
- Keep the wall-clock timestamp for humans, but use the sequence number for ordering. For causality across entities, a hybrid logical clock (a timestamp that never goes backwards and carries causal order: if event X caused event Y, Y's clock value is guaranteed to be greater than X's, so a reader can tell cause from effect even when the events sit in different partitions) is the usual addition.
In the example, seq 1, 2, 3 show that Bob's export happened while he had admin, even if the servers' clocks disagreed.
4. Operational controls during an incident
- Separation of duties: the people responding to an incident cannot delete or modify audit data; audit-store admin is a different role with its own approval.
- Legal hold: a flag placed on specific records that overrides the normal retention schedule, so declaring an incident and applying a hold guarantees those logs cannot be deleted, even by an expiry job, until the hold is explicitly lifted.
- Protect the pipeline, not just the store: agents buffer locally if the pipeline is down, and a gap in a service's sequence numbers raises an alert. An attacker's first move is often to stop logging, so "a service stopped sending audit events" is itself a high-priority alert.
- Continuous verification: a job re-verifies chains against anchored heads and pages on any mismatch.
Trade-offs and pitfalls
- Immutability versus data-protection deletion rights: you cannot delete a person's data from a compliance-locked log. Log IDs rather than personal data, or encrypt personal fields with a per-person key you can destroy.
- Per-entity ordering is cheap; global ordering is expensive. A single global sequence becomes a throughput bottleneck. Order per entity and accept that cross-entity order is approximate.
- Common wrong turn: relying on "only admins can delete". The threat model includes admins.
You operate a multi-tenant control plane where API keys and control APIs are currently shared across all tenants. Design an architecture that minimizes the blast radius if a single tenant is compromised, and explain how you would migrate off the shared-key model and what that isolation costs you versus what it buys you.
Sample Answer
Direct answer
With one shared API key, the blast radius (everything a single compromised credential can reach) of any tenant compromise is every tenant, because the control plane cannot tell tenants apart. The fix is to make the tenant a property of the credential, not of the request: issue per-tenant, short-lived, narrowly scoped credentials; derive the tenant ID from the authenticated identity and enforce it on every control-plane call; encrypt each tenant's data under its own key; and group tenants into cells (independent copies of the control plane, each serving a subset of tenants) so that even an infrastructure-level compromise reaches one cell, not everyone. I would migrate by running per-tenant credentials alongside the shared key, attributing and shadow-enforcing first (running the new tenant check and logging what it would have blocked, without actually blocking anything yet), moving tenants in waves, and revoking the shared key last. The cost is more keys, more deployment units and more operational work; what it buys is turning "one leak exposes all tenants" into "one leak exposes one tenant".
Why the shared model is so dangerous
Today, a request says "act on tenant 42" and carries the shared key. The control plane checks the key (valid) and trusts the tenant ID in the request. So:
- Anyone holding the key can act as any tenant by changing one parameter.
- A compromise of tenant 42's automation (a leaked CI (continuous integration) pipeline secret, a malicious insider at that customer) is a compromise of all tenants.
- You cannot revoke access for one tenant without breaking every tenant, so in an incident you are forced to choose between leaving the attacker in and a platform-wide outage.
- Audit logs cannot say which tenant actually made a call, which makes incident scoping guesswork.
Target architecture
flowchart LR
T[Tenant workload] -->|per-tenant short-lived token| GW[API gateway]
GW -->|tenant ID from token| CP[Control plane in cell 3]
CP --> AZ{Tenant check<br/>resource.tenant == token.tenant}
AZ -->|match| DS[(Tenant data<br/>encrypted with tenant key)]
AZ -->|mismatch| DENY[Deny and alert]
CP --> KMS[Key service<br/>per-tenant keys]
Layer 1: identity per tenant
- Each tenant gets its own credentials, ideally exchanged for short-lived tokens (for example 15-minute tokens from a token service) so a stolen token expires quickly. Long-lived API keys, where unavoidable, are per tenant, hashed at rest, scoped to specific actions, and rotatable without touching other tenants.
- Scopes per credential: a CI deploy key can deploy, but cannot read secrets or manage users.
Layer 2: authorization bound to the tenant
- The control plane takes the tenant ID from the verified token only. Any tenant ID in a URL or body must match it, or the call is denied and logged as a cross-tenant attempt.
- Every data-access path filters by tenant at the lowest layer possible (row-level security in the database, meaning the database itself filters rows by tenant, or per-tenant schemas) so an application bug that forgets the filter still cannot return another tenant's rows.
Layer 3: per-tenant encryption keys
- Envelope encryption: each tenant's data is encrypted with data keys that are themselves encrypted by a tenant-specific master key in the key service. A policy on each master key allows its use only by that tenant's workloads.
- A bug that returns the wrong ciphertext, or a stolen backup, does not yield another tenant's plaintext, and crypto-shredding (deleting the tenant's master key) is a clean offboarding.
Layer 4: cells and namespace isolation
- Cell-based architecture: split the control plane into N independent cells, each with its own database, queues and credentials to the underlying infrastructure. A small routing layer maps tenant to cell. A compromise of a cell's service account (a non-human identity that infrastructure code authenticates as, distinct from a human user's login), or a bad deploy, is limited to that cell's tenants.
- Tenant workloads run in their own namespaces (Kubernetes' isolated groupings of resources within one cluster; or full separate accounts, for high-value tenants) with network policies (Kubernetes rules that control which pods may talk to which, enforced by the cluster network itself) that deny cross-namespace traffic by default.
- This is the classic tenancy spectrum: pool (all tenants share infrastructure, separated by identity and data filters), silo (each tenant gets dedicated infrastructure) and bridge (a mix, for example pooled compute with siloed data). I would use pool within a cell for most tenants and silo for the few tenants whose contracts or regulators require it.
Layer 5: detection and containment
- Alert on cross-tenant denials (any occurrence is suspicious), unusual call volumes per tenant, and tokens used from new networks.
- Automated quarantine: a single action revokes all of one tenant's tokens and keys and freezes its control-plane writes, without affecting any other tenant. This is the capability the shared key made impossible.
Migration plan
- Inventory and attribute. Log every call made with the shared key with enough context (source network, user agent (the identifying string a client sends describing what software made the request), workload identity, the tenant ID it claimed) to know who uses it and for what. Typical finding: a few internal tools legitimately act across tenants; these become separately scoped operator identities with their own audit.
- Issue per-tenant credentials in parallel. Both the shared key and the new credentials work. SDKs (the client libraries customers install to call your API) and docs switch to the new model.
- Shadow enforcement. For calls with new credentials, evaluate the tenant check and log "would deny" without blocking. Fix false positives (usually legitimate cross-tenant operator workflows).
- Enforce per tenant, in waves. Once a tenant is fully on its own credentials, turn on enforcement for that tenant and disable the shared key for that tenant ID. Start with internal and low-risk tenants, then the rest.
- Deadline and cutoff. Publish a date; the remaining users get direct help. On the date, stop accepting the shared key, then rotate and destroy it (assume it has already leaked; it has been in many hands).
- Introduce cells afterwards as a separate project. Moving tenants between cells is a data migration; doing it together with the credential change doubles the risk in each step.
What it costs versus what it buys
Worked example with 500 tenants split into 10 cells of 50:
| Question | Shared key | Per-tenant credentials | Per-tenant credentials + 10 cells |
|---|---|---|---|
| Tenants exposed by one tenant's leaked credential | 500 | 1 | 1 |
| Tenants exposed by one compromised control-plane service account | 500 | 500 | 50 |
| Tenants affected by a bad deploy | 500 | 500 | 50 (if deploys roll cell by cell) |
| Credentials to manage | 1 | 500+ | 500+ |
| Control-plane deployments to operate | 1 | 1 | 10 |
Costs: credential lifecycle for hundreds of tenants (issuance, rotation, revocation, support tickets when a customer loses a key); per-tenant keys in the key service (per-key and per-request charges scale with tenant count); a fixed baseline cost for every cell (each cell's database and queues cost money even when mostly idle); deploy pipelines that must roll across cells; cross-cell features (global search, billing) get harder. Benefits: a single customer's compromise is contained to that customer; containment is a one-tenant action instead of a platform-wide outage; audit logs can scope an incident precisely; enterprise customers' security questionnaires get concrete answers.
My recommendation is per-tenant identity, tenant-bound authorization and per-tenant keys for everyone now, since they are cheap relative to the risk they remove, and cells once the tenant count or the value of the largest tenants makes a platform-wide incident unaffordable. Full silo per tenant only where a contract or regulator requires it, because at 500 tenants the fixed per-silo cost multiplies by 500.
Trade-offs and pitfalls
- Per-tenant keys with a shared admin path (one super-credential used by support tooling) recreates the original problem. Operator access must be scoped, time-bound and approved.
- Trusting the tenant ID in the request body after the migration is the most common way this design quietly fails.
- Cells without cell-aware deploys give no deploy-safety benefit; roll changes one cell at a time.
- Noisy neighbours (tenants that consume so much shared capacity, CPU, bandwidth, connections, that they degrade service for others on the same infrastructure) are an isolation problem too: without per-tenant rate limits, one compromised tenant can still degrade others even if it cannot read their data.
Design a distributed policy decision point (PDP) for authorization that must serve 1,000 decisions/sec with a 50ms latency SLA across three regions, while preventing privilege escalation during partitions. Describe how to store and replicate policies, your caching strategy at local PDPs, decision versioning, how to roll back a bad policy quickly, and the operational monitoring and audit logging compliance requires.
Sample Answer
Direct answer
At 1,000 decisions per second the hard part is not throughput, it is correctness across regions under a 50 ms budget. A cross-region round trip alone can use most of that budget, so every region must evaluate decisions locally from an in-memory copy of policy. I would keep policy as code in a single authoritative store, compile it into signed, immutable, versioned bundles, and distribute them to stateless (holding no durable state of its own between requests, so any instance can be replaced or added with no handoff) PDP (policy decision point: the component that answers "may subject S do action A on resource R?") instances in every region. Permission data (group memberships, relationships) replicates separately with a version token. Grants require a quorum (a majority, 2 of the 3 regions), revocations propagate on a priority path and "deny wins", and a region that loses contact with the quorum keeps serving only within a bounded staleness window before sensitive actions fail closed (are refused rather than risk granting something a fresher policy would deny). Rollback is repointing an "active version" pointer to the previous bundle, which takes one propagation cycle.
Requirements I am designing to
- 1,000 decisions/sec total, say split roughly evenly across three regions (about 333/sec each), with bursts of 3x.
- 50 ms p99 (99th percentile) decision latency, measured at the caller.
- No privilege escalation during partitions: a partitioned region must never grant something the authoritative policy would deny, beyond a small, measured staleness bound for low-risk reads.
- Fast rollback, and a full audit trail for compliance.
Architecture
flowchart TB
G[Policy repo<br/>reviewed and tested] --> B[Bundle builder<br/>compile, sign, version]
B --> S[(Bundle store<br/>immutable versions)]
P[Active-version pointer<br/>quorum across 3 regions] --> S
S --> R1[Region 1 PDPs]
S --> R2[Region 2 PDPs]
S --> R3[Region 3 PDPs]
D[(Permission data<br/>quorum writes)] --> R1
D --> R2
D --> R3
R1 --> L[Decision logs<br/>append-only audit]
R2 --> L
R3 --> L
Policy storage and replication
- Policy as code. Rules live in a version-controlled repository. Every change is reviewed, unit-tested against a suite of expected decisions, and compiled into a bundle.
- Bundles are immutable and signed. Each bundle has a content hash as its version and a signature from the build system. PDPs refuse unsigned or unverifiable bundles, so a compromised distribution channel cannot inject policy. Open Policy Agent (OPA, a widely used open-source policy engine) supports exactly this: signed bundles with a
.signatures.jsonfile and arevisionfield in the manifest. - One active-version pointer, written only through a quorum of the three regions, for example with etcd, a strongly consistent key-value store spanning them that uses Raft, a consensus protocol in which a majority must agree on each write. With 3 regions, a write needs 2, so a single partitioned region cannot change which policy is active.
- Permission data (who is in which group, who owns which document) is also written with quorum and carries a monotonically increasing data version (a number that only ever goes up, so comparing two version numbers always tells you unambiguously which is newer). Callers that just changed a permission can pass the returned version token, and the PDP answers only once it has data at least that fresh. Concretely: revoking Bob's access to a document returns
data_version: 4821; if the caller then shares a new secret with that same document a second later, it attachesdata_version: 4821(or later) to the follow-up read check, so a PDP replica still showingdata_version: 4820blocks and waits for fresher data (or forwards to one that has it) instead of answering from a copy that predates the revocation. This is the "new enemy" problem Google's Zanzibar paper (the 2019 paper describing the global authorization system behind Google Docs, Drive and similar products) describes: without it, a user removed from a document and then a new secret added to it could still read the secret from a stale replica.
Caching at local PDPs
Three layers, each bounded:
- The bundle itself is the cache. Each PDP holds the full compiled policy in memory. Evaluation is local computation with no network call.
- Permission data is held in a regional replica with a known replication lag, read from memory or a local store.
- Decision cache keyed by (subject, action, resource, bundle version, data version) with a short TTL (time-to-live), for example 5 seconds. Including both versions in the key means a new bundle or data change automatically invalidates old entries, with no invalidation protocol needed.
Latency budget (50 ms, illustrative split)
| Step | Budget |
|---|---|
| Caller to regional PDP (same region, one network hop) | 5 ms |
| Load-balancer and queueing headroom | 10 ms |
| Policy evaluation in memory | 5 ms |
| Permission-data lookup (regional replica) | 10 ms |
| Reserve for p99 tail | 20 ms |
These are budget allocations, not measurements. The point they make: there is no room for a cross-region call on the decision path, so any design where a PDP asks another region at decision time fails the SLA (service-level agreement).
Capacity
About 333 decisions/sec per region is small. Run at least 3 PDP instances per region across availability zones (physically separate data centers within the region, each with independent power and networking, so one zone failing does not take the others down with it) so losing a zone and doing a deploy at the same time still leaves capacity, and load-test to find the real per-instance ceiling rather than assuming one. A single modern instance could likely handle 333 decisions/sec alone, so the reason to run three per region is not throughput, it is that losing one zone or rolling a deploy must never drop capacity below what the SLA needs: capacity is set by availability, not by load.
Decision versioning
Every decision response and log record carries: decision_id, bundle_version (content hash), data_version, pdp_id, region, and the rule that matched. This lets you answer, months later, "why was this allowed?" by re-running the exact inputs against the exact bundle, and it lets monitoring detect a PDP answering on an old bundle.
Preventing privilege escalation during partitions
Privilege escalation during a partition happens in two ways, and the design blocks both:
- A stale allow: a revocation happened in region 1, region 3 has not heard. Mitigations: revocations replicate on a priority path separate from full bundles; deny rules override allow rules (so a PDP that has received a revocation applies it immediately); each PDP holds a lease (a time-limited authorization to keep serving) renewed from the quorum every 10 seconds. When the lease lapses, high-risk actions (grants, exports, admin) fail closed at once, and low-risk reads continue for a fixed budget (for example 5 minutes) with every such decision flagged in the log.
- A conflicting grant: the minority region accepts "make Bob admin" while cut off. Blocked because policy and permission writes require the 2-of-3 quorum; the minority side is read-only for security state.
Never fail open on PDP errors. If a PDP cannot evaluate, the answer is deny, and the fail-closed count is an alerted metric.
Rolling back a bad policy quickly
- Before it ships: policy unit tests, then shadow evaluation: the candidate bundle runs beside the active one on live traffic, and any decision that differs is logged. A diff rate above an agreed threshold (say 0.1% of decisions) blocks promotion.
- Staged rollout: activate in one region first, watch deny rate and error rate for 10 minutes, then the others.
- Rollback is writing the previous bundle hash to the active-version pointer (one quorum write). PDPs poll or are notified; with a 10-second poll interval, the fleet is back on the old policy within roughly one interval plus download time. Because bundles are immutable, "the old policy" is byte-for-byte what ran before.
- Break-glass: a documented, audited way to roll back even if the build pipeline is down, using a pre-signed last-known-good bundle.
Monitoring and audit logging
Operational metrics: decision latency p50 (median) and p99 (99th percentile, defined above) per region; bundle version per PDP and version skew across the fleet; lease age; deny rate and allow rate per application (a sudden change after a policy push is the fastest signal of a bad policy); fail-closed count; decisions served on stale data.
Audit log for compliance: one record per decision with timestamp, decision_id, subject, action, resource, result, matched rule, bundle_version, data_version, pdp_id. Sensitive input fields are masked before shipping (OPA's decision logs support masking via a policy). Records go to append-only storage with write-once retention, batched and hash-chained (each record's stored hash also incorporates the previous record's hash, so altering or deleting one entry breaks every hash that comes after it) so deletion or editing is detectable. Policy changes are audited separately: who changed what, who approved it, which tests ran.
Sizing the audit volume, assuming about 600 bytes per record:
1000 sdecisions×600 B×86,400 days=51.84 GB/day 51.84 GB/day×365≈18.9 TB/yearThat is manageable uncompressed and much smaller compressed. If a regulator needs every decision, log them all; sampling allows is a cost option only where policy permits it, and denies should always be kept.
Trade-offs and pitfalls
- Centralised PDP in one region is simpler and always consistent, but its cross-region latency breaks the 50 ms SLA and makes that region a single point of failure for all authorization.
- Long decision-cache TTLs make revocation slow; the cache TTL is part of the revocation SLO and must be sized from it.
- Treating the bundle as the only state forgets the permission data: most real escalations come from stale group membership, not stale rules.
- Skipping shadow evaluation means the first time you see a policy's effect is when users do.
Unlock Full Question Bank
Get access to all 12 Distributed Systems Security and Trust interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.