Files
cleveragents-core/robot/helper_entity_sync.py
T
2026-05-29 07:21:08 -04:00

295 lines
9.1 KiB
Python

"""Helper script for entity_sync.robot integration tests.
Each subcommand is a self-contained check that prints a sentinel on success.
"""
from __future__ import annotations
import sys
from pathlib import Path
# Ensure local source tree is importable
_SRC = str(Path(__file__).resolve().parents[1] / "src")
if _SRC not in sys.path:
sys.path.insert(0, _SRC)
from cleveragents.a2a.facade import A2aLocalFacade # noqa: E402
from cleveragents.a2a.models import A2aRequest # noqa: E402
from cleveragents.a2a.sync_models import ( # noqa: E402
ConflictResolution,
SyncDirection,
SyncEntitySnapshot,
SyncEntityType,
SyncOperationStatus,
SyncPullRequest,
VectorClock,
)
from cleveragents.application.services.sync_service import SyncService # noqa: E402
# ---------------------------------------------------------------------------
# Subcommands
# ---------------------------------------------------------------------------
def sync_pull() -> None:
"""Test sync pull via facade."""
svc = SyncService(node_id="robot-client")
entity = SyncEntitySnapshot(
entity_id="robot-actor",
entity_type=SyncEntityType.ACTOR,
namespace="team",
name="robot-actor",
vector_clock=VectorClock(entries={"server": 1}),
)
svc.register_server_entity(entity)
facade = A2aLocalFacade(services={"sync_service": svc})
request = A2aRequest(
method="_cleveragents/sync/pull",
params={"namespace": "team"},
)
response = facade.dispatch(request)
result = response.result or {}
if response.error is None and result.get("total_pulled", 0) >= 1:
print("sync-pull-ok")
else:
print(f"FAIL: {response.error or result}", file=sys.stderr)
sys.exit(1)
def sync_push() -> None:
"""Test sync push via facade."""
svc = SyncService(node_id="robot-client")
facade = A2aLocalFacade(services={"sync_service": svc})
params = {
"namespace": "team",
"entities": [
{
"entity_id": "push-actor",
"entity_type": "actor",
"namespace": "team",
"name": "push-test",
},
],
}
request = A2aRequest(
method="_cleveragents/sync/push",
params=params,
)
response = facade.dispatch(request)
result = response.result or {}
if response.error is None and result.get("total_accepted", 0) >= 1:
print("sync-push-ok")
else:
print(f"FAIL: {response.error or result}", file=sys.stderr)
sys.exit(1)
def sync_status() -> None:
"""Test sync status via facade."""
svc = SyncService(node_id="robot-client")
entity = SyncEntitySnapshot(
entity_id="status-actor",
entity_type=SyncEntityType.ACTOR,
namespace="team",
name="status-actor",
vector_clock=VectorClock(entries={"client": 1}),
)
svc.register_local_entity(entity)
facade = A2aLocalFacade(services={"sync_service": svc})
request = A2aRequest(
method="_cleveragents/sync/status",
params={"namespace": "team"},
)
response = facade.dispatch(request)
result = response.result or {}
if response.error is None and result.get("drift_detected") is True:
print("sync-status-ok")
else:
print(f"FAIL: {response.error or result}", file=sys.stderr)
sys.exit(1)
def sync_conflict_resolution() -> None:
"""Test conflict detection and resolution."""
svc = SyncService(node_id="robot-client")
local_ent = SyncEntitySnapshot(
entity_id="conflict-actor",
entity_type=SyncEntityType.ACTOR,
namespace="team",
name="conflict-actor",
vector_clock=VectorClock(entries={"client": 2, "server": 1}),
updated_at="2025-06-01T12:00:00+00:00",
)
server_ent = SyncEntitySnapshot(
entity_id="conflict-actor",
entity_type=SyncEntityType.ACTOR,
namespace="team",
name="conflict-actor",
vector_clock=VectorClock(entries={"client": 1, "server": 2}),
updated_at="2025-05-01T12:00:00+00:00",
)
svc.register_local_entity(local_ent)
svc.register_server_entity(server_ent)
pull_req = SyncPullRequest(namespace="team")
pull_resp = svc.pull(pull_req)
if len(pull_resp.conflicts) != 1:
print(
f"FAIL: expected 1 conflict, got {len(pull_resp.conflicts)}",
file=sys.stderr,
)
sys.exit(1)
conflict_id = svc.conflicts[0].conflict_id
resolved = svc.resolve_conflict(
conflict_id=conflict_id,
resolution=ConflictResolution.LAST_WRITER_WINS,
)
if resolved.resolved and resolved.winner == "local":
print("sync-conflict-resolution-ok")
else:
print(
f"FAIL: resolved={resolved.resolved}, winner={resolved.winner}",
file=sys.stderr,
)
sys.exit(1)
def sync_offline_queue() -> None:
"""Test offline queue operations."""
svc = SyncService(node_id="robot-client")
entity = SyncEntitySnapshot(
entity_id="offline-actor",
entity_type=SyncEntityType.ACTOR,
namespace="team",
name="offline-actor",
vector_clock=VectorClock(entries={"client": 1}),
)
entry = svc.enqueue_offline(
direction=SyncDirection.PUSH,
namespace="team",
entity=entity,
)
if entry.status != SyncOperationStatus.QUEUED:
print(f"FAIL: expected QUEUED, got {entry.status}", file=sys.stderr)
sys.exit(1)
processed = svc.process_offline_queue()
if len(processed) >= 1 and len(svc.offline_queue) == 0:
print("sync-offline-queue-ok")
else:
print(
f"FAIL: processed={len(processed)}, remaining={len(svc.offline_queue)}",
file=sys.stderr,
)
sys.exit(1)
def sync_vector_clock() -> None:
"""Test vector clock operations."""
c1 = VectorClock(entries={"a": 1, "b": 2})
c2 = VectorClock(entries={"a": 2, "b": 3})
if not c1.happens_before(c2):
print("FAIL: c1 should happen before c2", file=sys.stderr)
sys.exit(1)
c3 = VectorClock(entries={"a": 2, "b": 1})
if not c1.is_concurrent(c3):
print("FAIL: c1 and c3 should be concurrent", file=sys.stderr)
sys.exit(1)
merged = c1.merge(c2)
if merged.entries != {"a": 2, "b": 3}:
print(f"FAIL: unexpected merge result {merged.entries}", file=sys.stderr)
sys.exit(1)
incremented = c1.increment("a")
if incremented.entries["a"] != 2:
print(f"FAIL: expected a=2, got {incremented.entries['a']}", file=sys.stderr)
sys.exit(1)
print("sync-vector-clock-ok")
def sync_facade_no_service() -> None:
"""Test facade sync stubs when no service is registered."""
facade = A2aLocalFacade()
for op in (
"_cleveragents/sync/pull",
"_cleveragents/sync/push",
"_cleveragents/sync/status",
):
request = A2aRequest(method=op, params={"namespace": "test"})
response = facade.dispatch(request)
result = response.result or {}
if response.error is not None or not result.get("stub"):
print(
f"FAIL: {op} should return stub, got {response.error or result}",
file=sys.stderr,
)
sys.exit(1)
print("sync-facade-no-service-ok")
def sync_incremental() -> None:
"""Test incremental sync with since parameter."""
svc = SyncService(node_id="robot-client")
old_entity = SyncEntitySnapshot(
entity_id="old-actor",
entity_type=SyncEntityType.ACTOR,
namespace="team",
name="old-actor",
vector_clock=VectorClock(entries={"server": 1}),
updated_at="2024-01-01T00:00:00+00:00",
)
new_entity = SyncEntitySnapshot(
entity_id="new-actor",
entity_type=SyncEntityType.ACTOR,
namespace="team",
name="new-actor",
vector_clock=VectorClock(entries={"server": 2}),
updated_at="2025-06-01T00:00:00+00:00",
)
svc.register_server_entity(old_entity)
svc.register_server_entity(new_entity)
pull_req = SyncPullRequest(
namespace="team",
since="2025-01-01T00:00:00+00:00",
)
resp = svc.pull(pull_req)
if resp.total_pulled == 1 and resp.entities[0].entity_id == "new-actor":
print("sync-incremental-ok")
else:
print(f"FAIL: expected 1 entity, got {resp.total_pulled}", file=sys.stderr)
sys.exit(1)
# ---------------------------------------------------------------------------
# CLI dispatch
# ---------------------------------------------------------------------------
_COMMANDS: dict[str, object] = {
"sync-pull": sync_pull,
"sync-push": sync_push,
"sync-status": sync_status,
"sync-conflict-resolution": sync_conflict_resolution,
"sync-offline-queue": sync_offline_queue,
"sync-vector-clock": sync_vector_clock,
"sync-facade-no-service": sync_facade_no_service,
"sync-incremental": sync_incremental,
}
if __name__ == "__main__":
if len(sys.argv) < 2 or sys.argv[1] not in _COMMANDS:
print(f"Usage: {sys.argv[0]} <{'|'.join(sorted(_COMMANDS))}>>", file=sys.stderr)
sys.exit(2)
cmd = _COMMANDS[sys.argv[1]]
if callable(cmd):
cmd()