Repository navigation
Expand file tree
/
Copy pathMessengerMessageRepository.php
More file actions
125 lines (97 loc) · 2.95 KB
/
Copy pathMessengerMessageRepository.php
File metadata and controls
125 lines (97 loc) · 2.95 KB
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
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
<?php
namespace SyncEngine\Repository;
use Doctrine\DBAL\ArrayParameterType;
use Doctrine\ORM\EntityManagerInterface;
use Symfony\Component\DependencyInjection\Attribute\Autowire;
/**
* Custom Messenger repo to fetch messages directly from the DB.
*/
class MessengerMessageRepository
{
private EntityManagerInterface $em;
private string $table;
public function __construct(
#[Autowire( '%syncengine.db.table_prefix%' )] string $prefix,
EntityManagerInterface $em,
) {
$this->em = $em;
$this->table = $prefix . 'messenger_messages';
}
public function find( $id, $lockMode = null, $lockVersion = null ): ?array
{
$conn = $this->em->getConnection();
$sql = "SELECT * FROM {$this->table} WHERE id = :id LIMIT 1";
$stmt = $conn->prepare( $sql );
$stmt->bindValue( 'id', $id );
$row = $stmt->executeQuery()->fetchAssociative();
return $row ?: null;
}
public function findOneBy( array $criteria, array $orderBy = null ): ?array
{
$rows = $this->findBy( $criteria, $orderBy, 1, 0 );
return $rows[0] ?? null;
}
public function findAll(): array
{
$conn = $this->em->getConnection();
$sql = "SELECT * FROM {$this->table}";
$stmt = $conn->executeQuery( $sql );
return $stmt->fetchAllAssociative();
}
public function findBy( array $criteria, array $orderBy = null, $limit = null, $offset = null ): array
{
$conn = $this->em->getConnection();
$where = [];
$params = [];
$types = [];
foreach ( $criteria as $field => $value ) {
$field = $this->mapFieldName( $field );
if ( is_array( $value ) ) {
$where[] = "$field IN (:{$field})";
$params[ $field ] = $value;
$types[ $field ] = ArrayParameterType::STRING;
} else {
$where[] = "$field = :{$field}";
$params[ $field ] = $value;
}
}
$sql = "SELECT * FROM {$this->table}";
if ( $where ) {
$sql .= " WHERE " . implode( ' AND ', $where );
}
if ( $orderBy ) {
$orderParts = [];
foreach ( $orderBy as $field => $direction ) {
$field = $this->mapFieldName( $field );
$direction = strtoupper( $direction ) === 'DESC' ? 'DESC' : 'ASC';
$orderParts[] = "{$field} {$direction}";
}
$sql .= " ORDER BY " . implode( ', ', $orderParts );
}
if ( $limit !== null ) {
$sql .= " LIMIT " . (int) $limit;
}
if ( $offset !== null ) {
$sql .= " OFFSET " . (int) $offset;
}
$stmt = $conn->prepare( $sql );
foreach ( $params as $key => $value ) {
if ( isset( $types[ $key ] ) && $types[ $key ] === ArrayParameterType::STRING ) {
$stmt->bindValue( $key, $value, ArrayParameterType::STRING ); // @phpstan-ignore argument.type
} else {
$stmt->bindValue( $key, $value );
}
}
return $stmt->executeQuery()->fetchAllAssociative();
}
private function mapFieldName( $field ): string
{
return match( $field ) {
'createdAt' => 'created_at',
'availableAt' => 'available_at',
'deliveredAt' => 'delivered_at',
'redeliveredAt' => 'redelivered_at',
default => $field,
};
}
}