< Summary

Information
Class: Ice.Internal.OutgoingAsyncBase
Assembly: Ice
File(s): /_/csharp/src/Ice/Internal/OutgoingAsync.cs
Tag: 125_37167941578
Line coverage
85%
Covered lines: 127
Uncovered lines: 22
Coverable lines: 149
Total lines: 1588
Line coverage: 85.2%
Branch coverage
91%
Covered branches: 53
Total branches: 58
Branch coverage: 91.3%
Method coverage
84%
Covered methods: 22
Fully covered methods: 16
Total methods: 26
Method coverage: 84.6%
Full method coverage: 61.5%

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
sent()100%11100%
exception(...)100%11100%
response()100%210%
invokeSentAsync()100%1160%
invokeExceptionAsync()100%11100%
invokeResponseAsync()100%11100%
invokeSent()100%5466.67%
invokeException()100%2272.73%
invokeResponse()100%7672.22%
cancelable(...)100%22100%
cancel()100%11100%
attachRemoteObserver(...)100%44100%
attachCollocatedObserver(...)100%44100%
getOs()100%11100%
getIs()100%11100%
throwUserException()100%210%
cacheMessageBuffers()100%11100%
isSynchronous()100%11100%
.ctor(...)100%66100%
sentImpl(...)100%1010100%
exceptionImpl(...)87.5%9878.57%
responseImpl(...)100%66100%
cancel(...)100%22100%
warning(...)0%2040%
getObserver()100%11100%
sentSynchronously()100%210%

File(s)

/_/csharp/src/Ice/Internal/OutgoingAsync.cs

