self::acquireLock($token), static fn () => self::releaseLock($token), static fn (): int => (int) (Db::name('ej_pharmacy_sync_state')->where('id', self::STATE_ID)->value('cursor') ?? 0), static fn (int $cursor, int $pageLimit): array => $client->medicines($cursor, $pageLimit), static fn (array $items, int $nextCursor): array => self::mergePage($items, $nextCursor, $token), static function (int $cursor) use ($token): void { Db::name('ej_pharmacy_sync_state')->where('id', self::STATE_ID) ->where('lock_token', $token)->update([ 'cursor' => $cursor, 'last_success_time' => time(), 'last_error_summary' => '', 'update_time' => time(), ]); }, static function (int $cursor, string $error) use ($token): void { Db::name('ej_pharmacy_sync_state')->where('id', self::STATE_ID) ->where('lock_token', $token)->update([ 'cursor' => $cursor, 'last_failure_time' => time(), 'last_error_summary' => $error, 'update_time' => time(), ]); } ); return $workflow->sync($limit); } private static function ensureStateRow(): void { if (Db::name('ej_pharmacy_sync_state')->where('id', self::STATE_ID)->find()) { return; } try { Db::name('ej_pharmacy_sync_state')->insert([ 'id' => self::STATE_ID, 'cursor' => 0, 'lock_token' => '', 'lock_expires_at' => 0, 'create_time' => time(), 'update_time' => time(), ]); } catch (Throwable $exception) { if (!self::isDuplicateKey($exception)) { throw $exception; } } } private static function acquireLock(string $token): bool { $now = time(); $updated = Db::name('ej_pharmacy_sync_state') ->where('id', self::STATE_ID) ->where(function ($query) use ($now): void { $query->where('lock_token', '')->whereOr('lock_expires_at', '<', $now); }) ->update([ 'lock_token' => $token, 'lock_expires_at' => $now + self::LOCK_TTL, 'update_time' => $now, ]); return $updated === 1; } private static function releaseLock(string $token): void { Db::name('ej_pharmacy_sync_state') ->where('id', self::STATE_ID) ->where('lock_token', $token) ->update(['lock_token' => '', 'lock_expires_at' => 0, 'update_time' => time()]); } /** @return array{created:int,updated:int,unchanged:int,deactivated:int} */ private static function mergePage(array $items, int $nextCursor, string $token): array { return Db::transaction(function () use ($items, $nextCursor, $token): array { $stats = ['created' => 0, 'updated' => 0, 'unchanged' => 0, 'deactivated' => 0]; foreach ($items as $item) { $code = trim((string) ($item['medicine_code'] ?? '')); $model = $code === '' ? null : EjMedicineCatalog::where('medicine_code', $code)->lock(true)->find(); $result = EjMedicineCatalogSyncPolicy::merge($model ? $model->toArray() : null, $item); $values = $result['values']; if ($result['action'] === 'created') { EjMedicineCatalog::create($values); } elseif ($result['action'] === 'updated' && $model) { unset($values['medicine_code']); $model->save($values); } ++$stats[$result['action']]; $stats['deactivated'] += $result['deactivated']; if ($result['deactivated'] === 1) { $now = time(); Db::name('ej_medicine_mapping') ->where('medicine_code', $code) ->where('status', 1) ->update(['status' => 0, 'delete_time' => $now, 'update_time' => $now]); } } $state = Db::name('ej_pharmacy_sync_state') ->where('id', self::STATE_ID) ->lock(true) ->find(); if (!$state || !hash_equals((string) $state['lock_token'], $token)) { throw new RuntimeException('洛阳药房目录同步锁已失效,请重试'); } Db::name('ej_pharmacy_sync_state') ->where('id', self::STATE_ID) ->update([ 'cursor' => $nextCursor, 'lock_expires_at' => time() + self::LOCK_TTL, 'update_time' => time(), ]); return $stats; }); } private static function isDuplicateKey(Throwable $exception): bool { return (string) $exception->getCode() === '23000' || str_contains(strtolower($exception->getMessage()), 'duplicate'); } }