-
Notifications
You must be signed in to change notification settings - Fork 56
/
MySqlAggregateStreamStrategy.php
97 lines (82 loc) · 2.83 KB
/
MySqlAggregateStreamStrategy.php
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
<?php
/**
* This file is part of prooph/pdo-event-store.
* (c) 2016-2019 Alexander Miertsch <kontakt@codeliner.ws>
* (c) 2016-2019 Sascha-Oliver Prolic <saschaprolic@googlemail.com>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/
declare(strict_types=1);
namespace Prooph\EventStore\Pdo\PersistenceStrategy;
use Iterator;
use Prooph\Common\Messaging\MessageConverter;
use Prooph\EventStore\Pdo\DefaultMessageConverter;
use Prooph\EventStore\Pdo\Exception;
use Prooph\EventStore\Pdo\PersistenceStrategy;
use Prooph\EventStore\Pdo\Util\Json;
use Prooph\EventStore\StreamName;
final class MySqlAggregateStreamStrategy implements PersistenceStrategy
{
/**
* @var MessageConverter
*/
private $messageConverter;
public function __construct(?MessageConverter $messageConverter = null)
{
$this->messageConverter = $messageConverter ?? new DefaultMessageConverter();
}
/**
* @param string $tableName
* @return string[]
*/
public function createSchema(string $tableName): array
{
$statement = <<<EOT
CREATE TABLE `$tableName` (
`no` BIGINT(20) NOT NULL AUTO_INCREMENT,
`event_id` CHAR(36) COLLATE utf8mb4_bin NOT NULL,
`event_name` VARCHAR(100) COLLATE utf8mb4_bin NOT NULL,
`payload` JSON NOT NULL,
`metadata` JSON NOT NULL,
`created_at` DATETIME(6) NOT NULL,
`aggregate_version` INT(11) UNSIGNED GENERATED ALWAYS AS (JSON_EXTRACT(metadata, '$._aggregate_version')) STORED NOT NULL UNIQUE KEY,
PRIMARY KEY (`no`),
UNIQUE KEY `ix_event_id` (`event_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_bin;
EOT;
return [$statement];
}
public function columnNames(): array
{
return [
'no',
'event_id',
'event_name',
'payload',
'metadata',
'created_at',
];
}
public function prepareData(Iterator $streamEvents): array
{
$data = [];
foreach ($streamEvents as $event) {
$eventData = $this->messageConverter->convertToArray($event);
if (! isset($eventData['metadata']['_aggregate_version'])) {
throw new Exception\RuntimeException('_aggregate_version is missing in metadata');
}
$data[] = $eventData['metadata']['_aggregate_version'];
$data[] = $eventData['uuid'];
$data[] = $eventData['message_name'];
$data[] = Json::encode($eventData['payload']);
$data[] = Json::encode($eventData['metadata']);
$data[] = $eventData['created_at']->format('Y-m-d\TH:i:s.u');
}
return $data;
}
public function generateTableName(StreamName $streamName): string
{
return '_' . \sha1($streamName->toString());
}
}