#LineLine coverage
 1// Copyright (c) ZeroC, Inc.
 2
 3using System.Diagnostics;
 4
 5namespace Ice.Internal;
 6
 7public interface OutgoingAsyncCompletionCallback
 8{
 9    void init(OutgoingAsyncBase og);
 10
 11    bool handleSent(bool done, bool alreadySent, OutgoingAsyncBase og);
 12
 13    bool handleException(Ice.Exception ex, OutgoingAsyncBase og);
 14
 15    bool handleResponse(bool userThread, bool ok, OutgoingAsyncBase og);
 16
 17    void handleInvokeSent(bool sentSynchronously, bool done, bool alreadySent, OutgoingAsyncBase og);
 18
 19    void handleInvokeException(Ice.Exception ex, OutgoingAsyncBase og);
 20
 21    void handleInvokeResponse(bool ok, OutgoingAsyncBase og);
 22}
 23
 24public abstract class OutgoingAsyncBase
 25{
 126    public virtual bool sent() => sentImpl(true);
 27
 128    public virtual bool exception(Ice.Exception ex) => exceptionImpl(ex);
 29
 30    public virtual bool response()
 31    {
 32        Debug.Assert(false); // Must be overridden by request that can handle responses
 033        return false;
 34    }
 35
 36    public void invokeSentAsync()
 37    {
 38        //
 39        // This is called when it's not safe to call the sent callback
 40        // synchronously from this thread. Instead the exception callback
 41        // is called asynchronously from the client thread pool.
 42        //
 43        try
 44        {
 145            instance_.clientThreadPool().execute(invokeSent, cachedConnection_);
 146        }
 047        catch (Ice.CommunicatorDestroyedException)
 48        {
 049        }
 150    }
 51
 52    public void invokeExceptionAsync() =>
 53        //
 54        // CommunicatorDestroyedCompleted is the only exception that can propagate directly
 55        // from this method.
 56        //
 157        instance_.clientThreadPool().execute(invokeException, cachedConnection_);
 58
 59    public void invokeResponseAsync() =>
 60        //
 61        // CommunicatorDestroyedCompleted is the only exception that can propagate directly
 62        // from this method.
 63        //
 164        instance_.clientThreadPool().execute(invokeResponse, cachedConnection_);
 65
 66    public void invokeSent()
 67    {
 68        try
 69        {
 170            _completionCallback.handleInvokeSent(sentSynchronously_, _doneInSent, _alreadySent, this);
 171        }
 072        catch (System.Exception ex)
 73        {
 074            warning(ex);
 075        }
 76
 177        if (observer_ != null && _doneInSent)
 78        {
 179            observer_.detach();
 180            observer_ = null;
 81        }
 182    }
 83
 84    public void invokeException()
 85    {
 86        try
 87        {
 88            try
 89            {
 190                throw _ex;
 91            }
 192            catch (Ice.Exception ex)
 93            {
 194                _completionCallback.handleInvokeException(ex, this);
 195            }
 196        }
 097        catch (System.Exception ex)
 98        {
 099            warning(ex);
 0100        }
 101
 1102        observer_?.detach();
 1103        observer_ = null;
 1104    }
 105
 106    public void invokeResponse()
 107    {
 1108        if (_ex != null)
 109        {
 1110            invokeException();
 1111            return;
 112        }
 113
 114        try
 115        {
 116            try
 117            {
 1118                _completionCallback.handleInvokeResponse((state_ & StateOK) != 0, this);
 1119            }
 1120            catch (Ice.Exception ex)
 121            {
 1122                if (_completionCallback.handleException(ex, this))
 123                {
 1124                    _completionCallback.handleInvokeException(ex, this);
 125                }
 1126            }
 0127            catch (System.AggregateException ex)
 128            {
 0129                throw ex.InnerException;
 130            }
 1131        }
 0132        catch (System.Exception ex)
 133        {
 0134            warning(ex);
 0135        }
 136
 1137        observer_?.detach();
 1138        observer_ = null;
 1139    }
 140
 141    public virtual void cancelable(CancellationHandler handler)
 142    {
 1143        lock (mutex_)
 144        {
 1145            if (_cancellationException != null)
 146            {
 147                try
 148                {
 1149                    throw _cancellationException;
 150                }
 1151                catch (Ice.LocalException)
 152                {
 1153                    _cancellationException = null;
 1154                    throw;
 155                }
 156            }
 1157            _cancellationHandler = handler;
 1158        }
 1159    }
 160
 1161    public void cancel() => cancel(new Ice.InvocationCanceledException());
 162
 163    public void attachRemoteObserver(Ice.ConnectionInfo info, Ice.Endpoint endpt, int requestId)
 164    {
 1165        Ice.Instrumentation.InvocationObserver observer = getObserver();
 1166        if (observer != null)
 167        {
 1168            int size = os_.size() - Protocol.headerSize - 4;
 1169            childObserver_ = observer.getRemoteObserver(info, endpt, requestId, size);
 1170            childObserver_?.attach();
 171        }
 1172    }
 173
 174    public void attachCollocatedObserver(Ice.ObjectAdapter adapter, int requestId)
 175    {
 1176        Ice.Instrumentation.InvocationObserver observer = getObserver();
 1177        if (observer != null)
 178        {
 1179            int size = os_.size() - Protocol.headerSize - 4;
 1180            childObserver_ = observer.getCollocatedObserver(adapter, requestId, size);
 1181            childObserver_?.attach();
 182        }
 1183    }
 184
 1185    public Ice.OutputStream getOs() => os_;
 186
 1187    public Ice.InputStream getIs() => is_;
 188
 189    public virtual void throwUserException()
 190    {
 0191    }
 192
 193    public virtual void cacheMessageBuffers()
 194    {
 1195    }
 196
 1197    public bool isSynchronous() => synchronous_;
 198
 1199    protected OutgoingAsyncBase(
 1200        Instance instance,
 1201        OutgoingAsyncCompletionCallback completionCallback,
 1202        Ice.OutputStream os = null,
 1203        Ice.InputStream iss = null)
 204    {
 1205        instance_ = instance;
 1206        sentSynchronously_ = false;
 1207        synchronous_ = false;
 1208        _doneInSent = false;
 1209        _alreadySent = false;
 1210        state_ = 0;
 1211        os_ = os ?? new OutputStream(Protocol.currentProtocolEncoding, instance.defaultsAndOverrides().defaultFormat);
 1212        is_ = iss ?? new Ice.InputStream(instance, Protocol.currentProtocolEncoding);
 1213        _completionCallback = completionCallback;
 1214        _completionCallback?.init(this);
 1215    }
 216
 217    protected virtual bool sentImpl(bool done)
 218    {
 1219        lock (mutex_)
 220        {
 1221            _alreadySent = (state_ & StateSent) > 0;
 1222            state_ |= StateSent;
 1223            if (done)
 224            {
 1225                _doneInSent = true;
 1226                childObserver_?.detach();
 1227                childObserver_ = null;
 1228                _cancellationHandler = null;
 229
 230                //
 231                // For oneway requests after the data has been sent
 232                // the buffers can be reused unless this is a
 233                // collocated invocation. For collocated invocations
 234                // the buffer won't be reused because it has already
 235                // been marked as cached in invokeCollocated.
 236                //
 1237                cacheMessageBuffers();
 238            }
 239
 1240            bool invoke = _completionCallback.handleSent(done, _alreadySent, this);
 1241            if (!invoke && _doneInSent && observer_ != null)
 242            {
 1243                observer_.detach();
 1244                observer_ = null;
 245            }
 1246            return invoke;
 247        }
 1248    }
 249
 250    protected virtual bool exceptionImpl(Ice.Exception ex)
 251    {
 1252        lock (mutex_)
 253        {
 1254            _ex = ex;
 1255            if (childObserver_ != null)
 256            {
 0257                childObserver_.failed(ex.ice_id());
 0258                childObserver_.detach();
 0259                childObserver_ = null;
 260            }
 1261            _cancellationHandler = null;
 262
 1263            observer_?.failed(ex.ice_id());
 1264            bool invoke = _completionCallback.handleException(ex, this);
 1265            if (!invoke && observer_ != null)
 266            {
 1267                observer_.detach();
 1268                observer_ = null;
 269            }
 1270            return invoke;
 271        }
 1272    }
 273
 274    protected virtual bool responseImpl(bool userThread, bool ok, bool invoke)
 275    {
 1276        lock (mutex_)
 277        {
 1278            if (ok)
 279            {
 1280                state_ |= StateOK;
 281            }
 282
 1283            _cancellationHandler = null;
 284
 285            try
 286            {
 1287                invoke &= _completionCallback.handleResponse(userThread, ok, this);
 1288            }
 1289            catch (Ice.Exception ex)
 290            {
 1291                _ex = ex;
 1292                invoke = _completionCallback.handleException(ex, this);
 1293            }
 1294            if (!invoke && observer_ != null)
 295            {
 1296                observer_.detach();
 1297                observer_ = null;
 298            }
 1299            return invoke;
 300        }
 1301    }
 302
 303    protected void cancel(Ice.LocalException ex)
 304    {
 305        CancellationHandler handler;
 306        {
 1307            lock (mutex_)
 308            {
 1309                if (_cancellationHandler == null)
 310                {
 1311                    _cancellationException = ex;
 1312                    return;
 313                }
 1314                handler = _cancellationHandler;
 1315            }
 316        }
 1317        handler.asyncRequestCanceled(this, ex);
 1318    }
 319
 320    private void warning(System.Exception ex)
 321    {
 0322        if (instance_.initializationData().properties.getIcePropertyAsInt("Ice.Warn.AMICallback") > 0)
 323        {
 0324            instance_.initializationData().logger.warning("exception raised by AMI callback:\n" + ex);
 325        }
 0326    }
 327
 328    //
 329    // This virtual method is necessary for the communicator flush
 330    // batch requests implementation.
 331    //
 1332    protected virtual Ice.Instrumentation.InvocationObserver getObserver() => observer_;
 333
 0334    public bool sentSynchronously() => sentSynchronously_;
 335
 336    protected Instance instance_;
 337    protected Ice.Connection cachedConnection_;
 338    protected bool sentSynchronously_;
 339    protected bool synchronous_;
 340    protected int state_;
 341
 342    protected Ice.Instrumentation.InvocationObserver observer_;
 343    protected Ice.Instrumentation.ChildInvocationObserver childObserver_;
 344
 345    protected Ice.OutputStream os_;
 346    protected Ice.InputStream is_;
 347
 1348    protected readonly object mutex_ = new();
 349
 350    private bool _doneInSent;
 351    private bool _alreadySent;
 352    private Ice.Exception _ex;
 353    private Ice.LocalException _cancellationException;
 354    private CancellationHandler _cancellationHandler;
 355    private readonly OutgoingAsyncCompletionCallback _completionCallback;
 356
 357    protected const int StateOK = 0x1;
 358    protected const int StateSent = 0x4;
 359    protected const int StateCachedBuffers = 0x10;
 360
 361    public const int AsyncStatusQueued = 0;
 362    public const int AsyncStatusSent = 1;
 363    public const int AsyncStatusInvokeSentCallback = 2;
 364}
 365
 366//
 367// Base class for proxy based invocations. This class handles the
 368// retry for proxy invocations. It also ensures the child observer is
 369// correctly notified of failures and makes sure the retry task is
 370// correctly canceled when the invocation completes.
 371//
 372public abstract class ProxyOutgoingAsyncBase : OutgoingAsyncBase, TimerTask
 373{
 374    public abstract int invokeRemote(Ice.ConnectionI connection, bool compress, bool response);
 375
 376    public abstract int invokeCollocated(CollocatedRequestHandler handler);
 377
 378    public override bool exception(Ice.Exception ex)
 379    {
 380        if (childObserver_ != null)
 381        {
 382            childObserver_.failed(ex.ice_id());
 383            childObserver_.detach();
 384            childObserver_ = null;
 385        }
 386
 387        cachedConnection_ = null;
 388
 389        //
 390        // NOTE: at this point, synchronization isn't needed, no other threads should be
 391        // calling on the callback.
 392        //
 393        try
 394        {
 395            //
 396            // It's important to let the retry queue do the retry even if
 397            // the retry interval is 0. This method can be called with the
 398            // connection locked so we can't just retry here.
 399            //
 400            instance_.retryQueue().add(this, handleRetryAfterException(ex));
 401            return false;
 402        }
 403        catch (Ice.Exception retryEx)
 404        {
 405            return exceptionImpl(retryEx); // No retries, we're done
 406        }
 407    }
 408
 409    public void retryException()
 410    {
 411        try
 412        {
 413            // It's important to let the retry queue do the retry. This is
 414            // called from the connect request handler and the retry might
 415            // require could end up waiting for the flush of the
 416            // connection to be done.
 417
 418            proxy_.iceGetRequestHandlerCache().clearCachedRequestHandler(handler_);
 419            instance_.retryQueue().add(this, 0);
 420        }
 421        catch (Ice.Exception ex)
 422        {
 423            if (exception(ex))
 424            {
 425                invokeExceptionAsync();
 426            }
 427        }
 428    }
 429
 430    public void retry() => invokeImpl(false);
 431
 432    public void abort(Ice.Exception ex)
 433    {
 434        Debug.Assert(childObserver_ == null);
 435        if (exceptionImpl(ex))
 436        {
 437            invokeExceptionAsync();
 438        }
 439        else if (ex is Ice.CommunicatorDestroyedException)
 440        {
 441            //
 442            // If it's a communicator destroyed exception, swallow
 443            // it but instead notify the user thread. Even if no callback
 444            // was provided.
 445            //
 446            throw ex;
 447        }
 448    }
 449
 450    protected ProxyOutgoingAsyncBase(
 451        Ice.ObjectPrxHelperBase prx,
 452        OutgoingAsyncCompletionCallback completionCallback,
 453        Ice.OutputStream os = null,
 454        Ice.InputStream iss = null)
 455        : base(prx.iceReference().getInstance(), completionCallback, os, iss)
 456    {
 457        proxy_ = prx;
 458        mode_ = Ice.OperationMode.Normal;
 459        _cnt = 0;
 460        _sent = false;
 461    }
 462
 463    protected void invokeImpl(bool userThread)
 464    {
 465        try
 466        {
 467            if (userThread)
 468            {
 469                TimeSpan invocationTimeout = proxy_.iceReference().getInvocationTimeout();
 470                if (invocationTimeout > TimeSpan.Zero)
 471                {
 472                    instance_.timer().schedule(this, invocationTimeout);
 473                }
 474            }
 475            else
 476            {
 477                observer_?.retried();
 478            }
 479
 480            while (true)
 481            {
 482                try
 483                {
 484                    _sent = false;
 485                    handler_ = proxy_.iceGetRequestHandlerCache().requestHandler;
 486                    int status = handler_.sendAsyncRequest(this);
 487                    if ((status & AsyncStatusSent) != 0)
 488                    {
 489                        if (userThread)
 490                        {
 491                            sentSynchronously_ = true;
 492                            if ((status & AsyncStatusInvokeSentCallback) != 0)
 493                            {
 494                                invokeSent(); // Call the sent callback from the user thread.
 495                            }
 496                        }
 497                        else
 498                        {
 499                            if ((status & AsyncStatusInvokeSentCallback) != 0)
 500                            {
 501                                invokeSentAsync(); // Call the sent callback from a client thread pool thread.
 502                            }
 503                        }
 504                    }
 505                    return; // We're done!
 506                }
 507                catch (RetryException)
 508                {
 509                    // Clear request handler and always retry.
 510                    proxy_.iceGetRequestHandlerCache().clearCachedRequestHandler(handler_);
 511                }
 512                catch (Ice.Exception ex)
 513                {
 514                    if (childObserver_ != null)
 515                    {
 516                        childObserver_.failed(ex.ice_id());
 517                        childObserver_.detach();
 518                        childObserver_ = null;
 519                    }
 520                    int interval = handleRetryAfterException(ex);
 521                    if (interval > 0)
 522                    {
 523                        instance_.retryQueue().add(this, interval);
 524                        return;
 525                    }
 526                    else
 527                    {
 528                        observer_?.retried();
 529                    }
 530                }
 531            }
 532        }
 533        catch (Ice.Exception ex)
 534        {
 535            // If called from the user thread we re-throw, the exception will caught by the caller and handled using
 536            // abort.
 537            if (userThread)
 538            {
 539                throw;
 540            }
 541            else if (exceptionImpl(ex)) // No retries, we're done
 542            {
 543                invokeExceptionAsync();
 544            }
 545        }
 546    }
 547
 548    protected override bool sentImpl(bool done)
 549    {
 550        _sent = true;
 551        if (done)
 552        {
 553            if (proxy_.iceReference().getInvocationTimeout() > TimeSpan.Zero)
 554            {
 555                instance_.timer().cancel(this);
 556            }
 557        }
 558        return base.sentImpl(done);
 559    }
 560
 561    protected override bool exceptionImpl(Ice.Exception ex)
 562    {
 563        if (proxy_.iceReference().getInvocationTimeout() > TimeSpan.Zero)
 564        {
 565            instance_.timer().cancel(this);
 566        }
 567        return base.exceptionImpl(ex);
 568    }
 569
 570    protected override bool responseImpl(bool userThread, bool ok, bool invoke)
 571    {
 572        if (proxy_.iceReference().getInvocationTimeout() > TimeSpan.Zero)
 573        {
 574            instance_.timer().cancel(this);
 575        }
 576        return base.responseImpl(userThread, ok, invoke);
 577    }
 578
 579    public void runTimerTask() => cancel(new Ice.InvocationTimeoutException());
 580
 581    private int handleRetryAfterException(Ice.Exception ex)
 582    {
 583        // Clear the request handler
 584        proxy_.iceGetRequestHandlerCache().clearCachedRequestHandler(handler_);
 585
 586        // We only retry local exception.
 587        //
 588        // A CloseConnectionException indicates graceful server shutdown, and is therefore
 589        // always repeatable without violating "at-most-once". That's because by sending a
 590        // close connection message, the server guarantees that all outstanding requests
 591        // can safely be repeated.
 592        //
 593        // An ObjectNotExistException can always be retried as well without violating
 594        // "at-most-once" (see the implementation of the checkRetryAfterException method
 595        // below for the reasons why it can be useful).
 596        //
 597        // If the request didn't get sent or if it's non-mutating or idempotent it can
 598        // also always be retried if the retry count isn't reached.
 599        bool shouldRetry = ex is LocalException && (!_sent ||
 600            mode_ is not OperationMode.Normal ||
 601            ex is CloseConnectionException ||
 602            ex is ObjectNotExistException);
 603
 604        if (shouldRetry)
 605        {
 606            try
 607            {
 608                return checkRetryAfterException((LocalException)ex);
 609            }
 610            catch (CommunicatorDestroyedException)
 611            {
 612                throw ex; // The communicator is already destroyed, so we cannot retry.
 613            }
 614        }
 615        else
 616        {
 617            throw ex; // Retry could break at-most-once semantics, don't retry.
 618        }
 619    }
 620
 621    private int checkRetryAfterException(Ice.LocalException ex)
 622    {
 623        Reference @ref = proxy_.iceReference();
 624        Instance instance = @ref.getInstance();
 625
 626        TraceLevels traceLevels = instance.traceLevels();
 627        Ice.Logger logger = instance.initializationData().logger!;
 628
 629        // We don't retry batch requests because the exception might have caused
 630        // the all the requests batched with the connection to be aborted and we
 631        // want the application to be notified.
 632        if (@ref.getMode() == Reference.Mode.ModeBatchOneway || @ref.getMode() == Reference.Mode.ModeBatchDatagram)
 633        {
 634            throw ex;
 635        }
 636
 637        // If it's a fixed proxy, retrying isn't useful as the proxy is tied to
 638        // the connection and the request will fail with the exception.
 639        if (@ref is FixedReference)
 640        {
 641            throw ex;
 642        }
 643
 644        var one = ex as Ice.ObjectNotExistException;
 645        if (one is not null)
 646        {
 647            if (@ref.getRouterInfo() != null && one.operation == "ice_add_proxy")
 648            {
 649                // If we have a router, an ObjectNotExistException with an
 650                // operation name "ice_add_proxy" indicates to the client
 651                // that the router isn't aware of the proxy (for example,
 652                // because it was evicted by the router). In this case, we
 653                // must *always* retry, so that the missing proxy is added
 654                // to the router.
 655
 656                @ref.getRouterInfo().clearCache(@ref);
 657
 658                if (traceLevels.retry >= 1)
 659                {
 660                    string s = "retrying operation call to add proxy to router\n" + ex;
 661                    logger.trace(traceLevels.retryCat, s);
 662                }
 663                return 0; // We must always retry, so we don't look at the retry count.
 664            }
 665            else if (@ref.isIndirect())
 666            {
 667                // We retry ObjectNotExistException if the reference is indirect.
 668
 669                if (@ref.isWellKnown())
 670                {
 671                    LocatorInfo li = @ref.getLocatorInfo();
 672                    li?.clearCache(@ref);
 673                }
 674            }
 675            else
 676            {
 677                // For all other cases, we don't retry ObjectNotExistException.
 678                throw ex;
 679            }
 680        }
 681        else if (ex is Ice.RequestFailedException)
 682        {
 683            throw ex;
 684        }
 685
 686        // There is no point in retrying an operation that resulted in a
 687        // MarshalException. This must have been raised locally (because if
 688        // it happened in a server it would result in an UnknownLocalException
 689        // instead), which means there was a problem in this process that will
 690        // not change if we try again.
 691        //
 692        // A likely cause for a MarshalException is exceeding the
 693        // maximum message size. For example, a client can attempt to send a
 694        // message that exceeds the maximum memory size, or accumulate enough
 695        // batch requests without flushing that the maximum size is reached.
 696        //
 697        // This latter case is especially problematic, because if we were to
 698        // retry a batch request after a MarshalException, we would in fact
 699        // silently discard the accumulated requests and allow new batch
 700        // requests to accumulate. If the subsequent batched requests do not
 701        // exceed the maximum message size, it appears to the client that all
 702        // of the batched requests were accepted, when in reality only the
 703        // last few are actually sent.
 704        if (ex is Ice.MarshalException)
 705        {
 706            throw ex;
 707        }
 708
 709        // Don't retry if the communicator is destroyed, object adapter is deactivated,
 710        // or connection is closed by the application.
 711        if (ex is CommunicatorDestroyedException ||
 712           ex is ObjectAdapterDeactivatedException ||
 713           ex is ObjectAdapterDestroyedException ||
 714           (ex is ConnectionAbortedException connectionAbortedException &&
 715            connectionAbortedException.closedByApplication) ||
 716           (ex is ConnectionClosedException connectionClosedException &&
 717            connectionClosedException.closedByApplication))
 718        {
 719            throw ex;
 720        }
 721
 722        // Don't retry invocation timeouts.
 723        if (ex is Ice.InvocationTimeoutException || ex is Ice.InvocationCanceledException)
 724        {
 725            throw ex;
 726        }
 727
 728        ++_cnt;
 729        Debug.Assert(_cnt > 0);
 730
 731        int[] retryIntervals = instance.retryIntervals;
 732
 733        int interval;
 734        if (_cnt == (retryIntervals.Length + 1) && ex is Ice.CloseConnectionException)
 735        {
 736            // A close connection exception is always retried at least once, even if the retry
 737            // limit is reached.
 738            interval = 0;
 739        }
 740        else if (_cnt > retryIntervals.Length)
 741        {
 742            if (traceLevels.retry >= 1)
 743            {
 744                string s = "cannot retry operation call because retry limit has been exceeded\n" + ex;
 745                logger.trace(traceLevels.retryCat, s);
 746            }
 747            throw ex;
 748        }
 749        else
 750        {
 751            interval = retryIntervals[_cnt - 1];
 752        }
 753
 754        if (traceLevels.retry >= 1)
 755        {
 756            string s = "retrying operation call";
 757            if (interval > 0)
 758            {
 759                s += " in " + interval + "ms";
 760            }
 761            s += " because of exception\n" + ex;
 762            logger.trace(traceLevels.retryCat, s);
 763        }
 764
 765        return interval;
 766    }
 767
 768    protected readonly Ice.ObjectPrxHelperBase proxy_;
 769    protected RequestHandler handler_;
 770    protected Ice.OperationMode mode_;
 771
 772    private int _cnt;
 773    private bool _sent;
 774}
 775
 776//
 777// Class for handling Slice operation invocations
 778//
 779public class OutgoingAsync : ProxyOutgoingAsyncBase
 780{
 781    public OutgoingAsync(
 782        Ice.ObjectPrxHelperBase prx,
 783        OutgoingAsyncCompletionCallback completionCallback,
 784        Ice.OutputStream os = null,
 785        Ice.InputStream iss = null)
 786        : base(prx, completionCallback, os, iss)
 787    {
 788        encoding_ = proxy_.iceReference().getEncoding();
 789        synchronous_ = false;
 790    }
 791
 792    public void prepare(string operation, Ice.OperationMode mode, Dictionary<string, string> context)
 793    {
 794        if (proxy_.iceReference().getProtocol() != Protocol.currentProtocol)
 795        {
 796            throw new FeatureNotSupportedException(
 797                $"Cannot send request using protocol version {proxy_.iceReference().getProtocol()}.");
 798        }
 799
 800        mode_ = mode;
 801
 802        observer_ = ObserverHelper.get(proxy_, operation, context);
 803
 804        if (proxy_.iceReference().isBatch)
 805        {
 806            proxy_.iceReference().batchRequestQueue.prepareBatchRequest(os_);
 807        }
 808        else
 809        {
 810            os_.writeBlob(Protocol.requestHdr);
 811        }
 812
 813        Reference rf = proxy_.iceReference();
 814
 815        Identity.ice_write(os_, rf.getIdentity());
 816
 817        //
 818        // For compatibility with the old FacetPath.
 819        //
 820        string facet = rf.getFacet();
 821        if (facet == null || facet.Length == 0)
 822        {
 823            os_.writeStringSeq(null);
 824        }
 825        else
 826        {
 827            string[] facetPath = { facet };
 828            os_.writeStringSeq(facetPath);
 829        }
 830
 831        os_.writeString(operation);
 832
 833        os_.writeByte((byte)mode);
 834
 835        if (context != null)
 836        {
 837            //
 838            // Explicit context
 839            //
 840            Ice.ContextHelper.write(os_, context);
 841        }
 842        else
 843        {
 844            //
 845            // Implicit context
 846            //
 847            Ice.ImplicitContextI implicitContext = rf.getInstance().getImplicitContext();
 848            Dictionary<string, string> prxContext = rf.getContext();
 849
 850            if (implicitContext == null)
 851            {
 852                Ice.ContextHelper.write(os_, prxContext);
 853            }
 854            else
 855            {
 856                implicitContext.write(prxContext, os_);
 857            }
 858        }
 859    }
 860
 861    public override bool sent() => sentImpl(!proxy_.ice_isTwoway()); // done = true if it's not a two-way proxy
 862
 863    public override bool response()
 864    {
 865        //
 866        // NOTE: this method is called from ConnectionI.parseMessage
 867        // with the connection locked. Therefore, it must not invoke
 868        // any user callbacks.
 869        //
 870        Debug.Assert(proxy_.ice_isTwoway()); // Can only be called for twoways.
 871
 872        if (childObserver_ != null)
 873        {
 874            childObserver_.reply(is_.size() - Protocol.headerSize - 4);
 875            childObserver_.detach();
 876            childObserver_ = null;
 877        }
 878
 879        try
 880        {
 881            // We can't (shouldn't) use the generated code to unmarshal a possibly unknown reply status.
 882            var replyStatus = (ReplyStatus)is_.readByte();
 883
 884            switch (replyStatus)
 885            {
 886                case ReplyStatus.Ok:
 887                    break;
 888
 889                case ReplyStatus.UserException:
 890                    observer_?.userException();
 891                    break;
 892
 893                case ReplyStatus.ObjectNotExist:
 894                case ReplyStatus.FacetNotExist:
 895                case ReplyStatus.OperationNotExist:
 896                {
 897                    var ident = new Ice.Identity(is_);
 898
 899                    //
 900                    // For compatibility with the old FacetPath.
 901                    //
 902                    string[] facetPath = is_.readStringSeq();
 903                    string facet;
 904                    if (facetPath.Length > 0)
 905                    {
 906                        if (facetPath.Length > 1)
 907                        {
 908                            throw new MarshalException(
 909                                $"Received invalid facet path with {facetPath.Length} elements.");
 910                        }
 911                        facet = facetPath[0];
 912                    }
 913                    else
 914                    {
 915                        facet = "";
 916                    }
 917
 918                    string operation = is_.readString();
 919                    throw replyStatus switch
 920                    {
 921                        ReplyStatus.ObjectNotExist => new ObjectNotExistException(ident, facet, operation),
 922                        ReplyStatus.FacetNotExist => new FacetNotExistException(ident, facet, operation),
 923                        _ => new OperationNotExistException(ident, facet, operation)
 924                    };
 925                }
 926
 927                default:
 928                {
 929                    string message = is_.readString();
 930                    throw replyStatus switch
 931                    {
 932                        ReplyStatus.UnknownException => new UnknownException(message),
 933                        ReplyStatus.UnknownLocalException => new UnknownLocalException(message),
 934                        ReplyStatus.UnknownUserException => new UnknownUserException(message),
 935                        _ => new DispatchException(replyStatus, message),
 936                    };
 937                }
 938            }
 939
 940            return responseImpl(false, replyStatus == ReplyStatus.Ok, true);
 941        }
 942        catch (Ice.Exception ex)
 943        {
 944            return exception(ex);
 945        }
 946    }
 947
 948    public override int invokeRemote(Ice.ConnectionI connection, bool compress, bool response)
 949    {
 950        cachedConnection_ = connection;
 951        return connection.sendAsyncRequest(this, compress, response, 0);
 952    }
 953
 954    public override int invokeCollocated(CollocatedRequestHandler handler)
 955    {
 956        // The stream cannot be cached if the proxy is not a twoway or there is an invocation timeout set.
 957        if (!proxy_.ice_isTwoway() || proxy_.iceReference().getInvocationTimeout() > TimeSpan.Zero)
 958        {
 959            // Disable caching by marking the streams as cached!
 960            state_ |= StateCachedBuffers;
 961        }
 962        return handler.invokeAsyncRequest(this, 0, synchronous_);
 963    }
 964
 965    public new void abort(Ice.Exception ex)
 966    {
 967        if (proxy_.iceReference().isBatch)
 968        {
 969            proxy_.iceReference().batchRequestQueue.abortBatchRequest(os_);
 970        }
 971
 972        base.abort(ex);
 973    }
 974
 975    protected void invoke(string operation, bool synchronous)
 976    {
 977        synchronous_ = synchronous;
 978        if (proxy_.iceReference().isBatch)
 979        {
 980            sentSynchronously_ = true;
 981            proxy_.iceReference().batchRequestQueue.finishBatchRequest(os_, proxy_, operation);
 982            responseImpl(true, true, false); // Don't call sent/completed callback for batch AMI requests
 983            return;
 984        }
 985
 986        // invokeImpl can throw
 987        invokeImpl(true); // userThread = true
 988    }
 989
 990    public void invoke(
 991        string operation,
 992        Ice.OperationMode mode,
 993        Ice.FormatType? format,
 994        Dictionary<string, string> context,
 995        bool synchronous,
 996        System.Action<Ice.OutputStream> write)
 997    {
 998        try
 999        {
 1000            prepare(operation, mode, context);
 1001            if (write != null)
 1002            {
 1003                os_.startEncapsulation(encoding_, format);
 1004                write(os_);
 1005                os_.endEncapsulation();
 1006            }
 1007            else
 1008            {
 1009                os_.writeEmptyEncapsulation(encoding_);
 1010            }
 1011            invoke(operation, synchronous);
 1012        }
 1013        catch (Ice.Exception ex)
 1014        {
 1015            abort(ex);
 1016        }
 1017    }
 1018
 1019    public override void throwUserException()
 1020    {
 1021        try
 1022        {
 1023            is_.startEncapsulation();
 1024            is_.throwException();
 1025        }
 1026        catch (UserException ex)
 1027        {
 1028            is_.endEncapsulation();
 1029            userException_?.Invoke(ex);
 1030            throw UnknownUserException.fromTypeId(ex.ice_id());
 1031        }
 1032    }
 1033
 1034    public override void cacheMessageBuffers()
 1035    {
 1036        if (proxy_.iceReference().getInstance().cacheMessageBuffers() > 0)
 1037        {
 1038            lock (mutex_)
 1039            {
 1040                if ((state_ & StateCachedBuffers) > 0)
 1041                {
 1042                    return;
 1043                }
 1044                state_ |= StateCachedBuffers;
 1045            }
 1046
 1047            is_?.reset();
 1048            os_.reset();
 1049
 1050            proxy_.cacheMessageBuffers(is_, os_);
 1051
 1052            is_ = null;
 1053            os_ = null;
 1054        }
 1055    }
 1056
 1057    protected readonly Ice.EncodingVersion encoding_;
 1058    protected System.Action<Ice.UserException> userException_;
 1059}
 1060
 1061public class OutgoingAsyncT<T> : OutgoingAsync
 1062{
 1063    public OutgoingAsyncT(
 1064        Ice.ObjectPrxHelperBase prx,
 1065        OutgoingAsyncCompletionCallback completionCallback,
 1066        Ice.OutputStream os = null,
 1067        Ice.InputStream iss = null)
 1068        : base(prx, completionCallback, os, iss)
 1069    {
 1070    }
 1071
 1072    public void invoke(
 1073        string operation,
 1074        Ice.OperationMode mode,
 1075        Ice.FormatType? format,
 1076        Dictionary<string, string> context,
 1077        bool synchronous,
 1078        System.Action<Ice.OutputStream> write = null,
 1079        System.Action<Ice.UserException> userException = null,
 1080        System.Func<Ice.InputStream, T> read = null)
 1081    {
 1082        read_ = read;
 1083        userException_ = userException;
 1084        base.invoke(operation, mode, format, context, synchronous, write);
 1085    }
 1086
 1087    public T getResult(bool ok)
 1088    {
 1089        try
 1090        {
 1091            if (ok)
 1092            {
 1093                if (read_ == null)
 1094                {
 1095                    if (is_ == null || is_.isEmpty())
 1096                    {
 1097                        //
 1098                        // If there's no response (oneway, batch-oneway proxies), we just set the result
 1099                        // on completion without reading anything from the input stream. This is required for
 1100                        // batch invocations.
 1101                        //
 1102                    }
 1103                    else
 1104                    {
 1105                        is_.skipEmptyEncapsulation();
 1106                    }
 1107                    return default;
 1108                }
 1109                else
 1110                {
 1111                    is_.startEncapsulation();
 1112                    T r = read_(is_);
 1113                    is_.endEncapsulation();
 1114                    return r;
 1115                }
 1116            }
 1117            else
 1118            {
 1119                throwUserException();
 1120                return default; // make compiler happy
 1121            }
 1122        }
 1123        finally
 1124        {
 1125            cacheMessageBuffers();
 1126        }
 1127    }
 1128
 1129    protected System.Func<Ice.InputStream, T> read_;
 1130}
 1131
 1132//
 1133// Class for handling the proxy's begin_ice_flushBatchRequest request.
 1134//
 1135internal class ProxyFlushBatchAsync : ProxyOutgoingAsyncBase
 1136{
 1137    public ProxyFlushBatchAsync(Ice.ObjectPrxHelperBase prx, OutgoingAsyncCompletionCallback completionCallback)
 1138        : base(prx, completionCallback)
 1139    {
 1140    }
 1141
 1142    public override int invokeRemote(Ice.ConnectionI connection, bool compress, bool response)
 1143    {
 1144        if (_batchRequestNum == 0)
 1145        {
 1146            if (sent())
 1147            {
 1148                return AsyncStatusSent | AsyncStatusInvokeSentCallback;
 1149            }
 1150            else
 1151            {
 1152                return AsyncStatusSent;
 1153            }
 1154        }
 1155        cachedConnection_ = connection;
 1156        return connection.sendAsyncRequest(this, compress, false, _batchRequestNum);
 1157    }
 1158
 1159    public override int invokeCollocated(CollocatedRequestHandler handler)
 1160    {
 1161        if (_batchRequestNum == 0)
 1162        {
 1163            if (sent())
 1164            {
 1165                return AsyncStatusSent | AsyncStatusInvokeSentCallback;
 1166            }
 1167            else
 1168            {
 1169                return AsyncStatusSent;
 1170            }
 1171        }
 1172        return handler.invokeAsyncRequest(this, _batchRequestNum, false);
 1173    }
 1174
 1175    public void invoke(string operation, bool synchronous)
 1176    {
 1177        if (proxy_.iceReference().getProtocol().major != Protocol.currentProtocol.major)
 1178        {
 1179            throw new FeatureNotSupportedException(
 1180                $"Cannot send request using protocol version {proxy_.iceReference().getProtocol()}.");
 1181        }
 1182        try
 1183        {
 1184            synchronous_ = synchronous;
 1185            observer_ = ObserverHelper.get(proxy_, operation, null);
 1186            // Not used for proxy flush batch requests.
 1187            _batchRequestNum = proxy_.iceReference().batchRequestQueue.swap(os_, out _);
 1188            invokeImpl(true); // userThread = true
 1189        }
 1190        catch (Ice.Exception ex)
 1191        {
 1192            abort(ex);
 1193        }
 1194    }
 1195
 1196    private int _batchRequestNum;
 1197}
 1198
 1199//
 1200// Class for handling the proxy's begin_ice_getConnection request.
 1201//
 1202internal class ProxyGetConnection : ProxyOutgoingAsyncBase
 1203{
 1204    public ProxyGetConnection(Ice.ObjectPrxHelperBase prx, OutgoingAsyncCompletionCallback completionCallback)
 1205        : base(prx, completionCallback)
 1206    {
 1207    }
 1208
 1209    public override int invokeRemote(Ice.ConnectionI connection, bool compress, bool response)
 1210    {
 1211        // A fixed proxy never reaches the invocation machinery: invoke returns its bound connection directly.
 1212        Debug.Assert(!proxy_.ice_isFixed());
 1213        try
 1214        {
 1215            connection.throwException();
 1216        }
 1217        catch (Ice.LocalException ex)
 1218        {
 1219            // The connection is closed: throw RetryException so that the caller clears the cached request
 1220            // handler and calls invokeRemote again with a new connection.
 1221            throw new RetryException(ex);
 1222        }
 1223        cachedConnection_ = connection;
 1224        if (responseImpl(false, true, true))
 1225        {
 1226            invokeResponseAsync();
 1227        }
 1228        return AsyncStatusSent;
 1229    }
 1230
 1231    public override int invokeCollocated(CollocatedRequestHandler handler)
 1232    {
 1233        if (responseImpl(false, true, true))
 1234        {
 1235            invokeResponseAsync();
 1236        }
 1237        return AsyncStatusSent;
 1238    }
 1239
 1240    public Ice.Connection getConnection() => cachedConnection_;
 1241
 1242    public void invoke(string operation, bool synchronous)
 1243    {
 1244        try
 1245        {
 1246            synchronous_ = synchronous;
 1247            observer_ = ObserverHelper.get(proxy_, operation, null);
 1248            invokeImpl(true); // userThread = true
 1249        }
 1250        catch (Ice.Exception ex)
 1251        {
 1252            abort(ex);
 1253        }
 1254    }
 1255}
 1256
 1257internal class ConnectionFlushBatchAsync : OutgoingAsyncBase
 1258{
 1259    public ConnectionFlushBatchAsync(
 1260        Ice.ConnectionI connection,
 1261        Instance instance,
 1262        OutgoingAsyncCompletionCallback completionCallback)
 1263        : base(instance, completionCallback) => _connection = connection;
 1264
 1265    public void invoke(string operation, Ice.CompressBatch compressBatch, bool synchronous)
 1266    {
 1267        synchronous_ = synchronous;
 1268        observer_ = ObserverHelper.get(instance_, operation);
 1269        try
 1270        {
 1271            int status;
 1272            int batchRequestNum = _connection.getBatchRequestQueue().swap(os_, out bool compress);
 1273            if (batchRequestNum == 0)
 1274            {
 1275                status = AsyncStatusSent;
 1276                if (sent())
 1277                {
 1278                    status |= AsyncStatusInvokeSentCallback;
 1279                }
 1280            }
 1281            else
 1282            {
 1283                bool comp;
 1284                if (compressBatch == Ice.CompressBatch.Yes)
 1285                {
 1286                    comp = true;
 1287                }
 1288                else if (compressBatch == Ice.CompressBatch.No)
 1289                {
 1290                    comp = false;
 1291                }
 1292                else
 1293                {
 1294                    comp = compress;
 1295                }
 1296                status = _connection.sendAsyncRequest(this, comp, false, batchRequestNum);
 1297            }
 1298
 1299            if ((status & AsyncStatusSent) != 0)
 1300            {
 1301                sentSynchronously_ = true;
 1302                if ((status & AsyncStatusInvokeSentCallback) != 0)
 1303                {
 1304                    invokeSent();
 1305                }
 1306            }
 1307        }
 1308        catch (RetryException ex)
 1309        {
 1310            try
 1311            {
 1312                throw ex.get();
 1313            }
 1314            catch (Ice.LocalException ee)
 1315            {
 1316                if (exception(ee))
 1317                {
 1318                    invokeExceptionAsync();
 1319                }
 1320            }
 1321        }
 1322        catch (Ice.Exception ex)
 1323        {
 1324            if (exception(ex))
 1325            {
 1326                invokeExceptionAsync();
 1327            }
 1328        }
 1329    }
 1330
 1331    private readonly Ice.ConnectionI _connection;
 1332}
 1333
 1334public class CommunicatorFlushBatchAsync : OutgoingAsyncBase
 1335{
 1336    private class FlushBatch : OutgoingAsyncBase
 1337    {
 1338        public FlushBatch(
 1339            CommunicatorFlushBatchAsync outAsync,
 1340            Instance instance,
 1341            Ice.Instrumentation.InvocationObserver observer)
 1342            : base(instance, null)
 1343        {
 1344            _outAsync = outAsync;
 1345            _observer = observer;
 1346        }
 1347
 1348        public override bool
 1349        sent()
 1350        {
 1351            childObserver_?.detach();
 1352            childObserver_ = null;
 1353            _outAsync.check(false);
 1354            return false;
 1355        }
 1356
 1357        public override bool
 1358        exception(Ice.Exception ex)
 1359        {
 1360            if (childObserver_ != null)
 1361            {
 1362                childObserver_.failed(ex.ice_id());
 1363                childObserver_.detach();
 1364                childObserver_ = null;
 1365            }
 1366            _outAsync.check(false);
 1367            return false;
 1368        }
 1369
 1370        protected override Ice.Instrumentation.InvocationObserver
 1371        getObserver() => _observer;
 1372
 1373        private readonly CommunicatorFlushBatchAsync _outAsync;
 1374        private readonly Ice.Instrumentation.InvocationObserver _observer;
 1375    }
 1376
 1377    public CommunicatorFlushBatchAsync(Instance instance, OutgoingAsyncCompletionCallback callback)
 1378        : base(instance, callback) =>
 1379        //
 1380        // _useCount is initialized to 1 to prevent premature callbacks.
 1381        // The caller must invoke ready() after all flush requests have
 1382        // been initiated.
 1383        //
 1384        _useCount = 1;
 1385
 1386    internal void flushConnection(Ice.ConnectionI con, Ice.CompressBatch compressBatch)
 1387    {
 1388        lock (mutex_)
 1389        {
 1390            ++_useCount;
 1391        }
 1392
 1393        try
 1394        {
 1395            var flushBatch = new FlushBatch(this, instance_, observer_);
 1396            int batchRequestNum = con.getBatchRequestQueue().swap(flushBatch.getOs(), out bool compress);
 1397            if (batchRequestNum == 0)
 1398            {
 1399                flushBatch.sent();
 1400            }
 1401            else
 1402            {
 1403                bool comp;
 1404                if (compressBatch == Ice.CompressBatch.Yes)
 1405                {
 1406                    comp = true;
 1407                }
 1408                else if (compressBatch == Ice.CompressBatch.No)
 1409                {
 1410                    comp = false;
 1411                }
 1412                else
 1413                {
 1414                    comp = compress;
 1415                }
 1416                con.sendAsyncRequest(flushBatch, comp, false, batchRequestNum);
 1417            }
 1418        }
 1419        catch (RetryException ex)
 1420        {
 1421            // If the connection is closed and sendAsyncRequest throws a RetryException, we rethrow the underlying
 1422            // exception to the caller.
 1423            check(false);
 1424            throw ex.get();
 1425        }
 1426        catch (Ice.LocalException)
 1427        {
 1428            check(false);
 1429            throw;
 1430        }
 1431    }
 1432
 1433    public void invoke(string operation, Ice.CompressBatch compressBatch, bool synchronous)
 1434    {
 1435        synchronous_ = synchronous;
 1436        observer_ = ObserverHelper.get(instance_, operation);
 1437        instance_.outgoingConnectionFactory().flushAsyncBatchRequests(compressBatch, this);
 1438        instance_.objectAdapterFactory().flushAsyncBatchRequests(compressBatch, this);
 1439        check(true);
 1440    }
 1441
 1442    public void check(bool userThread)
 1443    {
 1444        lock (mutex_)
 1445        {
 1446            Debug.Assert(_useCount > 0);
 1447            if (--_useCount > 0)
 1448            {
 1449                return;
 1450            }
 1451        }
 1452
 1453        if (sentImpl(true))
 1454        {
 1455            if (userThread)
 1456            {
 1457                sentSynchronously_ = true;
 1458                invokeSent();
 1459            }
 1460            else
 1461            {
 1462                invokeSentAsync();
 1463            }
 1464        }
 1465    }
 1466
 1467    private int _useCount;
 1468}
 1469
 1470public abstract class TaskCompletionCallback<T> : TaskCompletionSource<T>, OutgoingAsyncCompletionCallback
 1471{
 1472    protected TaskCompletionCallback(System.IProgress<bool> progress, CancellationToken cancellationToken)
 1473        : base(TaskCreationOptions.RunContinuationsAsynchronously)
 1474    {
 1475        progress_ = progress;
 1476        _cancellationToken = cancellationToken;
 1477    }
 1478
 1479    public void init(OutgoingAsyncBase og)
 1480    {
 1481        if (_cancellationToken.CanBeCanceled)
 1482        {
 1483            CancellationTokenRegistration registration = _cancellationToken.Register(og.cancel);
 1484
 1485            // Dispose of the registration when the task of this TaskCompletionSource completes; otherwise, a
 1486            // long-lived cancellation token would retain a reference to every invocation created with it. We
 1487            // continue with (and don't await) this task because awaiting it would mark its exception as observed,
 1488            // and a faulted invocation would no longer reach TaskScheduler.UnobservedTaskException.
 1489            _ = this.Task.ContinueWith(
 1490                static (_, state) => ((CancellationTokenRegistration)state!).Dispose(),
 1491                registration,
 1492                CancellationToken.None,
 1493                TaskContinuationOptions.None,
 1494                TaskScheduler.Default);
 1495        }
 1496    }
 1497
 1498    public bool handleSent(bool done, bool alreadySent, OutgoingAsyncBase og)
 1499    {
 1500        if (done && og.isSynchronous())
 1501        {
 1502            Debug.Assert(progress_ == null);
 1503            handleInvokeSent(false, done, alreadySent, og);
 1504            return false;
 1505        }
 1506        return done || (progress_ != null && !alreadySent); // Invoke the sent callback only if not already invoked.
 1507    }
 1508
 1509    public bool handleException(Ice.Exception ex, OutgoingAsyncBase og)
 1510    {
 1511        //
 1512        // If this is a synchronous call, we can notify the task from this thread to avoid
 1513        // the thread context switch. We know there aren't any continuations setup with the
 1514        // task.
 1515        //
 1516        if (og.isSynchronous())
 1517        {
 1518            handleInvokeException(ex, og);
 1519            return false;
 1520        }
 1521        else
 1522        {
 1523            return true;
 1524        }
 1525    }
 1526
 1527    public bool handleResponse(bool userThread, bool ok, OutgoingAsyncBase og)
 1528    {
 1529        //
 1530        // If called from the user thread (only the case for batch requests) or if this
 1531        // is a synchronous call, we can notify the task from this thread to avoid the
 1532        // thread context switch. We know there aren't any continuations setup with the
 1533        // task.
 1534        //
 1535        if (userThread || og.isSynchronous())
 1536        {
 1537            handleInvokeResponse(ok, og);
 1538            return false;
 1539        }
 1540        else
 1541        {
 1542            return true;
 1543        }
 1544    }
 1545
 1546    public virtual void handleInvokeSent(bool sentSynchronously, bool done, bool alreadySent, OutgoingAsyncBase og)
 1547    {
 1548        if (progress_ != null && !alreadySent)
 1549        {
 1550            progress_.Report(sentSynchronously);
 1551        }
 1552        if (done)
 1553        {
 1554            SetResult(default);
 1555        }
 1556    }
 1557
 1558    public void handleInvokeException(Ice.Exception ex, OutgoingAsyncBase og) => SetException(ex);
 1559
 1560    public abstract void handleInvokeResponse(bool ok, OutgoingAsyncBase og);
 1561
 1562    private readonly CancellationToken _cancellationToken;
 1563
 1564    protected readonly System.IProgress<bool> progress_;
 1565}
 1566
 1567public class OperationTaskCompletionCallback<T> : TaskCompletionCallback<T>
 1568{
 1569    public OperationTaskCompletionCallback(System.IProgress<bool> progress, CancellationToken cancellationToken)
 1570        : base(progress, cancellationToken)
 1571    {
 1572    }
 1573
 1574    public override void handleInvokeResponse(bool ok, OutgoingAsyncBase og) =>
 1575        SetResult(((OutgoingAsyncT<T>)og).getResult(ok));
 1576}
 1577
 1578public class FlushBatchTaskCompletionCallback : TaskCompletionCallback<object>
 1579{
 1580    public FlushBatchTaskCompletionCallback(
 1581        IProgress<bool> progress = null,
 1582        CancellationToken cancellationToken = default)
 1583        : base(progress, cancellationToken)
 1584    {
 1585    }
 1586
 1587    public override void handleInvokeResponse(bool ok, OutgoingAsyncBase og) => SetResult(null);
 1588}