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) {