Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
35 commits
Select commit Hold shift + click to select a range
df71f84
refactor(storage): simplify document storage helpers
senamakel Oct 9, 2026
9da1ea1
feat(storage): expose documents module
senamakel Oct 9, 2026
c2916d1
feat(security): route device store through document backend
senamakel Oct 9, 2026
4300be3
style: apply rustfmt formatting to device and document store code
senamakel Oct 9, 2026
6980771
feat(desktop): add notification store documents
senamakel Oct 9, 2026
09cd1fb
feat(notifications): route store calls through document backend
senamakel Oct 9, 2026
e3fe2e5
chore(notifications): register store_documents module
senamakel Oct 9, 2026
877f5b7
style(notifications): reformat received_ms insert
senamakel Oct 9, 2026
e349430
fix(notifications): require sync on document edit closures
senamakel Oct 9, 2026
257739d
test(notifications): use explicit first version in store documents test
senamakel Oct 9, 2026
c818cf3
refactor(task_sources): extract patch application and route updates t…
senamakel Oct 9, 2026
46d1a85
docs(task_sources): move update_source doc comment to its function
senamakel Oct 9, 2026
bc95bd8
feat(integrations): add document store task source
senamakel Oct 9, 2026
4bfc66d
refactor(task_sources): remove unused storage error wrapper
senamakel Oct 9, 2026
10dbee4
feat(task_sources): route store operations through document backend
senamakel Oct 9, 2026
3be80fc
chore(task_sources): register store_documents module
senamakel Oct 9, 2026
ee8fefd
test(task_sources): add store documents tests
senamakel Oct 9, 2026
6343677
style(task_sources): reformat store_documents for rustfmt
senamakel Oct 9, 2026
ca03629
fix(task_sources): convert provider slug parse error to anyhow
senamakel Oct 9, 2026
fbc7cef
refactor(notifications): extract row conversion helpers into store_rows
senamakel Oct 9, 2026
8be0916
refactor(task_sources): extract row mapping into store_rows module
senamakel Oct 9, 2026
069dcc5
chore(notifications): drop unused DateTime import
senamakel Oct 9, 2026
1c8a3c5
test(cli): register storage_domains_e2e test binary
senamakel Oct 9, 2026
ab4bce6
docs: document storage backend paths for notifications, task sources,…
senamakel Oct 9, 2026
69e680f
docs(storage): document scope helpers, block_on and document port
senamakel Oct 9, 2026
628e1e0
Merge remote-tracking branch 'upstream/main' into storage-domains
senamakel Oct 9, 2026
1702d61
test(storage): collapse multi-line assertion into single line
senamakel Oct 9, 2026
8d3f196
feat(desktop): add notification store documents
senamakel Oct 9, 2026
0fce95c
refactor(notifications): split document storage out of notification s…
senamakel Oct 9, 2026
a8ae0be
test(desktop): add coverage for notification store documents
senamakel Oct 9, 2026
af4c134
feat(notifications): add document store notification tests
senamakel Oct 9, 2026
77cf5e4
feat(task_sources): add document store for task sources
senamakel Oct 9, 2026
482d788
feat(task_sources): add document store for task sources
senamakel Oct 9, 2026
f6ae6c9
style: reformat notification and task source store code
senamakel Oct 9, 2026
512275f
refactor(task_sources): simplify document query result handling
senamakel Oct 9, 2026
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
6 changes: 6 additions & 0 deletions crates/openhuman-cli/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,12 @@ path = "../../tests/cost_cap_removed.rs"
name = "storage_approvals_e2e"
path = "../../tests/storage_approvals_e2e.rs"

[[test]]
# Its own binary for the same reason as storage_approvals_e2e: it installs a
# storage backend into the process-wide slot.
name = "storage_domains_e2e"
path = "../../tests/storage_domains_e2e.rs"
Comment thread
senamakel marked this conversation as resolved.

