Skip to content
Merged
8 changes: 8 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,14 @@ API supports the OpenID implicit flow.

If `OIDC_SERVER_URL` and `OIDC_REALM` are not provided then the Django Auth is enabled by default.

#### Capacity limit on heavy calls
Semantic `$match` calls (kNN search and/or the in-request rerank) and `$rerank` calls are heavy: each holds an API worker for seconds. The capacity limit counts how many run at once in Redis lanes (cluster-wide, per API process host, semantic kNN, per plan tier and per user) and sends `X-OCL-Capacity-*` headers on those responses. It has three modes:
- `off`: nothing is counted.
- `shadow` (the default): calls are counted and logged, and none is refused.
- `enforce`: a call that finds a lane full gets a 429 with `Retry-After`, before any work or quota charge. With `enforce_for=aware` (the default) only clients that send `"capacity_aware": "true"` in `X-OCL-Event-Metadata` are refused; the rest stay in shadow mode.

The mode and all the numbers are runtime settings, changed by staff with `GET`/`PATCH`/`PUT /capacity/config/` or `python manage.py capacity show|set|reset|history|status`. Each change is kept in `/capacity/config/history/`, and every API process picks it up within `CAPACITY_CONFIG_CACHE_SECONDS` (10). `CAPACITY_LIMIT_MODE` sets the mode until staff first change it. If Redis can't be reached, calls go ahead uncounted.

### Run Checks
(use the `docker exec` command in a service started with `docker compose up -d`)
1. Pylint (pep8):
Expand Down
Empty file added core/capacity/__init__.py
Empty file.
214 changes: 214 additions & 0 deletions core/capacity/config.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,214 @@
"""
The capacity limit's settings: the mode and every number, changed at runtime by staff through
`/capacity/config/` or `manage.py capacity`, with no deploy or restart. Each change adds a CapacityConfig row
(who, when, the old config and the new one) and writes one log line. API processes re-read the newest row at
most every CAPACITY_CONFIG_CACHE_SECONDS, so a change applies across the cluster within that time.
"""
import copy
import logging
import threading
import time

from django.conf import settings
from django.core.exceptions import ValidationError
from django.db import connection, transaction

from core.capacity.constants import (
MODES, MODE_SHADOW, MODE_ENFORCE, ENFORCE_FOR, ENFORCE_FOR_AWARE, TIER_STAFF, TIER_CORE, TIER_EARLY_ACCESS,
TIER_PREVIEW, CONFIG_LOG_EVENT, SOURCE_API)
from core.capacity.logs import emit

logger = logging.getLogger('oclapi')

# The starting numbers (OpenConceptLab/ocl_online#275). Counts are heavy calls in flight at once.
DEFAULTS = {
'mode': MODE_SHADOW, # the boot default comes from settings.CAPACITY_LIMIT_MODE
'enforce_for': ENFORCE_FOR_AWARE,
'api_heavy': {
'cluster': 4, # across all API tasks: half their workers, so the rest stay free for everything else
'per_task': 2, # on each API task, since the load balancer doesn't see how busy a task is
},
'es2_knn': 3, # semantic $match calls, which run kNN searches in Elasticsearch
'reserve_single_row': 1, # cluster slots only single-row $match calls may use, so interactive calls don't wait
'tiers': {TIER_STAFF: 4, TIER_CORE: 4, TIER_EARLY_ACCESS: 3, TIER_PREVIEW: 2}, # ceilings, not reservations
'per_user': {TIER_STAFF: 4, TIER_CORE: 3, TIER_EARLY_ACCESS: 2, TIER_PREVIEW: 1},
'lease_seconds': 60, # a call's lease expires this long after its last renewal, so a dead worker's frees
'renew_seconds': 20,
'max_hold_seconds': 900, # stop renewing after this, in case a call never finishes
'retry_after': {
'base': 5, # seconds per call ahead in the fullest lane
'max': 30,
'paused': 120, # when a lane's limit is 0
},
}
CHOICES = {'mode': MODES, 'enforce_for': ENFORCE_FOR}
MAX_NUMBER = 100000
ADVISORY_LOCK_ID = 275275 # serializes config changes, so a PATCH never merges onto a stale version

_cache = {'config': None, 'expires_at': 0.0}
_cache_lock = threading.Lock()


def get_defaults():
defaults = copy.deepcopy(DEFAULTS)
mode = getattr(settings, 'CAPACITY_LIMIT_MODE', None)
if mode in MODES:
defaults['mode'] = mode
return defaults


def merge(base, changes, known_only=False):
"""`base` with `changes` merged in, dict by dict. With known_only, keys `base` doesn't have are dropped."""
result = copy.deepcopy(base)
for key, value in (changes or {}).items():
if known_only and key not in base:
continue
if isinstance(value, dict) and isinstance(result.get(key), dict):
result[key] = merge(result[key], value, known_only)
else:
result[key] = copy.deepcopy(value)
return result


