Fix audit storage migration and viewer races

Refs #1464
This commit is contained in:
rcourtman
2026-07-24 10:32:43 +01:00
parent 49217d284d
commit 2ac70dab70
26 changed files with 1938 additions and 165 deletions
@@ -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
@@ -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`,
@@ -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 <scope>` option wording wherever a product surface exposes filter selects
or segmented filter choices. Workloads filters, storage source
@@ -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"
]
}
],
@@ -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
@@ -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:
@@ -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<typeof import('@/utils/apiClient')>();
return {
...actual,
apiFetch: vi.fn(),
apiErrorFromResponse: vi.fn(),
};
});
type Deferred<T> = {
promise: Promise<T>;
resolve: (value: T) => void;
};
const deferred = <T,>(): Deferred<T> => {
let resolve!: (value: T) => void;
const promise = new Promise<T>((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<Response>();
const second = deferred<Response>();
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();
});
});
@@ -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<Record<string, AbortController>>(
{},
);
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;
+81 -51
View File
@@ -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 {
+109 -3
View File
@@ -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"),
+16
View File
@@ -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 {
+41
View File
@@ -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")
}
}
+25
View File
@@ -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
}
+61 -6
View File
@@ -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
+61
View File
@@ -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()
+11 -2
View File
@@ -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 {
+408 -41
View File
@@ -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 <DataDir>/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,
&notNull,
&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
}
+465
View File
@@ -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,
&notNull,
&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)
}
}
+5 -4
View File
@@ -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);
+60
View File
@@ -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()
+15
View File
@@ -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
+29 -9
View File
@@ -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")
+9
View File
@@ -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{}
@@ -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()}`,
@@ -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);
});
+25 -6
View File
@@ -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(&timestampType))
assert.Equal(t, "INTEGER", timestampType)
}