[[test]]
name = "embedded_server_shutdown_e2e"
path = "../../tests/embedded_server_shutdown_e2e.rs"
Expand Down
16 changes: 16 additions & 0 deletions crates/openhuman-core/src/desktop/notifications/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,22 @@ SQLite DB at `{workspace_dir}/notifications/notifications.db`, opened per-call v

`insert_if_not_recent` runs a `BEGIN IMMEDIATE` transaction so concurrent duplicate ingests collapse to a single insert.

### On a storage backend

When the host configured a storage backend (`OPENHUMAN_STORAGE_URL` /
`[storage] url`, see `crate::storage`), every `store` function uses
`store_documents.rs` instead of `notifications.db`: the same operations on the
`tinystoragedrivers` document port, under the current call's storage scope
(the acting agent; `local` on a single-user host; refused in SaaS mode with
no acting agent). Collections `integration_notifications`, `notification_dedup`,
`notification_settings` and `core_notifications`. Every insert first
advances a `notification_dedup` document (a hash of provider, account,
title and body) under compare-and-swap, which replaces the SQL store's
`BEGIN IMMEDIATE`: two processes ingesting the same content in the same
minute insert it once. `stats` folds the counts in the store, since the
port has no `GROUP BY`. With no backend configured (the desktop default)
`notifications.db` is used as described above.

## Dependencies

- `crate::core::bus::BUS` and `crate::core::events::DomainEvent` (`BUS.subscribe` for the bridge, `BUS.publish` for triage results); the `EventHandler` trait comes from `tinybus`.
Expand Down
1 change: 1 addition & 0 deletions crates/openhuman-core/src/desktop/notifications/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ pub mod bus;
pub mod rpc;
pub mod schemas;
pub mod store;
mod store_documents;
pub mod types;

