From bfec3e4f2b386c370442a999d095e596891f9adf Mon Sep 17 00:00:00 2001
From: Sean-Walker0 <309789691+Sean-Walker0@users.noreply.github.com>
Date: Wed, 30 Sep 2026 12:56:14 +0800
Subject: [PATCH] refactor(graph): clarify thread id naming in checkpoint
savers
The persistence savers accept a user-facing thread id through
RunnableConfig but store it in a thread_name column, while the
thread_id column holds an internal surrogate UUID generated to
support thread release and reuse. This split was introduced by the
saver rework in spring-ai-alibaba#3287 and kept the old column
names, so thread_id now names the internal key instead of the
user-facing id, which makes the storage layer easy to misread.
Make the identity model explicit without touching stored data or
SQL: document the two identities in every saver and in
AbstractJdbcCheckpointSaver, annotate the schema comments with what
each column actually holds, and rename internal Java identifiers so
threadId always means the user-facing id and persistedThreadId
always means the internal surrogate UUID.
---
.../graph/checkpoint/savers/h2/H2Saver.java | 36 ++++--
.../jdbc/AbstractJdbcCheckpointSaver.java | 23 +++-
.../checkpoint/savers/mongo/MongoSaver.java | 108 +++++++++-------
.../checkpoint/savers/mysql/MysqlSaver.java | 64 +++++----
.../checkpoint/savers/oracle/OracleSaver.java | 64 +++++----
.../savers/postgresql/PostgresSaver.java | 14 +-
.../checkpoint/savers/redis/RedisSaver.java | 121 ++++++++++--------
7 files changed, 251 insertions(+), 179 deletions(-)
diff --git a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/h2/H2Saver.java b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/h2/H2Saver.java
index b6d51bbfe..8467c5ed4 100644
--- a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/h2/H2Saver.java
+++ b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/h2/H2Saver.java
@@ -60,8 +60,8 @@
*
*
* CREATE TABLE GRAPH_THREAD (
- * thread_id VARCHAR(36) PRIMARY KEY,
- * thread_name VARCHAR(255) NOT NULL,
+ * thread_id VARCHAR(36) PRIMARY KEY, -- internal surrogate id, not the user-facing thread id
+ * thread_name VARCHAR(255) NOT NULL, -- user-facing thread id accepted by the API
* is_released BOOLEAN DEFAULT FALSE NOT NULL,
* active_thread_name VARCHAR(255) GENERATED ALWAYS AS (
* CASE WHEN is_released = FALSE THEN thread_name ELSE NULL END
@@ -73,7 +73,7 @@
* CREATE TABLE GRAPH_CHECKPOINT (
* checkpoint_seq BIGINT GENERATED BY DEFAULT AS IDENTITY UNIQUE,
* checkpoint_id VARCHAR(36) PRIMARY KEY,
- * thread_id VARCHAR(36) NOT NULL,
+ * thread_id VARCHAR(36) NOT NULL, -- references the internal surrogate id
* node_id VARCHAR(255),
* next_node_id VARCHAR(255),
* state_data CLOB NOT NULL,
@@ -88,6 +88,14 @@
*
*
*
+ * Thread identity: GRAPH_THREAD.thread_name stores the thread id supplied
+ * through {@code RunnableConfig}, while GRAPH_THREAD.thread_id stores an
+ * internally generated UUID that identifies one activation of that thread
+ * between a release and the next reuse of the same id. This column split is
+ * what allows a released thread id to be reused without orphaning the released
+ * checkpoint history.
+ *
+ *
* A builder can be used to create an instance of H2Saver. The builder allows
* configuring a DataSource or JDBC URL, CreateOption, StateSerializer, and the
* maximum number of latest checkpoints retained in memory.
@@ -388,30 +396,34 @@ protected void insertCheckpoint(String threadId, Checkpoint checkpoint) throws E
}
}
- private String activeThreadId(Connection conn, String threadName) throws SQLException {
- Optional activeThreadId = selectActiveThreadId(conn, threadName);
+ /**
+ * Returns the internal surrogate thread id of the active row for the given
+ * user-facing thread id, inserting a new row when none is active.
+ */
+ private String activeThreadId(Connection conn, String threadId) throws SQLException {
+ Optional activeThreadId = selectActiveThreadId(conn, threadId);
if (activeThreadId.isPresent()) {
return activeThreadId.get();
}
- String persistedThreadId = UUID.randomUUID().toString();
+ String newThreadId = UUID.randomUUID().toString();
try (PreparedStatement ps = conn.prepareStatement(INSERT_THREAD)) {
- ps.setString(1, persistedThreadId);
- ps.setString(2, threadName);
+ ps.setString(1, newThreadId);
+ ps.setString(2, threadId);
ps.executeUpdate();
- return persistedThreadId;
+ return newThreadId;
}
catch (SQLException ex) {
if (isUniqueConstraintViolation(ex)) {
- return selectActiveThreadId(conn, threadName).orElseThrow(() -> ex);
+ return selectActiveThreadId(conn, threadId).orElseThrow(() -> ex);
}
throw ex;
}
}
- private Optional selectActiveThreadId(Connection conn, String threadName) throws SQLException {
+ private Optional selectActiveThreadId(Connection conn, String threadId) throws SQLException {
try (PreparedStatement ps = conn.prepareStatement(SELECT_ACTIVE_THREAD)) {
- ps.setString(1, threadName);
+ ps.setString(1, threadId);
try (ResultSet rs = ps.executeQuery()) {
if (rs.next()) {
return Optional.of(rs.getString(1));
diff --git a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/jdbc/AbstractJdbcCheckpointSaver.java b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/jdbc/AbstractJdbcCheckpointSaver.java
index 6a379a1c8..e15336d24 100644
--- a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/jdbc/AbstractJdbcCheckpointSaver.java
+++ b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/jdbc/AbstractJdbcCheckpointSaver.java
@@ -35,6 +35,15 @@
* This class owns the common saver lifecycle and latest-checkpoint cache behavior.
* Subclasses keep database-specific SQL, transaction details and row mapping logic.
*
+ * Thread identity model: every saver operation is keyed on the user-facing thread
+ * id resolved from {@link RunnableConfig} (see
+ * {@link BaseCheckpointSaver#checkpointThreadId(RunnableConfig)}), which is the
+ * only thread identifier the public API exposes. Concrete schemas persist that
+ * id in their {@code thread_name} column, while their {@code thread_id} column
+ * holds an internally generated surrogate UUID identifying one activation of the
+ * thread between a release and the next reuse of the same id. The protected
+ * methods below always receive the user-facing thread id.
+ *
* Replacement: add artifact
* {@code io.github.agentic-spring-ai:agentic-spring-ai-graph-persistence-jdbc}
* and use
@@ -209,7 +218,7 @@ public final void latestCheckpointCacheEnabled(boolean enabled) {
/**
* Selects the active checkpoint history for a thread from the backing database.
*
- * @param threadId thread name/id used by the concrete saver schema
+ * @param threadId user-facing thread id; persisted by the concrete schema as its thread name
* @return checkpoint history in latest-first order
* @throws Exception when the concrete saver cannot read checkpoint history
*/
@@ -218,7 +227,7 @@ public final void latestCheckpointCacheEnabled(boolean enabled) {
/**
* Selects only the latest active checkpoint for a thread.
*
- * @param threadId thread name/id used by the concrete saver schema
+ * @param threadId user-facing thread id; persisted by the concrete schema as its thread name
* @return latest checkpoint when one exists
* @throws Exception when the concrete saver cannot read the latest checkpoint
*/
@@ -227,7 +236,7 @@ public final void latestCheckpointCacheEnabled(boolean enabled) {
/**
* Selects an active checkpoint by id for a thread.
*
- * @param threadId thread name/id used by the concrete saver schema
+ * @param threadId user-facing thread id; persisted by the concrete schema as its thread name
* @param checkpointId checkpoint id to look up
* @return matching checkpoint when one exists
* @throws Exception when the concrete saver cannot read the checkpoint
@@ -237,7 +246,7 @@ public final void latestCheckpointCacheEnabled(boolean enabled) {
/**
* Inserts a new active checkpoint for a thread.
*
- * @param threadId thread name/id used by the concrete saver schema
+ * @param threadId user-facing thread id; persisted by the concrete schema as its thread name
* @param checkpoint checkpoint to persist
* @throws Exception when the concrete saver cannot insert the checkpoint
*/
@@ -246,7 +255,7 @@ public final void latestCheckpointCacheEnabled(boolean enabled) {
/**
* Replaces an existing active checkpoint for a thread.
*
- * @param threadId thread name/id used by the concrete saver schema
+ * @param threadId user-facing thread id; persisted by the concrete schema as its thread name
* @param checkpointId checkpoint id to replace
* @param checkpoint replacement checkpoint data
* @throws Exception when the concrete saver cannot update the checkpoint
@@ -256,7 +265,7 @@ public final void latestCheckpointCacheEnabled(boolean enabled) {
/**
* Deletes active checkpoints by id for a thread.
*
- * @param threadId thread name/id used by the concrete saver schema
+ * @param threadId user-facing thread id; persisted by the concrete schema as its thread name
* @param checkpointIds checkpoint ids to delete
* @throws Exception when the concrete saver cannot delete checkpoints
*/
@@ -281,7 +290,7 @@ private void deleteRetainedCheckpoints(String threadId, RunnableConfig config) t
/**
* Marks the active thread as released in the backing database.
*
- * @param threadId thread name/id used by the concrete saver schema
+ * @param threadId user-facing thread id; persisted by the concrete schema as its thread name
* @throws Exception when the concrete saver cannot release the thread
*/
protected abstract void releaseThread(String threadId) throws Exception;
diff --git a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/mongo/MongoSaver.java b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/mongo/MongoSaver.java
index ccc97cb45..59b83968d 100644
--- a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/mongo/MongoSaver.java
+++ b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/mongo/MongoSaver.java
@@ -63,6 +63,16 @@
/**
* MongoDB checkpoint saver.
*
+ * Thread identity: every operation is keyed on the user-facing thread id
+ * resolved from {@code RunnableConfig}. That id is used verbatim as the
+ * {@code _id} of the {@code thread_meta} document
+ * ({@code mongo:thread:meta:}) and stored in its
+ * {@code thread_name} field, while the {@code thread_id} field holds an
+ * internally generated UUID identifying one activation of the thread between a
+ * release and the next reuse of the same id. Checkpoint documents are stored
+ * under that internal id so a released thread id can be reused without
+ * orphaning the released checkpoint history.
+ *
* Replacement: add artifact
* {@code io.github.agentic-spring-ai:agentic-spring-ai-graph-persistence-mongodb}
* and use {@code io.github.agentic.spring.ai.graph.persistence.mongodb.MongoSaver}.
@@ -155,20 +165,20 @@ private LinkedList deserializeCheckpoints(String content) throws IOE
}
/**
- * Gets or creates a thread_id for the given thread_name.
- * If an active thread exists, returns its thread_id.
- * If no active thread exists or the thread is released, creates a new thread_id.
+ * Returns the internal surrogate thread id of the active entry for the given
+ * user-facing thread id, creating a new one when no active entry exists or
+ * the previous one was released.
*
* This method uses atomic operations to prevent race conditions in concurrent scenarios.
* Uses findOneAndUpdate with conditional logic to ensure thread-safe creation.
*
- * @param threadName the thread name
+ * @param threadId the user-facing thread id
* @param clientSession the MongoDB client session for transaction
- * @return the thread_id (UUID string)
+ * @return the internal thread id (UUID string)
*/
- private String getOrCreateThreadId(String threadName, ClientSession clientSession) {
+ private String getOrCreateThreadId(String threadId, ClientSession clientSession) {
MongoCollection threadMetaCollection = database.getCollection(THREAD_META_COLLECTION);
- String metaId = THREAD_META_PREFIX + threadName;
+ String metaId = THREAD_META_PREFIX + threadId;
// Step 1: Try to atomically get an active thread
// Filter: _id matches AND is_released != true
@@ -187,10 +197,10 @@ private String getOrCreateThreadId(String threadName, ClientSession clientSessio
);
if (existingDoc != null) {
- String threadId = existingDoc.getString(FIELD_THREAD_ID);
- if (threadId != null) {
- // Active thread exists, return its thread_id
- return threadId;
+ String persistedThreadId = existingDoc.getString(FIELD_THREAD_ID);
+ if (persistedThreadId != null) {
+ // Active thread exists, return its internal thread id
+ return persistedThreadId;
}
}
@@ -271,53 +281,57 @@ private String getOrCreateThreadId(String threadName, ClientSession clientSessio
}
/**
- * Gets the active thread_id for the given thread_name.
+ * Gets the internal surrogate thread id of the active entry for the given
+ * user-facing thread id.
* Returns null if no active thread exists.
*
- * @param threadName the thread name
+ * @param threadId the user-facing thread id
* @param clientSession the MongoDB client session for transaction
- * @return the active thread_id, or null if not found
+ * @return the active internal thread id, or null if not found
*/
- private String getActiveThreadId(String threadName, ClientSession clientSession) {
+ private String getActiveThreadId(String threadId, ClientSession clientSession) {
MongoCollection threadMetaCollection = database.getCollection(THREAD_META_COLLECTION);
- String metaId = THREAD_META_PREFIX + threadName;
+ String metaId = THREAD_META_PREFIX + threadId;
Document metaDoc = threadMetaCollection.find(clientSession, new BasicDBObject("_id", metaId)).first();
if (metaDoc != null) {
- String threadId = metaDoc.getString(FIELD_THREAD_ID);
+ String persistedThreadId = metaDoc.getString(FIELD_THREAD_ID);
Boolean isReleased = metaDoc.getBoolean(FIELD_IS_RELEASED, false);
- if (threadId != null && !Boolean.TRUE.equals(isReleased)) {
- return threadId;
+ if (persistedThreadId != null && !Boolean.TRUE.equals(isReleased)) {
+ return persistedThreadId;
}
}
return null; // No active thread exists
}
- private String threadName(RunnableConfig config) {
+ /**
+ * Resolves the user-facing thread id every saver operation is keyed on.
+ */
+ private String threadId(RunnableConfig config) {
return checkpointThreadId(config);
}
@Override
public Collection list(RunnableConfig config) {
- String threadName = threadName(config);
+ String threadId = threadId(config);
ClientSession clientSession = this.client
.startSession(ClientSessionOptions.builder().defaultTransactionOptions(txnOptions).build());
clientSession.startTransaction();
List checkpoints = null;
try {
- // Get active thread_id for the thread_name
- String threadId = getActiveThreadId(threadName, clientSession);
- if (threadId == null) {
+ // Get the internal thread id of the active entry
+ String persistedThreadId = getActiveThreadId(threadId, clientSession);
+ if (persistedThreadId == null) {
clientSession.commitTransaction();
return Collections.emptyList();
}
- // Use thread_id to query checkpoints
+ // Use the internal thread id to query checkpoints
MongoCollection collection = database.getCollection(CHECKPOINT_COLLECTION);
- String checkpointId = CHECKPOINT_PREFIX + threadId;
+ String checkpointId = CHECKPOINT_PREFIX + persistedThreadId;
Document document = collection.find(clientSession, new BasicDBObject("_id", checkpointId)).first();
if (document == null) {
clientSession.commitTransaction();
@@ -339,23 +353,23 @@ public Collection list(RunnableConfig config) {
@Override
public Optional get(RunnableConfig config) {
- String threadName = threadName(config);
+ String threadId = threadId(config);
ClientSession clientSession = this.client
.startSession(ClientSessionOptions.builder().defaultTransactionOptions(txnOptions).build());
LinkedList checkpoints = null;
try {
clientSession.startTransaction();
- // Get active thread_id for the thread_name
- String threadId = getActiveThreadId(threadName, clientSession);
- if (threadId == null) {
+ // Get the internal thread id of the active entry
+ String persistedThreadId = getActiveThreadId(threadId, clientSession);
+ if (persistedThreadId == null) {
clientSession.commitTransaction();
return Optional.empty();
}
- // Use thread_id to query checkpoints
+ // Use the internal thread id to query checkpoints
MongoCollection collection = database.getCollection(CHECKPOINT_COLLECTION);
- String checkpointId = CHECKPOINT_PREFIX + threadId;
+ String checkpointId = CHECKPOINT_PREFIX + persistedThreadId;
Document document = collection.find(clientSession, new BasicDBObject("_id", checkpointId)).first();
if (document == null) {
clientSession.commitTransaction();
@@ -385,17 +399,17 @@ public Optional get(RunnableConfig config) {
@Override
public RunnableConfig put(RunnableConfig config, Checkpoint checkpoint) throws Exception {
- String threadName = threadName(config);
+ String threadId = threadId(config);
ClientSession clientSession = this.client
.startSession(ClientSessionOptions.builder().defaultTransactionOptions(txnOptions).build());
clientSession.startTransaction();
try {
- // Get or create thread_id
- String threadId = getOrCreateThreadId(threadName, clientSession);
+ // Get or create the internal thread id
+ String persistedThreadId = getOrCreateThreadId(threadId, clientSession);
- // Use thread_id as key for checkpoint storage
+ // Use the internal thread id as key for checkpoint storage
MongoCollection collection = database.getCollection(CHECKPOINT_COLLECTION);
- String checkpointDocId = CHECKPOINT_PREFIX + threadId;
+ String checkpointDocId = CHECKPOINT_PREFIX + persistedThreadId;
Document document = collection.find(clientSession, new BasicDBObject("_id", checkpointDocId)).first();
LinkedList checkpointLinkedList = null;
@@ -449,24 +463,24 @@ public RunnableConfig put(RunnableConfig config, Checkpoint checkpoint) throws E
@Override
public Tag release(RunnableConfig config) throws Exception {
- String threadName = threadName(config);
+ String threadId = threadId(config);
ClientSession clientSession = this.client
.startSession(ClientSessionOptions.builder().defaultTransactionOptions(txnOptions).build());
clientSession.startTransaction();
try {
MongoCollection threadMetaCollection = database.getCollection(THREAD_META_COLLECTION);
- String metaId = THREAD_META_PREFIX + threadName;
+ String metaId = THREAD_META_PREFIX + threadId;
Document metaDoc = threadMetaCollection.find(clientSession, new BasicDBObject("_id", metaId)).first();
if (metaDoc == null) {
clientSession.abortTransaction();
- throw new IllegalStateException("Thread not found: " + threadName);
+ throw new IllegalStateException("Thread not found: " + threadId);
}
- String threadId = metaDoc.getString(FIELD_THREAD_ID);
- if (threadId == null) {
+ String persistedThreadId = metaDoc.getString(FIELD_THREAD_ID);
+ if (persistedThreadId == null) {
clientSession.abortTransaction();
- throw new IllegalStateException("Thread not found: " + threadName);
+ throw new IllegalStateException("Thread not found: " + threadId);
}
// Mark thread as released atomically
@@ -484,12 +498,12 @@ public Tag release(RunnableConfig config) throws Exception {
if (updatedDoc == null) {
// Thread was already released or doesn't exist
clientSession.abortTransaction();
- throw new IllegalStateException("Thread is not active or already released: " + threadName);
+ throw new IllegalStateException("Thread is not active or already released: " + threadId);
}
- // Get checkpoints for Tag (using thread_id)
+ // Get checkpoints for Tag (using the internal thread id)
MongoCollection checkpointCollection = database.getCollection(CHECKPOINT_COLLECTION);
- String checkpointDocId = CHECKPOINT_PREFIX + threadId;
+ String checkpointDocId = CHECKPOINT_PREFIX + persistedThreadId;
Document checkpointDoc = checkpointCollection.find(clientSession, new BasicDBObject("_id", checkpointDocId))
.first();
@@ -502,7 +516,7 @@ public Tag release(RunnableConfig config) throws Exception {
}
clientSession.commitTransaction();
- return new Tag(threadName, checkpoints);
+ return new Tag(threadId, checkpoints);
}
catch (Exception e) {
diff --git a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/mysql/MysqlSaver.java b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/mysql/MysqlSaver.java
index 4b9394af0..3b6ecd3dc 100644
--- a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/mysql/MysqlSaver.java
+++ b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/mysql/MysqlSaver.java
@@ -53,8 +53,8 @@
*
*
* CREATE TABLE GRAPH_THREAD (
- * thread_id VARCHAR(36) PRIMARY KEY,
- * thread_name VARCHAR(255),
+ * thread_id VARCHAR(36) PRIMARY KEY, -- internal surrogate id, not the user-facing thread id
+ * thread_name VARCHAR(255), -- user-facing thread id accepted by the API
* is_released BOOLEAN DEFAULT FALSE NOT NULL,
* active_thread_name VARCHAR(255) GENERATED ALWAYS AS (
* CASE WHEN is_released = FALSE THEN thread_name ELSE NULL END
@@ -66,7 +66,7 @@
* CREATE TABLE GRAPH_CHECKPOINT (
* checkpoint_seq BIGINT NOT NULL AUTO_INCREMENT UNIQUE,
* checkpoint_id VARCHAR(36) PRIMARY KEY,
- * thread_id VARCHAR(36) NOT NULL,
+ * thread_id VARCHAR(36) NOT NULL, -- references the internal surrogate id
* node_id VARCHAR(255),
* next_node_id VARCHAR(255),
* state_data JSON NOT NULL,
@@ -80,6 +80,14 @@
*
*
*
+ * Thread identity: GRAPH_THREAD.thread_name stores the thread id supplied
+ * through {@code RunnableConfig}, while GRAPH_THREAD.thread_id stores an
+ * internally generated UUID that identifies one activation of that thread
+ * between a release and the next reuse of the same id. This column split is
+ * what allows a released thread id to be reused without orphaning the released
+ * checkpoint history.
+ *
+ *
* A builder can be used to create an instance of MysqlSaver. The builder
* allows to configure the following options:
* - DataSource: indicates which data source should be used to connect
@@ -428,12 +436,12 @@ private Checkpoint readCheckpoint(ResultSet resultSet)
* Loads full checkpoint history on demand without retaining it in cache.
*/
@Override
- protected LinkedList selectCheckpoints(String threadName) throws Exception {
+ protected LinkedList selectCheckpoints(String threadId) throws Exception {
LinkedList checkpoints = new LinkedList<>();
try (Connection connection = dataSource.getConnection();
PreparedStatement preparedStatement = connection.prepareStatement(SELECT_CHECKPOINTS)) {
- preparedStatement.setString(1, threadName);
+ preparedStatement.setString(1, threadId);
try (ResultSet resultSet = preparedStatement.executeQuery()) {
while (resultSet.next()) {
checkpoints.add(readCheckpoint(resultSet));
@@ -447,11 +455,11 @@ protected LinkedList selectCheckpoints(String threadName) throws Exc
}
@Override
- protected Optional selectLatestCheckpoint(String threadName) throws Exception {
+ protected Optional selectLatestCheckpoint(String threadId) throws Exception {
try (Connection connection = dataSource.getConnection();
PreparedStatement preparedStatement = connection.prepareStatement(SELECT_LATEST_CHECKPOINT)) {
- preparedStatement.setString(1, threadName);
+ preparedStatement.setString(1, threadId);
try (ResultSet resultSet = preparedStatement.executeQuery()) {
if (resultSet.next()) {
return Optional.of(readCheckpoint(resultSet));
@@ -465,11 +473,11 @@ protected Optional selectLatestCheckpoint(String threadName) throws
}
@Override
- protected Optional selectCheckpointById(String threadName, String checkpointId) throws Exception {
+ protected Optional selectCheckpointById(String threadId, String checkpointId) throws Exception {
try (Connection connection = dataSource.getConnection();
PreparedStatement preparedStatement = connection.prepareStatement(SELECT_CHECKPOINT_BY_ID)) {
- preparedStatement.setString(1, threadName);
+ preparedStatement.setString(1, threadId);
preparedStatement.setString(2, checkpointId);
try (ResultSet resultSet = preparedStatement.executeQuery()) {
if (resultSet.next()) {
@@ -484,7 +492,7 @@ protected Optional selectCheckpointById(String threadName, String ch
}
@Override
- protected void insertCheckpoint(String threadName, Checkpoint checkpoint) throws Exception {
+ protected void insertCheckpoint(String threadId, Checkpoint checkpoint) throws Exception {
Connection conn = null;
try (Connection ignored = conn = dataSource.getConnection()) {
conn.setAutoCommit(false);
@@ -493,29 +501,29 @@ protected void insertCheckpoint(String threadName, Checkpoint checkpoint) throws
PreparedStatement insertCheckpointStatement = conn.prepareStatement(INSERT_CHECKPOINT)) {
upsertStatement.setString(1, UUID.randomUUID().toString());
- upsertStatement.setString(2, threadName);
+ upsertStatement.setString(2, threadId);
upsertStatement.execute();
insertCheckpointStatement.setString(1, checkpoint.getId());
insertCheckpointStatement.setString(2, checkpoint.getNodeId());
insertCheckpointStatement.setString(3, checkpoint.getNextNodeId());
insertCheckpointStatement.setString(4, encodeState(checkpoint.getState()));
- insertCheckpointStatement.setString(5, threadName);
+ insertCheckpointStatement.setString(5, threadId);
insertCheckpointStatement.execute();
}
conn.commit();
- log.debug("Checkpoint {} for thread {} inserted successfully.", checkpoint.getId(), threadName);
+ log.debug("Checkpoint {} for thread {} inserted successfully.", checkpoint.getId(), threadId);
}
catch (SQLException | IOException ex) {
- log.error("Error inserting checkpoint with id {} in thread {}", checkpoint.getId(), threadName, ex);
- rollback(conn, checkpoint, threadName);
+ log.error("Error inserting checkpoint with id {} in thread {}", checkpoint.getId(), threadId, ex);
+ rollback(conn, checkpoint, threadId);
throw new Exception("Unable to insert checkpoint", ex);
}
}
@Override
- protected void updateCheckpoint(String threadName, String checkpointId, Checkpoint checkpoint) throws Exception {
+ protected void updateCheckpoint(String threadId, String checkpointId, Checkpoint checkpoint) throws Exception {
Connection conn = null;
try (Connection ignored = conn = dataSource.getConnection()) {
conn.setAutoCommit(false);
@@ -525,7 +533,7 @@ protected void updateCheckpoint(String threadName, String checkpointId, Checkpoi
preparedStatement.setString(2, checkpoint.getNodeId());
preparedStatement.setString(3, checkpoint.getNextNodeId());
preparedStatement.setString(4, encodeState(checkpoint.getState()));
- preparedStatement.setString(5, threadName);
+ preparedStatement.setString(5, threadId);
preparedStatement.setString(6, checkpointId);
int rowsAffected = preparedStatement.executeUpdate();
if (rowsAffected == 0) {
@@ -535,24 +543,24 @@ protected void updateCheckpoint(String threadName, String checkpointId, Checkpoi
}
conn.commit();
- log.debug("Checkpoint with id {} for thread {} updated successfully.", checkpoint.getId(), threadName);
+ log.debug("Checkpoint with id {} for thread {} updated successfully.", checkpoint.getId(), threadId);
}
catch (SQLException | IOException ex) {
- log.error("Error updating checkpoint with id {} in thread {}", checkpoint.getId(), threadName, ex);
- rollback(conn, checkpoint, threadName);
+ log.error("Error updating checkpoint with id {} in thread {}", checkpoint.getId(), threadId, ex);
+ rollback(conn, checkpoint, threadId);
throw new Exception("Unable to update checkpoint", ex);
}
}
@Override
- protected void deleteCheckpoints(String threadName, Collection checkpointIds) throws Exception {
+ protected void deleteCheckpoints(String threadId, Collection checkpointIds) throws Exception {
if (checkpointIds.isEmpty()) {
return;
}
try (Connection connection = dataSource.getConnection();
PreparedStatement preparedStatement = connection.prepareStatement(
DELETE_CHECKPOINTS.formatted(String.join(", ", Collections.nCopies(checkpointIds.size(), "?"))))) {
- preparedStatement.setString(1, threadName);
+ preparedStatement.setString(1, threadId);
int index = 2;
for (String checkpointId : checkpointIds) {
preparedStatement.setString(index++, checkpointId);
@@ -565,26 +573,26 @@ protected void deleteCheckpoints(String threadName, Collection checkpoin
}
@Override
- protected void releaseThread(String threadName) throws Exception {
+ protected void releaseThread(String threadId) throws Exception {
Connection conn = null;
try (Connection ignored = conn = dataSource.getConnection()) {
conn.setAutoCommit(false);
try (PreparedStatement preparedStatement = conn.prepareStatement(RELEASE_THREAD)) {
- preparedStatement.setString(1, threadName);
+ preparedStatement.setString(1, threadId);
int rowsAffected = preparedStatement.executeUpdate();
if (rowsAffected == 0) {
conn.rollback();
- throw new IllegalStateException(format("Thread '%s' not found or already released", threadName));
+ throw new IllegalStateException(format("Thread '%s' not found or already released", threadId));
}
}
conn.commit();
- log.debug("Thread {} released successfully.", threadName);
+ log.debug("Thread {} released successfully.", threadId);
}
catch (SQLException ex) {
- log.error("Error releasing thread {}", threadName, ex);
- rollback(conn, threadName);
+ log.error("Error releasing thread {}", threadId, ex);
+ rollback(conn, threadId);
throw new Exception("Unable to release checkpoint", ex);
}
}
diff --git a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/oracle/OracleSaver.java b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/oracle/OracleSaver.java
index fde34944a..9d04d6f26 100644
--- a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/oracle/OracleSaver.java
+++ b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/oracle/OracleSaver.java
@@ -63,8 +63,8 @@
*
*
* CREATE TABLE GRAPH_THREAD (
- * thread_id VARCHAR2(36) PRIMARY KEY,
- * thread_name VARCHAR(255),
+ * thread_id VARCHAR2(36) PRIMARY KEY, -- internal surrogate id, not the user-facing thread id
+ * thread_name VARCHAR(255), -- user-facing thread id accepted by the API
* is_released BOOLEAN DEFAULT FALSE NOT NULL
* )
* CREATE INDEX IDX_GRAPH_THREAD_NAME_RELEASED
@@ -72,7 +72,7 @@
*
* CREATE TABLE GRAPH_CHECKPOINT (
* checkpoint_id VARCHAR2(36) PRIMARY KEY,
- * thread_id VARCHAR2(36) NOT NULL,
+ * thread_id VARCHAR2(36) NOT NULL, -- references the internal surrogate id
* node_id VARCHAR(255),
* next_node_id VARCHAR(255),
* state_data JSON NOT NULL,
@@ -87,6 +87,14 @@
*
*
*
+ * Thread identity: GRAPH_THREAD.thread_name stores the thread id supplied
+ * through {@code RunnableConfig}, while GRAPH_THREAD.thread_id stores an
+ * internally generated UUID that identifies one activation of that thread
+ * between a release and the next reuse of the same id. This column split is
+ * what allows a released thread id to be reused without orphaning the released
+ * checkpoint history.
+ *
+ *
* A builder can be use to create an instance or OracleSaver. The builder
* allows to configure the following options:
* - DataSource: indicates which data source should be used to connect
@@ -388,14 +396,14 @@ private Checkpoint readCheckpoint(ResultSet resultSet, ObjectMapper objectMapper
* Loads full checkpoint history on demand without retaining it in cache.
*/
@Override
- protected LinkedList selectCheckpoints(String threadName) throws Exception {
+ protected LinkedList selectCheckpoints(String threadId) throws Exception {
LinkedList checkpoints = new LinkedList<>();
ObjectMapper objectMapper = osonObjectMapper();
try (Connection connection = dataSource.getConnection();
PreparedStatement preparedStatement = connection.prepareStatement(SELECT_CHECKPOINTS)) {
defineCheckpointColumns(preparedStatement);
- preparedStatement.setString(1, threadName);
+ preparedStatement.setString(1, threadId);
try (ResultSet resultSet = preparedStatement.executeQuery()) {
while (resultSet.next()) {
checkpoints.add(readCheckpoint(resultSet, objectMapper));
@@ -409,13 +417,13 @@ protected LinkedList selectCheckpoints(String threadName) throws Exc
}
@Override
- protected Optional selectLatestCheckpoint(String threadName) throws Exception {
+ protected Optional selectLatestCheckpoint(String threadId) throws Exception {
ObjectMapper objectMapper = osonObjectMapper();
try (Connection connection = dataSource.getConnection();
PreparedStatement preparedStatement = connection.prepareStatement(SELECT_LATEST_CHECKPOINT)) {
defineCheckpointColumns(preparedStatement);
- preparedStatement.setString(1, threadName);
+ preparedStatement.setString(1, threadId);
try (ResultSet resultSet = preparedStatement.executeQuery()) {
if (resultSet.next()) {
return Optional.of(readCheckpoint(resultSet, objectMapper));
@@ -429,13 +437,13 @@ protected Optional selectLatestCheckpoint(String threadName) throws
}
@Override
- protected Optional selectCheckpointById(String threadName, String checkpointId) throws Exception {
+ protected Optional selectCheckpointById(String threadId, String checkpointId) throws Exception {
ObjectMapper objectMapper = osonObjectMapper();
try (Connection connection = dataSource.getConnection();
PreparedStatement preparedStatement = connection.prepareStatement(SELECT_CHECKPOINT_BY_ID)) {
defineCheckpointColumns(preparedStatement);
- preparedStatement.setString(1, threadName);
+ preparedStatement.setString(1, threadId);
preparedStatement.setString(2, checkpointId);
try (ResultSet resultSet = preparedStatement.executeQuery()) {
if (resultSet.next()) {
@@ -450,7 +458,7 @@ protected Optional selectCheckpointById(String threadName, String ch
}
@Override
- protected void insertCheckpoint(String threadName, Checkpoint checkpoint) throws Exception {
+ protected void insertCheckpoint(String threadId, Checkpoint checkpoint) throws Exception {
Connection conn = null;
try (Connection ignored = conn = dataSource.getConnection()) {
@@ -460,7 +468,7 @@ protected void insertCheckpoint(String threadName, Checkpoint checkpoint) throws
PreparedStatement insertCheckpointStatement = conn.prepareStatement(INSERT_CHECKPOINT)) {
upsertStatement.setString(1, UUID.randomUUID().toString());
- upsertStatement.setString(2, threadName);
+ upsertStatement.setString(2, threadId);
upsertStatement.execute();
String encodedState = encodeState(checkpoint.getState());
@@ -469,23 +477,23 @@ protected void insertCheckpoint(String threadName, Checkpoint checkpoint) throws
insertCheckpointStatement.setString(3, checkpoint.getNextNodeId());
insertCheckpointStatement.setObject(4, encodedState, OracleType.JSON);
insertCheckpointStatement.setString(5, stateSerializer.contentType());
- insertCheckpointStatement.setString(6, threadName);
+ insertCheckpointStatement.setString(6, threadId);
insertCheckpointStatement.execute();
}
conn.commit();
- log.debug("Checkpoint {} for thread {} inserted successfully.", checkpoint.getId(), threadName);
+ log.debug("Checkpoint {} for thread {} inserted successfully.", checkpoint.getId(), threadId);
}
catch (SQLException | IOException ex) {
- log.error("Error inserting checkpoint with id {} in thread {}", checkpoint.getId(), threadName, ex);
- rollback(conn, checkpoint, threadName);
+ log.error("Error inserting checkpoint with id {} in thread {}", checkpoint.getId(), threadId, ex);
+ rollback(conn, checkpoint, threadId);
throw new Exception("Unable to insert checkpoint", ex);
}
}
@Override
- protected void updateCheckpoint(String threadName, String checkpointId, Checkpoint checkpoint) throws Exception {
+ protected void updateCheckpoint(String threadId, String checkpointId, Checkpoint checkpoint) throws Exception {
Connection conn = null;
try (Connection ignored = conn = dataSource.getConnection()) {
@@ -499,7 +507,7 @@ protected void updateCheckpoint(String threadName, String checkpointId, Checkpoi
preparedStatement.setObject(4, encodedState, OracleType.JSON);
preparedStatement.setString(5, stateSerializer.contentType());
preparedStatement.setString(6, checkpointId);
- preparedStatement.setString(7, threadName);
+ preparedStatement.setString(7, threadId);
int rowsAffected = preparedStatement.executeUpdate();
if (rowsAffected == 0) {
conn.rollback();
@@ -508,17 +516,17 @@ protected void updateCheckpoint(String threadName, String checkpointId, Checkpoi
}
conn.commit();
- log.debug("Checkpoint with id {} for thread {} updated successfully.", checkpoint.getId(), threadName);
+ log.debug("Checkpoint with id {} for thread {} updated successfully.", checkpoint.getId(), threadId);
}
catch (SQLException | IOException ex) {
- log.error("Error updating checkpoint with id {} in thread {}", checkpoint.getId(), threadName, ex);
- rollback(conn, checkpoint, threadName);
+ log.error("Error updating checkpoint with id {} in thread {}", checkpoint.getId(), threadId, ex);
+ rollback(conn, checkpoint, threadId);
throw new Exception("Unable to update checkpoint", ex);
}
}
@Override
- protected void deleteCheckpoints(String threadName, Collection checkpointIds) throws Exception {
+ protected void deleteCheckpoints(String threadId, Collection checkpointIds) throws Exception {
if (checkpointIds.isEmpty()) {
return;
}
@@ -529,7 +537,7 @@ protected void deleteCheckpoints(String threadName, Collection checkpoin
for (String checkpointId : checkpointIds) {
preparedStatement.setString(index++, checkpointId);
}
- preparedStatement.setString(index, threadName);
+ preparedStatement.setString(index, threadId);
preparedStatement.executeUpdate();
}
catch (SQLException ex) {
@@ -538,29 +546,29 @@ protected void deleteCheckpoints(String threadName, Collection checkpoin
}
@Override
- protected void releaseThread(String threadName) throws Exception {
+ protected void releaseThread(String threadId) throws Exception {
Connection conn = null;
try (Connection ignored = conn = dataSource.getConnection()) {
conn.setAutoCommit(false);
try (PreparedStatement preparedStatement = conn.prepareStatement(RELEASE_THREAD)) {
- preparedStatement.setString(1, threadName);
+ preparedStatement.setString(1, threadId);
int rowsAffected = preparedStatement.executeUpdate();
if (rowsAffected == 0) {
conn.rollback();
throw new IllegalStateException(
- format("Thread '%s' not found or already released", threadName));
+ format("Thread '%s' not found or already released", threadId));
}
}
conn.commit();
- log.debug("Thread {} released successfully.", threadName);
+ log.debug("Thread {} released successfully.", threadId);
}
catch (SQLException ex) {
- log.error("Error releasing thread {}", threadName, ex);
- rollback(conn, threadName);
+ log.error("Error releasing thread {}", threadId, ex);
+ rollback(conn, threadId);
throw new Exception("Unable to release checkpoint", ex);
}
}
diff --git a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/postgresql/PostgresSaver.java b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/postgresql/PostgresSaver.java
index 9120e5405..0202f531f 100644
--- a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/postgresql/PostgresSaver.java
+++ b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/postgresql/PostgresSaver.java
@@ -56,8 +56,8 @@
*
*
* CREATE TABLE GraphThread (
- * thread_id UUID PRIMARY KEY,
- * thread_name VARCHAR(255),
+ * thread_id UUID PRIMARY KEY, -- internal surrogate id, not the user-facing thread id
+ * thread_name VARCHAR(255), -- user-facing thread id accepted by the API
* is_released BOOLEAN DEFAULT FALSE NOT NULL
* )
* CREATE UNIQUE INDEX idx_unique_lg4jthread_thread_name_unreleased
@@ -66,7 +66,7 @@
* CREATE TABLE GraphCheckpoint (
* checkpoint_id UUID PRIMARY KEY,
* parent_checkpoint_id UUID,
- * thread_id UUID NOT NULL,
+ * thread_id UUID NOT NULL, -- references the internal surrogate id
* node_id VARCHAR(255),
* next_node_id VARCHAR(255),
* state_data JSONB NOT NULL,
@@ -81,6 +81,14 @@
*
*
*
+ * Thread identity: GraphThread.thread_name stores the thread id supplied
+ * through {@code RunnableConfig}, while GraphThread.thread_id stores an
+ * internally generated UUID that identifies one activation of that thread
+ * between a release and the next reuse of the same id. This column split is
+ * what allows a released thread id to be reused without orphaning the released
+ * checkpoint history.
+ *
+ *
* A builder can be used to create an instance of PostgresSaver. The builder
* allows to configure the following options:
* - DataSource: indicates which data source should be used to connect
diff --git a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/redis/RedisSaver.java b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/redis/RedisSaver.java
index 36dfaf3a9..d64b973a6 100644
--- a/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/redis/RedisSaver.java
+++ b/agentic-ai-graph-core/src/main/java/io/github/agentic/spring/ai/graph/checkpoint/savers/redis/RedisSaver.java
@@ -49,6 +49,15 @@
/**
* The type Redis saver.
*
+ * Thread identity: every operation is keyed on the user-facing thread id
+ * resolved from {@code RunnableConfig}. That id is used verbatim in the
+ * {@code graph:thread:meta:} key and stored in the
+ * {@code thread_name} field of the reverse mapping, while the {@code thread_id}
+ * hash field holds an internally generated UUID identifying one activation of
+ * the thread between a release and the next reuse of the same id. Checkpoint
+ * buckets are stored under that internal id so a released thread id can be
+ * reused without orphaning the released checkpoint history.
+ *
* Replacement: add artifact
* {@code io.github.agentic-spring-ai:agentic-spring-ai-graph-persistence-redis}
* and use {@code io.github.agentic.spring.ai.graph.persistence.redis.RedisSaver}.
@@ -131,27 +140,27 @@ private LinkedList deserializeCheckpoints(String content) throws IOE
}
/**
- * Gets or creates a thread_id for the given thread_name.
- * If an active thread exists, returns its thread_id.
- * If no active thread exists or the thread is released, creates a new thread_id.
+ * Returns the internal surrogate thread id of the active entry for the given
+ * user-facing thread id, creating a new one when no active entry exists or
+ * the previous one was released.
*
- * @param threadName the thread name
- * @return the thread_id (UUID string)
+ * @param threadId the user-facing thread id
+ * @return the internal thread id (UUID string)
*/
- private String getOrCreateThreadId(String threadName) {
- String metaKey = THREAD_META_PREFIX + threadName;
+ private String getOrCreateThreadId(String threadId) {
+ String metaKey = THREAD_META_PREFIX + threadId;
RMap meta = redisson.getMap(metaKey);
// Check if an active thread exists
- String threadId = meta.get(FIELD_THREAD_ID);
+ String persistedThreadId = meta.get(FIELD_THREAD_ID);
String isReleased = meta.get(FIELD_IS_RELEASED);
- if (threadId != null && !"true".equals(isReleased)) {
- // Active thread exists, return its thread_id
- return threadId;
+ if (persistedThreadId != null && !"true".equals(isReleased)) {
+ // Active thread exists, return its internal thread id
+ return persistedThreadId;
}
- // No active thread exists or thread is released, create a new thread_id
+ // No active thread exists or thread is released, create a new internal thread id
String newThreadId = UUID.randomUUID().toString();
meta.put(FIELD_THREAD_ID, newThreadId);
meta.put(FIELD_IS_RELEASED, "false");
@@ -162,7 +171,7 @@ private String getOrCreateThreadId(String threadName) {
// Set reverse mapping
String reverseKey = THREAD_REVERSE_PREFIX + newThreadId;
RMap reverse = redisson.getMap(reverseKey);
- reverse.put(FIELD_THREAD_NAME, threadName);
+ reverse.put(FIELD_THREAD_NAME, threadId);
reverse.put(FIELD_IS_RELEASED, "false");
if (ttl > 0) {
reverse.expire(java.time.Duration.ofMillis(ttlUnit.toMillis(ttl)));
@@ -172,34 +181,38 @@ private String getOrCreateThreadId(String threadName) {
}
/**
- * Gets the active thread_id for the given thread_name.
+ * Gets the internal surrogate thread id of the active entry for the given
+ * user-facing thread id.
* Returns null if no active thread exists.
*
- * @param threadName the thread name
- * @return the active thread_id, or null if not found
+ * @param threadId the user-facing thread id
+ * @return the active internal thread id, or null if not found
*/
- private String getActiveThreadId(String threadName) {
- String metaKey = THREAD_META_PREFIX + threadName;
+ private String getActiveThreadId(String threadId) {
+ String metaKey = THREAD_META_PREFIX + threadId;
RMap meta = redisson.getMap(metaKey);
- String threadId = meta.get(FIELD_THREAD_ID);
+ String persistedThreadId = meta.get(FIELD_THREAD_ID);
String isReleased = meta.get(FIELD_IS_RELEASED);
- if (threadId != null && !"true".equals(isReleased)) {
- return threadId;
+ if (persistedThreadId != null && !"true".equals(isReleased)) {
+ return persistedThreadId;
}
return null; // No active thread exists
}
- private String threadName(RunnableConfig config) {
+ /**
+ * Resolves the user-facing thread id every saver operation is keyed on.
+ */
+ private String threadId(RunnableConfig config) {
return checkpointThreadId(config);
}
@Override
public Collection list(RunnableConfig config) {
- String threadName = threadName(config);
- RLock lock = redisson.getLock(LOCK_PREFIX + threadName);
+ String threadId = threadId(config);
+ RLock lock = redisson.getLock(LOCK_PREFIX + threadId);
boolean tryLock = false;
try {
// 500ms timeout for read operations (list)
@@ -208,14 +221,14 @@ public Collection list(RunnableConfig config) {
return List.of();
}
- // Get active thread_id for the thread_name
- String threadId = getActiveThreadId(threadName);
- if (threadId == null) {
+ // Get the internal thread id of the active entry
+ String persistedThreadId = getActiveThreadId(threadId);
+ if (persistedThreadId == null) {
return List.of();
}
- // Use thread_id to query checkpoints
- RBucket bucket = redisson.getBucket(CHECKPOINT_PREFIX + threadId);
+ // Use the internal thread id to query checkpoints
+ RBucket bucket = redisson.getBucket(CHECKPOINT_PREFIX + persistedThreadId);
String content = bucket.get();
return deserializeCheckpoints(content);
@@ -235,8 +248,8 @@ public Collection list(RunnableConfig config) {
@Override
public Optional get(RunnableConfig config) {
- String threadName = threadName(config);
- RLock lock = redisson.getLock(LOCK_PREFIX + threadName);
+ String threadId = threadId(config);
+ RLock lock = redisson.getLock(LOCK_PREFIX + threadId);
boolean tryLock = false;
try {
// 500ms timeout for read operations (get)
@@ -245,14 +258,14 @@ public Optional get(RunnableConfig config) {
return Optional.empty();
}
- // Get active thread_id for the thread_name
- String threadId = getActiveThreadId(threadName);
- if (threadId == null) {
+ // Get the internal thread id of the active entry
+ String persistedThreadId = getActiveThreadId(threadId);
+ if (persistedThreadId == null) {
return Optional.empty();
}
- // Use thread_id to query checkpoints
- RBucket bucket = redisson.getBucket(CHECKPOINT_PREFIX + threadId);
+ // Use the internal thread id to query checkpoints
+ RBucket bucket = redisson.getBucket(CHECKPOINT_PREFIX + persistedThreadId);
String content = bucket.get();
LinkedList checkpoints = deserializeCheckpoints(content);
@@ -280,21 +293,21 @@ public Optional get(RunnableConfig config) {
@Override
public RunnableConfig put(RunnableConfig config, Checkpoint checkpoint) throws Exception {
- String threadName = threadName(config);
- RLock lock = redisson.getLock(LOCK_PREFIX + threadName);
+ String threadId = threadId(config);
+ RLock lock = redisson.getLock(LOCK_PREFIX + threadId);
boolean tryLock = false;
try {
// 3 seconds timeout for write operations (put) - longer timeout for concurrent scenarios
tryLock = lock.tryLock(3, TimeUnit.SECONDS);
if (!tryLock) {
- throw new RuntimeException("Failed to acquire lock for thread: " + threadName);
+ throw new RuntimeException("Failed to acquire lock for thread: " + threadId);
}
- // Get or create thread_id
- String threadId = getOrCreateThreadId(threadName);
+ // Get or create the internal thread id
+ String persistedThreadId = getOrCreateThreadId(threadId);
- // Use thread_id as key for checkpoint storage
- RBucket bucket = redisson.getBucket(CHECKPOINT_PREFIX + threadId);
+ // Use the internal thread id as key for checkpoint storage
+ RBucket bucket = redisson.getBucket(CHECKPOINT_PREFIX + persistedThreadId);
String content = bucket.get();
LinkedList checkpoints = deserializeCheckpoints(content);
@@ -335,41 +348,41 @@ public RunnableConfig put(RunnableConfig config, Checkpoint checkpoint) throws E
@Override
public Tag release(RunnableConfig config) throws Exception {
- String threadName = threadName(config);
- RLock lock = redisson.getLock(LOCK_PREFIX + threadName);
+ String threadId = threadId(config);
+ RLock lock = redisson.getLock(LOCK_PREFIX + threadId);
boolean tryLock = false;
try {
// 3 seconds timeout for write operations (release) - longer timeout for concurrent scenarios
tryLock = lock.tryLock(3, TimeUnit.SECONDS);
if (!tryLock) {
- throw new RuntimeException("Failed to acquire lock for thread: " + threadName);
+ throw new RuntimeException("Failed to acquire lock for thread: " + threadId);
}
- String metaKey = THREAD_META_PREFIX + threadName;
+ String metaKey = THREAD_META_PREFIX + threadId;
RMap meta = redisson.getMap(metaKey);
- String threadId = meta.get(FIELD_THREAD_ID);
- if (threadId == null) {
- throw new IllegalStateException("Thread not found: " + threadName);
+ String persistedThreadId = meta.get(FIELD_THREAD_ID);
+ if (persistedThreadId == null) {
+ throw new IllegalStateException("Thread not found: " + threadId);
}
// Mark thread as released
meta.put(FIELD_IS_RELEASED, "true");
// Update reverse mapping
- String reverseKey = THREAD_REVERSE_PREFIX + threadId;
+ String reverseKey = THREAD_REVERSE_PREFIX + persistedThreadId;
RMap reverse = redisson.getMap(reverseKey);
if (reverse != null) {
reverse.put(FIELD_IS_RELEASED, "true");
}
- // Get checkpoints for Tag (using thread_id)
- String contentKey = CHECKPOINT_PREFIX + threadId;
+ // Get checkpoints for Tag (using the internal thread id)
+ String contentKey = CHECKPOINT_PREFIX + persistedThreadId;
RBucket bucket = redisson.getBucket(contentKey);
String content = bucket.get();
Collection checkpoints = deserializeCheckpoints(content);
- return new Tag(threadName, checkpoints);
+ return new Tag(threadId, checkpoints);
}
catch (InterruptedException e) {