Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
168 changes: 69 additions & 99 deletions lib/Data.php
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,9 @@ class Data {
*/
public const MAX_HISTOGRAM_ROWS = 100000;

/** @var */
private const ACTIVITY_COLUMNS = ['app', 'subject', 'subjectparams', 'message', 'messageparams', 'file', 'link', 'user', 'affecteduser', 'timestamp', 'priority', 'type', 'object_type', 'object_id'];
private const MAIL_COLUMNS = ['amq_appid', 'amq_subject', 'amq_subjectparams', 'amq_affecteduser', 'amq_timestamp', 'amq_type', 'amq_latest_send', 'object_type', 'object_id'];

protected ?IQueryBuilder $insertActivity = null;
protected ?IQueryBuilder $insertMail = null;

Expand Down Expand Up @@ -87,47 +89,10 @@ public function send(IEvent $event): int {
return 0;
}

if ($this->insertActivity === null) {
$this->insertActivity = $this->connection->getQueryBuilder();
$this->insertActivity->insert('activity')
->values([
'app' => $this->insertActivity->createParameter('app'),
'subject' => $this->insertActivity->createParameter('subject'),
'subjectparams' => $this->insertActivity->createParameter('subjectparams'),
'message' => $this->insertActivity->createParameter('message'),
'messageparams' => $this->insertActivity->createParameter('messageparams'),
'file' => $this->insertActivity->createParameter('object_name'),
'link' => $this->insertActivity->createParameter('link'),
'user' => $this->insertActivity->createParameter('user'),
'affecteduser' => $this->insertActivity->createParameter('affecteduser'),
'timestamp' => $this->insertActivity->createParameter('timestamp'),
'priority' => $this->insertActivity->createParameter('priority'),
'type' => $this->insertActivity->createParameter('type'),
'object_type' => $this->insertActivity->createParameter('object_type'),
'object_id' => $this->insertActivity->createParameter('object_id'),
]);
}

// store in DB
$this->insertActivity->setParameters([
'app' => $event->getApp(),
'type' => $event->getType(),
'affecteduser' => $event->getAffectedUser(),
'user' => $event->getAuthor(),
'timestamp' => $event->getTimestamp(),
'subject' => $event->getSubject(),
'subjectparams' => json_encode($event->getSubjectParameters()),
'message' => $event->getMessage(),
'messageparams' => json_encode($event->getMessageParameters()),
'priority' => IExtension::PRIORITY_MEDIUM,
'object_type' => $event->getObjectType(),
'object_id' => $event->getObjectId(),
'object_name' => $event->getObjectName(),
'link' => $event->getLink(),
]);
$this->insertActivity->executeStatement();
$insert = $this->getActivityInsert();
$this->executeInsert($insert, self::ACTIVITY_COLUMNS, $this->getActivityRow($event, $event->getAffectedUser()));

return $this->insertActivity->getLastInsertId();
return $insert->getLastInsertId();
}

