feat(server): wire onPeerDeactivated at federationWorker peer-deactivation sites
This commit is contained in:
@@ -11,7 +11,7 @@ import { getDmMessageWithUser } from '../routes/dm.js';
|
|||||||
import { connectionManager } from '../ws/handler.js';
|
import { connectionManager } from '../ws/handler.js';
|
||||||
import { generateThumbnail } from './thumbnail.js';
|
import { generateThumbnail } from './thumbnail.js';
|
||||||
import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent } from '@backspace/shared';
|
import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent } from '@backspace/shared';
|
||||||
import { onPeerActivated, startupBootstrapSync } from './federationPeerActivation.js';
|
import { onPeerActivated, startupBootstrapSync, onPeerDeactivated } from './federationPeerActivation.js';
|
||||||
import fs from 'node:fs';
|
import fs from 'node:fs';
|
||||||
import path from 'node:path';
|
import path from 'node:path';
|
||||||
import crypto from 'node:crypto';
|
import crypto from 'node:crypto';
|
||||||
@@ -310,6 +310,9 @@ export async function processOutboxTick(): Promise<void> {
|
|||||||
})
|
})
|
||||||
.where(eq(schema.federationPeers.id, peerId))
|
.where(eq(schema.federationPeers.id, peerId))
|
||||||
.run();
|
.run();
|
||||||
|
onPeerDeactivated(peerId, 'auth_threshold').catch(err =>
|
||||||
|
console.error('[federation-worker] onPeerDeactivated from auth threshold failed:', err)
|
||||||
|
);
|
||||||
console.warn(
|
console.warn(
|
||||||
`[federation-worker] Peer ${peerOrigin} transitioned to needs_attention after ${decision.newAuthFailures} consecutive ${response.status} responses`,
|
`[federation-worker] Peer ${peerOrigin} transitioned to needs_attention after ${decision.newAuthFailures} consecutive ${response.status} responses`,
|
||||||
);
|
);
|
||||||
@@ -416,6 +419,12 @@ function handleOutboxDeliveryFailure(
|
|||||||
.set(updates)
|
.set(updates)
|
||||||
.where(eq(schema.federationPeers.id, peerId))
|
.where(eq(schema.federationPeers.id, peerId))
|
||||||
.run();
|
.run();
|
||||||
|
|
||||||
|
if (newFailures >= PEER_UNREACHABLE_THRESHOLD) {
|
||||||
|
onPeerDeactivated(peerId, 'network_threshold').catch(err =>
|
||||||
|
console.error('[federation-worker] onPeerDeactivated from unreachable threshold failed:', err)
|
||||||
|
);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ─── Pending Peer Resolution ───────────────────────────────────────────────
|
// ─── Pending Peer Resolution ───────────────────────────────────────────────
|
||||||
@@ -468,6 +477,10 @@ async function resolvePendingPeers(): Promise<void> {
|
|||||||
.where(eq(schema.federationOutbox.peerId, peerId))
|
.where(eq(schema.federationOutbox.peerId, peerId))
|
||||||
.run();
|
.run();
|
||||||
|
|
||||||
|
onPeerDeactivated(peerId, 'remote_rejected').catch(err =>
|
||||||
|
console.error('[federation-worker] onPeerDeactivated from resolvePendingPeers rejected failed:', err)
|
||||||
|
);
|
||||||
|
|
||||||
// Push federation_peer_rejected WS event to affected users
|
// Push federation_peer_rejected WS event to affected users
|
||||||
pushPeerRejectedEvent(peerOrigin, contextMap);
|
pushPeerRejectedEvent(peerOrigin, contextMap);
|
||||||
// Notify admins of peer state change
|
// Notify admins of peer state change
|
||||||
|
|||||||
Reference in New Issue
Block a user