Compare commits
No commits in common. "257bc0110ef4332bbbb21e98528650644a2b3ed4" and "c94fc5308fa3cd5b676542c5357884f47fd00924" have entirely different histories.
257bc0110e
...
c94fc5308f
|
|
@ -190,8 +190,7 @@ export STACKCHAIN_TRUSTED_PROXY_CIDRS='127.0.0.0/8'
|
||||||
# Optional Web Push. Generate a VAPID key pair outside the repo and inject it.
|
# Optional Web Push. Generate a VAPID key pair outside the repo and inject it.
|
||||||
# The feature stays disabled unless all three values are present. Privacy-safe update
|
# The feature stays disabled unless all three values are present. Privacy-safe update
|
||||||
# alerts offer Mark read and Tomorrow; Tomorrow syncs the unread item to Later at
|
# alerts offer Mark read and Tomorrow; Tomorrow syncs the unread item to Later at
|
||||||
# 09:00 in the device's local timezone without opening the dashboard. Bursts send
|
# 09:00 in the device's local timezone without opening the dashboard.
|
||||||
# three individual alerts followed by one private digest that opens Updates.
|
|
||||||
export STACKCHAIN_VAPID_PUBLIC_KEY='<url-safe-public-key>'
|
export STACKCHAIN_VAPID_PUBLIC_KEY='<url-safe-public-key>'
|
||||||
export STACKCHAIN_VAPID_PRIVATE_KEY='<private-key-from-secret-manager>'
|
export STACKCHAIN_VAPID_PRIVATE_KEY='<private-key-from-secret-manager>'
|
||||||
export STACKCHAIN_VAPID_SUBJECT='mailto:ops@example.com'
|
export STACKCHAIN_VAPID_SUBJECT='mailto:ops@example.com'
|
||||||
|
|
@ -202,8 +201,6 @@ export STACKCHAIN_VAPID_SUBJECT='mailto:ops@example.com'
|
||||||
export STACKCHAIN_PUSH_POLL_SECONDS=30
|
export STACKCHAIN_PUSH_POLL_SECONDS=30
|
||||||
export STACKCHAIN_PUSH_SEND_TIMEOUT_SECONDS=10
|
export STACKCHAIN_PUSH_SEND_TIMEOUT_SECONDS=10
|
||||||
export STACKCHAIN_PUSH_MAX_CONCURRENCY=8
|
export STACKCHAIN_PUSH_MAX_CONCURRENCY=8
|
||||||
# Maximum individual alerts per device and poll before one digest covers the rest.
|
|
||||||
export STACKCHAIN_PUSH_MAX_INDIVIDUAL_NOTIFICATIONS=3
|
|
||||||
export STACKCHAIN_PUSH_LEASE_SECONDS=60
|
export STACKCHAIN_PUSH_LEASE_SECONDS=60
|
||||||
export STACKCHAIN_PUSH_DB='/var/lib/stackchain-dashboard/push-subscriptions.sqlite3'
|
export STACKCHAIN_PUSH_DB='/var/lib/stackchain-dashboard/push-subscriptions.sqlite3'
|
||||||
uvicorn src.main:app --host 127.0.0.1 --port 8000
|
uvicorn src.main:app --host 127.0.0.1 --port 8000
|
||||||
|
|
|
||||||
|
|
@ -278,21 +278,6 @@ self.addEventListener('push', event => {
|
||||||
const route = String(payload.route || '');
|
const route = String(payload.route || '');
|
||||||
const tag = String(payload.tag || '');
|
const tag = String(payload.tag || '');
|
||||||
const notificationId = Number(payload.notification_id);
|
const notificationId = Number(payload.notification_id);
|
||||||
const updateCount = Number(payload.update_count);
|
|
||||||
if (
|
|
||||||
route === '#/my-work/updates'
|
|
||||||
&& tag === 'stackchain-update-digest'
|
|
||||||
&& Number.isSafeInteger(updateCount)
|
|
||||||
&& updateCount > 0
|
|
||||||
&& updateCount <= 50
|
|
||||||
) {
|
|
||||||
event.waitUntil(self.registration.showNotification(updateCount + ' new work updates', {
|
|
||||||
body: 'Tap to review them in Stackchain.',
|
|
||||||
tag,
|
|
||||||
data: {route},
|
|
||||||
}));
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
if (!/^#\/my-work\/update\/\d+$/.test(route) || !/^stackchain-update-\d+$/.test(tag)) return;
|
if (!/^#\/my-work\/update\/\d+$/.test(route) || !/^stackchain-update-\d+$/.test(tag)) return;
|
||||||
const options = {
|
const options = {
|
||||||
body: 'Tap to review it in Stackchain.',
|
body: 'Tap to review it in Stackchain.',
|
||||||
|
|
|
||||||
|
|
@ -90,10 +90,6 @@ async def _push_poll_loop() -> None:
|
||||||
max_concurrency = max(
|
max_concurrency = max(
|
||||||
1, int(os.getenv("STACKCHAIN_PUSH_MAX_CONCURRENCY", "8"))
|
1, int(os.getenv("STACKCHAIN_PUSH_MAX_CONCURRENCY", "8"))
|
||||||
)
|
)
|
||||||
max_individual_notifications = max(
|
|
||||||
0,
|
|
||||||
int(os.getenv("STACKCHAIN_PUSH_MAX_INDIVIDUAL_NOTIFICATIONS", "3")),
|
|
||||||
)
|
|
||||||
while True:
|
while True:
|
||||||
await asyncio.sleep(interval)
|
await asyncio.sleep(interval)
|
||||||
try:
|
try:
|
||||||
|
|
@ -105,7 +101,6 @@ async def _push_poll_loop() -> None:
|
||||||
lease_seconds=lease_seconds,
|
lease_seconds=lease_seconds,
|
||||||
send_timeout_seconds=send_timeout,
|
send_timeout_seconds=send_timeout,
|
||||||
max_concurrency=max_concurrency,
|
max_concurrency=max_concurrency,
|
||||||
max_individual_notifications=max_individual_notifications,
|
|
||||||
)
|
)
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
raise
|
raise
|
||||||
|
|
|
||||||
|
|
@ -44,7 +44,6 @@ async def dispatch_unread_updates(
|
||||||
lease_seconds: float = 60.0,
|
lease_seconds: float = 60.0,
|
||||||
send_timeout_seconds: float = 10.0,
|
send_timeout_seconds: float = 10.0,
|
||||||
max_concurrency: int = 8,
|
max_concurrency: int = 8,
|
||||||
max_individual_notifications: int = 3,
|
|
||||||
) -> int:
|
) -> int:
|
||||||
if not configuration.enabled:
|
if not configuration.enabled:
|
||||||
return 0
|
return 0
|
||||||
|
|
@ -88,16 +87,7 @@ async def dispatch_unread_updates(
|
||||||
if ownership_lost.is_set():
|
if ownership_lost.is_set():
|
||||||
return 0
|
return 0
|
||||||
count = 0
|
count = 0
|
||||||
digest_pending = set(delivery.digest_ids)
|
for thread_id in delivery.thread_ids:
|
||||||
new_ids = tuple(
|
|
||||||
thread_id
|
|
||||||
for thread_id in delivery.thread_ids
|
|
||||||
if thread_id not in digest_pending
|
|
||||||
)
|
|
||||||
individual_ids = new_ids[:max(0, max_individual_notifications)]
|
|
||||||
overflow_ids = delivery.digest_ids + new_ids[len(individual_ids):]
|
|
||||||
can_send_digest = True
|
|
||||||
for thread_id in individual_ids:
|
|
||||||
still_owner = await asyncio.to_thread(
|
still_owner = await asyncio.to_thread(
|
||||||
store.acquire_dispatch_lease,
|
store.acquire_dispatch_lease,
|
||||||
owner,
|
owner,
|
||||||
|
|
@ -133,7 +123,7 @@ async def dispatch_unread_updates(
|
||||||
await asyncio.to_thread(
|
await asyncio.to_thread(
|
||||||
store.delete_session, delivery.session_id
|
store.delete_session, delivery.session_id
|
||||||
)
|
)
|
||||||
can_send_digest = False
|
break
|
||||||
# Leave this device's transient failures unseen for a later
|
# Leave this device's transient failures unseen for a later
|
||||||
# poll instead of paying the endpoint deadline repeatedly.
|
# poll instead of paying the endpoint deadline repeatedly.
|
||||||
break
|
break
|
||||||
|
|
@ -141,53 +131,6 @@ async def dispatch_unread_updates(
|
||||||
store.mark_delivered, delivery.session_id, (thread_id,)
|
store.mark_delivered, delivery.session_id, (thread_id,)
|
||||||
)
|
)
|
||||||
count += 1
|
count += 1
|
||||||
if overflow_ids and can_send_digest:
|
|
||||||
still_owner = await asyncio.to_thread(
|
|
||||||
store.acquire_dispatch_lease,
|
|
||||||
owner,
|
|
||||||
now=time.time(),
|
|
||||||
lease_seconds=lease_seconds,
|
|
||||||
)
|
|
||||||
if not still_owner:
|
|
||||||
ownership_lost.set()
|
|
||||||
return count
|
|
||||||
payload = json.dumps(
|
|
||||||
{
|
|
||||||
"title": f"{len(overflow_ids)} new work updates",
|
|
||||||
"body": "Tap to review them in Stackchain.",
|
|
||||||
"route": "#/my-work/updates",
|
|
||||||
"tag": "stackchain-update-digest",
|
|
||||||
"update_count": len(overflow_ids),
|
|
||||||
},
|
|
||||||
separators=(",", ":"),
|
|
||||||
)
|
|
||||||
try:
|
|
||||||
if send is None:
|
|
||||||
operation = send_web_push(
|
|
||||||
delivery.subscription, payload, configuration
|
|
||||||
)
|
|
||||||
else:
|
|
||||||
operation = send(delivery.subscription, payload)
|
|
||||||
await asyncio.wait_for(operation, timeout=send_timeout_seconds)
|
|
||||||
except Exception as error:
|
|
||||||
status = getattr(
|
|
||||||
getattr(error, "response", None), "status_code", None
|
|
||||||
)
|
|
||||||
if status in {404, 410}:
|
|
||||||
await asyncio.to_thread(
|
|
||||||
store.delete_session, delivery.session_id
|
|
||||||
)
|
|
||||||
else:
|
|
||||||
await asyncio.to_thread(
|
|
||||||
store.mark_digest_pending,
|
|
||||||
delivery.session_id,
|
|
||||||
overflow_ids,
|
|
||||||
)
|
|
||||||
return count
|
|
||||||
await asyncio.to_thread(
|
|
||||||
store.mark_delivered, delivery.session_id, overflow_ids
|
|
||||||
)
|
|
||||||
count += 1
|
|
||||||
return count
|
return count
|
||||||
|
|
||||||
counts = await asyncio.gather(
|
counts = await asyncio.gather(
|
||||||
|
|
|
||||||
|
|
@ -11,7 +11,6 @@ class PushDelivery:
|
||||||
session_id: str
|
session_id: str
|
||||||
subscription: dict
|
subscription: dict
|
||||||
thread_ids: tuple[int, ...]
|
thread_ids: tuple[int, ...]
|
||||||
digest_ids: tuple[int, ...] = ()
|
|
||||||
|
|
||||||
|
|
||||||
class PushSubscriptionStore:
|
class PushSubscriptionStore:
|
||||||
|
|
@ -41,13 +40,6 @@ class PushSubscriptionStore:
|
||||||
owner TEXT NOT NULL,
|
owner TEXT NOT NULL,
|
||||||
expires_at REAL NOT NULL
|
expires_at REAL NOT NULL
|
||||||
);
|
);
|
||||||
CREATE TABLE IF NOT EXISTS push_digest_pending (
|
|
||||||
session_id TEXT NOT NULL,
|
|
||||||
thread_id INTEGER NOT NULL,
|
|
||||||
PRIMARY KEY (session_id, thread_id),
|
|
||||||
FOREIGN KEY (session_id) REFERENCES push_subscriptions(session_id)
|
|
||||||
ON DELETE CASCADE
|
|
||||||
);
|
|
||||||
"""
|
"""
|
||||||
)
|
)
|
||||||
os.chmod(self.path, 0o600)
|
os.chmod(self.path, 0o600)
|
||||||
|
|
@ -128,34 +120,12 @@ class PushSubscriptionStore:
|
||||||
}
|
}
|
||||||
unseen = tuple(value for value in candidates if value not in delivered)
|
unseen = tuple(value for value in candidates if value not in delivered)
|
||||||
if unseen:
|
if unseen:
|
||||||
pending = {
|
deliveries.append(PushDelivery(session_id, json.loads(encoded), unseen))
|
||||||
row[0]
|
|
||||||
for row in connection.execute(
|
|
||||||
"SELECT thread_id FROM push_digest_pending WHERE session_id = ?",
|
|
||||||
(session_id,),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
digest_ids = tuple(value for value in unseen if value in pending)
|
|
||||||
deliveries.append(
|
|
||||||
PushDelivery(session_id, json.loads(encoded), unseen, digest_ids)
|
|
||||||
)
|
|
||||||
return deliveries
|
return deliveries
|
||||||
|
|
||||||
def mark_digest_pending(self, session_id: str, thread_ids: Iterable[int]) -> None:
|
|
||||||
with self._connect() as connection:
|
|
||||||
connection.executemany(
|
|
||||||
"INSERT OR IGNORE INTO push_digest_pending(session_id, thread_id) VALUES (?, ?)",
|
|
||||||
((session_id, int(thread_id)) for thread_id in thread_ids),
|
|
||||||
)
|
|
||||||
|
|
||||||
def mark_delivered(self, session_id: str, thread_ids: Iterable[int]) -> None:
|
def mark_delivered(self, session_id: str, thread_ids: Iterable[int]) -> None:
|
||||||
values = tuple(int(thread_id) for thread_id in thread_ids)
|
|
||||||
with self._connect() as connection:
|
with self._connect() as connection:
|
||||||
connection.executemany(
|
connection.executemany(
|
||||||
"INSERT OR IGNORE INTO push_deliveries(session_id, thread_id) VALUES (?, ?)",
|
"INSERT OR IGNORE INTO push_deliveries(session_id, thread_id) VALUES (?, ?)",
|
||||||
((session_id, thread_id) for thread_id in values),
|
((session_id, int(thread_id)) for thread_id in thread_ids),
|
||||||
)
|
|
||||||
connection.executemany(
|
|
||||||
"DELETE FROM push_digest_pending WHERE session_id = ? AND thread_id = ?",
|
|
||||||
((session_id, thread_id) for thread_id in values),
|
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -85,7 +85,6 @@ async def test_push_poll_authorizes_delivery_against_managed_session(monkeypatch
|
||||||
raise asyncio.CancelledError
|
raise asyncio.CancelledError
|
||||||
|
|
||||||
monkeypatch.setenv("STACKCHAIN_PUSH_MAX_CONCURRENCY", "3")
|
monkeypatch.setenv("STACKCHAIN_PUSH_MAX_CONCURRENCY", "3")
|
||||||
monkeypatch.setenv("STACKCHAIN_PUSH_MAX_INDIVIDUAL_NOTIFICATIONS", "4")
|
|
||||||
monkeypatch.setattr(main.asyncio, "sleep", no_wait)
|
monkeypatch.setattr(main.asyncio, "sleep", no_wait)
|
||||||
monkeypatch.setattr(main, "dispatch_unread_updates", stop_after_capture)
|
monkeypatch.setattr(main, "dispatch_unread_updates", stop_after_capture)
|
||||||
|
|
||||||
|
|
@ -94,7 +93,6 @@ async def test_push_poll_authorizes_delivery_against_managed_session(monkeypatch
|
||||||
|
|
||||||
assert captured["session_active"] is dashboard_auth.managed_session_active
|
assert captured["session_active"] is dashboard_auth.managed_session_active
|
||||||
assert captured["max_concurrency"] == 3
|
assert captured["max_concurrency"] == 3
|
||||||
assert captured["max_individual_notifications"] == 4
|
|
||||||
|
|
||||||
|
|
||||||
def test_replacing_subscription_resets_delivery_cursor_and_revocation_removes_device(tmp_path):
|
def test_replacing_subscription_resets_delivery_cursor_and_revocation_removes_device(tmp_path):
|
||||||
|
|
@ -150,108 +148,6 @@ async def test_dispatch_sends_one_privacy_safe_deep_link_per_new_thread(tmp_path
|
||||||
assert "Secret title" not in json.dumps(sent)
|
assert "Secret title" not in json.dumps(sent)
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.anyio
|
|
||||||
async def test_update_burst_sends_bounded_individual_pushes_and_one_private_digest(tmp_path):
|
|
||||||
store = PushSubscriptionStore(tmp_path / "push.sqlite3")
|
|
||||||
store.upsert("session-a", {
|
|
||||||
"endpoint": "https://push.example/device-a",
|
|
||||||
"keys": {"p256dh": "public-key", "auth": "auth-secret"},
|
|
||||||
})
|
|
||||||
sent = []
|
|
||||||
|
|
||||||
async def unread():
|
|
||||||
return {"items": [
|
|
||||||
{
|
|
||||||
"id": notification_id,
|
|
||||||
"repository": "private/repo",
|
|
||||||
"title": f"Secret update {notification_id}",
|
|
||||||
}
|
|
||||||
for notification_id in range(41, 46)
|
|
||||||
]}
|
|
||||||
|
|
||||||
async def send(_subscription, payload):
|
|
||||||
sent.append(json.loads(payload))
|
|
||||||
|
|
||||||
config = PushConfiguration("public", "private", "mailto:ops@example.com")
|
|
||||||
delivered = await dispatch_unread_updates(
|
|
||||||
store,
|
|
||||||
config,
|
|
||||||
unread,
|
|
||||||
send,
|
|
||||||
max_individual_notifications=2,
|
|
||||||
)
|
|
||||||
|
|
||||||
assert delivered == 3
|
|
||||||
assert [payload["tag"] for payload in sent] == [
|
|
||||||
"stackchain-update-41",
|
|
||||||
"stackchain-update-42",
|
|
||||||
"stackchain-update-digest",
|
|
||||||
]
|
|
||||||
assert sent[-1] == {
|
|
||||||
"title": "3 new work updates",
|
|
||||||
"body": "Tap to review them in Stackchain.",
|
|
||||||
"route": "#/my-work/updates",
|
|
||||||
"tag": "stackchain-update-digest",
|
|
||||||
"update_count": 3,
|
|
||||||
}
|
|
||||||
assert "private/repo" not in json.dumps(sent)
|
|
||||||
assert "Secret update" not in json.dumps(sent)
|
|
||||||
assert await dispatch_unread_updates(
|
|
||||||
store,
|
|
||||||
config,
|
|
||||||
unread,
|
|
||||||
send,
|
|
||||||
max_individual_notifications=2,
|
|
||||||
) == 0
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.anyio
|
|
||||||
async def test_failed_digest_retries_only_overflow_after_individual_checkpoints(tmp_path):
|
|
||||||
store = PushSubscriptionStore(tmp_path / "push.sqlite3")
|
|
||||||
store.upsert("session-a", {
|
|
||||||
"endpoint": "https://push.example/device-a",
|
|
||||||
"keys": {"p256dh": "public-key", "auth": "auth-secret"},
|
|
||||||
})
|
|
||||||
attempts = []
|
|
||||||
|
|
||||||
async def unread():
|
|
||||||
return {"items": [{"id": notification_id} for notification_id in range(51, 56)]}
|
|
||||||
|
|
||||||
async def fail_digest(_subscription, payload):
|
|
||||||
decoded = json.loads(payload)
|
|
||||||
attempts.append(decoded["tag"])
|
|
||||||
if decoded["tag"] == "stackchain-update-digest":
|
|
||||||
raise RuntimeError("push provider unavailable")
|
|
||||||
|
|
||||||
config = PushConfiguration("public", "private", "mailto:ops@example.com")
|
|
||||||
assert await dispatch_unread_updates(
|
|
||||||
store,
|
|
||||||
config,
|
|
||||||
unread,
|
|
||||||
fail_digest,
|
|
||||||
max_individual_notifications=2,
|
|
||||||
) == 2
|
|
||||||
assert store.claim_unseen(range(51, 56))[0].thread_ids == (53, 54, 55)
|
|
||||||
|
|
||||||
async def succeed(_subscription, payload):
|
|
||||||
attempts.append(json.loads(payload)["tag"])
|
|
||||||
|
|
||||||
assert await dispatch_unread_updates(
|
|
||||||
store,
|
|
||||||
config,
|
|
||||||
unread,
|
|
||||||
succeed,
|
|
||||||
max_individual_notifications=2,
|
|
||||||
) == 1
|
|
||||||
assert attempts == [
|
|
||||||
"stackchain-update-51",
|
|
||||||
"stackchain-update-52",
|
|
||||||
"stackchain-update-digest",
|
|
||||||
"stackchain-update-digest",
|
|
||||||
]
|
|
||||||
assert store.claim_unseen(range(51, 56)) == []
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.anyio
|
@pytest.mark.anyio
|
||||||
async def test_production_notification_page_dispatches_one_unread_push(tmp_path):
|
async def test_production_notification_page_dispatches_one_unread_push(tmp_path):
|
||||||
def upstream(_request):
|
def upstream(_request):
|
||||||
|
|
|
||||||
|
|
@ -417,32 +417,6 @@ def test_inbound_push_renders_generic_update_and_tap_opens_existing_workflow():
|
||||||
assert "must-not-render" not in json.dumps(result["notifications"])
|
assert "must-not-render" not in json.dumps(result["notifications"])
|
||||||
|
|
||||||
|
|
||||||
def test_update_digest_push_opens_unread_inbox_without_item_actions_or_private_copy():
|
|
||||||
result = run_worker_scenario(
|
|
||||||
"""
|
|
||||||
await dispatchPush({
|
|
||||||
title:'must-not-render', body:'private details must-not-render',
|
|
||||||
tag:'stackchain-update-digest', route:'#/my-work/updates', update_count:7,
|
|
||||||
});
|
|
||||||
await dispatchNotificationClick('#/my-work/updates');
|
|
||||||
process.stdout.write(JSON.stringify(state));
|
|
||||||
"""
|
|
||||||
)
|
|
||||||
|
|
||||||
assert result["notifications"] == [{
|
|
||||||
"title": "7 new work updates",
|
|
||||||
"options": {
|
|
||||||
"body": "Tap to review them in Stackchain.",
|
|
||||||
"tag": "stackchain-update-digest",
|
|
||||||
"data": {"route": "#/my-work/updates"},
|
|
||||||
},
|
|
||||||
}]
|
|
||||||
assert result["opened"] == [
|
|
||||||
"https://forge.example/dashboard/#/my-work/updates"
|
|
||||||
]
|
|
||||||
assert "must-not-render" not in json.dumps(result["notifications"])
|
|
||||||
|
|
||||||
|
|
||||||
def test_push_mark_read_action_confirms_authenticated_mutation_without_opening_app():
|
def test_push_mark_read_action_confirms_authenticated_mutation_without_opening_app():
|
||||||
result = run_worker_scenario(
|
result = run_worker_scenario(
|
||||||
"""
|
"""
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue
Block a user