Skip to content
Open
Show file tree
Hide file tree
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
37 changes: 37 additions & 0 deletions lib/Db/ImipData.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
<?php

declare(strict_types=1);

/**
* SPDX-FileCopyrightText: 2026 Nextcloud GmbH and Nextcloud contributors
* SPDX-License-Identifier: AGPL-3.0-or-later
*/

namespace OCA\Mail\Db;

use OCP\AppFramework\Db\Entity;

/**
* @method void setImipMessageId(int $messageId)
* @method int getImipMessageId()
* @method void setError(bool $error)
* @method bool getError()
* @method void setProcessedAt(?int $processedAt)
* @method int|null getProcessedAt()
*/
class ImipData extends Entity {
/** @var int */
protected $imipMessageId;

/** @var bool */
protected $error;

/** @var int|null */
protected $processedAt;

public function __construct() {
$this->addType('imipMessageId', 'integer');
$this->addType('error', 'boolean');
$this->addType('processedAt', 'integer');
}
}
107 changes: 107 additions & 0 deletions lib/Db/ImipDataMapper.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,107 @@
<?php

declare(strict_types=1);

/**
* SPDX-FileCopyrightText: 2026 Nextcloud GmbH and Nextcloud contributors
* SPDX-License-Identifier: AGPL-3.0-or-later
*/

namespace OCA\Mail\Db;

use OCP\AppFramework\Db\DoesNotExistException;
use OCP\AppFramework\Db\QBMapper;
use OCP\AppFramework\Utility\ITimeFactory;
use OCP\DB\Exception;
use OCP\DB\QueryBuilder\IQueryBuilder;
use OCP\IDBConnection;
use Throwable;

