LEN: Add batch mode and some optimizations

This commit is contained in:
DazedAnon 2026-06-12 16:39:46 -05:00
parent 892a8b0574
commit de56782391
3 changed files with 841 additions and 166 deletions

View file

@ -16,6 +16,7 @@ An AI-powered game translation tool with a GUI. Translate RPG Maker, Ren'Py, Tyr
- [Using the GUI](#using-the-gui)
- [Vocab & Prompt](#vocab--prompt)
- [Tips](#tips)
- [Batch Translation (Anthropic, 50% off)](#batch-translation-anthropic-50-off)
- [Folder Structure](#folder-structure)
- [Finding Untranslated Text (Snipping Tool OCR)](#finding-untranslated-text-snipping-tool-ocr)
- [RPG Maker Translation Workflow](#rpg-maker-translation-workflow)
@ -185,6 +186,53 @@ This is the system prompt sent to the AI. A default `prompt.txt` is included and
---
## Batch Translation (Anthropic, 50% off)
When using a Claude model, the CLI offers a third mode that translates through the
[Anthropic Message Batches API](https://platform.claude.com/docs/en/build-with-claude/batch-processing.md)
— every token (input, output, and prompt-cache reads/writes) is billed at **50% of the live price**.
Batches usually finish within an hour (24h worst case), so use it for large jobs where you don't
need results immediately.
```
python start.py
-> 3. Batch Translate (Anthropic Batches API, 50% off)
```
How it works (all engine modules are supported automatically):
1. **Pass 1 (collect)** — files are processed normally, but instead of calling the API each
request is queued to `log/batch_requests.json`. Requests are byte-identical to live ones:
the static `prompt.txt` block is cached with a 1h TTL, matched vocab and translation
history ride along per request, and structured output enforces the exact line count.
Speaker/variable names still translate live during this pass (they get embedded into the
dialogue payloads, so both passes must resolve them identically) — they're a tiny share
of the volume.
2. **Cost estimate** — before anything is submitted you get a cost breakdown
(batch + cache / batch worst-case / live price) and a y/n confirmation.
3. **Submit / poll / fetch** — the batch is submitted, polled until it ends
(`batchPollInterval` env var controls the interval, default 60s), and the results are
saved to `log/batch_results.json`. Ctrl-C while polling is safe — the batch keeps
processing server-side.
4. **Pass 2 (consume)** — files are processed again; every payload is filled from the batch
results through the normal validation pipeline (line counts, placeholders, content
checks). Anything the batch missed or that fails validation falls back to the live API
automatically, so the output is always complete.
Context note: in live mode the rolling translation history contains the previous batch's
English lines; in batch mode requests are independent, so the history carries the previous
batch's *source* lines instead. The model still sees the surrounding scene, and `vocab.txt`
keeps names and terms consistent — so keep your vocab file in good shape for batch runs.
Cost tracking is exact: per-file and total costs printed after the consume pass use the real
billed token counts (cache reads at 0.1x, cache writes at 2x, output at the output rate) with
the 50% batch discount applied.
`python selftest_batch.py` round-trips the whole flow offline (no API key needed) if you want
to verify the pipeline after making changes.
---
## Folder Structure
| Folder | Purpose |

View file

@ -88,22 +88,50 @@ tqdm.write(
def main():
# Clear the translation cache at the start of the run
from util.translation import clear_cache
clear_cache()
estimate = ""
batch_mode = False
speaker_parse = False # Deferred until after engine select
while estimate == "":
estimate = input("Select Mode:\n\n 1. Translate\n 2. Estimate\n")
estimate = input("Select Mode:\n\n 1. Translate\n 2. Estimate\n 3. Batch Translate (Anthropic Batches API, 50% off)\n")
match estimate:
case "1":
estimate = False
case "2":
estimate = True
case "3":
estimate = False
batch_mode = True
case _:
estimate = ""
resume_state = None
if batch_mode:
from util.translation import isClaudeNative, batchRunState
if not isClaudeNative(os.getenv("model", "")):
tqdm.write(
Fore.RED
+ "Batch Translate requires a Claude model with the 'api' env var unset or pointing at anthropic.com."
+ Fore.RESET
)
return
# An interrupted batch run can be resumed instead of re-collecting
# (a second submission would be billed again).
resume_state = batchRunState()
if resume_state:
confirm = ""
while confirm not in ("y", "n"):
confirm = input(f"A previous batch run was interrupted ({resume_state}). Resume it? (y/n)\n").strip().lower()
if confirm == "n":
resume_state = None
# Clear the translation cache at the start of the run. Kept when resuming a
# batch so names translated during collect stay consistent with the queued
# payloads in the consume pass.
if not resume_state:
clear_cache()
version = ""
while True:
tqdm.write("Select game engine:\n")
@ -125,7 +153,7 @@ files to translate are in the /files folder and that you picked the right game e
# If translating RPGMaker MV/MZ, prompt for speaker parse mode
speaker_parse = False
if version == 0 and not estimate:
if version == 0 and not estimate and not batch_mode:
sub = ""
while sub == "":
sub = input("RPGMaker MV/MZ options:\n\n 1. Standard Translate\n 2. Parse Speakers (collect speaker names only)\n")
@ -180,78 +208,133 @@ files to translate are in the /files folder and that you picked the right game e
except Exception:
pass
# Use single worker for estimate mode to prevent race conditions
max_workers = 1 if estimate else THREADS
with ThreadPoolExecutor(max_workers=max_workers) as executor:
futures = []
files_root = "files"
def runFiles(estimate):
runCost = totalCost
# Use single worker for estimate mode to prevent race conditions
max_workers = 1 if estimate else THREADS
with ThreadPoolExecutor(max_workers=max_workers) as executor:
futures = []
files_root = "files"
# Special-case: Images engine expects a folder, not a file; schedule per directory containing assets
if MODULES[version][0] == "Images":
for root, dirs, filenames in os.walk(files_root):
# Skip hidden/system directories
dirs[:] = [d for d in dirs if d not in {".git", "__pycache__"}]
# Special-case: Images engine expects a folder, not a file; schedule per directory containing assets
if MODULES[version][0] == "Images":
for root, dirs, filenames in os.walk(files_root):
# Skip hidden/system directories
dirs[:] = [d for d in dirs if d not in {".git", "__pycache__"}]
# Skip the root 'files' itself to avoid processing everything twice
# We'll still allow scheduling for root if it contains assets
# Skip the root 'files' itself to avoid processing everything twice
# We'll still allow scheduling for root if it contains assets
# Only schedule directories that contain potential assets
has_assets = any(fn.lower().endswith((".png", ".txt")) for fn in filenames)
if not has_assets:
continue
# Compute relative directory path and ensure translated mirror exists
rel_dir = os.path.relpath(root, files_root).replace(os.sep, "/")
if rel_dir == ".":
# Represent root as empty string so handler creates files under translated/ directly
rel_dir = ""
try:
target_dir = os.path.join("translated", rel_dir.replace("/", os.sep)) if rel_dir else "translated"
os.makedirs(target_dir, exist_ok=True)
except Exception:
pass
futures.append(
executor.submit(MODULES[version][2], rel_dir, estimate)
)
else:
# Gather all candidate files recursively
for root, dirs, filenames in os.walk(files_root):
# Skip hidden/system directories if any
dirs[:] = [d for d in dirs if d not in {".git", "__pycache__"}]
for fname in filenames:
if fname == ".gitkeep":
# Only schedule directories that contain potential assets
has_assets = any(fn.lower().endswith((".png", ".txt")) for fn in filenames)
if not has_assets:
continue
abs_path = os.path.join(root, fname)
# Build relative path from 'files' root using POSIX-style separators so handlers can do 'files/' + rel
rel_path = os.path.relpath(abs_path, files_root)
rel_path_posix = rel_path.replace(os.sep, "/")
# Compute relative directory path and ensure translated mirror exists
rel_dir = os.path.relpath(root, files_root).replace(os.sep, "/")
if rel_dir == ".":
# Represent root as empty string so handler creates files under translated/ directly
rel_dir = ""
try:
target_dir = os.path.join("translated", rel_dir.replace("/", os.sep)) if rel_dir else "translated"
os.makedirs(target_dir, exist_ok=True)
except Exception:
pass
# Check extension match for the selected module version
for m in MODULES[version][1]:
if rel_path_posix.endswith(m):
# Ensure the corresponding directory exists under 'translated'
rel_dir = os.path.dirname(rel_path_posix)
if rel_dir:
try:
os.makedirs(os.path.join("translated", rel_dir.replace("/", os.sep)), exist_ok=True)
except Exception:
# Best-effort; handler may attempt write and fail if permissions are insufficient
pass
futures.append(
executor.submit(MODULES[version][2], rel_dir, estimate)
)
else:
# Gather all candidate files recursively
for root, dirs, filenames in os.walk(files_root):
# Skip hidden/system directories if any
dirs[:] = [d for d in dirs if d not in {".git", "__pycache__"}]
futures.append(
executor.submit(MODULES[version][2], rel_path_posix, estimate)
)
break # Avoid double-adding if multiple ext entries match
for fname in filenames:
if fname == ".gitkeep":
continue
for future in as_completed(futures):
abs_path = os.path.join(root, fname)
# Build relative path from 'files' root using POSIX-style separators so handlers can do 'files/' + rel
rel_path = os.path.relpath(abs_path, files_root)
rel_path_posix = rel_path.replace(os.sep, "/")
# Check extension match for the selected module version
for m in MODULES[version][1]:
if rel_path_posix.endswith(m):
# Ensure the corresponding directory exists under 'translated'
rel_dir = os.path.dirname(rel_path_posix)
if rel_dir:
try:
os.makedirs(os.path.join("translated", rel_dir.replace("/", os.sep)), exist_ok=True)
except Exception:
# Best-effort; handler may attempt write and fail if permissions are insufficient
pass
futures.append(
executor.submit(MODULES[version][2], rel_path_posix, estimate)
)
break # Avoid double-adding if multiple ext entries match
for future in as_completed(futures):
try:
runCost = future.result()
except Exception as e:
tracebackLineNo = str(traceback.extract_tb(sys.exc_info()[2])[-1].lineno)
tqdm.write(Fore.RED + str(e) + "|" + tracebackLineNo + Fore.RESET)
return runCost
if batch_mode:
from util.translation import (
set_batch_phase,
clearBatchFiles,
pendingBatchRequests,
estimateBatchCost,
runTranslationBatches,
)
poll = int(os.getenv("batchPollInterval", "60") or 60)
run_consume = True
if resume_state is None:
clearBatchFiles()
# Pass 1 — collect every needed request without calling the API.
tqdm.write(Fore.CYAN + "[BATCH] Pass 1/2: collecting requests..." + Fore.RESET)
set_batch_phase("collect")
try:
totalCost = future.result()
except Exception as e:
tracebackLineNo = str(traceback.extract_tb(sys.exc_info()[2])[-1].lineno)
tqdm.write(Fore.RED + str(e) + "|" + tracebackLineNo + Fore.RESET)
totalCost = runFiles(False)
finally:
set_batch_phase(None)
if pendingBatchRequests() == 0:
tqdm.write("[BATCH] No requests queued — nothing needed the API.")
run_consume = False
else:
estimateBatchCost()
confirm = ""
while confirm not in ("y", "n"):
confirm = input("Submit batch? (y/n)\n").strip().lower()
if confirm == "n":
tqdm.write("[BATCH] Not submitted. The queue is kept in log/batch_requests.json.")
return
runTranslationBatches(poll)
elif resume_state == "submitted":
tqdm.write(Fore.CYAN + "[BATCH] Resuming the submitted batch..." + Fore.RESET)
runTranslationBatches(poll)
else: # "fetched" — results already downloaded, just write the files
tqdm.write(Fore.CYAN + "[BATCH] Resuming from fetched results..." + Fore.RESET)
if run_consume:
# Pass 2 — write the translated files from the fetched results.
# Anything the batch missed falls back to the live API.
tqdm.write(Fore.CYAN + "[BATCH] Pass 2/2: writing translated files..." + Fore.RESET)
set_batch_phase("consume")
try:
totalCost = runFiles(False)
finally:
set_batch_phase(None)
else:
totalCost = runFiles(estimate)
# Finalize speaker parse mode by writing collected speakers to vocab
if speaker_parse:

View file

@ -95,6 +95,29 @@ def _normalize_openai_base_url(url: str) -> str:
_url += "/"
return _url
def isClaudeModel(model):
"""True when the model name looks like an Anthropic Claude model."""
return bool(model) and any(x in model.lower() for x in ("claude", "sonnet", "haiku", "opus"))
def isClaudeNative(model):
"""True when this model routes to the native Anthropic SDK.
Mirrors the routing check in translateText: the model must look like Claude
AND the configured API URL must be unset or point at anthropic.com. Any
other custom URL (e.g. DeepSeek, OpenAI proxy) uses the OpenAI-compatible
path even for Claude-named models.
"""
live_api = os.getenv("api", "").strip()
return isClaudeModel(model) and (not live_api or "anthropic" in live_api.lower())
# Models that REJECT sampling params (temperature/top_p/top_k) with a 400 —
# Claude Opus 4.7 and up retired them in favour of adaptive thinking. Matches
# opus-4-7, opus-4-8, opus-4-10+ and fable, but NOT opus-4-6 or Sonnet/Haiku.
_NO_SAMPLING_RE = re.compile(r"opus-4-(?:[7-9]\b|[1-9]\d)|fable", re.I)
# Tracks which distinct batch sizes have already been cache-written during this estimate run.
# Each unique numLines value maps to a distinct output_config schema → one write per size.
# Persisted to disk so sequential GUI subprocesses share state.
@ -620,6 +643,429 @@ def set_var_translations_batch(pairs):
_save_var_map()
# ===== Anthropic Message Batches (50% off all token usage) =====
# Batch translation is a two-pass flow driven by the batch phase (kept in the
# BATCH_PHASE env var so GUI subprocesses inherit it):
# collect: translateAI builds each cache-missed request (byte-identical to a
# live request, including the cached system block) and queues it
# instead of calling the API. Text is left untranslated.
# consume: translateAI feeds the fetched batch responses through the normal
# validation/restore path; anything missing or invalid falls back
# to the live API automatically.
# Between the passes, submit/poll/fetch the queue with runTranslationBatches().
BATCH_QUEUE_FILE = Path("log/batch_requests.json")
BATCH_STATE_FILE = Path("log/batch_state.json")
BATCH_RESULTS_FILE = Path("log/batch_results.json")
BATCH_LOCK_FILE = Path("log/batch_files.lock")
BATCH_LOCK = threading.RLock()
# API limits are 100,000 requests / 256 MB per batch; stay safely under both.
BATCH_MAX_REQUESTS = 100_000
BATCH_MAX_BYTES = 200 * 1024 * 1024
_batch_phase = None
_batch_results = None # in-memory copy of BATCH_RESULTS_FILE (read-only during consume)
_batch_queue_pending = {} # process-local queue entries not yet flushed to disk
def set_batch_phase(phase):
"""Set the batch phase ('collect', 'consume' or None) for this process and
any subprocesses it spawns (the GUI runs one per file)."""
global _batch_phase, _batch_results
_batch_phase = phase if phase in ("collect", "consume") else None
if _batch_phase:
os.environ["BATCH_PHASE"] = _batch_phase
else:
os.environ.pop("BATCH_PHASE", None)
_batch_results = None # phase change invalidates the in-memory results copy
def get_batch_phase():
"""Current batch phase, or None when batch translation is off."""
if _batch_phase:
return _batch_phase
phase = os.getenv("BATCH_PHASE", "").strip().lower()
return phase if phase in ("collect", "consume") else None
@contextmanager
def _batch_file_lock():
"""Cross-process lock for the batch queue/state/results files."""
BATCH_LOCK_FILE.parent.mkdir(parents=True, exist_ok=True)
with open(BATCH_LOCK_FILE, "a+b") as lock_file:
if os.name == "nt":
import msvcrt
lock_file.seek(0)
msvcrt.locking(lock_file.fileno(), msvcrt.LK_LOCK, 1)
try:
yield
finally:
lock_file.seek(0)
msvcrt.locking(lock_file.fileno(), msvcrt.LK_UNLCK, 1)
else:
import fcntl
fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX)
try:
yield
finally:
fcntl.flock(lock_file.fileno(), fcntl.LOCK_UN)
def _read_batch_file(path):
"""Read a batch JSON file; return an empty dict if it is unavailable."""
try:
if path.exists():
with open(path, "r", encoding="utf-8") as f:
data = json.load(f)
return data if isinstance(data, dict) else {}
except Exception:
pass
return {}
def _write_batch_file(path, data):
"""Atomically write a batch JSON file (no indent — queues can be large)."""
path.parent.mkdir(parents=True, exist_ok=True)
tmp_file = path.with_name(f"{path.name}.{os.getpid()}.{threading.get_ident()}.tmp")
with open(tmp_file, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False)
tmp_file.replace(path)
def peek_cached_translation(payload, language):
"""Cache lookup that never blocks or writes a pending marker.
The collect pass uses this instead of get_cached_translation so abandoned
pending markers can't stall the consume pass for CACHE_PENDING_TTL."""
key = get_cache_key(payload, language)
with CACHE_LOCK:
with _translation_cache_file_lock():
cache = _read_cache_from_disk()
if _cache:
cache = _merge_translation_caches(cache, _cache)
entry = cache.get(key)
if entry is None or _is_pending_cache_entry(entry):
return None
return entry
def queue_batch_request(payload, language, params):
"""Queue one Batches API request during the collect pass.
Deduped by the same key the translation cache uses, so identical payloads
across files are only paid for once. Entries are buffered in memory and
merged to disk by flush_batch_queue() at the end of each translateAI call.
"""
key = get_cache_key(payload, language)
with BATCH_LOCK:
_batch_queue_pending[key] = {"payload": payload, "language": language, "params": params}
return key
def flush_batch_queue():
"""Merge this process's pending queue entries into the on-disk queue."""
global _batch_queue_pending
with BATCH_LOCK:
if not _batch_queue_pending:
return
pending, _batch_queue_pending = _batch_queue_pending, {}
try:
with _batch_file_lock():
queue = _read_batch_file(BATCH_QUEUE_FILE)
for key, entry in pending.items():
queue.setdefault(key, entry)
_write_batch_file(BATCH_QUEUE_FILE, queue)
except Exception:
_batch_queue_pending.update(pending) # keep entries for the next flush
def take_batch_result(payload, language):
"""Return the fetched batch response dict for a payload, or None.
The results file is loaded once per process it is written before the
consume pass starts and never changes mid-consume."""
global _batch_results
if _batch_results is None:
with BATCH_LOCK:
if _batch_results is None:
with _batch_file_lock():
_batch_results = _read_batch_file(BATCH_RESULTS_FILE)
return _batch_results.get(get_cache_key(payload, language))
def pendingBatchRequests():
"""Number of queued batch requests (call after the collect pass)."""
flush_batch_queue()
with _batch_file_lock():
return len(_read_batch_file(BATCH_QUEUE_FILE))
def batchRunState():
"""'submitted' when a batch is still in flight, 'fetched' when results are
waiting to be consumed, else None. Lets an interrupted batch run resume
instead of re-collecting and paying for a second submission."""
with _batch_file_lock():
if _read_batch_file(BATCH_STATE_FILE).get("batches"):
return "submitted"
if _read_batch_file(BATCH_RESULTS_FILE):
return "fetched"
return None
def clearBatchFiles():
"""Remove queue/state/results left over from any previous batch run."""
global _batch_results, _batch_queue_pending
with BATCH_LOCK:
_batch_results = None
_batch_queue_pending = {}
with _batch_file_lock():
for path in (BATCH_QUEUE_FILE, BATCH_STATE_FILE, BATCH_RESULTS_FILE):
try:
if path.exists():
path.unlink()
except Exception:
pass
def _get_anthropic_client():
key = os.getenv("key", "").strip()
if not key:
raise Exception("Batch translation requires the 'key' env var (see .env).")
return anthropic.Anthropic(api_key=key)
def estimateBatchCost(model=None):
"""Print a cost estimate for the queued batch requests and return it.
Cached-prefix accounting mirrors what Anthropic bills: each distinct cached
prefix is written once (2x input rate at the 1h TTL) and read by every other
request that shares it (0.10x); everything is then halved by the batch
discount. Cache hits inside a batch are best-effort, so the no-cache batch
figure is the worst-case bound.
"""
flush_batch_queue()
with _batch_file_lock():
queue = _read_batch_file(BATCH_QUEUE_FILE)
if not queue:
print("[BATCH] No batch requests queued.", flush=True)
return None
enc = tiktoken.encoding_for_model("gpt-4")
prefix_count = {} # cached prefix text -> how many requests reuse it
prefix_tokens = {} # cached prefix text -> token count
dynamic_tokens = 0
output_tokens = 0
models = set()
for entry in queue.values():
params = entry.get("params", {})
if params.get("model"):
models.add(params["model"])
blocks = params.get("system") or []
cut = 0 # split system blocks at the cache breakpoint (0 = nothing cached)
for i, b in enumerate(blocks):
if "cache_control" in b:
cut = i + 1
break
prefix = "".join(b.get("text", "") for b in blocks[:cut])
dyn = "".join(b.get("text", "") for b in blocks[cut:])
if prefix:
prefix_count[prefix] = prefix_count.get(prefix, 0) + 1
if prefix not in prefix_tokens:
prefix_tokens[prefix] = len(enc.encode(prefix))
msg_text = "\n".join(str(m.get("content", "")) for m in params.get("messages", []))
dynamic_tokens += len(enc.encode(dyn)) + len(enc.encode(msg_text)) + 8
# Output heuristic mirrors countTokens(): payload tokens x 2.5 covers
# the echoed JSON scaffold plus EN expansion.
output_tokens += round(len(enc.encode(str(entry.get("payload", "")))) * 2.5)
est_model = model or next(iter(models), None) or os.getenv("model", "")
# Use Anthropic's count_tokens for the exact cached-prefix size when possible.
if prefix_tokens:
try:
client = _get_anthropic_client()
for prefix in list(prefix_tokens.keys()):
resp = client.messages.count_tokens(
model=est_model,
system=[{"type": "text", "text": prefix}],
messages=[{"role": "user", "content": "x"}],
)
prefix_tokens[prefix] = resp.input_tokens
except Exception:
pass # tiktoken estimate already in place
cache_write_tok = sum(prefix_tokens.values())
cache_read_tok = sum(prefix_tokens[p] * (prefix_count[p] - 1) for p in prefix_count)
raw_input_tok = sum(prefix_tokens[p] * prefix_count[p] for p in prefix_count) + dynamic_tokens
pricing = getPricingConfig(est_model)
in_rate = pricing["inputAPICost"] / 1_000_000
out_rate = pricing["outputAPICost"] / 1_000_000
# Batch = 50% off every token, including cache writes (2x, 1h TTL) and reads (0.10x).
batch_cached = (cache_write_tok * in_rate * 2.00
+ cache_read_tok * in_rate * 0.10
+ dynamic_tokens * in_rate
+ output_tokens * out_rate) * 0.50
batch_nocache = (raw_input_tok * in_rate + output_tokens * out_rate) * 0.50
live_cost = raw_input_tok * in_rate + output_tokens * out_rate
n_reread = sum(prefix_count.values()) - len(prefix_count)
print(f"[BATCH] {len(queue)} requests queued for {est_model}", flush=True)
print(f"[BATCH] cached prefix: {cache_write_tok:,} tokens (written once, re-read by {n_reread:,} requests)", flush=True)
print(f"[BATCH] dynamic input: {dynamic_tokens:,} tokens | estimated output: {output_tokens:,} tokens", flush=True)
print(f"[BATCH] estimated cost: ${batch_cached:.2f} (batch + prompt cache)", flush=True)
print(f"[BATCH] ${batch_nocache:.2f} (batch, worst case no cache hits)", flush=True)
print(f"[BATCH] ${live_cost:.2f} (live API, no batch discount)", flush=True)
return {
"requests": len(queue),
"model": est_model,
"cache_write_tokens": cache_write_tok,
"cache_read_tokens": cache_read_tok,
"dynamic_tokens": dynamic_tokens,
"output_tokens": output_tokens,
"batch_cached_cost": batch_cached,
"batch_nocache_cost": batch_nocache,
"live_cost": live_cost,
}
def submitTranslationBatches():
"""Submit the queued requests to the Anthropic Message Batches API.
Splits at the API limits and saves the custom_id -> cache-key mapping so
fetchTranslationBatches can route results back. Returns the batch ids."""
flush_batch_queue()
with _batch_file_lock():
queue = _read_batch_file(BATCH_QUEUE_FILE)
if not queue:
print("[BATCH] No batch requests queued.", flush=True)
return []
client = _get_anthropic_client()
batches = []
requests, id_map, size = [], {}, 0
def _submit():
nonlocal requests, id_map, size
if not requests:
return
batch = client.messages.batches.create(requests=requests)
batches.append({"id": batch.id, "custom_ids": id_map})
print(f"[BATCH] submitted {batch.id} ({len(requests)} requests)", flush=True)
requests, id_map, size = [], {}, 0
for i, (key, entry) in enumerate(queue.items()):
custom_id = f"req-{i:06d}"
params = entry["params"]
requests.append({"custom_id": custom_id, "params": params})
id_map[custom_id] = key
size += len(json.dumps(params, ensure_ascii=False).encode("utf-8"))
if len(requests) >= BATCH_MAX_REQUESTS or size >= BATCH_MAX_BYTES:
_submit()
_submit()
with BATCH_LOCK:
with _batch_file_lock():
_write_batch_file(BATCH_STATE_FILE, {"batches": batches})
return [b["id"] for b in batches]
def checkTranslationBatches():
"""Print the processing status of submitted batches. True when all ended."""
with _batch_file_lock():
state = _read_batch_file(BATCH_STATE_FILE)
if not state.get("batches"):
print("[BATCH] No submitted batches — submit the queue first.", flush=True)
return False
client = _get_anthropic_client()
all_ended = True
for info in state["batches"]:
b = client.messages.batches.retrieve(info["id"])
counts = getattr(b, "request_counts", None)
suffix = f" counts: {counts}" if counts else ""
print(f"[BATCH] {time.strftime('%H:%M:%S')} {b.id}: {b.processing_status}{suffix}", flush=True)
if b.processing_status != "ended":
all_ended = False
return all_ended
def fetchTranslationBatches():
"""Download finished batch results into the local results store.
Successes are stored keyed by the payload cache key for the consume pass;
errored/expired requests are reported and simply fall back to the live API
during consume. Returns (succeeded, errored) counts."""
global _batch_results
with _batch_file_lock():
state = _read_batch_file(BATCH_STATE_FILE)
if not state.get("batches"):
print("[BATCH] No submitted batches — nothing to fetch.", flush=True)
return 0, 0
client = _get_anthropic_client()
results, errored = {}, []
for info in state["batches"]:
id_map = info.get("custom_ids", {})
for r in client.messages.batches.results(info["id"]):
key = id_map.get(r.custom_id)
if key is None:
continue
res = r.result
if res.type != "succeeded":
detail = res.type
err = getattr(res, "error", None)
if err is not None:
inner = getattr(err, "error", err)
detail = f"{res.type} | {getattr(inner, 'type', '')}: {str(getattr(inner, 'message', '') or err)[:200]}"
errored.append((r.custom_id, detail))
continue
msg = res.message
text = "".join(getattr(b, "text", "") or "" for b in msg.content)
u = msg.usage
cr = getattr(u, "cache_read_input_tokens", 0) or 0
cw = getattr(u, "cache_creation_input_tokens", 0) or 0
inp = getattr(u, "input_tokens", 0) or 0
out = getattr(u, "output_tokens", 0) or 0
results[key] = {
"text": text,
# prompt_tokens matches _AnthropicCompat: total incl. cache fields.
"prompt_tokens": inp + cr + cw,
"completion_tokens": out,
"cache_read_input_tokens": cr,
"cache_creation_input_tokens": cw,
}
with BATCH_LOCK:
with _batch_file_lock():
merged = _read_batch_file(BATCH_RESULTS_FILE)
merged.update(results)
_write_batch_file(BATCH_RESULTS_FILE, merged)
# Queue and state are consumed; only the results store remains.
for path in (BATCH_QUEUE_FILE, BATCH_STATE_FILE):
try:
if path.exists():
path.unlink()
except Exception:
pass
_batch_results = None
print(f"[BATCH] fetched {len(results)} results ({len(errored)} errored).", flush=True)
for cid, why in errored[:20]:
print(f"[BATCH] ! {cid}: {why}", flush=True)
if len(errored) > 20:
print(f"[BATCH] ... ({len(errored) - 20} more)", flush=True)
return len(results), len(errored)
def runTranslationBatches(poll=60):
"""Submit the queue (unless already submitted), poll to completion, fetch."""
with _batch_file_lock():
state = _read_batch_file(BATCH_STATE_FILE)
if not state.get("batches"):
if not submitTranslationBatches():
return 0, 0
print(f"[BATCH] polling every {poll}s (Ctrl-C is safe — resume later with fetchTranslationBatches)...", flush=True)
while not checkTranslationBatches():
time.sleep(poll)
return fetchTranslationBatches()
class TranslationConfig:
"""Configuration class to hold all translation settings"""
@ -1103,6 +1549,85 @@ def createTranslationSchema(numLines):
}
class _AnthropicCompat:
"""OpenAI-shaped wrapper around an Anthropic response (text + usage)."""
class _Usage:
def __init__(self, prompt, completion, cr, cw):
self.prompt_tokens = prompt
self.completion_tokens = completion
self.cache_read_input_tokens = cr
self.cache_creation_input_tokens = cw
@property
def model_extra(self):
return {
"cache_read_input_tokens": self.cache_read_input_tokens,
"cache_creation_input_tokens": self.cache_creation_input_tokens,
}
class _Choice:
class _Msg:
def __init__(self, c): self.content = c
def __init__(self, c):
self.message = _AnthropicCompat._Choice._Msg(c)
def __init__(self, text, prompt, output, cr, cw):
self.choices = [_AnthropicCompat._Choice(text)]
self.usage = _AnthropicCompat._Usage(prompt, output, cr, cw)
def buildClaudeRequest(system, user, history, formatType, model, numLines=None, vocab_text=""):
"""Build the native Anthropic request kwargs.
Shared by live calls (translateText) and batch collection (translateAI) so
both produce byte-identical requests and share the same prompt cache. Only
the static prompt is cached 1h TTL so async batch requests processed
minutes apart still hit it; vocab and history ride in messages so they
never bust the prefix cache (the entire system parameter is the cache key).
"""
if DISABLE_CACHE:
# No cache_control — sends as a plain content block for a real uncached run.
combined_system = system + vocab_text
ant_system = [{"type": "text", "text": f"```\n{combined_system}\n```"}]
else:
ant_system = [{"type": "text", "text": f"```\n{system}\n```", "cache_control": {"type": "ephemeral", "ttl": "1h"}}]
native_msgs = [{"role": "user", "content": f"```\n{user}\n```"}]
# Vocab goes into messages as a user turn so it doesn't bust the prefix cache.
if vocab_text and vocab_text.strip():
native_msgs.insert(0, {"role": "user", "content": vocab_text.strip()})
native_msgs.insert(1, {"role": "assistant", "content": "Understood."})
# History also goes into messages, NOT ant_system.
if isinstance(history, list):
history_items = [str(h) for h in history if h and str(h).strip()]
else:
history_items = [str(history)] if history and str(history).strip() else []
if history_items:
history_block = "Translation History:\n```\n" + "\n".join(history_items) + "\n```"
native_msgs.insert(0, {"role": "user", "content": history_block})
native_msgs.insert(1, {"role": "assistant", "content": "Understood."})
ant_kwargs = dict(
model=model,
max_tokens=16384,
system=ant_system,
messages=native_msgs,
)
if formatType == "json" and numLines is not None:
ant_kwargs["output_config"] = {
"format": {
"type": "json_schema",
"schema": createTranslationSchema(numLines),
}
}
elif not _NO_SAMPLING_RE.search(model or ""):
# Plain completions still allow explicit sampling params (except on
# Opus 4.7+ which rejects them with a 400).
ant_kwargs["temperature"] = 0
# Do not pass temperature with output_config: newer Claude (e.g. Opus 4.7)
# returns errors such as "temperature is not supported" for structured outputs.
return ant_kwargs
def translateText(system, user, history, penalty, formatType, model, numLines=None, vocab_text=""):
"""Send translation request to the selected API.
@ -1239,59 +1764,12 @@ def translateText(system, user, history, penalty, formatType, model, numLines=No
# Use native Anthropic SDK — the OpenAI compat endpoint strips cache_control
# and never returns cache_read/creation_input_tokens.
if _is_claude:
# system blocks were built above (with cache_control on static prompt).
ant_system = list(msg[0]["content"])
# Convert remaining messages to native format.
# Skip system wrapper lines; collect assistant history and user content.
history_items = []
native_msgs = []
for m in msg[1:]:
role = m["role"]
content = m.get("content", "")
if role == "system":
pass # skip "Translation History:\n```" / "```" wrapper lines
elif role == "assistant":
history_items.append(str(content))
elif role == "user":
native_msgs.append({"role": "user", "content": str(content)})
# Vocab goes into messages as a user turn so it doesn't bust the
# Anthropic prefix cache (the entire system parameter is the key).
if vocab_text and vocab_text.strip():
native_msgs.insert(0, {"role": "user", "content": vocab_text.strip()})
native_msgs.insert(1, {"role": "assistant", "content": "Understood."})
# History also goes into messages, NOT ant_system.
if history_items:
history_block = "Translation History:\n```\n" + "\n".join(history_items) + "\n```"
native_msgs.insert(0, {"role": "user", "content": history_block})
native_msgs.insert(1, {"role": "assistant", "content": "Understood."})
if not native_msgs:
raise ValueError("No user message found for Anthropic native call")
# Built by the shared builder so live and batch requests are
# byte-identical and share the same prompt cache.
ant_kwargs = buildClaudeRequest(system, user, history, formatType, model, numLines, vocab_text=vocab_text)
ant_client = anthropic.Anthropic(api_key=openai.api_key)
try:
ant_kwargs = dict(
model=model,
max_tokens=16384,
system=ant_system,
messages=native_msgs,
)
if formatType == "json" and numLines is not None:
ant_kwargs["output_config"] = {
"format": {
"type": "json_schema",
"schema": createTranslationSchema(numLines),
}
}
else:
# Plain completions still allow explicit sampling params.
ant_kwargs["temperature"] = 0
# Do not pass temperature with output_config: newer Claude (e.g. Opus 4.7)
# returns errors such as "temperature is not supported" for structured outputs.
ant_resp = ant_client.messages.create(**ant_kwargs)
except Exception as e:
raise Exception(f"Anthropic API error: {e}")
@ -1306,28 +1784,6 @@ def translateText(system, user, history, penalty, formatType, model, numLines=No
# input_tokens (native SDK) = non-cached portion; add cache fields for true total.
_total_prompt = _inp + _cr + _cw
class _AnthropicCompat:
class _Usage:
def __init__(self, prompt, completion, cr, cw):
self.prompt_tokens = prompt
self.completion_tokens = completion
self.cache_read_input_tokens = cr
self.cache_creation_input_tokens = cw
@property
def model_extra(self):
return {
"cache_read_input_tokens": self.cache_read_input_tokens,
"cache_creation_input_tokens": self.cache_creation_input_tokens,
}
class _Choice:
class _Msg:
def __init__(self, c): self.content = c
def __init__(self, c):
self.message = _AnthropicCompat._Choice._Msg(c)
def __init__(self, text, prompt, output, cr, cw):
self.choices = [_AnthropicCompat._Choice(text)]
self.usage = _AnthropicCompat._Usage(prompt, output, cr, cw)
compat_response = _AnthropicCompat(_ant_text, _total_prompt, _out, _cr, _cw)
_write_request_debug_log("anthropic", ant_kwargs, compat_response.usage)
return compat_response
@ -1578,14 +2034,24 @@ def calculateCost(inputTokens, outputTokens, model):
cw = getattr(_thread_local, 'file_cache_write', 0)
reg = getattr(_thread_local, 'file_regular', 0)
out = getattr(_thread_local, 'file_output', 0)
bcr = getattr(_thread_local, 'file_batch_read', 0)
bcw = getattr(_thread_local, 'file_batch_write', 0)
breg = getattr(_thread_local, 'file_batch_regular', 0)
bout = getattr(_thread_local, 'file_batch_output', 0)
pricing = getPricingConfig(model)
br = pricing["inputAPICost"] / 1_000_000
orr = pricing["outputAPICost"] / 1_000_000
cost = cr * br * 0.10 + cw * br * 2.00 + reg * br + out * orr
cost = (cr * br * 0.10 + cw * br * 2.00 + reg * br + out * orr
# Batch API tokens: same rates, then the 50% batch discount.
+ (bcr * br * 0.10 + bcw * br * 2.00 + breg * br + bout * orr) * 0.50)
_thread_local.file_cache_read = 0
_thread_local.file_cache_write = 0
_thread_local.file_regular = 0
_thread_local.file_output = 0
_thread_local.file_batch_read = 0
_thread_local.file_batch_write = 0
_thread_local.file_batch_regular = 0
_thread_local.file_batch_output = 0
_thread_local.file_cost_ready = False
return cost
# TOTAL call (flag already cleared): return the cross-thread running total.
@ -1693,12 +2159,28 @@ def translateAI(text, history, config, filename=None, pbar=None, lock=None, mism
_thread_local.file_cache_write = 0
_thread_local.file_regular = 0
_thread_local.file_output = 0
# Batch API usage is billed at 50% so it is accumulated separately.
if not hasattr(_thread_local, 'file_batch_read'):
_thread_local.file_batch_read = 0
_thread_local.file_batch_write = 0
_thread_local.file_batch_regular = 0
_thread_local.file_batch_output = 0
# Snapshot accumulators so end-of-call delta only counts tokens from this call.
_prev_cr = _thread_local.file_cache_read
_prev_cw = _thread_local.file_cache_write
_prev_reg = _thread_local.file_regular
_prev_out = _thread_local.file_output
_prev_bcr = _thread_local.file_batch_read
_prev_bcw = _thread_local.file_batch_write
_prev_breg = _thread_local.file_batch_regular
_prev_bout = _thread_local.file_batch_output
_thread_local.file_cost_ready = False # will be set True at end of translateAI
# Anthropic batch phase ('collect'/'consume'); None when off, in estimate
# mode, or when the model doesn't route to the native Anthropic SDK.
batch_phase = get_batch_phase()
if batch_phase and (config.estimateMode or not isClaudeNative(config.model)):
batch_phase = None
if isinstance(text, list):
formatType = "json"
@ -1835,8 +2317,20 @@ def translateAI(text, history, config, filename=None, pbar=None, lock=None, mism
else:
subbedT = json.dumps({"Line1": protected_items}, indent=4, ensure_ascii=False)
# Check cache for this exact payload
cached_result = get_cached_translation(subbedT, config.language)
# Batch collect queues list payloads only. Single strings (speaker and
# variable names) translate live — modules memoize them and embed the
# results into later payloads, so they must resolve identically in both
# passes or the consume pass couldn't match the queued payload keys.
# This is the names-first phase; names are a tiny share of the volume.
queue_for_batch = batch_phase == "collect" and isinstance(tItem, list)
# Check cache for this exact payload (the collect pass uses a
# non-blocking peek so no pending markers are left behind for the
# consume pass to wait on)
if queue_for_batch:
cached_result = peek_cached_translation(subbedT, config.language)
else:
cached_result = get_cached_translation(subbedT, config.language)
if cached_result is not None:
# In estimate mode, never replace tList[index] from cache — the cached value
# may have been stored for a batch with a different number of skip_indices,
@ -1865,6 +2359,21 @@ def translateAI(text, history, config, filename=None, pbar=None, lock=None, mism
# vocab_text is the per-batch matched vocabulary (dynamic).
static_system, vocab_text, user = createContext(config, subbedT, formatType, history)
# Batch collect pass: queue the request (built exactly like a live one)
# instead of calling the API. The text stays untranslated; the consume
# pass fills it in from the fetched batch results. History carries the
# preceding source lines so the model still sees scene context.
if queue_for_batch:
numLines = len(clean_tItem) if isinstance(tItem, list) else 1
params = buildClaudeRequest(static_system, user, history, formatType,
config.model, numLines, vocab_text=vocab_text)
queue_batch_request(subbedT, config.language, params)
if lock and pbar is not None:
with lock:
pbar.update(len(tItem))
history = tItem[-config.maxHistory:]
continue
# Calculate estimate if in estimate mode
if config.estimateMode:
estimate = countTokens(static_system + vocab_text, user, history)
@ -1944,24 +2453,39 @@ def translateAI(text, history, config, filename=None, pbar=None, lock=None, mism
if pbar:
pbar.write(f"Retrying translation... (Attempt {attempt + 1}/{max_retries + 1})")
# Translate
try:
response = translateText(static_system, current_user, history, 0.05, formatType, config.model, numLines, vocab_text=vocab_text)
except Exception as api_err:
err_msg = f"[API_ERROR] {api_err}"
# Print to stdout so the GUI captures it immediately
print(err_msg, flush=True)
if pbar:
pbar.write(err_msg)
# Also write to the translation log file for persistence
# Translate — the consume pass tries the fetched batch result first;
# a missing or invalid result falls through to the live API.
from_batch = False
if batch_phase == "consume" and attempt == 0:
batch_result = take_batch_result(subbedT, config.language)
if batch_result is not None:
response = _AnthropicCompat(
batch_result.get("text", ""),
batch_result.get("prompt_tokens", 0) or 0,
batch_result.get("completion_tokens", 0) or 0,
batch_result.get("cache_read_input_tokens", 0) or 0,
batch_result.get("cache_creation_input_tokens", 0) or 0,
)
from_batch = True
_write_request_debug_log("anthropic-batch", {"payload": subbedT}, response.usage)
if not from_batch:
try:
Path(config.logFilePath).parent.mkdir(parents=True, exist_ok=True)
with open(config.logFilePath, "a", encoding="utf-8") as _lf:
_lf.write(f"{err_msg}\n")
_lf.flush()
except Exception:
pass
raise # Let retry decorator handle it
response = translateText(static_system, current_user, history, 0.05, formatType, config.model, numLines, vocab_text=vocab_text)
except Exception as api_err:
err_msg = f"[API_ERROR] {api_err}"
# Print to stdout so the GUI captures it immediately
print(err_msg, flush=True)
if pbar:
pbar.write(err_msg)
# Also write to the translation log file for persistence
try:
Path(config.logFilePath).parent.mkdir(parents=True, exist_ok=True)
with open(config.logFilePath, "a", encoding="utf-8") as _lf:
_lf.write(f"{err_msg}\n")
_lf.flush()
except Exception:
pass
raise # Let retry decorator handle it
translatedText = response.choices[0].message.content
last_raw_translation = translatedText
@ -1987,11 +2511,18 @@ def translateAI(text, history, config, filename=None, pbar=None, lock=None, mism
batch_regular = max(0, batch_prompt_total - batch_cache_read - batch_cache_write)
batch_output = getattr(usage, "completion_tokens", 0) or 0
# Accumulate into per-file thread-local counters.
_thread_local.file_cache_read += batch_cache_read
_thread_local.file_cache_write += batch_cache_write
_thread_local.file_regular += batch_regular
_thread_local.file_output += batch_output
# Accumulate into per-file thread-local counters. Batch API
# usage is billed at 50% so it goes into its own counters.
if from_batch:
_thread_local.file_batch_read += batch_cache_read
_thread_local.file_batch_write += batch_cache_write
_thread_local.file_batch_regular += batch_regular
_thread_local.file_batch_output += batch_output
else:
_thread_local.file_cache_read += batch_cache_read
_thread_local.file_cache_write += batch_cache_write
_thread_local.file_regular += batch_regular
_thread_local.file_output += batch_output
# --- Debug Token Logging ---
if DEBUG:
@ -2217,6 +2748,10 @@ def translateAI(text, history, config, filename=None, pbar=None, lock=None, mism
# Save cache after processing (for both estimate and translation modes)
save_cache()
# Batch collect pass: merge this call's queued requests into the disk queue.
if batch_phase == "collect":
flush_batch_queue()
# For Claude: accumulate only this call's delta into the cross-thread total.
# file_* accumulators hold full per-file totals; calculateCost() reads them.
_is_claude_final = config.model and any(x in config.model.lower() for x in ("claude", "sonnet", "haiku", "opus"))
@ -2229,11 +2764,20 @@ def translateAI(text, history, config, filename=None, pbar=None, lock=None, mism
_delta_cw = getattr(_thread_local, 'file_cache_write', 0) - _prev_cw
_delta_reg = getattr(_thread_local, 'file_regular', 0) - _prev_reg
_delta_out = getattr(_thread_local, 'file_output', 0) - _prev_out
_delta_bcr = getattr(_thread_local, 'file_batch_read', 0) - _prev_bcr
_delta_bcw = getattr(_thread_local, 'file_batch_write', 0) - _prev_bcw
_delta_breg = getattr(_thread_local, 'file_batch_regular', 0) - _prev_breg
_delta_bout = getattr(_thread_local, 'file_batch_output', 0) - _prev_bout
_call_cost = (
_delta_cr * _br * 0.10 +
_delta_cw * _br * 2.00 +
_delta_reg * _br +
_delta_out * _or
_delta_out * _or +
# Batch API tokens: same rates, then the 50% batch discount.
(_delta_bcr * _br * 0.10 +
_delta_bcw * _br * 2.00 +
_delta_breg * _br +
_delta_bout * _or) * 0.50
)
global _global_accurate_cost
with _global_accurate_cost_lock: