diff --git a/agentscope-extensions/agentscope-extensions-jdbc/src/main/java/io/agentscope/extensions/jdbc/state/JdbcAgentStateStore.java b/agentscope-extensions/agentscope-extensions-jdbc/src/main/java/io/agentscope/extensions/jdbc/state/JdbcAgentStateStore.java index daa94a4e37..a5b0c808bb 100644 --- a/agentscope-extensions/agentscope-extensions-jdbc/src/main/java/io/agentscope/extensions/jdbc/state/JdbcAgentStateStore.java +++ b/agentscope-extensions/agentscope-extensions-jdbc/src/main/java/io/agentscope/extensions/jdbc/state/JdbcAgentStateStore.java @@ -233,7 +233,7 @@ public long saveIfVersion( String userId, String sessionId, String key, State value, long expectedVersion) { if (expectedVersion == UNVERSIONED) { save(userId, sessionId, key, value); - return getVersioned(userId, sessionId, key, State.class).version(); + return readVersionOnly(userId, sessionId, key); } String slotId = slotId(userId, sessionId); validateSlotId(slotId); @@ -276,6 +276,26 @@ public long saveIfVersion( } } + private long readVersionOnly(String userId, String sessionId, String key) { + String slotId = slotId(userId, sessionId); + validateSlotId(slotId); + validateStateKey(key); + + BoundSql boundSql = dialect.sessionStateSelectVersioned(slotId, key, SINGLE_STATE_INDEX); + try (Connection conn = dataSource.getConnection(); + PreparedStatement stmt = conn.prepareStatement(boundSql.sql())) { + bindParams(stmt, boundSql.params()); + try (ResultSet rs = stmt.executeQuery()) { + if (!rs.next()) { + return 0L; + } + return rs.getLong("version"); + } + } catch (Exception e) { + throw new RuntimeException("Failed to read version for state: " + key, e); + } + } + @Override public Optional get( String userId, String sessionId, String key, Class type) { diff --git a/agentscope-extensions/agentscope-extensions-mysql/src/main/java/io/agentscope/extensions/mysql/state/MysqlAgentStateStore.java b/agentscope-extensions/agentscope-extensions-mysql/src/main/java/io/agentscope/extensions/mysql/state/MysqlAgentStateStore.java index 6580e077fe..25a0a8f3c8 100644 --- a/agentscope-extensions/agentscope-extensions-mysql/src/main/java/io/agentscope/extensions/mysql/state/MysqlAgentStateStore.java +++ b/agentscope-extensions/agentscope-extensions-mysql/src/main/java/io/agentscope/extensions/mysql/state/MysqlAgentStateStore.java @@ -425,7 +425,7 @@ public long saveIfVersion( String userId, String sessionId, String key, State value, long expectedVersion) { if (expectedVersion == UNVERSIONED) { save(userId, sessionId, key, value); - return getVersioned(userId, sessionId, key, State.class).version(); + return readVersionOnly(userId, sessionId, key); } String slotId = slotId(userId, sessionId); @@ -449,6 +449,29 @@ public long saveIfVersion( } } + private long readVersionOnly(String userId, String sessionId, String key) { + String slotId = slotId(userId, sessionId); + validateSessionId(slotId); + validateStateKey(key); + + String sql = "SELECT version FROM " + getFullTableName() + + " WHERE session_id = ? AND state_key = ? AND item_index = 0"; + + try (Connection conn = dataSource.getConnection(); + PreparedStatement stmt = conn.prepareStatement(sql)) { + stmt.setString(1, slotId); + stmt.setString(2, key); + try (ResultSet rs = stmt.executeQuery()) { + if (!rs.next()) { + return 0L; + } + return rs.getLong("version"); + } + } catch (Exception e) { + throw new RuntimeException("Failed to read version for state: " + key, e); + } + } + private long insertIfAbsent(Connection conn, String slotId, String key, State value) throws Exception { String insertSql = diff --git a/agentscope-extensions/agentscope-extensions-postgresql/src/main/java/io/agentscope/extensions/postgresql/state/PostgresAgentStateStore.java b/agentscope-extensions/agentscope-extensions-postgresql/src/main/java/io/agentscope/extensions/postgresql/state/PostgresAgentStateStore.java index 24323f130b..cf32c9c397 100644 --- a/agentscope-extensions/agentscope-extensions-postgresql/src/main/java/io/agentscope/extensions/postgresql/state/PostgresAgentStateStore.java +++ b/agentscope-extensions/agentscope-extensions-postgresql/src/main/java/io/agentscope/extensions/postgresql/state/PostgresAgentStateStore.java @@ -330,7 +330,7 @@ public long saveIfVersion( String userId, String sessionId, String key, State value, long expectedVersion) { if (expectedVersion == UNVERSIONED) { save(userId, sessionId, key, value); - return getVersioned(userId, sessionId, key, State.class).version(); + return readVersionOnly(userId, sessionId, key); } String slotId = slotId(userId, sessionId); @@ -354,6 +354,29 @@ public long saveIfVersion( } } + private long readVersionOnly(String userId, String sessionId, String key) { + String slotId = slotId(userId, sessionId); + validateSessionId(slotId); + validateStateKey(key); + + String sql = "SELECT version FROM " + getFullTableName() + + " WHERE session_id = ? AND state_key = ? AND item_index = 0"; + + try (Connection conn = dataSource.getConnection(); + PreparedStatement stmt = conn.prepareStatement(sql)) { + stmt.setString(1, slotId); + stmt.setString(2, key); + try (ResultSet rs = stmt.executeQuery()) { + if (!rs.next()) { + return 0L; + } + return rs.getLong("version"); + } + } catch (Exception e) { + throw new RuntimeException("Failed to read version for state: " + key, e); + } + } + private long insertIfAbsent(Connection conn, String slotId, String key, State value) throws Exception { String insertSql = diff --git a/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/RedisAgentStateStore.java b/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/RedisAgentStateStore.java index 7679fd8b5e..a7c43303b8 100644 --- a/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/RedisAgentStateStore.java +++ b/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/RedisAgentStateStore.java @@ -270,8 +270,7 @@ public long saveIfVersion( String userId, String sessionId, String key, State value, long expectedVersion) { if (expectedVersion == UNVERSIONED) { save(userId, sessionId, key, value); - VersionedState after = getVersioned(userId, sessionId, key, State.class); - return after.version(); + return readVersionOnly(userId, sessionId, key); } String slotId = slotId(userId, sessionId); String redisKey = getStateKey(slotId, key); @@ -290,6 +289,25 @@ public long saveIfVersion( } } + /** + * Read only the version number without deserializing the payload. + * This avoids Jackson's inability to deserialize the State marker interface. + */ + private long readVersionOnly(String userId, String sessionId, String key) { + String slotId = slotId(userId, sessionId); + String redisKey = getStateKey(slotId, key); + String versionKey = RedisStateVersionSupport.versionKey(redisKey); + try { + String json = client.get(redisKey); + if (json == null) { + return 0L; + } + return RedisStateVersionSupport.parseVersion(json, client.get(versionKey)); + } catch (Exception e) { + throw new RuntimeException("Failed to read version for state: " + key, e); + } + } + @Override public void save(String userId, String sessionId, String key, List values) { String slotId = slotId(userId, sessionId); diff --git a/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/jedis/JedisAgentStateStore.java b/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/jedis/JedisAgentStateStore.java index 0513656db5..458f326124 100644 --- a/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/jedis/JedisAgentStateStore.java +++ b/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/jedis/JedisAgentStateStore.java @@ -119,11 +119,26 @@ public long saveIfVersion( String userId, String sessionId, String key, State value, long expectedVersion) { if (expectedVersion == UNVERSIONED) { save(userId, sessionId, key, value); - return getVersioned(userId, sessionId, key, State.class).version(); + return readVersionOnly(userId, sessionId, key); } return evalSave(userId, sessionId, key, value, Long.toString(expectedVersion)); } + private long readVersionOnly(String userId, String sessionId, String key) { + String slotId = slotId(userId, sessionId); + String redisKey = getStateKey(slotId, key); + String versionKey = RedisStateVersionSupport.versionKey(redisKey); + try (Jedis jedis = jedisPool.getResource()) { + String json = jedis.get(redisKey); + if (json == null) { + return 0L; + } + return RedisStateVersionSupport.parseVersion(json, jedis.get(versionKey)); + } catch (Exception e) { + throw new RuntimeException("Failed to read version for state: " + key, e); + } + } + private long evalSave( String userId, String sessionId, String key, State value, String expectedVersionArg) { String slotId = slotId(userId, sessionId); diff --git a/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/redisson/RedissonAgentStateStore.java b/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/redisson/RedissonAgentStateStore.java index f497da3818..b123cc75b5 100644 --- a/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/redisson/RedissonAgentStateStore.java +++ b/agentscope-extensions/agentscope-extensions-redis/src/main/java/io/agentscope/extensions/redis/state/redisson/RedissonAgentStateStore.java @@ -125,11 +125,28 @@ public long saveIfVersion( String userId, String sessionId, String key, State value, long expectedVersion) { if (expectedVersion == UNVERSIONED) { save(userId, sessionId, key, value); - return getVersioned(userId, sessionId, key, State.class).version(); + return readVersionOnly(userId, sessionId, key); } return evalSave(userId, sessionId, key, value, Long.toString(expectedVersion)); } + private long readVersionOnly(String userId, String sessionId, String key) { + String slotId = slotId(userId, sessionId); + String redisKey = getStateKey(slotId, key); + String versionKey = RedisStateVersionSupport.versionKey(redisKey); + try { + RBucket bucket = redissonClient.getBucket(redisKey, StringCodec.INSTANCE); + String json = bucket.get(); + if (json == null) { + return 0L; + } + RBucket versionBucket = redissonClient.getBucket(versionKey, StringCodec.INSTANCE); + return RedisStateVersionSupport.parseVersion(json, versionBucket.get()); + } catch (Exception e) { + throw new RuntimeException("Failed to read version for state: " + key, e); + } + } + private long evalSave( String userId, String sessionId, String key, State value, String expectedVersionArg) { String slotId = slotId(userId, sessionId);