fix(fleet-sync): make control-identity-mismatch sticky and surface in UI (#1117)

Treat 409 CONTROL_IDENTITY_MISMATCH from a replica as a non-retriable
failure instead of looping the same 409 through the 5-minute retry
service forever and silently writing identical failure rows.

Backend
- DatabaseService: add `sticky_error_code`, `sticky_error_expected`,
  `sticky_error_got` columns to `fleet_sync_status` via an idempotent
  migration. New methods setFleetSyncSticky, getFleetSyncStickyCode,
  clearFleetSyncStickyForNode. recordFleetSyncSuccess clears the sticky
  flag on a clean push. getFailedSyncTargets SQL adds
  `AND sticky_error_code IS NULL` so the retry loop skips sticky rows.
- FleetSyncService.executePushToNode: short-circuits at the top when
  sticky is set (covers event-driven pushResourceAsync calls). On a 409
  with code CONTROL_IDENTITY_MISMATCH, records the failure once and
  pins sticky with the expected/got fingerprints carried in the 409 body.
- routes/nodes.ts: new POST /api/nodes/:id/fleet-sync/reset-anchor.
  Admin + paid + node:manage. Proxies POST /api/fleet/role/reanchor to
  the peer with `{override:true}` using the stored Bearer node_proxy
  token. On peer 200, clears every sticky row for the node so the next
  push re-anchors and resumes replication. Distinct 502 / 504 responses
  for peer-rejected / peer-unreachable so the UI can show a useful toast.

Frontend
- New lib/fleetSyncApi.ts + hooks/useFleetSyncStatus.ts. Polling hook
  (30s visibilityInterval) skips fetch when !isPaid.
- NodeManager.tsx: destructive banner per affected node listing both
  fingerprints, with `Reset anchor on peer` and `Remove node` buttons.
  Hidden for community-tier users via empty hook data.
- FleetConfiguration.tsx (Fleet -> Status): read-only `Policy sync`
  SummaryRow per remote node card. In sync / degraded / paused with
  a tooltip; no action buttons (the action lives in NodeManager).

Tests
- fleet-sync-service.test.ts: 4 new cases for sticky-set on first
  mismatch, short-circuit on subsequent pushes, null fingerprints,
  and non-mismatch failures not setting sticky.
- database-fleet-sync-sticky.test.ts (new): 6 cases pinning the DB
  contract incl. retry-loop SQL filter and migration idempotency.
- nodes-fleet-sync-reset-anchor.test.ts (new): 6 cases covering
  happy path, peer 401 -> 502, peer unreachable -> 504, local-node
  rejection, unknown node id, and community-tier 403.

Gate parity (Directive 30): the new POST .../reset-anchor enforces
requireAdmin + requirePaid + node:manage (matches the existing read at
GET /api/fleet/sync-status). UI banner + SummaryRow only render when the
hook returns data, which it only does for paid-tier authed users. No
existing tier-gate file moved; this is greenfield parity.

Auth audit: the peer's POST /api/fleet/role/reanchor route already uses
requireAdmin, which accepts the central's stored node_proxy Bearer
token because authMiddleware maps `scope === 'node_proxy'` to
`req.user = { username: 'node-proxy', role: 'admin', userId: 0 }`.
No widening required.

Backend tsc clean. Frontend tsc -b clean. 59 fleet-sync tests pass; full
backend suite green minus the pre-existing Windows-only file-lock flake
on filesystem-backup.test.ts that reproduces unchanged on main.
This commit is contained in:
Anso
2026-05-19 19:49:16 -04:00
committed by GitHub
parent 69bc955c3b
commit e05099f2a1
10 changed files with 876 additions and 11 deletions
@@ -0,0 +1,129 @@
/**
* Pins the DatabaseService sticky-error wiring used by the F-16 fix:
* - setFleetSyncSticky writes the code + expected + got fingerprints.
* - getFleetSyncStickyCode reads them back.
* - getFailedSyncTargets excludes sticky rows (the retry loop must not pick them up).
* - recordFleetSyncSuccess clears the sticky on success (operator reset → push resumes).
* - clearFleetSyncStickyForNode clears every resource for one node id (used by the reset endpoint).
*/
import { describe, it, expect, beforeAll, afterAll, beforeEach } from 'vitest';
import { setupTestDb, cleanupTestDb } from './helpers/setupTestDb';
let tmpDir: string;
let DatabaseService: typeof import('../services/DatabaseService').DatabaseService;
let nodeId: number;
let siblingId: number;
beforeAll(async () => {
tmpDir = await setupTestDb();
({ DatabaseService } = await import('../services/DatabaseService'));
const db = DatabaseService.getInstance();
nodeId = db.addNode({
name: 'sticky-target',
type: 'remote',
compose_dir: '/app/compose',
is_default: false,
api_url: 'https://sticky.example',
api_token: 'tok',
mode: 'proxy',
});
siblingId = db.addNode({
name: 'sticky-sibling',
type: 'remote',
compose_dir: '/app/compose',
is_default: false,
api_url: 'https://sibling-sticky.example',
api_token: 'tok',
mode: 'proxy',
});
});
afterAll(() => {
cleanupTestDb(tmpDir);
});
beforeEach(() => {
const db = DatabaseService.getInstance();
// Wipe stale rows from prior tests in this file so each case starts clean.
db.getDb().prepare('DELETE FROM fleet_sync_status WHERE node_id IN (?, ?)').run(nodeId, siblingId);
});
describe('fleet_sync_status sticky-error column', () => {
it('setFleetSyncSticky persists the code and fingerprints', () => {
const db = DatabaseService.getInstance();
db.setFleetSyncSticky(nodeId, 'scan_policies', 'CONTROL_IDENTITY_MISMATCH', 'aaa111', 'bbb222');
const row = db.getFleetSyncStatuses().find(
(s) => s.node_id === nodeId && s.resource === 'scan_policies',
);
expect(row).toBeDefined();
expect(row!.sticky_error_code).toBe('CONTROL_IDENTITY_MISMATCH');
expect(row!.sticky_error_expected).toBe('aaa111');
expect(row!.sticky_error_got).toBe('bbb222');
expect(db.getFleetSyncStickyCode(nodeId, 'scan_policies')).toBe('CONTROL_IDENTITY_MISMATCH');
});
it('setFleetSyncSticky upserts when no row exists yet', () => {
const db = DatabaseService.getInstance();
// Pre-state: no row.
expect(db.getFleetSyncStickyCode(nodeId, 'cve_suppressions')).toBeNull();
db.setFleetSyncSticky(nodeId, 'cve_suppressions', 'CONTROL_IDENTITY_MISMATCH', null, null);
expect(db.getFleetSyncStickyCode(nodeId, 'cve_suppressions')).toBe('CONTROL_IDENTITY_MISMATCH');
});
it('getFailedSyncTargets excludes rows where sticky_error_code is set', () => {
const db = DatabaseService.getInstance();
db.recordFleetSyncFailure(nodeId, 'scan_policies', 'timeout');
db.recordFleetSyncFailure(siblingId, 'scan_policies', 'connection refused');
// Mark only `nodeId` as sticky; the sibling stays retriable.
db.setFleetSyncSticky(nodeId, 'scan_policies', 'CONTROL_IDENTITY_MISMATCH', null, null);
const retriable = db.getFailedSyncTargets('scan_policies', 24 * 60 * 60_000);
const retriableIds = retriable.map((r) => r.node_id);
expect(retriableIds).toContain(siblingId);
expect(retriableIds).not.toContain(nodeId);
});
it('recordFleetSyncSuccess clears the sticky flag (operator-reset round-trip)', () => {
const db = DatabaseService.getInstance();
db.setFleetSyncSticky(nodeId, 'scan_policies', 'CONTROL_IDENTITY_MISMATCH', 'aaa', 'bbb');
expect(db.getFleetSyncStickyCode(nodeId, 'scan_policies')).toBe('CONTROL_IDENTITY_MISMATCH');
db.recordFleetSyncSuccess(nodeId, 'scan_policies');
expect(db.getFleetSyncStickyCode(nodeId, 'scan_policies')).toBeNull();
const row = db.getFleetSyncStatuses().find(
(s) => s.node_id === nodeId && s.resource === 'scan_policies',
);
expect(row!.sticky_error_expected).toBeNull();
expect(row!.sticky_error_got).toBeNull();
});
it('clearFleetSyncStickyForNode clears every resource for one node, leaves siblings untouched', () => {
const db = DatabaseService.getInstance();
db.setFleetSyncSticky(nodeId, 'scan_policies', 'CONTROL_IDENTITY_MISMATCH', null, null);
db.setFleetSyncSticky(nodeId, 'cve_suppressions', 'CONTROL_IDENTITY_MISMATCH', null, null);
db.setFleetSyncSticky(siblingId, 'scan_policies', 'CONTROL_IDENTITY_MISMATCH', null, null);
db.clearFleetSyncStickyForNode(nodeId);
expect(db.getFleetSyncStickyCode(nodeId, 'scan_policies')).toBeNull();
expect(db.getFleetSyncStickyCode(nodeId, 'cve_suppressions')).toBeNull();
expect(db.getFleetSyncStickyCode(siblingId, 'scan_policies')).toBe('CONTROL_IDENTITY_MISMATCH');
});
it('migrateFleetSyncStickyError is idempotent (running twice does not error)', () => {
// The constructor already runs the migration once at boot. Manually
// invoke the private method twice via index access to confirm
// tryAddColumn's idempotency contract holds for this migration.
const db = DatabaseService.getInstance() as unknown as {
migrateFleetSyncStickyError: () => void;
};
expect(() => {
db.migrateFleetSyncStickyError();
db.migrateFleetSyncStickyError();
}).not.toThrow();
});
});
@@ -16,6 +16,8 @@ const {
mockInsertAuditLog,
mockRecordFleetSyncSuccess,
mockRecordFleetSyncFailure,
mockSetFleetSyncSticky,
mockGetFleetSyncStickyCode,
mockGetSystemState,
mockSetSystemState,
mockTransaction,
@@ -33,6 +35,8 @@ const {
mockInsertAuditLog: vi.fn(),
mockRecordFleetSyncSuccess: vi.fn(),
mockRecordFleetSyncFailure: vi.fn(),
mockSetFleetSyncSticky: vi.fn(),
mockGetFleetSyncStickyCode: vi.fn().mockReturnValue(null),
mockGetSystemState: vi.fn().mockReturnValue(null),
mockSetSystemState: vi.fn(),
mockTransaction: vi.fn().mockImplementation((fn: () => unknown) => fn()),
@@ -53,6 +57,8 @@ vi.mock('../services/DatabaseService', () => ({
insertAuditLog: mockInsertAuditLog,
recordFleetSyncSuccess: mockRecordFleetSyncSuccess,
recordFleetSyncFailure: mockRecordFleetSyncFailure,
setFleetSyncSticky: mockSetFleetSyncSticky,
getFleetSyncStickyCode: mockGetFleetSyncStickyCode,
getSystemState: mockGetSystemState,
setSystemState: mockSetSystemState,
transaction: mockTransaction,
@@ -90,6 +96,7 @@ import { FleetSyncService, LOCAL_IDENTITY_SENTINEL } from '../services/FleetSync
beforeEach(() => {
vi.clearAllMocks();
mockGetSystemState.mockReturnValue(null);
mockGetFleetSyncStickyCode.mockReturnValue(null);
});
describe('FleetSyncService.getRole', () => {
@@ -594,3 +601,102 @@ describe('FleetSyncService.formatError redaction', () => {
expect(failure[2]).toContain('[redacted-jwt]');
});
});
describe('FleetSyncService CONTROL_IDENTITY_MISMATCH sticky handling', () => {
function makeMismatchError(expected: string, got: string) {
return async () => {
const { AxiosError } = await import('axios');
const err = new AxiosError('Request failed with status code 409');
(err as unknown as { response: unknown }).response = {
status: 409,
statusText: 'Conflict',
data: {
error: `Control identity mismatch: replica is anchored to "${expected}", push from "${got}"`,
code: 'CONTROL_IDENTITY_MISMATCH',
expected,
got,
},
};
throw err;
};
}
it('sets the sticky flag carrying the expected/got fingerprints on first mismatch', async () => {
mockGetNodes.mockReturnValue([
{ id: 7, type: 'remote', api_url: 'https://peer.example', api_token: 'tok', name: 'peer', mode: 'proxy' },
]);
mockGetLocalScanPolicies.mockReturnValue([]);
mockGetFleetSyncStickyCode.mockReturnValue(null);
mockAxiosPost.mockImplementation(makeMismatchError('cb45a2eff9db81d8', '555f8d1f7e7e71e3'));
await FleetSyncService.getInstance().pushResource('scan_policies');
expect(mockRecordFleetSyncFailure).toHaveBeenCalledTimes(1);
expect(mockSetFleetSyncSticky).toHaveBeenCalledTimes(1);
expect(mockSetFleetSyncSticky).toHaveBeenCalledWith(
7,
'scan_policies',
'CONTROL_IDENTITY_MISMATCH',
'cb45a2eff9db81d8',
'555f8d1f7e7e71e3',
);
});
it('short-circuits subsequent pushes when sticky is already set; no HTTP call, no failure record', async () => {
mockGetNodes.mockReturnValue([
{ id: 7, type: 'remote', api_url: 'https://peer.example', api_token: 'tok', name: 'peer', mode: 'proxy' },
]);
mockGetLocalScanPolicies.mockReturnValue([]);
// Sticky already set from a prior push.
mockGetFleetSyncStickyCode.mockReturnValue('CONTROL_IDENTITY_MISMATCH');
await FleetSyncService.getInstance().pushResource('scan_policies');
expect(mockAxiosPost).not.toHaveBeenCalled();
expect(mockRecordFleetSyncFailure).not.toHaveBeenCalled();
expect(mockRecordFleetSyncSuccess).not.toHaveBeenCalled();
expect(mockSetFleetSyncSticky).not.toHaveBeenCalled();
});
it('tolerates a missing expected/got payload (passes null through)', async () => {
mockGetNodes.mockReturnValue([
{ id: 7, type: 'remote', api_url: 'https://peer.example', api_token: 'tok', name: 'peer', mode: 'proxy' },
]);
mockGetLocalScanPolicies.mockReturnValue([]);
mockGetFleetSyncStickyCode.mockReturnValue(null);
mockAxiosPost.mockImplementation(async () => {
const { AxiosError } = await import('axios');
const err = new AxiosError('Request failed with status code 409');
(err as unknown as { response: unknown }).response = {
status: 409,
statusText: 'Conflict',
data: { error: 'mismatch', code: 'CONTROL_IDENTITY_MISMATCH' },
};
throw err;
});
await FleetSyncService.getInstance().pushResource('scan_policies');
expect(mockSetFleetSyncSticky).toHaveBeenCalledWith(
7,
'scan_policies',
'CONTROL_IDENTITY_MISMATCH',
null,
null,
);
});
it('does not set sticky for non-mismatch failures (network errors, 500s)', async () => {
mockGetNodes.mockReturnValue([
{ id: 7, type: 'remote', api_url: 'https://peer.example', api_token: 'tok', name: 'peer', mode: 'proxy' },
]);
mockGetLocalScanPolicies.mockReturnValue([]);
mockGetFleetSyncStickyCode.mockReturnValue(null);
mockAxiosPost.mockRejectedValue(new Error('ECONNREFUSED'));
await FleetSyncService.getInstance().pushResource('scan_policies');
expect(mockRecordFleetSyncFailure).toHaveBeenCalledTimes(1);
expect(mockSetFleetSyncSticky).not.toHaveBeenCalled();
});
});
@@ -0,0 +1,158 @@
/**
* Tests for POST /api/nodes/:id/fleet-sync/reset-anchor (F-16 fix).
*
* The endpoint proxies the peer's reanchor endpoint and clears every
* sticky-error row for the node on success. Covers:
* - happy path: 200 from peer → sticky rows cleared, 200 returned.
* - peer 401/403 → 502 with helpful message.
* - peer unreachable → 504.
* - missing/non-proxy node → 400.
* - non-paid tier → 403.
*/
import { describe, it, expect, beforeAll, afterAll, beforeEach, vi } from 'vitest';
import request from 'supertest';
import jwt from 'jsonwebtoken';
import { setupTestDb, cleanupTestDb, TEST_USERNAME, TEST_JWT_SECRET } from './helpers/setupTestDb';
let tmpDir: string;
let app: import('express').Express;
let authHeader: string;
let peerNodeId: number;
const originalFetch = globalThis.fetch;
beforeAll(async () => {
tmpDir = await setupTestDb();
({ app } = await import('../index'));
const token = jwt.sign({ username: TEST_USERNAME }, TEST_JWT_SECRET, { expiresIn: '1m' });
authHeader = `Bearer ${token}`;
// Seed a proxy-mode remote node and a sticky row for it. Tests then drive
// the route handler and assert side effects on the test DB.
const { DatabaseService } = await import('../services/DatabaseService');
const db = DatabaseService.getInstance();
peerNodeId = db.addNode({
name: 'sticky-peer',
type: 'remote',
compose_dir: '/app/compose',
is_default: false,
api_url: 'http://192.168.1.99:1852',
api_token: 'peer-token',
mode: 'proxy',
});
db.setFleetSyncSticky(
peerNodeId,
'scan_policies',
'CONTROL_IDENTITY_MISMATCH',
'cb45a2eff9db81d8',
'555f8d1f7e7e71e3',
);
});
afterAll(() => {
globalThis.fetch = originalFetch;
cleanupTestDb(tmpDir);
});
beforeEach(async () => {
vi.restoreAllMocks();
globalThis.fetch = originalFetch;
// Re-establish the paid-tier spy after restoreAllMocks. Individual tests
// can override with `mockReturnValue('community')` to exercise the tier
// gate's deny path.
const { LicenseService } = await import('../services/LicenseService');
vi.spyOn(LicenseService.getInstance(), 'getTier').mockReturnValue('paid');
});
describe('POST /api/nodes/:id/fleet-sync/reset-anchor', () => {
it('proxies to the peer reanchor, clears sticky rows, returns 200', async () => {
const fetchSpy = vi.fn().mockResolvedValue(
new Response(JSON.stringify({ success: true }), { status: 200 }),
);
globalThis.fetch = fetchSpy as unknown as typeof fetch;
const res = await request(app)
.post(`/api/nodes/${peerNodeId}/fleet-sync/reset-anchor`)
.set('Authorization', authHeader)
.send({});
expect(res.status).toBe(200);
expect(res.body.success).toBe(true);
expect(fetchSpy).toHaveBeenCalledTimes(1);
const [url, init] = fetchSpy.mock.calls[0] as [string, RequestInit];
expect(url).toBe('http://192.168.1.99:1852/api/fleet/role/reanchor');
expect(init.method).toBe('POST');
expect((init.headers as Record<string, string>).Authorization).toBe('Bearer peer-token');
expect(JSON.parse(init.body as string)).toEqual({ override: true });
const { DatabaseService } = await import('../services/DatabaseService');
const sticky = DatabaseService.getInstance().getFleetSyncStickyCode(peerNodeId, 'scan_policies');
expect(sticky).toBeNull();
});
it('returns 502 with a helpful message when peer responds 401', async () => {
const { DatabaseService } = await import('../services/DatabaseService');
DatabaseService.getInstance().setFleetSyncSticky(
peerNodeId, 'cve_suppressions', 'CONTROL_IDENTITY_MISMATCH', 'aaa', 'bbb',
);
const fetchSpy = vi.fn().mockResolvedValue(
new Response(JSON.stringify({ error: 'Admin access required.' }), { status: 401 }),
);
globalThis.fetch = fetchSpy as unknown as typeof fetch;
const res = await request(app)
.post(`/api/nodes/${peerNodeId}/fleet-sync/reset-anchor`)
.set('Authorization', authHeader)
.send({});
expect(res.status).toBe(502);
expect(res.body.error).toMatch(/Admin access required/);
// Sticky rows must remain set so the operator can retry.
const sticky = DatabaseService.getInstance().getFleetSyncStickyCode(peerNodeId, 'cve_suppressions');
expect(sticky).toBe('CONTROL_IDENTITY_MISMATCH');
});
it('returns 504 when the peer is unreachable', async () => {
const fetchSpy = vi.fn().mockRejectedValue(new Error('connect ECONNREFUSED'));
globalThis.fetch = fetchSpy as unknown as typeof fetch;
const res = await request(app)
.post(`/api/nodes/${peerNodeId}/fleet-sync/reset-anchor`)
.set('Authorization', authHeader)
.send({});
expect(res.status).toBe(504);
expect(res.body.error).toMatch(/unreachable/i);
});
it('returns 400 for a local node', async () => {
const { DatabaseService } = await import('../services/DatabaseService');
const local = DatabaseService.getInstance().getNodes().find((n) => n.type === 'local');
expect(local).toBeTruthy();
const res = await request(app)
.post(`/api/nodes/${local!.id}/fleet-sync/reset-anchor`)
.set('Authorization', authHeader)
.send({});
expect(res.status).toBe(400);
});
it('returns 404 for an unknown node id', async () => {
const res = await request(app)
.post('/api/nodes/9999/fleet-sync/reset-anchor')
.set('Authorization', authHeader)
.send({});
expect(res.status).toBe(404);
});
it('returns 403 (PAID_REQUIRED) when the license is community-tier', async () => {
const { LicenseService } = await import('../services/LicenseService');
vi.spyOn(LicenseService.getInstance(), 'getTier').mockReturnValue('community');
const res = await request(app)
.post(`/api/nodes/${peerNodeId}/fleet-sync/reset-anchor`)
.set('Authorization', authHeader)
.send({});
expect(res.status).toBe(403);
expect(res.body.code).toBe('PAID_REQUIRED');
});
});
+88 -1
View File
@@ -4,7 +4,7 @@ import crypto from 'crypto';
import { authMiddleware } from '../middleware/auth';
import { requirePermission } from '../middleware/permissions';
import { rejectApiTokenScope } from '../middleware/apiTokenScope';
import { requireAdmiral } from '../middleware/tierGates';
import { requireAdmin, requireAdmiral, requirePaid } from '../middleware/tierGates';
import { enrollmentLimiter } from '../middleware/rateLimiters';
import { DatabaseService } from '../services/DatabaseService';
import { NodeRegistry } from '../services/NodeRegistry';
@@ -354,6 +354,93 @@ nodesRouter.post('/:id/uncordon', (req: Request, res: Response) => {
}
});
/**
* Reset the FleetSync control anchor on a remote peer.
*
* Proxies POST /api/fleet/role/reanchor to the peer using its stored
* Bearer token. A successful reanchor clears every sticky-error row for
* this node so the next push (event-driven or via the 5-minute retry
* service) re-attempts cleanly and the peer accepts the central's
* fingerprint as the new anchor.
*
* Surfaces UI affordance for the F-16 audit (mesh-e2e-2026-05-17.md):
* when a peer was previously enrolled by a different central, FleetSync
* keeps 409'ing every reconcile tick; the sticky flag halts retries and
* this endpoint is the single one-click recovery for the operator.
*/
nodesRouter.post('/:id/fleet-sync/reset-anchor', async (req: Request, res: Response) => {
if (rejectApiTokenScope(req, res, NODE_SCOPE_MESSAGE)) return;
const nodeIdParam = req.params.id as string;
if (!requirePermission(req, res, 'node:manage', 'node', nodeIdParam)) return;
if (!requirePaid(req, res)) return;
// Reset-anchor is symmetric with `/api/fleet/sync-status` (admin-only).
// Keeping read and write gated at the same role avoids a banner-invisible-to-the-actor
// gap where a node-admin could call reset without ever seeing why.
if (!requireAdmin(req, res)) return;
try {
const id = parseInt(nodeIdParam, 10);
if (!Number.isFinite(id) || id <= 0) {
res.status(400).json({ error: 'Invalid node id' });
return;
}
const node = DatabaseService.getInstance().getNode(id);
if (!node) {
res.status(404).json({ error: 'Node not found' });
return;
}
if (node.type !== 'remote' || node.mode !== 'proxy') {
res.status(400).json({ error: 'Reset anchor only applies to proxy-mode remote nodes' });
return;
}
if (!node.api_url || !node.api_token) {
res.status(400).json({ error: 'Node is missing api_url or api_token' });
return;
}
const baseUrl = node.api_url.replace(/\/$/, '');
let peerResponse: globalThis.Response;
try {
peerResponse = await fetch(`${baseUrl}/api/fleet/role/reanchor`, {
method: 'POST',
headers: {
Authorization: `Bearer ${node.api_token}`,
'Content-Type': 'application/json',
},
body: JSON.stringify({ override: true }),
signal: AbortSignal.timeout(15_000),
});
} catch (networkErr) {
const message = getErrorMessage(networkErr, 'Failed to reach peer');
console.warn(`[Nodes] Reset anchor unreachable for node ${id}: ${message}`);
res.status(504).json({ error: `Peer unreachable: ${message}` });
return;
}
if (!peerResponse.ok) {
const status = peerResponse.status;
const body = await peerResponse.json().catch(() => ({}));
const peerError = (body as { error?: string })?.error
?? `Peer returned HTTP ${status}`;
if (status === 401 || status === 403) {
console.warn(`[Nodes] Reset anchor rejected by peer ${id}: ${peerError}`);
res.status(502).json({
error: `Peer rejected the reanchor request: ${peerError}. The node's API token may need to be regenerated.`,
});
return;
}
res.status(502).json({ error: peerError });
return;
}
DatabaseService.getInstance().clearFleetSyncStickyForNode(id);
console.log(`[Nodes] Fleet-sync anchor reset on node ${id} ("${node.name}")`);
res.json({ success: true });
} catch (error: unknown) {
console.error('Failed to reset fleet-sync anchor:', error);
res.status(500).json({ error: getErrorMessage(error, 'Failed to reset fleet-sync anchor') });
}
});
nodesRouter.post('/:id/test', async (req: Request, res: Response) => {
try {
const id = parseInt(req.params.id as string);
+77 -4
View File
@@ -555,6 +555,9 @@ export interface FleetSyncStatus {
last_success_at: number | null;
last_failure_at: number | null;
last_error: string | null;
sticky_error_code: string | null;
sticky_error_expected: string | null;
sticky_error_got: string | null;
}
export interface CveSuppression {
@@ -659,6 +662,7 @@ export class DatabaseService {
this.migrateAddNodeCordonFields();
this.migrateAddBlueprintPinnedNode();
this.migrateAutoHealNodeId();
this.migrateFleetSyncStickyError();
// Reset the cache once at end of constructor in case any migration
// populated it via getGlobalSettings() and a subsequent migration
@@ -1581,6 +1585,12 @@ export class DatabaseService {
this.tryAddColumn('blueprints', 'pinned_node_id', 'INTEGER');
}
private migrateFleetSyncStickyError(): void {
this.tryAddColumn('fleet_sync_status', 'sticky_error_code', 'TEXT');
this.tryAddColumn('fleet_sync_status', 'sticky_error_expected', 'TEXT');
this.tryAddColumn('fleet_sync_status', 'sticky_error_got', 'TEXT');
}
private migrateAutoHealNodeId(): void {
const markerKey = 'migration_auto_heal_node_scope_v1';
const markerDone = this.getGlobalSettings()[markerKey] === '1';
@@ -3931,11 +3941,15 @@ export class DatabaseService {
const now = Date.now();
this.db
.prepare(
`INSERT INTO fleet_sync_status (node_id, resource, last_success_at, last_failure_at, last_error)
VALUES (?, ?, ?, NULL, NULL)
`INSERT INTO fleet_sync_status (node_id, resource, last_success_at, last_failure_at, last_error,
sticky_error_code, sticky_error_expected, sticky_error_got)
VALUES (?, ?, ?, NULL, NULL, NULL, NULL, NULL)
ON CONFLICT(node_id, resource) DO UPDATE SET
last_success_at = excluded.last_success_at,
last_error = NULL`,
last_error = NULL,
sticky_error_code = NULL,
sticky_error_expected = NULL,
sticky_error_got = NULL`,
)
.run(nodeId, resource, now);
}
@@ -3953,6 +3967,64 @@ export class DatabaseService {
.run(nodeId, resource, now, error);
}
/**
* Mark a (node, resource) pair as having hit a non-retriable failure. The
* retry service skips sticky rows and the push paths short-circuit before
* any HTTP call. The first such failure still records `last_failure_at` +
* `last_error` via `recordFleetSyncFailure`; the sticky write is additive.
*
* `expected` and `got` carry the fingerprints from a 409
* CONTROL_IDENTITY_MISMATCH response so the UI can render
* "anchored to <expected>, this central is <got>" without parsing the
* error string.
*/
public setFleetSyncSticky(
nodeId: number,
resource: string,
code: string,
expected: string | null,
got: string | null,
): void {
this.db
.prepare(
`INSERT INTO fleet_sync_status (node_id, resource, sticky_error_code,
sticky_error_expected, sticky_error_got)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT(node_id, resource) DO UPDATE SET
sticky_error_code = excluded.sticky_error_code,
sticky_error_expected = excluded.sticky_error_expected,
sticky_error_got = excluded.sticky_error_got`,
)
.run(nodeId, resource, code, expected, got);
}
public getFleetSyncStickyCode(nodeId: number, resource: string): string | null {
const row = this.db
.prepare(
`SELECT sticky_error_code FROM fleet_sync_status
WHERE node_id = ? AND resource = ?`,
)
.get(nodeId, resource) as { sticky_error_code: string | null } | undefined;
return row?.sticky_error_code ?? null;
}
/**
* Clear every sticky-error row for one node id. Used by the
* reset-anchor endpoint after the peer has acknowledged a reanchor; the
* next push attempt re-tries normally.
*/
public clearFleetSyncStickyForNode(nodeId: number): void {
this.db
.prepare(
`UPDATE fleet_sync_status
SET sticky_error_code = NULL,
sticky_error_expected = NULL,
sticky_error_got = NULL
WHERE node_id = ?`,
)
.run(nodeId);
}
public getFailedSyncTargets(resource: string, maxAgeMs: number): FleetSyncStatus[] {
const cutoff = Date.now() - maxAgeMs;
return this.db
@@ -3960,7 +4032,8 @@ export class DatabaseService {
`SELECT * FROM fleet_sync_status
WHERE resource = ?
AND (last_failure_at IS NOT NULL AND last_failure_at > ?)
AND (last_success_at IS NULL OR last_success_at < last_failure_at)`,
AND (last_success_at IS NULL OR last_success_at < last_failure_at)
AND sticky_error_code IS NULL`,
)
.all(resource, cutoff) as FleetSyncStatus[];
}
+47 -1
View File
@@ -449,12 +449,38 @@ export class FleetSyncService {
return next;
}
/**
* Concurrency note: the sticky-set in the CONTROL_IDENTITY_MISMATCH catch
* branch is best-effort against an operator-initiated reset that lands
* during a push's HTTP round-trip. The reset endpoint clears
* `sticky_error_code` to NULL; if this push's 409 arrives after that
* clear, it will re-pin the row. The operator clicks Reset again. The
* window is bounded by one HTTP round-trip per resource; no lock or
* generation counter is justified.
*/
private async executePushToNode(
node: Node & { id: number },
resource: FleetResource,
partial: Omit<FleetSyncPayload, 'targetIdentity'>,
): Promise<void> {
const db = DatabaseService.getInstance();
// Sticky-error short-circuit. When a previous push hit a non-retriable
// failure (today: CONTROL_IDENTITY_MISMATCH), every subsequent push
// would re-issue the same 409 and re-spam the log every 5 minutes. The
// sticky flag is cleared by either (a) the operator resetting the
// anchor via POST /api/nodes/:id/fleet-sync/reset-anchor, or (b) a
// successful push to the same node (recordFleetSyncSuccess clears
// it). Until then, skip the HTTP call entirely.
if (db.getFleetSyncStickyCode(node.id, resource)) {
if (isDebugEnabled()) {
console.debug(
`[FleetSync:debug] Skipping ${resource} push to "${node.name}": sticky error blocks retries.`,
);
}
return;
}
const apiUrl = node.api_url ?? '';
const baseUrl = apiUrl.replace(/\/$/, '');
const payload: FleetSyncPayload = { ...partial, targetIdentity: apiUrl };
@@ -479,7 +505,7 @@ export class FleetSyncService {
// healthy; suppress the failure record so it does not surface as
// an alert in the sync-status panel.
if (err instanceof AxiosError && err.response?.status === 409) {
const data = err.response.data as { code?: string } | undefined;
const data = err.response.data as { code?: string; expected?: string; got?: string } | undefined;
if (data?.code === SYNC_ERROR_CODES.staleSyncPush) {
if (isDebugEnabled()) {
console.debug(
@@ -488,6 +514,26 @@ export class FleetSyncService {
}
return;
}
// CONTROL_IDENTITY_MISMATCH is permanent until the operator
// explicitly resets the peer anchor. Record the failure once
// (last_error + last_failure_at) plus a sticky flag so the
// retry service and event-driven push path skip subsequent
// attempts.
if (data?.code === SYNC_ERROR_CODES.controlIdentityMismatch) {
const message = this.formatError(err);
console.warn(
`[FleetSync] Failed to push ${resource} to "${node.name}" (${baseUrl}): ${message}`,
);
db.recordFleetSyncFailure(node.id, resource, message);
db.setFleetSyncSticky(
node.id,
resource,
SYNC_ERROR_CODES.controlIdentityMismatch,
typeof data.expected === 'string' ? data.expected : null,
typeof data.got === 'string' ? data.got : null,
);
return;
}
}
const message = this.formatError(err);
console.warn(