def diff(old, new, prefix=''):
"""{dotted.path: [old, new]} for every value that differs."""
changes = {}
for key in sorted(set(old or {}) | set(new or {})):
before, after = (old or {}).get(key), (new or {}).get(key)
path = f'{prefix}{key}'
if isinstance(before, dict) and isinstance(after, dict):
changes.update(diff(before, after, f'{path}.'))
elif before != after:
changes[path] = [before, after]
return changes


def _validate_shape(config, template, prefix, errors):
if not isinstance(config, dict):
errors.append(f'{prefix.rstrip(".") or "config"} must be an object.')
return
for key in config:
if key not in template:
errors.append(f'Unknown setting "{prefix}{key}".')
for key, default in template.items():
path = f'{prefix}{key}'
if key not in config:
errors.append(f'"{path}" is required.')
elif isinstance(default, dict):
_validate_shape(config[key], default, f'{path}.', errors)
elif key in CHOICES:
if config[key] not in CHOICES[key]:
errors.append(f'"{path}" must be one of {", ".join(CHOICES[key])}.')
elif isinstance(config[key], bool) or not isinstance(config[key], int) or not (
0 <= config[key] <= MAX_NUMBER):
errors.append(f'"{path}" must be a whole number from 0 to {MAX_NUMBER}.')


def validate(config):
"""Raise ValidationError unless `config` is a complete, consistent config."""
errors = []
_validate_shape(config, DEFAULTS, '', errors)
if not errors:
if config['reserve_single_row'] > config['api_heavy']['cluster']:
errors.append('"reserve_single_row" can\'t be more than "api_heavy.cluster".')
if config['lease_seconds'] < 15:
errors.append('"lease_seconds" must be at least 15.')
if not 1 <= config['renew_seconds'] <= config['lease_seconds'] // 3:
errors.append('"renew_seconds" must be at least 1 and at most a third of "lease_seconds".')
if config['max_hold_seconds'] < config['lease_seconds']:
errors.append('"max_hold_seconds" can\'t be less than "lease_seconds".')
retry_after = config['retry_after']
if retry_after['base'] < 1 or retry_after['paused'] < 1 or retry_after['max'] < retry_after['base']:
errors.append('"retry_after" needs "base" and "paused" of at least 1, and "max" of at least "base".')
if errors:
raise ValidationError(errors)
return config


def resolve(stored):
"""The config a stored version puts in force: the defaults, overlaid with what it sets."""
config = merge(get_defaults(), stored, known_only=True)
try:
return validate(config)
except ValidationError as ex:
logger.error('Capacity config is invalid (%s); using the defaults', '; '.join(ex.messages))
return get_defaults()


def get_current():
"""(the newest CapacityConfig row or None, the config in force), read from the database."""
from core.capacity.models import CapacityConfig
latest = CapacityConfig.get_latest()
return latest, resolve(latest.config if latest else None)


def read_current():
"""
get_current() for the request path, bounded: a locked or stalled table costs the call at most
CAPACITY_CONFIG_READ_TIMEOUT_MS, then raises. Inside a transaction (tests) the timeout would outlive this read,
so it isn't set there; requests don't run in one.
"""
if connection.vendor != 'postgresql' or connection.in_atomic_block:
return get_current()
with transaction.atomic():
with connection.cursor() as cursor:
cursor.execute('SET LOCAL statement_timeout = %s', [int(settings.CAPACITY_CONFIG_READ_TIMEOUT_MS)])
return get_current()


def get_config():
"""The config in force, cached per process for CAPACITY_CONFIG_CACHE_SECONDS. Never raises."""
now = time.monotonic()
config = _cache['config']
if config is not None and now < _cache['expires_at']:
return config
with _cache_lock:
if _cache['config'] is not None and time.monotonic() < _cache['expires_at']:
return _cache['config'] # another thread just refreshed it
try:
_, config = read_current()
except Exception as ex:
# Keep the last good config (or the defaults) rather than fail the request, but never refuse on a
# config that couldn't be confirmed: enforce falls back to shadow until a read succeeds.
logger.warning('Capacity config could not be read (%s); using the last known config', ex)
config = copy.deepcopy(_cache['config'] or get_defaults())
if config['mode'] == MODE_ENFORCE:
config['mode'] = MODE_SHADOW
_cache['config'] = config
_cache['expires_at'] = now + settings.CAPACITY_CONFIG_CACHE_SECONDS
return config


def clear_cache():
with _cache_lock:
_cache['config'] = None
_cache['expires_at'] = 0.0


