feat(server): onPeerDeactivated utility mirrors onPeerActivated (TDD)
This commit is contained in:
@@ -41,6 +41,7 @@ vi.mock('../ws/handler.js', () => ({
|
|||||||
getAllOnlineUserIds: () => [],
|
getAllOnlineUserIds: () => [],
|
||||||
sendToUser: vi.fn(),
|
sendToUser: vi.fn(),
|
||||||
sendToDmMembers: vi.fn(),
|
sendToDmMembers: vi.fn(),
|
||||||
|
evictFederatedCallsForHost: vi.fn().mockReturnValue(0),
|
||||||
},
|
},
|
||||||
}));
|
}));
|
||||||
|
|
||||||
@@ -348,3 +349,106 @@ describe('onPeerActivated', () => {
|
|||||||
await expect(onPeerActivated('peer-err', 'ensure_peered')).resolves.toBeUndefined();
|
await expect(onPeerActivated('peer-err', 'ensure_peered')).resolves.toBeUndefined();
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe('onPeerDeactivated', () => {
|
||||||
|
beforeEach(() => {
|
||||||
|
sqlite = new Database(':memory:');
|
||||||
|
testDb = drizzle(sqlite, { schema });
|
||||||
|
applyMigrations(sqlite);
|
||||||
|
vi.clearAllMocks();
|
||||||
|
});
|
||||||
|
|
||||||
|
it('evicts federated calls for the peer origin with peer_transient_failure on network_threshold', async () => {
|
||||||
|
seedPeer('peer-net', 'unreachable');
|
||||||
|
testDb.update(schema.federationPeers)
|
||||||
|
.set({ instanceName: 'NetPeer' })
|
||||||
|
.where(eq(schema.federationPeers.id, 'peer-net'))
|
||||||
|
.run();
|
||||||
|
|
||||||
|
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
|
||||||
|
await onPeerDeactivated('peer-net', 'network_threshold');
|
||||||
|
|
||||||
|
const { connectionManager } = await import('../ws/handler.js');
|
||||||
|
expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledWith(
|
||||||
|
'https://peer-net.example',
|
||||||
|
{ reason: 'peer_transient_failure', peerLabel: 'NetPeer' },
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('maps rejected status to peer_rejected reason', async () => {
|
||||||
|
seedPeer('peer-rej', 'rejected');
|
||||||
|
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
|
||||||
|
await onPeerDeactivated('peer-rej', 'remote_rejected');
|
||||||
|
|
||||||
|
const { connectionManager } = await import('../ws/handler.js');
|
||||||
|
expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledWith(
|
||||||
|
'https://peer-rej.example',
|
||||||
|
{ reason: 'peer_rejected', peerLabel: undefined },
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('maps revoked status to peer_rejected reason', async () => {
|
||||||
|
seedPeer('peer-rev', 'revoked');
|
||||||
|
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
|
||||||
|
await onPeerDeactivated('peer-rev', 'admin_revoked');
|
||||||
|
|
||||||
|
const { connectionManager } = await import('../ws/handler.js');
|
||||||
|
expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledWith(
|
||||||
|
'https://peer-rev.example',
|
||||||
|
{ reason: 'peer_rejected', peerLabel: undefined },
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('broadcasts federation_peers_changed to admins', async () => {
|
||||||
|
seedPeer('peer-broad', 'unreachable');
|
||||||
|
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
|
||||||
|
await onPeerDeactivated('peer-broad', 'network_threshold');
|
||||||
|
|
||||||
|
const { connectionManager } = await import('../ws/handler.js');
|
||||||
|
expect(connectionManager.sendToAdmins).toHaveBeenCalledWith({ type: 'federation_peers_changed' });
|
||||||
|
});
|
||||||
|
|
||||||
|
it('aborts silently when the peer row is missing', async () => {
|
||||||
|
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
|
||||||
|
await expect(onPeerDeactivated('peer-missing', 'network_threshold')).resolves.toBeUndefined();
|
||||||
|
|
||||||
|
const { connectionManager } = await import('../ws/handler.js');
|
||||||
|
expect(connectionManager.evictFederatedCallsForHost).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
|
||||||
|
it('deduplicates concurrent calls for the same peerId', async () => {
|
||||||
|
seedPeer('peer-x', 'unreachable');
|
||||||
|
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
|
||||||
|
const { connectionManager } = await import('../ws/handler.js');
|
||||||
|
|
||||||
|
const p1 = onPeerDeactivated('peer-x', 'network_threshold');
|
||||||
|
const p2 = onPeerDeactivated('peer-x', 'network_threshold');
|
||||||
|
await Promise.all([p1, p2]);
|
||||||
|
|
||||||
|
// Exactly one eviction call, not two
|
||||||
|
expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledTimes(1);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('uses a dedup map SEPARATE from onPeerActivated', async () => {
|
||||||
|
seedPeer('peer-flap', 'active');
|
||||||
|
const { onPeerActivated, onPeerDeactivated } = await import('./federationPeerActivation.js');
|
||||||
|
|
||||||
|
// Simulate an activation already in flight — spawn onPeerActivated then
|
||||||
|
// immediately kick off a deactivation for the same peer id. The latter
|
||||||
|
// must NOT be swallowed as a dedup hit against the activation.
|
||||||
|
const actPromise = onPeerActivated('peer-flap', 'ensure_peered');
|
||||||
|
|
||||||
|
// Mark peer non-active now — otherwise the deactivation utility's
|
||||||
|
// own status guard would skip it.
|
||||||
|
testDb.update(schema.federationPeers)
|
||||||
|
.set({ status: 'unreachable' })
|
||||||
|
.where(eq(schema.federationPeers.id, 'peer-flap'))
|
||||||
|
.run();
|
||||||
|
|
||||||
|
const deactPromise = onPeerDeactivated('peer-flap', 'network_threshold');
|
||||||
|
await Promise.all([actPromise, deactPromise]);
|
||||||
|
|
||||||
|
const { connectionManager } = await import('../ws/handler.js');
|
||||||
|
expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
@@ -191,3 +191,88 @@ export async function startupBootstrapSync(): Promise<void> {
|
|||||||
await onPeerActivated(peer.id, 'startup_bootstrap');
|
await onPeerActivated(peer.id, 'startup_bootstrap');
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export type PeerDeactivationReason =
|
||||||
|
| 'network_threshold' // outbox worker hit PEER_UNREACHABLE_THRESHOLD
|
||||||
|
| 'auth_threshold' // outbox worker hit AUTH_FAILURE_THRESHOLD
|
||||||
|
| 'remote_rejected' // auto-peer handshake got 403 PEERING_REQUIRES_APPROVAL
|
||||||
|
| 'admin_revoked' // admin revoked peering from this side
|
||||||
|
| 'admin_reset'; // admin reset peer clearing to non-active status
|
||||||
|
|
||||||
|
// Dedup: concurrent deactivations for the same peerId share one promise.
|
||||||
|
// SEPARATE from inFlightActivation — a flapping peer's activate-then-deactivate
|
||||||
|
// sequence must not collapse into one slot.
|
||||||
|
const inFlightDeactivation = new Map<string, Promise<void>>();
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Called whenever federation_peers.status transitions OUT OF 'active' for any reason.
|
||||||
|
* Sweeps connectionManager.federatedCalls for entries whose federatedCallHost matches
|
||||||
|
* the peer origin, emitting dm_call_undeliverable { phase: 'host_unreachable', terminal: true }
|
||||||
|
* to stranded ringed users and clearing the entries.
|
||||||
|
*
|
||||||
|
* Call sites (must remain exhaustive — grep `onPeerDeactivated(` to audit):
|
||||||
|
* - utils/federationWorker.ts handleOutboxDeliveryFailure when status flips to 'unreachable'
|
||||||
|
* - utils/federationWorker.ts auth-failure path when status flips to 'needs_attention'
|
||||||
|
* - utils/federationWorker.ts resolvePendingPeers case 'rejected'
|
||||||
|
* - routes/federation.ts admin revoke endpoint
|
||||||
|
* - routes/federation.ts admin reset endpoint (when it transitions to a non-active status)
|
||||||
|
* - utils/federationPeering.ts performHandshake 403 PEERING_REQUIRES_APPROVAL path
|
||||||
|
*
|
||||||
|
* Deduplicated by peerId — concurrent calls share one promise. Separate map from
|
||||||
|
* onPeerActivated so flapping peers don't collapse transitions.
|
||||||
|
*/
|
||||||
|
export async function onPeerDeactivated(
|
||||||
|
peerId: string,
|
||||||
|
reason: PeerDeactivationReason,
|
||||||
|
): Promise<void> {
|
||||||
|
const existing = inFlightDeactivation.get(peerId);
|
||||||
|
if (existing) return existing;
|
||||||
|
|
||||||
|
const promise = (async () => {
|
||||||
|
try {
|
||||||
|
const db = getDb();
|
||||||
|
const peer = db.select({
|
||||||
|
origin: schema.federationPeers.origin,
|
||||||
|
status: schema.federationPeers.status,
|
||||||
|
instanceName: schema.federationPeers.instanceName,
|
||||||
|
})
|
||||||
|
.from(schema.federationPeers)
|
||||||
|
.where(eq(schema.federationPeers.id, peerId))
|
||||||
|
.get();
|
||||||
|
|
||||||
|
if (!peer) {
|
||||||
|
// Peer row gone — nothing to sweep against.
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
const { connectionManager } = await import('../ws/handler.js');
|
||||||
|
|
||||||
|
// Map status to user-facing reason.
|
||||||
|
const isRejectedLike = peer.status === 'rejected' || peer.status === 'revoked';
|
||||||
|
const mappedReason: 'peer_rejected' | 'peer_transient_failure' =
|
||||||
|
isRejectedLike ? 'peer_rejected' : 'peer_transient_failure';
|
||||||
|
|
||||||
|
const evicted = connectionManager.evictFederatedCallsForHost(peer.origin, {
|
||||||
|
reason: mappedReason,
|
||||||
|
peerLabel: peer.instanceName ?? undefined,
|
||||||
|
});
|
||||||
|
|
||||||
|
if (evicted > 0) {
|
||||||
|
console.log(
|
||||||
|
`[federation] onPeerDeactivated(${peerId}, ${reason}) evicted ${evicted} FederatedCallEntry object${evicted === 1 ? '' : 's'} for ${peer.origin}`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const });
|
||||||
|
} catch (err) {
|
||||||
|
console.error(`[federation] onPeerDeactivated(${peerId}, ${reason}) failed:`, err);
|
||||||
|
}
|
||||||
|
})();
|
||||||
|
|
||||||
|
inFlightDeactivation.set(peerId, promise);
|
||||||
|
try {
|
||||||
|
await promise;
|
||||||
|
} finally {
|
||||||
|
inFlightDeactivation.delete(peerId);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user