martinweiler commented on code in PR #2333:
URL:
https://github.com/apache/incubator-kie-kogito-apps/pull/2333#discussion_r3483478661
##########
jobs/jobs-common-embedded/src/main/java/org/kie/kogito/app/jobs/impl/VertxJobScheduler.java:
##########
@@ -435,6 +453,84 @@ public JobTimeoutExecution call() throws Exception {
return current;
}
+ /**
+ * Handles retry scheduling for RETRY and ERROR job statuses.
+ * If there's an active transaction that was marked for rollback,
schedules the retry
+ * in a new transaction. Otherwise, persists the job immediately.
+ */
+ private void scheduleRetry(JobContext jobContext, JobDetails jobDetails) {
+ String jobId = jobDetails.getId();
+ JobStatus status = jobDetails.getStatus();
+
+ // Check if we need to execute in a new transaction
+ if (transactionRollbackMarker != null &&
transactionRollbackMarker.isTransactionActive()) {
+ // Current transaction will be rolled back, so schedule retry in a
new transaction
+ LOG.trace("Job {} with status {} will be scheduled in new
transaction", jobId, status);
+
+ vertx.executeBlocking(promise -> {
+ try {
+ JobContext newJobContext = jobContextFactory.newContext();
+
+ Callable<Void> newTransactionTask = () -> {
+ jobStore.update(newJobContext, jobDetails);
+ fireEvents(jobDetails); // Fire events after rollback
+ // Schedule timer in new transaction after persistence
+ if (status == JobStatus.RETRY) {
+ addOrUpdateTxTimer(jobDetails);
+ } else {
+ removeTxTimer(jobDetails);
+ }
+ return null;
+ };
+
+ // Apply transaction interceptors to ensure this runs in a
new transaction
+ Callable<Void> wrappedTask = newTransactionTask;
+ for (JobTimeoutInterceptor interceptor : interceptors) {
+ final Callable<Void> currentTask = wrappedTask;
+ wrappedTask = () -> {
+ interceptor.chainIntercept(() -> {
+ currentTask.call();
+ return new JobTimeoutExecution(jobDetails);
+ }).call();
+ return null;
+ };
+ }
+
+ wrappedTask.call();
+ promise.complete();
+ } catch (Exception e) {
+ LOG.error("Failed to schedule retry for job {} in new
transaction", jobId, e);
+ promise.fail(e);
+ }
+ }, false, result -> {
+ if (result.failed()) {
+ LOG.error("Async retry scheduling failed for job {}",
jobId, result.cause());
+ }
+ });
+ } else {
+ // No active transaction or transaction not marked for rollback
+ LOG.trace("Job {} with status {} will be scheduled", jobId,
status);
+
+ try {
+ jobStore.update(jobContext, jobDetails);
+ // Fire events AFTER persistence to ensure Data Index can
query the updated job
+ fireEvents(jobDetails);
+ // Schedule timer AFTER persistence (no transaction
synchronization needed)
+ if (status == JobStatus.RETRY) {
+ addTimerInfo(jobDetails);
+ } else {
+ String mapKey = getMapKey(jobDetails);
+ TimerInfo timerInfo = jobsScheduled.remove(mapKey);
+ if (timerInfo != null) {
+ removeTimerInfo(timerInfo);
+ }
+ }
+ } catch (Exception e) {
+ LOG.error("Failed to schedule retry for job {}", jobId, e);
+ }
Review Comment:
You are right, both branches of the if should be aligned. I guess I went
back and forth a couple of times to get the tests working correctly and must
have changed this part as well. Thanks for spotting this!
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]