Token-Bucket Sync CLI: Mediating Asynchronous API Rate Limits Across Multiple Platforms in a Single File
Why We Abandoned Redis: The Despair of Early Development
When I first designed this pipeline, I tried to be modern by synchronizing rate limits using a distributed cache (Redis).
I was bound by the preconceived notion that atomic counter operations were absolutely necessary when hitting APIs simultaneously from multiple processes or containers. However, this was the beginning of a nightmare.
- Bloated Environmental Dependencies: Even for small local batch scripts, CI/CD containers, or standalone VPS environments, starting and constantly running a Redis server became a prerequisite. I was plagued by connection errors with every deployment.
- The Quagmire of "Broken Tests": Every time I ran unit tests or local concurrency validations, the codebase was polluted with Redis mocks or the instantiation of in-memory versions.
- Excessive Complexity: What we truly needed was simply a boolean answer to: "Is it safe to execute this batch at this very moment?" We did not need a global distributed lock shared across a massive cluster of microservices.
"Can't we distill this into a more primitive, robust mechanism with zero external dependencies?"
This line of thought led me to an approach combining OS-level file locking (fcntl / msvcrt) with a local JSON state file—a "token-bucket synchronization storage."
To illustrate the paradigm shift, consider the architectural difference:
flowchart TD
subgraph RedisArchitecture ["Previous: Redis-based Architecture"]
P1["Process 1"]
P2["Process 2"]
P3["Process 3"]
Redis["Redis Server (Network/Memory Overhead)"]
P1 -- "Network Request" --> Redis
P2 -- "Network Request" --> Redis
P3 -- "Network Request" --> Redis
end
subgraph FileLockArchitecture ["Current: File Lock Architecture"]
F1["Process 1"]
F2["Process 2"]
F3["Process 3"]
StateFile["Local JSON State File"]
F1 -- "OS File Lock (fcntl/msvcrt)" --> StateFile
F2 -- "OS File Lock (fcntl/msvcrt)" --> StateFile
F3 -- "OS File Lock (fcntl/msvcrt)" --> StateFile
end
The Core of the Implementation: A 10-Second Time Limit and the Battle for Mutual Exclusion
The tool I created, the "Token-Bucket Sync CLI," is exceedingly simple. Based on the capacity and refill rate of each service defined in a configuration file (JSON), it determines whether a request is permitted and atomically updates the state file.
However, to make this work in a "multi-process environment with frequent contention," I had to overcome several gritty technical hurdles.
1. The Quagmire of Cross-Platform File Locking
In a Linux environment, fcntl.flock is the undisputed choice. However, considering that some team members develop on Windows machines and acknowledging the diversity of local testing environments, a fallback to msvcrt.locking was absolutely necessary.
Furthermore, to prevent a race condition in the initial state where the file does not exist (the exact moment multiple processes simultaneously execute state_path.write_text("{}")), the sequence of checking for existence, opening the file, and acquiring the lock required meticulous attention.
2. The 10-Second Timeout and the Philosophy of Retries
When multiple processes attempt to update the state file simultaneously under high load, lock waiting naturally occurs. Throwing an exception immediately in this scenario is nonsensical.
Batch processing pipelines are resilient to "being made to wait" but fragile against "immediate failure." Therefore, if file lock contention or a temporary json.JSONDecodeError (reading while another process is writing) occurs, the system is designed to retry in 0.1-second intervals, holding out for a maximum of 10 seconds.
try:
if not state_path.exists():
state_path.write_text("{}")
with open(state_path, "r+") as f:
lock_file(f)
try:
# Token calculation and atomic write-back
...
finally:
unlock_file(f)
except (IOError, json.JSONDecodeError):
time.sleep(0.1)
Through this gritty combination of retries and file locks, we successfully arbitrated simultaneous access from multiple processes completely, without relying on any external middleware.
Here is the flowchart representing the atomic lock and retry mechanism:
flowchart TD
Start["Start Request"]
CheckTimeout{"Is Timeout Exceeded?"}
TryLock["Attempt OS File Lock"]
LockSuccess{"Lock Acquired?"}
ReadJSON["Read & Parse JSON"]
CalcToken["Calculate Token Refill & Consumption"]
HasToken{"Sufficient Tokens?"}
WriteJSON["Truncate & Write JSON"]
Unlock["Release Lock"]
Wait["Sleep 0.1s"]
Approve["Return Approved"]
Reject["Return Rejected (Wait Time)"]
Error["Raise TimeoutError"]
Start -- "Initialize" --> CheckTimeout
CheckTimeout -- "Yes" --> Error
CheckTimeout -- "No" --> TryLock
TryLock -- "Execute" --> LockSuccess
LockSuccess -- "No (IOError)" --> Wait
LockSuccess -- "Yes" --> ReadJSON
ReadJSON -- "JSONDecodeError" --> Unlock
ReadJSON -- "Success" --> CalcToken
CalcToken -- "Evaluate" --> HasToken
HasToken -- "Yes" --> WriteJSON
HasToken -- "No" --> WriteJSON
WriteJSON -- "Flush to Disk" --> Unlock
Unlock -- "from IOError/JSONDecodeError" --> Wait
Unlock -- "from Success (Has Token)" --> Approve
Unlock -- "from Success (No Token)" --> Reject
Wait -- "Retry" --> CheckTimeout
The Final Code: Token-Bucket Sync CLI
By simply saving the following code as token_bucket_sync.py and granting it execution permissions, you can embed a robust rate-limit arbitration mechanism into your pipeline.
#!/usr/bin/env python3
import argparse
import json
import os
import sys
import time
from pathlib import Path
if os.name == 'nt':
import msvcrt
def lock_file(f):
msvcrt.locking(f.fileno(), msvcrt.LK_LOCK, 1)
def unlock_file(f):
msvcrt.locking(f.fileno(), msvcrt.LK_UNLCK, 1)
else:
import fcntl
def lock_file(f):
fcntl.flock(f.fileno(), fcntl.LOCK_EX)
def unlock_file(f):
fcntl.flock(f.fileno(), fcntl.LOCK_UN)
DEFAULT_STATE_FILE = ".token_bucket_state.json"
def load_and_update_buckets(services_config, state_path: Path, timeout: float = 10.0):
start_time = time.time()
while True:
if time.time() - start_time > timeout:
raise TimeoutError("Rate limit synchronization timeout exceeded.")
try:
if not state_path.exists():
state_path.write_text("{}")
with open(state_path, "r+") as f:
lock_file(f)
try:
content = f.read()
state = json.loads(content) if content.strip() else {}
now = time.time()
plan = {}
all_approved = True
for svc, cfg in services_config.items():
capacity = cfg.get("capacity", 10)
refill_rate = cfg.get("refill_rate", 1)
requested = cfg.get("requested", 1)
svc_state = state.get(svc, {"tokens": capacity, "last_update": now})
tokens = min(capacity, svc_state["tokens"] + (now - svc_state["last_update"]) * refill_rate)
if tokens >= requested:
tokens -= requested
approved, wait_time = True, 0.0
else:
approved, wait_time = False, (requested - tokens) / refill_rate
all_approved = False
state[svc] = {"tokens": tokens, "last_update": now}
plan[svc] = {"approved": approved, "current_tokens": round(tokens, 2), "wait_time_seconds": round(wait_time, 2)}
f.seek(0)
f.truncate()
json.dump(state, f)
f.flush()
os.fsync(f.fileno())
return {"all_approved": all_approved, "plan": plan, "timestamp": now}
finally:
unlock_file(f)
except (IOError, json.JSONDecodeError):
time.sleep(0.1)
def main():
parser = argparse.ArgumentParser(description="Token-Bucket Sync CLI")
parser.add_argument("--config", required=True, help="Path to JSON config")
parser.add_argument("--state-file", default=DEFAULT_STATE_FILE, help="Path to state file")
parser.add_argument("--timeout", type=float, default=10.0, help="Timeout in seconds")
args = parser.parse_args()
try:
config_path = Path(args.config)
if not config_path.exists():
raise FileNotFoundError(f"Config file not found: {args.config}")
with open(config_path, "r", encoding="utf-8") as cf:
services_config = json.load(cf)
result = load_and_update_buckets(services_config, Path(args.state_file), timeout=args.timeout)
print(json.dumps(result, indent=2))
sys.exit(0 if result["all_approved"] else 2)
except Exception as e:
print(json.dumps({"error": str(e)}))
sys.exit(1)
if __name__ == "__main__":
main()
💡 For immediate deployment: The complete source code suite (ZIP) for this architecture is available on Gumroad for $0+ (Pay What You Want).
Usage and Integration into Your Pipeline
Usage is incredibly intuitive. First, prepare a configuration JSON (e.g., config.json) defining the parameters for each service you want to validate.
{
"twitter_api": {
"capacity": 5,
"refill_rate": 0.5,
"requested": 1
},
"discord_webhook": {
"capacity": 30,
"refill_rate": 5,
"requested": 1
}
}
Insert this right before your batch processing or within your shell script pipeline.
#!/bin/bash
# Evaluate rate limit availability
python3 token_bucket_sync.py --config config.json
EXIT_CODE=$?
if [ $EXIT_CODE -eq 0 ]; then
echo "All API requests approved. Starting execution."
# Proceed to actual API calls
elif [ $EXIT_CODE -eq 2 ]; then
echo "API limit reached for some services. Execution deferred or skipped."
else
echo "An unexpected error occurred."
exit 1
fi
An exit code of 0 means "All Approved," while 2 means "Partial Limit Exceeded (Defer Burst)." Because it directly utilizes standard exit codes, it integrates seamlessly into shell scripts and CI flow control architectures.
Conclusion
Over-engineering architecture often ends up strangling the developer.
By discarding the assumption that "it's a distributed system, so Redis is required" and returning to the pure simplicity of algorithms and OS primitive functions, I acquired an astonishingly robust and portable tool.
I hope this single-file mediator serves the pipelines of all engineers struggling with multi-API rate limits.
If this engineering log saved your production server (and your sanity), consider supporting our architecture on GitHub Sponsors.
Top comments (0)