matab-panel/app/Services/SyncService.php

451 lines
18 KiB
PHP

<?php
declare(strict_types=1);
namespace App\Services;
use App\Models\OsurgInitial;
use App\Models\SyncLog;
use Illuminate\Support\Facades\DB;
use Illuminate\Support\Facades\Http;
use Illuminate\Support\Facades\Log;
use Illuminate\Support\Facades\Storage;
class SyncService
{
public function __construct(
private readonly string $peerIp,
private readonly int $peerPort,
private readonly string $token,
private readonly string $ftpUser,
private readonly string $ftpPass,
private readonly int $ftpPort,
private readonly string $ftpRemotePath,
) {}
public static function fromConfig(): static
{
return new static(
peerIp: OsurgInitial::val('sync_peer_ip', config('sync.peer_ip', '')),
peerPort: (int) OsurgInitial::val('sync_peer_port', config('sync.peer_port', 80)),
token: OsurgInitial::val('sync_token', config('sync.token', '')),
ftpUser: OsurgInitial::val('ftp_user', config('sync.ftp_user', '')),
ftpPass: OsurgInitial::val('ftp_pass', config('sync.ftp_pass', '')),
ftpPort: (int) OsurgInitial::val('ftp_port', config('sync.ftp_port', 2121)),
ftpRemotePath: OsurgInitial::val('ftp_remote_path', config('sync.ftp_remote_path', '')),
);
}
public function isConfigured(): bool
{
return $this->peerIp !== '' && $this->token !== '';
}
public function isFtpConfigured(): bool
{
return $this->peerIp !== '' && $this->ftpUser !== '' && $this->ftpRemotePath !== '';
}
public function sync(): array
{
$report = [];
$baseUrl = "http://{$this->peerIp}:{$this->peerPort}";
$peerCursor = (int) \App\Models\OsurgInitial::val('sync_peer_cursor', 0);
$ourLastSync = SyncLog::lastSyncId();
try {
$pullResponse = Http::timeout(30)
->withHeader('X-Sync-Token', $this->token)
->get("{$baseUrl}/api/sync/export", ['since' => $peerCursor]);
if ($pullResponse->successful()) {
$peerChanges = $pullResponse->json('changes', []);
if (! empty($peerChanges)) {
app()->instance('sync.applying', true);
DB::transaction(function () use ($peerChanges, &$report) {
foreach ($peerChanges as $change) {
$this->applyChange($change, 'دریافت', $report);
}
});
$maxPeerChangeId = max(array_column($peerChanges, 'id'));
$this->savePeerCursor(max($peerCursor, $maxPeerChangeId));
if ($this->isFtpConfigured()) {
$this->transferFilesViaFtp($peerChanges, 's2c', $report);
} else {
$this->syncFilesViaHttp($peerChanges, $baseUrl, $report);
}
}
$report[] = count($peerChanges) . ' تغییر از سیستم مقابل دریافت شد.';
} else {
$report[] = 'خطا در دریافت تغییرات از سیستم مقابل: ' . $pullResponse->status();
}
} catch (\Throwable $e) {
$report[] = 'خطا در اتصال به سیستم مقابل: ' . $e->getMessage();
return ['success' => false, 'report' => $report];
}
$ourChanges = SyncLog::changesSince($ourLastSync)->toArray();
try {
$pushResponse = Http::timeout(120)
->withHeader('X-Sync-Token', $this->token)
->post("{$baseUrl}/api/sync/apply", ['changes' => $ourChanges]);
if ($pushResponse->successful()) {
foreach ($pushResponse->json('report', []) as $line) {
$report[] = 'ارسال: ' . $line;
}
$report[] = count($ourChanges) . ' تغییر به سیستم مقابل ارسال شد.';
if ($this->isFtpConfigured()) {
$this->transferFilesViaFtp($ourChanges, 'c2s', $report);
} else {
$this->pushFilesViaHttp($ourChanges, $baseUrl, $report);
}
} else {
$report[] = 'خطا در ارسال تغییرات به سیستم مقابل: ' . $pushResponse->status();
}
} catch (\Throwable $e) {
$report[] = 'خطا در ارسال به سیستم مقابل: ' . $e->getMessage();
}
SyncLog::markSynced(request()->ip() ?? '127.0.0.1');
return ['success' => true, 'report' => $report];
}
public function applyChanges(array $changes): array
{
$report = [];
app()->instance('sync.applying', true);
DB::transaction(function () use ($changes, &$report) {
foreach ($changes as $change) {
$this->applyChange($change, '', $report);
}
});
SyncLog::markSynced(request()->ip() ?? '');
return $report;
}
private function transferFilesViaFtp(array $changes, string $side, array &$report): void
{
if (empty($changes)) {
return;
}
$ftp = $this->connectFtp();
if (! $ftp) {
$report[] = 'خطا در اتصال FTP به سیستم مقابل';
Log::error('FTP connection failed in transferFilesViaFtp');
return;
}
$fileFields = ['photos_before', 'photos_after', 'videos', 'audio_files'];
$localBase = rtrim(Storage::disk('public')->path(''), '/\\');
$remoteBase = rtrim($this->ftpRemotePath, '/\\');
foreach ($changes as $change) {
$data = $change['changed_data'] ?? [];
foreach ($fileFields as $field) {
if (empty($data[$field])) {
continue;
}
$files = is_array($data[$field]) ? $data[$field] : json_decode($data[$field], true) ?? [];
foreach ($files as $filePath) {
if (! $filePath) {
continue;
}
$filePath = ltrim(str_replace('\\', '/', $filePath), '/');
$localFile = $localBase . '/' . $filePath;
$remoteFile = $remoteBase . '/' . $filePath;
$remoteFolder = dirname($remoteFile);
if ($side === 'c2s') {
if (! file_exists($localFile)) {
$errorMsg = "فایل [{$filePath}] در سیستم محلی یافت نشد";
$report[] = $errorMsg . ' ✗';
Log::warning('FTP upload skipped - local file not found', [
'file' => $filePath,
'local_path' => $localFile,
]);
continue;
}
if (! $this->createFtpDirectory($ftp, $remoteFolder)) {
$errorMsg = "ایجاد پوشه [{$remoteFolder}] برای فایل [{$filePath}] شکست خورد";
$report[] = $errorMsg . ' ✗';
Log::error('FTP mkdir failed', [
'file' => $filePath,
'remote_folder' => $remoteFolder,
]);
continue;
}
if (ftp_put($ftp, $remoteFile, $localFile, FTP_BINARY)) {
$report[] = "فایل [{$filePath}] ارسال شد (FTP) ✓";
Log::info('FTP file uploaded', [
'file' => $filePath,
'size' => filesize($localFile),
]);
} else {
$errorDetail = error_get_last();
$errorMsg = "فایل [{$filePath}] ارسال نشد (FTP)";
if ($errorDetail) {
$errorMsg .= " - " . $errorDetail['message'];
}
$report[] = $errorMsg . ' ✗';
Log::error('FTP upload failed', [
'file' => $filePath,
'local_file' => $localFile,
'remote_file' => $remoteFile,
'error' => $errorDetail,
]);
}
} else {
if (file_exists($localFile)) {
Log::debug('FTP download skipped - file already exists', [
'file' => $filePath,
]);
continue;
}
$localFolder = dirname($localFile);
if (! is_dir($localFolder)) {
if (! mkdir($localFolder, 0755, true)) {
$report[] = "ایجاد پوشه [{$localFolder}] شکست خورد ✗";
Log::error('Local mkdir failed', [
'folder' => $localFolder,
'error' => error_get_last(),
]);
continue;
}
}
if (ftp_get($ftp, $localFile, $remoteFile, FTP_BINARY)) {
$report[] = "فایل [{$filePath}] دریافت شد (FTP) ✓";
Log::info('FTP file downloaded', [
'file' => $filePath,
'size' => filesize($localFile),
]);
} else {
$errorDetail = error_get_last();
$errorMsg = "فایل [{$filePath}] دریافت نشد (FTP)";
if ($errorDetail) {
$errorMsg .= " - " . $errorDetail['message'];
}
$report[] = $errorMsg . ' ✗';
Log::error('FTP download failed', [
'file' => $filePath,
'local_file' => $localFile,
'remote_file' => $remoteFile,
'error' => $errorDetail,
]);
}
}
}
}
}
ftp_close($ftp);
Log::info('FTP connection closed');
}
private function connectFtp(): mixed
{
if (! function_exists('ftp_connect')) {
Log::error('FTP extension is not available');
return null;
}
$ftp = ftp_connect($this->peerIp, $this->ftpPort, 30);
if (! $ftp) {
Log::error('FTP connection failed', [
'host' => $this->peerIp,
'port' => $this->ftpPort,
'error' => error_get_last(),
]);
return null;
}
if (! ftp_login($ftp, $this->ftpUser, $this->ftpPass)) {
Log::error('FTP login failed', [
'host' => $this->peerIp,
'port' => $this->ftpPort,
'user' => $this->ftpUser,
'error' => error_get_last(),
]);
ftp_close($ftp);
return null;
}
ftp_pasv($ftp, true);
Log::info('FTP connection established', [
'host' => $this->peerIp,
'port' => $this->ftpPort,
]);
return $ftp;
}
private function createFtpDirectory($ftp, string $path): bool
{
$parts = explode('/', trim($path, '/'));
$currentPath = '';
foreach ($parts as $part) {
if (empty($part)) {
continue;
}
$currentPath .= '/' . $part;
$originalDir = ftp_pwd($ftp);
if (@ftp_chdir($ftp, $currentPath)) {
@ftp_chdir($ftp, $originalDir);
continue;
}
if (! @ftp_mkdir($ftp, $currentPath)) {
Log::error('FTP mkdir failed', [
'path' => $currentPath,
'full_path' => $path,
'error' => error_get_last(),
]);
return false;
}
Log::debug('FTP directory created', ['path' => $currentPath]);
}
return true;
}
private function pushFilesViaHttp(array $changes, string $baseUrl, array &$report): void
{
$fileFields = ['photos_before', 'photos_after', 'videos', 'audio_files'];
foreach ($changes as $change) {
$data = $change['changed_data'] ?? [];
foreach ($fileFields as $field) {
if (empty($data[$field])) {
continue;
}
$files = is_array($data[$field]) ? $data[$field] : json_decode($data[$field], true) ?? [];
foreach ($files as $filePath) {
if (! $filePath || ! Storage::disk('public')->exists($filePath)) {
continue;
}
try {
$response = Http::timeout(600)
->withHeader('X-Sync-Token', $this->token)
->attach('file', Storage::disk('public')->get($filePath), basename($filePath))
->post("{$baseUrl}/api/sync/receive-file", ['path' => $filePath]);
if ($response->successful()) {
$report[] = "فایل [{$filePath}] ارسال شد (HTTP) ✓";
} else {
$report[] = "فایل [{$filePath}] ارسال نشد (HTTP) ✗ " . $response->status();
}
} catch (\Throwable $e) {
$report[] = "فایل [{$filePath}] ارسال نشد (HTTP) ✗ " . $e->getMessage();
}
}
}
}
}
private function syncFilesViaHttp(array $changes, string $baseUrl, array &$report): void
{
$fileFields = ['photos_before', 'photos_after', 'videos', 'audio_files'];
foreach ($changes as $change) {
$data = $change['changed_data'] ?? [];
foreach ($fileFields as $field) {
if (empty($data[$field])) {
continue;
}
$files = is_array($data[$field]) ? $data[$field] : json_decode($data[$field], true) ?? [];
foreach ($files as $filePath) {
if (! $filePath || Storage::disk('public')->exists($filePath)) {
continue;
}
try {
$response = Http::timeout(60)
->withHeader('X-Sync-Token', $this->token)
->get("{$baseUrl}/api/sync/file", ['path' => $filePath]);
if ($response->successful()) {
Storage::disk('public')->put($filePath, $response->body());
$report[] = "فایل [{$filePath}] دریافت شد (HTTP) ✓";
}
} catch (\Throwable $e) {
$report[] = "فایل [{$filePath}] دریافت نشد (HTTP) ✗ " . $e->getMessage();
}
}
}
}
}
private function applyChange(array $change, string $prefix, array &$report): void
{
$table = $change['table_name'];
$recordId = $change['record_id'];
$action = $change['action'];
$data = $change['changed_data'] ?? [];
$label = $prefix ? "{$prefix} [{$table}#{$recordId}]" : "[{$table}#{$recordId}]";
try {
if ($action === 'deleted') {
DB::table($table)->where('id', $recordId)->delete();
$report[] = "{$label} حذف ✓";
} elseif ($action === 'created') {
$exists = DB::table($table)->where('id', $recordId)->exists();
if ($exists) {
DB::table($table)->where('id', $recordId)->update($this->sanitize($data));
$report[] = "{$label} ایجاد→بروزرسانی ✓";
} else {
DB::table($table)->insert($this->sanitize(array_merge(['id' => $recordId], $data)));
$report[] = "{$label} ایجاد ✓";
}
} elseif ($action === 'updated') {
DB::table($table)->where('id', $recordId)->update($this->sanitize($data));
$report[] = "{$label} بروزرسانی ✓";
}
} catch (\Throwable $e) {
$report[] = "{$label} {$action}" . $e->getMessage();
}
}
private function savePeerCursor(int $cursor): void
{
\App\Models\OsurgInitial::updateOrCreate(
['init_parameter' => 'sync_peer_cursor'],
['init_value' => $cursor],
);
cache()->forget('osurg_initial.sync_peer_cursor');
}
private function sanitize(array $data): array
{
$result = [];
foreach ($data as $key => $value) {
if (is_array($value)) {
$result[$key] = json_encode($value, JSON_UNESCAPED_UNICODE);
} elseif (is_scalar($value) || $value === null) {
$result[$key] = $value;
}
}
return $result;
}
}