/**
* @template-extends QBMapper<ImipData>
*/
class ImipDataMapper extends QBMapper {

/**
* @param IDBConnection $db
*/
public function __construct(
IDBConnection $db,
private ITimeFactory $timeFactory,
) {
parent::__construct($db, 'mail_messages_imip');
}

public function findByMessageId(int $messageId): ?ImipData {
$qb = $this->db->getQueryBuilder();
$qb->select('*')
->from($this->getTableName())
->where(
$qb->expr()->eq('imip_message_id', $qb->createNamedParameter($messageId, IQueryBuilder::PARAM_INT))
);

try {
return $this->findEntity($qb);
} catch (DoesNotExistException) {
return null;
}
}

/**
* @throws Exception
*/
public function markAsImipMessage(int $messageId): void {
if ($this->findByMessageId($messageId) !== null) {
return;
}

$imipData = new ImipData();
$imipData->setImipMessageId($messageId);
$imipData->setError(false);
$imipData->setProcessedAt(null);
$this->insert($imipData);
}

/**
* @throws Exception
*/
public function markProcessed(int $messageId, bool $error): void {
$qb = $this->db->getQueryBuilder();
$update = $qb->update($this->getTableName())
->set('error', $qb->createNamedParameter($error, IQueryBuilder::PARAM_BOOL))
->set('processed_at', $qb->createNamedParameter($this->timeFactory->getTime(), IQueryBuilder::PARAM_INT))
->where(
$qb->expr()->eq('imip_message_id', $qb->createNamedParameter($messageId, IQueryBuilder::PARAM_INT))
);

$update->executeStatement();

}

/**
* @throws Exception
* @throws Throwable
*/
public function markProcessedBulk(Message ...$messages): array {
$this->db->beginTransaction();

try {
foreach ($messages as $message) {
if (empty($message->getUpdatedFields())) {
continue;
}

$this->markProcessed($message->getId(), $message->isImipError());
}

$this->db->commit();
} catch (Throwable $e) {
$this->db->rollBack();

throw $e;
}

return $messages;
}

}
18 changes: 8 additions & 10 deletions lib/Db/MessageMapper.php
Original file line number Diff line number Diff line change
Expand Up @@ -585,7 +585,6 @@ public function updatePreviewDataBulk(Message ...$messages): array {
->set('preview_text', $query->createParameter('preview_text'))
->set('structure_analyzed', $query->createNamedParameter(true, IQueryBuilder::PARAM_BOOL))
->set('updated_at', $query->createNamedParameter($this->timeFactory->getTime(), IQueryBuilder::PARAM_INT))
->set('imip_message', $query->createParameter('imip_message'))
->set('encrypted', $query->createParameter('encrypted'))
->set('mentions_me', $query->createParameter('mentions_me'))
->where($query->expr()->andX(
Expand Down Expand Up @@ -616,7 +615,6 @@ public function updatePreviewDataBulk(Message ...$messages): array {
$previewText,
$previewText === null ? IQueryBuilder::PARAM_NULL : IQueryBuilder::PARAM_STR
);
$query->setParameter('imip_message', $message->isImipMessage(), IQueryBuilder::PARAM_BOOL);
$query->setParameter('encrypted', $message->isEncrypted(), IQueryBuilder::PARAM_BOOL);
$query->setParameter('mentions_me', $message->getMentionsMe(), IQueryBuilder::PARAM_BOOL);

Expand Down Expand Up @@ -1564,15 +1562,15 @@ public function findIMipMessagesAscending(): array {
$time = $this->timeFactory->getTime() - 60 * 60 * 24 * 14;
$qb = $this->db->getQueryBuilder();

$select = $qb->select('*')
->from($this->getTableName())
$select = $qb->select('m.*')
->from($this->getTableName(), 'm')
->join('m', 'mail_messages_imip', 'i', $qb->expr()->eq('i.imip_message_id', 'm.id'))
->where(
$qb->expr()->eq('imip_message', $qb->createNamedParameter(true, IQueryBuilder::PARAM_BOOL), IQueryBuilder::PARAM_BOOL),
$qb->expr()->eq('imip_processed', $qb->createNamedParameter(false, IQueryBuilder::PARAM_BOOL), IQueryBuilder::PARAM_BOOL),
$qb->expr()->eq('imip_error', $qb->createNamedParameter(false, IQueryBuilder::PARAM_BOOL), IQueryBuilder::PARAM_BOOL),
$qb->expr()->eq('flag_junk', $qb->createNamedParameter(false, IQueryBuilder::PARAM_BOOL), IQueryBuilder::PARAM_BOOL),
$qb->expr()->gt('sent_at', $qb->createNamedParameter($time, IQueryBuilder::PARAM_INT)),
)->orderBy('sent_at', 'ASC'); // make sure we don't process newer messages first
$qb->expr()->eq('i.error', $qb->createNamedParameter(false, IQueryBuilder::PARAM_BOOL), IQueryBuilder::PARAM_BOOL),
$qb->expr()->isNull('i.processed_at'),
$qb->expr()->eq('m.flag_junk', $qb->createNamedParameter(false, IQueryBuilder::PARAM_BOOL), IQueryBuilder::PARAM_BOOL),
$qb->expr()->gt('m.sent_at', $qb->createNamedParameter($time, IQueryBuilder::PARAM_INT)),
)->orderBy('m.sent_at', 'ASC'); // make sure we don't process newer messages first

return $this->findEntities($select);
}
Expand Down
24 changes: 23 additions & 1 deletion lib/IMAP/PreviewEnhancer.php
Original file line number Diff line number Diff line change
Expand Up @@ -11,13 +11,15 @@

use Horde_Imap_Client_Exception;
use OCA\Mail\Account;
use OCA\Mail\Db\ImipDataMapper;
use OCA\Mail\Db\Mailbox;
use OCA\Mail\Db\Message;
use OCA\Mail\Db\MessageMapper as DbMapper;
use OCA\Mail\IMAP\MessageMapper as ImapMapper;
use OCA\Mail\Service\Attachment\AttachmentService;
use OCA\Mail\Service\Avatar\Avatar;
use OCA\Mail\Service\AvatarService;
use OCP\DB\Exception;
use Psr\Log\LoggerInterface;
use function array_key_exists;
use function array_map;
Expand All @@ -40,19 +42,24 @@ class PreviewEnhancer {
/** @var AvatarService */
private $avatarService;

/** @var ImipDataMapper */
private $imipDataMapper;

public function __construct(
IMAPClientFactory $clientFactory,
ImapMapper $imapMapper,
DbMapper $dbMapper,
LoggerInterface $logger,
AvatarService $avatarService,
private AttachmentService $attachmentService,
ImipDataMapper $imipDataMapper,
) {
$this->clientFactory = $clientFactory;
$this->imapMapper = $imapMapper;
$this->mapper = $dbMapper;
$this->logger = $logger;
$this->avatarService = $avatarService;
$this->imipDataMapper = $imipDataMapper;
}

/**
Expand Down Expand Up @@ -116,6 +123,22 @@ public function process(Account $account, Mailbox $mailbox, array $messages, boo
$client->logout();
}

foreach ($messages as $message) {
if (!array_key_exists($message->getUid(), $data)) {
continue;
}

if ($data[$message->getUid()]->isImipMessage()) {
try {
$this->imipDataMapper->markAsImipMessage($message->getId());
} catch (Exception $e) {
$this->logger->warning('Could not mark message as imip: ' . $e->getMessage(), [
'exception' => $e,
]);
}
}
}

return $this->mapper->updatePreviewDataBulk(...array_map(static function (Message $message) use ($data) {
if (!array_key_exists($message->getUid(), $data)) {
// Nothing to do
Expand All @@ -126,7 +149,6 @@ public function process(Account $account, Mailbox $mailbox, array $messages, boo
$message->setFlagAttachments($structureData->hasAttachments());
$message->setPreviewText($structureData->getPreviewText());
$message->setStructureAnalyzed(true);
$message->setImipMessage($structureData->isImipMessage());
$message->setEncrypted($structureData->isEncrypted());
$message->setMentionsMe($structureData->getMentionsMe());

Expand Down
66 changes: 66 additions & 0 deletions lib/Migration/Version5201Date20260720120000.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
<?php

declare(strict_types=1);

/**
* SPDX-FileCopyrightText: 2026 Nextcloud GmbH and Nextcloud contributors
* SPDX-License-Identifier: AGPL-3.0-or-later
*/

namespace OCA\Mail\Migration;

use Closure;
use OCP\DB\ISchemaWrapper;
use OCP\DB\Types;
use OCP\Migration\IOutput;
use OCP\Migration\SimpleMigrationStep;
use Override;

class Version5201Date20260720120000 extends SimpleMigrationStep {

/**
* @param IOutput $output
* @param Closure(): ISchemaWrapper $schemaClosure
* @param array $options
* @return null|ISchemaWrapper
*/
#[Override]
public function changeSchema(IOutput $output, Closure $schemaClosure, array $options): ?ISchemaWrapper {
$schema = $schemaClosure();

if (!$schema->hasTable('mail_messages_imip')) {
$table = $schema->createTable('mail_messages_imip');
$table->addColumn('id', Types::INTEGER, [
'autoincrement' => true,
'notnull' => true,
]);
$table->addColumn('imip_message_id', Types::BIGINT, [
'notnull' => true,
]);
$table->addColumn('error', Types::BOOLEAN, [
'notnull' => true,
'default' => false,
]);
$table->addColumn('processed_at', Types::INTEGER, [
'notnull' => false,
'default' => null,
]);
$table->setPrimaryKey(['id']);
$table->addUniqueIndex(['imip_message_id'], 'mail_msg_imip_msg_uniq');
$table->addIndex(['error', 'processed_at'], 'mail_msg_imip_unproc_idx');

if ($schema->hasTable('mail_messages')) {
$table->addForeignKeyConstraint(
$schema->getTable('mail_messages'),
['imip_message_id'],
['id'],
[
'onDelete' => 'CASCADE',
]
);
}
}

return $schema;
}
}
9 changes: 7 additions & 2 deletions lib/Service/IMipService.php
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
namespace OCA\Mail\Service;

use OCA\Mail\Account;
use OCA\Mail\Db\ImipDataMapper;
use OCA\Mail\Db\Mailbox;
use OCA\Mail\Db\MailboxMapper;
use OCA\Mail\Db\Message;
Expand All @@ -32,6 +33,7 @@ class IMipService {
private MailManager $mailManager;
private MessageMapper $messageMapper;
private ServerVersion $serverVersion;
private ImipDataMapper $imipDataMapper;

public function __construct(
AccountService $accountService,
Expand All @@ -41,6 +43,7 @@ public function __construct(
MailManager $mailManager,
MessageMapper $messageMapper,
ServerVersion $serverVersion,
ImipDataMapper $imipDataMapper,
) {
$this->accountService = $accountService;
$this->calendarManager = $manager;
Expand All @@ -49,6 +52,7 @@ public function __construct(
$this->mailManager = $mailManager;
$this->messageMapper = $messageMapper;
$this->serverVersion = $serverVersion;
$this->imipDataMapper = $imipDataMapper;
}

public function process(): void {
Expand Down Expand Up @@ -108,7 +112,8 @@ public function process(): void {
$message->setImipProcessed(true);
return $message;
}, $filteredMessages); // Silently drop from passing to DAV and mark as processed, so we won't run into these messages again.
$this->messageMapper->updateImipData(...$processedMessages);
$this->imipDataMapper->markProcessedBulk(...$processedMessages);

continue;
}

Expand Down Expand Up @@ -184,7 +189,7 @@ public function process(): void {
$message->setImipError(true);
}
}
$this->messageMapper->updateImipData(...$filteredMessages);
$this->imipDataMapper->markProcessedBulk(...$filteredMessages);
}
}
}
Loading