From 544304765f7f1c7115db1c609dfc7ad7adff0564 Mon Sep 17 00:00:00 2001 From: igor04091968 Date: Mon, 11 May 2026 21:57:23 +0300 Subject: [PATCH] feat(dlp): add WAL transport and collector health telemetry --- windows/ActivityWatch.Windows.Common.psm1 | 54 ++++++- windows/dlp-endpoint-signals-collector.ps1 | 129 ++++++++++++++++- windows/file-operations-collector.ps1 | 161 ++++++++++++++++++++- 3 files changed, 339 insertions(+), 5 deletions(-) diff --git a/windows/ActivityWatch.Windows.Common.psm1 b/windows/ActivityWatch.Windows.Common.psm1 index e9c854f..360aca7 100755 --- a/windows/ActivityWatch.Windows.Common.psm1 +++ b/windows/ActivityWatch.Windows.Common.psm1 @@ -83,6 +83,41 @@ function Get-ActivityWatchPackageRoot { return (Split-Path -Path (Split-Path -Path $afkBinary.FullName -Parent) -Parent) } +function Expand-ActivityWatchArchiveSafe { + param( + [Parameter(Mandatory = $true)] + [string]$ArchivePath, + [Parameter(Mandatory = $true)] + [string]$DestinationPath, + [int]$Attempts = 3 + ) + + for ($attempt = 1; $attempt -le $Attempts; $attempt++) { + try { + if (Test-Path -LiteralPath $DestinationPath) { + Remove-Item -LiteralPath $DestinationPath -Recurse -Force -ErrorAction SilentlyContinue + } + New-ActivityWatchDirectory -Path $DestinationPath + Expand-Archive -Path $ArchivePath -DestinationPath $DestinationPath -Force -ErrorAction Stop + return + } + catch { + if ($attempt -lt $Attempts) { + Start-Sleep -Milliseconds (500 * $attempt) + continue + } + } + } + + # Fallback for intermittent Expand-Archive issues in Windows PowerShell. + if (Test-Path -LiteralPath $DestinationPath) { + Remove-Item -LiteralPath $DestinationPath -Recurse -Force -ErrorAction SilentlyContinue + } + New-ActivityWatchDirectory -Path $DestinationPath + Add-Type -AssemblyName System.IO.Compression.FileSystem + [System.IO.Compression.ZipFile]::ExtractToDirectory($ArchivePath, $DestinationPath) +} + function Install-ActivityWatchPackage { param( [Parameter(Mandatory = $true)] @@ -98,6 +133,13 @@ function Install-ActivityWatchPackage { New-ActivityWatchDirectory -Path $WorkingRoot New-ActivityWatchDirectory -Path $BackupRoot + # Cleanup stale extraction directories from previous failed deployments. + Get-ChildItem -LiteralPath $WorkingRoot -Directory -ErrorAction SilentlyContinue | + Where-Object { $_.Name -like 'extract-*' } | + ForEach-Object { + try { Remove-Item -LiteralPath $_.FullName -Recurse -Force -ErrorAction SilentlyContinue } catch {} + } + # Ensure nothing is holding locks inside InstallRoot during upgrade. foreach ($procName in @('aw-watcher-afk', 'aw-watcher-window', 'aw-server', 'aw-qt')) { try { @@ -114,7 +156,17 @@ function Install-ActivityWatchPackage { } New-ActivityWatchDirectory -Path $extractRoot - Expand-Archive -Path $ArchivePath -DestinationPath $extractRoot -Force + $archiveSize = (Get-Item -LiteralPath $ArchivePath -ErrorAction Stop).Length + $workDrive = (Get-PSDrive -Name ([System.IO.Path]::GetPathRoot($WorkingRoot).TrimEnd('\').TrimEnd(':')) -ErrorAction SilentlyContinue) + if ($workDrive) { + # Require at least ~2.5x archive size to handle extraction + copy safely. + $required = [int64]([Math]::Ceiling($archiveSize * 2.5)) + if ([int64]$workDrive.Free -lt $required) { + throw ("Недостаточно свободного места на {0}: free={1} bytes, required>={2} bytes" -f $workDrive.Name, $workDrive.Free, $required) + } + } + + Expand-ActivityWatchArchiveSafe -ArchivePath $ArchivePath -DestinationPath $extractRoot $packageRoot = Get-ActivityWatchPackageRoot -ExpandedRoot $extractRoot if (Test-Path -LiteralPath $InstallRoot) { diff --git a/windows/dlp-endpoint-signals-collector.ps1 b/windows/dlp-endpoint-signals-collector.ps1 index 61c4960..9aa0325 100644 --- a/windows/dlp-endpoint-signals-collector.ps1 +++ b/windows/dlp-endpoint-signals-collector.ps1 @@ -28,6 +28,15 @@ try { catch { } +$script:TransportQueuePath = $null +$script:TransportQueueLockPath = $null +$script:TransportMetrics = @{ + eventsEnqueued = 0 + eventsFlushed = 0 + sendFailures = 0 + queueDepth = 0 +} + $policyClientModulePath = Join-Path $PSScriptRoot 'dlp-policy-client.ps1' if (Test-Path -LiteralPath $policyClientModulePath) { try { @@ -118,6 +127,98 @@ function Invoke-AwJsonPost { } } +function Initialize-TransportQueue { + param([Parameter(Mandatory = $true)][string]$StateRoot) + $script:TransportQueuePath = Join-Path $StateRoot 'dlp-endpoint-signals-queue.jsonl' + $script:TransportQueueLockPath = Join-Path $StateRoot 'dlp-endpoint-signals-queue.lock' + if (-not (Test-Path -LiteralPath $script:TransportQueuePath)) { + New-Item -Path $script:TransportQueuePath -ItemType File -Force | Out-Null + } +} + +function Get-TransportQueueLock { + $tries = 0 + while ($tries -lt 50) { + try { + return [System.IO.File]::Open($script:TransportQueueLockPath, [System.IO.FileMode]::OpenOrCreate, [System.IO.FileAccess]::ReadWrite, [System.IO.FileShare]::None) + } + catch { + Start-Sleep -Milliseconds 50 + $tries++ + } + } + throw "Failed to acquire transport queue lock: $script:TransportQueueLockPath" +} + +function Add-TransportQueueItem { + param( + [Parameter(Mandatory = $true)][string]$Uri, + [Parameter(Mandatory = $true)][string]$Payload, + [string]$Kind = 'endpoint' + ) + $lock = Get-TransportQueueLock + try { + $line = @{ + ts = (Get-Date).ToUniversalTime().ToString('o') + uri = $Uri + payload = $Payload + kind = $Kind + } | ConvertTo-Json -Compress + Add-Content -LiteralPath $script:TransportQueuePath -Value $line -Encoding UTF8 + $script:TransportMetrics.eventsEnqueued++ + } + finally { + $lock.Dispose() + } +} + +function Read-TransportQueueItems { + if (-not (Test-Path -LiteralPath $script:TransportQueuePath)) { return @() } + $items = @() + foreach ($line in @(Get-Content -LiteralPath $script:TransportQueuePath -ErrorAction SilentlyContinue)) { + if ([string]::IsNullOrWhiteSpace($line)) { continue } + try { $items += ($line | ConvertFrom-Json) } catch {} + } + return $items +} + +function Flush-TransportQueue { + param([int]$MaxItems = 200) + if (-not (Test-Path -LiteralPath $script:TransportQueuePath)) { return } + $lock = Get-TransportQueueLock + try { + $items = Read-TransportQueueItems + $script:TransportMetrics.queueDepth = $items.Count + if ($items.Count -eq 0) { return } + $left = New-Object System.Collections.Generic.List[object] + $sent = 0 + foreach ($item in $items) { + if ($sent -ge $MaxItems) { + $left.Add($item) + continue + } + try { + Invoke-AwJsonPost -Uri ([string]$item.uri) -Json ([string]$item.payload) + $sent++ + $script:TransportMetrics.eventsFlushed++ + } + catch { + $script:TransportMetrics.sendFailures++ + $left.Add($item) + } + } + foreach ($item in $items | Select-Object -Skip ($sent + $left.Count)) { + $left.Add($item) + } + $lines = @($left | ForEach-Object { $_ | ConvertTo-Json -Compress }) + Set-Content -LiteralPath $script:TransportQueuePath -Value $lines -Encoding UTF8 + $script:TransportMetrics.queueDepth = $left.Count + } + finally { + $lock.Dispose() + } +} + function Ensure-Bucket { param( [string]$BucketId, @@ -186,7 +287,8 @@ function Send-EndpointSignalHeartbeat { } + $Data } | ConvertTo-Json -Depth 6 -Compress - Invoke-AwJsonPost -Uri "$($script:ApiBase)/buckets/$bucketId/heartbeat?pulsetime=$script:PulseSeconds" -Json $payload + Add-TransportQueueItem -Uri "$($script:ApiBase)/buckets/$bucketId/heartbeat?pulsetime=$script:PulseSeconds" -Payload $payload -Kind 'endpoint_signal' + Flush-TransportQueue -MaxItems 50 } function Send-DlpIncidentHeartbeat { @@ -227,7 +329,8 @@ function Send-DlpIncidentHeartbeat { } + $Data + $captureData } | ConvertTo-Json -Depth 7 -Compress - Invoke-AwJsonPost -Uri "$($script:ApiBase)/buckets/$bucketId/heartbeat?pulsetime=$script:PulseSeconds" -Json $payload + Add-TransportQueueItem -Uri "$($script:ApiBase)/buckets/$bucketId/heartbeat?pulsetime=$script:PulseSeconds" -Payload $payload -Kind 'dlp_incident' + Flush-TransportQueue -MaxItems 100 } function Get-FileSha256Hex { @@ -1007,6 +1110,7 @@ $script:PolicyRefreshSeconds = [Math]::Max($resolvedPolicyRefreshSeconds, 60) $script:PolicyCachePath = $resolvedPolicyCachePath $script:LocalPolicyPath = $resolvedPolicyPath $script:LastPolicyRefreshAt = [datetime]::MinValue +$script:TransportBackoffSeconds = 1 # Integration test flag (backward compatible - defaults to false) $script:IntegrationTestEnabled = if ($deploymentConfig -and $deploymentConfig.PSObject.Properties.Name -contains 'integrationTestEnabled') { [bool]$deploymentConfig.integrationTestEnabled } else { $false } @@ -1014,11 +1118,21 @@ $script:IntegrationTestEnabled = if ($deploymentConfig -and $deploymentConfig.PS $script:TotalEventsProcessed = 0 $script:LastEventTime = $null +Initialize-TransportQueue -StateRoot $resolvedStateRoot Initialize-DlpPolicy Write-EndpointLog ("endpoint collector started against {0}" -f $script:ApiBase) while ($true) { try { + try { + Flush-TransportQueue -MaxItems 200 + $script:TransportBackoffSeconds = 1 + } + catch { + $script:TransportBackoffSeconds = [Math]::Min($script:TransportBackoffSeconds * 2, 60) + Write-EndpointLog ("transport flush failed, backoff={0}s err={1}" -f $script:TransportBackoffSeconds, $_.Exception.Message) + } + if ($script:PolicyMode -eq 'server' -and (($nowUtc = (Get-Date).ToUniversalTime()) - $script:LastPolicyRefreshAt).TotalSeconds -ge $script:PolicyRefreshSeconds) { [void](Refresh-DlpPolicyFromServer) } @@ -1032,6 +1146,10 @@ while ($true) { policySource = $script:PolicySource policyVersion = $script:PolicyVersion policyChecksum = $script:PolicyChecksum + queueDepth = [int]$script:TransportMetrics.queueDepth + eventsEnqueued = [int]$script:TransportMetrics.eventsEnqueued + eventsFlushed = [int]$script:TransportMetrics.eventsFlushed + sendFailures = [int]$script:TransportMetrics.sendFailures } $script:LastSelfTestAt = $nowUtc } @@ -1209,5 +1327,10 @@ while ($true) { } } - Start-Sleep -Seconds $resolvedPollSeconds + if ($script:TransportBackoffSeconds -gt $resolvedPollSeconds) { + Start-Sleep -Seconds $script:TransportBackoffSeconds + } + else { + Start-Sleep -Seconds $resolvedPollSeconds + } } diff --git a/windows/file-operations-collector.ps1 b/windows/file-operations-collector.ps1 index c40b810..ad46eda 100644 --- a/windows/file-operations-collector.ps1 +++ b/windows/file-operations-collector.ps1 @@ -22,6 +22,14 @@ Add-Type -AssemblyName System.Net.Http $script:KnownBuckets = @{} $script:Hostname = $env:COMPUTERNAME $script:SessionId = [System.Diagnostics.Process]::GetCurrentProcess().SessionId +$script:TransportQueuePath = $null +$script:TransportQueueLockPath = $null +$script:TransportMetrics = @{ + eventsEnqueued = 0 + eventsFlushed = 0 + sendFailures = 0 + queueDepth = 0 +} # Настройка логирования $script:LogPath = $LogPath @@ -58,9 +66,11 @@ function Invoke-AwJsonPost { $reason = [string]$response.ReasonPhrase $body = $response.Content.ReadAsStringAsync().Result Write-FileCollectorLog ("POST failed: uri={0} status={1} reason={2} body={3}" -f $Uri, $status, $reason, $body) + throw "HTTP POST failed status=$status" } } catch { Write-FileCollectorLog "POST Error: $($_.Exception.Message)" + throw } finally { if ($null -ne $httpClient) { $httpClient.Dispose() @@ -68,6 +78,105 @@ function Invoke-AwJsonPost { } } +function Initialize-TransportQueue { + param( + [Parameter(Mandatory = $true)][string]$StateRoot + ) + $script:TransportQueuePath = Join-Path $StateRoot 'file-operations-queue.jsonl' + $script:TransportQueueLockPath = Join-Path $StateRoot 'file-operations-queue.lock' + if (-not (Test-Path -LiteralPath $script:TransportQueuePath)) { + New-Item -Path $script:TransportQueuePath -ItemType File -Force | Out-Null + } +} + +function Get-TransportQueueLock { + $tries = 0 + while ($tries -lt 50) { + try { + $fs = [System.IO.File]::Open($script:TransportQueueLockPath, [System.IO.FileMode]::OpenOrCreate, [System.IO.FileAccess]::ReadWrite, [System.IO.FileShare]::None) + return $fs + } + catch { + Start-Sleep -Milliseconds 50 + $tries++ + } + } + throw "Failed to acquire transport queue lock: $script:TransportQueueLockPath" +} + +function Add-TransportQueueItem { + param( + [Parameter(Mandatory = $true)][string]$Uri, + [Parameter(Mandatory = $true)][string]$Payload, + [string]$Kind = 'file_op' + ) + $lock = Get-TransportQueueLock + try { + $line = @{ + ts = (Get-Date).ToUniversalTime().ToString('o') + uri = $Uri + payload = $Payload + kind = $Kind + } | ConvertTo-Json -Compress + Add-Content -LiteralPath $script:TransportQueuePath -Value $line -Encoding UTF8 + $script:TransportMetrics.eventsEnqueued++ + } + finally { + $lock.Dispose() + } +} + +function Read-TransportQueueItems { + if (-not (Test-Path -LiteralPath $script:TransportQueuePath)) { return @() } + $items = @() + foreach ($line in @(Get-Content -LiteralPath $script:TransportQueuePath -ErrorAction SilentlyContinue)) { + if ([string]::IsNullOrWhiteSpace($line)) { continue } + try { $items += ($line | ConvertFrom-Json) } catch {} + } + return $items +} + +function Flush-TransportQueue { + param( + [int]$MaxItems = 100 + ) + if (-not (Test-Path -LiteralPath $script:TransportQueuePath)) { return } + $lock = Get-TransportQueueLock + try { + $items = Read-TransportQueueItems + $script:TransportMetrics.queueDepth = $items.Count + if ($items.Count -eq 0) { return } + + $left = New-Object System.Collections.Generic.List[object] + $sent = 0 + foreach ($item in $items) { + if ($sent -ge $MaxItems) { + $left.Add($item) + continue + } + try { + Invoke-AwJsonPost -Uri ([string]$item.uri) -Json ([string]$item.payload) + $sent++ + $script:TransportMetrics.eventsFlushed++ + } + catch { + $script:TransportMetrics.sendFailures++ + $left.Add($item) + } + } + foreach ($item in $items | Select-Object -Skip ($sent + $left.Count)) { + $left.Add($item) + } + + $lines = @($left | ForEach-Object { $_ | ConvertTo-Json -Compress }) + Set-Content -LiteralPath $script:TransportQueuePath -Value $lines -Encoding UTF8 + $script:TransportMetrics.queueDepth = $left.Count + } + finally { + $lock.Dispose() + } +} + function Ensure-Bucket { param( [string]$BucketId, @@ -140,7 +249,29 @@ function Send-FileOperationEvent { data = $data } | ConvertTo-Json -Depth 5 -Compress - Invoke-AwJsonPost -Uri "$($script:ApiBase)/buckets/$bucketId/heartbeat?pulsetime=15" -Json $payload + Add-TransportQueueItem -Uri "$($script:ApiBase)/buckets/$bucketId/heartbeat?pulsetime=15" -Payload $payload -Kind 'file_op' + Flush-TransportQueue -MaxItems 20 +} + +function Send-CollectorHealthEvent { + $bucketId = 'aw-file-operations_' + $script:Hostname + Ensure-Bucket -BucketId $bucketId -ClientName 'aw-file-operations' -BucketType 'aw.file.operation' + $payload = @{ + timestamp = (Get-Date).ToUniversalTime().ToString('yyyy-MM-ddTHH:mm:ss.fffZ') + duration = 0 + data = @{ + signalType = 'collector_health' + username = $env:USERNAME + hostname = $script:Hostname + sessionId = $script:SessionId + queueDepth = [int]$script:TransportMetrics.queueDepth + eventsEnqueued = [int]$script:TransportMetrics.eventsEnqueued + eventsFlushed = [int]$script:TransportMetrics.eventsFlushed + sendFailures = [int]$script:TransportMetrics.sendFailures + } + } | ConvertTo-Json -Depth 5 -Compress + Add-TransportQueueItem -Uri "$($script:ApiBase)/buckets/$bucketId/heartbeat?pulsetime=30" -Payload $payload -Kind 'health' + Flush-TransportQueue -MaxItems 50 } $config = Get-DeploymentConfig -Path $ConfigPath @@ -151,6 +282,8 @@ $scheme = if ($ServerScheme) { $ServerScheme } elseif ($config.server.scheme) { $hostName = if ($ServerHost) { $ServerHost } elseif ($config.server.host) { $config.server.host } else { 'localhost' } $port = if ($ServerPort) { $ServerPort } elseif ($config.server.port) { $config.server.port } else { 5600 } $script:ApiBase = "{0}://{1}:{2}/api/0" -f $scheme, $hostName, $port +$stateRoot = if ($config.paths -and $config.paths.stateRoot) { [string]$config.paths.stateRoot } else { 'C:\ProgramData\AWatch-rus' } +Initialize-TransportQueue -StateRoot $stateRoot $bucketId = 'aw-file-operations_' + $script:Hostname Ensure-Bucket -BucketId $bucketId -ClientName 'aw-file-operations' -BucketType 'aw.file.operation' @@ -206,7 +339,33 @@ foreach ($path in $resolvedPaths) { Write-FileCollectorLog "Collector started. Waiting for events..." try { + $lastHealth = [datetime]::UtcNow.AddMinutes(-5) + $backoffSeconds = 1 while ($true) { + try { + Flush-TransportQueue -MaxItems 100 + $backoffSeconds = 1 + } + catch { + $script:TransportMetrics.sendFailures++ + $backoffSeconds = [Math]::Min($backoffSeconds * 2, 60) + Write-FileCollectorLog ("Queue flush failed, backoff={0}s err={1}" -f $backoffSeconds, $_.Exception.Message) + } + + if ((New-TimeSpan -Start $lastHealth -End ([datetime]::UtcNow)).TotalSeconds -ge ([Math]::Max($PollSeconds * 3, 30))) { + try { + Send-CollectorHealthEvent + } + catch { + $script:TransportMetrics.sendFailures++ + } + $lastHealth = [datetime]::UtcNow + } + + if ($backoffSeconds -gt $PollSeconds) { + Start-Sleep -Seconds $backoffSeconds + continue + } Start-Sleep -Seconds $PollSeconds } }