Files
coolify/app/Jobs/DatabaseStartJob.php
T
Andras BacsaiandClaude Opus 5.5 a8117f9114 fix(databases): unblock interrupted imports and prevent double starts
- A Coolify restart now fails queued or running database imports, so
  they no longer block all later imports with 409.
- An import without progress for the SSH command timeout plus 30
  minutes (at least 2 hours) is stale and no longer blocks.
- The import sets its operation property when the activity is created,
  so a worker cannot save the activity without it.
- Start and restart refuse a second operation while a start, restart,
  or import is in progress, in the UI, the API (409), and actions.
- DatabaseStartJob runs only while its activity is still queued and
  the newest one, and holds a lock per database during the commands.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 10:40:53 +02:00

149 lines
5.2 KiB
PHP

<?php
namespace App\Jobs;
use App\Actions\Database\StartClickhouse;
use App\Actions\Database\StartDragonfly;
use App\Actions\Database\StartKeydb;
use App\Actions\Database\StartMariadb;
use App\Actions\Database\StartMongodb;
use App\Actions\Database\StartMysql;
use App\Actions\Database\StartPostgresql;
use App\Actions\Database\StartRedis;
use App\Enums\ProcessStatus;
use App\Events\DatabaseStatusChanged;
use App\Exceptions\DatabaseStartException;
use App\Models\StandaloneClickhouse;
use App\Models\StandaloneDragonfly;
use App\Models\StandaloneKeydb;
use App\Models\StandaloneMariadb;
use App\Models\StandaloneMongodb;
use App\Models\StandaloneMysql;
use App\Models\StandalonePostgresql;
use App\Models\StandaloneRedis;
use App\Support\ResourceStartActivity;
use Illuminate\Bus\Queueable;
use Illuminate\Contracts\Queue\ShouldBeEncrypted;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Database\Eloquent\Model;
use Illuminate\Foundation\Bus\Dispatchable;
use Illuminate\Queue\InteractsWithQueue;
use Illuminate\Queue\SerializesModels;
use Illuminate\Support\Facades\Cache;
use Spatie\Activitylog\Models\Activity;
use Throwable;
class DatabaseStartJob implements ShouldBeEncrypted, ShouldQueue
{
use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;
public int $tries = 1;
public int $timeout = 600;
public function __construct(
public string $databaseClass,
public int $databaseId,
public int $teamId,
public int $activityId,
public ?int $userId,
) {
$this->onQueue(deployment_queue());
}
public function handle(): void
{
$database = $this->databaseClass::query()->findOrFail($this->databaseId);
abort_unless((int) $database->team()->id === $this->teamId, 403);
$activity = Activity::query()->findOrFail($this->activityId);
if (! $this->isStillRequested($activity, $database->uuid)) {
event(new DatabaseStatusChanged($this->userId));
return;
}
$lock = Cache::lock(ResourceStartActivity::databaseStartLockKey($database->uuid), ResourceStartActivity::DATABASE_START_LOCK_SECONDS);
if (! $lock->get()) {
ResourceStartActivity::markFailed($activity, ResourceStartActivity::ALREADY_STARTING_MESSAGE);
event(new DatabaseStatusChanged($this->userId));
return;
}
try {
$this->runStart($database, $activity);
} finally {
$lock->release();
}
}
/**
* A start runs only while its activity is still queued and no newer start of the same
* database was requested. A start that waited in a busy queue can be treated as stale and
* replaced by a new one; without this check both would run one after the other.
*/
private function isStillRequested(Activity $activity, string $databaseUuid): bool
{
if (data_get($activity, 'properties.status') !== ProcessStatus::QUEUED->value) {
return false;
}
$newerStartExists = Activity::query()
->where('properties->type_uuid', $databaseUuid)
->where('properties->operation', ResourceStartActivity::DATABASE_START_OPERATION)
->whereIn('properties->status', [
ProcessStatus::QUEUED->value,
ProcessStatus::IN_PROGRESS->value,
ProcessStatus::FINISHED->value,
])
->where('id', '>', $activity->getKey())
->exists();
if ($newerStartExists) {
ResourceStartActivity::markFailed($activity, ResourceStartActivity::SUPERSEDED_MESSAGE);
return false;
}
return true;
}
private function runStart(Model $database, Activity $activity): void
{
$result = match ($database->getMorphClass()) {
StandalonePostgresql::class => StartPostgresql::run($database, $activity),
StandaloneRedis::class => StartRedis::run($database, $activity),
StandaloneMongodb::class => StartMongodb::run($database, $activity),
StandaloneMysql::class => StartMysql::run($database, $activity),
StandaloneMariadb::class => StartMariadb::run($database, $activity),
StandaloneKeydb::class => StartKeydb::run($database, $activity),
StandaloneDragonfly::class => StartDragonfly::run($database, $activity),
StandaloneClickhouse::class => StartClickhouse::run($database, $activity),
};
if (! $result instanceof Activity || data_get($result, 'properties.status') !== ProcessStatus::FINISHED->value) {
throw DatabaseStartException::startCommandsDidNotRun();
}
event(new DatabaseStatusChanged($this->userId));
}
public function failed(?Throwable $exception): void
{
try {
$activity = Activity::query()->find($this->activityId);
if (! $activity) {
return;
}
ResourceStartActivity::markFailed(
$activity,
$exception instanceof DatabaseStartException ? $exception->getMessage() : 'Database start failed.',
);
} finally {
event(new DatabaseStatusChanged($this->userId));
}
}
}