def save_config(changes, user=None, source=SOURCE_API, note='', replace=False):
"""
Put a new version in force: `changes` merged onto the config in force, or onto the defaults with replace.
Raises ValidationError if the result isn't valid. Returns (row, config, {path: [old, new]}); when nothing
changes, no row is added and row is the current one (or None).
"""
from core.capacity.models import CapacityConfig
with transaction.atomic():
if connection.vendor == 'postgresql':
with connection.cursor() as cursor:
cursor.execute('SELECT pg_advisory_xact_lock(%s)', [ADVISORY_LOCK_ID])
latest, previous = get_current()
config = validate(merge(get_defaults() if replace else previous, changes))
changed = diff(previous, config)
if not changed:
return latest, previous, {}
row = CapacityConfig.objects.create(
config=config, previous_config=previous, created_by=user, source=source, note=note or '')
clear_cache()
try:
emit({
'event': CONFIG_LOG_EVENT, 'version': row.id, 'changed_by': getattr(user, 'username', None),
'source': source, 'note': note or None, 'changes': changed, 'mode': config['mode'],
})
except Exception as ex: # the change is saved, and its row is the record; don't fail the request over a log line
logger.warning('Capacity config version %s was saved, but its log line failed: %s', row.id, ex)
return row, config, changed
56 changes: 56 additions & 0 deletions core/capacity/constants.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
MODE_OFF = 'off' # no counting, no headers
MODE_SHADOW = 'shadow' # count, log and send headers, but never refuse
MODE_ENFORCE = 'enforce' # refuse with 429 + Retry-After when a lane is full
MODES = (MODE_OFF, MODE_SHADOW, MODE_ENFORCE)

# Which clients enforce mode refuses. The rest stay in shadow: counted and logged, never refused.
ENFORCE_FOR_AWARE = 'aware' # only clients whose X-OCL-Event-Metadata carries "capacity_aware": "true"
ENFORCE_FOR_ALL = 'all'
ENFORCE_FOR = (ENFORCE_FOR_AWARE, ENFORCE_FOR_ALL)
CAPACITY_AWARE_METADATA_KEY = 'capacity_aware'

# A user's tier is their highest plan group, highest first.
TIER_STAFF = 'staff'
TIER_CORE = 'core'
TIER_EARLY_ACCESS = 'early_access'
TIER_PREVIEW = 'preview'
TIERS = (TIER_STAFF, TIER_CORE, TIER_EARLY_ACCESS, TIER_PREVIEW)

# Lanes: each counts the heavy calls in flight in one scope.
LANE_API_HEAVY = 'api_heavy' # every heavy call, cluster-wide
LANE_API_HEAVY_TASK = 'api_heavy_task' # every heavy call on this API task
LANE_ES2_KNN = 'es2_knn' # semantic $match calls, which run kNN searches in Elasticsearch
LANE_TIER = 'tier' # the calls of every user in this tier
LANE_USER = 'user' # this user's calls
LANES = (LANE_API_HEAVY, LANE_API_HEAVY_TASK, LANE_ES2_KNN, LANE_TIER, LANE_USER)

DECISION_ADMITTED = 'admitted'
DECISION_SHADOW_REFUSED = 'shadow-refused' # a lane was full; enforce mode would have refused it
DECISION_REFUSED = 'refused'
DECISION_UNAVAILABLE = 'unavailable' # Redis couldn't be reached, so the call went ahead uncounted

ENDPOINT_MATCH = '$match'
ENDPOINT_RERANK = '$rerank'

HEADER_DECISION = 'X-OCL-Capacity-Decision'
HEADER_LIMIT = 'X-OCL-Capacity-Limit'
HEADER_IN_FLIGHT = 'X-OCL-Capacity-In-Flight'
HEADER_TIER = 'X-OCL-Capacity-Tier'
HEADER_TIER_LIMIT = 'X-OCL-Capacity-Tier-Limit'
HEADER_TIER_IN_FLIGHT = 'X-OCL-Capacity-Tier-In-Flight'
HEADER_SUGGESTED_CONCURRENCY = 'X-OCL-Capacity-Suggested-Concurrency'
HEADERS = (
HEADER_DECISION, HEADER_LIMIT, HEADER_IN_FLIGHT, HEADER_TIER, HEADER_TIER_LIMIT, HEADER_TIER_IN_FLIGHT,
HEADER_SUGGESTED_CONCURRENCY,
)

CAPACITY_EXCEEDED_ERROR_CODE = 'capacity_exceeded'

# The "event" of the JSON log lines, which CloudWatch metric filters match on.
LOG_EVENT = 'ocl_capacity'
CONFIG_LOG_EVENT = 'ocl_capacity_config'

REDIS_KEY_PREFIX = 'ocl:capacity'

SOURCE_API = 'api'
SOURCE_COMMAND = 'command'
Loading
Loading