From 56dabe84560927f59d64eb37a797ac1ae0a46b1b Mon Sep 17 00:00:00 2001
From: Jake Vanderwerf <get@jakevanderwerf.ca>
Date: Tue, 12 May 2026 18:58:21 +0000
Subject: [PATCH] =added state schema.org class

---
 inc/managers/queue/Processor.php |   96 ++++++++++++++++++------------------------------
 1 files changed, 36 insertions(+), 60 deletions(-)

diff --git a/inc/managers/queue/Processor.php b/inc/managers/queue/Processor.php
index a3eac9d..ccc4ddb 100644
--- a/inc/managers/queue/Processor.php
+++ b/inc/managers/queue/Processor.php
@@ -14,14 +14,23 @@
 
 	public function run(): void
 	{
+		if (get_transient(BASE.'queue_running')) {
+			return;
+		}
+		set_transient(BASE.'queue_running', true, 60);
 		if (!$this->hasAdequateResources()) {
 			error_log('[Processor] Insufficient resources to start processing');
 			return;
 		}
 
-		$ops = $this->storage->fetchRunnable(3);
-
+		$ops = $this->storage->fetchRunnable();
+		if (empty($ops)) {
+			return;
+		}
 		foreach ($ops as $op) {
+			if ($op->state === 'completed') {
+				return;
+			}
 			if (!$this->dependenciesSatisfied($op)) {
 				continue;
 			}
@@ -37,6 +46,10 @@
 
 	private function processOne(Operation $op): void
 	{
+		if (get_transient(BASE.$op->id)) {
+			return;
+		}
+		set_transient(BASE.$op->id, true, 500);
 		$progress = new Progress($op);
 
 		$executor = $this->registry->getExecutor($op->type) ?? $this->defaultExecutor;
@@ -45,28 +58,13 @@
 
 		$this->storage->saveProgress($op);
 
-		//Check to see if we can merge first
-		$mergeable = $this->registry->getMergeable($op->type);
-
-		if ($mergeable) {
-			$existing = $this->storage->findMergeable(
-				$op->type,
-				$op->userId
-			);
-
-			if ($existing && $mergeable->canMerge($existing, $op)) {
-				$this->applyMerge($mergeable, $existing, $op);
-				return;
-			}
-		}
-
 		try {
-			// Check if this operation should be chunked
 			$chunkKey = $op->metadata['chunk_key'] ?? null;
 
+			// No transaction wrapping — executor handles its own
 			$result = $chunkKey
 				? $this->executeChunked($op, $executor, $progress, $chunkKey)
-				: $this->storage->withTransaction(fn() => $executor->execute($op, $progress));
+				: $executor->execute($op, $progress);
 
 			$op->state = 'completed';
 			$op->outcome = $result->outcome;
@@ -78,7 +76,16 @@
 			$this->handleFailure($op, $e);
 		}
 
-		$this->storage->saveFinal($op);
+		$this->saveOperation($op);
+	}
+	private function saveOperation(Operation $op): void
+	{
+		if ($op->state === 'completed') {
+			$this->storage->saveFinal($op);
+		} else {
+			// Retryable failure — save as scheduled/failed without requiring 'completed' state
+			$this->storage->save($op);
+		}
 	}
 
 	private function executeChunked(Operation $op, Executor $executor, Progress $progress, string|array $chunkKey): Result
@@ -116,23 +123,20 @@
 			}
 
 			try {
-				$chunkResult = $this->storage->withTransaction(function () use ($op, $executor, $progress, $chunk, $keys, $index) {
-					// Clone operation with only this chunk's data
-					$chunkOp = clone $op;
-					$chunkOp->requestData = array_merge(
-						array_diff_key($op->requestData, array_flip($keys)),
-						$chunk['data']
-					);
+				$chunkOp = clone $op;
+				$chunkOp->requestData = array_merge(
+					array_diff_key($op->requestData, array_flip($keys)),
+					$chunk['data']
+				);
 
-					// Execute this chunk
+				$executeChunk = function () use ($op, $executor, $progress, $chunkOp, $index) {
 					$result = $executor->execute($chunkOp, $progress);
-
-					// Update progress
 					$op->metadata['chunk_offset'] = $index + 1;
 					$this->storage->saveProgress($op);
-
 					return $result;
-				});
+				};
+
+				$chunkResult = $executeChunk();
 
 				// Merge results
 				if (!empty($chunkResult->result)) {
@@ -243,33 +247,6 @@
 		return date('Y-m-d H:i:s', time() + $delay + $jitter);
 	}
 
-	private function applyMerge(
-		Mergeable $mergeable,
-		Operation $target,
-		Operation $incoming
-	): void {
-		// Safety: only merge into actively processing ops
-		if ($target->state !== 'processing') {
-			return;
-		}
-
-		$this->storage->withTransaction(function () use ($mergeable, $target, $incoming) {
-			$mergeable->merge($target, $incoming);
-
-			$target->dependencies[] = $incoming->id;
-			$target->dependencies = array_values(array_unique($target->dependencies));
-
-			$this->storage->saveProgress($target);
-
-			$incoming->state = 'completed';
-			$incoming->outcome = 'success';
-			$incoming->completedAt = current_time('mysql');
-			$incoming->result = ['merged_into' => $target->id];
-
-			$this->storage->saveFinal($incoming);
-		});
-	}
-
 	private function checkResourceLimits(): bool
 	{
 		// Check memory (leave 20% buffer)
@@ -336,7 +313,6 @@
 		if (empty($op->dependencies)) {
 			return true;
 		}
-
 		foreach ($op->dependencies as $depId) {
 			$dep = $this->storage->find($depId);
 

--
Gitblit v1.10.0