diff --git a/lib-cmd-queue-redis/changelog.txt b/lib-cmd-queue-redis/changelog.txt index 9f893f41..a816e073 100644 --- a/lib-cmd-queue-redis/changelog.txt +++ b/lib-cmd-queue-redis/changelog.txt @@ -1,6 +1,20 @@ # lib-cmd-queue-redis changelog 0.4.1 - 31 Jul 2026 +- Add error-tracking fields to CommandState for observability of a retry storm on a command that + stays retryable (i.e. errors that do not terminally fail it, since a thrown handler is now + retried — see the fix below): errorsCount (count of consecutive processing errors since the last + successful processing) and modifiedAt (last-write timestamp). The existing `error` field now also + holds the message of a transient (non-terminal) processing error — it carries the most recent + error message, transient or terminal; the terminal failure is identified by status == FAILED, not + by error being non-null. CommandServiceImpl increments the count / records the message on each + caught handler exception (best-effort — a failed record never changes control flow), and resets + the streak on any successful transition or recovery. Backward-compatible: the new fields default + to 0/null when an older serialized CommandState is read. +- This restores the error tracking originally released as 0.6.0 (#89, e54229e), which was dropped + by the revert to 0.4.0 below. It is the companion to the retry fix below: once a thrown handler + is retried rather than terminal-failed, a command can retry indefinitely, and these fields are + what make that visible. - Fix: a handler that throws is no longer treated as a terminal command failure. The catch in CommandServiceImpl.processCommandWithHandler now returns false — keeping the message unacked so the stream layer re-polls it — instead of persisting FAILED and acking (which removed the diff --git a/lib-cmd-queue-redis/src/main/java/io/seqera/data/command/CommandServiceImpl.java b/lib-cmd-queue-redis/src/main/java/io/seqera/data/command/CommandServiceImpl.java index 06eff28e..698b1e11 100644 --- a/lib-cmd-queue-redis/src/main/java/io/seqera/data/command/CommandServiceImpl.java +++ b/lib-cmd-queue-redis/src/main/java/io/seqera/data/command/CommandServiceImpl.java @@ -281,6 +281,10 @@ private
boolean processCommandWithHandler( // Ensure state reflects RUNNING status for accurate reporting if (state.status() != CommandStatus.RUNNING) { store.save(state.started()); + } else if (state.errorsCount() > 0) { + // Recovered after one or more transient errors — reset the streak. Single write, + // and only when there is something to reset, so healthy re-polls stay write-free. + store.save(state.clearErrors()); } return false; // Keep in queue - will retry and call checkStatus() } @@ -302,6 +306,7 @@ private
boolean processCommandWithHandler( // non-terminal, stranding the work. Deciding a command has *permanently* failed is // delegated to the domain layer that owns the entity state (see seqeralabs/sched#712). log.error("Command processing errored, will retry: id={}", msg.commandId(), e); + recordError(state, e); return false; // Keep in queue - redelivered / re-polled } } @@ -342,4 +347,34 @@ private
CommandResult handler,
throw new RuntimeException("Command execution failed", e);
}
}
+
+ /**
+ * Best-effort: record a non-terminal processing error on the command state — increment the
+ * consecutive-error count and capture the message — for observability of a retry storm on a
+ * command that stays retryable. A failure to persist this must not change control flow: the
+ * command is kept in the queue and retried regardless.
+ */
+ private void recordError(CommandState state, Exception e) {
+ try {
+ store.save(state.withError(rootMessage(e)));
+ } catch (Exception fail) {
+ log.warn("Failed to record command error state: id={}", state.id(), fail);
+ }
+ }
+
+ /**
+ * The most specific message available for a processing error. {@link #executeWithTimeout}
+ * wraps a handler exception in a generic {@code RuntimeException("Command execution failed")},
+ * so the root cause's message is recorded instead — otherwise every transient error on the
+ * execute() path would read "Command execution failed" and the field would be useless for
+ * diagnosing a retry storm. The checkStatus() path throws directly, where the root cause is
+ * the exception itself.
+ */
+ private static String rootMessage(Throwable e) {
+ Throwable root = e;
+ while (root.getCause() != null && root.getCause() != root) {
+ root = root.getCause();
+ }
+ return root.getMessage() != null ? root.getMessage() : root.toString();
+ }
}
diff --git a/lib-cmd-queue-redis/src/main/java/io/seqera/data/command/CommandState.java b/lib-cmd-queue-redis/src/main/java/io/seqera/data/command/CommandState.java
index c922eb01..8a79a248 100644
--- a/lib-cmd-queue-redis/src/main/java/io/seqera/data/command/CommandState.java
+++ b/lib-cmd-queue-redis/src/main/java/io/seqera/data/command/CommandState.java
@@ -14,7 +14,6 @@
* limitations under the License.
*
*/
-
package io.seqera.data.command;
import java.time.Instant;
@@ -26,6 +25,28 @@
* Persistent state of a command, stored as JSON in the database.
* Uses @JsonTypeInfo to preserve type information for params and result
* during serialization, enabling proper deserialization without explicit type knowledge.
+ *
+ * {@code errorsCount} counts processing errors that did not terminally fail the
+ * command — a handler that threw is retried (see {@code CommandServiceImpl}), so this records how
+ * many consecutive times it has thrown, for observability of a retry storm on an otherwise
+ * non-terminal command. {@code error} holds the message of the most recent error, transient or
+ * terminal — check {@code status == FAILED} to tell a terminal failure from a transient one, not
+ * {@code error != null}. {@code modifiedAt} is refreshed on every state write, giving a
+ * last-touched timestamp.
+ *
+ * @param id command id
+ * @param type command type discriminator
+ * @param status current lifecycle status
+ * @param params command parameters (polymorphic, type preserved via {@code @JsonTypeInfo})
+ * @param result terminal result payload, if any (polymorphic)
+ * @param error message of the most recent error, transient or terminal (nullable); terminal only
+ * when {@code status == FAILED}
+ * @param errorsCount number of consecutive processing errors since the last successful
+ * processing; reset to 0 on any successful transition or recovery
+ * @param createdAt when the command was first submitted
+ * @param startedAt when the command first transitioned to RUNNING (nullable)
+ * @param modifiedAt when the command state was last written (nullable for pre-existing records)
+ * @param completedAt when the command reached a terminal state (nullable)
*/
public record CommandState(
String id,
@@ -36,8 +57,10 @@ public record CommandState(
@JsonTypeInfo(use = JsonTypeInfo.Id.CLASS)
@Nullable Object result,
@Nullable String error,
+ int errorsCount,
Instant createdAt,
@Nullable Instant startedAt,
+ @Nullable Instant modifiedAt,
@Nullable Instant completedAt
) {
@@ -45,19 +68,22 @@ public record CommandState(
* Create a new submitted command state.
*/
public static CommandState submitted(String id, String type, Object params) {
+ final Instant now = Instant.now();
return new CommandState(
id, type, CommandStatus.SUBMITTED, params,
- null, null, Instant.now(), null, null
+ null, null, 0, now, null, now, null
);
}
/**
- * Transition to RUNNING status.
+ * Transition to RUNNING status. A successful (non-throwing) transition, so the
+ * consecutive-error streak is reset.
*/
public CommandState started() {
+ final Instant now = Instant.now();
return new CommandState(
id, type, CommandStatus.RUNNING, params,
- result, error, createdAt, Instant.now(), completedAt
+ result, error, 0, createdAt, now, now, completedAt
);
}
@@ -65,9 +91,10 @@ public CommandState started() {
* Transition to SUCCEEDED status with result.
*/
public CommandState completed(Object result) {
+ final Instant now = Instant.now();
return new CommandState(
id, type, CommandStatus.SUCCEEDED, params,
- result, null, createdAt, startedAt, Instant.now()
+ result, null, 0, createdAt, startedAt, now, now
);
}
@@ -75,9 +102,10 @@ public CommandState completed(Object result) {
* Transition to FAILED status with error.
*/
public CommandState failed(String error) {
+ final Instant now = Instant.now();
return new CommandState(
id, type, CommandStatus.FAILED, params,
- null, error, createdAt, startedAt, Instant.now()
+ null, error, errorsCount, createdAt, startedAt, now, now
);
}
@@ -85,9 +113,33 @@ public CommandState failed(String error) {
* Transition to CANCELLED status.
*/
public CommandState cancelled() {
+ final Instant now = Instant.now();
return new CommandState(
id, type, CommandStatus.CANCELLED, params,
- null, null, createdAt, startedAt, Instant.now()
+ null, null, 0, createdAt, startedAt, now, now
+ );
+ }
+
+ /**
+ * Record a non-terminal processing error: keep the current status (the command stays retryable),
+ * increment the consecutive-error count, capture the message, and refresh {@code modifiedAt}.
+ * Called when a handler throws and the command is kept in the queue for retry.
+ */
+ public CommandState withError(String message) {
+ return new CommandState(
+ id, type, status, params,
+ result, message, errorsCount + 1, createdAt, startedAt, Instant.now(), completedAt
+ );
+ }
+
+ /**
+ * Clear the consecutive-error streak after a recovery, without changing status. Refreshes
+ * {@code modifiedAt}. {@code error} is retained as a historical marker of the last error seen.
+ */
+ public CommandState clearErrors() {
+ return new CommandState(
+ id, type, status, params,
+ result, error, 0, createdAt, startedAt, Instant.now(), completedAt
);
}
diff --git a/lib-cmd-queue-redis/src/test/groovy/io/seqera/data/command/CommandServiceTest.groovy b/lib-cmd-queue-redis/src/test/groovy/io/seqera/data/command/CommandServiceTest.groovy
index fbcfac60..8ccaad65 100644
--- a/lib-cmd-queue-redis/src/test/groovy/io/seqera/data/command/CommandServiceTest.groovy
+++ b/lib-cmd-queue-redis/src/test/groovy/io/seqera/data/command/CommandServiceTest.groovy
@@ -188,6 +188,26 @@ class CommandServiceTest extends Specification implements TestPropertyProvider {
then: 'the throw was treated as transient and retried to success, not persisted as FAILED'
state.status() == CommandStatus.SUCCEEDED
commandService.getResult(command.id(), TestResult).orElseThrow().message == 'Recovered'
+
+ and: 'the consecutive-error streak is reset once the command recovers'
+ state.errorsCount() == 0
+ }
+
+ def 'should track consecutive errors and last message without failing a still-retryable command'() {
+ given: 'a handler that always throws'
+ def params = new TestParams(0, 'always-throw')
+ def command = new TestCommand(TsidCreator.getTsid().toLowerCase(), 'test', params)
+
+ when: 'command is submitted and retried a few times'
+ commandService.submit(command)
+ sleep(2000)
+ def state = commandService.getState(command.id()).orElseThrow()
+
+ then: 'the command stays retryable while the error streak and last message are recorded'
+ !state.status().isTerminal()
+ state.errorsCount() >= 1
+ state.error() == 'Persistent boom'
+ state.modifiedAt() != null
}
def 'should handle unknown command type'() {
@@ -278,6 +298,10 @@ class TestCommandHandler implements CommandHandler