Compare commits

...

2 Commits

Author SHA1 Message Date
Patrick Buckley b3926c372a chore: bump version to 1.0.1 2026-04-02 20:29:02 -07:00
Patrick Buckley 5efe52d433 fix: chunk IN clauses to stay within DB parameter limits (#286)
* fix: chunk IN clauses to stay within DB parameter limits

psycopg caps query parameters at 65 535 and SQLite defaults to 999.
assign_buckets, prune_workstreams, and count_skill_resources_bulk were
passing unbounded lists into single IN(...) clauses, causing
OperationalError during rebalancer runs on full-size hash rings.

Chunk sizes: 10 000 (PostgreSQL), 500 (SQLite).

* fix: deduplicate assign_buckets input, add chunking regression tests

Address review feedback: deduplicate bucket list before chunking to
prevent inflated rowcount from cross-chunk duplicates. Add tests that
exercise the multi-chunk path (1200 buckets > SQLite chunk_size of 500)
and verify dedup preserves accurate counts.
2026-04-02 20:28:54 -07:00
6 changed files with 118 additions and 67 deletions
+1 -1
View File
@@ -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"
+15
View File
@@ -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 -1
View File
@@ -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"
+52 -32
View File
@@ -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 --------------------------------------------------------
+48 -32
View File
@@ -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
Generated
+1 -1
View File
@@ -2496,7 +2496,7 @@ wheels = [
[[package]]
name = "turnstone"
version = "1.0.0"
version = "1.0.1"
source = { editable = "." }
dependencies = [
{ name = "alembic" },