/**
Expand All @@ -147,47 +112,10 @@ public function bulkSend(IEvent $event, array $affectedUsers): array {

$activityIds = [];
try {
if ($this->insertActivity === null) {
$this->insertActivity = $this->connection->getQueryBuilder();
}
$this->insertActivity->insert('activity')
->values([
'app' => $this->insertActivity->createParameter('app'),
'subject' => $this->insertActivity->createParameter('subject'),
'subjectparams' => $this->insertActivity->createParameter('subjectparams'),
'message' => $this->insertActivity->createParameter('message'),
'messageparams' => $this->insertActivity->createParameter('messageparams'),
'file' => $this->insertActivity->createParameter('object_name'),
'link' => $this->insertActivity->createParameter('link'),
'user' => $this->insertActivity->createParameter('user'),
'affecteduser' => $this->insertActivity->createParameter('affecteduser'),
'timestamp' => $this->insertActivity->createParameter('timestamp'),
'priority' => $this->insertActivity->createParameter('priority'),
'type' => $this->insertActivity->createParameter('type'),
'object_type' => $this->insertActivity->createParameter('object_type'),
'object_id' => $this->insertActivity->createParameter('object_id'),
]);

$this->insertActivity->setParameters([
'app' => $event->getApp(),
'type' => $event->getType(),
'user' => $event->getAuthor(),
'timestamp' => $event->getTimestamp(),
'subject' => $event->getSubject(),
'subjectparams' => json_encode($event->getSubjectParameters()),
'message' => $event->getMessage(),
'messageparams' => json_encode($event->getMessageParameters()),
'priority' => IExtension::PRIORITY_MEDIUM,
'object_type' => $event->getObjectType(),
'object_id' => $event->getObjectId(),
'object_name' => $event->getObjectName(),
'link' => $event->getLink(),
]);

$insert = $this->getActivityInsert();
foreach ($affectedUsers as $affectedUser) {
$this->insertActivity->setParameter('affecteduser', $affectedUser);
$this->insertActivity->executeStatement();
$activityIds[$this->insertActivity->getLastInsertId()] = (string)$affectedUser;
$this->executeInsert($insert, self::ACTIVITY_COLUMNS, $this->getActivityRow($event, (string)$affectedUser));
$activityIds[$insert->getLastInsertId()] = (string)$affectedUser;
}

$this->connection->commit();
Expand All @@ -212,22 +140,10 @@ public function storeMail(IEvent $event, int $latestSendTime): bool {
}

if ($this->insertMail === null) {
$this->insertMail = $this->connection->getQueryBuilder();
$this->insertMail->insert('activity_mq')
->values([
'amq_appid' => $this->insertMail->createParameter('amq_appid'),
'amq_subject' => $this->insertMail->createParameter('amq_subject'),
'amq_subjectparams' => $this->insertMail->createParameter('amq_subjectparams'),
'amq_affecteduser' => $this->insertMail->createParameter('amq_affecteduser'),
'amq_timestamp' => $this->insertMail->createParameter('amq_timestamp'),
'amq_type' => $this->insertMail->createParameter('amq_type'),
'amq_latest_send' => $this->insertMail->createParameter('amq_latest_send'),
'object_type' => $this->insertMail->createParameter('object_type'),
'object_id' => $this->insertMail->createParameter('object_id'),
]);
}

$this->insertMail->setParameters([
$this->insertMail = $this->createInsert('activity_mq', self::MAIL_COLUMNS);
}

$this->executeInsert($this->insertMail, self::MAIL_COLUMNS, [
'amq_appid' => $event->getApp(),
'amq_subject' => $event->getSubject(),
'amq_subjectparams' => json_encode($event->getSubjectParameters()),
Expand All @@ -239,11 +155,65 @@ public function storeMail(IEvent $event, int $latestSendTime): bool {
'object_id' => $event->getObjectId(),
]);

$this->insertMail->executeStatement();

return true;
}

private function getActivityInsert(): IQueryBuilder {
if ($this->insertActivity === null) {
$this->insertActivity = $this->createInsert('activity', self::ACTIVITY_COLUMNS);
}
return $this->insertActivity;
}

/**
* @return array<string, mixed>
*/
private function getActivityRow(IEvent $event, string $affectedUser): array {
return [
'app' => $event->getApp(),
'subject' => $event->getSubject(),
'subjectparams' => json_encode($event->getSubjectParameters()),
'message' => $event->getMessage(),
'messageparams' => json_encode($event->getMessageParameters()),
'file' => $event->getObjectName(),
'link' => $event->getLink(),
'user' => $event->getAuthor(),
'affecteduser' => $affectedUser,
'timestamp' => $event->getTimestamp(),
'priority' => IExtension::PRIORITY_MEDIUM,
'type' => $event->getType(),
'object_type' => $event->getObjectType(),
'object_id' => $event->getObjectId(),
];
}

/**
* Build an insert that is executed many times per request
*
* Positional parameters are used, as named ones make the database layer
* parse the statement again on every execution.
*
* @param list<string> $columns
*/
private function createInsert(string $table, array $columns): IQueryBuilder {
$query = $this->connection->getQueryBuilder();
$values = [];
foreach ($columns as $column) {
$values[$column] = $query->createPositionalParameter(null);
}
$query->insert($table)->values($values);
return $query;
}

/**
* @param list<string> $columns
* @param array<string, mixed> $row
*/
private function executeInsert(IQueryBuilder $query, array $columns, array $row): void {
$query->setParameters(array_map(static fn (string $column): mixed => $row[$column], $columns));
$query->executeStatement();
}

/**
* Read a list of events from the activity stream
*
Expand Down
Loading