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