< Summary

Information
Class: IceDiscovery.AdapterRequest
Assembly: IceDiscovery
File(s): /_/csharp/src/IceDiscovery/LookupI.cs
Tag: 125_37167941578
Line coverage
100%
Covered lines: 48
Uncovered lines: 0
Coverable lines: 48
Total lines: 579
Line coverage: 100%
Branch coverage
100%
Covered branches: 20
Total branches: 20
Branch coverage: 100%
Method coverage
100%
Covered methods: 7
Fully covered methods: 7
Total methods: 7
Method coverage: 100%
Full method coverage: 100%

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
retry()100%44100%
response(...)100%44100%
finished(...)100%1010100%
runTimerTask()100%11100%
invokeWithLookup(...)100%11100%
sendResponse(...)100%22100%

File(s)

/_/csharp/src/IceDiscovery/LookupI.cs

#LineLine coverage
 1// Copyright (c) ZeroC, Inc.
 2
 3using System.Diagnostics;
 4using System.Text;
 5
 6namespace IceDiscovery;
 7
 8internal abstract class Request<T>
 9{
 10    protected Request(LookupI lookup, T id, int retryCount)
 11    {
 12        lookup_ = lookup;
 13        retryCount_ = retryCount;
 14        _id = id;
 15        _requestId = Guid.NewGuid().ToString();
 16    }
 17
 18    public T getId() => _id;
 19
 20    public bool addCallback(TaskCompletionSource<Ice.ObjectPrx> cb)
 21    {
 22        callbacks_.Add(cb);
 23        return callbacks_.Count == 1;
 24    }
 25
 26    public virtual bool retry() => --retryCount_ >= 0;
 27
 28    public void invoke(string domainId, Dictionary<LookupPrx, LookupReplyPrx> lookups)
 29    {
 30        ++_generation;
 31        _lookupCount = lookups.Count;
 32        _failureCount = 0;
 33        var id = new Ice.Identity(_requestId, "");
 34        foreach (KeyValuePair<LookupPrx, LookupReplyPrx> entry in lookups)
 35        {
 36            invokeWithLookup(
 37                domainId,
 38                entry.Key,
 39                LookupReplyPrxHelper.uncheckedCast(entry.Value.ice_identity(id)));
 40        }
 41    }
 42
 43    public bool exception(int generation)
 44    {
 45        // Ignore a delayed failure from an earlier invocation round: it must not count against the current round, which
 46        // reset _failureCount and may still have outstanding lookups.
 47        if (generation != _generation)
 48        {
 49            return false;
 50        }
 51
 52        if (++_failureCount == _lookupCount)
 53        {
 54            finished(null);
 55            return true;
 56        }
 57        return false;
 58    }
 59
 60    public string getRequestId() => _requestId;
 61
 62    public abstract void finished(Ice.ObjectPrx proxy);
 63
 64    protected abstract void invokeWithLookup(string domainId, LookupPrx lookup, LookupReplyPrx lookupReply);
 65
 66    private readonly string _requestId;
 67
 68    protected LookupI lookup_;
 69    protected int retryCount_;
 70    protected int _lookupCount;
 71    protected int _failureCount;
 72
 73    // Incremented on each invoke() (i.e. each retry round). The exception callbacks capture the value current when they
 74    // were sent, so a delayed failure from an earlier round is ignored instead of counting against the current round.
 75    protected int _generation;
 76    protected List<TaskCompletionSource<Ice.ObjectPrx>> callbacks_ = new List<TaskCompletionSource<Ice.ObjectPrx>>();
 77
 78    protected T _id;
 79}
 80
 81internal class AdapterRequest : Request<string>, Ice.Internal.TimerTask
 82{
 83    public AdapterRequest(LookupI lookup, string id, int retryCount)
 184        : base(lookup, id, retryCount) => _start = Stopwatch.GetTimestamp();
 85
 86    public override bool retry()
 87    {
 188        if (_proxies.Count == 0 && --retryCount_ >= 0)
 89        {
 90            // We only reach here with _proxies empty, i.e. no replica-group response arrived. _latency is set only
 91            // together with inserting into _proxies, so the latency timer is not running.
 92            Debug.Assert(_latency == TimeSpan.Zero);
 93
 94            // Restart the replica-group latency window so it's measured from this round rather than from request
 95            // creation.
 196            _start = Stopwatch.GetTimestamp();
 197            return true;
 98        }
 199        return false;
 100    }
 101
 102    public bool response(Ice.ObjectPrx proxy, bool isReplicaGroup)
 103    {
 1104        if (isReplicaGroup)
 105        {
 1106            _proxies.Add(proxy);
 1107            if (_latency == TimeSpan.Zero)
 108            {
 109                // The aggregation window is the measured response time, scaled by IceDiscovery.LatencyMultiplier,
 110                // clamped to the timer's valid delay range, with a 1ms floor so we never schedule a degenerate
 111                // zero-length window.
 1112                double responseTimeMs = Stopwatch.GetElapsedTime(_start).TotalMilliseconds;
 1113                _latency = TimeSpan.FromMilliseconds(
 1114                    Math.Clamp(responseTimeMs * lookup_.latencyMultiplier(), 1, int.MaxValue));
 1115                lookup_.timer().cancel(this);
 1116                lookup_.timer().schedule(this, _latency);
 117            }
 1118            return false;
 119        }
 1120        finished(proxy);
 1121        return true;
 122    }
 123
 124    public override void finished(Ice.ObjectPrx proxy)
 125    {
 1126        if (proxy != null || _proxies.Count == 0)
 127        {
 1128            sendResponse(proxy);
 129        }
 1130        else if (_proxies.Count == 1)
 131        {
 1132            sendResponse(_proxies.First());
 133        }
 134        else
 135        {
 1136            var endpoints = new List<Ice.Endpoint>();
 1137            Ice.ObjectPrx result = null;
 1138            foreach (Ice.ObjectPrx prx in _proxies)
 139            {
 1140                result ??= prx;
 1141                endpoints.AddRange(prx.ice_getEndpoints());
 142            }
 1143            sendResponse(result.ice_endpoints(endpoints.ToArray()));
 144        }
 1145    }
 146
 1147    public void runTimerTask() => lookup_.adapterRequestTimedOut(this);
 148
 149    protected override void invokeWithLookup(string domainId, LookupPrx lookup, LookupReplyPrx lookupReply)
 150    {
 1151        int generation = _generation;
 1152        lookup.findAdapterByIdAsync(domainId, _id, lookupReply).ContinueWith(
 1153            task =>
 1154            {
 1155                try
 1156                {
 1157                    task.Wait();
 1158                }
 1159                catch (AggregateException ex)
 1160                {
 1161                    lookup_.adapterRequestException(this, ex.InnerException, generation);
 1162                }
 1163            },
 1164            lookup.ice_scheduler());
 1165    }
 166
 167    private void sendResponse(Ice.ObjectPrx proxy)
 168    {
 1169        foreach (TaskCompletionSource<Ice.ObjectPrx> cb in callbacks_)
 170        {
 1171            cb.SetResult(proxy);
 172        }
 1173        callbacks_.Clear();
 1174    }
 175
 176    //
 177    // We use a HashSet because the same IceDiscovery plugin might return multiple times
 178    // the same proxy if it's accessible through multiple network interfaces and if we
 179    // also sent the request to multiple interfaces.
 180    //
 1181    private readonly HashSet<Ice.ObjectPrx> _proxies = new HashSet<Ice.ObjectPrx>();
 182    private long _start;
 183    private TimeSpan _latency;
 184}
 185
 186internal class ObjectRequest : Request<Ice.Identity>, Ice.Internal.TimerTask
 187{
 188    public ObjectRequest(LookupI lookup, Ice.Identity id, int retryCount)
 189        : base(lookup, id, retryCount)
 190    {
 191    }
 192
 193    public void response(Ice.ObjectPrx proxy) => finished(proxy);
 194
 195    public override void finished(Ice.ObjectPrx proxy)
 196    {
 197        foreach (TaskCompletionSource<Ice.ObjectPrx> cb in callbacks_)
 198        {
 199            cb.SetResult(proxy);
 200        }
 201        callbacks_.Clear();
 202    }
 203
 204    public void runTimerTask() => lookup_.objectRequestTimedOut(this);
 205
 206    protected override void invokeWithLookup(string domainId, LookupPrx lookup, LookupReplyPrx lookupReply)
 207    {
 208        int generation = _generation;
 209        lookup.findObjectByIdAsync(domainId, _id, lookupReply).ContinueWith(
 210            task =>
 211            {
 212                try
 213                {
 214                    task.Wait();
 215                }
 216                catch (AggregateException ex)
 217                {
 218                    lookup_.objectRequestException(this, ex.InnerException, generation);
 219                }
 220            },
 221            lookup.ice_scheduler());
 222    }
 223}
 224
 225internal class LookupI : LookupDisp_
 226{
 227    public LookupI(LocatorRegistryI registry, LookupPrx lookup, Ice.Properties properties)
 228    {
 229        _registry = registry;
 230        _lookup = lookup;
 231        _timeout = TimeSpan.FromMilliseconds(properties.getIcePropertyAsInt("IceDiscovery.Timeout"));
 232        if (_timeout <= TimeSpan.Zero)
 233        {
 234            throw new Ice.PropertyException("property 'IceDiscovery.Timeout' must be greater than 0");
 235        }
 236        _retryCount = properties.getIcePropertyAsInt("IceDiscovery.RetryCount");
 237        if (_retryCount < 0)
 238        {
 239            throw new Ice.PropertyException("property 'IceDiscovery.RetryCount' must be greater than or equal to 0");
 240        }
 241        _latencyMultiplier = properties.getIcePropertyAsInt("IceDiscovery.LatencyMultiplier");
 242        if (_latencyMultiplier < 1)
 243        {
 244            throw new Ice.PropertyException(
 245                "property 'IceDiscovery.LatencyMultiplier' must be greater than or equal to 1");
 246        }
 247        _domainId = properties.getIceProperty("IceDiscovery.DomainId");
 248        _timer = Ice.Internal.Util.getInstance(lookup.ice_getCommunicator()).timer();
 249
 250        //
 251        // Create one lookup proxy per endpoint from the given proxy. We want to send a multicast
 252        // datagram on each endpoint.
 253        //
 254        var single = new Ice.Endpoint[1];
 255        foreach (Ice.Endpoint endpt in lookup.ice_getEndpoints())
 256        {
 257            single[0] = endpt;
 258            _lookups[(LookupPrx)lookup.ice_endpoints(single)] = null;
 259        }
 260        Debug.Assert(_lookups.Count > 0);
 261    }
 262
 263    public void setLookupReply(LookupReplyPrx lookupReply)
 264    {
 265        //
 266        // Use a lookup reply proxy whose address matches the interface used to send multicast datagrams.
 267        //
 268        var single = new Ice.Endpoint[1];
 269        foreach (LookupPrx key in new List<LookupPrx>(_lookups.Keys))
 270        {
 271            var info = (Ice.UDPEndpointInfo)key.ice_getEndpoints()[0].getInfo();
 272            if (info.mcastInterface.Length > 0)
 273            {
 274                foreach (Ice.Endpoint q in lookupReply.ice_getEndpoints())
 275                {
 276                    Ice.EndpointInfo r = q.getInfo();
 277                    if (r is Ice.IPEndpointInfo &&
 278                        ((Ice.IPEndpointInfo)r).host.Equals(info.mcastInterface, StringComparison.Ordinal))
 279                    {
 280                        single[0] = q;
 281                        _lookups[key] = (LookupReplyPrx)lookupReply.ice_endpoints(single);
 282                        break;
 283                    }
 284                }
 285            }
 286
 287            if (_lookups[key] == null)
 288            {
 289                // Fallback: just use the given lookup reply proxy if no matching endpoint found.
 290                _lookups[key] = lookupReply;
 291            }
 292        }
 293    }
 294
 295    public override void findObjectById(string domainId, Ice.Identity id, LookupReplyPrx reply, Ice.Current current)
 296    {
 297        Ice.ObjectPrx.checkNotNull(reply, current);
 298        if (!domainId.Equals(_domainId, StringComparison.Ordinal))
 299        {
 300            return; // Ignore
 301        }
 302
 303        Ice.ObjectPrx proxy = _registry.findObject(id);
 304        if (proxy != null)
 305        {
 306            //
 307            // Reply to the multicast request using the given proxy.
 308            //
 309            try
 310            {
 311                reply.foundObjectByIdAsync(id, proxy);
 312            }
 313            catch (Ice.LocalException)
 314            {
 315                // Ignore.
 316            }
 317        }
 318    }
 319
 320    public override void findAdapterById(string domainId, string adapterId, LookupReplyPrx reply, Ice.Current current)
 321    {
 322        Ice.ObjectPrx.checkNotNull(reply, current);
 323        if (!domainId.Equals(_domainId, StringComparison.Ordinal))
 324        {
 325            return; // Ignore
 326        }
 327
 328        Ice.ObjectPrx proxy = _registry.findAdapter(adapterId, out bool isReplicaGroup);
 329        if (proxy != null)
 330        {
 331            //
 332            // Reply to the multicast request using the given proxy.
 333            //
 334            try
 335            {
 336                reply.foundAdapterByIdAsync(adapterId, proxy, isReplicaGroup);
 337            }
 338            catch (Ice.LocalException)
 339            {
 340                // Ignore.
 341            }
 342        }
 343    }
 344
 345    internal Task<Ice.ObjectPrx> findObject(Ice.Identity id)
 346    {
 347        lock (_mutex)
 348        {
 349            if (!_objectRequests.TryGetValue(id, out ObjectRequest request))
 350            {
 351                request = new ObjectRequest(this, id, _retryCount);
 352                _objectRequests.Add(id, request);
 353            }
 354
 355            var task = new TaskCompletionSource<Ice.ObjectPrx>(TaskCreationOptions.RunContinuationsAsynchronously);
 356            if (request.addCallback(task))
 357            {
 358                try
 359                {
 360                    request.invoke(_domainId, _lookups);
 361                    _timer.schedule(request, _timeout);
 362                }
 363                catch (Ice.LocalException)
 364                {
 365                    request.finished(null);
 366                    _objectRequests.Remove(id);
 367                }
 368            }
 369            return task.Task;
 370        }
 371    }
 372
 373    internal Task<Ice.ObjectPrx> findAdapter(string adapterId)
 374    {
 375        lock (_mutex)
 376        {
 377            if (!_adapterRequests.TryGetValue(adapterId, out AdapterRequest request))
 378            {
 379                request = new AdapterRequest(this, adapterId, _retryCount);
 380                _adapterRequests.Add(adapterId, request);
 381            }
 382
 383            var task = new TaskCompletionSource<Ice.ObjectPrx>(TaskCreationOptions.RunContinuationsAsynchronously);
 384            if (request.addCallback(task))
 385            {
 386                try
 387                {
 388                    request.invoke(_domainId, _lookups);
 389                    _timer.schedule(request, _timeout);
 390                }
 391                catch (Ice.LocalException)
 392                {
 393                    request.finished(null);
 394                    _adapterRequests.Remove(adapterId);
 395                }
 396            }
 397            return task.Task;
 398        }
 399    }
 400
 401    internal void foundObject(Ice.Identity id, string requestId, Ice.ObjectPrx proxy)
 402    {
 403        lock (_mutex)
 404        {
 405            if (_objectRequests.TryGetValue(id, out ObjectRequest request) && request.getRequestId() == requestId)
 406            {
 407                request.response(proxy);
 408                _timer.cancel(request);
 409                _objectRequests.Remove(id);
 410            }
 411            // else ignore responses from old requests
 412        }
 413    }
 414
 415    internal void foundAdapter(string adapterId, string requestId, Ice.ObjectPrx proxy, bool isReplicaGroup)
 416    {
 417        lock (_mutex)
 418        {
 419            if (
 420                _adapterRequests.TryGetValue(adapterId, out AdapterRequest request) &&
 421                request.getRequestId() == requestId)
 422            {
 423                if (request.response(proxy, isReplicaGroup))
 424                {
 425                    _timer.cancel(request);
 426                    _adapterRequests.Remove(request.getId());
 427                }
 428            }
 429            // else ignore responses from old requests
 430        }
 431    }
 432
 433    internal void objectRequestTimedOut(ObjectRequest request)
 434    {
 435        lock (_mutex)
 436        {
 437            if (!_objectRequests.TryGetValue(request.getId(), out ObjectRequest r) || r != request)
 438            {
 439                return;
 440            }
 441
 442            if (request.retry())
 443            {
 444                try
 445                {
 446                    request.invoke(_domainId, _lookups);
 447                    _timer.schedule(request, _timeout);
 448                    return;
 449                }
 450                catch (Ice.LocalException)
 451                {
 452                }
 453            }
 454
 455            request.finished(null);
 456            _objectRequests.Remove(request.getId());
 457            _timer.cancel(request);
 458        }
 459    }
 460
 461    internal void objectRequestException(ObjectRequest request, Exception ex, int generation)
 462    {
 463        lock (_mutex)
 464        {
 465            if (!_objectRequests.TryGetValue(request.getId(), out ObjectRequest r) || r != request)
 466            {
 467                return;
 468            }
 469
 470            if (request.exception(generation))
 471            {
 472                if (_warnOnce)
 473                {
 474                    var s = new StringBuilder();
 475                    s.Append("failed to lookup object `");
 476                    s.Append(_lookup.ice_getCommunicator().identityToString(request.getId()));
 477                    s.Append("' with lookup proxy `");
 478                    s.Append(_lookup);
 479                    s.Append("':\n");
 480                    s.Append(ex.ToString());
 481                    _lookup.ice_getCommunicator().getLogger().warning(s.ToString());
 482                    _warnOnce = false;
 483                }
 484                _timer.cancel(request);
 485                _objectRequests.Remove(request.getId());
 486            }
 487        }
 488    }
 489
 490    internal void adapterRequestTimedOut(AdapterRequest request)
 491    {
 492        lock (_mutex)
 493        {
 494            if (!_adapterRequests.TryGetValue(request.getId(), out AdapterRequest r) || r != request)
 495            {
 496                return;
 497            }
 498
 499            if (request.retry())
 500            {
 501                try
 502                {
 503                    request.invoke(_domainId, _lookups);
 504                    _timer.schedule(request, _timeout);
 505                    return;
 506                }
 507                catch (Ice.LocalException)
 508                {
 509                }
 510            }
 511
 512            request.finished(null);
 513            _adapterRequests.Remove(request.getId());
 514            _timer.cancel(request);
 515        }
 516    }
 517
 518    internal void adapterRequestException(AdapterRequest request, Exception ex, int generation)
 519    {
 520        lock (_mutex)
 521        {
 522            if (!_adapterRequests.TryGetValue(request.getId(), out AdapterRequest r) || r != request)
 523            {
 524                return;
 525            }
 526
 527            if (request.exception(generation))
 528            {
 529                if (_warnOnce)
 530                {
 531                    var s = new StringBuilder();
 532                    s.Append("failed to lookup adapter `");
 533                    s.Append(request.getId());
 534                    s.Append("' with lookup proxy `");
 535                    s.Append(_lookup);
 536                    s.Append("':\n");
 537                    s.Append(ex.ToString());
 538                    _lookup.ice_getCommunicator().getLogger().warning(s.ToString());
 539                    _warnOnce = false;
 540                }
 541                _timer.cancel(request);
 542                _adapterRequests.Remove(request.getId());
 543            }
 544        }
 545    }
 546
 547    internal Ice.Internal.Timer timer() => _timer;
 548
 549    internal int latencyMultiplier() => _latencyMultiplier;
 550
 551    private readonly LocatorRegistryI _registry;
 552    private readonly LookupPrx _lookup;
 553    private readonly Dictionary<LookupPrx, LookupReplyPrx> _lookups = new Dictionary<LookupPrx, LookupReplyPrx>();
 554    private readonly TimeSpan _timeout;
 555    private readonly int _retryCount;
 556    private readonly int _latencyMultiplier;
 557    private readonly string _domainId;
 558
 559    private readonly Ice.Internal.Timer _timer;
 560    private bool _warnOnce = true;
 561    private readonly Dictionary<Ice.Identity, ObjectRequest> _objectRequests =
 562        new Dictionary<Ice.Identity, ObjectRequest>();
 563
 564    private readonly Dictionary<string, AdapterRequest> _adapterRequests = new Dictionary<string, AdapterRequest>();
 565    private readonly object _mutex = new();
 566}
 567
 568internal class LookupReplyI : LookupReplyDisp_
 569{
 570    public LookupReplyI(LookupI lookup) => _lookup = lookup;
 571
 572    public override void foundObjectById(Ice.Identity id, Ice.ObjectPrx proxy, Ice.Current c) =>
 573        _lookup.foundObject(id, c.id.name, proxy);
 574
 575    public override void foundAdapterById(string adapterId, Ice.ObjectPrx proxy, bool isReplicaGroup, Ice.Current c) =>
 576        _lookup.foundAdapter(adapterId, c.id.name, proxy, isReplicaGroup);
 577
 578    private readonly LookupI _lookup;
 579}