Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 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
1 change: 1 addition & 0 deletions CHANGES/7919.bugfix
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Fixed RedisWorker being blocked by its own stale locks after a worker restarts under the same name, by releasing them at startup.
1 change: 1 addition & 0 deletions CHANGES/7920.bugfix
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Fixed orphaned Redis task/resource locks left behind when a worker's AppStatus was already deleted, via periodic reconciliation.
213 changes: 213 additions & 0 deletions pulpcore/tasking/redis_locks.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,25 @@
# Redis key prefix for resource locks
REDIS_LOCK_PREFIX = "pulp:resource_lock:"

# Redis key prefix for the per-owner lock registry. Each owner has a SET listing the
# lock keys it currently holds so cleanup is O(locks held by owner), not O(all locks).
REDIS_OWNER_REGISTRY_PREFIX = "pulp:owner_locks:"

# Redis SET of owner names that currently hold at least one lock. Enumerating owners
# via SMEMBERS on this key is O(#owners); it avoids a full-keyspace SCAN (a SCAN with
# MATCH still walks every key). Kept in sync inside the acquire/release/cleanup Lua
# scripts, which hardcode this literal -- keep them matching this constant.
ACTIVE_OWNERS_KEY = "pulp:active_owners"

# Throttle key + interval (seconds) for the legacy full-keyspace SCAN fallback used
# during rolling upgrades (locks acquired before the registry existed).
LEGACY_OWNER_SCAN_KEY = "pulp:last_legacy_owner_scan"
LEGACY_OWNER_SCAN_INTERVAL = 900 # ~15 min, fleet-wide

# Owner name prefix used by safe_release_task_locks for immediate tasks that run in an
# API process without an AppStatus. These owners never have an AppStatus row.
IMMEDIATE_OWNER_PREFIX = "immediate-"

REDIS_ACQUIRE_LOCKS_SCRIPT = """
-- KEYS[1]: task_lock_key
-- KEYS[2...]: exclusive_lock_keys, then shared_lock_keys
Expand All @@ -29,6 +48,7 @@
local lock_owner = ARGV[1]
local num_exclusive = tonumber(ARGV[2])
local blocked_resources = {}
local owner_registry_key = "pulp:owner_locks:" .. lock_owner

-- Check task lock first (fail fast)
if redis.call("exists", task_lock_key) == 1 then
Expand Down Expand Up @@ -89,6 +109,15 @@
redis.call("sadd", key, lock_owner)
end

-- Register every held lock key under the owner registry (atomic with acquisition).
-- This lets cleanup enumerate an owner's locks without scanning the whole keyspace.
-- One variadic SADD: all of KEYS are held lock keys (task lock + resource locks).
redis.call("sadd", owner_registry_key, unpack(KEYS))

-- Track this owner in the global active-owners set so reconcile can enumerate
-- lock owners with SMEMBERS instead of a full-keyspace SCAN.
redis.call("sadd", "pulp:active_owners", lock_owner)

-- Return empty table to indicate success
return {}
"""
Expand All @@ -107,6 +136,7 @@
local not_owned_exclusive = {}
local not_in_shared = {}
local task_lock_not_owned = false
local owner_registry_key = "pulp:owner_locks:" .. lock_owner

-- Release exclusive locks
-- Resource keys start at KEYS[2]
Expand All @@ -118,6 +148,7 @@
local current_owner = redis.call("get", key)
if current_owner == lock_owner then
redis.call("del", key)
redis.call("srem", owner_registry_key, key)
elseif current_owner ~= false then
-- Lock exists but we don't own it
table.insert(not_owned_exclusive, resource_name)
Expand All @@ -127,12 +158,18 @@

-- Release shared locks
-- Shared keys start at KEYS[2 + num_exclusive]
-- INVARIANT: an owner runs one task at a time (see RedisWorker.handle_tasks), so it
-- never holds the same shared resource for two concurrent tasks. That lets us drop
-- the registry entry on release unconditionally. If workers ever become concurrent,
-- this must become reference-counted or the shared lock could be released early.
for i = num_exclusive + 1, #KEYS - 1 do
local key = KEYS[1 + i]
local resource_name = ARGV[2 + i]

