feat(notification): add persistent notification
This commit is contained in:
@@ -14,6 +14,7 @@ mod favorites_pg_repository;
|
||||
pub mod file_metadata_repository;
|
||||
mod magic_link_token_pg_repository;
|
||||
mod nextcloud_object_id_repository;
|
||||
mod notification_pg_repository;
|
||||
mod opaque_pg_repository;
|
||||
pub mod playlist_pg_repository;
|
||||
mod recent_items_pg_repository;
|
||||
@@ -48,6 +49,7 @@ pub use file_metadata_repository::FileMetadataRepository;
|
||||
pub use folder_db_repository::FolderDbRepository;
|
||||
pub use magic_link_token_pg_repository::MagicLinkTokenPgRepository;
|
||||
pub use nextcloud_object_id_repository::NextcloudObjectIdRepository;
|
||||
pub use notification_pg_repository::NotificationPgRepository;
|
||||
pub use opaque_pg_repository::OpaquePgRepository;
|
||||
pub use playlist_pg_repository::{
|
||||
AudioMetadataPgRepository, PlaylistItemPgRepository, PlaylistPgRepository,
|
||||
|
||||
@@ -0,0 +1,281 @@
|
||||
//! PostgreSQL implementation of [`NotificationRepository`].
|
||||
//!
|
||||
//! Backs the bell UI plus the daily retention job. All queries scope on
|
||||
//! `user_id` at the SQL layer so a row misroute in the caller can't
|
||||
//! leak another user's data through mark_read / delete. Schema lives
|
||||
//! in `migrations/20261026000000_notifications.sql`.
|
||||
|
||||
use async_trait::async_trait;
|
||||
use chrono::{DateTime, Utc};
|
||||
use sqlx::{PgPool, Row};
|
||||
use std::sync::Arc;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::common::errors::{DomainError, ErrorKind};
|
||||
use crate::domain::entities::notification::{NewNotification, Notification};
|
||||
use crate::domain::repositories::notification_repository::{
|
||||
NotificationListFilter, NotificationRepository,
|
||||
};
|
||||
|
||||
pub struct NotificationPgRepository {
|
||||
pool: Arc<PgPool>,
|
||||
}
|
||||
|
||||
impl NotificationPgRepository {
|
||||
pub fn new(pool: Arc<PgPool>) -> Self {
|
||||
Self { pool }
|
||||
}
|
||||
|
||||
fn map_row(row: &sqlx::postgres::PgRow) -> Result<Notification, DomainError> {
|
||||
let map_err = |field: &str, e: sqlx::Error| {
|
||||
DomainError::new(
|
||||
ErrorKind::DatabaseError,
|
||||
"Notification",
|
||||
format!("read {field}: {e}"),
|
||||
)
|
||||
};
|
||||
Ok(Notification {
|
||||
id: row.try_get("id").map_err(|e| map_err("id", e))?,
|
||||
user_id: row.try_get("user_id").map_err(|e| map_err("user_id", e))?,
|
||||
kind: row.try_get("kind").map_err(|e| map_err("kind", e))?,
|
||||
payload: row.try_get("payload").map_err(|e| map_err("payload", e))?,
|
||||
created_at: row
|
||||
.try_get("created_at")
|
||||
.map_err(|e| map_err("created_at", e))?,
|
||||
read_at: row.try_get("read_at").ok(),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
fn db_err(op: &'static str, e: sqlx::Error) -> DomainError {
|
||||
DomainError::new(
|
||||
ErrorKind::DatabaseError,
|
||||
"Notification",
|
||||
format!("{op}: {e}"),
|
||||
)
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl NotificationRepository for NotificationPgRepository {
|
||||
async fn create(&self, new_notif: &NewNotification) -> Result<Notification, DomainError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
INSERT INTO notif.notifications (user_id, kind, payload)
|
||||
VALUES ($1::uuid, $2, $3)
|
||||
RETURNING id, user_id, kind, payload, created_at, read_at
|
||||
"#,
|
||||
)
|
||||
.bind(new_notif.user_id)
|
||||
.bind(&new_notif.kind)
|
||||
.bind(&new_notif.payload)
|
||||
.fetch_one(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("create", e))?;
|
||||
Self::map_row(&row)
|
||||
}
|
||||
|
||||
async fn list_for_user(
|
||||
&self,
|
||||
user_id: Uuid,
|
||||
filter: &NotificationListFilter,
|
||||
) -> Result<Vec<Notification>, DomainError> {
|
||||
// Dynamic-shape query built to still hit the
|
||||
// notifications_user_created_read index — every branch keys
|
||||
// on (user_id, created_at DESC).
|
||||
let limit: i64 = filter.limit.unwrap_or(50).min(500) as i64;
|
||||
let rows = match (filter.unread_only, filter.before) {
|
||||
(None, None) => {
|
||||
sqlx::query(
|
||||
r#"
|
||||
SELECT id, user_id, kind, payload, created_at, read_at
|
||||
FROM notif.notifications
|
||||
WHERE user_id = $1::uuid
|
||||
ORDER BY created_at DESC
|
||||
LIMIT $2
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(limit)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
}
|
||||
(Some(true), None) => {
|
||||
sqlx::query(
|
||||
r#"
|
||||
SELECT id, user_id, kind, payload, created_at, read_at
|
||||
FROM notif.notifications
|
||||
WHERE user_id = $1::uuid AND read_at IS NULL
|
||||
ORDER BY created_at DESC
|
||||
LIMIT $2
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(limit)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
}
|
||||
(Some(false), None) => {
|
||||
sqlx::query(
|
||||
r#"
|
||||
SELECT id, user_id, kind, payload, created_at, read_at
|
||||
FROM notif.notifications
|
||||
WHERE user_id = $1::uuid AND read_at IS NOT NULL
|
||||
ORDER BY created_at DESC
|
||||
LIMIT $2
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(limit)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
}
|
||||
(None, Some(before)) => {
|
||||
sqlx::query(
|
||||
r#"
|
||||
SELECT id, user_id, kind, payload, created_at, read_at
|
||||
FROM notif.notifications
|
||||
WHERE user_id = $1::uuid AND created_at < $2
|
||||
ORDER BY created_at DESC
|
||||
LIMIT $3
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(before)
|
||||
.bind(limit)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
}
|
||||
(Some(true), Some(before)) => {
|
||||
sqlx::query(
|
||||
r#"
|
||||
SELECT id, user_id, kind, payload, created_at, read_at
|
||||
FROM notif.notifications
|
||||
WHERE user_id = $1::uuid AND read_at IS NULL AND created_at < $2
|
||||
ORDER BY created_at DESC
|
||||
LIMIT $3
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(before)
|
||||
.bind(limit)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
}
|
||||
(Some(false), Some(before)) => {
|
||||
sqlx::query(
|
||||
r#"
|
||||
SELECT id, user_id, kind, payload, created_at, read_at
|
||||
FROM notif.notifications
|
||||
WHERE user_id = $1::uuid AND read_at IS NOT NULL AND created_at < $2
|
||||
ORDER BY created_at DESC
|
||||
LIMIT $3
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(before)
|
||||
.bind(limit)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
}
|
||||
}
|
||||
.map_err(|e| db_err("list_for_user", e))?;
|
||||
|
||||
rows.iter().map(Self::map_row).collect()
|
||||
}
|
||||
|
||||
async fn count_unread_for_user(&self, user_id: Uuid) -> Result<i64, DomainError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
SELECT COUNT(*)::bigint AS c
|
||||
FROM notif.notifications
|
||||
WHERE user_id = $1::uuid AND read_at IS NULL
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.fetch_one(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("count_unread_for_user", e))?;
|
||||
row.try_get::<i64, _>("c")
|
||||
.map_err(|e| db_err("count_unread_for_user.map", e))
|
||||
}
|
||||
|
||||
async fn mark_read(
|
||||
&self,
|
||||
notification_id: Uuid,
|
||||
user_id: Uuid,
|
||||
at: DateTime<Utc>,
|
||||
) -> Result<bool, DomainError> {
|
||||
// Guard on read_at IS NULL so a re-issued call from a client
|
||||
// that's already ack'd the row is a no-op instead of stamping
|
||||
// a later timestamp over the earlier one.
|
||||
let res = sqlx::query(
|
||||
r#"
|
||||
UPDATE notif.notifications
|
||||
SET read_at = $3
|
||||
WHERE id = $1::uuid
|
||||
AND user_id = $2::uuid
|
||||
AND read_at IS NULL
|
||||
"#,
|
||||
)
|
||||
.bind(notification_id)
|
||||
.bind(user_id)
|
||||
.bind(at)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("mark_read", e))?;
|
||||
Ok(res.rows_affected() == 1)
|
||||
}
|
||||
|
||||
async fn mark_all_read_for_user(
|
||||
&self,
|
||||
user_id: Uuid,
|
||||
at: DateTime<Utc>,
|
||||
) -> Result<u64, DomainError> {
|
||||
let res = sqlx::query(
|
||||
r#"
|
||||
UPDATE notif.notifications
|
||||
SET read_at = $2
|
||||
WHERE user_id = $1::uuid AND read_at IS NULL
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(at)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("mark_all_read_for_user", e))?;
|
||||
Ok(res.rows_affected())
|
||||
}
|
||||
|
||||
async fn delete_by_id(
|
||||
&self,
|
||||
notification_id: Uuid,
|
||||
user_id: Uuid,
|
||||
) -> Result<bool, DomainError> {
|
||||
let res = sqlx::query(
|
||||
r#"
|
||||
DELETE FROM notif.notifications
|
||||
WHERE id = $1::uuid AND user_id = $2::uuid
|
||||
"#,
|
||||
)
|
||||
.bind(notification_id)
|
||||
.bind(user_id)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("delete_by_id", e))?;
|
||||
Ok(res.rows_affected() == 1)
|
||||
}
|
||||
|
||||
async fn purge_read_before(&self, cutoff: DateTime<Utc>) -> Result<u64, DomainError> {
|
||||
let res = sqlx::query(
|
||||
r#"
|
||||
DELETE FROM notif.notifications
|
||||
WHERE read_at IS NOT NULL AND read_at < $1
|
||||
"#,
|
||||
)
|
||||
.bind(cutoff)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("purge_read_before", e))?;
|
||||
Ok(res.rows_affected())
|
||||
}
|
||||
}
|
||||
@@ -199,6 +199,7 @@ fn event_kind(event: &MessageBusEvent) -> &'static str {
|
||||
MessageBusEvent::FolderMoved { .. } => "folder_moved",
|
||||
MessageBusEvent::FolderDeleted { .. } => "folder_deleted",
|
||||
MessageBusEvent::AuthzChanged { .. } => "authz_changed",
|
||||
MessageBusEvent::NotificationReceived { .. } => "notification_received",
|
||||
MessageBusEvent::JobRunStarted { .. } => "job_run_started",
|
||||
MessageBusEvent::JobRunProgress { .. } => "job_run_progress",
|
||||
MessageBusEvent::JobRunEnded { .. } => "job_run_ended",
|
||||
|
||||
@@ -39,6 +39,7 @@ pub mod mock_email_sender;
|
||||
pub mod mount_provider_factory;
|
||||
pub mod nextcloud_chunked_upload_service;
|
||||
pub mod noop_face_analyzer;
|
||||
pub mod notifications_cleanup_service;
|
||||
pub mod oidc_service;
|
||||
#[cfg(feature = "faces-onnx")]
|
||||
pub mod onnx_face_analyzer;
|
||||
|
||||
@@ -0,0 +1,136 @@
|
||||
//! `notifications_cleanup` scheduled job — daily retention sweep.
|
||||
//!
|
||||
//! Deletes rows from `notif.notifications` where `read_at IS NOT NULL`
|
||||
//! and older than the retention window. Unread rows are preserved
|
||||
//! unconditionally (the whole point of the durable table is that a
|
||||
//! user offline for a month still sees the share-granted notice on
|
||||
//! next login).
|
||||
//!
|
||||
//! Retention window comes from `OXICLOUD_NOTIFICATIONS_RETENTION_DAYS`
|
||||
//! (default 30), applied at job dispatch — one env var maps to one
|
||||
//! `retention_days` parameter so an operator can override the default
|
||||
//! at trigger time without a redeploy.
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use chrono::Utc;
|
||||
use tracing::info;
|
||||
|
||||
use crate::application::services::notification_application_service::NotificationApplicationService;
|
||||
use crate::infrastructure::scheduler::{JobHandler, JobOutcome, JobRegistry, JobRunArgs, Mutates};
|
||||
|
||||
/// Parameter declaration table. Kept at module scope so
|
||||
/// `JobHandler::parameters` can return a `'static` slice without
|
||||
/// stack-allocating each call.
|
||||
static PARAMETERS: [crate::infrastructure::scheduler::JobParam; 1] =
|
||||
[crate::infrastructure::scheduler::JobParam::number(
|
||||
"retention_days",
|
||||
30,
|
||||
"Delete read notifications older than this many days.",
|
||||
)];
|
||||
|
||||
pub struct NotificationsCleanupService {
|
||||
service: Arc<NotificationApplicationService>,
|
||||
/// Default retention window in days when the trigger call did NOT
|
||||
/// supply an explicit `retention_days` parameter. Read from
|
||||
/// `OXICLOUD_NOTIFICATIONS_RETENTION_DAYS` at boot; the constructor
|
||||
/// clamps to a minimum of 1 day (0 would purge every read row on
|
||||
/// every tick).
|
||||
default_retention_days: i64,
|
||||
}
|
||||
|
||||
impl NotificationsCleanupService {
|
||||
pub const JOB_NAME: &'static str = "notifications_cleanup";
|
||||
|
||||
pub fn new(service: Arc<NotificationApplicationService>, default_retention_days: u32) -> Self {
|
||||
Self {
|
||||
service,
|
||||
default_retention_days: default_retention_days.max(1) as i64,
|
||||
}
|
||||
}
|
||||
|
||||
/// Interval — daily. Same tier as `trash_cleanup`; retention is a
|
||||
/// "days" concept, so a finer cadence buys nothing.
|
||||
fn interval() -> Duration {
|
||||
Duration::from_secs(24 * 3600)
|
||||
}
|
||||
|
||||
/// Register self with the scheduler. Chained DI helper, same shape
|
||||
/// as [`TrashCleanupService::register`].
|
||||
pub async fn register(self: Arc<Self>, registry: &JobRegistry) -> Arc<Self> {
|
||||
registry
|
||||
.register(self.clone(), Some(Self::interval()), None)
|
||||
.await;
|
||||
self
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl JobHandler for NotificationsCleanupService {
|
||||
fn name(&self) -> &str {
|
||||
Self::JOB_NAME
|
||||
}
|
||||
|
||||
fn description(&self) -> &'static str {
|
||||
"Deletes read notifications older than the retention window \
|
||||
(default 30 days, override via `retention_days` parameter or \
|
||||
OXICLOUD_NOTIFICATIONS_RETENTION_DAYS). Unread rows are \
|
||||
preserved unconditionally."
|
||||
}
|
||||
|
||||
fn mutates(&self) -> Mutates {
|
||||
Mutates::Always
|
||||
}
|
||||
|
||||
fn parameters(&self) -> &'static [crate::infrastructure::scheduler::JobParam] {
|
||||
// Declared default of 30 days is the SAME literal the config
|
||||
// block's env fallback uses (`OXICLOUD_NOTIFICATIONS_RETENTION_DAYS`
|
||||
// default), so an operator who never sets the env sees 30
|
||||
// everywhere. The env-derived `default_retention_days` on
|
||||
// this struct only diverges from 30 when the operator DID
|
||||
// set the env — see the guard in `run()` below.
|
||||
&PARAMETERS
|
||||
}
|
||||
|
||||
async fn run(&self, args: &JobRunArgs) -> JobOutcome {
|
||||
// `get_number` returns the fallback ONLY when the arg is
|
||||
// absent — but declared defaults are seeded by the engine
|
||||
// before `run` runs (see JobRunArgs::normalized_for), so the
|
||||
// param is always present with either the caller's value or
|
||||
// the declared 30. We treat "declared default AND env
|
||||
// override differs" as "use env override" to keep the
|
||||
// OXICLOUD_NOTIFICATIONS_RETENTION_DAYS knob effective
|
||||
// without teaching the engine per-instance defaults.
|
||||
let declared_default = 30_i64;
|
||||
let raw = args.get_number("retention_days", declared_default);
|
||||
let retention_days = if raw == declared_default {
|
||||
self.default_retention_days
|
||||
} else {
|
||||
raw
|
||||
}
|
||||
.max(1);
|
||||
let cutoff = Utc::now() - chrono::Duration::days(retention_days);
|
||||
|
||||
match self.service.purge_read_before_cutoff(cutoff).await {
|
||||
Ok(removed) => {
|
||||
info!(
|
||||
target: "audit",
|
||||
event = "notifications.retention_sweep",
|
||||
retention_days,
|
||||
removed,
|
||||
"🧹 notifications retention sweep: {removed} row(s) purged (retention {retention_days} d)"
|
||||
);
|
||||
JobOutcome::ok_with(
|
||||
removed,
|
||||
serde_json::json!({
|
||||
"retention_days": retention_days,
|
||||
"removed": removed,
|
||||
}),
|
||||
)
|
||||
}
|
||||
Err(e) => JobOutcome::err(format!("notifications cleanup failed: {e}")),
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user