Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
108 changes: 89 additions & 19 deletions packages/grpc-js/test/test-weighted-round-robin.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,20 +28,19 @@ const protoFile = path.join(__dirname, 'fixtures', 'echo_service.proto');
const EchoService = (loadProtoFile(protoFile) as unknown as ProtoGrpcType).EchoService;

function makeNCalls(client: EchoServiceClient, count: number): Promise<{[serverId: string]: number}> {
return new Promise((resolve, reject) => {
const result: {[serverId: string]: number} = {};
function makeOneCall(callsLeft: number) {
if (callsLeft <= 0) {
resolve(result);
} else {
const result: {[serverId: string]: number} = {};
const promises: Promise<void>[] = [];
for (let index = 0; index < count; index++) {
promises.push(
new Promise<void>((resolve, reject) => {
const deadline = new Date();
deadline.setMilliseconds(deadline.getMilliseconds() + 100);
const call= client.echo({}, {deadline}, (error, value) => {
deadline.setSeconds(deadline.getSeconds() + 2);
const call = client.echo({}, {deadline}, (error, value) => {
if (error) {
reject(error);
return;
}
makeOneCall(callsLeft - 1);
resolve();
});
call.on('metadata', metadata => {
const serverEntry = metadata.get('server');
Expand All @@ -53,9 +52,60 @@ function makeNCalls(client: EchoServiceClient, count: number): Promise<{[serverI
result[serverId] += 1;
}
});
})
);
}
return Promise.all(promises).then(() => result);
}

function warmupAllServers(client: EchoServiceClient, serverCount = 2): Promise<void> {
const seenServers = new Set<string>();
const deadline = new Date();
deadline.setSeconds(deadline.getSeconds() + 5);

return new Promise<void>((resolve, reject) => {
function sendNextCall() {
if (seenServers.size >= serverCount) {
resolve();
return;
}
const call = client.echo({}, {deadline}, error => {
if (error) {
reject(error);
return;
}
sendNextCall();
});
call.on('metadata', metadata => {
const serverEntry = metadata.get('server');
if (serverEntry.length > 0) {
const serverId = serverEntry[0] as string;
if (serverId) {
seenServers.add(serverId);
}
}
});
}
makeOneCall(count);
sendNextCall();
});
}

function waitForPickerUpdate(client: EchoServiceClient, timeoutMs = 2000): Promise<void> {
const internalChannel = (client.getChannel() as any).internalChannel;
const initialPicker = internalChannel.currentPicker;
const startTime = Date.now();

return new Promise<void>((resolve, reject) => {
const checkPicker = () => {
if (internalChannel.currentPicker !== initialPicker) {
resolve();
} else if (Date.now() - startTime > timeoutMs) {
reject(new Error('Timed out waiting for WeightedRoundRobinPicker to update'));
} else {
setTimeout(checkPicker, 10);
}
};
checkPicker();
});
}

Expand All @@ -72,13 +122,28 @@ function createClient(ports: number[], serviceConfig: grpc.ServiceConfig) {
return new EchoService(`ipv4:${ports.map(port => `127.0.0.1:${port}`).join(',')}`, grpc.credentials.createInsecure(), {'grpc.service_config': JSON.stringify(serviceConfig)});
}

function waitForClientReady(client: EchoServiceClient): Promise<void> {
return new Promise((resolve, reject) => {
const deadline = new Date();
deadline.setSeconds(deadline.getSeconds() + 5);
client.waitForReady(deadline, error => {
if (error) {
reject(error);
} else {
resolve();
}
});
});
}

function asyncTimeout(delay: number): Promise<void> {
return new Promise(resolve => {
setTimeout(resolve, delay);
});
}

describe('Weighted round robin LB policy', () => {
describe('Weighted round robin LB policy', function () {
this.timeout(process.platform === 'win32' ? 10000 : 5000);
describe('Config parsing', () => {
it('Should have default values with an empty object', () => {
const config = WeightedRoundRobinLoadBalancingConfig.createFromJson({});
Expand Down Expand Up @@ -261,7 +326,8 @@ describe('Weighted round robin LB policy', () => {
it('Should evenly balance among endpoints with no weight', async () => {
const serviceConfig = createServiceConfig({});
client = createClient([port1, port2], serviceConfig);
await makeNCalls(client, 10);
await waitForClientReady(client);
await warmupAllServers(client);
const result = await makeNCalls(client, 30);
assert(Math.abs(result['1'] - result['2']) < 3, `server1: ${result['1']}, server2: ${result[2]}`);
});
Expand All @@ -271,12 +337,13 @@ describe('Weighted round robin LB policy', () => {
weight_update_period: '0.1s'
});
client = createClient([port1, port2], serviceConfig);
await waitForClientReady(client);
server1Metrics.qps = 3;
server1Metrics.utilization = 1;
server2Metrics.qps = 1;
server2Metrics.utilization = 1;
await makeNCalls(client, 10);
await asyncTimeout(200);
await warmupAllServers(client);
await waitForPickerUpdate(client);
const result = await makeNCalls(client, 40);
assert(Math.abs(result['1'] - 30) < 3, `server1: ${result['1']}, server2: ${result['2']}`);
});
Expand Down Expand Up @@ -324,14 +391,15 @@ describe('Weighted round robin LB policy', () => {
error_utilization_penalty: 1
});
client = createClient([port1, port2], serviceConfig);
await waitForClientReady(client);
server1Metrics.qps = 2;
server1Metrics.utilization = 1;
server1Metrics.eps = 0;
server2Metrics.qps = 2;
server2Metrics.utilization = 1;
server2Metrics.eps = 2;
await makeNCalls(client, 10);
await asyncTimeout(100);
await warmupAllServers(client);
await waitForPickerUpdate(client);
const result = await makeNCalls(client, 30);
assert(Math.abs(result['1'] - 20) < 3, `server1: ${result['1']}, server2: ${result['2']}`);
});
Expand Down Expand Up @@ -411,7 +479,8 @@ describe('Weighted round robin LB policy', () => {
blackout_period: '0.01s'
});
client = createClient([port1, port2], serviceConfig);
await makeNCalls(client, 10);
await waitForClientReady(client);
await warmupAllServers(client);
const result = await makeNCalls(client, 30);
assert(Math.abs(result['1'] - result['2']) < 3, `server1: ${result['1']}, server2: ${result[2]}`);
});
Expand All @@ -423,12 +492,13 @@ describe('Weighted round robin LB policy', () => {
weight_update_period: '0.1s'
});
client = createClient([port1, port2], serviceConfig);
await waitForClientReady(client);
server1MetricRecorder.setQpsMetric(3);
server1MetricRecorder.setApplicationUtilizationMetric(1);
server2MetricRecorder.setQpsMetric(1);
server2MetricRecorder.setApplicationUtilizationMetric(1);
await makeNCalls(client, 10);
await asyncTimeout(200);
await warmupAllServers(client);
await waitForPickerUpdate(client);
const result = await makeNCalls(client, 40);
assert(Math.abs(result['1'] - 30) < 3, `server1: ${result['1']}, server2: ${result['2']}`);
});
Expand Down
Loading