fix(file-sync): serialize concurrent sync cycles
This commit is contained in:
@@ -1,8 +1,10 @@
|
||||
"""Tests for FileSyncManager — mtime tracking, deletion detection, transactional rollback."""
|
||||
|
||||
import concurrent.futures
|
||||
import io
|
||||
import os
|
||||
import tarfile
|
||||
import threading
|
||||
import time
|
||||
from pathlib import Path
|
||||
from unittest.mock import MagicMock, patch
|
||||
@@ -261,6 +263,59 @@ class TestEdgeCases:
|
||||
upload.assert_not_called() # _file_mtime_key returns None, skipped
|
||||
|
||||
|
||||
class TestConcurrency:
|
||||
def test_sync_back_waits_for_active_sync_transaction(self, tmp_path):
|
||||
initial_file = tmp_path / "initial.png"
|
||||
new_file = tmp_path / "new.png"
|
||||
initial_file.write_bytes(b"initial")
|
||||
upload_started = threading.Event()
|
||||
release_upload = threading.Event()
|
||||
sync_back_transport_started = threading.Event()
|
||||
overlap_detected = threading.Event()
|
||||
download_calls = []
|
||||
|
||||
def get_files():
|
||||
return [
|
||||
(str(path), f"/root/.hermes/cache/images/{path.name}")
|
||||
for path in sorted(tmp_path.glob("*.png"))
|
||||
]
|
||||
|
||||
def upload(host_path, _remote_path):
|
||||
if host_path == str(new_file):
|
||||
upload_started.set()
|
||||
sync_back_transport_started.wait(timeout=1.0)
|
||||
release_upload.set()
|
||||
|
||||
def bulk_download(destination):
|
||||
if not release_upload.is_set():
|
||||
overlap_detected.set()
|
||||
sync_back_transport_started.set()
|
||||
download_calls.append(destination)
|
||||
with tarfile.open(destination, "w"):
|
||||
pass
|
||||
|
||||
mgr = FileSyncManager(
|
||||
get_files_fn=get_files,
|
||||
upload_fn=upload,
|
||||
delete_fn=MagicMock(),
|
||||
bulk_download_fn=bulk_download,
|
||||
)
|
||||
mgr.sync(force=True)
|
||||
new_file.write_bytes(b"new")
|
||||
|
||||
with concurrent.futures.ThreadPoolExecutor(max_workers=2) as executor:
|
||||
sync_future = executor.submit(mgr.sync, force=True)
|
||||
assert upload_started.wait(timeout=2.0)
|
||||
|
||||
sync_back_future = executor.submit(mgr.sync_back, hermes_home=tmp_path)
|
||||
|
||||
sync_future.result(timeout=3.0)
|
||||
sync_back_future.result(timeout=3.0)
|
||||
|
||||
assert len(download_calls) == 1
|
||||
assert not overlap_detected.is_set()
|
||||
|
||||
|
||||
class TestSyncBackSecurity:
|
||||
def test_sync_back_does_not_overwrite_uploaded_credential_files(self, tmp_path, monkeypatch):
|
||||
credential = tmp_path / "token.json"
|
||||
|
||||
@@ -1,4 +1,6 @@
|
||||
import concurrent.futures
|
||||
import json
|
||||
import threading
|
||||
from types import SimpleNamespace
|
||||
|
||||
|
||||
@@ -38,6 +40,79 @@ def test_postprocess_adds_agent_visible_image_for_active_ssh_env(monkeypatch, tm
|
||||
assert sync_calls == [True]
|
||||
|
||||
|
||||
def test_concurrent_image_results_preserve_shared_remote_sync_state(monkeypatch, tmp_path):
|
||||
from tools import image_generation_tool
|
||||
from tools.environments import file_sync
|
||||
|
||||
hermes_home = tmp_path / ".hermes"
|
||||
image_dir = hermes_home / "cache" / "images"
|
||||
image_dir.mkdir(parents=True)
|
||||
first_image = image_dir / "first.png"
|
||||
second_image = image_dir / "second.png"
|
||||
first_image.write_bytes(b"first")
|
||||
|
||||
first_upload_started = threading.Event()
|
||||
second_sync_finished = threading.Event()
|
||||
worker = threading.local()
|
||||
|
||||
def get_files():
|
||||
return [
|
||||
(
|
||||
str(path),
|
||||
f"/home/remote/.hermes/cache/images/{path.name}",
|
||||
)
|
||||
for path in sorted(image_dir.iterdir())
|
||||
]
|
||||
|
||||
def upload(_host_path, _remote_path):
|
||||
if worker.label == "first":
|
||||
first_upload_started.set()
|
||||
# Without transaction serialization, the second sync commits its
|
||||
# newer snapshot before this first sync resumes and overwrites it.
|
||||
second_sync_finished.wait(timeout=1.0)
|
||||
|
||||
sync_manager = file_sync.FileSyncManager(
|
||||
get_files_fn=get_files,
|
||||
upload_fn=upload,
|
||||
delete_fn=lambda _paths: None,
|
||||
)
|
||||
env = SimpleNamespace(
|
||||
_remote_home="/home/remote",
|
||||
_sync_manager=sync_manager,
|
||||
)
|
||||
|
||||
monkeypatch.setenv("HERMES_HOME", str(hermes_home))
|
||||
monkeypatch.setattr(file_sync, "_credential_host_paths", lambda: set())
|
||||
monkeypatch.setattr(image_generation_tool, "_active_terminal_env", lambda _task_id: env)
|
||||
|
||||
def postprocess(label, image_path):
|
||||
worker.label = label
|
||||
try:
|
||||
raw = json.dumps({"success": True, "image": str(image_path)})
|
||||
return image_generation_tool._postprocess_image_generate_result(
|
||||
raw,
|
||||
task_id="shared-task",
|
||||
)
|
||||
finally:
|
||||
if label == "second":
|
||||
second_sync_finished.set()
|
||||
|
||||
with concurrent.futures.ThreadPoolExecutor(max_workers=2) as executor:
|
||||
first_future = executor.submit(postprocess, "first", first_image)
|
||||
assert first_upload_started.wait(timeout=2.0)
|
||||
|
||||
second_image.write_bytes(b"second")
|
||||
second_future = executor.submit(postprocess, "second", second_image)
|
||||
|
||||
first_future.result(timeout=3.0)
|
||||
second_future.result(timeout=3.0)
|
||||
|
||||
assert set(sync_manager._synced_files) == {
|
||||
"/home/remote/.hermes/cache/images/first.png",
|
||||
"/home/remote/.hermes/cache/images/second.png",
|
||||
}
|
||||
|
||||
|
||||
def test_handle_image_generate_postprocesses_plugin_result(monkeypatch, tmp_path):
|
||||
from tools import image_generation_tool
|
||||
|
||||
|
||||
@@ -156,6 +156,7 @@ class FileSyncManager:
|
||||
self._bulk_upload_fn = bulk_upload_fn
|
||||
self._bulk_download_fn = bulk_download_fn
|
||||
self._delete_fn = delete_fn
|
||||
self._transaction_lock = threading.Lock()
|
||||
self._synced_files: dict[str, tuple[float, int]] = {} # remote_path -> (mtime, size)
|
||||
self._pushed_hashes: dict[str, str] = {} # remote_path -> sha256 hex digest
|
||||
self._upload_only_host_paths: set[str] = set()
|
||||
@@ -171,6 +172,11 @@ class FileSyncManager:
|
||||
Transactional: state only committed if ALL operations succeed.
|
||||
On failure, state rolls back so the next cycle retries everything.
|
||||
"""
|
||||
with self._transaction_lock:
|
||||
self._sync_transaction(force=force)
|
||||
|
||||
def _sync_transaction(self, *, force: bool = False) -> None:
|
||||
"""Execute one sync cycle while holding the per-manager lock."""
|
||||
if not force and not os.environ.get(_FORCE_SYNC_ENV):
|
||||
now = _monotonic()
|
||||
if now - self._last_sync_time < self._sync_interval:
|
||||
@@ -257,6 +263,11 @@ class FileSyncManager:
|
||||
Protected against SIGINT (defers the signal until complete) and
|
||||
serialized across concurrent gateway sandboxes via file lock.
|
||||
"""
|
||||
with self._transaction_lock:
|
||||
self._sync_back_transaction(hermes_home=hermes_home)
|
||||
|
||||
def _sync_back_transaction(self, hermes_home: Path | None = None) -> None:
|
||||
"""Execute sync-back against a stable snapshot of manager state."""
|
||||
if self._bulk_download_fn is None:
|
||||
return
|
||||
|
||||
|
||||
@@ -64,8 +64,13 @@ Your selection is saved to `config.yaml`:
|
||||
image_gen:
|
||||
model: fal-ai/flux-2/klein/9b
|
||||
use_gateway: false # true if using Nous Subscription
|
||||
max_parallel_requests: 4 # concurrent images in one tool-call batch
|
||||
```
|
||||
|
||||
`max_parallel_requests` defaults to `4`. Hermes clamps it to at least one and
|
||||
to the global tool-worker limit, so image providers receive bounded parallel
|
||||
requests without allowing an image batch to bypass the agent's concurrency cap.
|
||||
|
||||
### GPT-Image Quality
|
||||
|
||||
The `fal-ai/gpt-image-1.5` and `fal-ai/gpt-image-2` request quality is pinned to `medium` (~$0.034–$0.06/image at 1024×1024). We don't expose the `low` / `high` tiers as a user-facing option so that Nous Portal billing stays predictable across all users — the cost spread between tiers is 3–22×. If you want a cheaper option, pick Klein 9B or Z-Image Turbo; if you want higher quality, use Nano Banana Pro or Recraft V4 Pro.
|
||||
|
||||
Reference in New Issue
Block a user