5 use Drupal\Core\Queue\QueueWorkerBase;
6 use Drupal\Core\Logger\RfcLogLevel;
7 use Procrastinator\Result;
8 use Drupal\Core\Plugin\ContainerFactoryPluginInterface;
9 use Symfony\Component\DependencyInjection\ContainerInterface;
20 class Import extends QueueWorkerBase implements ContainerFactoryPluginInterface {
22 use \Drupal\Core\Logger\LoggerChannelTrait;
32 public static function create(ContainerInterface $container, array $configuration, $plugin_id, $plugin_definition) {
33 return new Import($configuration, $plugin_id, $plugin_definition, $container);
48 public function __construct(array $configuration, $plugin_id, $plugin_definition, ContainerInterface $container) {
49 parent::__construct($configuration, $plugin_id, $plugin_definition);
50 $this->container = $container;
56 public function processItem($data) {
60 $datastore = $this->container->get(
'dkan_datastore.service');
62 $results = $datastore->import($data[
'uuid']);
64 foreach ($results as $result) {
65 $this->processResult($result, $data);
68 catch (\Exception $e) {
69 $this->
log(RfcLogLevel::ERROR,
70 "Import for {$data['uuid']} returned an error: {$e->getMessage()}");
77 private function processResult(Result $result, $data) {
78 $level = RfcLogLevel::INFO;
80 $status = $result->getStatus();
83 $newQueueItemId = $this->
requeue($data);
84 $message =
"Import for {$data['uuid']} is requeueing for iteration No. {$data['queue_iteration']}. (ID:{$newQueueItemId}).";
87 case Result::IN_PROGRESS:
89 $level = RfcLogLevel::ERROR;
90 $message =
"Import for {$data['uuid']} returned an error: {$result->getError()}";
94 $message =
"Import for {$data['uuid']} completed.";
97 $this->
log($level, $message);
103 protected function log($level, $message, array $context = []) {
104 $this->getLogger($this->getPluginId())
105 ->log($level, $message, $context);
120 return $this->container->get(
'queue')
121 ->get($this->getPluginId())