-- Remove from set
local removed = redis.call("srem", key, lock_owner)
-- No longer a member, so drop the registry entry for this shared key.
redis.call("srem", owner_registry_key, key)
if removed == 0 then
-- We weren't in the set
table.insert(not_in_shared, resource_name)
Expand All @@ -143,15 +180,77 @@
local task_lock_owner = redis.call("get", task_lock_key)
if task_lock_owner == lock_owner then
redis.call("del", task_lock_key)
redis.call("srem", owner_registry_key, task_lock_key)
elseif task_lock_owner ~= false then
-- Task lock exists but we don't own it
task_lock_not_owned = true
end

-- If this owner no longer holds any locks, drop it from the active-owners set
-- (the registry key auto-deletes once empty, so scard == 0 means "no locks left").
if redis.call("scard", owner_registry_key) == 0 then
redis.call("srem", "pulp:active_owners", lock_owner)
end

return {not_owned_exclusive, not_in_shared, task_lock_not_owned}
"""


REDIS_CLEANUP_OWNER_LOCKS_SCRIPT = """
-- Release every lock held by an owner, using the per-owner registry set.
-- ARGV[1]: lock_owner
-- Returns: number of locks released (best effort)
local lock_owner = ARGV[1]
local owner_registry_key = "pulp:owner_locks:" .. lock_owner
local keys = redis.call("smembers", owner_registry_key)
local released = 0

for _, key in ipairs(keys) do
local key_type = redis.call("type", key)["ok"]
if key_type == "string" then
-- Only delete if we still own it (a successor may have re-taken the name).
if redis.call("get", key) == lock_owner then
redis.call("del", key)
released = released + 1
end
elseif key_type == "set" then
-- srem; the set auto-deletes once its last member leaves.
released = released + redis.call("srem", key, lock_owner)
end
-- key_type == "none": stale registry entry, nothing to release.
end

redis.call("del", owner_registry_key)
redis.call("srem", "pulp:active_owners", lock_owner)
return released
"""


REDIS_DELETE_STRING_IF_OWNER_SCRIPT = """
-- Atomically delete a string lock only if it is owned by lock_owner.
-- KEYS[1]: lock key
-- ARGV[1]: lock_owner
-- ARGV[2]: owner_registry_key
if redis.call("get", KEYS[1]) == ARGV[1] then
redis.call("del", KEYS[1])
redis.call("srem", ARGV[2], KEYS[1])
return 1
end
return 0
"""


REDIS_SREM_OWNER_SCRIPT = """
-- Atomically remove lock_owner from a shared set (auto-deletes when empty).
-- KEYS[1]: shared set key
-- ARGV[1]: lock_owner
-- ARGV[2]: owner_registry_key
local removed = redis.call("srem", KEYS[1], ARGV[1])
redis.call("srem", ARGV[2], KEYS[1])
return removed
"""


def resource_to_lock_key(resource_name):
"""
Convert a resource name to a Redis lock key.
Expand All @@ -178,6 +277,120 @@ def get_task_lock_key(task_id):
return f"task:{task_id}"


def get_owner_registry_key(owner):
"""Return the Redis key for an owner's lock registry SET."""
return f"{REDIS_OWNER_REGISTRY_PREFIX}{owner}"


def _decode(value):
"""Decode a redis bytes value to str (redis-py returns bytes by default)."""
return value.decode() if isinstance(value, bytes) else value


