From b9dffff661f8f85a504fdf7d52bcdaabf395aa4d Mon Sep 17 00:00:00 2001 From: AptS-1547 Date: Thu, 20 Aug 2026 21:58:13 +0800 Subject: [PATCH 1/3] feat(remote): audit tunnel connection lifecycle --- CHANGELOG.md | 4 + crates/aster_drive_model/src/types/audit.rs | 28 ++ .../src/i18n/locales/en/admin/audit.json | 5 + .../src/i18n/locales/zh/admin/audit.json | 5 + frontend-panel/src/services/api.generated.ts | 2 +- src/runtime/startup/primary.rs | 1 + src/services/ops/audit/details.rs | 20 + src/services/ops/audit/mod.rs | 21 +- src/services/ops/audit/presentation.rs | 46 ++ src/services/ops/audit/tests.rs | 50 ++ src/storage/remote_protocol/runtime.rs | 5 + .../remote_protocol/tunnel/server/mod.rs | 65 ++- .../tunnel/server/registry/mod.rs | 464 ++++++++++++++++++ .../tunnel/server/registry/streaming.rs | 15 +- 14 files changed, 703 insertions(+), 28 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0d1ca5b28..7185ecac9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added + +- **远端节点连接生命周期审计** — reverse tunnel 连接、正常下线、异常断线和心跳超时现在会按 remote node / binding 聚合写入系统 audit;四条 streaming lane 的同时变化只产生一次节点级状态转换,并记录连接代际、故障代际、lane 数量、transport 和稳定 reason code,不包含 access key、secret、signature、URL 凭据或 token。 + ## [v0.5.0] - 2026-08-20 ### Changed diff --git a/crates/aster_drive_model/src/types/audit.rs b/crates/aster_drive_model/src/types/audit.rs index 3f182ebd8..f243f5633 100644 --- a/crates/aster_drive_model/src/types/audit.rs +++ b/crates/aster_drive_model/src/types/audit.rs @@ -169,6 +169,10 @@ macro_rules! define_audit_action_list { TagDelete, TagAttach, TagDetach, + RemoteNodeConnected, + RemoteNodeGracefulDisconnect, + RemoteNodeUnexpectedDisconnect, + RemoteNodeHeartbeatTimeout, } }; } @@ -689,6 +693,14 @@ pub enum AuditAction { TagAttach, #[sea_orm(string_value = "tag_detach")] TagDetach, + #[sea_orm(string_value = "remote_node_connected")] + RemoteNodeConnected, + #[sea_orm(string_value = "remote_node_graceful_disconnect")] + RemoteNodeGracefulDisconnect, + #[sea_orm(string_value = "remote_node_unexpected_disconnect")] + RemoteNodeUnexpectedDisconnect, + #[sea_orm(string_value = "remote_node_heartbeat_timeout")] + RemoteNodeHeartbeatTimeout, } impl AuditAction { @@ -846,6 +858,10 @@ impl AuditAction { Self::TagDelete => 146, Self::TagAttach => 147, Self::TagDetach => 148, + Self::RemoteNodeConnected => 149, + Self::RemoteNodeGracefulDisconnect => 150, + Self::RemoteNodeUnexpectedDisconnect => 151, + Self::RemoteNodeHeartbeatTimeout => 152, } } @@ -1002,6 +1018,10 @@ impl AuditAction { Self::TagDelete => "tag_delete", Self::TagAttach => "tag_attach", Self::TagDetach => "tag_detach", + Self::RemoteNodeConnected => "remote_node_connected", + Self::RemoteNodeGracefulDisconnect => "remote_node_graceful_disconnect", + Self::RemoteNodeUnexpectedDisconnect => "remote_node_unexpected_disconnect", + Self::RemoteNodeHeartbeatTimeout => "remote_node_heartbeat_timeout", } } @@ -1158,6 +1178,10 @@ impl AuditAction { "tag_delete" => Some(Self::TagDelete), "tag_attach" => Some(Self::TagAttach), "tag_detach" => Some(Self::TagDetach), + "remote_node_connected" => Some(Self::RemoteNodeConnected), + "remote_node_graceful_disconnect" => Some(Self::RemoteNodeGracefulDisconnect), + "remote_node_unexpected_disconnect" => Some(Self::RemoteNodeUnexpectedDisconnect), + "remote_node_heartbeat_timeout" => Some(Self::RemoteNodeHeartbeatTimeout), _ => None, } } @@ -1303,6 +1327,10 @@ impl AuditAction { | Self::TagDelete | Self::TagAttach | Self::TagDetach => "tag", + Self::RemoteNodeConnected + | Self::RemoteNodeGracefulDisconnect + | Self::RemoteNodeUnexpectedDisconnect + | Self::RemoteNodeHeartbeatTimeout => "remote", } } } diff --git a/frontend-panel/src/i18n/locales/en/admin/audit.json b/frontend-panel/src/i18n/locales/en/admin/audit.json index c7b411cfe..30228d2a5 100644 --- a/frontend-panel/src/i18n/locales/en/admin/audit.json +++ b/frontend-panel/src/i18n/locales/en/admin/audit.json @@ -112,6 +112,10 @@ "audit_action_trash_purge_all": "Emptied trash", "audit_action_remote_enrollment_redeem": "Redeemed remote enrollment", "audit_action_remote_enrollment_ack": "Acknowledged remote enrollment", + "audit_action_remote_node_connected": "Remote node connected", + "audit_action_remote_node_graceful_disconnect": "Remote node disconnected gracefully", + "audit_action_remote_node_unexpected_disconnect": "Remote node disconnected unexpectedly", + "audit_action_remote_node_heartbeat_timeout": "Remote node heartbeat timed out", "audit_action_user_revoke_other_sessions": "Revoked other sessions", "audit_action_user_revoke_session": "Revoked session", "audit_action_user_update_preferences": "Updated user preferences", @@ -225,6 +229,7 @@ "audit_presentation_mfa_email_code_sent": "{{method}} flow {{flow_id}}, expires in {{expires_in}}s", "audit_presentation_wopi_user_info_updated": "{{app_key}} for file {{file_id}}, {{user_info_len}} bytes", "audit_presentation_remote_enrollment_changed": "{{phase}} {{remote_node_name}}, enabled {{is_enabled}}", + "audit_presentation_remote_node_connection_lifecycle": "{{reason}} on {{transport}} (generation {{generation}}, outage {{outage_generation}}, {{active_lanes}} active / {{lane_count}} lane(s))", "audit_presentation_invitation_snapshot": "{{email}} status {{status}}, expires {{expires_at}}", "audit_presentation_external_auth_unlinked": "{{provider_key}} identity {{subject}}", "audit_presentation_follower_binding_synced": "{{name}} enabled {{is_enabled}}", diff --git a/frontend-panel/src/i18n/locales/zh/admin/audit.json b/frontend-panel/src/i18n/locales/zh/admin/audit.json index a65534fc6..3803e85f9 100644 --- a/frontend-panel/src/i18n/locales/zh/admin/audit.json +++ b/frontend-panel/src/i18n/locales/zh/admin/audit.json @@ -112,6 +112,10 @@ "audit_action_trash_purge_all": "清空回收站", "audit_action_remote_enrollment_redeem": "兑换远程节点注册", "audit_action_remote_enrollment_ack": "确认远程节点注册", + "audit_action_remote_node_connected": "远程节点已连接", + "audit_action_remote_node_graceful_disconnect": "远程节点正常下线", + "audit_action_remote_node_unexpected_disconnect": "远程节点异常断线", + "audit_action_remote_node_heartbeat_timeout": "远程节点心跳超时", "audit_action_user_revoke_other_sessions": "撤销其他登录会话", "audit_action_user_revoke_session": "撤销登录会话", "audit_action_user_update_preferences": "更新用户偏好", @@ -225,6 +229,7 @@ "audit_presentation_mfa_email_code_sent": "{{method}} flow {{flow_id}},{{expires_in}} 秒后过期", "audit_presentation_wopi_user_info_updated": "{{app_key}} 更新文件 {{file_id}} 的用户信息,{{user_info_len}} 字节", "audit_presentation_remote_enrollment_changed": "{{phase}} {{remote_node_name}},启用 {{is_enabled}}", + "audit_presentation_remote_node_connection_lifecycle": "{{transport}} 上发生 {{reason}}(连接代际 {{generation}},故障代际 {{outage_generation}},{{active_lanes}} 条活跃 / 共 {{lane_count}} 条 lane)", "audit_presentation_invitation_snapshot": "{{email}} 状态 {{status}},过期 {{expires_at}}", "audit_presentation_external_auth_unlinked": "{{provider_key}} 身份 {{subject}}", "audit_presentation_follower_binding_synced": "{{name}} 启用 {{is_enabled}}", diff --git a/frontend-panel/src/services/api.generated.ts b/frontend-panel/src/services/api.generated.ts index 80adda448..704452be9 100644 --- a/frontend-panel/src/services/api.generated.ts +++ b/frontend-panel/src/services/api.generated.ts @@ -4863,7 +4863,7 @@ export interface components { * @description 审计日志动作 * @enum {string} */ - AuditAction: "admin_create_user" | "admin_force_delete_user" | "admin_create_team" | "admin_create_policy_group" | "admin_archive_team" | "admin_restore_team" | "admin_revoke_user_sessions" | "admin_reset_user_password" | "admin_reset_user_mfa" | "admin_update_team" | "admin_update_user" | "admin_delete_policy_group" | "admin_migrate_policy_group_users" | "admin_update_policy_group" | "admin_create_policy" | "admin_update_policy" | "admin_delete_policy" | "admin_trigger_storage_action" | "admin_delete_config" | "admin_delete_share" | "admin_force_unlock" | "admin_cleanup_expired_locks" | "admin_cleanup_tasks" | "admin_create_blob_maintenance_task" | "admin_create_remote_node" | "admin_update_remote_node" | "admin_delete_remote_node" | "admin_test_remote_node" | "admin_create_remote_node_enrollment_token" | "admin_create_remote_ingress_profile" | "admin_update_remote_ingress_profile" | "admin_delete_remote_ingress_profile" | "admin_create_external_auth_provider" | "admin_update_external_auth_provider" | "admin_delete_external_auth_provider" | "admin_test_external_auth_provider" | "batch_copy" | "batch_delete" | "batch_move" | "config_action_execute" | "config_update" | "file_copy" | "file_create" | "file_delete" | "file_download" | "file_direct_link_create" | "file_edit" | "file_move" | "file_rename" | "file_upload" | "file_preview_link_create" | "file_wopi_open" | "file_upload_cancel" | "file_restore" | "file_purge" | "file_lock" | "file_unlock" | "file_version_restore" | "file_version_delete" | "folder_copy" | "folder_create" | "folder_delete" | "folder_move" | "folder_policy_change" | "folder_rename" | "folder_restore" | "folder_purge" | "folder_lock" | "folder_unlock" | "property_set" | "property_delete" | "share_batch_delete" | "share_create" | "share_delete" | "share_update" | "system_setup" | "server_start" | "server_shutdown" | "team_archive" | "team_cleanup_expired" | "team_create" | "team_member_add" | "team_member_remove" | "team_member_update" | "team_restore" | "team_update" | "task_retry" | "archive_compress" | "archive_extract" | "archive_download" | "offline_download" | "trash_purge_all" | "remote_enrollment_redeem" | "remote_enrollment_ack" | "user_revoke_other_sessions" | "user_revoke_session" | "user_update_preferences" | "user_update_profile" | "user_upload_avatar" | "user_set_avatar_source" | "user_update_wopi_info" | "webdav_account_create" | "webdav_account_delete" | "webdav_account_toggle" | "team_webdav_account_create" | "team_webdav_account_delete" | "team_webdav_account_toggle" | "user_change_password" | "user_confirm_password_reset" | "user_confirm_email_change" | "user_confirm_registration" | "user_login" | "user_logout" | "user_mfa_enable" | "user_mfa_disable" | "user_mfa_recovery_codes_regenerate" | "user_mfa_email_code_send" | "user_mfa_challenge_success" | "user_mfa_challenge_failed" | "user_passkey_delete" | "user_passkey_login" | "user_passkey_register" | "user_passkey_rename" | "user_external_auth_login" | "user_external_auth_link" | "user_external_auth_unlink" | "user_refresh_token_reuse_detected" | "user_request_email_change" | "user_request_password_reset" | "user_register" | "user_resend_email_change" | "user_resend_registration" | "follower_binding_sync" | "follower_object_read" | "follower_object_write" | "follower_object_delete" | "follower_object_compose" | "follower_ingress_profile_create" | "follower_ingress_profile_update" | "follower_ingress_profile_delete" | "mail_send" | "mail_delivery_failed" | "admin_create_invitation" | "admin_revoke_invitation" | "tag_create" | "tag_update" | "tag_delete" | "tag_attach" | "tag_detach"; + AuditAction: "admin_create_user" | "admin_force_delete_user" | "admin_create_team" | "admin_create_policy_group" | "admin_archive_team" | "admin_restore_team" | "admin_revoke_user_sessions" | "admin_reset_user_password" | "admin_reset_user_mfa" | "admin_update_team" | "admin_update_user" | "admin_delete_policy_group" | "admin_migrate_policy_group_users" | "admin_update_policy_group" | "admin_create_policy" | "admin_update_policy" | "admin_delete_policy" | "admin_trigger_storage_action" | "admin_delete_config" | "admin_delete_share" | "admin_force_unlock" | "admin_cleanup_expired_locks" | "admin_cleanup_tasks" | "admin_create_blob_maintenance_task" | "admin_create_remote_node" | "admin_update_remote_node" | "admin_delete_remote_node" | "admin_test_remote_node" | "admin_create_remote_node_enrollment_token" | "admin_create_remote_ingress_profile" | "admin_update_remote_ingress_profile" | "admin_delete_remote_ingress_profile" | "admin_create_external_auth_provider" | "admin_update_external_auth_provider" | "admin_delete_external_auth_provider" | "admin_test_external_auth_provider" | "batch_copy" | "batch_delete" | "batch_move" | "config_action_execute" | "config_update" | "file_copy" | "file_create" | "file_delete" | "file_download" | "file_direct_link_create" | "file_edit" | "file_move" | "file_rename" | "file_upload" | "file_preview_link_create" | "file_wopi_open" | "file_upload_cancel" | "file_restore" | "file_purge" | "file_lock" | "file_unlock" | "file_version_restore" | "file_version_delete" | "folder_copy" | "folder_create" | "folder_delete" | "folder_move" | "folder_policy_change" | "folder_rename" | "folder_restore" | "folder_purge" | "folder_lock" | "folder_unlock" | "property_set" | "property_delete" | "share_batch_delete" | "share_create" | "share_delete" | "share_update" | "system_setup" | "server_start" | "server_shutdown" | "team_archive" | "team_cleanup_expired" | "team_create" | "team_member_add" | "team_member_remove" | "team_member_update" | "team_restore" | "team_update" | "task_retry" | "archive_compress" | "archive_extract" | "archive_download" | "offline_download" | "trash_purge_all" | "remote_enrollment_redeem" | "remote_enrollment_ack" | "user_revoke_other_sessions" | "user_revoke_session" | "user_update_preferences" | "user_update_profile" | "user_upload_avatar" | "user_set_avatar_source" | "user_update_wopi_info" | "webdav_account_create" | "webdav_account_delete" | "webdav_account_toggle" | "team_webdav_account_create" | "team_webdav_account_delete" | "team_webdav_account_toggle" | "user_change_password" | "user_confirm_password_reset" | "user_confirm_email_change" | "user_confirm_registration" | "user_login" | "user_logout" | "user_mfa_enable" | "user_mfa_disable" | "user_mfa_recovery_codes_regenerate" | "user_mfa_email_code_send" | "user_mfa_challenge_success" | "user_mfa_challenge_failed" | "user_passkey_delete" | "user_passkey_login" | "user_passkey_register" | "user_passkey_rename" | "user_external_auth_login" | "user_external_auth_link" | "user_external_auth_unlink" | "user_refresh_token_reuse_detected" | "user_request_email_change" | "user_request_password_reset" | "user_register" | "user_resend_email_change" | "user_resend_registration" | "follower_binding_sync" | "follower_object_read" | "follower_object_write" | "follower_object_delete" | "follower_object_compose" | "follower_ingress_profile_create" | "follower_ingress_profile_update" | "follower_ingress_profile_delete" | "mail_send" | "mail_delivery_failed" | "admin_create_invitation" | "admin_revoke_invitation" | "tag_create" | "tag_update" | "tag_delete" | "tag_attach" | "tag_detach" | "remote_node_connected" | "remote_node_graceful_disconnect" | "remote_node_unexpected_disconnect" | "remote_node_heartbeat_timeout"; /** * @description 审计日志实体类型 * @enum {string} diff --git a/src/runtime/startup/primary.rs b/src/runtime/startup/primary.rs index f1b9e4794..318240a55 100644 --- a/src/runtime/startup/primary.rs +++ b/src/runtime/startup/primary.rs @@ -34,6 +34,7 @@ pub async fn prepare_primary() -> Result { let remote_protocol = crate::runtime::PrimaryAppState::new_remote_protocol(); remote_protocol.set_persistence_db(common.database.clone()); + remote_protocol.set_audit_runtime_config(runtime_config.clone()); remote_protocol.configure_tunnel_owner_directory( common.database.clone(), &common.cfg.deployment, diff --git a/src/services/ops/audit/details.rs b/src/services/ops/audit/details.rs index bd04b0f26..2e7c958bd 100644 --- a/src/services/ops/audit/details.rs +++ b/src/services/ops/audit/details.rs @@ -456,6 +456,26 @@ pub struct FollowerBindingAuditDetails<'a> { pub is_enabled: bool, } +/// Stable, redacted context for a remote node reverse-tunnel lifecycle event. +/// +/// The tunnel has multiple polling/streaming lanes, but the audit contract is +/// node/binding scoped. Generation counters let operators correlate a recovery +/// with the outage it follows without persisting credentials or URLs. +#[derive(Serialize)] +pub struct RemoteNodeConnectionAuditDetails<'a> { + pub remote_node_id: i64, + pub binding_id: i64, + pub transport: &'a str, + pub reason: &'a str, + pub generation: u64, + pub outage_generation: u64, + pub active_lanes: usize, + pub lane_count: usize, + pub observed_at: DateTime, + #[serde(skip_serializing_if = "Option::is_none")] + pub first_lane_id: Option<&'a str>, +} + #[derive(Serialize)] pub struct FollowerObjectAuditDetails<'a> { pub binding_id: i64, diff --git a/src/services/ops/audit/mod.rs b/src/services/ops/audit/mod.rs index 4d8af8bdc..007b0278f 100644 --- a/src/services/ops/audit/mod.rs +++ b/src/services/ops/audit/mod.rs @@ -23,16 +23,17 @@ pub use details::{ MfaEmailCodeAuditDetails, PolicyGroupAuditDetails, PolicyGroupMigrationDetails, PropertyAuditDetails, RemoteEnrollmentAuditDetails, RemoteIngressProfileAuditDetails, RemoteIngressProfileDeleteAuditDetails, RemoteNodeAuditDetails, - RemoteNodeEnrollmentTokenAuditDetails, RemoteNodeParamTestAuditDetails, - ShareBatchDeleteDetails, ShareCreateAuditDetails, ShareDeleteAuditDetails, ShareUpdateDetails, - StoragePolicyActionAuditDetails, StoragePolicyAuditDetails, StoragePolicyPromotionAuditDetails, - TagAssignmentAuditDetails, TagAuditDetails, TaskRetryAuditDetails, TeamAuditDetails, - TeamCleanupAuditDetails, TeamMemberAddAuditDetails, TeamMemberRemoveAuditDetails, - TeamMemberUpdateAuditDetails, TrashPurgeAllAuditDetails, UploadCancelAuditDetails, - UserAvatarSourceAuditDetails, UserAvatarUploadAuditDetails, UserLoginAuditDetails, - UserMfaManageAuditDetails, UserPreferencesAuditDetails, UserProfileAuditDetails, - UserWopiInfoAuditDetails, WorkspaceTransferCopyDetails, WorkspaceTransferMoveDetails, - WorkspaceTransferScopeDetails, details, + RemoteNodeConnectionAuditDetails, RemoteNodeEnrollmentTokenAuditDetails, + RemoteNodeParamTestAuditDetails, ShareBatchDeleteDetails, ShareCreateAuditDetails, + ShareDeleteAuditDetails, ShareUpdateDetails, StoragePolicyActionAuditDetails, + StoragePolicyAuditDetails, StoragePolicyPromotionAuditDetails, TagAssignmentAuditDetails, + TagAuditDetails, TaskRetryAuditDetails, TeamAuditDetails, TeamCleanupAuditDetails, + TeamMemberAddAuditDetails, TeamMemberRemoveAuditDetails, TeamMemberUpdateAuditDetails, + TrashPurgeAllAuditDetails, UploadCancelAuditDetails, UserAvatarSourceAuditDetails, + UserAvatarUploadAuditDetails, UserLoginAuditDetails, UserMfaManageAuditDetails, + UserPreferencesAuditDetails, UserProfileAuditDetails, UserWopiInfoAuditDetails, + WorkspaceTransferCopyDetails, WorkspaceTransferMoveDetails, WorkspaceTransferScopeDetails, + details, }; pub use filters::{AuditLogFilterQuery, AuditLogFilters}; pub use manager::{ diff --git a/src/services/ops/audit/presentation.rs b/src/services/ops/audit/presentation.rs index 5fa3f5837..00e1c05e6 100644 --- a/src/services/ops/audit/presentation.rs +++ b/src/services/ops/audit/presentation.rs @@ -764,6 +764,28 @@ fn detail_message( ); Some(message("remote_enrollment_changed", params)) } + AuditAction::RemoteNodeConnected + | AuditAction::RemoteNodeGracefulDisconnect + | AuditAction::RemoteNodeUnexpectedDisconnect + | AuditAction::RemoteNodeHeartbeatTimeout => { + copy_params( + details, + &mut params, + &[ + "remote_node_id", + "binding_id", + "transport", + "reason", + "generation", + "outage_generation", + "active_lanes", + "lane_count", + "observed_at", + "first_lane_id", + ], + ); + Some(message("remote_node_connection_lifecycle", params)) + } AuditAction::AdminCreateInvitation | AuditAction::AdminRevokeInvitation => { copy_params( details, @@ -1339,6 +1361,30 @@ mod tests { ); } + #[test] + fn presentation_includes_remote_node_connection_lifecycle_detail() { + let presentation = build_audit_presentation( + AuditAction::RemoteNodeHeartbeatTimeout, + AuditEntityType::RemoteNode, + Some(42), + Some("edge-a"), + Some( + r#"{"remote_node_id":42,"binding_id":42,"transport":"reverse_tunnel","reason":"heartbeat_timeout","generation":3,"outage_generation":2,"active_lanes":0,"lane_count":4,"observed_at":"2026-08-20T10:00:00Z","first_lane_id":"lane-0"}"#, + ), + ) + .expect("presentation should be built"); + + let detail = presentation.detail.as_ref().unwrap(); + assert_eq!(detail.code, "remote_node_connection_lifecycle"); + assert_eq!( + detail.params.get("reason"), + Some(&Value::String("heartbeat_timeout".to_string())) + ); + assert_eq!(detail.params.get("generation"), Some(&Value::from(3))); + assert_eq!(detail.params.get("active_lanes"), Some(&Value::from(0))); + assert_eq!(detail.params.get("lane_count"), Some(&Value::from(4))); + } + #[test] fn presentation_includes_invitation_detail() { let presentation = build_audit_presentation( diff --git a/src/services/ops/audit/tests.rs b/src/services/ops/audit/tests.rs index 925995ea6..8f7ee3deb 100644 --- a/src/services/ops/audit/tests.rs +++ b/src/services/ops/audit/tests.rs @@ -270,6 +270,19 @@ fn audit_action_strings_match_existing_contract() { "remote_enrollment_redeem", ), (AuditAction::RemoteEnrollmentAck, "remote_enrollment_ack"), + (AuditAction::RemoteNodeConnected, "remote_node_connected"), + ( + AuditAction::RemoteNodeGracefulDisconnect, + "remote_node_graceful_disconnect", + ), + ( + AuditAction::RemoteNodeUnexpectedDisconnect, + "remote_node_unexpected_disconnect", + ), + ( + AuditAction::RemoteNodeHeartbeatTimeout, + "remote_node_heartbeat_timeout", + ), ( AuditAction::UserRevokeOtherSessions, "user_revoke_other_sessions", @@ -436,6 +449,43 @@ async fn log_writes_synchronously_without_global_manager() { .await .expect("audit query should succeed"); assert_eq!(count, 1); + + super::log_with_details( + &state, + &AuditContext::system(), + AuditAction::RemoteNodeConnected, + crate::services::ops::audit::AuditEntityType::RemoteNode, + Some(42), + Some("edge-a"), + || { + super::details(super::RemoteNodeConnectionAuditDetails { + remote_node_id: 42, + binding_id: 42, + transport: "reverse_tunnel", + reason: "connected", + generation: 1, + outage_generation: 0, + active_lanes: 1, + lane_count: 1, + observed_at: chrono::Utc::now(), + first_lane_id: Some("lane-0"), + }) + }, + ) + .await; + + let entry = audit_log::Entity::find() + .filter(audit_log::Column::Action.eq(AuditAction::RemoteNodeConnected)) + .one(&db) + .await + .expect("remote node audit query should succeed") + .expect("remote node audit row should exist"); + assert_eq!(entry.user_id, 0); + assert_eq!(entry.entity_id, Some(42)); + let details = entry.details.expect("remote node details should exist"); + assert!(details.contains("\"transport\":\"reverse_tunnel\"")); + assert!(!details.contains("access-key")); + assert!(!details.contains("secret-key")); } #[tokio::test] diff --git a/src/storage/remote_protocol/runtime.rs b/src/storage/remote_protocol/runtime.rs index 1baee392a..135f13870 100644 --- a/src/storage/remote_protocol/runtime.rs +++ b/src/storage/remote_protocol/runtime.rs @@ -30,6 +30,11 @@ impl RemoteProtocolRuntime { self.tunnel_registry.set_persistence_db(db); } + pub fn set_audit_runtime_config(&self, runtime_config: Arc) { + self.tunnel_registry + .set_audit_runtime_config(runtime_config); + } + pub fn tunnel_registry(&self) -> &Arc { &self.tunnel_registry } diff --git a/src/storage/remote_protocol/tunnel/server/mod.rs b/src/storage/remote_protocol/tunnel/server/mod.rs index eec12d8e9..3774b3e36 100644 --- a/src/storage/remote_protocol/tunnel/server/mod.rs +++ b/src/storage/remote_protocol/tunnel/server/mod.rs @@ -38,6 +38,7 @@ pub use proxy::{ ClusterRemoteTunnelBroker, REMOTE_TUNNEL_PROXY_PATH_PREFIX, RemoteTunnelProxyQuery, proxy_tunnel_request, }; +pub(crate) use registry::TunnelDisconnectReason; pub use registry::{ RemoteTunnelBroker, RemoteTunnelHttpResponse, RemoteTunnelRegistry, RemoteTunnelStreamHttpResponse, reverse_tunnel_offline_error, @@ -104,6 +105,7 @@ pub async fn poll( let registry = state.remote_protocol().tunnel_registry(); let (request_rx, _registration) = registry.register_poll(remote_node); + registry.record_handshake(remote_node, None); managed_follower_repo::touch_tunnel_result( state.writer_db(), remote_node.id, @@ -210,7 +212,8 @@ async fn run_connected_stream( shutdown_token: CancellationToken, ) -> Result<()> { let registry = state.remote_protocol().tunnel_registry().clone(); - let (lane_id, mut request_rx, _registration) = registry.register_stream_lane(&remote_node); + let (lane_id, mut request_rx, registration) = registry.register_stream_lane(&remote_node); + registry.record_handshake(&remote_node, Some(&lane_id)); tracing::info!( remote_node_id = remote_node.id, lane_id = %lane_id, @@ -235,6 +238,7 @@ async fn run_connected_stream( heartbeat.tick().await; let mut liveness = TunnelHeartbeat::new(Instant::now()); let mut draining = false; + let disconnect_reason: Option; let mut drain_deadline = Box::pin(tokio::time::sleep(Duration::from_secs(24 * 60 * 60))); loop { @@ -258,6 +262,7 @@ async fn run_connected_stream( lane_id = %lane_id, "reverse tunnel streaming lane closing for primary shutdown" ); + disconnect_reason = Some(TunnelDisconnectReason::GracefulShutdown); break; } } @@ -268,6 +273,7 @@ async fn run_connected_stream( timeout_secs = REMOTE_TUNNEL_SHUTDOWN_DRAIN_TIMEOUT.as_secs(), "reverse tunnel shutdown drain timed out; closing streaming lane with in-flight request" ); + disconnect_reason = Some(TunnelDisconnectReason::GracefulShutdown); break; } _ = owner_renewal.tick(), if owner_directory.is_some() => { @@ -283,6 +289,7 @@ async fn run_connected_stream( runtime_id = %directory.runtime_id(), "reverse tunnel owner lease was fenced by another primary" ); + disconnect_reason = Some(TunnelDisconnectReason::OwnerFenced); break; } Err(error) => { @@ -291,6 +298,7 @@ async fn run_connected_stream( runtime_id = %directory.runtime_id(), "reverse tunnel owner lease renewal failed: {error}" ); + disconnect_reason = Some(TunnelDisconnectReason::OwnerRenewalFailed); break; } } @@ -304,6 +312,7 @@ async fn run_connected_stream( timeout_secs = REMOTE_TUNNEL_HEARTBEAT_TIMEOUT.as_secs(), "reverse tunnel streaming lane heartbeat timed out waiting for follower activity" ); + disconnect_reason = Some(TunnelDisconnectReason::HeartbeatTimeout); break; } if let Err(error) = session.ping(b"aster-tunnel-heartbeat").await { @@ -312,11 +321,13 @@ async fn run_connected_stream( lane_id = %lane_id, "failed to send reverse tunnel heartbeat ping: {error}" ); + disconnect_reason = Some(TunnelDisconnectReason::HeartbeatSendFailed); break; } } message = stream.next() => { let Some(message) = message else { + disconnect_reason = Some(TunnelDisconnectReason::Eof); break; }; let message = match message { @@ -327,6 +338,11 @@ async fn run_connected_stream( lane_id = %lane_id, "reverse tunnel streaming lane read failed: {error}" ); + disconnect_reason = Some(if error.to_string().to_ascii_lowercase().contains("reset") { + TunnelDisconnectReason::ConnectionReset + } else { + TunnelDisconnectReason::ProtocolReadError + }); break; } }; @@ -353,12 +369,14 @@ async fn run_connected_stream( lane_id = %lane_id, "failed to decode reverse tunnel streaming frame: {error}" ); + disconnect_reason = Some(TunnelDisconnectReason::ProtocolDecodeError); break; } } } actix_ws::Message::Ping(bytes) => { if session.pong(&bytes).await.is_err() { + disconnect_reason = Some(TunnelDisconnectReason::ProtocolReadError); break; } liveness.record_activity(Instant::now()); @@ -375,6 +393,7 @@ async fn run_connected_stream( close_reason = ?reason, "reverse tunnel streaming lane closed by follower" ); + disconnect_reason = Some(TunnelDisconnectReason::PeerClose); break; } _ => {} @@ -382,10 +401,12 @@ async fn run_connected_stream( } frame = request_rx.recv() => { let Some(frame) = frame else { + disconnect_reason = Some(TunnelDisconnectReason::Eof); break; }; let bytes = encode_stream_frame(&frame)?; if session.binary(bytes).await.is_err() { + disconnect_reason = Some(TunnelDisconnectReason::ConnectionReset); break; } } @@ -396,11 +417,12 @@ async fn run_connected_stream( lane_id = %lane_id, "reverse tunnel streaming lane drained before primary shutdown" ); + disconnect_reason = Some(TunnelDisconnectReason::GracefulShutdown); break; } } - close_connected_stream( + let close_handshake_ok = close_connected_stream( session, stream, remote_node.id, @@ -415,6 +437,11 @@ async fn run_connected_stream( }, ) .await; + let mut final_disconnect_reason = disconnect_reason.unwrap_or(TunnelDisconnectReason::Eof); + if !close_handshake_ok { + final_disconnect_reason = TunnelDisconnectReason::CloseHandshakeFailed; + } + registration.set_disconnect_reason(final_disconnect_reason); Ok(()) } @@ -445,14 +472,14 @@ async fn close_connected_stream( remote_node_id: i64, lane_id: &str, reason: Option, -) { +) -> bool { if let Err(error) = session.close(reason).await { tracing::warn!( remote_node_id, lane_id, "failed to send reverse tunnel streaming lane close frame: {error}" ); - return; + return false; } let handshake = tokio::time::timeout(REMOTE_TUNNEL_CLOSE_HANDSHAKE_TIMEOUT, async { while let Some(message) = stream.next().await { @@ -466,18 +493,24 @@ async fn close_connected_stream( }) .await; match handshake { - Ok(Ok(())) => {} - Ok(Err(error)) => tracing::warn!( - remote_node_id, - lane_id, - "reverse tunnel streaming lane close handshake read failed: {error}" - ), - Err(_) => tracing::warn!( - remote_node_id, - lane_id, - timeout_secs = REMOTE_TUNNEL_CLOSE_HANDSHAKE_TIMEOUT.as_secs(), - "reverse tunnel streaming lane close handshake timed out" - ), + Ok(Ok(())) => true, + Ok(Err(error)) => { + tracing::warn!( + remote_node_id, + lane_id, + "reverse tunnel streaming lane close handshake read failed: {error}" + ); + false + } + Err(_) => { + tracing::warn!( + remote_node_id, + lane_id, + timeout_secs = REMOTE_TUNNEL_CLOSE_HANDSHAKE_TIMEOUT.as_secs(), + "reverse tunnel streaming lane close handshake timed out" + ); + false + } } } diff --git a/src/storage/remote_protocol/tunnel/server/registry/mod.rs b/src/storage/remote_protocol/tunnel/server/registry/mod.rs index b0ead1359..1bbd92ceb 100644 --- a/src/storage/remote_protocol/tunnel/server/registry/mod.rs +++ b/src/storage/remote_protocol/tunnel/server/registry/mod.rs @@ -5,7 +5,10 @@ use dashmap::DashMap; use sea_orm::DatabaseConnection; use tokio::sync::Notify; +use crate::config::RuntimeConfig; +use crate::services::ops::audit::{self, AuditContext, AuditLogInput}; use aster_drive_model::entities::managed_follower; +use aster_drive_model::types::{AuditAction, AuditEntityType}; use aster_drive_storage::StorageErrorKind; mod broker; @@ -25,6 +28,82 @@ const REMOTE_TUNNEL_REQUEST_TIMEOUT: Duration = Duration::from_secs(60 * 60); const REMOTE_TUNNEL_ONLINE_TTL: Duration = Duration::from_secs(75); const REMOTE_TUNNEL_STREAM_CHANNEL_CAPACITY: usize = 16; +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum TunnelDisconnectReason { + GracefulShutdown, + PeerClose, + Eof, + ConnectionReset, + ProtocolReadError, + ProtocolDecodeError, + HeartbeatTimeout, + HeartbeatSendFailed, + OwnerFenced, + OwnerRenewalFailed, + CloseHandshakeFailed, +} + +impl TunnelDisconnectReason { + fn action(self) -> AuditAction { + match self { + Self::GracefulShutdown | Self::PeerClose => AuditAction::RemoteNodeGracefulDisconnect, + Self::HeartbeatTimeout => AuditAction::RemoteNodeHeartbeatTimeout, + Self::Eof + | Self::ConnectionReset + | Self::ProtocolReadError + | Self::ProtocolDecodeError + | Self::HeartbeatSendFailed + | Self::OwnerFenced + | Self::OwnerRenewalFailed + | Self::CloseHandshakeFailed => AuditAction::RemoteNodeUnexpectedDisconnect, + } + } + + const fn as_str(self) -> &'static str { + match self { + Self::GracefulShutdown => "primary_shutdown", + Self::PeerClose => "peer_close", + Self::Eof => "eof", + Self::ConnectionReset => "connection_reset", + Self::ProtocolReadError => "protocol_read_error", + Self::ProtocolDecodeError => "protocol_decode_error", + Self::HeartbeatTimeout => "heartbeat_timeout", + Self::HeartbeatSendFailed => "heartbeat_send_failed", + Self::OwnerFenced => "owner_fenced", + Self::OwnerRenewalFailed => "owner_renewal_failed", + Self::CloseHandshakeFailed => "close_handshake_failed", + } + } + + const fn priority(self) -> u8 { + match self { + Self::GracefulShutdown | Self::PeerClose => 0, + Self::Eof + | Self::ConnectionReset + | Self::ProtocolReadError + | Self::ProtocolDecodeError + | Self::HeartbeatSendFailed + | Self::OwnerFenced + | Self::OwnerRenewalFailed + | Self::CloseHandshakeFailed => 1, + Self::HeartbeatTimeout => 2, + } + } +} + +#[derive(Debug, Clone, Default)] +struct ConnectionLifecycleState { + online: bool, + generation: u64, + outage_generation: u64, + active_lanes: usize, + lane_count: usize, + first_lane_id: Option, + pending_disconnect_reason: Option, + last_disconnect_reason: Option, + observation_revision: u64, +} + #[derive(Default)] pub struct RemoteTunnelRegistry { connections: DashMap, @@ -33,7 +112,9 @@ pub struct RemoteTunnelRegistry { stream_pending: DashMap, last_errors: DashMap, last_seen_at: DashMap>, + lifecycle: DashMap, persistence_db: parking_lot::RwLock>, + audit_runtime_config: parking_lot::RwLock>>, connection_notify: Notify, } @@ -46,6 +127,10 @@ impl RemoteTunnelRegistry { *self.persistence_db.write() = Some(db); } + pub fn set_audit_runtime_config(&self, runtime_config: Arc) { + *self.audit_runtime_config.write() = Some(runtime_config); + } + pub fn is_online(&self, remote_node: &managed_follower::Model) -> bool { let local_last_seen = self .last_seen_at @@ -60,6 +145,183 @@ impl RemoteTunnelRegistry { self.last_seen_at.insert(remote_node_id, chrono::Utc::now()); } + pub(crate) fn record_handshake( + self: &Arc, + remote_node: &managed_follower::Model, + lane_id: Option<&str>, + ) { + let (event, observation_revision) = { + let mut state = self.lifecycle.entry(remote_node.id).or_default(); + state.observation_revision = state.observation_revision.saturating_add(1); + state.pending_disconnect_reason = None; + if let Some(lane_id) = lane_id { + state.active_lanes = state.active_lanes.saturating_add(1); + state.lane_count = state.lane_count.max(state.active_lanes); + if state.first_lane_id.is_none() { + state.first_lane_id = Some(lane_id.to_string()); + } + } + let event = if state.online { + None + } else { + state.online = true; + state.generation = state.generation.saturating_add(1); + Some((state.clone(), lane_id.map(ToOwned::to_owned))) + }; + (event, state.observation_revision) + }; + if let Some((state, lane_id)) = event { + self.record_lifecycle_audit( + remote_node, + AuditAction::RemoteNodeConnected, + "connected", + &state, + lane_id.as_deref().or(state.first_lane_id.as_deref()), + chrono::Utc::now(), + ); + } + self.schedule_lifecycle_expiry(remote_node.clone(), observation_revision); + } + + pub(crate) fn record_stream_disconnect( + self: &Arc, + remote_node: &managed_follower::Model, + reason: TunnelDisconnectReason, + ) { + let graceful_event = { + let mut state = self.lifecycle.entry(remote_node.id).or_default(); + if !state.online { + return; + } + state.active_lanes = state.active_lanes.saturating_sub(1); + state.pending_disconnect_reason = Some( + state + .pending_disconnect_reason + .filter(|current| current.priority() >= reason.priority()) + .unwrap_or(reason), + ); + if state.active_lanes == 0 && reason.priority() == 0 { + state.observation_revision = state.observation_revision.saturating_add(1); + state.online = false; + state.outage_generation = state.outage_generation.saturating_add(1); + let reason = state.pending_disconnect_reason.take().unwrap_or(reason); + state.last_disconnect_reason = Some(reason); + let event = state.clone(); + state.first_lane_id = None; + state.lane_count = 0; + Some((event, reason)) + } else { + None + } + }; + if let Some((state, reason)) = graceful_event { + self.record_lifecycle_audit( + remote_node, + reason.action(), + reason.as_str(), + &state, + state.first_lane_id.as_deref(), + chrono::Utc::now(), + ); + } + } + + fn schedule_lifecycle_expiry( + self: &Arc, + remote_node: managed_follower::Model, + observation_revision: u64, + ) { + let registry = self.clone(); + tokio::spawn(async move { + tokio::time::sleep(REMOTE_TUNNEL_ONLINE_TTL).await; + registry.expire_lifecycle_if_stale(&remote_node, observation_revision); + }); + } + + fn expire_lifecycle_if_stale( + &self, + remote_node: &managed_follower::Model, + observation_revision: u64, + ) { + let event = { + let mut state = self.lifecycle.entry(remote_node.id).or_default(); + if !state.online || state.observation_revision != observation_revision { + return; + } + state.online = false; + state.outage_generation = state.outage_generation.saturating_add(1); + let reason = state + .pending_disconnect_reason + .take() + .unwrap_or(TunnelDisconnectReason::Eof); + state.last_disconnect_reason = Some(reason); + let event = state.clone(); + state.active_lanes = 0; + state.first_lane_id = None; + state.lane_count = 0; + (event, reason) + }; + let (state, reason) = event; + self.record_lifecycle_audit( + remote_node, + reason.action(), + reason.as_str(), + &state, + state.first_lane_id.as_deref(), + chrono::Utc::now(), + ); + } + + fn record_lifecycle_audit( + &self, + remote_node: &managed_follower::Model, + action: AuditAction, + reason: &'static str, + state: &ConnectionLifecycleState, + first_lane_id: Option<&str>, + observed_at: chrono::DateTime, + ) { + let Some(db) = self.persistence_db.read().clone() else { + return; + }; + let Some(runtime_config) = self.audit_runtime_config.read().clone() else { + return; + }; + let details = audit::details(audit::RemoteNodeConnectionAuditDetails { + remote_node_id: remote_node.id, + binding_id: remote_node.id, + transport: remote_node + .transport_mode + .resolve(&remote_node.base_url) + .as_str(), + reason, + generation: state.generation, + outage_generation: state.outage_generation, + active_lanes: state.active_lanes, + lane_count: state.lane_count, + observed_at, + first_lane_id, + }); + let entity_id = remote_node.id; + let entity_name = remote_node.name.clone(); + tokio::spawn(async move { + let ctx = AuditContext::system(); + audit::log_with_db_and_config( + &db, + &runtime_config, + AuditLogInput { + ctx: &ctx, + action, + entity_type: AuditEntityType::RemoteNode, + entity_id: Some(entity_id), + entity_name: Some(&entity_name), + }, + || details, + ) + .await; + }); + } + pub fn last_error(&self, remote_node_id: i64) -> Option { self.last_errors .get(&remote_node_id) @@ -119,6 +381,29 @@ pub fn reverse_tunnel_offline_error(remote_node_id: i64) -> crate::errors::Aster #[cfg(test)] mod tests { use super::*; + use aster_drive_model::types::RemoteNodeTransportMode; + + fn test_remote_node() -> managed_follower::Model { + let now = chrono::Utc::now(); + managed_follower::Model { + id: 42, + name: "edge-a".to_string(), + base_url: String::new(), + access_key: "access-key".to_string(), + secret_key: "secret-key".to_string(), + is_enabled: true, + transport_mode: RemoteNodeTransportMode::ReverseTunnel, + last_capabilities: "{}".to_string(), + last_error: String::new(), + last_checked_at: None, + tunnel_last_error: String::new(), + tunnel_last_seen_at: None, + binding_revision: 1, + binding_applied_revision: 1, + created_at: now, + updated_at: now, + } + } #[test] fn persisted_tunnel_seen_time_obeys_online_ttl_boundary() { @@ -128,4 +413,183 @@ mod tests { - chrono::Duration::milliseconds(1); assert!(!is_recent_tunnel_seen_at(expired)); } + + #[tokio::test] + async fn lifecycle_aggregates_four_lanes_into_one_outage_generation() { + let registry = Arc::new(RemoteTunnelRegistry::new()); + let node = test_remote_node(); + + for lane in ["lane-0", "lane-1", "lane-2", "lane-3"] { + registry.record_handshake(&node, Some(lane)); + } + registry.record_handshake(&node, None); + registry.record_handshake(&node, None); + let state = registry.lifecycle.get(&node.id).expect("lifecycle state"); + assert!(state.online); + assert_eq!(state.generation, 1); + assert_eq!(state.outage_generation, 0); + assert_eq!(state.active_lanes, 4); + assert_eq!(state.lane_count, 4); + drop(state); + + for _ in 0..3 { + registry.record_stream_disconnect(&node, TunnelDisconnectReason::Eof); + } + let state = registry.lifecycle.get(&node.id).expect("lifecycle state"); + assert!(state.online); + assert_eq!(state.active_lanes, 1); + drop(state); + + registry.record_stream_disconnect(&node, TunnelDisconnectReason::HeartbeatTimeout); + let state = registry.lifecycle.get(&node.id).expect("lifecycle state"); + assert!(state.online); + assert_eq!(state.active_lanes, 0); + assert_eq!(state.generation, 1); + assert_eq!(state.outage_generation, 0); + assert_eq!(state.lane_count, 4); + assert_eq!(state.last_disconnect_reason, None); + drop(state); + + registry.record_stream_disconnect(&node, TunnelDisconnectReason::Eof); + let state = registry.lifecycle.get(&node.id).expect("lifecycle state"); + assert_eq!(state.outage_generation, 0); + let revision = state.observation_revision; + drop(state); + + registry.expire_lifecycle_if_stale(&node, revision); + let state = registry.lifecycle.get(&node.id).expect("lifecycle state"); + assert!(!state.online); + assert_eq!(state.outage_generation, 1); + assert_eq!(state.lane_count, 0); + assert_eq!( + state.last_disconnect_reason, + Some(TunnelDisconnectReason::HeartbeatTimeout) + ); + drop(state); + + registry.record_handshake(&node, Some("lane-recovered")); + let state = registry.lifecycle.get(&node.id).expect("lifecycle state"); + assert!(state.online); + assert_eq!(state.generation, 2); + assert_eq!(state.outage_generation, 1); + assert_eq!(state.active_lanes, 1); + } + + #[test] + fn lifecycle_reason_codes_keep_graceful_and_failure_actions_distinct() { + assert_eq!( + TunnelDisconnectReason::GracefulShutdown.action(), + AuditAction::RemoteNodeGracefulDisconnect + ); + assert_eq!( + TunnelDisconnectReason::PeerClose.action(), + AuditAction::RemoteNodeGracefulDisconnect + ); + assert_eq!( + TunnelDisconnectReason::HeartbeatTimeout.action(), + AuditAction::RemoteNodeHeartbeatTimeout + ); + assert_eq!( + TunnelDisconnectReason::ConnectionReset.action(), + AuditAction::RemoteNodeUnexpectedDisconnect + ); + assert_eq!( + TunnelDisconnectReason::CloseHandshakeFailed.action(), + AuditAction::RemoteNodeUnexpectedDisconnect + ); + assert_eq!(TunnelDisconnectReason::OwnerFenced.as_str(), "owner_fenced"); + } + + #[tokio::test] + async fn lifecycle_keeps_strongest_lane_failure_until_aggregate_disconnect() { + let registry = Arc::new(RemoteTunnelRegistry::new()); + let node = test_remote_node(); + registry.record_handshake(&node, Some("lane-0")); + registry.record_handshake(&node, Some("lane-1")); + + registry.record_stream_disconnect(&node, TunnelDisconnectReason::HeartbeatTimeout); + registry.record_stream_disconnect(&node, TunnelDisconnectReason::PeerClose); + + let state = registry.lifecycle.get(&node.id).expect("lifecycle state"); + assert_eq!(state.outage_generation, 1); + assert_eq!( + state.last_disconnect_reason, + Some(TunnelDisconnectReason::HeartbeatTimeout) + ); + } + + #[tokio::test] + async fn stale_expiry_cannot_disconnect_after_a_new_poll_handshake() { + let registry = Arc::new(RemoteTunnelRegistry::new()); + let node = test_remote_node(); + registry.record_handshake(&node, Some("lane-0")); + let stale_revision = registry + .lifecycle + .get(&node.id) + .expect("lifecycle state") + .observation_revision; + registry.record_stream_disconnect(&node, TunnelDisconnectReason::Eof); + registry.record_handshake(&node, None); + + registry.expire_lifecycle_if_stale(&node, stale_revision); + + let state = registry.lifecycle.get(&node.id).expect("lifecycle state"); + assert!(state.online); + assert_eq!(state.generation, 1); + assert_eq!(state.outage_generation, 0); + assert_eq!(state.pending_disconnect_reason, None); + } + + #[tokio::test] + async fn graceful_last_lane_disconnects_immediately() { + let registry = Arc::new(RemoteTunnelRegistry::new()); + let node = test_remote_node(); + registry.record_handshake(&node, Some("lane-0")); + + registry.record_stream_disconnect(&node, TunnelDisconnectReason::GracefulShutdown); + + let state = registry.lifecycle.get(&node.id).expect("lifecycle state"); + assert!(!state.online); + assert_eq!(state.generation, 1); + assert_eq!(state.outage_generation, 1); + assert_eq!( + state.last_disconnect_reason, + Some(TunnelDisconnectReason::GracefulShutdown) + ); + } + + #[test] + fn lifecycle_details_are_redacted_and_do_not_include_credentials() { + let node = test_remote_node(); + let state = ConnectionLifecycleState { + online: true, + generation: 3, + outage_generation: 2, + active_lanes: 1, + lane_count: 4, + first_lane_id: Some("lane-0".to_string()), + pending_disconnect_reason: None, + last_disconnect_reason: Some(TunnelDisconnectReason::ConnectionReset), + observation_revision: 3, + }; + let details = audit::details(audit::RemoteNodeConnectionAuditDetails { + remote_node_id: node.id, + binding_id: node.id, + transport: node.transport_mode.resolve(&node.base_url).as_str(), + reason: "connected", + generation: state.generation, + outage_generation: state.outage_generation, + active_lanes: state.active_lanes, + lane_count: state.lane_count, + observed_at: chrono::Utc::now(), + first_lane_id: state.first_lane_id.as_deref(), + }) + .expect("details should serialize"); + let encoded = details.to_string(); + assert!(encoded.contains("binding_id")); + assert!(!encoded.contains(&node.access_key)); + assert!(!encoded.contains(&node.secret_key)); + assert!(!encoded.contains("signature")); + assert!(!encoded.contains("token")); + } } diff --git a/src/storage/remote_protocol/tunnel/server/registry/streaming.rs b/src/storage/remote_protocol/tunnel/server/registry/streaming.rs index ef589da26..d36efbc15 100644 --- a/src/storage/remote_protocol/tunnel/server/registry/streaming.rs +++ b/src/storage/remote_protocol/tunnel/server/registry/streaming.rs @@ -13,7 +13,8 @@ use super::broker::RemoteTunnelStreamHttpResponse; use super::headers::request_headers; use super::{ REMOTE_TUNNEL_CONNECT_WAIT_TIMEOUT, REMOTE_TUNNEL_REQUEST_TIMEOUT, - REMOTE_TUNNEL_STREAM_CHANNEL_CAPACITY, RemoteTunnelRegistry, reverse_tunnel_offline_error, + REMOTE_TUNNEL_STREAM_CHANNEL_CAPACITY, RemoteTunnelRegistry, TunnelDisconnectReason, + reverse_tunnel_offline_error, }; use crate::errors::{AsterError, Result}; use crate::storage::remote_protocol::tunnel::server::response::header_pairs_to_map; @@ -112,7 +113,15 @@ pub(crate) struct RemoteTunnelStreamRegistrationGuard { registry: Arc, access_key: String, remote_node_id: i64, + remote_node: managed_follower::Model, lane_id: String, + disconnect_reason: parking_lot::Mutex, +} + +impl RemoteTunnelStreamRegistrationGuard { + pub(crate) fn set_disconnect_reason(&self, reason: TunnelDisconnectReason) { + *self.disconnect_reason.lock() = reason; + } } impl Drop for RemoteTunnelStreamRegistrationGuard { @@ -124,6 +133,8 @@ impl Drop for RemoteTunnelStreamRegistrationGuard { ); self.registry .fail_stream_requests_for_lane(&self.lane_id, "reverse tunnel streaming lane closed"); + self.registry + .record_stream_disconnect(&self.remote_node, *self.disconnect_reason.lock()); } } @@ -158,7 +169,9 @@ impl RemoteTunnelRegistry { registry: self.clone(), access_key: remote_node.access_key.clone(), remote_node_id: remote_node.id, + remote_node: remote_node.clone(), lane_id: lane_id.clone(), + disconnect_reason: parking_lot::Mutex::new(TunnelDisconnectReason::Eof), }; (lane_id, request_rx, guard) } From 16b2decc7f5d73e31568e9a2999a806a01e0ffb1 Mon Sep 17 00:00:00 2001 From: AptS-1547 Date: Thu, 20 Aug 2026 23:45:31 +0800 Subject: [PATCH 2/3] fix(remote): address lifecycle audit review --- .../src/i18n/locales/en/admin/audit.json | 2 +- .../src/i18n/locales/zh/admin/audit.json | 2 +- src/services/ops/audit/mod.rs | 1 + src/services/ops/audit/presentation.rs | 93 +++++++++++++++++++ src/services/ops/audit/query.rs | 18 ++-- .../remote_protocol/tunnel/server/mod.rs | 16 +++- .../remote_protocol/tunnel/server/tests.rs | 16 ++++ 7 files changed, 135 insertions(+), 13 deletions(-) diff --git a/frontend-panel/src/i18n/locales/en/admin/audit.json b/frontend-panel/src/i18n/locales/en/admin/audit.json index 30228d2a5..fc6ad84a7 100644 --- a/frontend-panel/src/i18n/locales/en/admin/audit.json +++ b/frontend-panel/src/i18n/locales/en/admin/audit.json @@ -229,7 +229,7 @@ "audit_presentation_mfa_email_code_sent": "{{method}} flow {{flow_id}}, expires in {{expires_in}}s", "audit_presentation_wopi_user_info_updated": "{{app_key}} for file {{file_id}}, {{user_info_len}} bytes", "audit_presentation_remote_enrollment_changed": "{{phase}} {{remote_node_name}}, enabled {{is_enabled}}", - "audit_presentation_remote_node_connection_lifecycle": "{{reason}} on {{transport}} (generation {{generation}}, outage {{outage_generation}}, {{active_lanes}} active / {{lane_count}} lane(s))", + "audit_presentation_remote_node_connection_lifecycle": "{{reason}} on {{transport}} (connection #{{generation}}, interruption #{{outage_generation}}, {{active_lanes}} active / {{lane_count}} lane(s))", "audit_presentation_invitation_snapshot": "{{email}} status {{status}}, expires {{expires_at}}", "audit_presentation_external_auth_unlinked": "{{provider_key}} identity {{subject}}", "audit_presentation_follower_binding_synced": "{{name}} enabled {{is_enabled}}", diff --git a/frontend-panel/src/i18n/locales/zh/admin/audit.json b/frontend-panel/src/i18n/locales/zh/admin/audit.json index 3803e85f9..dbab7cb14 100644 --- a/frontend-panel/src/i18n/locales/zh/admin/audit.json +++ b/frontend-panel/src/i18n/locales/zh/admin/audit.json @@ -229,7 +229,7 @@ "audit_presentation_mfa_email_code_sent": "{{method}} flow {{flow_id}},{{expires_in}} 秒后过期", "audit_presentation_wopi_user_info_updated": "{{app_key}} 更新文件 {{file_id}} 的用户信息,{{user_info_len}} 字节", "audit_presentation_remote_enrollment_changed": "{{phase}} {{remote_node_name}},启用 {{is_enabled}}", - "audit_presentation_remote_node_connection_lifecycle": "{{transport}} 上发生 {{reason}}(连接代际 {{generation}},故障代际 {{outage_generation}},{{active_lanes}} 条活跃 / 共 {{lane_count}} 条 lane)", + "audit_presentation_remote_node_connection_lifecycle": "{{transport}} 上发生 {{reason}}(第 {{generation}} 次连接,第 {{outage_generation}} 次中断,{{active_lanes}} 条活跃 / 共 {{lane_count}} 条 lane)", "audit_presentation_invitation_snapshot": "{{email}} 状态 {{status}},过期 {{expires_at}}", "audit_presentation_external_auth_unlinked": "{{provider_key}} 身份 {{subject}}", "audit_presentation_follower_binding_synced": "{{name}} 启用 {{is_enabled}}", diff --git a/src/services/ops/audit/mod.rs b/src/services/ops/audit/mod.rs index 007b0278f..43456e1f3 100644 --- a/src/services/ops/audit/mod.rs +++ b/src/services/ops/audit/mod.rs @@ -41,4 +41,5 @@ pub use manager::{ should_record, should_record_with_config, }; pub use models::{AuditLogEntry, AuditPresentation, AuditPresentationMessage, TeamAuditEntryInfo}; +pub use presentation::{sanitize_details, sanitize_entity_name}; pub use query::{cleanup_expired, query, query_team_entries}; diff --git a/src/services/ops/audit/presentation.rs b/src/services/ops/audit/presentation.rs index 00e1c05e6..42441bade 100644 --- a/src/services/ops/audit/presentation.rs +++ b/src/services/ops/audit/presentation.rs @@ -6,6 +6,77 @@ use aster_drive_model::types::{AuditAction, AuditEntityType}; use super::models::{AuditPresentation, AuditPresentationMessage}; +/// Neutralize spreadsheet formula prefixes before user-controlled audit names +/// reach the admin UI or an exported representation. +fn neutralize_formula(value: String) -> String { + if value.starts_with(['=', '+', '-', '@', '\t', '\r']) { + format!("'{value}") + } else { + value + } +} + +fn sensitive_detail_key(key: &str) -> bool { + let normalized = key + .chars() + .filter(|character| character.is_ascii_alphanumeric()) + .flat_map(char::to_lowercase) + .collect::(); + [ + "token", + "secret", + "password", + "credential", + "authorization", + "signature", + "accesskey", + "apikey", + "session", + "mfa", + "otp", + "totp", + "bearer", + "appkey", + "wopikey", + "sharetoken", + "privatekey", + ] + .iter() + .any(|needle| normalized.contains(needle)) +} + +fn redact_sensitive_details(value: &mut Value) { + match value { + Value::Object(object) => { + object.retain(|key, value| { + if sensitive_detail_key(key) { + return false; + } + redact_sensitive_details(value); + true + }); + } + Value::Array(values) => values.iter_mut().for_each(redact_sensitive_details), + _ => {} + } +} + +pub fn sanitize_details(raw: Option<&str>) -> Option { + let mut details = serde_json::from_str::(raw?).ok()?; + redact_sensitive_details(&mut details); + Some(details.to_string()) +} + +pub fn sanitize_entity_name(entity_type: &str, value: Option) -> Option { + value.map(|value| { + if entity_type == "share" { + value + } else { + neutralize_formula(value) + } + }) +} + pub fn build_audit_presentation( action: AuditAction, entity_type: AuditEntityType, @@ -1385,6 +1456,28 @@ mod tests { assert_eq!(detail.params.get("lane_count"), Some(&Value::from(4))); } + #[test] + fn audit_query_sanitizers_redact_sensitive_details_and_formula_names() { + assert_eq!( + sanitize_entity_name("remote_node", Some("=HYPERLINK(\"x\")".to_string())), + Some("'=HYPERLINK(\"x\")".to_string()) + ); + assert_eq!( + sanitize_entity_name("share", Some("=share-token".to_string())), + Some("=share-token".to_string()) + ); + let safe = sanitize_details(Some( + r#"{"remote_node_id":42,"access_key":"access","secret_key":"secret","nested":[{"token":"tok"}],"transport":"reverse_tunnel"}"#, + )) + .expect("valid details should sanitize"); + assert!(safe.contains("remote_node_id")); + assert!(safe.contains("transport")); + assert!(!safe.contains("access_key")); + assert!(!safe.contains("secret_key")); + assert!(!safe.contains("token")); + assert!(sanitize_details(Some("not-json")).is_none()); + } + #[test] fn presentation_includes_invitation_detail() { let presentation = build_audit_presentation( diff --git a/src/services/ops/audit/query.rs b/src/services/ops/audit/query.rs index 7ac53796d..cb974c7c8 100644 --- a/src/services/ops/audit/query.rs +++ b/src/services/ops/audit/query.rs @@ -12,7 +12,7 @@ use aster_forge_api::{OffsetPage, SortOrder}; use super::filters::AuditLogFilters; use super::models::{AuditLogEntry, TeamAuditEntryInfo}; -use super::presentation::build_audit_presentation; +use super::presentation::{build_audit_presentation, sanitize_details, sanitize_entity_name}; const DEFAULT_RETENTION_DAYS: i64 = 90; @@ -76,12 +76,14 @@ async fn build_audit_entries( // Build presentation at read time so historical audit rows keep using // the current presentation rules without duplicating display snapshots. + let safe_entity_name = sanitize_entity_name(&model.entity_type, model.entity_name.clone()); + let safe_details = sanitize_details(model.details.as_deref()); let presentation = build_audit_presentation( model.action, entity_type, model.entity_id, - model.entity_name.as_deref(), - model.details.as_deref(), + safe_entity_name.as_deref(), + safe_details.as_deref(), ); items.push(AuditLogEntry { @@ -90,8 +92,8 @@ async fn build_audit_entries( action: model.action, entity_type, entity_id: model.entity_id, - entity_name: model.entity_name, - details: model.details, + entity_name: safe_entity_name, + details: safe_details, presentation, ip_address: model.ip_address, user_agent: model.user_agent, @@ -132,6 +134,8 @@ fn build_team_audit_entry( .details .as_deref() .and_then(|raw| serde_json::from_str::(raw).ok()); + let safe_entity_name = sanitize_entity_name(&entry.entity_type, entry.entity_name.clone()); + let safe_details = sanitize_details(entry.details.as_deref()); let member_user_id = parsed_details .as_ref() @@ -161,8 +165,8 @@ fn build_team_audit_entry( entry.action, entity_type, entry.entity_id, - entry.entity_name.as_deref(), - entry.details.as_deref(), + safe_entity_name.as_deref(), + safe_details.as_deref(), ) }), actor: users.get(&entry.user_id).cloned(), diff --git a/src/storage/remote_protocol/tunnel/server/mod.rs b/src/storage/remote_protocol/tunnel/server/mod.rs index 3774b3e36..83e1e187e 100644 --- a/src/storage/remote_protocol/tunnel/server/mod.rs +++ b/src/storage/remote_protocol/tunnel/server/mod.rs @@ -437,14 +437,22 @@ async fn run_connected_stream( }, ) .await; - let mut final_disconnect_reason = disconnect_reason.unwrap_or(TunnelDisconnectReason::Eof); - if !close_handshake_ok { - final_disconnect_reason = TunnelDisconnectReason::CloseHandshakeFailed; - } + let final_disconnect_reason = finalize_disconnect_reason(disconnect_reason, close_handshake_ok); registration.set_disconnect_reason(final_disconnect_reason); Ok(()) } +fn finalize_disconnect_reason( + reason: Option, + close_handshake_ok: bool, +) -> TunnelDisconnectReason { + match (reason, close_handshake_ok) { + (Some(reason), _) => reason, + (None, true) => TunnelDisconnectReason::Eof, + (None, false) => TunnelDisconnectReason::CloseHandshakeFailed, + } +} + #[derive(Debug, Clone, Copy)] struct TunnelHeartbeat { last_activity_at: Instant, diff --git a/src/storage/remote_protocol/tunnel/server/tests.rs b/src/storage/remote_protocol/tunnel/server/tests.rs index 8d8a6d499..f1d297a48 100644 --- a/src/storage/remote_protocol/tunnel/server/tests.rs +++ b/src/storage/remote_protocol/tunnel/server/tests.rs @@ -52,6 +52,22 @@ fn tunnel_heartbeat_tracks_activity_and_missed_pongs() { assert!(heartbeat.is_timed_out(last_pong_at + REMOTE_TUNNEL_HEARTBEAT_TIMEOUT)); } +#[test] +fn close_handshake_failure_does_not_override_recorded_disconnect_reason() { + assert_eq!( + finalize_disconnect_reason(Some(TunnelDisconnectReason::HeartbeatTimeout), false), + TunnelDisconnectReason::HeartbeatTimeout + ); + assert_eq!( + finalize_disconnect_reason(Some(TunnelDisconnectReason::Eof), false), + TunnelDisconnectReason::Eof + ); + assert_eq!( + finalize_disconnect_reason(None, false), + TunnelDisconnectReason::CloseHandshakeFailed + ); +} + #[test] fn reverse_tunnel_transport_gate_uses_effective_mode() { let mut node = build_remote_node(1, "transport-gate"); From 7d0a20ae22b3083fcb1b7c170dc94a3b0b30690a Mon Sep 17 00:00:00 2001 From: AptS-1547 Date: Fri, 21 Aug 2026 00:21:59 +0800 Subject: [PATCH 3/3] fix(audit): preserve safe boolean metadata --- src/services/ops/audit/presentation.rs | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/src/services/ops/audit/presentation.rs b/src/services/ops/audit/presentation.rs index 42441bade..283bb0780 100644 --- a/src/services/ops/audit/presentation.rs +++ b/src/services/ops/audit/presentation.rs @@ -49,7 +49,9 @@ fn redact_sensitive_details(value: &mut Value) { match value { Value::Object(object) => { object.retain(|key, value| { - if sensitive_detail_key(key) { + if sensitive_detail_key(key) + && matches!(value, Value::String(_) | Value::Object(_) | Value::Array(_)) + { return false; } redact_sensitive_details(value); @@ -1467,7 +1469,7 @@ mod tests { Some("=share-token".to_string()) ); let safe = sanitize_details(Some( - r#"{"remote_node_id":42,"access_key":"access","secret_key":"secret","nested":[{"token":"tok"}],"transport":"reverse_tunnel"}"#, + r#"{"remote_node_id":42,"access_key":"access","secret_key":"secret","nested":[{"token":"tok"}],"temporary_password_generated":true,"transport":"reverse_tunnel"}"#, )) .expect("valid details should sanitize"); assert!(safe.contains("remote_node_id")); @@ -1475,6 +1477,7 @@ mod tests { assert!(!safe.contains("access_key")); assert!(!safe.contains("secret_key")); assert!(!safe.contains("token")); + assert!(safe.contains("temporary_password_generated")); assert!(sanitize_details(Some("not-json")).is_none()); }