diff --git a/docs/release-control/v6/internal/subsystems/agent-lifecycle.md b/docs/release-control/v6/internal/subsystems/agent-lifecycle.md index 0edcb7a54..1aa2704e1 100644 --- a/docs/release-control/v6/internal/subsystems/agent-lifecycle.md +++ b/docs/release-control/v6/internal/subsystems/agent-lifecycle.md @@ -1418,6 +1418,11 @@ the intentionally sparse public response. 4. Keep shared agent-side TLS identity fail-closed across `cmd/pulse-agent/main.go`, `internal/hostagent/`, `internal/agentupdate/`, `internal/remoteconfig/client.go`, and `internal/agenttls/config.go`. Self-signed deployments may use a canonical pinned Pulse server certificate fingerprint, but lifecycle transport must route that pin through reporting, enrollment, command websocket, remote-config, and self-update clients instead of widening `PULSE_INSECURE_SKIP_VERIFY` into a blanket MITM carve-out. A configured custom CA bundle is part of that same trust boundary: if the bundle is unreadable or invalid, lifecycle transport must refuse the connection path rather than silently downgrading back to system roots. 5. Keep release-grade updater trust fail-closed across `internal/agentupdate/`, `internal/dockeragent/`, and the shared `internal/api/unified_agent.go` download helpers. When release builds embed trusted update signing keys, published agent binaries and installer assets must carry detached `.sig` plus `.sshsig` sidecars; updater/runtime paths must require `X-Signature-Ed25519` in addition to `X-Checksum-Sha256`, and installer-owned download flows must require the matching base64-encoded `X-Signature-SSHSIG`, instead of silently downgrading to checksum-only trust. 6. Keep shared `internal/api/` helper edits isolated from agent lifecycle semantics: Patrol-specific status transport or alert-trigger wiring changes in shared handlers must not bleed into auto-register, installer, or fleet-control behavior unless this contract moves in the same slice. + Audit-log list, filter, pagination, and storage-error contract changes in + `internal/api/activity_audit_handlers.go` remain security/API behavior. + They must not change agent enrollment, reporting, command transport, + lifecycle event creation, or fleet-control authority merely because the + handler shares the `internal/api/` package. SSO browser-session display labels in shared auth/session helpers are likewise API/security presentation state, not lifecycle identity. Agent enrollment, installer reruns, update recovery, token binding, and fleet diff --git a/docs/release-control/v6/internal/subsystems/api-contracts.md b/docs/release-control/v6/internal/subsystems/api-contracts.md index fd8add889..53a90ec86 100644 --- a/docs/release-control/v6/internal/subsystems/api-contracts.md +++ b/docs/release-control/v6/internal/subsystems/api-contracts.md @@ -49,6 +49,7 @@ singular `availability` field remains only a compatibility projection. 4. `internal/api/alerts.go` 4a. `internal/api/attention_handlers.go` 5. `internal/api/activity_audit_handlers.go` + 5a. `pkg/extensions/audit_admin.go` 6. `internal/api/actions.go` 5a. `internal/api/action_executor.go` 5b. `internal/api/docker_container_action_executor.go` @@ -1757,9 +1758,10 @@ payload shape change when the portal presents compact client rows. 88. `pkg/aicontracts/investigation.go` shared with `ai-runtime`: the public Patrol investigation record and finding contract is both an AI runtime handoff boundary and a canonical API payload contract for Patrol, Assistant, unified findings, persistence, and audit surfaces. 89. `pkg/aicontracts/orchestrator_deps.go` shared with `ai-runtime`: the public investigation orchestrator dependency contract is both an AI runtime handoff boundary and a canonical API payload contract for Assistant and Patrol tool-call history. 90. `pkg/extensions/ai_autofix.go` shared with `ai-runtime`: the enterprise auto-fix extension dependency seam is both an AI runtime approved-action boundary and a canonical API extension contract over Assistant and Patrol execution dependencies. -91. `pulse-mobile:config/mobile-api-surface.json` shared with `relay-runtime`: the Pulse Mobile consumer minimum and released-line probe inventory are both API compatibility and relay runtime boundaries. -92. `pulse-mobile:src/generated/coreCompatibility.ts` shared with `relay-runtime`: the generated mobile route, pairing, and push projection is both an API consumer contract and relay runtime boundary. -91. `scripts/generate-pulse-intelligence-docs.go` shared with `ai-runtime`: the Pulse Intelligence manifest docs generator is both an AI runtime docs/onboarding projection and a canonical API contract projection over the agent capabilities manifest and Pulse MCP surface tool contract. +91. `pkg/extensions/audit_admin.go` shared with `security-privacy`: the enterprise audit endpoint and canonical store configuration seam is both a security persistence trust boundary and a canonical API extension contract. +92. `pulse-mobile:config/mobile-api-surface.json` shared with `relay-runtime`: the Pulse Mobile consumer minimum and released-line probe inventory are both API compatibility and relay runtime boundaries. +93. `pulse-mobile:src/generated/coreCompatibility.ts` shared with `relay-runtime`: the generated mobile route, pairing, and push projection is both an API consumer contract and relay runtime boundary. +94. `scripts/generate-pulse-intelligence-docs.go` shared with `ai-runtime`: the Pulse Intelligence manifest docs generator is both an AI runtime docs/onboarding projection and a canonical API contract projection over the agent capabilities manifest and Pulse MCP surface tool contract. Update-plan responses own the structured readiness verdict for server updater capability, rollback support, agent continuity, v5 agent migration transport security, and agent reporting token scope. That verdict is part @@ -1988,8 +1990,17 @@ a new API state machine, queue contract, or verification-accounting field. semantics from `pkg/audit`: transient store pressure returns `503` `audit_store_busy` with `Retry-After`, unavailable or corrupt audit storage returns `503` `audit_store_unavailable`, and unrelated query failures remain - `query_failed`. List pagination must stay bounded so audit history reads - cannot become unbounded table scans through API parameters. + `query_failed`. List response sizes must stay bounded while valid + non-negative offsets remain able to traverse the retained audit history. + Invalid + booleans, timestamps, limits, and offsets fail with a stable `400` + validation response instead of being ignored. Successful list responses + always encode `events` as an array, including an empty history, and obtain + `events` plus `total` from one storage snapshot with stable + `(timestamp, id)` ordering. Enterprise endpoint binders may extend export, + summary, and webhook behavior, but list and verification replacements must + delegate these canonical handlers so Pro cannot discard or remap the + storage contract. Unified action planning is part of that same API-first action contract: `POST /api/actions/plan` must route through `internal/api/actions.go`, `internal/actionlifecycle/service.go`, diff --git a/docs/release-control/v6/internal/subsystems/frontend-primitives.md b/docs/release-control/v6/internal/subsystems/frontend-primitives.md index f08a41b84..df7ea687a 100644 --- a/docs/release-control/v6/internal/subsystems/frontend-primitives.md +++ b/docs/release-control/v6/internal/subsystems/frontend-primitives.md @@ -3834,6 +3834,12 @@ through `apiErrorFromResponse` and render customer-facing copy from `frontend-modern/src/utils/auditLogPresentation.ts`; the settings shell may own refresh and pagination state, but it must not show raw `Internal Server Error` strings or unbounded page sizes as local hook behavior. +Audit-log page loads are latest-request-wins. Changing the page size atomically +resets the offset and starts one replacement request, and any superseded +request is aborted or ignored so an older response cannot overwrite the new +page. A failed load clears previously rendered events and totals, and a +successful payload must contain an event array rather than treating `null` or +an absent list as an empty audit history. That shared filter-option primitive is also the canonical owner for default `All ` option wording wherever a product surface exposes filter selects or segmented filter choices. Workloads filters, storage source diff --git a/docs/release-control/v6/internal/subsystems/registry.json b/docs/release-control/v6/internal/subsystems/registry.json index 6b2752320..fc865525f 100644 --- a/docs/release-control/v6/internal/subsystems/registry.json +++ b/docs/release-control/v6/internal/subsystems/registry.json @@ -1101,6 +1101,14 @@ "api-contracts" ] }, + { + "path": "pkg/extensions/audit_admin.go", + "rationale": "the enterprise audit endpoint and canonical store configuration seam is both a security persistence trust boundary and a canonical API extension contract", + "subsystems": [ + "api-contracts", + "security-privacy" + ] + }, { "path": "pulse-mobile:config/mobile-api-surface.json", "rationale": "the Pulse Mobile consumer minimum and released-line probe inventory are both API compatibility and relay runtime boundaries", @@ -2485,6 +2493,7 @@ "pkg/aicontracts/investigation.go", "pkg/aicontracts/orchestrator_deps.go", "pkg/extensions/ai_autofix.go", + "pkg/extensions/audit_admin.go", "pkg/pulsecli/actions.go", "pkg/pulsecli/api_client.go", "pkg/pulsecli/fleet.go", @@ -2844,7 +2853,9 @@ "match_prefixes": [ "internal/api/" ], - "match_files": [], + "match_files": [ + "pkg/extensions/audit_admin.go" + ], "allow_same_subsystem_tests": false, "test_prefixes": [ "frontend-modern/src/api/__tests__/" @@ -2853,11 +2864,13 @@ "frontend-modern/src/types/api.ts", "internal/api/ai_handlers_more_test.go", "internal/api/ai_handlers_patrol_actions_additional_test.go", + "internal/api/audit_handlers_test.go", "internal/api/contract_test.go", "internal/api/docker_agents_report_size_test.go", "internal/api/host_agent_removal_lifecycle_integration_test.go", "internal/api/metadata_handlers_test.go", - "internal/api/patrol_autopilot_test.go" + "internal/api/patrol_autopilot_test.go", + "pulse-enterprise:test/extensions_contract_test.go" ] }, { @@ -4779,7 +4792,8 @@ "test_prefixes": [], "exact_files": [ "frontend-modern/src/components/Settings/__tests__/dataHandlingPanelModel.test.ts", - "frontend-modern/src/components/Settings/__tests__/settingsArchitecture.test.ts" + "frontend-modern/src/components/Settings/__tests__/settingsArchitecture.test.ts", + "frontend-modern/src/components/Settings/__tests__/useAuditLogPanelState.test.tsx" ] }, { @@ -6552,6 +6566,7 @@ "pkg/audit/sqlite_logger.go", "pkg/auth/rbac.go", "pkg/auth/sqlite_manager.go", + "pkg/extensions/audit_admin.go", "pkg/server/server.go", "pkg/server/telemetry_pulse_intelligence.go", "pkg/tlsutil/fingerprint.go", @@ -6677,7 +6692,9 @@ "match_prefixes": [ "pkg/audit/" ], - "match_files": [], + "match_files": [ + "pkg/extensions/audit_admin.go" + ], "allow_same_subsystem_tests": true, "test_prefixes": [], "exact_files": [ @@ -6690,7 +6707,9 @@ "pkg/audit/sqlite_logger_test.go", "pkg/audit/tenant_logger_manager_test.go", "pkg/audit/webhook_delivery_test.go", - "pkg/audit/webhook_validation_test.go" + "pkg/audit/webhook_validation_test.go", + "pulse-enterprise:internal/enterpriseruntime/runtime_test.go", + "pulse-enterprise:test/extensions_contract_test.go" ] }, { @@ -6779,7 +6798,8 @@ "internal/api/rbac_tenant_provider_test.go", "internal/api/security_regression_test.go", "pkg/auth/rbac_manager_test.go", - "pkg/auth/sqlite_manager_test.go" + "pkg/auth/sqlite_manager_test.go", + "pkg/server/server_test.go" ] } ], diff --git a/docs/release-control/v6/internal/subsystems/security-privacy.md b/docs/release-control/v6/internal/subsystems/security-privacy.md index e8db08c0f..077c74d21 100644 --- a/docs/release-control/v6/internal/subsystems/security-privacy.md +++ b/docs/release-control/v6/internal/subsystems/security-privacy.md @@ -64,14 +64,17 @@ controls as normal product settings. 36. `pkg/audit/audit.go` 37. `pkg/audit/async_logger.go` 38. `pkg/audit/sqlite_logger.go` -39. `scripts/telemetry_adoption_report.py` -40. `frontend-modern/src/components/Settings/DataHandlingPanel.tsx` -41. `frontend-modern/src/components/Settings/dataHandlingPanelModel.ts` -42. `internal/api/agent_exec_token_binding.go` -43. `internal/logging/logging.go` -44. `pkg/auth/rbac.go` -45. `pkg/auth/sqlite_manager.go` -46. `pkg/server/server.go` +39. `pkg/audit/signer.go` +40. `pkg/audit/sqlite_factory.go` +41. `pkg/extensions/audit_admin.go` +42. `scripts/telemetry_adoption_report.py` +43. `frontend-modern/src/components/Settings/DataHandlingPanel.tsx` +44. `frontend-modern/src/components/Settings/dataHandlingPanelModel.ts` +45. `internal/api/agent_exec_token_binding.go` +46. `internal/logging/logging.go` +47. `pkg/auth/rbac.go` +48. `pkg/auth/sqlite_manager.go` +49. `pkg/server/server.go` ## Shared Boundaries @@ -166,6 +169,8 @@ with missing, unknown, or unrelated scopes fail closed. storing a brand setting never becomes a free branding bypass. 16. `internal/cloudcp/auth/magiclink.go` shared with `cloud-paid`: control-plane magic-link HMAC handling is both a Pulse Cloud account-access boundary and a security/privacy token-secrecy boundary. 17. `internal/cloudcp/auth/magiclink_store.go` shared with `cloud-paid`: control-plane magic-link persistence is both a Pulse Cloud account-access boundary and a security/privacy storage-hardening boundary. +18. `pkg/extensions/audit_admin.go` shared with `api-contracts`: the enterprise audit endpoint and canonical store configuration seam is both a security persistence trust boundary and a canonical API extension contract. + ## Extension Points Catalog edits in `frontend-modern/src/i18n/` that add or promote Patrol-trigger @@ -755,6 +760,25 @@ were previously written into `audit_events.timestamp`, including Unix seconds, SQLite datetime values, and Go wall-clock strings carrying a monotonic `m=+...` suffix, so valid historical audit rows cannot make `/api/audit` return `query_failed`. +The runtime now has one canonical persistent audit owner across Community and +Pro. Enterprise startup may provide directory, signing-key, and explicit +retention settings through `ResolveAuditStoreConfig`, but it must not open a +second database, create a competing signing key, replace the canonical list or +verify handlers, or make storage behavior depend on install history. Store +startup must transactionally normalize legacy `DATETIME`, textual, and mixed +timestamp rows into the schema-v2 integer contract before serving reads. A +malformed legacy row must roll the migration back and surface a sanitized +`audit_store_unavailable` diagnostic rather than deleting, skipping, or +silently returning an empty history. +Canonical reads use stable `(timestamp DESC, id DESC)` ordering, and a paged +read must derive rows plus total from one SQLite snapshot. Async wrapping must +retain persistent-reader and signature-verification capabilities. Retention +`0` is an explicit keep-forever setting, persisted retention is restored at +startup, the configured cleanup cadence is preserved across the enterprise +configuration seam, and cleanup must tolerate concurrent writers without +exposing partial results. Existing core and Pro signing keys remain valid through upgrade, and +legacy Pro signature encodings remain verifiable while all new writes use the +canonical signed event representation. That shared token-management boundary now also includes `frontend-modern/src/utils/apiTokenPresentation.ts`, so API-token load, generate, and revoke errors stay on one governed customer-facing wording path diff --git a/docs/release-control/v6/internal/subsystems/storage-recovery.md b/docs/release-control/v6/internal/subsystems/storage-recovery.md index 78ae70df7..19c2a681f 100644 --- a/docs/release-control/v6/internal/subsystems/storage-recovery.md +++ b/docs/release-control/v6/internal/subsystems/storage-recovery.md @@ -2192,6 +2192,17 @@ request context, while invalid credentials fail before any action audit write. It does not add recovery storage, alter `ActionAuditRecord`, or let response headers participate in approval identity. +The security audit database remains owned by `security-privacy`, but its +upgrade and recovery behavior is a shared storage boundary. Startup must +transactionally rebuild legacy v5 and early-v6 `audit_events` schemas into the +canonical integer-timestamp schema without partial commits, and malformed or +corrupt storage must fail closed without replacing history with an empty +result. Offline backup and restore must preserve both `audit.db` and its +signing key, while a restored store must retain ordering, tenant filtering, +signature verification, and explicit retention behavior. Concurrent writes +and SQLite busy conditions may be retried within bounded limits, but they must +never weaken snapshot consistency or cause destructive migration fallback. + The Patrol finding lifecycle endpoints advertised in the agent capabilities manifest (`acknowledge_finding`, `snooze_finding`, `dismiss_finding`, `resolve_finding`) follow the same storage boundary: diff --git a/frontend-modern/src/components/Settings/__tests__/useAuditLogPanelState.test.tsx b/frontend-modern/src/components/Settings/__tests__/useAuditLogPanelState.test.tsx new file mode 100644 index 000000000..3d6e5ee46 --- /dev/null +++ b/frontend-modern/src/components/Settings/__tests__/useAuditLogPanelState.test.tsx @@ -0,0 +1,180 @@ +import { renderHook, waitFor } from '@solidjs/testing-library'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; + +import { apiErrorFromResponse, apiFetch } from '@/utils/apiClient'; +import { useAuditLogPanelState } from '../useAuditLogPanelState'; + +const navigate = vi.fn(); + +vi.mock('@solidjs/router', () => ({ + useLocation: () => ({ pathname: '/settings/security', search: '' }), + useNavigate: () => navigate, +})); + +vi.mock('@/stores/license', () => ({ + getRuntimeCapabilityBlock: () => undefined, + hasFeature: (feature: string) => feature === 'audit_logging', + loadRuntimeCapabilities: vi.fn().mockResolvedValue(undefined), + runtimeCapabilitiesLoaded: () => true, +})); + +vi.mock('@/stores/licenseCommercial', () => ({ + getUpgradeActionDestination: () => '/upgrade', +})); + +vi.mock('@/stores/sessionPresentationPolicy', () => ({ + presentationPolicyHidesUpgradePrompts: () => false, +})); + +vi.mock('@/utils/upgradeNavigation', () => ({ + resolveUpgradeDestination: () => '/upgrade', +})); + +vi.mock('@/utils/toast', () => ({ + showSuccess: vi.fn(), + showToast: vi.fn(), + showWarning: vi.fn(), +})); + +vi.mock('@/utils/apiClient', async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + apiFetch: vi.fn(), + apiErrorFromResponse: vi.fn(), + }; +}); + +type Deferred = { + promise: Promise; + resolve: (value: T) => void; +}; + +const deferred = (): Deferred => { + let resolve!: (value: T) => void; + const promise = new Promise((next) => { + resolve = next; + }); + return { promise, resolve }; +}; + +const auditResponse = (ids: string[], total = ids.length): Response => + new Response( + JSON.stringify({ + events: ids.map((id) => ({ + id, + timestamp: '2026-07-24T09:00:00Z', + event: 'login', + user: 'operator', + ip: '127.0.0.1', + path: '/api/auth', + success: true, + details: id, + })), + total, + persistentLogging: true, + }), + { status: 200, headers: { 'Content-Type': 'application/json' } }, + ); + +describe('useAuditLogPanelState audit request lifecycle', () => { + beforeEach(() => { + window.localStorage.clear(); + navigate.mockReset(); + vi.mocked(apiFetch).mockReset(); + vi.mocked(apiErrorFromResponse).mockReset(); + }); + + afterEach(() => { + vi.restoreAllMocks(); + }); + + it('keeps the newest response when an older request resolves last', async () => { + const first = deferred(); + const second = deferred(); + vi.mocked(apiFetch) + .mockImplementationOnce(() => first.promise) + .mockImplementationOnce(() => second.promise); + + const { result, cleanup } = renderHook(() => useAuditLogPanelState()); + await waitFor(() => expect(apiFetch).toHaveBeenCalledTimes(1)); + + result.refresh(); + await waitFor(() => expect(apiFetch).toHaveBeenCalledTimes(2)); + second.resolve(auditResponse(['newest'])); + await waitFor(() => expect(result.events().map((event) => event.id)).toEqual(['newest'])); + + first.resolve(auditResponse(['stale'])); + await Promise.resolve(); + expect(result.events().map((event) => event.id)).toEqual(['newest']); + expect(result.loading()).toBe(false); + cleanup(); + }); + + it('clears prior rows when a refresh fails', async () => { + vi.mocked(apiFetch) + .mockResolvedValueOnce(auditResponse(['existing'])) + .mockResolvedValueOnce( + new Response(JSON.stringify({ code: 'audit_store_busy' }), { status: 503 }), + ); + vi.mocked(apiErrorFromResponse).mockResolvedValue( + Object.assign(new Error('Audit log storage is temporarily busy.'), { + code: 'audit_store_busy', + }), + ); + + const { result, cleanup } = renderHook(() => useAuditLogPanelState()); + await waitFor(() => expect(result.events()).toHaveLength(1)); + result.refresh(); + await waitFor(() => expect(result.error()).toContain('storage is busy')); + + expect(result.events()).toEqual([]); + expect(result.totalEvents()).toBe(0); + expect(result.isPersistent()).toBe(false); + cleanup(); + }); + + it('resets and refetches atomically when page size changes', async () => { + vi.mocked(apiFetch) + .mockResolvedValueOnce( + auditResponse( + Array.from({ length: 100 }, (_, index) => `old-${index}`), + 250, + ), + ) + .mockResolvedValueOnce( + auditResponse( + Array.from({ length: 25 }, (_, index) => `new-${index}`), + 250, + ), + ); + + const { result, cleanup } = renderHook(() => useAuditLogPanelState()); + await waitFor(() => expect(result.events()).toHaveLength(100)); + result.setPageSize(25); + await waitFor(() => expect(result.events()).toHaveLength(25)); + + const lastURL = vi.mocked(apiFetch).mock.calls.at(-1)?.[0]; + expect(String(lastURL)).toContain('limit=25'); + expect(String(lastURL)).toContain('offset=0'); + expect(result.pageSize()).toBe(25); + expect(result.pageNumber()).toBe(1); + expect(result.pageRangeText()).toBe('Showing 1-25 of 250'); + cleanup(); + }); + + it('rejects null event collections instead of masking the API contract violation', async () => { + vi.mocked(apiFetch).mockResolvedValueOnce( + new Response(JSON.stringify({ events: null, total: 0, persistentLogging: true }), { + status: 200, + headers: { 'Content-Type': 'application/json' }, + }), + ); + + const { result, cleanup } = renderHook(() => useAuditLogPanelState()); + await waitFor(() => expect(result.error()).toBe('Audit log returned an invalid response')); + expect(result.events()).toEqual([]); + expect(result.isPersistent()).toBe(false); + cleanup(); + }); +}); diff --git a/frontend-modern/src/components/Settings/useAuditLogPanelState.ts b/frontend-modern/src/components/Settings/useAuditLogPanelState.ts index 9b2f6d123..79319379c 100644 --- a/frontend-modern/src/components/Settings/useAuditLogPanelState.ts +++ b/frontend-modern/src/components/Settings/useAuditLogPanelState.ts @@ -128,7 +128,11 @@ export const useAuditLogPanelState = () => { ); const pageSize = () => normalizeAuditPageSize(storedPageSize()); const setPageSize = (value: number): void => { - setStoredPageSize(normalizeAuditPageSize(value)); + const normalized = normalizeAuditPageSize(value); + setStoredPageSize(normalized); + setPageOffset(0); + setPageInput(''); + void fetchAuditEvents({ limit: normalized, offset: 0 }); }; const [pageOffset, setPageOffset] = createLocalStorageNumberSignal( STORAGE_KEYS.AUDIT_PAGE_OFFSET, @@ -154,6 +158,8 @@ export const useAuditLogPanelState = () => { const [verifyControllers, setVerifyControllers] = createSignal>( {}, ); + let auditRequestController: AbortController | null = null; + let auditRequestGeneration = 0; const auditLoggingEnabled = createMemo( () => runtimeCapabilitiesLoaded() && hasFeature('audit_logging'), @@ -176,6 +182,10 @@ export const useAuditLogPanelState = () => { const upgradeActionLabel = () => (paidRuntimeRequired() ? 'Download Pulse Pro' : 'View plans'); const fetchAuditEvents = async (options?: { limit?: number; offset?: number }) => { + auditRequestController?.abort(); + auditRequestController = null; + const requestGeneration = ++auditRequestGeneration; + if (!auditLoggingEnabled()) { setEvents([]); setTotalEvents(0); @@ -187,6 +197,8 @@ export const useAuditLogPanelState = () => { const limit = options?.limit ?? pageSize(); const offset = options?.offset ?? pageOffset(); + const controller = new AbortController(); + auditRequestController = controller; setLoading(true); setError(null); @@ -204,7 +216,10 @@ export const useAuditLogPanelState = () => { if (successFilter() === 'success') params.set('success', 'true'); if (successFilter() === 'failed') params.set('success', 'false'); - const response = await apiFetch(`/api/audit?${params.toString()}`); + const response = await apiFetch(`/api/audit?${params.toString()}`, { + signal: controller.signal, + }); + if (requestGeneration !== auditRequestGeneration) return; if (response.status === 402) { setEvents([]); setTotalEvents(0); @@ -220,28 +235,43 @@ export const useAuditLogPanelState = () => { } const data: AuditResponse = await response.json(); - setEvents(data.events || []); - setIsPersistent(data.persistentLogging); - setTotalEvents(data.total ?? 0); + if (requestGeneration !== auditRequestGeneration) return; + if ( + !Array.isArray(data.events) || + typeof data.total !== 'number' || + typeof data.persistentLogging !== 'boolean' + ) { + throw new Error('Audit log returned an invalid response'); + } if (data.total && offset >= data.total) { const maxOffset = Math.max(0, Math.floor((data.total - 1) / limit) * limit); if (maxOffset !== offset) { setPageOffset(maxOffset); - void fetchAuditEvents({ limit, offset: maxOffset }); + await fetchAuditEvents({ limit, offset: maxOffset }); return; } } + setEvents(data.events); + setIsPersistent(data.persistentLogging); + setTotalEvents(data.total); + if (data.persistentLogging && autoVerifyEnabled()) { const verificationLimit = autoVerifyLimit(); if (verificationLimit <= 0) return; setTimeout(() => { - if (!isMounted()) return; + if (!isMounted() || requestGeneration !== auditRequestGeneration) return; void verifyAllEvents({ limit: verificationLimit, showToast: false }); }, 0); } } catch (err) { + if ( + requestGeneration !== auditRequestGeneration || + (err as { name?: string })?.name === 'AbortError' + ) { + return; + } const message = getAuditLogFetchErrorMessage(err); if (typeof message === 'string' && /feature not included in license/i.test(message)) { setEvents([]); @@ -250,10 +280,16 @@ export const useAuditLogPanelState = () => { setError(null); return; } + setEvents([]); + setTotalEvents(0); + setIsPersistent(false); setError(message); showWarning('Audit Log Error', message); } finally { - setLoading(false); + if (requestGeneration === auditRequestGeneration) { + auditRequestController = null; + setLoading(false); + } } }; @@ -509,10 +545,12 @@ export const useAuditLogPanelState = () => { const clearFilters = () => { const hadFilters = activeFilterCount() > 0; - setEventFilter(''); - setUserFilter(''); - setSuccessFilter('all'); - setVerificationFilter('all'); + updateSearchParam((params) => { + params.delete('event'); + params.delete('user'); + params.delete('success'); + params.delete('verification'); + }); if (hadFilters) { showSuccess('Audit filters cleared'); } @@ -610,6 +648,9 @@ export const useAuditLogPanelState = () => { return; } if (!hasFeature('audit_logging')) { + auditRequestController?.abort(); + auditRequestController = null; + auditRequestGeneration += 1; setEvents([]); setTotalEvents(0); setIsPersistent(false); @@ -672,6 +713,9 @@ export const useAuditLogPanelState = () => { onCleanup(() => { setIsMounted(false); setCancelVerifyAll(true); + auditRequestController?.abort(); + auditRequestController = null; + auditRequestGeneration += 1; if (userFilterDebounceHandle !== null) { window.clearTimeout(userFilterDebounceHandle); userFilterDebounceHandle = null; diff --git a/internal/api/activity_audit_handlers.go b/internal/api/activity_audit_handlers.go index 4ec9647d9..eb100f7b6 100644 --- a/internal/api/activity_audit_handlers.go +++ b/internal/api/activity_audit_handlers.go @@ -67,64 +67,20 @@ func (h *AuditHandlers) HandleListAuditEvents(w http.ResponseWriter, r *http.Req return } - query := r.URL.Query() - - filter := audit.QueryFilter{ - EventType: query.Get("event"), - User: query.Get("user"), + filter, validationErr := auditEventListFilterFromRequest(r) + if validationErr != nil { + writeErrorResponse(w, http.StatusBadRequest, validationErr.code, validationErr.message, nil) + return } - filter.Limit = defaultAuditEventListLimit - if limitStr := query.Get("limit"); limitStr != "" { - if limit, err := strconv.Atoi(limitStr); err == nil && limit > 0 { - filter.Limit = min(limit, maxAuditEventListLimit) - } - } - - // Parse offset - if offsetStr := query.Get("offset"); offsetStr != "" { - if offset, err := strconv.Atoi(offsetStr); err == nil && offset >= 0 { - filter.Offset = offset - } - } - - // Parse startTime - if startStr := query.Get("startTime"); startStr != "" { - if t, err := time.Parse(time.RFC3339, startStr); err == nil { - filter.StartTime = &t - } - } - - // Parse endTime - if endStr := query.Get("endTime"); endStr != "" { - if t, err := time.Parse(time.RFC3339, endStr); err == nil { - filter.EndTime = &t - } - } - - // Parse success - if successStr := query.Get("success"); successStr != "" { - success := successStr == "true" - filter.Success = &success - } - - // Query events from the current logger - events, err := logger.Query(filter) + events, totalCount, err := audit.QueryPage(logger, filter) if err != nil { log.Error().Err(err).Str("org_id", orgID).Msg("audit list: query failed") writeAuditReadErrorResponse(w, err, "Failed to query audit events") return } - - countFilter := filter - countFilter.Limit = 0 - countFilter.Offset = 0 - - totalCount, err := logger.Count(countFilter) - if err != nil { - log.Error().Err(err).Str("org_id", orgID).Msg("audit list: count failed") - writeAuditReadErrorResponse(w, err, "Failed to count audit events") - return + if events == nil { + events = make([]audit.Event, 0) } response := map[string]interface{}{ @@ -137,6 +93,80 @@ func (h *AuditHandlers) HandleListAuditEvents(w http.ResponseWriter, r *http.Req json.NewEncoder(w).Encode(response) } +type auditQueryValidationError struct { + code string + message string +} + +func auditEventListFilterFromRequest(r *http.Request) (audit.QueryFilter, *auditQueryValidationError) { + query := r.URL.Query() + filter := audit.QueryFilter{ + EventType: query.Get("event"), + User: query.Get("user"), + Limit: defaultAuditEventListLimit, + } + + if query.Has("limit") { + limit, err := strconv.Atoi(query.Get("limit")) + if err != nil || limit <= 0 { + return audit.QueryFilter{}, &auditQueryValidationError{ + code: "invalid_limit", + message: "Invalid limit; expected a positive integer", + } + } + filter.Limit = min(limit, maxAuditEventListLimit) + } + if query.Has("offset") { + offset, err := strconv.Atoi(query.Get("offset")) + if err != nil || offset < 0 { + return audit.QueryFilter{}, &auditQueryValidationError{ + code: "invalid_offset", + message: "Invalid offset; expected a non-negative integer", + } + } + filter.Offset = offset + } + if query.Has("startTime") { + start, err := time.Parse(time.RFC3339, query.Get("startTime")) + if err != nil { + return audit.QueryFilter{}, &auditQueryValidationError{ + code: "invalid_start_time", + message: "Invalid startTime; expected an RFC3339 timestamp", + } + } + filter.StartTime = &start + } + if query.Has("endTime") { + end, err := time.Parse(time.RFC3339, query.Get("endTime")) + if err != nil { + return audit.QueryFilter{}, &auditQueryValidationError{ + code: "invalid_end_time", + message: "Invalid endTime; expected an RFC3339 timestamp", + } + } + filter.EndTime = &end + } + if filter.StartTime != nil && filter.EndTime != nil && !filter.StartTime.Before(*filter.EndTime) { + return audit.QueryFilter{}, &auditQueryValidationError{ + code: "invalid_time_range", + message: "startTime must be before endTime", + } + } + if query.Has("success") { + rawSuccess := query.Get("success") + if rawSuccess != "true" && rawSuccess != "false" { + return audit.QueryFilter{}, &auditQueryValidationError{ + code: "invalid_success", + message: "Invalid success; expected a boolean", + } + } + success := rawSuccess == "true" + filter.Success = &success + } + + return filter, nil +} + // HandleVerifyAuditEvent handles GET /api/audit/{id}/verify func (h *AuditHandlers) HandleVerifyAuditEvent(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { diff --git a/internal/api/audit_handlers_test.go b/internal/api/audit_handlers_test.go index 55dc0189b..117f882d0 100644 --- a/internal/api/audit_handlers_test.go +++ b/internal/api/audit_handlers_test.go @@ -306,11 +306,11 @@ func TestHandleListAuditEvents(t *testing.T) { t.Errorf("expected total 2, got %v", resp["total"]) } - // Test parse error for startTime/endTime + // Invalid time filters fail closed instead of silently widening the query. req = httptest.NewRequest(http.MethodGet, "/api/audit?startTime=invalid&endTime=invalid", nil) rec = httptest.NewRecorder() handler.HandleListAuditEvents(rec, req) - assert.Equal(t, http.StatusOK, rec.Code) // It just ignores invalid times + assert.Equal(t, http.StatusBadRequest, rec.Code) // Test method not allowed req = httptest.NewRequest(http.MethodPost, "/api/audit", nil) @@ -532,7 +532,7 @@ func TestHandleListAuditEvents_PaginationDefaults(t *testing.T) { setAuditLogger(t, logger) handler := NewAuditHandlers() - req := httptest.NewRequest(http.MethodGet, "/api/audit?limit=0&offset=-4", nil) + req := httptest.NewRequest(http.MethodGet, "/api/audit", nil) rec := httptest.NewRecorder() handler.HandleListAuditEvents(rec, req) @@ -546,6 +546,15 @@ func TestHandleListAuditEvents_PaginationDefaults(t *testing.T) { t.Fatalf("offset = %d, want 0", logger.lastQuery.Offset) } + for _, query := range []string{"?limit=0", "?offset=-4"} { + req = httptest.NewRequest(http.MethodGet, "/api/audit"+query, nil) + rec = httptest.NewRecorder() + handler.HandleListAuditEvents(rec, req) + if rec.Code != http.StatusBadRequest { + t.Fatalf("query %q status = %d, want 400", query, rec.Code) + } + } + req = httptest.NewRequest(http.MethodGet, "/api/audit?limit=5000&offset=10", nil) rec = httptest.NewRecorder() handler.HandleListAuditEvents(rec, req) @@ -564,6 +573,103 @@ func TestHandleListAuditEvents_PaginationDefaults(t *testing.T) { } } +func TestHandleListAuditEvents_RejectsMalformedFilters(t *testing.T) { + setAuditLogger(t, &testAuditLogger{}) + handler := NewAuditHandlers() + tests := []struct { + query string + code string + }{ + {query: "limit=abc", code: "invalid_limit"}, + {query: "limit=0", code: "invalid_limit"}, + {query: "offset=-1", code: "invalid_offset"}, + {query: "startTime=not-a-time", code: "invalid_start_time"}, + {query: "endTime=not-a-time", code: "invalid_end_time"}, + { + query: "startTime=2026-07-05T12%3A00%3A00Z&endTime=2026-07-05T11%3A00%3A00Z", + code: "invalid_time_range", + }, + {query: "success=maybe", code: "invalid_success"}, + } + for _, tt := range tests { + t.Run(tt.code, func(t *testing.T) { + req := httptest.NewRequest(http.MethodGet, "/api/audit?"+tt.query, nil) + rec := httptest.NewRecorder() + handler.HandleListAuditEvents(rec, req) + if rec.Code != http.StatusBadRequest { + t.Fatalf("status = %d, want 400", rec.Code) + } + var payload struct { + Code string `json:"code"` + } + if err := json.Unmarshal(rec.Body.Bytes(), &payload); err != nil { + t.Fatalf("decode error: %v", err) + } + if payload.Code != tt.code { + t.Fatalf("code = %q, want %q", payload.Code, tt.code) + } + }) + } +} + +func TestHandleListAuditEvents_EmptyEventsIsArray(t *testing.T) { + setAuditLogger(t, &testAuditLogger{events: nil}) + req := httptest.NewRequest(http.MethodGet, "/api/audit", nil) + rec := httptest.NewRecorder() + NewAuditHandlers().HandleListAuditEvents(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("status = %d, want 200", rec.Code) + } + var payload struct { + Events json.RawMessage `json:"events"` + } + if err := json.Unmarshal(rec.Body.Bytes(), &payload); err != nil { + t.Fatalf("decode response: %v", err) + } + if string(payload.Events) != "[]" { + t.Fatalf("events = %s, want []", payload.Events) + } +} + +func TestHandleListAuditEvents_TenantDataIsolation(t *testing.T) { + previousManager := GetTenantAuditManager() + manager := audit.NewTenantLoggerManager(t.TempDir(), &audit.SQLiteLoggerFactory{ + RetentionDays: 0, + RetentionConfigured: true, + }) + SetTenantAuditManager(manager) + t.Cleanup(func() { + manager.Close() + SetTenantAuditManager(previousManager) + }) + + for orgID, eventID := range map[string]string{"org-a": "event-a", "org-b": "event-b"} { + if err := manager.Log(orgID, "login", orgID, "127.0.0.1", "/api/auth", true, eventID); err != nil { + t.Fatalf("log %s event: %v", orgID, err) + } + } + + handler := NewAuditHandlers() + for orgID, wantID := range map[string]string{"org-a": "event-a", "org-b": "event-b"} { + req := httptest.NewRequest(http.MethodGet, "/api/audit", nil) + req = req.WithContext(context.WithValue(req.Context(), OrgIDContextKey, orgID)) + rec := httptest.NewRecorder() + handler.HandleListAuditEvents(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("%s status = %d, want 200", orgID, rec.Code) + } + var payload struct { + Events []audit.Event `json:"events"` + } + if err := json.Unmarshal(rec.Body.Bytes(), &payload); err != nil { + t.Fatalf("decode %s response: %v", orgID, err) + } + if len(payload.Events) != 1 || payload.Events[0].Details != wantID { + t.Fatalf("%s events = %#v, want only %q", orgID, payload.Events, wantID) + } + } +} + func TestHandleListAuditEvents_QueryError(t *testing.T) { setAuditLogger(t, &testAuditLogger{ queryErr: fmt.Errorf("db error"), diff --git a/pkg/audit/async_logger.go b/pkg/audit/async_logger.go index b1ec8c9d2..6976dba9e 100644 --- a/pkg/audit/async_logger.go +++ b/pkg/audit/async_logger.go @@ -71,6 +71,11 @@ func (l *AsyncLogger) Count(filter QueryFilter) (int, error) { return l.backend.Count(filter) } +// QueryPage preserves snapshot-consistent pagination when the backend supports it. +func (l *AsyncLogger) QueryPage(filter QueryFilter) ([]Event, int, error) { + return QueryPage(l.backend, filter) +} + // GetWebhookURLs delegates to the backend logger. func (l *AsyncLogger) GetWebhookURLs() []string { return l.backend.GetWebhookURLs() @@ -89,6 +94,17 @@ func (l *AsyncLogger) IsPersistentAuditLogger() bool { return IsPersistentLogger(l.backend) } +// VerifySignature delegates verification to persistent backends. +func (l *AsyncLogger) VerifySignature(event Event) bool { + if l == nil { + return false + } + verifier, ok := l.backend.(interface { + VerifySignature(Event) bool + }) + return ok && verifier.VerifySignature(event) +} + // Close drains queued events, stops the worker, and closes the backend logger. func (l *AsyncLogger) Close() error { if l == nil { diff --git a/pkg/audit/async_logger_test.go b/pkg/audit/async_logger_test.go new file mode 100644 index 000000000..a28b9ae00 --- /dev/null +++ b/pkg/audit/async_logger_test.go @@ -0,0 +1,41 @@ +package audit + +import ( + "testing" + "time" +) + +func TestAsyncLoggerPreservesPersistentReadAndVerificationContracts(t *testing.T) { + backend, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: t.TempDir(), + CryptoMgr: newMockCryptoManager(), + }) + if err != nil { + t.Fatalf("NewSQLiteLogger: %v", err) + } + event := Event{ + ID: "async-contract", + Timestamp: time.Now().UTC(), + EventType: "test", + Success: true, + } + if err := backend.Record(event); err != nil { + t.Fatalf("record event: %v", err) + } + + logger := NewAsyncLogger(backend, AsyncLoggerConfig{BufferSize: 8}) + defer logger.Close() + events, total, err := logger.QueryPage(QueryFilter{ID: event.ID, Limit: 1}) + if err != nil { + t.Fatalf("QueryPage: %v", err) + } + if len(events) != 1 || total != 1 { + t.Fatalf("QueryPage = %d events, total %d, want 1/1", len(events), total) + } + if !logger.VerifySignature(events[0]) { + t.Fatal("async wrapper did not delegate signature verification") + } + if _, ok := any(logger).(PersistentLogger); !ok { + t.Fatal("async persistent logger no longer satisfies export contract") + } +} diff --git a/pkg/audit/audit.go b/pkg/audit/audit.go index f1a34a902..dd67f46e5 100644 --- a/pkg/audit/audit.go +++ b/pkg/audit/audit.go @@ -65,6 +65,31 @@ type Logger interface { Close() error } +type pageLogger interface { + QueryPage(filter QueryFilter) ([]Event, int, error) +} + +// QueryPage returns a page and matching count from one storage snapshot when +// the logger supports it. Other logger implementations retain the legacy +// Query-then-Count contract. +func QueryPage(logger Logger, filter QueryFilter) ([]Event, int, error) { + if pager, ok := logger.(pageLogger); ok { + return pager.QueryPage(filter) + } + events, err := logger.Query(filter) + if err != nil { + return nil, 0, err + } + countFilter := filter + countFilter.Limit = 0 + countFilter.Offset = 0 + total, err := logger.Count(countFilter) + if err != nil { + return nil, 0, err + } + return events, total, nil +} + type persistentAuditLogger interface { IsPersistentAuditLogger() bool } diff --git a/pkg/audit/signer.go b/pkg/audit/signer.go index d89825630..70ae4b75d 100644 --- a/pkg/audit/signer.go +++ b/pkg/audit/signer.go @@ -1,6 +1,7 @@ package audit import ( + "bytes" "crypto/hmac" "crypto/rand" "crypto/sha256" @@ -11,6 +12,7 @@ import ( "os" "path/filepath" "strconv" + "time" "github.com/rs/zerolog/log" ) @@ -45,8 +47,8 @@ func NewSigner(dataDir string, cryptoMgr CryptoEncryptor) (*Signer, error) { if err != nil { return nil, err } - if len(key) != 32 { - return nil, fmt.Errorf("invalid audit signing key length: got %d, want 32", len(key)) + if len(key) < 32 { + return nil, fmt.Errorf("invalid audit signing key length: got %d, want at least 32", len(key)) } if migratedPlaintext { rewritten, err := cryptoMgr.Encrypt(key) @@ -85,13 +87,29 @@ func NewSigner(dataDir string, cryptoMgr CryptoEncryptor) (*Signer, error) { return &Signer{key: key}, nil } +// NewSignerWithKey creates a signer backed by externally managed key material. +func NewSignerWithKey(key []byte) (*Signer, error) { + if len(key) < 32 { + return nil, fmt.Errorf("invalid audit signing key length: got %d, want at least 32", len(key)) + } + return &Signer{key: append([]byte(nil), key...)}, nil +} + func loadAuditSigningKey(cryptoMgr CryptoEncryptor, data []byte) ([]byte, bool, error) { key, err := cryptoMgr.Decrypt(data) if err == nil { return key, false, nil } - if len(data) == 32 { - return append([]byte(nil), data...), true, nil + plaintext := bytes.TrimSpace(data) + if len(plaintext) == 32 { + return append([]byte(nil), plaintext...), true, nil + } + if len(plaintext) == 64 { + if _, decodeErr := hex.DecodeString(string(plaintext)); decodeErr == nil { + // The former Pro store used the printable hex value itself as the + // HMAC key. Preserve those bytes while encrypting the file in place. + return append([]byte(nil), plaintext...), true, nil + } } return nil, false, fmt.Errorf("failed to decrypt audit signing key: %w", err) } @@ -116,8 +134,23 @@ func (s *Signer) Verify(event Event) bool { return false } - expected := s.Sign(event) - return hmac.Equal([]byte(expected), []byte(event.Signature)) + for _, canonical := range []string{ + s.canonicalForm(event), + s.legacyUnixCanonicalForm(event), + s.legacyTimeCanonicalForm(event), + } { + expected := s.signCanonical(canonical) + if hmac.Equal([]byte(expected), []byte(event.Signature)) { + return true + } + } + return false +} + +func (s *Signer) signCanonical(canonical string) string { + mac := hmac.New(sha256.New, s.key) + mac.Write([]byte(canonical)) + return hex.EncodeToString(mac.Sum(nil)) } // canonicalForm creates a deterministic string representation of an event for signing. @@ -138,6 +171,28 @@ func (s *Signer) canonicalForm(event Event) string { event.Details } +func (s *Signer) legacyUnixCanonicalForm(event Event) string { + return event.ID + "|" + + strconv.FormatInt(event.Timestamp.Unix(), 10) + "|" + + event.EventType + "|" + + event.User + "|" + + event.IP + "|" + + event.Path + "|" + + strconv.FormatBool(event.Success) + "|" + + event.Details +} + +func (s *Signer) legacyTimeCanonicalForm(event Event) string { + return event.ID + "|" + + event.Timestamp.UTC().Format(time.RFC3339Nano) + "|" + + event.EventType + "|" + + event.User + "|" + + event.IP + "|" + + event.Path + "|" + + strconv.FormatBool(event.Success) + "|" + + event.Details +} + // SigningEnabled returns true if the signer has a valid key. func (s *Signer) SigningEnabled() bool { return s.key != nil diff --git a/pkg/audit/signer_test.go b/pkg/audit/signer_test.go index 16c58b4cf..d9a662fda 100644 --- a/pkg/audit/signer_test.go +++ b/pkg/audit/signer_test.go @@ -136,6 +136,41 @@ func TestNewSigner_MigratesPlaintextKeyFile(t *testing.T) { } } +func TestNewSigner_MigratesLegacyProHexKeyFile(t *testing.T) { + tempDir := t.TempDir() + crypto := taggedMockCryptoManager{} + legacyKey := []byte("0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef") + keyPath := filepath.Join(tempDir, ".audit-signing.key") + if err := os.WriteFile(keyPath, append(append([]byte(nil), legacyKey...), '\n'), 0o600); err != nil { + t.Fatalf("write legacy Pro key: %v", err) + } + + signer, err := NewSigner(tempDir, crypto) + if err != nil { + t.Fatalf("migrate legacy Pro key: %v", err) + } + expected, err := NewSignerWithKey(legacyKey) + if err != nil { + t.Fatalf("create expected signer: %v", err) + } + event := Event{ + ID: "legacy-key", + Timestamp: time.Date(2026, 7, 5, 10, 30, 0, 0, time.UTC), + EventType: "startup", + Success: true, + } + if signer.Sign(event) != expected.Sign(event) { + t.Fatal("legacy Pro key bytes changed during migration") + } + reloaded, err := NewSigner(tempDir, crypto) + if err != nil { + t.Fatalf("reload migrated legacy Pro key: %v", err) + } + if signer.Sign(event) != reloaded.Sign(event) { + t.Fatal("migrated legacy Pro key did not survive restart") + } +} + func TestNewSignerWithoutCrypto(t *testing.T) { tempDir := t.TempDir() @@ -253,6 +288,32 @@ func TestSignerVerify(t *testing.T) { } } +func TestSignerVerifyAcceptsLegacyProCanonicalForms(t *testing.T) { + signer, err := NewSignerWithKey([]byte("0123456789abcdef0123456789abcdef")) + if err != nil { + t.Fatalf("NewSignerWithKey: %v", err) + } + event := Event{ + ID: "legacy-event", + Timestamp: time.Date(2026, 7, 5, 10, 30, 0, 123456789, time.UTC), + EventType: "startup", + User: "admin", + IP: "127.0.0.1", + Path: "/api/audit", + Success: true, + Details: "legacy", + } + + event.Signature = signer.signCanonical(signer.legacyUnixCanonicalForm(event)) + if !signer.Verify(event) { + t.Fatal("current Pro Unix/boolean signature did not verify") + } + event.Signature = signer.signCanonical(signer.legacyTimeCanonicalForm(event)) + if !signer.Verify(event) { + t.Fatal("v5 Pro RFC3339/boolean signature did not verify") + } +} + func TestSignerCanonicalForm(t *testing.T) { tempDir := t.TempDir() crypto := newMockCryptoManager() diff --git a/pkg/audit/sqlite_factory.go b/pkg/audit/sqlite_factory.go index 213257ee5..72285744c 100644 --- a/pkg/audit/sqlite_factory.go +++ b/pkg/audit/sqlite_factory.go @@ -3,6 +3,7 @@ package audit import ( "fmt" "path/filepath" + "time" ) // SQLiteLoggerFactory creates SQLite-backed audit loggers for tenant databases. @@ -24,6 +25,12 @@ type SQLiteLoggerFactory struct { // RetentionDays controls how long audit events are retained (0 uses SQLiteLogger defaults). RetentionDays int + + // RetentionConfigured treats RetentionDays=0 as an explicit forever setting. + RetentionConfigured bool + + // CleanupInterval controls the retention cleanup cadence. + CleanupInterval time.Duration } func (f *SQLiteLoggerFactory) CreateLogger(dbPath string) (Logger, error) { @@ -34,8 +41,10 @@ func (f *SQLiteLoggerFactory) CreateLogger(dbPath string) (Logger, error) { dataDir := filepath.Dir(dbPath) cfg := SQLiteLoggerConfig{ - DataDir: dataDir, - RetentionDays: f.RetentionDays, + DataDir: dataDir, + RetentionDays: f.RetentionDays, + RetentionConfigured: f.RetentionConfigured, + CleanupInterval: f.CleanupInterval, } if f.CryptoMgrForDataDir != nil { diff --git a/pkg/audit/sqlite_logger.go b/pkg/audit/sqlite_logger.go index 9ca56e4f2..58c25ee9d 100644 --- a/pkg/audit/sqlite_logger.go +++ b/pkg/audit/sqlite_logger.go @@ -18,9 +18,13 @@ import ( // SQLiteLoggerConfig configures the SQLite audit logger. type SQLiteLoggerConfig struct { - DataDir string // Directory for audit.db - CryptoMgr CryptoEncryptor // For encrypting the signing key (optional) - RetentionDays int // Days to keep events (default: 90, 0 = forever) + DataDir string // Base data directory; audit data is stored in /audit. + AuditDir string // Optional exact audit directory, overriding DataDir. + CryptoMgr CryptoEncryptor // For encrypting the signing key (optional). + SigningKey []byte // Optional externally managed HMAC key. + RetentionDays int // Days to keep events. + RetentionConfigured bool // Treat RetentionDays=0 as an explicit forever setting. + CleanupInterval time.Duration // Retention cleanup cadence (default: 24 hours). } // SQLiteLogger implements Logger with persistent SQLite storage and HMAC signing. @@ -31,6 +35,7 @@ type SQLiteLogger struct { signer *Signer webhookDelivery *WebhookDelivery retentionDays int + cleanupInterval time.Duration stopChan chan struct{} wg sync.WaitGroup closeOnce sync.Once @@ -60,6 +65,20 @@ const ( var auditSQLiteRetryDelays = []time.Duration{25 * time.Millisecond, 75 * time.Millisecond} var auditSQLiteRetrySleep = time.Sleep +const ( + auditSchemaVersion = 2 + auditTimestampMigrationBatchSize = 512 + defaultAuditCleanupInterval = 24 * time.Hour +) + +type auditStoreDataError struct { + field string +} + +func (e *auditStoreDataError) Error() string { + return "audit store contains an invalid " + e.field + " encoding" +} + // IsStoreBusyError reports whether an audit store error is a transient SQLite lock. func IsStoreBusyError(err error) bool { if err == nil { @@ -96,6 +115,10 @@ func IsStoreUnavailableError(err error) bool { if IsStoreBusyError(err) { return false } + var dataErr *auditStoreDataError + if errors.As(err, &dataErr) { + return true + } var sqliteErr *sqlite.Error if errors.As(err, &sqliteErr) { @@ -135,12 +158,15 @@ func withSQLiteRetry(operation string, run func() error) error { // NewSQLiteLogger creates a new SQLite-backed audit logger. func NewSQLiteLogger(cfg SQLiteLoggerConfig) (*SQLiteLogger, error) { - if cfg.DataDir == "" { + if cfg.DataDir == "" && cfg.AuditDir == "" { return nil, fmt.Errorf("data directory is required") } // Ensure directory exists - auditDir := filepath.Join(cfg.DataDir, "audit") + auditDir := cfg.AuditDir + if auditDir == "" { + auditDir = filepath.Join(cfg.DataDir, "audit") + } if err := os.MkdirAll(auditDir, 0700); err != nil { return nil, fmt.Errorf("failed to create audit directory: %w", err) } @@ -168,23 +194,33 @@ func NewSQLiteLogger(cfg SQLiteLoggerConfig) (*SQLiteLogger, error) { db.SetConnMaxLifetime(0) // Initialize signer - signer, err := NewSigner(auditDir, cfg.CryptoMgr) + var signer *Signer + if len(cfg.SigningKey) > 0 { + signer, err = NewSignerWithKey(cfg.SigningKey) + } else { + signer, err = NewSigner(auditDir, cfg.CryptoMgr) + } if err != nil { db.Close() return nil, fmt.Errorf("failed to initialize audit signer: %w", err) } retentionDays := cfg.RetentionDays - if retentionDays == 0 { + if retentionDays == 0 && !cfg.RetentionConfigured { retentionDays = 90 // Default } + cleanupInterval := cfg.CleanupInterval + if cleanupInterval <= 0 { + cleanupInterval = defaultAuditCleanupInterval + } l := &SQLiteLogger{ - db: db, - dbPath: dbPath, - signer: signer, - retentionDays: retentionDays, - stopChan: make(chan struct{}), + db: db, + dbPath: dbPath, + signer: signer, + retentionDays: retentionDays, + cleanupInterval: cleanupInterval, + stopChan: make(chan struct{}), } // Initialize schema @@ -192,6 +228,12 @@ func NewSQLiteLogger(cfg SQLiteLoggerConfig) (*SQLiteLogger, error) { db.Close() return nil, fmt.Errorf("failed to initialize schema: %w", err) } + if cfg.RetentionDays == 0 && !cfg.RetentionConfigured { + if persisted, ok := l.loadRetentionDays(); ok { + l.retentionDays = persisted + retentionDays = persisted + } + } l.migrateAutoVacuum() @@ -238,6 +280,7 @@ func (l *SQLiteLogger) initSchema() error { ); CREATE INDEX IF NOT EXISTS idx_audit_timestamp ON audit_events(timestamp); + CREATE INDEX IF NOT EXISTS idx_audit_timestamp_id ON audit_events(timestamp DESC, id DESC); CREATE INDEX IF NOT EXISTS idx_audit_event_type ON audit_events(event_type); CREATE INDEX IF NOT EXISTS idx_audit_user ON audit_events(user) WHERE user != ''; CREATE INDEX IF NOT EXISTS idx_audit_success ON audit_events(success); @@ -269,7 +312,235 @@ func (l *SQLiteLogger) initSchema() error { time.Now().Unix()) return err }) - return err + if err != nil { + return err + } + + return l.migrateLegacyTimestamps() +} + +type canonicalAuditRow struct { + rowID int64 + id string + timestamp int64 + eventType string + user sql.NullString + ip sql.NullString + path sql.NullString + success int + details sql.NullString + signature sql.NullString +} + +func (l *SQLiteLogger) migrateLegacyTimestamps() error { + var migrated int + err := withSQLiteRetry("migrate_audit_timestamps", func() error { + tx, err := l.db.Begin() + if err != nil { + return err + } + defer tx.Rollback() + + needsRebuild, err := auditEventsNeedRebuild(tx) + if err != nil { + return err + } + if !needsRebuild { + if _, err := tx.Exec( + `INSERT OR IGNORE INTO schema_version (version, applied_at) VALUES (?, ?)`, + auditSchemaVersion, + time.Now().Unix(), + ); err != nil { + return err + } + return tx.Commit() + } + + if err := tx.QueryRow( + `SELECT COUNT(*) FROM audit_events WHERE typeof(timestamp) != 'integer'`, + ).Scan(&migrated); err != nil { + return err + } + if _, err := tx.Exec(`DROP TABLE IF EXISTS audit_events_v2`); err != nil { + return err + } + if _, err := tx.Exec(` + CREATE TABLE audit_events_v2 ( + id TEXT PRIMARY KEY, + timestamp INTEGER NOT NULL, + event_type TEXT NOT NULL, + user TEXT, + ip TEXT, + path TEXT, + success INTEGER NOT NULL, + details TEXT, + signature TEXT NOT NULL + )`); err != nil { + return err + } + + lastRowID := int64(-1 << 63) + for { + rows, err := tx.Query(` + SELECT rowid, id, timestamp, event_type, user, ip, path, success, details, signature + FROM audit_events + WHERE rowid > ? + ORDER BY rowid + LIMIT ?`, lastRowID, auditTimestampMigrationBatchSize) + if err != nil { + return err + } + + batch := make([]canonicalAuditRow, 0, auditTimestampMigrationBatchSize) + for rows.Next() { + var item canonicalAuditRow + var value any + if err := rows.Scan( + &item.rowID, + &item.id, + &value, + &item.eventType, + &item.user, + &item.ip, + &item.path, + &item.success, + &item.details, + &item.signature, + ); err != nil { + rows.Close() + return &auditStoreDataError{field: "row"} + } + timestamp, err := parseAuditTimestamp(value) + if err != nil { + rows.Close() + return &auditStoreDataError{field: "timestamp"} + } + if item.id == "" || item.eventType == "" || (item.success != 0 && item.success != 1) { + rows.Close() + return &auditStoreDataError{field: "row"} + } + item.timestamp = timestamp.Unix() + batch = append(batch, item) + } + if err := rows.Err(); err != nil { + rows.Close() + return err + } + if err := rows.Close(); err != nil { + return err + } + if len(batch) == 0 { + break + } + + for _, item := range batch { + if _, err := tx.Exec( + `INSERT INTO audit_events_v2 + (id, timestamp, event_type, user, ip, path, success, details, signature) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, + item.id, + item.timestamp, + item.eventType, + item.user, + item.ip, + item.path, + item.success, + item.details, + item.signature.String, + ); err != nil { + return err + } + } + lastRowID = batch[len(batch)-1].rowID + } + + if _, err := tx.Exec(`DROP TABLE audit_events`); err != nil { + return err + } + if _, err := tx.Exec(`ALTER TABLE audit_events_v2 RENAME TO audit_events`); err != nil { + return err + } + if _, err := tx.Exec(` + CREATE INDEX idx_audit_timestamp ON audit_events(timestamp); + CREATE INDEX idx_audit_timestamp_id ON audit_events(timestamp DESC, id DESC); + CREATE INDEX idx_audit_event_type ON audit_events(event_type); + CREATE INDEX idx_audit_user ON audit_events(user) WHERE user != ''; + CREATE INDEX idx_audit_success ON audit_events(success); + `); err != nil { + return err + } + + if _, err := tx.Exec( + `INSERT OR REPLACE INTO schema_version (version, applied_at) VALUES (?, ?)`, + auditSchemaVersion, + time.Now().Unix(), + ); err != nil { + return err + } + return tx.Commit() + }) + if err != nil { + return fmt.Errorf("normalize audit timestamp storage: %w", err) + } + if migrated > 0 { + log.Info(). + Int("rows", migrated). + Msg("Normalized legacy audit timestamp storage") + } + return nil +} + +func auditEventsNeedRebuild(tx *sql.Tx) (bool, error) { + rows, err := tx.Query(`PRAGMA table_info(audit_events)`) + if err != nil { + return false, err + } + defer rows.Close() + + columns := make(map[string]struct { + columnType string + notNull bool + }) + for rows.Next() { + var position int + var name, columnType string + var notNull, primaryKey int + var defaultValue any + if err := rows.Scan( + &position, + &name, + &columnType, + ¬Null, + &defaultValue, + &primaryKey, + ); err != nil { + return false, &auditStoreDataError{field: "schema"} + } + columns[name] = struct { + columnType string + notNull bool + }{columnType: strings.ToUpper(columnType), notNull: notNull != 0} + } + if err := rows.Err(); err != nil { + return false, err + } + + timestamp, timestampOK := columns["timestamp"] + signature, signatureOK := columns["signature"] + if !timestampOK || !signatureOK { + return false, &auditStoreDataError{field: "schema"} + } + if timestamp.columnType != "INTEGER" || !timestamp.notNull || !signature.notNull { + return true, nil + } + + var hasNonInteger int + if err := tx.QueryRow( + `SELECT EXISTS(SELECT 1 FROM audit_events WHERE typeof(timestamp) != 'integer')`, + ).Scan(&hasNonInteger); err != nil { + return false, err + } + return hasNonInteger != 0, nil } // Log records an audit event with HMAC signature. @@ -351,6 +622,15 @@ func (l *SQLiteLogger) Query(filter QueryFilter) ([]Event, error) { } func (l *SQLiteLogger) queryLocked(filter QueryFilter) ([]Event, error) { + return queryAuditEvents(l.db, filter) +} + +type auditSQLReader interface { + Query(query string, args ...any) (*sql.Rows, error) + QueryRow(query string, args ...any) *sql.Row +} + +func queryAuditEvents(reader auditSQLReader, filter QueryFilter) ([]Event, error) { query := "SELECT id, timestamp, event_type, user, ip, path, success, details, signature FROM audit_events WHERE 1=1" args := []interface{}{} @@ -383,7 +663,7 @@ func (l *SQLiteLogger) queryLocked(filter QueryFilter) ([]Event, error) { args = append(args, success) } - query += " ORDER BY timestamp DESC" + query += " ORDER BY timestamp DESC, id DESC" if filter.Limit > 0 { query += " LIMIT ?" @@ -398,13 +678,13 @@ func (l *SQLiteLogger) queryLocked(filter QueryFilter) ([]Event, error) { args = append(args, filter.Offset) } - rows, err := l.db.Query(query, args...) + rows, err := reader.Query(query, args...) if err != nil { return nil, fmt.Errorf("failed to query audit events: %w", err) } defer rows.Close() - var events []Event + events := make([]Event, 0) for rows.Next() { var e Event var timestampValue any @@ -434,6 +714,39 @@ func (l *SQLiteLogger) queryLocked(filter QueryFilter) ([]Event, error) { return events, rows.Err() } +// QueryPage reads events and their matching count from one SQLite snapshot. +func (l *SQLiteLogger) QueryPage(filter QueryFilter) ([]Event, int, error) { + var events []Event + var total int + err := withSQLiteRetry("query_audit_page", func() error { + l.mu.RLock() + defer l.mu.RUnlock() + + tx, err := l.db.Begin() + if err != nil { + return err + } + defer tx.Rollback() + + events, err = queryAuditEvents(tx, filter) + if err != nil { + return err + } + countFilter := filter + countFilter.Limit = 0 + countFilter.Offset = 0 + total, err = countAuditEvents(tx, countFilter) + if err != nil { + return err + } + return tx.Commit() + }) + if err != nil { + return nil, 0, err + } + return events, total, nil +} + func parseAuditTimestamp(value any) (time.Time, error) { switch v := value.(type) { case time.Time: @@ -493,7 +806,7 @@ func parseAuditTimestampString(raw string) (time.Time, error) { return timestamp, nil } } - return time.Time{}, fmt.Errorf("unsupported timestamp value %q", raw) + return time.Time{}, errors.New("unsupported timestamp encoding") } // Count returns the number of events matching the filter. @@ -517,6 +830,10 @@ func (l *SQLiteLogger) Count(filter QueryFilter) (int, error) { } func (l *SQLiteLogger) countLocked(filter QueryFilter) (int, error) { + return countAuditEvents(l.db, filter) +} + +func countAuditEvents(reader auditSQLReader, filter QueryFilter) (int, error) { query := "SELECT COUNT(*) FROM audit_events WHERE 1=1" args := []interface{}{} @@ -550,7 +867,7 @@ func (l *SQLiteLogger) countLocked(filter QueryFilter) (int, error) { } var count int - err := l.db.QueryRow(query, args...).Scan(&count) + err := reader.QueryRow(query, args...).Scan(&count) if err != nil { return 0, fmt.Errorf("failed to count audit events: %w", err) } @@ -638,12 +955,26 @@ func (l *SQLiteLogger) loadWebhookURLs() []string { return strings.Split(value, ",") } +func (l *SQLiteLogger) loadRetentionDays() (int, bool) { + var value string + if err := l.db.QueryRow( + `SELECT value FROM audit_config WHERE key = 'retention_days'`, + ).Scan(&value); err != nil { + return 0, false + } + days, err := strconv.Atoi(value) + if err != nil || days < 0 { + return 0, false + } + return days, true +} + // retentionWorker runs periodically to clean up old events. func (l *SQLiteLogger) retentionWorker() { defer l.wg.Done() // Run at 3 AM daily - ticker := time.NewTicker(24 * time.Hour) + ticker := time.NewTicker(l.cleanupInterval) defer ticker.Stop() // Also run once at startup after a short delay. Keep timer stoppable @@ -734,43 +1065,79 @@ func (l *SQLiteLogger) reclaimFreePages() { // cleanupOldEvents deletes events older than the retention period. func (l *SQLiteLogger) cleanupOldEvents() { - if l.retentionDays <= 0 { - return - } + var deleted int64 + var retentionDays int + err := withSQLiteRetry("cleanup_audit_events", func() error { + l.mu.Lock() + defer l.mu.Unlock() - l.mu.Lock() - defer l.mu.Unlock() + retentionDays = l.retentionDays + if retentionDays <= 0 { + return nil + } + cutoff := time.Now().AddDate(0, 0, -retentionDays).Unix() + tx, err := l.db.Begin() + if err != nil { + return err + } + defer tx.Rollback() - cutoff := time.Now().AddDate(0, 0, -l.retentionDays).Unix() - - result, err := l.db.Exec(`DELETE FROM audit_events WHERE timestamp < ?`, cutoff) + result, err := tx.Exec(`DELETE FROM audit_events WHERE timestamp < ?`, cutoff) + if err != nil { + return err + } + deleted, err = result.RowsAffected() + if err != nil { + return err + } + if deleted > 0 { + now := time.Now().UTC().Truncate(time.Second) + event := Event{ + ID: fmt.Sprintf("cleanup-%d", now.UnixNano()), + Timestamp: now, + EventType: "audit_cleanup", + User: "system", + Success: true, + Details: fmt.Sprintf("Deleted %d events older than %d days", deleted, retentionDays), + } + event.Signature = l.signer.Sign(event) + if _, err := tx.Exec(` + INSERT INTO audit_events (id, timestamp, event_type, user, ip, path, success, details, signature) + VALUES (?, ?, ?, ?, ?, ?, 1, ?, ?)`, + event.ID, + event.Timestamp.Unix(), + event.EventType, + event.User, + event.IP, + event.Path, + event.Details, + event.Signature, + ); err != nil { + return err + } + } + if err := tx.Commit(); err != nil { + return err + } + l.reclaimFreePages() + return nil + }) if err != nil { log.Error().Err(err).Msg("Failed to cleanup old audit events") return } - - deleted, _ := result.RowsAffected() if deleted > 0 { log.Info(). Int64("deleted", deleted). - Int("retentionDays", l.retentionDays). + Int("retentionDays", retentionDays). Msg("Cleaned up old audit events") - - // Log the cleanup as an audit event (without recursion - direct insert) - _, _ = l.db.Exec(` - INSERT INTO audit_events (id, timestamp, event_type, user, ip, path, success, details, signature) - VALUES (?, ?, 'audit_cleanup', 'system', '', '', 1, ?, '')`, - fmt.Sprintf("cleanup-%d", time.Now().Unix()), - time.Now().Unix(), - fmt.Sprintf("Deleted %d events older than %d days", deleted, l.retentionDays), - ) } - - l.reclaimFreePages() } // GetRetentionDays returns the current retention period. func (l *SQLiteLogger) GetRetentionDays() int { + l.mu.RLock() + defer l.mu.RUnlock() return l.retentionDays } diff --git a/pkg/audit/sqlite_logger_migration_test.go b/pkg/audit/sqlite_logger_migration_test.go new file mode 100644 index 000000000..2f88dc6f2 --- /dev/null +++ b/pkg/audit/sqlite_logger_migration_test.go @@ -0,0 +1,465 @@ +package audit + +import ( + "database/sql" + "errors" + "fmt" + "net/url" + "os" + "path/filepath" + "strings" + "sync" + "testing" + "time" + + _ "modernc.org/sqlite" +) + +type legacyAuditFixtureRow struct { + id string + timestamp any + user any + signature any +} + +func createLegacyAuditFixture( + t *testing.T, + dataDir string, + rows []legacyAuditFixtureRow, +) string { + t.Helper() + + auditDir := filepath.Join(dataDir, "audit") + if err := os.MkdirAll(auditDir, 0o700); err != nil { + t.Fatalf("create legacy audit directory: %v", err) + } + dbPath := filepath.Join(auditDir, "audit.db") + db, err := sql.Open("sqlite", dbPath) + if err != nil { + t.Fatalf("open legacy audit database: %v", err) + } + defer db.Close() + + if _, err := db.Exec(` + CREATE TABLE audit_events ( + id TEXT PRIMARY KEY, + timestamp DATETIME NOT NULL, + event_type TEXT NOT NULL, + user TEXT, + ip TEXT, + path TEXT, + success INTEGER NOT NULL, + details TEXT, + signature TEXT, + created_at DATETIME DEFAULT CURRENT_TIMESTAMP + ); + CREATE INDEX idx_audit_timestamp ON audit_events(timestamp); + `); err != nil { + t.Fatalf("create legacy audit schema: %v", err) + } + for _, row := range rows { + if _, err := db.Exec(` + INSERT INTO audit_events + (id, timestamp, event_type, user, ip, path, success, details, signature) + VALUES (?, ?, 'startup', ?, '127.0.0.1', '/api/audit', 1, 'fixture', ?)`, + row.id, + row.timestamp, + row.user, + row.signature, + ); err != nil { + t.Fatalf("insert legacy row %q: %v", row.id, err) + } + } + return dbPath +} + +func TestSQLiteLoggerMigratesLegacySchemaAndQueryContract(t *testing.T) { + dataDir := t.TempDir() + now := time.Now().UTC().Truncate(time.Second) + old := now.AddDate(0, 0, -120) + recent := now.Add(-time.Hour) + sameSecond := now.Add(-30 * time.Minute) + dbPath := createLegacyAuditFixture(t, dataDir, []legacyAuditFixtureRow{ + { + id: "legacy-old", + timestamp: old.Format("2006-01-02 15:04:05 -0700 MST") + " m=+0.009025344", + user: nil, + signature: nil, + }, + {id: "legacy-recent", timestamp: recent, user: "alice", signature: "legacy-signature"}, + {id: "same-a", timestamp: sameSecond, user: "alice", signature: "legacy-signature"}, + {id: "same-z", timestamp: sameSecond, user: "alice", signature: "legacy-signature"}, + }) + + logger, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: dataDir, + CryptoMgr: newMockCryptoManager(), + RetentionDays: 30, + }) + if err != nil { + t.Fatalf("migrate legacy logger: %v", err) + } + + start := now.Add(-2 * time.Hour) + events, total, err := logger.QueryPage(QueryFilter{ + StartTime: &start, + User: "alice", + Limit: 10, + }) + if err != nil { + t.Fatalf("query migrated page: %v", err) + } + if total != 3 || len(events) != 3 { + t.Fatalf("migrated page = %d events, total %d, want 3/3", len(events), total) + } + if events[0].ID != "same-z" || events[1].ID != "same-a" || events[2].ID != "legacy-recent" { + t.Fatalf("deterministic order = %v", []string{events[0].ID, events[1].ID, events[2].ID}) + } + + var timestampType string + var timestampNotNull int + rows, err := logger.db.Query(`PRAGMA table_info(audit_events)`) + if err != nil { + t.Fatalf("read migrated schema: %v", err) + } + for rows.Next() { + var position, primaryKey int + var name, columnType string + var notNull int + var defaultValue any + if err := rows.Scan( + &position, + &name, + &columnType, + ¬Null, + &defaultValue, + &primaryKey, + ); err != nil { + t.Fatalf("scan migrated schema: %v", err) + } + if name == "timestamp" { + timestampType = columnType + timestampNotNull = notNull + } + } + if err := rows.Close(); err != nil { + t.Fatalf("close schema rows: %v", err) + } + if timestampType != "INTEGER" || timestampNotNull != 1 { + t.Fatalf("timestamp schema = %q/%d, want INTEGER/1", timestampType, timestampNotNull) + } + var nonInteger int + if err := logger.db.QueryRow( + `SELECT COUNT(*) FROM audit_events WHERE typeof(timestamp) != 'integer'`, + ).Scan(&nonInteger); err != nil { + t.Fatalf("count non-integer timestamps: %v", err) + } + if nonInteger != 0 { + t.Fatalf("non-integer timestamps = %d, want 0", nonInteger) + } + var schemaVersion int + if err := logger.db.QueryRow(`SELECT MAX(version) FROM schema_version`).Scan(&schemaVersion); err != nil { + t.Fatalf("read schema version: %v", err) + } + if schemaVersion != auditSchemaVersion { + t.Fatalf("schema version = %d, want %d", schemaVersion, auditSchemaVersion) + } + + logger.cleanupOldEvents() + oldEvents, err := logger.Query(QueryFilter{ID: "legacy-old"}) + if err != nil { + t.Fatalf("query retained legacy row: %v", err) + } + if len(oldEvents) != 0 { + t.Fatal("legacy row older than retention window was not removed") + } + cleanupEvents, err := logger.Query(QueryFilter{EventType: "audit_cleanup", Limit: 1}) + if err != nil || len(cleanupEvents) != 1 { + t.Fatalf("cleanup event = %d, err %v", len(cleanupEvents), err) + } + if !logger.VerifySignature(cleanupEvents[0]) { + t.Fatal("cleanup event must carry a valid signature") + } + + if err := logger.Close(); err != nil { + t.Fatalf("close migrated logger: %v", err) + } + restarted, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: dataDir, + CryptoMgr: newMockCryptoManager(), + RetentionDays: 30, + }) + if err != nil { + t.Fatalf("restart migrated logger: %v", err) + } + defer restarted.Close() + if got := restarted.dbPath; got != dbPath { + t.Fatalf("restart database path = %q, want %q", got, dbPath) + } +} + +func TestSQLiteLoggerEmptyQueryReturnsArrayShape(t *testing.T) { + logger, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: t.TempDir(), + CryptoMgr: newMockCryptoManager(), + }) + if err != nil { + t.Fatalf("NewSQLiteLogger: %v", err) + } + defer logger.Close() + events, total, err := logger.QueryPage(QueryFilter{Limit: 100}) + if err != nil { + t.Fatalf("QueryPage: %v", err) + } + if events == nil || len(events) != 0 || total != 0 { + t.Fatalf("empty page = %#v, total %d, want non-nil empty/0", events, total) + } +} + +func TestSQLiteLoggerLegacyMigrationRollsBackMalformedRows(t *testing.T) { + dataDir := t.TempDir() + const sensitiveValue = "customer-secret-invalid-timestamp" + dbPath := createLegacyAuditFixture(t, dataDir, []legacyAuditFixtureRow{ + {id: "valid", timestamp: time.Now().UTC(), user: "alice", signature: nil}, + {id: "invalid", timestamp: sensitiveValue, user: "bob", signature: nil}, + }) + + logger, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: dataDir, + CryptoMgr: newMockCryptoManager(), + }) + if logger != nil || err == nil { + t.Fatal("malformed legacy row must fail logger initialization") + } + if !IsStoreUnavailableError(err) { + t.Fatalf("malformed row error must be store unavailable: %v", err) + } + if strings.Contains(err.Error(), sensitiveValue) { + t.Fatalf("malformed row error leaked stored value: %v", err) + } + + db, openErr := sql.Open("sqlite", dbPath) + if openErr != nil { + t.Fatalf("reopen rolled-back database: %v", openErr) + } + defer db.Close() + var tableType string + if err := db.QueryRow(` + SELECT type FROM pragma_table_info('audit_events') WHERE name = 'timestamp' + `).Scan(&tableType); err != nil { + t.Fatalf("read rolled-back schema: %v", err) + } + if tableType != "DATETIME" { + t.Fatalf("migration was partially applied, timestamp type = %q", tableType) + } + var scratchTables int + if err := db.QueryRow(` + SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = 'audit_events_v2' + `).Scan(&scratchTables); err != nil { + t.Fatalf("check migration scratch table: %v", err) + } + if scratchTables != 0 { + t.Fatal("migration scratch table survived rollback") + } +} + +func TestSQLiteLoggerMigratesLargeLegacyHistory(t *testing.T) { + dataDir := t.TempDir() + const rowCount = auditTimestampMigrationBatchSize*4 + 17 + rows := make([]legacyAuditFixtureRow, 0, rowCount) + base := time.Date(2025, 1, 1, 0, 0, 0, 0, time.UTC) + for i := 0; i < rowCount; i++ { + rows = append(rows, legacyAuditFixtureRow{ + id: fmt.Sprintf("legacy-%05d", i), + timestamp: base.Add(time.Duration(i) * time.Second), + user: "operator", + signature: nil, + }) + } + createLegacyAuditFixture(t, dataDir, rows) + + logger, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: dataDir, + CryptoMgr: newMockCryptoManager(), + RetentionDays: 0, + RetentionConfigured: true, + }) + if err != nil { + t.Fatalf("migrate large legacy history: %v", err) + } + defer logger.Close() + events, total, err := logger.QueryPage(QueryFilter{Limit: 100, Offset: rowCount - 100}) + if err != nil { + t.Fatalf("query deep page: %v", err) + } + if len(events) != 100 || total != rowCount { + t.Fatalf("deep page = %d events, total %d, want 100/%d", len(events), total, rowCount) + } +} + +func TestSQLiteLoggerRealBusyLockFailsClosedAndRecovers(t *testing.T) { + dataDir := t.TempDir() + logger, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: dataDir, + CryptoMgr: newMockCryptoManager(), + }) + if err != nil { + t.Fatalf("NewSQLiteLogger: %v", err) + } + defer logger.Close() + if _, err := logger.db.Exec(`PRAGMA busy_timeout = 1`); err != nil { + t.Fatalf("set short busy timeout: %v", err) + } + + blockerDSN := logger.dbPath + "?" + url.Values{ + "_pragma": []string{"busy_timeout(1)", "journal_mode(WAL)"}, + }.Encode() + blocker, err := sql.Open("sqlite", blockerDSN) + if err != nil { + t.Fatalf("open blocking connection: %v", err) + } + defer blocker.Close() + tx, err := blocker.Begin() + if err != nil { + t.Fatalf("begin blocking transaction: %v", err) + } + if _, err := tx.Exec(`INSERT INTO audit_config (key, value, updated_at) VALUES ('lock', '1', 1)`); err != nil { + t.Fatalf("acquire write lock: %v", err) + } + + previousSleep := auditSQLiteRetrySleep + auditSQLiteRetrySleep = func(time.Duration) {} + t.Cleanup(func() { auditSQLiteRetrySleep = previousSleep }) + event := Event{ID: "busy", Timestamp: time.Now(), EventType: "test", Success: true} + err = logger.Record(event) + if !IsStoreBusyError(err) { + t.Fatalf("write-lock error = %v, want store busy", err) + } + if err := tx.Rollback(); err != nil { + t.Fatalf("release write lock: %v", err) + } + if err := logger.Record(event); err != nil { + t.Fatalf("write did not recover after lock release: %v", err) + } +} + +func TestSQLiteLoggerConcurrentPageSnapshots(t *testing.T) { + logger, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: t.TempDir(), + CryptoMgr: newMockCryptoManager(), + }) + if err != nil { + t.Fatalf("NewSQLiteLogger: %v", err) + } + defer logger.Close() + + var writers sync.WaitGroup + for writer := 0; writer < 4; writer++ { + writer := writer + writers.Add(1) + go func() { + defer writers.Done() + for i := 0; i < 50; i++ { + _ = logger.Record(Event{ + ID: fmt.Sprintf("%d-%03d", writer, i), + Timestamp: time.Unix(1_800_000_000+int64(i), 0), + EventType: "concurrent", + Success: true, + }) + } + }() + } + for i := 0; i < 40; i++ { + events, total, err := logger.QueryPage(QueryFilter{EventType: "concurrent", Limit: 25}) + if err != nil { + t.Fatalf("concurrent QueryPage: %v", err) + } + if total < len(events) { + t.Fatalf("snapshot total %d smaller than page %d", total, len(events)) + } + seen := make(map[string]struct{}, len(events)) + for _, event := range events { + if _, exists := seen[event.ID]; exists { + t.Fatalf("duplicate event %q in snapshot page", event.ID) + } + seen[event.ID] = struct{}{} + } + } + writers.Wait() +} + +func TestSQLiteLoggerOfflineBackupRestore(t *testing.T) { + sourceDir := t.TempDir() + logger, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: sourceDir, + CryptoMgr: newMockCryptoManager(), + }) + if err != nil { + t.Fatalf("NewSQLiteLogger: %v", err) + } + event := Event{ + ID: "backup-event", + Timestamp: time.Now().UTC(), + EventType: "backup", + Success: true, + } + if err := logger.Record(event); err != nil { + t.Fatalf("record backup event: %v", err) + } + if err := logger.Close(); err != nil { + t.Fatalf("close source logger: %v", err) + } + + restoreDir := t.TempDir() + restoreAuditDir := filepath.Join(restoreDir, "audit") + if err := os.MkdirAll(restoreAuditDir, 0o700); err != nil { + t.Fatalf("create restore directory: %v", err) + } + for _, name := range []string{"audit.db", ".audit-signing.key"} { + source := filepath.Join(sourceDir, "audit", name) + data, err := os.ReadFile(source) + if err != nil { + t.Fatalf("read backup member %q: %v", name, err) + } + if err := os.WriteFile(filepath.Join(restoreAuditDir, name), data, 0o600); err != nil { + t.Fatalf("restore backup member %q: %v", name, err) + } + } + + restored, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: restoreDir, + CryptoMgr: newMockCryptoManager(), + }) + if err != nil { + t.Fatalf("open restored logger: %v", err) + } + defer restored.Close() + events, err := restored.Query(QueryFilter{ID: event.ID}) + if err != nil || len(events) != 1 { + t.Fatalf("restored events = %d, err %v", len(events), err) + } + if !restored.VerifySignature(events[0]) { + t.Fatal("restored event signature did not verify") + } +} + +func TestSQLiteLoggerCorruptStorageFailsUnavailable(t *testing.T) { + dataDir := t.TempDir() + auditDir := filepath.Join(dataDir, "audit") + if err := os.MkdirAll(auditDir, 0o700); err != nil { + t.Fatalf("create audit directory: %v", err) + } + if err := os.WriteFile(filepath.Join(auditDir, "audit.db"), []byte("not a sqlite database"), 0o600); err != nil { + t.Fatalf("write corrupt audit database: %v", err) + } + logger, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: dataDir, + CryptoMgr: newMockCryptoManager(), + }) + if logger != nil || err == nil { + t.Fatal("corrupt database must fail logger initialization") + } + if !IsStoreUnavailableError(err) && !errors.Is(err, os.ErrInvalid) { + t.Fatalf("corrupt database error was not unavailable: %v", err) + } +} diff --git a/pkg/audit/sqlite_logger_queryplan_test.go b/pkg/audit/sqlite_logger_queryplan_test.go index 0d777ef7f..69470c696 100644 --- a/pkg/audit/sqlite_logger_queryplan_test.go +++ b/pkg/audit/sqlite_logger_queryplan_test.go @@ -39,10 +39,10 @@ func TestQueryPlansUseIndexes(t *testing.T) { query: `SELECT id, timestamp, event_type, user, ip, path, success, details, signature FROM audit_events WHERE 1=1 AND timestamp >= ? AND timestamp <= ? - ORDER BY timestamp DESC + ORDER BY timestamp DESC, id DESC LIMIT ?`, args: []any{int64(0), int64(9999999999), 100}, - wantIndex: "idx_audit_timestamp", + wantIndex: "idx_audit_timestamp_id", }, { name: "query by event_type + timestamp", @@ -50,7 +50,7 @@ func TestQueryPlansUseIndexes(t *testing.T) { FROM audit_events WHERE 1=1 AND timestamp >= ? AND timestamp <= ? AND event_type = ? - ORDER BY timestamp DESC + ORDER BY timestamp DESC, id DESC LIMIT ?`, args: []any{int64(0), int64(9999999999), "auth_login", 100}, // Planner may use idx_audit_event_type or idx_audit_timestamp — @@ -62,7 +62,7 @@ func TestQueryPlansUseIndexes(t *testing.T) { FROM audit_events WHERE 1=1 AND timestamp >= ? AND timestamp <= ? AND user = ? - ORDER BY timestamp DESC + ORDER BY timestamp DESC, id DESC LIMIT ?`, args: []any{int64(0), int64(9999999999), "admin", 100}, // Planner may use idx_audit_user or idx_audit_timestamp. @@ -160,6 +160,7 @@ func newAuditPlanTestDB(t *testing.T) *sql.DB { ); CREATE INDEX IF NOT EXISTS idx_audit_timestamp ON audit_events(timestamp); + CREATE INDEX IF NOT EXISTS idx_audit_timestamp_id ON audit_events(timestamp DESC, id DESC); CREATE INDEX IF NOT EXISTS idx_audit_event_type ON audit_events(event_type); CREATE INDEX IF NOT EXISTS idx_audit_user ON audit_events(user) WHERE user != ''; CREATE INDEX IF NOT EXISTS idx_audit_success ON audit_events(success); diff --git a/pkg/audit/sqlite_logger_test.go b/pkg/audit/sqlite_logger_test.go index aa2fdabb6..e8b6e992b 100644 --- a/pkg/audit/sqlite_logger_test.go +++ b/pkg/audit/sqlite_logger_test.go @@ -46,6 +46,39 @@ func TestNewSQLiteLoggerDefaultRetention(t *testing.T) { } } +func TestNewSQLiteLoggerExplicitForeverRetention(t *testing.T) { + logger, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: t.TempDir(), + CryptoMgr: newMockCryptoManager(), + RetentionDays: 0, + RetentionConfigured: true, + }) + if err != nil { + t.Fatalf("NewSQLiteLogger failed: %v", err) + } + defer logger.Close() + if logger.GetRetentionDays() != 0 { + t.Fatalf("retention days = %d, want forever (0)", logger.GetRetentionDays()) + } +} + +func TestNewSQLiteLoggerPreservesConfiguredCleanupInterval(t *testing.T) { + logger, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: t.TempDir(), + RetentionDays: 30, + RetentionConfigured: true, + CleanupInterval: 6 * time.Hour, + }) + if err != nil { + t.Fatalf("NewSQLiteLogger failed: %v", err) + } + defer logger.Close() + + if logger.cleanupInterval != 6*time.Hour { + t.Fatalf("cleanup interval = %s, want 6h", logger.cleanupInterval) + } +} + func TestIsPersistentLoggerUnwrapsAsyncConsoleBackend(t *testing.T) { console := NewConsoleLogger() if IsPersistentLogger(console) { @@ -566,6 +599,33 @@ func TestSQLiteLoggerSetRetentionDays(t *testing.T) { } } +func TestSQLiteLoggerLoadsPersistedRetentionDays(t *testing.T) { + dataDir := t.TempDir() + first, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: dataDir, + CryptoMgr: newMockCryptoManager(), + }) + if err != nil { + t.Fatalf("NewSQLiteLogger failed: %v", err) + } + first.SetRetentionDays(45) + if err := first.Close(); err != nil { + t.Fatalf("close first logger: %v", err) + } + + restarted, err := NewSQLiteLogger(SQLiteLoggerConfig{ + DataDir: dataDir, + CryptoMgr: newMockCryptoManager(), + }) + if err != nil { + t.Fatalf("restart logger: %v", err) + } + defer restarted.Close() + if restarted.GetRetentionDays() != 45 { + t.Fatalf("persisted retention days = %d, want 45", restarted.GetRetentionDays()) + } +} + func TestSQLiteLoggerCloseIsIdempotent(t *testing.T) { tempDir := t.TempDir() diff --git a/pkg/extensions/audit_admin.go b/pkg/extensions/audit_admin.go index 79364fb2a..28f3cb235 100644 --- a/pkg/extensions/audit_admin.go +++ b/pkg/extensions/audit_admin.go @@ -3,6 +3,7 @@ package extensions import ( "context" "net/http" + "time" "github.com/rcourtman/pulse-go-rewrite/pkg/audit" ) @@ -29,6 +30,20 @@ type AuditAdminRuntime struct { WriteError WriteAuditErrorFunc } +// AuditStoreConfig customizes the canonical SQLite audit store for a runtime. +// SigningKey is externally managed key material and must never be logged. +type AuditStoreConfig struct { + Directory string + SigningKey []byte + RetentionDays int + RetentionConfigured bool + CleanupInterval time.Duration +} + +// ResolveAuditStoreConfigFunc resolves runtime-specific audit persistence +// settings after the server has established its canonical base data directory. +type ResolveAuditStoreConfigFunc func(baseDataDir string) AuditStoreConfig + // BindAuditAdminEndpointsFunc allows enterprise modules to bind replacement // audit admin endpoints while retaining access to default handlers. type BindAuditAdminEndpointsFunc func(defaults AuditAdminEndpoints, runtime AuditAdminRuntime) AuditAdminEndpoints diff --git a/pkg/server/server.go b/pkg/server/server.go index 70b7ad7ff..c47dc04a5 100644 --- a/pkg/server/server.go +++ b/pkg/server/server.go @@ -58,6 +58,10 @@ type BusinessHooks struct { // audit admin endpoints without importing internal API packages. BindAuditAdminEndpoints extensions.BindAuditAdminEndpointsFunc + // ResolveAuditStoreConfig allows enterprise runtimes to configure the + // canonical audit store without opening a competing database connection. + ResolveAuditStoreConfig extensions.ResolveAuditStoreConfigFunc + // BindSSOAdminEndpoints allows enterprise modules to replace or decorate // SSO admin endpoints without importing internal API packages. BindSSOAdminEndpoints extensions.BindSSOAdminEndpointsFunc @@ -117,6 +121,7 @@ func SetBusinessHooks(h BusinessHooks) { func runtimeIdentityForBusinessHooks(h BusinessHooks) pkglicensing.RuntimeIdentity { if h.BindAuditAdminEndpoints != nil || + h.ResolveAuditStoreConfig != nil || h.BindRBACAdminEndpoints != nil || h.BindSSOAdminEndpoints != nil || h.BindReportingAdminEndpoints != nil || @@ -214,19 +219,31 @@ func Run(ctx context.Context, version string) error { // Always capture audit events to SQLite (defense in depth). Read/export endpoints are license-gated. // For the default org, TenantLoggerManager routes to the global logger, so initialize it as SQLite too. + globalHooksMu.Lock() + resolveAuditStoreConfig := globalHooks.ResolveAuditStoreConfig + globalHooksMu.Unlock() + resolvedAuditStore := extensions.AuditStoreConfig{} + if resolveAuditStoreConfig != nil { + resolvedAuditStore = resolveAuditStoreConfig(baseDataDir) + } var globalCrypto audit.CryptoEncryptor if cm, err := crypto.NewCryptoManagerAt(baseDataDir); err != nil { log.Warn().Err(err).Str("data_dir", baseDataDir).Msg("Failed to initialize crypto manager for audit signing; signatures will be disabled") } else { globalCrypto = cm - if sqliteLogger, err := audit.NewSQLiteLogger(audit.SQLiteLoggerConfig{ - DataDir: baseDataDir, - CryptoMgr: cm, - }); err != nil { - log.Warn().Err(err).Str("data_dir", baseDataDir).Msg("Failed to initialize global SQLite audit logger; falling back to console logger") - } else { - audit.SetLogger(sqliteLogger) - } + } + if sqliteLogger, err := audit.NewSQLiteLogger(audit.SQLiteLoggerConfig{ + DataDir: baseDataDir, + AuditDir: resolvedAuditStore.Directory, + CryptoMgr: globalCrypto, + SigningKey: resolvedAuditStore.SigningKey, + RetentionDays: resolvedAuditStore.RetentionDays, + RetentionConfigured: resolvedAuditStore.RetentionConfigured, + CleanupInterval: resolvedAuditStore.CleanupInterval, + }); err != nil { + log.Warn().Err(err).Str("data_dir", baseDataDir).Msg("Failed to initialize global SQLite audit logger; falling back to console logger") + } else { + audit.SetLogger(sqliteLogger) } // Initialize tenant audit manager for per-tenant audit logging @@ -236,7 +253,10 @@ func Run(ctx context.Context, version string) error { return crypto.NewCryptoManagerAt(dataDir) }, // Fallback for environments where per-tenant crypto initialization fails. - CryptoMgr: globalCrypto, + CryptoMgr: globalCrypto, + RetentionDays: resolvedAuditStore.RetentionDays, + RetentionConfigured: resolvedAuditStore.RetentionConfigured, + CleanupInterval: resolvedAuditStore.CleanupInterval, }) api.SetTenantAuditManager(tenantAuditManager) log.Info().Msg("Tenant audit manager initialized") diff --git a/pkg/server/server_test.go b/pkg/server/server_test.go index 8609d438d..188acb5b7 100644 --- a/pkg/server/server_test.go +++ b/pkg/server/server_test.go @@ -98,6 +98,15 @@ func TestRuntimeIdentityForBusinessHooks(t *testing.T) { t.Fatalf("enterprise hooks runtime build=%q, want pro", got.Build) } + got = runtimeIdentityForBusinessHooks(BusinessHooks{ + ResolveAuditStoreConfig: func(string) extensions.AuditStoreConfig { + return extensions.AuditStoreConfig{} + }, + }) + if got.Build != pkglicensing.RuntimeBuildPro { + t.Fatalf("audit store config hook runtime build=%q, want pro", got.Build) + } + got = runtimeIdentityForBusinessHooks(BusinessHooks{ ResolveMonitoredSystemAdmissionPolicy: func(context.Context, extensions.MonitoredSystemAdmissionInput) extensions.MonitoredSystemAdmissionDecision { return extensions.MonitoredSystemAdmissionDecision{} diff --git a/tests/integration/tests/journeys/06-audit-log-reporting.spec.ts b/tests/integration/tests/journeys/06-audit-log-reporting.spec.ts index 172333ac7..96c477b4e 100644 --- a/tests/integration/tests/journeys/06-audit-log-reporting.spec.ts +++ b/tests/integration/tests/journeys/06-audit-log-reporting.spec.ts @@ -105,9 +105,29 @@ test.describe.serial( ).toBeTruthy(); const body = await res.json(); - expect(body).toHaveProperty('events'); + expect(Array.isArray(body.events)).toBe(true); expect(body).toHaveProperty('total'); expect(body).toHaveProperty('persistentLogging'); + expect(typeof body.total).toBe('number'); + expect(typeof body.persistentLogging).toBe('boolean'); + + if (body.persistentLogging && body.events.length > 0) { + const signed = body.events.find( + (event: { signature?: string }) => event.signature, + ); + expect(signed, 'Persistent audit rows should be signed').toBeTruthy(); + const verifyRes = await apiRequest( + page, + `/api/audit/${signed.id}/verify`, + ); + expect( + verifyRes.ok(), + `Audit signature verification failed: ${verifyRes.status()}`, + ).toBeTruthy(); + const verification = await verifyRes.json(); + expect(verification.available).toBe(true); + expect(verification.verified).toBe(true); + } }); test('audit export endpoint returns a file', async ({ page }, testInfo) => { @@ -118,12 +138,6 @@ test.describe.serial( const res = await apiRequest(page, '/api/audit/export?format=json'); - // 501 = licensed but no persistent logger (OSS backend). - if (res.status() === 501) { - // This is expected on community/dev instances. - return; - } - expect( res.ok(), `Audit export failed: ${res.status()}`, @@ -147,11 +161,6 @@ test.describe.serial( const res = await apiRequest(page, '/api/audit/summary'); - // 501 = licensed but no persistent logger. - if (res.status() === 501) { - return; - } - expect( res.ok(), `Audit summary failed: ${res.status()}`, diff --git a/tests/integration/tests/journeys/07-audit-log-resilience.spec.ts b/tests/integration/tests/journeys/07-audit-log-resilience.spec.ts new file mode 100644 index 000000000..64cd866c5 --- /dev/null +++ b/tests/integration/tests/journeys/07-audit-log-resilience.spec.ts @@ -0,0 +1,164 @@ +import { test, expect, ensureJourneyReady } from "./journeyAuth"; + +const capabilities = { + capabilities: ["audit_logging"], + limits: [], + hosted_mode: false, + max_history_days: 90, + runtime: { + build: "pro", + label: "Pulse Pro runtime", + }, + blocked_capabilities: [], +}; + +const event = (id: string, user = "operator") => ({ + id, + timestamp: "2026-07-24T09:00:00Z", + event: "login", + user, + ip: "127.0.0.1", + path: "/api/auth", + success: true, + details: id, +}); + +test("audit log fails closed, recovers, pages atomically, and ignores stale responses", async ({ + page, +}, testInfo) => { + test.skip( + testInfo.project.name.startsWith("mobile-"), + "Desktop audit resilience journey", + ); + + let mode: + "initial" | "busy" | "recovered" | "page-size" | "race" | "invalid" = + "initial"; + const auditRequests: URL[] = []; + + await page.route("**/api/license/runtime-capabilities", async (route) => { + await route.fulfill({ + status: 200, + contentType: "application/json", + body: JSON.stringify(capabilities), + }); + }); + await page.route("**/api/audit?*", async (route) => { + const url = new URL(route.request().url()); + auditRequests.push(url); + if (mode === "busy") { + await route.fulfill({ + status: 503, + headers: { "Retry-After": "2" }, + contentType: "application/json", + body: JSON.stringify({ + code: "audit_store_busy", + error: "Audit log storage is temporarily busy. Try again shortly.", + }), + }); + return; + } + if (mode === "invalid") { + await route.fulfill({ + status: 200, + contentType: "application/json", + body: JSON.stringify({ + events: null, + total: 0, + persistentLogging: true, + }), + }); + return; + } + if (mode === "race" && !url.searchParams.get("user")) { + await new Promise((resolve) => setTimeout(resolve, 850)); + await route.fulfill({ + status: 200, + contentType: "application/json", + body: JSON.stringify({ + events: [event("stale-row")], + total: 1, + persistentLogging: true, + }), + }); + return; + } + if (mode === "race") { + await route.fulfill({ + status: 200, + contentType: "application/json", + body: JSON.stringify({ + events: [event("latest-row", "alice")], + total: 1, + persistentLogging: true, + }), + }); + return; + } + if (mode === "page-size") { + const size = Number(url.searchParams.get("limit")); + await route.fulfill({ + status: 200, + contentType: "application/json", + body: JSON.stringify({ + events: Array.from({ length: size }, (_, index) => + event(`page-${index}`), + ), + total: 250, + persistentLogging: true, + }), + }); + return; + } + const id = mode === "recovered" ? "recovered-row" : "initial-row"; + await route.fulfill({ + status: 200, + contentType: "application/json", + body: JSON.stringify({ + events: [event(id)], + total: 1, + persistentLogging: true, + }), + }); + }); + + await ensureJourneyReady(page); + await page.goto("/settings/security-audit", { + waitUntil: "domcontentloaded", + }); + await expect(page.getByText("initial-row", { exact: true })).toBeVisible(); + + mode = "busy"; + await page.getByRole("button", { name: "Refresh" }).click(); + await expect(page.locator("main").getByRole("alert")).toContainText("storage is busy"); + await expect(page.getByText("initial-row", { exact: true })).toHaveCount(0); + + mode = "recovered"; + await page.getByRole("button", { name: "Refresh" }).click(); + await expect(page.getByText("recovered-row", { exact: true })).toBeVisible(); + await expect(page.locator("main").getByRole("alert")).toHaveCount(0); + + mode = "page-size"; + await page.getByLabel("Audit page size").selectOption("25"); + await expect(page.getByText("Showing 1-25 of 250")).toBeVisible(); + expect( + auditRequests.some( + (url) => + url.searchParams.get("limit") === "25" && + url.searchParams.get("offset") === "0", + ), + ).toBe(true); + + mode = "race"; + await page.getByRole("button", { name: "Refresh" }).click(); + await page.getByPlaceholder("Filter by user...").fill("alice"); + await expect(page.getByText("latest-row", { exact: true })).toBeVisible(); + await page.waitForTimeout(700); + await expect(page.getByText("latest-row", { exact: true })).toBeVisible(); + await expect(page.getByText("stale-row", { exact: true })).toHaveCount(0); + + mode = "invalid"; + await page.getByRole("button", { name: "Refresh" }).click(); + await expect(page.locator("main").getByRole("alert")).toContainText("invalid response"); + await expect(page.getByText("latest-row", { exact: true })).toHaveCount(0); +}); diff --git a/tests/migration/v5_session_db_test.go b/tests/migration/v5_session_db_test.go index 9282605fc..78be86d40 100644 --- a/tests/migration/v5_session_db_test.go +++ b/tests/migration/v5_session_db_test.go @@ -449,7 +449,7 @@ func TestV5DataDir_AuditDBSchemaAutoMigration(t *testing.T) { assert.Equal(t, "admin", events[0].User) assert.True(t, events[0].Success) - // Verify the schema_version table exists and has version 1 + // Verify the schema_version table records canonical timestamp storage. dbPath := filepath.Join(dataDir, "audit", "audit.db") dsn := dbPath + "?" + url.Values{ "_pragma": []string{"busy_timeout(5000)"}, @@ -460,7 +460,7 @@ func TestV5DataDir_AuditDBSchemaAutoMigration(t *testing.T) { var version int require.NoError(t, db.QueryRow("SELECT version FROM schema_version ORDER BY version DESC LIMIT 1").Scan(&version)) - assert.Equal(t, 1, version, "schema_version should be 1") + assert.Equal(t, 2, version, "schema_version should be 2") } // TestV5DataDir_AuditDBPreExistingData verifies that the v6 audit logger @@ -484,14 +484,15 @@ func TestV5DataDir_AuditDBPreExistingData(t *testing.T) { _, err = rawDB.Exec(` CREATE TABLE IF NOT EXISTS audit_events ( id TEXT PRIMARY KEY, - timestamp INTEGER NOT NULL, + timestamp DATETIME NOT NULL, event_type TEXT NOT NULL, user TEXT, ip TEXT, path TEXT, success INTEGER NOT NULL, details TEXT, - signature TEXT NOT NULL + signature TEXT, + created_at DATETIME DEFAULT CURRENT_TIMESTAMP ); CREATE TABLE IF NOT EXISTS audit_config ( key TEXT PRIMARY KEY, @@ -502,14 +503,22 @@ func TestV5DataDir_AuditDBPreExistingData(t *testing.T) { require.NoError(t, err) // Insert v5 audit events - ts := time.Now().Unix() + ts := time.Now().UTC().Truncate(time.Second) _, err = rawDB.Exec(`INSERT INTO audit_events (id, timestamp, event_type, user, ip, path, success, details, signature) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, "v5-event-001", ts, "login", "admin", "192.168.1.1", "/api/auth/login", 1, "successful login", "v5-sig-placeholder") require.NoError(t, err) _, err = rawDB.Exec(`INSERT INTO audit_events (id, timestamp, event_type, user, ip, path, success, details, signature) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, - "v5-event-002", ts-3600, "config_change", "admin", "10.0.0.1", "/api/settings", 1, "changed polling interval", "v5-sig-placeholder-2") + "v5-event-002", + ts.Add(-time.Hour).Format("2006-01-02 15:04:05 -0700 MST")+" m=+0.009025344", + "config_change", + "admin", + "10.0.0.1", + "/api/settings", + 1, + "changed polling interval", + nil) require.NoError(t, err) rawDB.Close() @@ -533,6 +542,7 @@ func TestV5DataDir_AuditDBPreExistingData(t *testing.T) { require.NoError(t, err) require.Len(t, events2, 1) assert.Equal(t, "config_change", events2[0].EventType) + assert.True(t, events2[0].Timestamp.Equal(ts.Add(-time.Hour))) // Verify v6 can write new events alongside v5 data newEvent := audit.Event{ @@ -550,4 +560,13 @@ func TestV5DataDir_AuditDBPreExistingData(t *testing.T) { allEvents, err := logger.Query(audit.QueryFilter{Limit: 100}) require.NoError(t, err) assert.GreaterOrEqual(t, len(allEvents), 3, "v5 + v6 events must coexist") + + schemaDB, err := sql.Open("sqlite", dsn) + require.NoError(t, err) + defer schemaDB.Close() + var timestampType string + require.NoError(t, schemaDB.QueryRow(` + SELECT type FROM pragma_table_info('audit_events') WHERE name = 'timestamp' + `).Scan(×tampType)) + assert.Equal(t, "INTEGER", timestampType) }