commit 5aa6c81c78652e503fa2160356ca67d8d83c2a38
parent 1e494bcc07875cfce2e028bf7529ef36a5919cbe
Author: triesap <tyson@radroots.org>
Date: Fri, 17 Jul 2026 22:29:49 +0000
nip46: recover logout acknowledgement finalization
- persist logout acknowledgements as a distinct delivery job
- retry failed acknowledgement publication during startup
- revoke sessions only after durable publish confirmation
- verify crash recovery without duplicate publication
Diffstat:
7 files changed, 322 insertions(+), 15 deletions(-)
diff --git a/src/app/runtime.rs b/src/app/runtime.rs
@@ -328,7 +328,16 @@ impl MycRuntime {
let published_records = self
.delivery_outbox_store
.list_by_status(MycDeliveryOutboxStatus::PublishedPendingFinalize)?;
- if queued_records.is_empty() && published_records.is_empty() {
+ let failed_logout_records = self
+ .delivery_outbox_store
+ .list_by_status(MycDeliveryOutboxStatus::Failed)?
+ .into_iter()
+ .filter(|record| record.kind == MycDeliveryOutboxKind::LogoutAcknowledgementPublish)
+ .collect::<Vec<_>>();
+ if queued_records.is_empty()
+ && published_records.is_empty()
+ && failed_logout_records.is_empty()
+ {
if let Err(error) = self.ensure_no_orphaned_publish_workflows() {
self.record_delivery_recovery_summary(
MycOperationAuditOutcome::Rejected,
@@ -343,6 +352,7 @@ impl MycRuntime {
}
queued_records.extend(published_records);
+ queued_records.extend(failed_logout_records);
queued_records.sort_by(|left, right| {
left.created_at_unix
.cmp(&right.created_at_unix)
@@ -445,7 +455,12 @@ impl MycRuntime {
);
match record.status {
- MycDeliveryOutboxStatus::Queued => {
+ MycDeliveryOutboxStatus::Queued | MycDeliveryOutboxStatus::Failed => {
+ if record.status == MycDeliveryOutboxStatus::Failed
+ && record.kind != MycDeliveryOutboxKind::LogoutAcknowledgementPublish
+ {
+ return Ok(false);
+ }
if record.signer_publish_workflow_id.is_some() && workflow.is_none() {
return Err(self.wrap_recovery_error(
&record,
@@ -520,7 +535,7 @@ impl MycRuntime {
self.finalize_recovered_delivery_job(manager, record, workflow.as_ref(), None)?;
Ok(false)
}
- MycDeliveryOutboxStatus::Finalized | MycDeliveryOutboxStatus::Failed => Ok(false),
+ MycDeliveryOutboxStatus::Finalized => Ok(false),
}
}
@@ -551,6 +566,21 @@ impl MycRuntime {
)),
)
})?;
+ } else if record.kind == MycDeliveryOutboxKind::LogoutAcknowledgementPublish {
+ let connection = self.recovery_connection_record(manager, &record)?;
+ manager
+ .revoke_connection(
+ &connection.connection_id,
+ Some("NIP-46 logout acknowledged".to_owned()),
+ )
+ .map_err(|error| {
+ self.wrap_recovery_error(
+ &record,
+ MycError::InvalidOperation(format!(
+ "failed to finalize NIP-46 logout during startup recovery: {error}"
+ )),
+ )
+ })?;
} else {
self.ensure_record_is_already_finalized_without_workflow(manager, &record)?;
}
@@ -609,6 +639,15 @@ impl MycRuntime {
)),
));
}
+ MycDeliveryOutboxKind::LogoutAcknowledgementPublish => {
+ return Err(self.wrap_recovery_error(
+ record,
+ MycError::InvalidOperation(format!(
+ "logout acknowledgement delivery outbox job `{}` unexpectedly references signer workflow `{workflow_id}`",
+ record.job_id
+ )),
+ ));
+ }
}
Ok(())
@@ -643,12 +682,24 @@ impl MycRuntime {
record: &MycDeliveryOutboxRecord,
) -> Result<(), MycError> {
match record.kind {
- MycDeliveryOutboxKind::DiscoveryHandlerPublish => {
+ MycDeliveryOutboxKind::DiscoveryHandlerPublish
+ | MycDeliveryOutboxKind::LogoutAcknowledgementPublish => {
if record.signer_publish_workflow_id.is_some() {
return Err(self.wrap_recovery_error(
record,
+ MycError::InvalidOperation(format!(
+ "{:?} delivery outbox jobs must not reference signer publish workflows",
+ record.kind
+ )),
+ ));
+ }
+ if record.kind == MycDeliveryOutboxKind::LogoutAcknowledgementPublish
+ && record.connection_id.is_none()
+ {
+ return Err(self.wrap_recovery_error(
+ record,
MycError::InvalidOperation(
- "discovery delivery outbox jobs must not reference signer publish workflows"
+ "logout acknowledgement delivery outbox jobs require a connection id"
.to_owned(),
),
));
@@ -688,7 +739,8 @@ impl MycRuntime {
MycDeliveryOutboxKind::AuthReplayPublish => {
RadrootsNostrSignerPublishWorkflowKind::AuthReplayFinalization
}
- MycDeliveryOutboxKind::DiscoveryHandlerPublish => unreachable!(),
+ MycDeliveryOutboxKind::DiscoveryHandlerPublish
+ | MycDeliveryOutboxKind::LogoutAcknowledgementPublish => unreachable!(),
};
if workflow.kind != kind_label {
return Err(self.wrap_recovery_error(
@@ -892,6 +944,9 @@ impl MycRuntime {
fn recovery_operation_label(kind: MycDeliveryOutboxKind) -> &'static str {
match kind {
MycDeliveryOutboxKind::ListenerResponsePublish => "listener response recovery publish",
+ MycDeliveryOutboxKind::LogoutAcknowledgementPublish => {
+ "logout acknowledgement recovery publish"
+ }
MycDeliveryOutboxKind::ConnectAcceptPublish => "connect accept recovery publish",
MycDeliveryOutboxKind::AuthReplayPublish => "auth replay recovery publish",
MycDeliveryOutboxKind::DiscoveryHandlerPublish => "discovery handler recovery publish",
@@ -903,6 +958,9 @@ fn recovery_operation_audit_kind(kind: MycDeliveryOutboxKind) -> MycOperationAud
MycDeliveryOutboxKind::ListenerResponsePublish => {
MycOperationAuditKind::ListenerResponsePublish
}
+ MycDeliveryOutboxKind::LogoutAcknowledgementPublish => {
+ MycOperationAuditKind::ListenerResponsePublish
+ }
MycDeliveryOutboxKind::ConnectAcceptPublish => MycOperationAuditKind::ConnectAcceptPublish,
MycDeliveryOutboxKind::AuthReplayPublish => MycOperationAuditKind::AuthReplayPublish,
MycDeliveryOutboxKind::DiscoveryHandlerPublish => {
diff --git a/src/operability/mod.rs b/src/operability/mod.rs
@@ -1204,7 +1204,8 @@ fn is_delivery_outbox_unfinished(record: &MycDeliveryOutboxRecord) -> bool {
matches!(
record.status,
MycDeliveryOutboxStatus::Queued | MycDeliveryOutboxStatus::PublishedPendingFinalize
- )
+ ) || (record.status == MycDeliveryOutboxStatus::Failed
+ && record.kind == crate::outbox::MycDeliveryOutboxKind::LogoutAcknowledgementPublish)
}
fn is_critical_delivery_outbox_job(record: &MycDeliveryOutboxRecord) -> bool {
@@ -1239,6 +1240,11 @@ fn classify_blocked_delivery_outbox_record(
}
}
crate::outbox::MycDeliveryOutboxKind::ListenerResponsePublish => {}
+ crate::outbox::MycDeliveryOutboxKind::LogoutAcknowledgementPublish => {
+ if record.signer_publish_workflow_id.is_some() || record.connection_id.is_none() {
+ return Some(true);
+ }
+ }
}
if let Some(workflow_id) = record.signer_publish_workflow_id.as_ref() {
diff --git a/src/outbox.rs b/src/outbox.rs
@@ -18,6 +18,7 @@ pub struct MycDeliveryOutboxJobId(String);
#[serde(rename_all = "snake_case")]
pub enum MycDeliveryOutboxKind {
ListenerResponsePublish,
+ LogoutAcknowledgementPublish,
ConnectAcceptPublish,
AuthReplayPublish,
DiscoveryHandlerPublish,
diff --git a/src/outbox_sqlite.rs b/src/outbox_sqlite.rs
@@ -442,6 +442,7 @@ fn parse_json_field<T: DeserializeOwned>(
fn kind_label(kind: MycDeliveryOutboxKind) -> &'static str {
match kind {
MycDeliveryOutboxKind::ListenerResponsePublish => "listener_response_publish",
+ MycDeliveryOutboxKind::LogoutAcknowledgementPublish => "logout_acknowledgement_publish",
MycDeliveryOutboxKind::ConnectAcceptPublish => "connect_accept_publish",
MycDeliveryOutboxKind::AuthReplayPublish => "auth_replay_publish",
MycDeliveryOutboxKind::DiscoveryHandlerPublish => "discovery_handler_publish",
@@ -451,6 +452,7 @@ fn kind_label(kind: MycDeliveryOutboxKind) -> &'static str {
fn parse_kind(value: &str) -> Result<MycDeliveryOutboxKind, MycError> {
match value {
"listener_response_publish" => Ok(MycDeliveryOutboxKind::ListenerResponsePublish),
+ "logout_acknowledgement_publish" => Ok(MycDeliveryOutboxKind::LogoutAcknowledgementPublish),
"connect_accept_publish" => Ok(MycDeliveryOutboxKind::ConnectAcceptPublish),
"auth_replay_publish" => Ok(MycDeliveryOutboxKind::AuthReplayPublish),
"discovery_handler_publish" => Ok(MycDeliveryOutboxKind::DiscoveryHandlerPublish),
diff --git a/src/persistence.rs b/src/persistence.rs
@@ -1218,10 +1218,13 @@ fn verify_restored_delivery_state(
for record in outbox_records {
verify_discovery_restore_author(record, signer_public_key, discovery_app_public_key)?;
+ let recoverable_logout_failure = record.status == MycDeliveryOutboxStatus::Failed
+ && record.kind == MycDeliveryOutboxKind::LogoutAcknowledgementPublish;
if !matches!(
record.status,
MycDeliveryOutboxStatus::Queued | MycDeliveryOutboxStatus::PublishedPendingFinalize
- ) {
+ ) && !recoverable_logout_failure
+ {
continue;
}
@@ -1293,6 +1296,26 @@ fn verify_restore_outbox_record<'a>(
)));
}
}
+ MycDeliveryOutboxKind::LogoutAcknowledgementPublish => {
+ if record.signer_publish_workflow_id.is_some() {
+ return Err(MycError::InvalidOperation(format!(
+ "persistence verify-restore found logout acknowledgement delivery outbox job `{}` that incorrectly references a signer publish workflow",
+ record.job_id
+ )));
+ }
+ let connection_id = record.connection_id.as_ref().ok_or_else(|| {
+ MycError::InvalidOperation(format!(
+ "persistence verify-restore found logout acknowledgement delivery outbox job `{}` without a connection id",
+ record.job_id
+ ))
+ })?;
+ if !connections_by_id.contains_key(connection_id.as_str()) {
+ return Err(MycError::InvalidOperation(format!(
+ "persistence verify-restore found logout acknowledgement delivery outbox job `{}` referencing missing connection `{connection_id}`",
+ record.job_id
+ )));
+ }
+ }
MycDeliveryOutboxKind::ConnectAcceptPublish | MycDeliveryOutboxKind::AuthReplayPublish => {
if record.signer_publish_workflow_id.is_none() {
return Err(MycError::InvalidOperation(format!(
@@ -1314,7 +1337,8 @@ fn verify_restore_outbox_record<'a>(
MycDeliveryOutboxKind::AuthReplayPublish => {
RadrootsNostrSignerPublishWorkflowKind::AuthReplayFinalization
}
- MycDeliveryOutboxKind::DiscoveryHandlerPublish => unreachable!(),
+ MycDeliveryOutboxKind::DiscoveryHandlerPublish
+ | MycDeliveryOutboxKind::LogoutAcknowledgementPublish => unreachable!(),
};
if workflow.kind != expected_kind {
return Err(MycError::InvalidOperation(format!(
@@ -1414,6 +1438,12 @@ fn verify_already_finalized_without_workflow(
record.job_id
)));
}
+ MycDeliveryOutboxKind::LogoutAcknowledgementPublish => {
+ return Err(MycError::InvalidOperation(format!(
+ "persistence verify-restore found logout acknowledgement delivery outbox job `{}` unexpectedly referencing signer workflow `{workflow_id}`",
+ record.job_id
+ )));
+ }
}
Ok(())
diff --git a/src/transport/nip46.rs b/src/transport/nip46.rs
@@ -435,6 +435,11 @@ impl MycNip46Service {
}
let outbox_record = match self.build_listener_outbox_record(
+ if revoke_logout_connection.is_some() {
+ MycDeliveryOutboxKind::LogoutAcknowledgementPublish
+ } else {
+ MycDeliveryOutboxKind::ListenerResponsePublish
+ },
response_event.clone(),
connection_id.as_ref(),
request_id.as_str(),
@@ -590,17 +595,15 @@ impl MycNip46Service {
fn build_listener_outbox_record(
&self,
+ kind: MycDeliveryOutboxKind,
response_event: RadrootsNostrEvent,
connection_id: Option<&RadrootsNostrSignerConnectionId>,
request_id: &str,
workflow_id: Option<&RadrootsNostrSignerWorkflowId>,
) -> Result<MycDeliveryOutboxRecord, MycError> {
- let mut record = MycDeliveryOutboxRecord::new(
- MycDeliveryOutboxKind::ListenerResponsePublish,
- response_event,
- self.transport.relays().to_vec(),
- )?
- .with_request_id(request_id.to_owned());
+ let mut record =
+ MycDeliveryOutboxRecord::new(kind, response_event, self.transport.relays().to_vec())?
+ .with_request_id(request_id.to_owned());
if let Some(connection_id) = connection_id {
record = record.with_connection_id(connection_id);
}
diff --git a/tests/nip46_e2e.rs b/tests/nip46_e2e.rs
@@ -1566,6 +1566,19 @@ async fn live_listener_acknowledges_logout_before_revoking_session() -> TestResu
RadrootsNostrSignerConnectionStatus::Revoked,
)
.await?;
+ let logout_outbox = wait_for_delivery_outbox_records(&runtime, |records| {
+ records.iter().any(|record| {
+ record.kind == MycDeliveryOutboxKind::LogoutAcknowledgementPublish
+ && record.status == MycDeliveryOutboxStatus::Finalized
+ })
+ })
+ .await?;
+ let logout_outbox = logout_outbox
+ .iter()
+ .find(|record| record.kind == MycDeliveryOutboxKind::LogoutAcknowledgementPublish)
+ .expect("logout acknowledgement outbox record");
+ assert_eq!(logout_outbox.request_id.as_deref(), Some("logout-request"));
+ assert!(logout_outbox.signer_publish_workflow_id.is_none());
let repeated_logout = build_request_event(
&client_identity,
@@ -1654,6 +1667,200 @@ async fn live_listener_acknowledges_logout_before_revoking_session() -> TestResu
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
+async fn failed_logout_publish_is_retried_and_revoked_during_startup_recovery() -> TestResult<()> {
+ let relay = TestRelay::spawn().await?;
+ let test_runtime = MycTestRuntime::new(relay.url(), MycConnectionApproval::NotRequired);
+ let MycTestRuntime {
+ _temp: _tempdir,
+ runtime,
+ } = test_runtime;
+ let config = runtime.config().clone();
+ let signer_public_key = runtime.signer_identity().public_key();
+ let client_identity =
+ identity("3737373737373737373737373737373737373737373737373737373737373737");
+ let base_created_at = Timestamp::now().as_secs();
+ relay
+ .queue_publish_outcomes(signer_public_key, &[true, false, true])
+ .await;
+
+ let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
+ let service_runtime = runtime.clone();
+ let listener_task = tokio::spawn(async move {
+ service_runtime
+ .run_until(async {
+ let _ = shutdown_rx.await;
+ })
+ .await
+ });
+ relay.wait_for_subscription_count(1).await?;
+
+ let connect = build_request_event(
+ &client_identity,
+ signer_public_key,
+ connect_request_message("recovery-connect", signer_public_key, "recovery-secret"),
+ base_created_at,
+ );
+ publish_event(relay.url(), &connect).await?;
+ relay
+ .wait_for_published_events_by_author(signer_public_key, 1)
+ .await?;
+
+ let logout = build_request_event(
+ &client_identity,
+ signer_public_key,
+ RadrootsNostrConnectRequestMessage::new(
+ "recovery-logout",
+ RadrootsNostrConnectRequest::Logout,
+ ),
+ base_created_at + 1,
+ );
+ publish_event(relay.url(), &logout).await?;
+ let failed_outbox = wait_for_delivery_outbox_records(&runtime, |records| {
+ records.iter().any(|record| {
+ record.kind == MycDeliveryOutboxKind::LogoutAcknowledgementPublish
+ && record.status == MycDeliveryOutboxStatus::Failed
+ })
+ })
+ .await?;
+ assert_eq!(
+ failed_outbox
+ .iter()
+ .find(|record| record.kind == MycDeliveryOutboxKind::LogoutAcknowledgementPublish)
+ .and_then(|record| record.request_id.as_deref()),
+ Some("recovery-logout")
+ );
+ assert_eq!(
+ runtime
+ .signer_manager()?
+ .find_connections_by_client_public_key(&client_identity.public_key())?[0]
+ .status,
+ RadrootsNostrSignerConnectionStatus::Active
+ );
+ assert_eq!(
+ relay
+ .published_events_by_author(signer_public_key)
+ .await
+ .len(),
+ 1
+ );
+
+ let _ = shutdown_tx.send(());
+ listener_task.await??;
+
+ let restarted_runtime = MycRuntime::bootstrap(config)?;
+ restarted_runtime.clone().run_until(async {}).await?;
+ let responses = relay
+ .wait_for_published_events_by_author(signer_public_key, 2)
+ .await?;
+ let recovered_response = decrypt_response(&client_identity, signer_public_key, &responses[1]);
+ let recovered_response = RadrootsNostrConnectResponse::from_envelope(
+ &RadrootsNostrConnectRequest::Logout.method(),
+ recovered_response,
+ )?;
+ assert_eq!(
+ recovered_response,
+ RadrootsNostrConnectResponse::LogoutAcknowledged
+ );
+ let recovered_connection = restarted_runtime
+ .signer_manager()?
+ .find_connections_by_client_public_key(&client_identity.public_key())?
+ .into_iter()
+ .next()
+ .expect("recovered connection");
+ assert_eq!(
+ recovered_connection.status,
+ RadrootsNostrSignerConnectionStatus::Revoked
+ );
+ let recovered_outbox = restarted_runtime.delivery_outbox_store().list_all()?;
+ assert!(recovered_outbox.iter().any(|record| {
+ record.kind == MycDeliveryOutboxKind::LogoutAcknowledgementPublish
+ && record.status == MycDeliveryOutboxStatus::Finalized
+ }));
+ Ok(())
+}
+
+#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
+async fn published_logout_acknowledgement_is_finalized_without_republish_on_restart()
+-> TestResult<()> {
+ let relay = TestRelay::spawn().await?;
+ let test_runtime = MycTestRuntime::new(relay.url(), MycConnectionApproval::NotRequired);
+ let MycTestRuntime {
+ _temp: _tempdir,
+ runtime,
+ } = test_runtime;
+ let config = runtime.config().clone();
+ let signer_public_key = runtime.signer_identity().public_key();
+ let client_identity =
+ identity("3838383838383838383838383838383838383838383838383838383838383838");
+ let relay_url: RadrootsNostrRelayUrl = relay.url().parse()?;
+ let connection = runtime.signer_manager()?.register_connection(
+ RadrootsNostrSignerConnectionDraft::new(
+ client_identity.public_key(),
+ runtime.user_public_identity(),
+ )
+ .with_relays(vec![relay_url.clone()])
+ .with_approval_requirement(RadrootsNostrSignerApprovalRequirement::NotRequired),
+ )?;
+ let acknowledgement_event = runtime
+ .signer_identity()
+ .sign_event_builder(
+ RadrootsNostrEventBuilder::new(
+ RadrootsNostrKind::Custom(RADROOTS_NOSTR_CONNECT_RPC_KIND),
+ "published logout acknowledgement fixture",
+ ),
+ "published logout acknowledgement fixture",
+ )
+ .map_err(|error| format!("failed to sign logout acknowledgement fixture: {error}"))?;
+ publish_event(relay.url(), &acknowledgement_event).await?;
+ let outbox_record = MycDeliveryOutboxRecord::new(
+ MycDeliveryOutboxKind::LogoutAcknowledgementPublish,
+ acknowledgement_event,
+ vec![relay_url],
+ )?
+ .with_connection_id(&connection.connection_id)
+ .with_request_id("published-logout-ack");
+ runtime.delivery_outbox_store().enqueue(&outbox_record)?;
+ runtime
+ .delivery_outbox_store()
+ .mark_published_pending_finalize(&outbox_record.job_id, 1)?;
+ assert_eq!(
+ runtime
+ .signer_manager()?
+ .get_connection(&connection.connection_id)?
+ .expect("active connection")
+ .status,
+ RadrootsNostrSignerConnectionStatus::Active
+ );
+
+ let restarted_runtime = MycRuntime::bootstrap(config)?;
+ restarted_runtime.clone().run_until(async {}).await?;
+ let recovered_connection = restarted_runtime
+ .signer_manager()?
+ .get_connection(&connection.connection_id)?
+ .expect("recovered connection");
+ assert_eq!(
+ recovered_connection.status,
+ RadrootsNostrSignerConnectionStatus::Revoked
+ );
+ assert_eq!(
+ restarted_runtime
+ .delivery_outbox_store()
+ .get(&outbox_record.job_id)?
+ .expect("recovered outbox")
+ .status,
+ MycDeliveryOutboxStatus::Finalized
+ );
+ assert_eq!(
+ relay
+ .published_events_by_author(signer_public_key)
+ .await
+ .len(),
+ 1
+ );
+ Ok(())
+}
+
+#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn live_listener_works_with_sqlite_signer_state_and_runtime_audit() -> TestResult<()> {
let relay = TestRelay::spawn().await?;
let test_runtime = MycTestRuntime::new_with_transport_config(