def _legacy_scan_cleanup_for_owner(redis_conn, owner):
"""
Release an owner's locks by scanning the keyspace (no registry available).

Used only for locks acquired before the per-owner registry existed (rolling
upgrade). Uses SCAN (never KEYS) and atomic per-key Lua so a concurrent worker
that re-took a key by the same name is not clobbered.

Returns:
int: Number of locks released (best effort).
"""
registry_key = get_owner_registry_key(owner)
delete_if_owner = redis_conn.register_script(REDIS_DELETE_STRING_IF_OWNER_SCRIPT)
srem_owner = redis_conn.register_script(REDIS_SREM_OWNER_SCRIPT)
released = 0

for key in redis_conn.scan_iter(match="task:*", count=500):
if _decode(redis_conn.get(key)) == owner:
released += delete_if_owner(keys=[key], args=[owner, registry_key])

for key in redis_conn.scan_iter(match=f"{REDIS_LOCK_PREFIX}*", count=500):
if _decode(redis_conn.type(key)) == "string":
if _decode(redis_conn.get(key)) == owner:
released += delete_if_owner(keys=[key], args=[owner, registry_key])
else:
released += srem_owner(keys=[key], args=[owner, registry_key])

redis_conn.delete(registry_key)
redis_conn.srem(ACTIVE_OWNERS_KEY, owner)
return released


def cleanup_locks_for_owner(redis_conn, owner, allow_legacy_scan=False):
"""
Release all Redis locks held by ``owner``.

Prefers the per-owner registry (O(locks held by owner)). Falls back to a legacy
keyspace SCAN only when the registry is missing and ``allow_legacy_scan`` is set.

Args:
redis_conn: Redis connection
owner (str): The lock owner (worker name or ``immediate-{task_pk}``)
allow_legacy_scan (bool): Permit the legacy SCAN fallback for pre-registry locks

Returns:
bool: True if cleanup completed (including a no-op), False on error so the
caller can retain state and retry on a later pass.
"""
registry_key = get_owner_registry_key(owner)
try:
released = 0
if redis_conn.exists(registry_key):
cleanup_script = redis_conn.register_script(REDIS_CLEANUP_OWNER_LOCKS_SCRIPT)
released = cleanup_script(keys=[], args=[owner])
elif allow_legacy_scan:
released = _legacy_scan_cleanup_for_owner(redis_conn, owner)
else:
# No registry and no scan: nothing to release, but drop any stale
# active-owners marker so reconcile stops re-visiting this owner.
redis_conn.srem(ACTIVE_OWNERS_KEY, owner)
if released:
_logger.info("Reclaimed %d Redis lock(s) held by owner %s", released, owner)
return True
except Exception as e:
_logger.error("Error cleaning up Redis locks for owner %s: %s", owner, e)
return False


def collect_lock_owners(redis_conn, allow_legacy_scan=False):
"""
Return the set of owner names that currently hold Redis locks.

Owners are read from the global active-owners SET via SMEMBERS -- O(#owners) and
scan-free (a SCAN with MATCH still walks the whole keyspace, so scanning for
registry keys would be O(all locks)). The legacy SCAN of the full ``task:*`` /
``pulp:resource_lock:*`` keyspace is expensive and is only run when
``allow_legacy_scan`` is set (throttled by the caller) to catch pre-registry locks.

Args:
redis_conn: Redis connection
allow_legacy_scan (bool): Also discover owners of pre-registry (legacy) locks

Returns:
set: Owner names holding at least one lock.
"""
owners = {_decode(member) for member in redis_conn.smembers(ACTIVE_OWNERS_KEY)}

if allow_legacy_scan:
for key in redis_conn.scan_iter(match="task:*", count=500):
value = redis_conn.get(key)
if value:
owners.add(_decode(value))
for key in redis_conn.scan_iter(match=f"{REDIS_LOCK_PREFIX}*", count=500):
if _decode(redis_conn.type(key)) == "string":
value = redis_conn.get(key)
if value:
owners.add(_decode(value))
else:
for member in redis_conn.smembers(key):
owners.add(_decode(member))

return owners


def extract_task_resources(task):
"""
Extract exclusive and shared resources from a task.
Expand Down
Loading
Loading