mirror of
https://github.com/turnstonelabs/turnstone.git
synced 2026-08-13 15:32:24 -06:00
Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| b3926c372a | |||
| 5efe52d433 |
+1
-1
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
||||
|
||||
[project]
|
||||
name = "turnstone"
|
||||
version = "1.0.0"
|
||||
version = "1.0.1"
|
||||
description = "Multi-node AI orchestration platform with tool use, agent routing, and cluster simulation."
|
||||
readme = "README.md"
|
||||
license = "BUSL-1.1"
|
||||
|
||||
@@ -41,6 +41,21 @@ class TestHashRingBuckets:
|
||||
# Empty list returns 0
|
||||
assert storage.assign_buckets([], "node-x") == 0
|
||||
|
||||
def test_assign_large_list_exceeds_chunk_size(self, storage):
|
||||
"""Regression: lists larger than chunk_size must not hit param limits."""
|
||||
n = 1200 # exceeds SQLite chunk_size (500) and exercises multi-chunk path
|
||||
storage.seed_ring_buckets([(i, "node-a") for i in range(n)])
|
||||
count = storage.assign_buckets(list(range(n)), "node-b")
|
||||
assert count == n
|
||||
rows = storage.list_ring_buckets()
|
||||
assert all(r["node_id"] == "node-b" for r in rows)
|
||||
|
||||
def test_assign_deduplicates_input(self, storage):
|
||||
"""Duplicates in the input list should not inflate rowcount."""
|
||||
storage.seed_ring_buckets([(0, "node-a"), (1, "node-a")])
|
||||
count = storage.assign_buckets([0, 1, 0, 1, 0], "node-b")
|
||||
assert count == 2
|
||||
|
||||
|
||||
class TestBucketStats:
|
||||
def test_increment_creates_row(self, storage):
|
||||
|
||||
@@ -1,3 +1,3 @@
|
||||
"""turnstone - Multi-node AI orchestration platform with tool use, agent routing, and cluster simulation."""
|
||||
|
||||
__version__ = "1.0.0"
|
||||
__version__ = "1.0.1"
|
||||
|
||||
@@ -216,13 +216,16 @@ class PostgreSQLBackend:
|
||||
).fetchall()
|
||||
orphan_ids = [r[0] for r in orphan_rows]
|
||||
if orphan_ids:
|
||||
conn.execute(
|
||||
sa.delete(workstream_config).where(workstream_config.c.ws_id.in_(orphan_ids))
|
||||
)
|
||||
result = conn.execute(
|
||||
sa.delete(workstreams).where(workstreams.c.ws_id.in_(orphan_ids))
|
||||
)
|
||||
orphans = result.rowcount
|
||||
chunk_size = 10_000
|
||||
for i in range(0, len(orphan_ids), chunk_size):
|
||||
chunk = orphan_ids[i : i + chunk_size]
|
||||
conn.execute(
|
||||
sa.delete(workstream_config).where(workstream_config.c.ws_id.in_(chunk))
|
||||
)
|
||||
result = conn.execute(
|
||||
sa.delete(workstreams).where(workstreams.c.ws_id.in_(chunk))
|
||||
)
|
||||
orphans += result.rowcount
|
||||
|
||||
# 2. Remove old unnamed workstreams
|
||||
if retention_days > 0:
|
||||
@@ -237,16 +240,19 @@ class PostgreSQLBackend:
|
||||
).fetchall()
|
||||
stale_ids = [r[0] for r in stale_rows]
|
||||
if stale_ids:
|
||||
conn.execute(
|
||||
sa.delete(conversations).where(conversations.c.ws_id.in_(stale_ids))
|
||||
)
|
||||
conn.execute(
|
||||
sa.delete(workstream_config).where(workstream_config.c.ws_id.in_(stale_ids))
|
||||
)
|
||||
result = conn.execute(
|
||||
sa.delete(workstreams).where(workstreams.c.ws_id.in_(stale_ids))
|
||||
)
|
||||
stale = result.rowcount
|
||||
chunk_size = 10_000
|
||||
for i in range(0, len(stale_ids), chunk_size):
|
||||
chunk = stale_ids[i : i + chunk_size]
|
||||
conn.execute(
|
||||
sa.delete(conversations).where(conversations.c.ws_id.in_(chunk))
|
||||
)
|
||||
conn.execute(
|
||||
sa.delete(workstream_config).where(workstream_config.c.ws_id.in_(chunk))
|
||||
)
|
||||
result = conn.execute(
|
||||
sa.delete(workstreams).where(workstreams.c.ws_id.in_(chunk))
|
||||
)
|
||||
stale += result.rowcount
|
||||
|
||||
conn.commit()
|
||||
return (orphans, stale)
|
||||
@@ -1254,14 +1260,22 @@ class PostgreSQLBackend:
|
||||
def assign_buckets(self, buckets: list[int], node_id: str) -> int:
|
||||
if not buckets:
|
||||
return 0
|
||||
# De-duplicate so rowcount stays accurate across chunks.
|
||||
buckets = list(dict.fromkeys(buckets))
|
||||
# psycopg limits query parameters to 65 535; chunk to stay well under.
|
||||
chunk_size = 10_000
|
||||
total = 0
|
||||
with self._engine.connect() as conn:
|
||||
result = conn.execute(
|
||||
sa.update(hash_ring_buckets)
|
||||
.where(hash_ring_buckets.c.bucket.in_(buckets))
|
||||
.values(node_id=node_id)
|
||||
)
|
||||
for i in range(0, len(buckets), chunk_size):
|
||||
chunk = buckets[i : i + chunk_size]
|
||||
result = conn.execute(
|
||||
sa.update(hash_ring_buckets)
|
||||
.where(hash_ring_buckets.c.bucket.in_(chunk))
|
||||
.values(node_id=node_id)
|
||||
)
|
||||
total += result.rowcount
|
||||
conn.commit()
|
||||
return result.rowcount
|
||||
return total
|
||||
|
||||
def increment_bucket_count(self, bucket: int, active: bool = False) -> None:
|
||||
from sqlalchemy.dialects.postgresql import insert as pg_insert
|
||||
@@ -1966,16 +1980,22 @@ class PostgreSQLBackend:
|
||||
def count_skill_resources_bulk(self, skill_ids: list[str]) -> dict[str, int]:
|
||||
if not skill_ids:
|
||||
return {}
|
||||
chunk_size = 10_000
|
||||
result: dict[str, int] = {}
|
||||
with self._engine.connect() as conn:
|
||||
rows = conn.execute(
|
||||
sa.select(
|
||||
skill_resources.c.skill_id,
|
||||
sa.func.count().label("cnt"),
|
||||
)
|
||||
.where(skill_resources.c.skill_id.in_(skill_ids))
|
||||
.group_by(skill_resources.c.skill_id)
|
||||
).fetchall()
|
||||
return {r[0]: r[1] for r in rows}
|
||||
for i in range(0, len(skill_ids), chunk_size):
|
||||
chunk = skill_ids[i : i + chunk_size]
|
||||
rows = conn.execute(
|
||||
sa.select(
|
||||
skill_resources.c.skill_id,
|
||||
sa.func.count().label("cnt"),
|
||||
)
|
||||
.where(skill_resources.c.skill_id.in_(chunk))
|
||||
.group_by(skill_resources.c.skill_id)
|
||||
).fetchall()
|
||||
for r in rows:
|
||||
result[r[0]] = r[1]
|
||||
return result
|
||||
|
||||
# -- Skill versions --------------------------------------------------------
|
||||
|
||||
|
||||
@@ -297,17 +297,20 @@ class SQLiteBackend:
|
||||
).fetchall()
|
||||
]
|
||||
if orphan_ids:
|
||||
placeholders = ",".join([":p" + str(i) for i in range(len(orphan_ids))])
|
||||
params = {f"p{i}": oid for i, oid in enumerate(orphan_ids)}
|
||||
conn.execute(
|
||||
sa.text(f"DELETE FROM workstream_config WHERE ws_id IN ({placeholders})"),
|
||||
params,
|
||||
)
|
||||
result = conn.execute(
|
||||
sa.text(f"DELETE FROM workstreams WHERE ws_id IN ({placeholders})"),
|
||||
params,
|
||||
)
|
||||
orphans = result.rowcount
|
||||
chunk_size = 500
|
||||
for i in range(0, len(orphan_ids), chunk_size):
|
||||
chunk = orphan_ids[i : i + chunk_size]
|
||||
placeholders = ",".join([":p" + str(j) for j in range(len(chunk))])
|
||||
params = {f"p{j}": oid for j, oid in enumerate(chunk)}
|
||||
conn.execute(
|
||||
sa.text(f"DELETE FROM workstream_config WHERE ws_id IN ({placeholders})"),
|
||||
params,
|
||||
)
|
||||
result = conn.execute(
|
||||
sa.text(f"DELETE FROM workstreams WHERE ws_id IN ({placeholders})"),
|
||||
params,
|
||||
)
|
||||
orphans += result.rowcount
|
||||
|
||||
# 2. Remove old unnamed workstreams
|
||||
if retention_days > 0:
|
||||
@@ -325,21 +328,26 @@ class SQLiteBackend:
|
||||
).fetchall()
|
||||
]
|
||||
if stale_ids:
|
||||
placeholders = ",".join([":p" + str(i) for i in range(len(stale_ids))])
|
||||
params = {f"p{i}": sid for i, sid in enumerate(stale_ids)}
|
||||
conn.execute(
|
||||
sa.text(f"DELETE FROM workstream_config WHERE ws_id IN ({placeholders})"),
|
||||
params,
|
||||
)
|
||||
conn.execute(
|
||||
sa.text(f"DELETE FROM conversations WHERE ws_id IN ({placeholders})"),
|
||||
params,
|
||||
)
|
||||
result = conn.execute(
|
||||
sa.text(f"DELETE FROM workstreams WHERE ws_id IN ({placeholders})"),
|
||||
params,
|
||||
)
|
||||
stale = result.rowcount
|
||||
chunk_size = 500
|
||||
for i in range(0, len(stale_ids), chunk_size):
|
||||
chunk = stale_ids[i : i + chunk_size]
|
||||
placeholders = ",".join([":p" + str(j) for j in range(len(chunk))])
|
||||
params = {f"p{j}": sid for j, sid in enumerate(chunk)}
|
||||
conn.execute(
|
||||
sa.text(
|
||||
f"DELETE FROM workstream_config WHERE ws_id IN ({placeholders})"
|
||||
),
|
||||
params,
|
||||
)
|
||||
conn.execute(
|
||||
sa.text(f"DELETE FROM conversations WHERE ws_id IN ({placeholders})"),
|
||||
params,
|
||||
)
|
||||
result = conn.execute(
|
||||
sa.text(f"DELETE FROM workstreams WHERE ws_id IN ({placeholders})"),
|
||||
params,
|
||||
)
|
||||
stale += result.rowcount
|
||||
|
||||
conn.commit()
|
||||
return (orphans, stale)
|
||||
@@ -1329,14 +1337,22 @@ class SQLiteBackend:
|
||||
def assign_buckets(self, buckets: list[int], node_id: str) -> int:
|
||||
if not buckets:
|
||||
return 0
|
||||
# De-duplicate so rowcount stays accurate across chunks.
|
||||
buckets = list(dict.fromkeys(buckets))
|
||||
# SQLite default SQLITE_MAX_VARIABLE_NUMBER is 999; chunk conservatively.
|
||||
chunk_size = 500
|
||||
total = 0
|
||||
with self._engine.connect() as conn:
|
||||
result = conn.execute(
|
||||
sa.update(hash_ring_buckets)
|
||||
.where(hash_ring_buckets.c.bucket.in_(buckets))
|
||||
.values(node_id=node_id)
|
||||
)
|
||||
for i in range(0, len(buckets), chunk_size):
|
||||
chunk = buckets[i : i + chunk_size]
|
||||
result = conn.execute(
|
||||
sa.update(hash_ring_buckets)
|
||||
.where(hash_ring_buckets.c.bucket.in_(chunk))
|
||||
.values(node_id=node_id)
|
||||
)
|
||||
total += result.rowcount
|
||||
conn.commit()
|
||||
return result.rowcount
|
||||
return total
|
||||
|
||||
def increment_bucket_count(self, bucket: int, active: bool = False) -> None:
|
||||
from sqlalchemy.dialects.sqlite import insert as sqlite_insert
|
||||
|
||||
Reference in New Issue
Block a user