(@NotNull Iterable<Long> roleIds, long typeId, @NotNull Binary fullEncodedProtocol, boolean trySend)
| 1311 | |
| 1312 | // 可在事务外执行 |
| 1313 | public int sendDirect(@NotNull Iterable<Long> roleIds, long typeId, @NotNull Binary fullEncodedProtocol, |
| 1314 | boolean trySend) { |
| 1315 | var roleIdSet = new LongHashSet(); |
| 1316 | for (var roleId : roleIds) |
| 1317 | roleIdSet.add(roleId); // 去重 |
| 1318 | if (roleIdSet.isEmpty()) |
| 1319 | return 0; |
| 1320 | var groups = new HashMap<String, LinkRoles>(); |
| 1321 | var links = providerApp.providerService.getLinks(); |
| 1322 | for (var it = roleIdSet.iterator(); it.moveToNext(); ) { |
| 1323 | var roleId = it.value(); |
| 1324 | var onlineShared = _tOnlineShared.selectDirty(roleId); |
| 1325 | if (onlineShared == null) { |
| 1326 | if (!trySend) { |
| 1327 | logger.info("sendDirects({}): not found roleId={} in _tonline", |
| 1328 | getTypeId(fullEncodedProtocol), roleId); |
| 1329 | } |
| 1330 | continue; |
| 1331 | } |
| 1332 | var link = onlineShared.getLink(); |
| 1333 | var state = link.getState(); |
| 1334 | if (state != eLogined) { |
| 1335 | if (!trySend) { |
| 1336 | logger.info("sendDirects({}): state={} != eLogined for roleId={}", |
| 1337 | getTypeId(fullEncodedProtocol), state, roleId); |
| 1338 | } |
| 1339 | continue; |
| 1340 | } |
| 1341 | var linkName = link.getLinkName(); |
| 1342 | // 后面保存connector.socket并使用,如果之后连接被关闭,以后发送协议失败。 |
| 1343 | var group = groups.get(linkName); |
| 1344 | if (group == null) { |
| 1345 | var linkSocket = getLinkSocket(links, linkName, roleId); // maybe null |
| 1346 | groups.put(linkName, group = new LinkRoles(linkName, linkSocket, typeId, fullEncodedProtocol)); |
| 1347 | } |
| 1348 | group.send.Argument.getLinkSids().add(link.getLinkSid()); |
| 1349 | group.roleIds.add(roleId); |
| 1350 | } |
| 1351 | |
| 1352 | int sendCount = 0; |
| 1353 | for (var group : groups.values()) { |
| 1354 | if (group.linkSocket == null) { |
| 1355 | processErrorSids(group.send.Argument.getLinkSids(), group); |
| 1356 | continue; // link miss process done |
| 1357 | } |
| 1358 | |
| 1359 | group.roleIds.foreach(this::setLocalActiveTimeIfPresent); |
| 1360 | if (group.send.Send(group.linkSocket, rpc -> { |
| 1361 | var send = group.send; |
| 1362 | var errorSids = send.isTimeout() ? send.Argument.getLinkSids() : send.Result.getErrorLinkSids(); |
| 1363 | processErrorSids(errorSids, group); |
| 1364 | return Procedure.Success; |
| 1365 | })) |
| 1366 | sendCount++; |
| 1367 | } |
| 1368 | return sendCount; |
| 1369 | } |
| 1370 |
no test coverage detected