pub use bus::{
Expand Down
2 changes: 1 addition & 1 deletion crates/openhuman-core/src/desktop/notifications/rpc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -88,7 +88,7 @@ pub async fn handle_ingest(params: Map<String, Value>) -> Result<Value, String>
// Spawn background triage — the ingest RPC returns immediately.
let id_for_triage = id.clone();
let config_for_triage = config.clone();
tokio::spawn(async move {
crate::core::runtime::spawn_scoped(async move {
let envelope = TriggerEnvelope {
source: TriggerSource::WebviewIntegration {
provider: req.provider.clone(),
Expand Down
133 changes: 66 additions & 67 deletions crates/openhuman-core/src/desktop/notifications/store.rs
Original file line number Diff line number Diff line change
@@ -1,10 +1,13 @@
//! SQLite persistence for `IntegrationNotification` records.
//!
//! With a storage backend configured ([`crate::storage`]) every function
//! here is served from the document port instead (`store_documents.rs`).
//!
//! Uses a synchronous `rusqlite::Connection` opened per call, following the
//! same `with_connection` pattern as the cron domain.

use anyhow::{Context, Result};
use chrono::{DateTime, Utc};
use chrono::Utc;
use rusqlite::{params, Connection};

use crate::config::Config;
Expand Down Expand Up @@ -106,6 +109,9 @@ fn with_connection<T>(config: &Config, f: impl FnOnce(&Connection) -> Result<T>)

/// Persist a new notification to the store.
pub fn insert(config: &Config, n: &IntegrationNotification) -> Result<()> {
if let Some(docs) = super::store_documents::current()? {
return docs.insert(n, false).map(|_| ());
}
with_connection(config, |conn| {
conn.execute(
"INSERT INTO integration_notifications
Expand Down Expand Up @@ -144,6 +150,9 @@ pub fn insert(config: &Config, n: &IntegrationNotification) -> Result<()> {
/// event is ignored (no duplicates). Returns `true` when a new row was written,
/// `false` when an event with the same id already existed.
pub fn insert_core_notification(config: &Config, event: &CoreNotificationEvent) -> Result<bool> {
if let Some(docs) = super::store_documents::current()? {
return docs.insert_core_notification(&config.workspace_dir.to_string_lossy(), event);
}
with_connection(config, |conn| {
let payload = serde_json::to_string(event)
.context("[notifications::store] serialize core notification failed")?;
Expand Down Expand Up @@ -172,6 +181,13 @@ pub fn list_core_notifications(
only_unread: bool,
limit: usize,
) -> Result<Vec<CoreNotificationEvent>> {
if let Some(docs) = super::store_documents::current()? {
return docs.list_core_notifications(
&config.workspace_dir.to_string_lossy(),
only_unread,
limit,
);
}
with_connection(config, |conn| {
let sql = if only_unread {
"SELECT payload FROM core_notifications WHERE read = 0
Expand Down Expand Up @@ -207,6 +223,9 @@ pub fn list_core_notifications(
/// Mark a persisted core notification as read so it isn't re-surfaced on the
/// next sync-down. Returns `true` when a row was updated.
pub fn mark_core_notification_read(config: &Config, id: &str) -> Result<bool> {
if let Some(docs) = super::store_documents::current()? {
return docs.mark_core_notification_read(&config.workspace_dir.to_string_lossy(), id);
}
with_connection(config, |conn| {
let affected = conn
.execute(
Expand All @@ -220,6 +239,9 @@ pub fn mark_core_notification_read(config: &Config, id: &str) -> Result<bool> {

/// Count unread persisted core notifications.
pub fn unread_core_notification_count(config: &Config) -> Result<i64> {
if let Some(docs) = super::store_documents::current()? {
return docs.unread_core_notification_count(&config.workspace_dir.to_string_lossy());
}
with_connection(config, |conn| {
let count: i64 = conn
.query_row(
Expand All @@ -236,6 +258,9 @@ pub fn unread_core_notification_count(config: &Config) -> Result<i64> {
///
/// Returns `true` when inserted, `false` when skipped as duplicate.
pub fn insert_if_not_recent(config: &Config, n: &IntegrationNotification) -> Result<bool> {
if let Some(docs) = super::store_documents::current()? {
return docs.insert(n, true);
}
with_connection(config, |conn| {
conn.execute_batch("BEGIN IMMEDIATE")
.context("[notifications::store] begin insert_if_not_recent tx failed")?;
Expand Down Expand Up @@ -313,6 +338,9 @@ pub fn list(
provider_filter: Option<&str>,
min_score: Option<f32>,
) -> Result<Vec<IntegrationNotification>> {
if let Some(docs) = super::store_documents::current()? {
return docs.list(limit, offset, provider_filter, min_score);
}
with_connection(config, |conn| {
// Build a dynamic query instead of relying on nullable-aware WHERE
// logic so the SQL stays readable for future contributors.
Expand Down Expand Up @@ -360,6 +388,13 @@ pub fn update_triage(
action: &str,
reason: &str,
) -> Result<()> {
if let Some(docs) = super::store_documents::current()? {
Comment thread
senamakel marked this conversation as resolved.
return docs.update_triage(id, score, action, reason).map(|found| {
if !found {
tracing::warn!(id = %id, action = %action, "[notifications::store] update_triage matched no rows");
}
});
}
with_connection(config, |conn| {
let now = Utc::now().to_rfc3339();
let updated = conn
Expand Down Expand Up @@ -393,6 +428,13 @@ pub fn update_triage(

/// Transition a notification from `unread` to `read`.
pub fn mark_read(config: &Config, id: &str) -> Result<()> {
if let Some(docs) = super::store_documents::current()? {
return docs.set_status(id, NotificationStatus::Read).map(|found| {
if !found {
tracing::warn!(id = %id, "[notifications::store] mark_read matched no rows");
}
});
}
with_connection(config, |conn| {
let updated = conn
.execute(
Expand All @@ -414,6 +456,9 @@ pub fn mark_read(config: &Config, id: &str) -> Result<()> {

/// Count unread notifications.
pub fn unread_count(config: &Config) -> Result<i64> {
if let Some(docs) = super::store_documents::current()? {
return docs.unread_count();
}
with_connection(config, |conn| {
let count: i64 = conn
.query_row(
Expand All @@ -435,6 +480,9 @@ pub fn exists_recent(
title: &str,
body: &str,
) -> Result<bool> {
if let Some(docs) = super::store_documents::current()? {
return docs.exists_recent(provider, account_id, title, body);
}
with_connection(config, |conn| {
let count: i64 = match account_id {
Some(aid) => conn.query_row(
Expand Down Expand Up @@ -463,6 +511,9 @@ pub fn exists_recent(
///
/// Returns `true` when at least one row matched and was updated.
pub fn mark_dismissed(config: &Config, id: &str) -> Result<bool> {
if let Some(docs) = super::store_documents::current()? {
return docs.set_status(id, NotificationStatus::Dismissed);
}
with_connection(config, |conn| {
let updated = conn
.execute(
Expand All @@ -484,6 +535,9 @@ pub fn mark_dismissed(config: &Config, id: &str) -> Result<bool> {
///
/// Returns `true` when at least one row matched and was updated.
pub fn mark_acted(config: &Config, id: &str) -> Result<bool> {
if let Some(docs) = super::store_documents::current()? {
return docs.set_status(id, NotificationStatus::Acted);
}
with_connection(config, |conn| {
let updated = conn
.execute(
Expand All @@ -503,6 +557,9 @@ pub fn mark_acted(config: &Config, id: &str) -> Result<bool> {

/// Return aggregate statistics for the notification intelligence pipeline.
pub fn stats(config: &Config) -> Result<super::types::NotificationStats> {
if let Some(docs) = super::store_documents::current()? {
return docs.stats();
}
use std::collections::HashMap;
with_connection(config, |conn| {
let total: i64 = conn
Expand Down Expand Up @@ -587,6 +644,9 @@ pub fn stats(config: &Config) -> Result<super::types::NotificationStats> {

/// Upsert provider-level notification settings.
pub fn upsert_settings(config: &Config, settings: &NotificationSettings) -> Result<()> {
if let Some(docs) = super::store_documents::current()? {
return docs.upsert_settings(settings);
}
with_connection(config, |conn| {
conn.execute(
"INSERT INTO notification_settings (provider, enabled, importance_threshold, route_to_orchestrator)
Expand All @@ -609,6 +669,9 @@ pub fn upsert_settings(config: &Config, settings: &NotificationSettings) -> Resu

/// Read provider-level notification settings with defaults when missing.
pub fn get_settings(config: &Config, provider: &str) -> Result<NotificationSettings> {
if let Some(docs) = super::store_documents::current()? {
return docs.get_settings(provider);
}
with_connection(config, |conn| {
let mut stmt = conn
.prepare(
Expand Down Expand Up @@ -638,72 +701,8 @@ pub fn get_settings(config: &Config, provider: &str) -> Result<NotificationSetti
})
}

// ─────────────────────────────────────────────────────────────────────────────
// Row conversion helpers
// ─────────────────────────────────────────────────────────────────────────────

fn rows_to_notifications(mut rows: rusqlite::Rows<'_>) -> Result<Vec<IntegrationNotification>> {
let mut out = Vec::new();
while let Some(row) = rows
.next()
.context("[notifications::store] row iteration failed")?
{
out.push(row_to_notification(row)?);
}
Ok(out)
}

fn row_to_notification(row: &rusqlite::Row<'_>) -> Result<IntegrationNotification> {
let raw_payload_str: String = row.get(5)?;
let raw_payload: serde_json::Value = serde_json::from_str(&raw_payload_str)
.unwrap_or(serde_json::Value::String(raw_payload_str));

let status_str: String = row.get(9)?;
let status = match status_str.as_str() {
"read" => NotificationStatus::Read,
"acted" => NotificationStatus::Acted,
"dismissed" => NotificationStatus::Dismissed,
_ => NotificationStatus::Unread,
};

let received_at_str: String = row.get(10)?;
let received_at: DateTime<Utc> = received_at_str.parse().unwrap_or_else(|e| {
tracing::warn!(
raw = %received_at_str,
error = %e,
"[notifications::store] invalid received_at, using now"
);
Utc::now()
});

let scored_at_str: Option<String> = row.get(11)?;
let scored_at: Option<DateTime<Utc>> = scored_at_str.and_then(|s| match s.parse() {
Ok(t) => Some(t),
Err(e) => {
tracing::warn!(
raw = %s,
error = %e,
"[notifications::store] invalid scored_at, treating as unscored"
);
None
}
});

Ok(IntegrationNotification {
id: row.get(0)?,
provider: row.get(1)?,
account_id: row.get(2)?,
title: row.get(3)?,
body: row.get(4)?,
raw_payload,
importance_score: row.get(6)?,
triage_action: row.get(7)?,
triage_reason: row.get(8)?,
status,
received_at,
scored_at,
})
}
mod store_rows;
Comment thread
senamakel marked this conversation as resolved.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

priority critical security confident

Add the missing store\_rows module source

This declaration requires store_rows.rs (or an inline module body), but the pull request does not provide that source. The crate therefore fails to compile. Add the module containing rows_to_notifications, or remove the declaration and restore the helper implementation here.


Additional critique observation

priority critical confident

Add the missing store_rows module source

[RULE] build-break

Rust resolves this declaration to store_rows.rs or store_rows/mod.rs beside store.rs. Neither source file is present in the complete diff, so this change fails to compile unless an existing unshown file supplies it. Add the module source (including rows_to_notifications) with the change.

[RULE] missing-module ·

use store_rows::rows_to_notifications;

#[cfg(test)]
#[path = "store_tests.rs"]
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
//! Row conversion for the SQLite notification store.

use anyhow::{Context, Result};
use chrono::{DateTime, Utc};

use super::super::types::{IntegrationNotification, NotificationStatus};

pub(super) fn rows_to_notifications(
mut rows: rusqlite::Rows<'_>,
) -> Result<Vec<IntegrationNotification>> {
let mut out = Vec::new();
while let Some(row) = rows
.next()
.context("[notifications::store] row iteration failed")?
{
out.push(row_to_notification(row)?);
}
Ok(out)
}

fn row_to_notification(row: &rusqlite::Row<'_>) -> Result<IntegrationNotification> {
let raw_payload_str: String = row.get(5)?;
let raw_payload: serde_json::Value = serde_json::from_str(&raw_payload_str)
.unwrap_or(serde_json::Value::String(raw_payload_str));

let status_str: String = row.get(9)?;
let status = match status_str.as_str() {
"read" => NotificationStatus::Read,
"acted" => NotificationStatus::Acted,
"dismissed" => NotificationStatus::Dismissed,
_ => NotificationStatus::Unread,
Comment thread
senamakel marked this conversation as resolved.
};

let received_at_str: String = row.get(10)?;
let received_at: DateTime<Utc> = received_at_str.parse().unwrap_or_else(|e| {
tracing::warn!(
raw = %received_at_str,
error = %e,
"[notifications::store] invalid received_at, using now"
);
Utc::now()
Comment thread
senamakel marked this conversation as resolved.
});

let scored_at_str: Option<String> = row.get(11)?;
let scored_at: Option<DateTime<Utc>> = scored_at_str.and_then(|s| match s.parse() {
Comment thread
senamakel marked this conversation as resolved.
Ok(t) => Some(t),
Err(e) => {
tracing::warn!(
raw = %s,
error = %e,
"[notifications::store] invalid scored_at, treating as unscored"
);
None
Comment thread
senamakel marked this conversation as resolved.
}
});

Ok(IntegrationNotification {
id: row.get(0)?,
provider: row.get(1)?,
account_id: row.get(2)?,
title: row.get(3)?,
body: row.get(4)?,
raw_payload,
importance_score: row.get(6)?,
triage_action: row.get(7)?,
triage_reason: row.get(8)?,
status,
received_at,
scored_at,
})
}
Loading
Loading