DEV Community

PhenoX
PhenoX

Posted on

Token-Bucket Sync CLI: Mediating Asynchronous API Rate Limits Across Multiple Platforms in a Single File

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.

  1. 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.
  2. 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.
  3. 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)
Enter fullscreen mode Exit fullscreen mode

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()
Enter fullscreen mode Exit fullscreen mode

💡 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
  }
}
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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.
Sponsor on GitHub

Top comments (0)